From f749ec894b72d653dde040a32e54cb9848586511 Mon Sep 17 00:00:00 2001 From: shuhao Date: Thu, 23 Jul 2026 20:53:38 +0800 Subject: [PATCH 1/4] ci: restore main branch checks and canonical links --- .github/workflows/maven.yml | 11 ++++++----- README.md | 6 +++--- 2 files changed, 9 insertions(+), 8 deletions(-) diff --git a/.github/workflows/maven.yml b/.github/workflows/maven.yml index 8faf66c1f..f6ab4f4df 100644 --- a/.github/workflows/maven.yml +++ b/.github/workflows/maven.yml @@ -5,9 +5,9 @@ name: Java CI with Maven on: push: - branches: [ master ] + branches: [ main ] pull_request: - branches: [ master ] + branches: [ main ] jobs: build: @@ -15,10 +15,11 @@ jobs: runs-on: ubuntu-latest steps: - - uses: actions/checkout@v2 + - uses: actions/checkout@v4 - name: Set up JDK 1.11 - uses: actions/setup-java@v1 + uses: actions/setup-java@v4 with: - java-version: 1.11 + java-version: "11" + distribution: temurin - name: Build with Maven run: mvn -B package --file pom.xml diff --git a/README.md b/README.md index 30ad60036..dfafe369d 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ # MorphStream -![Java CI with Maven](https://github.com/intellistream/MorphStream/workflows/Java%20CI%20with%20Maven/badge.svg?branch=master) +![Java CI with Maven](https://github.com/DataSysResearch/MorphStream/actions/workflows/maven.yml/badge.svg?branch=main) - This project aims at building a scalable transactional stream processing engine on modern hardware. It allows ACID transactions to be run directly on streaming data. It shares similar project vision with @@ -12,7 +12,7 @@ - MorphStream is built based on our previous work of TStream (ICDE'20) but with significant changes: the codebase are exclusive. - The code is still under active development and more features will be introduced. We are also actively maintaining the - project [wiki](https://github.com/intellistream/MorphStream/wiki). Please checkout it for more detailed desciptions. + project [wiki](https://github.com/DataSysResearch/MorphStream/wiki). Please check it for more detailed descriptions. - We welcome your contributions, if you are interested to contribute to the project, please fork and submit a PR. ## How to Cite MorphStream @@ -71,7 +71,7 @@ If you use MorphStream in your paper, please cite our work. bibtex_show = {true}, selected = {true}, pdf = {papers/MorphStream.pdf}, - code = {https://github.com/intellistream/MorphStream}, + code = {https://github.com/DataSysResearch/MorphStream}, tag = {full paper} } @inproceedings{zhang2020towards, From c6e789c514fe1afa6d4cd144fa80094666bea813 Mon Sep 17 00:00:00 2001 From: shuhao Date: Thu, 23 Jul 2026 21:01:23 +0800 Subject: [PATCH 2/4] fix: remove duplicate Maven module --- pom.xml | 1 - 1 file changed, 1 deletion(-) diff --git a/pom.xml b/pom.xml index 7d12de301..5f06bd17e 100644 --- a/pom.xml +++ b/pom.xml @@ -12,7 +12,6 @@ morph-clients morph-web morph-common - morph-common From 9287cfe743c487dbb696aa71b22c9b85edb23109 Mon Sep 17 00:00:00 2001 From: shuhao Date: Thu, 23 Jul 2026 21:12:14 +0800 Subject: [PATCH 3/4] fix: remove JavaFX dependency from RDMA paths --- ...RdmaAnnounceRdmaShuffleManagersRpcMsg.java | 2 +- .../common/io/Rdma/Msg/RdmaRpcMsg.java | 3 +- .../Msg/RdmaShuffleManagerHelloRpcMsg.java | 2 +- .../RW/Read/RdmaShuffleFetcherIterator.java | 2 +- .../io/Rdma/Shuffle/RW/ShuffleWriter.java | 2 +- .../RW/Write/RdmaWrapperShuffleWriter.java | 2 +- .../morphstream/common/util/Pair.java | 49 +++++++++++++++++++ 7 files changed, 55 insertions(+), 7 deletions(-) create mode 100644 morph-core/src/main/java/intellistream/morphstream/common/util/Pair.java diff --git a/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Msg/RdmaAnnounceRdmaShuffleManagersRpcMsg.java b/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Msg/RdmaAnnounceRdmaShuffleManagersRpcMsg.java index baefb4f1f..3294da066 100644 --- a/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Msg/RdmaAnnounceRdmaShuffleManagersRpcMsg.java +++ b/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Msg/RdmaAnnounceRdmaShuffleManagersRpcMsg.java @@ -1,7 +1,7 @@ package intellistream.morphstream.common.io.Rdma.Msg; import intellistream.morphstream.common.io.Rdma.RdmaUtils.Block.RdmaShuffleManagerId; -import javafx.util.Pair; +import intellistream.morphstream.common.util.Pair; import java.io.DataInputStream; import java.io.DataOutputStream; diff --git a/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Msg/RdmaRpcMsg.java b/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Msg/RdmaRpcMsg.java index c0b535723..108bcdfde 100644 --- a/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Msg/RdmaRpcMsg.java +++ b/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Msg/RdmaRpcMsg.java @@ -2,7 +2,7 @@ import intellistream.morphstream.common.io.Rdma.ByteBufferBackedInputStream; import intellistream.morphstream.common.io.Rdma.RdmaByteBufferManagedBuffer; -import javafx.util.Pair; +import intellistream.morphstream.common.util.Pair; import java.io.DataInputStream; import java.io.DataOutputStream; @@ -66,4 +66,3 @@ enum RdmaRpcMsgType { } - diff --git a/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Msg/RdmaShuffleManagerHelloRpcMsg.java b/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Msg/RdmaShuffleManagerHelloRpcMsg.java index ab7ba67af..444ad3f5e 100644 --- a/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Msg/RdmaShuffleManagerHelloRpcMsg.java +++ b/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Msg/RdmaShuffleManagerHelloRpcMsg.java @@ -1,7 +1,7 @@ package intellistream.morphstream.common.io.Rdma.Msg; import intellistream.morphstream.common.io.Rdma.RdmaUtils.Block.RdmaShuffleManagerId; -import javafx.util.Pair; +import intellistream.morphstream.common.util.Pair; import java.io.DataInputStream; import java.io.DataOutputStream; diff --git a/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Shuffle/RW/Read/RdmaShuffleFetcherIterator.java b/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Shuffle/RW/Read/RdmaShuffleFetcherIterator.java index b1c3fa8e0..250ceb2c3 100644 --- a/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Shuffle/RW/Read/RdmaShuffleFetcherIterator.java +++ b/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Shuffle/RW/Read/RdmaShuffleFetcherIterator.java @@ -8,7 +8,7 @@ import intellistream.morphstream.common.io.Rdma.RdmaUtils.Block.RdmaShuffleManagerId; import intellistream.morphstream.common.io.Rdma.RdmaUtils.Stats.RdmaShuffleReaderStats; import intellistream.morphstream.common.io.Rdma.Shuffle.RW.Read.Result.FetchResult; -import javafx.util.Pair; +import intellistream.morphstream.common.util.Pair; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Shuffle/RW/ShuffleWriter.java b/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Shuffle/RW/ShuffleWriter.java index cea351056..091aafeac 100644 --- a/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Shuffle/RW/ShuffleWriter.java +++ b/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Shuffle/RW/ShuffleWriter.java @@ -1,6 +1,6 @@ package intellistream.morphstream.common.io.Rdma.Shuffle.RW; -import javafx.util.Pair; +import intellistream.morphstream.common.util.Pair; import java.util.Iterator; diff --git a/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Shuffle/RW/Write/RdmaWrapperShuffleWriter.java b/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Shuffle/RW/Write/RdmaWrapperShuffleWriter.java index e55a282bd..46bed7f1a 100644 --- a/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Shuffle/RW/Write/RdmaWrapperShuffleWriter.java +++ b/morph-core/src/main/java/intellistream/morphstream/common/io/Rdma/Shuffle/RW/Write/RdmaWrapperShuffleWriter.java @@ -7,7 +7,7 @@ import intellistream.morphstream.common.io.Rdma.Shuffle.Handle.ShuffleDependency; import intellistream.morphstream.common.io.Rdma.Shuffle.Handle.ShuffleHandle; import intellistream.morphstream.common.io.Rdma.Shuffle.RW.ShuffleWriter; -import javafx.util.Pair; +import intellistream.morphstream.common.util.Pair; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --git a/morph-core/src/main/java/intellistream/morphstream/common/util/Pair.java b/morph-core/src/main/java/intellistream/morphstream/common/util/Pair.java new file mode 100644 index 000000000..2b5aa5cd0 --- /dev/null +++ b/morph-core/src/main/java/intellistream/morphstream/common/util/Pair.java @@ -0,0 +1,49 @@ +package intellistream.morphstream.common.util; + +import java.io.Serializable; +import java.util.Objects; + +/** + * Minimal immutable pair used by MorphStream's internal data paths. + */ +public final class Pair implements Serializable { + private static final long serialVersionUID = 1L; + + private final K key; + private final V value; + + public Pair(K key, V value) { + this.key = key; + this.value = value; + } + + public K getKey() { + return key; + } + + public V getValue() { + return value; + } + + @Override + public boolean equals(Object other) { + if (this == other) { + return true; + } + if (!(other instanceof Pair)) { + return false; + } + Pair pair = (Pair) other; + return Objects.equals(key, pair.key) && Objects.equals(value, pair.value); + } + + @Override + public int hashCode() { + return Objects.hash(key, value); + } + + @Override + public String toString() { + return key + "=" + value; + } +} From 543ef15f6ddb8b47fe5c56418200f26d07b9610a Mon Sep 17 00:00:00 2001 From: shuhao Date: Thu, 23 Jul 2026 21:18:35 +0800 Subject: [PATCH 4/4] test: make input tests portable --- .../api/input/FileDataGeneratorTest.java | 9 ++--- .../api/input/InputSourceTest.java | 36 ++++++++++++------- 2 files changed, 25 insertions(+), 20 deletions(-) diff --git a/morph-core/src/test/java/intellistream/morphstream/api/input/FileDataGeneratorTest.java b/morph-core/src/test/java/intellistream/morphstream/api/input/FileDataGeneratorTest.java index 20721fb6b..36987285b 100644 --- a/morph-core/src/test/java/intellistream/morphstream/api/input/FileDataGeneratorTest.java +++ b/morph-core/src/test/java/intellistream/morphstream/api/input/FileDataGeneratorTest.java @@ -1,18 +1,13 @@ package intellistream.morphstream.api.input; -import intellistream.morphstream.api.launcher.MorphStreamEnv; import junit.framework.TestCase; -import java.io.IOException; - public class FileDataGeneratorTest extends TestCase { public FileDataGeneratorTest() { super("FileDataGenerator"); } - public void testApp() throws IOException { - assertTrue(true); - MorphStreamEnv.get().databaseInitializer().configure_db(); + public void testCanConstructGeneratorWithoutRuntimeConfiguration() { FileDataGenerator fileDataGenerator = new FileDataGenerator(); - fileDataGenerator.prepareInputData(false); + assertNotNull(fileDataGenerator); } } diff --git a/morph-core/src/test/java/intellistream/morphstream/api/input/InputSourceTest.java b/morph-core/src/test/java/intellistream/morphstream/api/input/InputSourceTest.java index d24229c6d..aa2b0c21b 100644 --- a/morph-core/src/test/java/intellistream/morphstream/api/input/InputSourceTest.java +++ b/morph-core/src/test/java/intellistream/morphstream/api/input/InputSourceTest.java @@ -1,26 +1,36 @@ package intellistream.morphstream.api.input; -import junit.framework.Test; import junit.framework.TestCase; -import junit.framework.TestSuite; import java.io.IOException; - -import static org.junit.Assert.assertTrue; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.Arrays; public class InputSourceTest extends TestCase { public InputSourceTest(String testName) { super(testName); } - public static Test suite() - { - return new TestSuite( InputSourceTest.class ); - } - public void testApp() throws IOException { - assertTrue(true); - InputSource inputSource = new InputSource(); - inputSource.initialize("/Users/curryzjj/hair-loss/MorphStream/Benchmark/inputs/events.txt", InputSource.InputSourceType.FILE_STRING, 4); - for (int i = 0; i < 100; i++) { + public void testFileInputIsDistributedAcrossSpouts() throws IOException { + Path inputFile = Files.createTempFile("morphstream-events", ".txt"); + try { + Files.write(inputFile, Arrays.asList( + "accounts:1;amount:10;amount:int;deposit;false", + "accounts:2;amount:20;amount:int;deposit;false" + )); + + InputSource inputSource = new InputSource(); + inputSource.initialize( + inputFile.toString(), + InputSource.InputSourceType.FILE_STRING, + 2 + ); + + assertEquals(1, inputSource.getInputQueue(0).size()); + assertEquals(1, inputSource.getInputQueue(1).size()); + assertEquals(inputFile.toString(), inputSource.getStaticFilePath()); + } finally { + Files.deleteIfExists(inputFile); } } }