diff --git a/.github/workflows/jacoco-badge.yml b/.github/workflows/jacoco-badge.yml
new file mode 100644
index 0000000..c61e164
--- /dev/null
+++ b/.github/workflows/jacoco-badge.yml
@@ -0,0 +1,39 @@
+name: Update JaCoCo Badge
+
+on:
+ push:
+ branches: [ "main" ]
+
+jobs:
+ badge:
+ runs-on: ubuntu-latest
+
+ steps:
+ - uses: actions/checkout@v4
+ - uses: actions/setup-java@v4
+ with:
+ java-version: '21'
+ distribution: 'temurin'
+ cache: maven
+
+ - name: Test with JUnit
+ run: mvn -B test
+
+ - name: Generate JaCoCo Badge
+ id: jacoco
+ uses: cicirello/jacoco-badge-generator@v2
+ with:
+ generate-branches-badge: true
+ jacoco-csv-file: >
+ meteor-jedis/target/site/jacoco/jacoco.csv
+ meteor-core/target/site/jacoco/jacoco.csv
+
+ - name: Log coverage percentage
+ run: |
+ echo "coverage = ${{ steps.jacoco.outputs.coverage }}"
+ echo "branch coverage = ${{ steps.jacoco.outputs.branches }}"
+
+ - uses: EndBug/add-and-commit@v9
+ with:
+ default_author: github_actions
+ message: "Update JaCoCo badge"
diff --git a/.github/workflows/maven.yml b/.github/workflows/maven.yml
index 6376ff8..d132b1f 100644
--- a/.github/workflows/maven.yml
+++ b/.github/workflows/maven.yml
@@ -1,54 +1,19 @@
-# This workflow will build a Java project with Maven, and cache/restore any dependencies to improve the workflow execution time
-# For more information see: https://docs.github.com/en/actions/automating-builds-and-tests/building-and-testing-java-with-maven
-
-# This workflow uses actions that are not certified by GitHub.
-# They are provided by a third-party and are governed by
-# separate terms of service, privacy policy, and support
-# documentation.
-
name: Java CI with Maven
on:
- push:
- branches: [ "main" ]
pull_request:
branches: [ "main" ]
jobs:
build:
-
runs-on: ubuntu-latest
-
steps:
- - uses: actions/checkout@v3
- with:
- ref: ${{ github.event.pull_request.head.ref }}
- - name: Set up JDK 20
- uses: actions/setup-java@v3
+ - uses: actions/checkout@v4
+ - uses: actions/setup-java@v4
with:
- java-version: '17.0.8+7'
+ java-version: '21'
distribution: 'temurin'
cache: maven
- # run tests with junit
- name: Test with JUnit
- run: mvn -B test --file pom.xml
-
- - name: Generate JaCoCo Badge
- id: jacoco
- uses: cicirello/jacoco-badge-generator@v2
- with:
- generate-branches-badge: true
- jacoco-csv-file: >
- meteor-jedis/target/site/jacoco/jacoco.csv
- meteor-core/target/site/jacoco/jacoco.csv
-
- - name: Log coverage percentage
- run: |
- echo "coverage = ${{ steps.jacoco.outputs.coverage }}"
- echo "branch coverage = ${{ steps.jacoco.outputs.branches }}"
-
- - uses: EndBug/add-and-commit@v9 # You can change this to use a specific version.
- with:
- default_author: github_actions
- message: "Add JaCoCo badge"
\ No newline at end of file
+ run: mvn -B test
diff --git a/.github/workflows/publish.yml b/.github/workflows/publish.yml
index ad26b36..36ceb9b 100644
--- a/.github/workflows/publish.yml
+++ b/.github/workflows/publish.yml
@@ -1,26 +1,56 @@
-name: Publish package to the Maven Central Repository
+name: Release to Maven Central
+
on:
- release:
- types: [created,edited]
+ push:
+ tags:
+ - '[0-9]*.[0-9]*.[0-9]*'
jobs:
+ create-release:
+ runs-on: ubuntu-latest
+ permissions:
+ contents: write
+ outputs:
+ tag: ${{ steps.tag.outputs.tag }}
+ steps:
+ - uses: actions/checkout@v4
+ - id: tag
+ run: echo "tag=${GITHUB_REF_NAME}" >> "$GITHUB_OUTPUT"
+ - uses: softprops/action-gh-release@v2
+ with:
+ tag_name: ${{ github.ref_name }}
+ name: Release ${{ github.ref_name }}
+ generate_release_notes: true
+
publish:
+ needs: create-release
runs-on: ubuntu-latest
+ permissions:
+ contents: write
steps:
- - uses: actions/checkout@v3
- - name: Set up Maven Central Repository
- uses: actions/setup-java@v3
+ - uses: actions/checkout@v4
+ - uses: actions/setup-java@v4
with:
- java-version: '17.0.8+7'
+ java-version: '21'
distribution: 'temurin'
- server-id: ossrh
+ server-id: central
server-username: MAVEN_USERNAME
server-password: MAVEN_PASSWORD
gpg-private-key: ${{ secrets.MAVEN_GPG_PRIVATE_KEY }}
gpg-passphrase: MAVEN_GPG_PASSPHRASE
- - name: Publish package
- run: mvn -Drevision=${{ github.event.release.tag_name }} --batch-mode deploy -P release
+ - run: mvn -Drevision=${{ needs.create-release.outputs.tag }} --batch-mode deploy -P release
env:
- MAVEN_USERNAME: ${{ secrets.OSSRH_USERNAME }}
- MAVEN_PASSWORD: ${{ secrets.OSSRH_TOKEN }}
- MAVEN_GPG_PASSPHRASE: ${{ secrets.MAVEN_GPG_PASSPHRASE }}
\ No newline at end of file
+ MAVEN_USERNAME: ${{ secrets.SONATYPE_USERNAME }}
+ MAVEN_PASSWORD: ${{ secrets.SONATYPE_PASSWORD }}
+ MAVEN_GPG_PASSPHRASE: ${{ secrets.MAVEN_GPG_PASSPHRASE }}
+ - uses: actions/checkout@v4
+ with:
+ ref: main
+ - name: Update README version
+ run: |
+ TAG=${{ needs.create-release.outputs.tag }}
+ sed -i "s|[^<]*|${TAG}|" README.md
+ - uses: EndBug/add-and-commit@v9
+ with:
+ default_author: github_actions
+ message: "chore: update README version to ${{ needs.create-release.outputs.tag }}"
diff --git a/.gitignore b/.gitignore
index d323035..16acae6 100644
--- a/.gitignore
+++ b/.gitignore
@@ -15,4 +15,9 @@ target/
### Mac OS ###
.DS_Store
**/.flattened-pom.xml
-**/dependency-reduced-pom.xml
\ No newline at end of file
+**/dependency-reduced-pom.xml
+
+.settings/
+**/.project
+**/.factorypath
+**/.classpath
\ No newline at end of file
diff --git a/benchmarks/pom.xml b/benchmarks/pom.xml
index e983fb8..3650657 100644
--- a/benchmarks/pom.xml
+++ b/benchmarks/pom.xml
@@ -12,10 +12,6 @@
benchmarks
- 17
- 17
- UTF-8
-
1.37
@@ -23,17 +19,17 @@
dev.pixelib.meteor
meteor-core
- ${revision}
+ ${project.version}
dev.pixelib.meteor
meteor-common
- ${revision}
+ ${project.version}
- dev.pixelib.meteor.transport
+ dev.pixelib.meteor
meteor-jedis
- ${revision}
+ ${project.version}
@@ -58,8 +54,6 @@
org.apache.maven.plugins
maven-compiler-plugin
- 17
- 17
org.openjdk.jmh
@@ -73,6 +67,7 @@
org.apache.maven.plugins
maven-shade-plugin
+ 3.6.0
package
diff --git a/benchmarks/src/main/java/dev/pixelib/meteor/benchmarks/ScoresWithMaps.java b/benchmarks/src/main/java/dev/pixelib/meteor/benchmarks/ScoresWithMaps.java
index 0de342b..8bcacbe 100644
--- a/benchmarks/src/main/java/dev/pixelib/meteor/benchmarks/ScoresWithMaps.java
+++ b/benchmarks/src/main/java/dev/pixelib/meteor/benchmarks/ScoresWithMaps.java
@@ -22,7 +22,6 @@ public class ScoresWithMaps {
@Setup
public void setup() {
- System.out.println("workerThreads: " + workerThreads);
RpcOptions rpcOptions = new RpcOptions();
rpcOptions.setExecutorThreads(workerThreads);
meteor = new Meteor(new LoopbackTransport(), rpcOptions);
diff --git a/benchmarks/src/main/java/dev/pixelib/meteor/benchmarks/SimpleIncrement.java b/benchmarks/src/main/java/dev/pixelib/meteor/benchmarks/SimpleIncrement.java
index 013d5a0..111ebed 100644
--- a/benchmarks/src/main/java/dev/pixelib/meteor/benchmarks/SimpleIncrement.java
+++ b/benchmarks/src/main/java/dev/pixelib/meteor/benchmarks/SimpleIncrement.java
@@ -21,7 +21,6 @@ public class SimpleIncrement {
@Setup
public void setup() {
- System.out.println("workerThreads: " + workerThreads);
RpcOptions rpcOptions = new RpcOptions();
rpcOptions.setExecutorThreads(workerThreads);
meteor = new Meteor(new LoopbackTransport(), rpcOptions);
diff --git a/examples/pom.xml b/examples/pom.xml
index a973764..e284c11 100644
--- a/examples/pom.xml
+++ b/examples/pom.xml
@@ -11,12 +11,6 @@
examples
-
- 17
- 17
- UTF-8
-
-
dev.pixelib.meteor
@@ -24,7 +18,7 @@
${revision}
- dev.pixelib.meteor.transport
+ dev.pixelib.meteor
meteor-jedis
${revision}
diff --git a/examples/src/main/java/dev/pixelib/meteor/sender/ScoreboardExample.java b/examples/src/main/java/dev/pixelib/meteor/sender/ScoreboardExample.java
index e757188..1cc7fe1 100644
--- a/examples/src/main/java/dev/pixelib/meteor/sender/ScoreboardExample.java
+++ b/examples/src/main/java/dev/pixelib/meteor/sender/ScoreboardExample.java
@@ -2,10 +2,13 @@
import dev.pixelib.meteor.base.defaults.LoopbackTransport;
import dev.pixelib.meteor.core.Meteor;
+import lombok.extern.java.Log;
import java.util.HashMap;
import java.util.Map;
+import java.util.logging.Level;
+@Log
public class ScoreboardExample {
public static void main(String[] args) throws Exception{
@@ -19,10 +22,10 @@ public static void main(String[] args) throws Exception{
scoreboard.setScoreForPlayer("player3", 30);
Map scores = scoreboard.getAllScores();
- System.out.println("scores: " + scores);
+ log.log(Level.INFO, "scores: {0}", scores);
int player1Score = scoreboard.getScoreForPlayer("player1");
- System.out.println("player1 score: " + player1Score);
+ log.log(Level.INFO, "player1 score: {0}", player1Score);
meteor.stop();
}
diff --git a/examples/src/main/java/dev/pixelib/meteor/sender/SendingUpdateMethod.java b/examples/src/main/java/dev/pixelib/meteor/sender/SendingUpdateMethod.java
index b293e1f..b2dfb50 100644
--- a/examples/src/main/java/dev/pixelib/meteor/sender/SendingUpdateMethod.java
+++ b/examples/src/main/java/dev/pixelib/meteor/sender/SendingUpdateMethod.java
@@ -2,7 +2,11 @@
import dev.pixelib.meteor.base.defaults.LoopbackTransport;
import dev.pixelib.meteor.core.Meteor;
+import lombok.extern.java.Log;
+import java.util.logging.Level;
+
+@Log
public class SendingUpdateMethod {
public static void main(String[] args) throws Exception{
@@ -17,13 +21,13 @@ public static void main(String[] args) throws Exception{
meteor.registerImplementation(new MathFunctionsImpl());
int subResult = mathSubstract.substract(10, 1, 2, 3, 4, 5);
- System.out.println("10 - 1 - 2 - 3 - 4 - 5 = " + subResult);
+ log.log(Level.INFO, "10 - 1 - 2 - 3 - 4 - 5 = {0}", subResult);
int addResult = mathAdd.add(1, 2, 3, 4, 5);
- System.out.println("1 + 2 + 3 + 4 + 5 = " + addResult);
+ log.log(Level.INFO, "1 + 2 + 3 + 4 + 5 = {0}", addResult);
int multiResult = mathMultiply.multiply(5, 5);
- System.out.println("5 * 5 = " + multiResult);
+ log.log(Level.INFO, "5 * 5 = {0}", multiResult);
meteor.stop();
}
diff --git a/examples/src/main/java/dev/pixelib/meteor/sender/SendingUpdateMethodRedis.java b/examples/src/main/java/dev/pixelib/meteor/sender/SendingUpdateMethodRedis.java
index 17e7902..e385b24 100644
--- a/examples/src/main/java/dev/pixelib/meteor/sender/SendingUpdateMethodRedis.java
+++ b/examples/src/main/java/dev/pixelib/meteor/sender/SendingUpdateMethodRedis.java
@@ -2,11 +2,15 @@
import dev.pixelib.meteor.core.Meteor;
import dev.pixelib.meteor.transport.redis.RedisTransport;
+import lombok.extern.java.Log;
+import java.util.logging.Level;
+
+@Log
public class SendingUpdateMethodRedis {
public static void main(String[] args) throws Exception{
- Meteor meteor = new Meteor(new RedisTransport("192.168.178.46", 6379, "test"));
+ Meteor meteor = new Meteor(new RedisTransport("localhost", 6379, "test"));
MathAdd mathAdd = meteor.registerProcedure(MathAdd.class);
MathSubstract mathSubstract = meteor.registerProcedure(MathSubstract.class);
@@ -18,13 +22,13 @@ public static void main(String[] args) throws Exception{
int subResult = mathSubstract.substract(10, 1, 2, 3, 4, 5);
- System.out.println("10 - 1 - 2 - 3 - 4 - 5 = " + subResult);
+ log.log(Level.INFO, "10 - 1 - 2 - 3 - 4 - 5 = {0}", subResult);
int addResult = mathAdd.add(1, 2, 3, 4, 5);
- System.out.println("1 + 2 + 3 + 4 + 5 = " + addResult);
+ log.log(Level.INFO, "1 + 2 + 3 + 4 + 5 = {0}", addResult);
int multiResult = mathMultiply.multiply(5, 5);
- System.out.println("5 * 5 = " + multiResult);
+ log.log(Level.INFO, "5 * 5 = {0}", multiResult);
meteor.stop();
}
diff --git a/meteor-common/pom.xml b/meteor-common/pom.xml
index 8b8dfa7..2c2213d 100644
--- a/meteor-common/pom.xml
+++ b/meteor-common/pom.xml
@@ -4,7 +4,6 @@
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
4.0.0
-
dev.pixelib.meteor
meteor-parent
@@ -14,9 +13,6 @@
meteor-common
-
-
-
com.google.code.gson
diff --git a/meteor-common/src/main/java/dev/pixelib/meteor/base/defaults/GsonSerializer.java b/meteor-common/src/main/java/dev/pixelib/meteor/base/defaults/GsonSerializer.java
index 66fbaf3..8e49953 100644
--- a/meteor-common/src/main/java/dev/pixelib/meteor/base/defaults/GsonSerializer.java
+++ b/meteor-common/src/main/java/dev/pixelib/meteor/base/defaults/GsonSerializer.java
@@ -9,7 +9,7 @@ public class GsonSerializer implements RpcSerializer {
* the Gson instance to use for serialization/deserialization
* this is a static field so that it is shared between all instances of this class and can be swapped out at runtime
*/
- public static Gson GSON = new Gson();
+ public static final Gson GSON = new Gson();
/**
* @param obj the object to serialize
diff --git a/meteor-common/src/main/java/dev/pixelib/meteor/base/defaults/LoopbackTransport.java b/meteor-common/src/main/java/dev/pixelib/meteor/base/defaults/LoopbackTransport.java
index ec6ffb7..6106481 100644
--- a/meteor-common/src/main/java/dev/pixelib/meteor/base/defaults/LoopbackTransport.java
+++ b/meteor-common/src/main/java/dev/pixelib/meteor/base/defaults/LoopbackTransport.java
@@ -9,9 +9,13 @@
import java.util.EnumMap;
import java.util.List;
import java.util.Map;
+import java.util.logging.Level;
+import java.util.logging.Logger;
public class LoopbackTransport implements RpcTransport {
+ private final Logger logger = Logger.getLogger(LoopbackTransport.class.getSimpleName());
+
private final Map> onReceiveFunctions = new EnumMap<>(Direction.class);
/**
@@ -27,8 +31,7 @@ public void send(Direction direction, byte[] bytes) {
boolean matched = onReceiveFunction.onPacket(bytes);
if (matched) break;
} catch (Exception e) {
- // TODO: Add Logger
- e.printStackTrace();
+ logger.log(Level.SEVERE, "Error occurred while processing packet", e);
}
}
}
diff --git a/meteor-core/pom.xml b/meteor-core/pom.xml
index 40d244f..f08a321 100644
--- a/meteor-core/pom.xml
+++ b/meteor-core/pom.xml
@@ -14,11 +14,6 @@
meteor-core
-
- io.netty
- netty-buffer
-
-
dev.pixelib.meteor
meteor-common
@@ -26,6 +21,12 @@
compile
+
+ io.netty
+ netty-buffer
+ 4.2.15.Final
+
+
org.junit.jupiter
junit-jupiter
diff --git a/meteor-core/src/main/java/dev/pixelib/meteor/core/proxy/PendingInvocation.java b/meteor-core/src/main/java/dev/pixelib/meteor/core/proxy/PendingInvocation.java
index 2c64ec9..089dbee 100644
--- a/meteor-core/src/main/java/dev/pixelib/meteor/core/proxy/PendingInvocation.java
+++ b/meteor-core/src/main/java/dev/pixelib/meteor/core/proxy/PendingInvocation.java
@@ -54,16 +54,16 @@ public void complete(Object response) throws IllegalStateException {
boolean isVoidOrNullable = invocationDescriptor.getReturnType().equals(Void.TYPE) || !invocationDescriptor.getReturnType().isPrimitive();
// check instance of response
- if (!isVoidOrNullable && !invocationDescriptor.getReturnType().isInstance(response)) {
- // is the normal return type primitive? then check if its still assignable as a boxed
- if (invocationDescriptor.getReturnType().isPrimitive()) {
- if (!ArgumentMapper.ensureBoxedClass(invocationDescriptor.getReturnType()).isAssignableFrom(response.getClass())) {
+ // is the normal return type primitive? then check if its still assignable as a boxed
+ if (!isVoidOrNullable && !invocationDescriptor.getReturnType().isInstance(response)
+ && invocationDescriptor.getReturnType().isPrimitive()
+ && !ArgumentMapper.ensureBoxedClass(invocationDescriptor.getReturnType()).isAssignableFrom(response.getClass())
+ ) {
throw new IllegalStateException("Response is not an instance of the expected return type. " +
"Expected: " + invocationDescriptor.getReturnType().getName() + ", " +
"Actual: " + response.getClass().getName());
}
- }
- }
+
isComplete.set(true);
this.completable.complete((T) response);
diff --git a/meteor-core/src/main/java/dev/pixelib/meteor/core/trackers/IncomingInvocationTracker.java b/meteor-core/src/main/java/dev/pixelib/meteor/core/trackers/IncomingInvocationTracker.java
index 54adc61..255795d 100644
--- a/meteor-core/src/main/java/dev/pixelib/meteor/core/trackers/IncomingInvocationTracker.java
+++ b/meteor-core/src/main/java/dev/pixelib/meteor/core/trackers/IncomingInvocationTracker.java
@@ -3,6 +3,7 @@
import dev.pixelib.meteor.core.executor.ImplementationWrapper;
import java.util.Collection;
+import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
public class IncomingInvocationTracker {
@@ -28,7 +29,7 @@ public void registerImplementation(Object implementation, String namespace) {
}
}
- public ConcurrentHashMap, Collection> getImplementations() {
+ public Map, Collection> getImplementations() {
return implementations;
}
}
diff --git a/meteor-core/src/main/java/dev/pixelib/meteor/core/transport/TransportHandler.java b/meteor-core/src/main/java/dev/pixelib/meteor/core/transport/TransportHandler.java
index b1dea18..b40f93d 100644
--- a/meteor-core/src/main/java/dev/pixelib/meteor/core/transport/TransportHandler.java
+++ b/meteor-core/src/main/java/dev/pixelib/meteor/core/transport/TransportHandler.java
@@ -93,7 +93,7 @@ private boolean handleInvocationRequest(byte[] bytes) throws ClassNotFoundExcept
Object response = matchedImplementation.invokeOn(invocationDescriptor, invocationDescriptor.getReturnType());
InvocationResponse invocationResponse = new InvocationResponse(invocationDescriptor.getUniqueInvocationId(), response);
transport.send(Direction.METHOD_PROXY, invocationResponse.toBytes(serializer));
- } catch (Throwable e) {
+ } catch (Exception e) {
logger.log(Level.SEVERE, "An error occurred while invoking a method", e);
}
});
diff --git a/meteor-core/src/main/java/dev/pixelib/meteor/core/utils/ArgumentMapper.java b/meteor-core/src/main/java/dev/pixelib/meteor/core/utils/ArgumentMapper.java
index 5f9d000..c076162 100644
--- a/meteor-core/src/main/java/dev/pixelib/meteor/core/utils/ArgumentMapper.java
+++ b/meteor-core/src/main/java/dev/pixelib/meteor/core/utils/ArgumentMapper.java
@@ -1,10 +1,13 @@
package dev.pixelib.meteor.core.utils;
+import lombok.experimental.UtilityClass;
+
import java.lang.reflect.Array;
import java.lang.reflect.Method;
import java.util.Map;
import java.util.Objects;
+@UtilityClass
public class ArgumentMapper {
/**
@@ -74,9 +77,8 @@ public static Object[] overflowArguments(Method method, Object[] allArguments) {
return output;
}
- for (int i = 0; i < method.getParameterCount() - 1; i++) {
- output[i] = allArguments[i];
- }
+ int fixedParamCount = Math.max(0, method.getParameterCount() - 1);
+ System.arraycopy(allArguments, 0, output, 0, fixedParamCount);
Class> lastParameterType = method.getParameterTypes()[method.getParameterCount() - 1];
if (lastParameterType.isArray()) {
diff --git a/meteor-core/src/test/java/dev/pixelib/meteor/core/LogicTest.java b/meteor-core/src/test/java/dev/pixelib/meteor/core/LogicTest.java
index 8d7ded0..48d4ee6 100644
--- a/meteor-core/src/test/java/dev/pixelib/meteor/core/LogicTest.java
+++ b/meteor-core/src/test/java/dev/pixelib/meteor/core/LogicTest.java
@@ -1,21 +1,18 @@
package dev.pixelib.meteor.core;
import dev.pixelib.meteor.base.RpcSerializer;
-import dev.pixelib.meteor.base.RpcTransport;
import dev.pixelib.meteor.base.defaults.GsonSerializer;
-import dev.pixelib.meteor.base.defaults.LoopbackTransport;
import dev.pixelib.meteor.core.executor.ImplementationWrapper;
import dev.pixelib.meteor.core.transport.packets.InvocationDescriptor;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertEquals;
-public class LogicTest {
+class LogicTest {
@Test
- public void testSerializedReflectionWithPrimitiveArray() throws ClassNotFoundException, NoSuchMethodException {
+ void testSerializedReflectionWithPrimitiveArray() throws ClassNotFoundException, NoSuchMethodException {
// Confirmation that core logic of barebones reflection invocations over a serialized array still works
- RpcTransport transport = new LoopbackTransport();
RpcSerializer serializer = new GsonSerializer();
// this should return 30
@@ -40,9 +37,8 @@ public void testSerializedReflectionWithPrimitiveArray() throws ClassNotFoundExc
}
@Test
- public void testSerializedReflectionWithPrimitiveArrayAndNull() throws ClassNotFoundException, NoSuchMethodException {
+ void testSerializedReflectionWithPrimitiveArrayAndNull() throws ClassNotFoundException, NoSuchMethodException {
// Confirmation that core logic of barebones reflection invocations over a serialized array still works
- RpcTransport transport = new LoopbackTransport();
RpcSerializer serializer = new GsonSerializer();
// this should return 30
@@ -67,9 +63,8 @@ public void testSerializedReflectionWithPrimitiveArrayAndNull() throws ClassNotF
}
@Test
- public void testSerializedReflectionWithPrimitiveArrayAndNullAndNull() throws ClassNotFoundException, NoSuchMethodException {
+ void testSerializedReflectionWithPrimitiveArrayAndNullAndNull() throws ClassNotFoundException, NoSuchMethodException {
// Confirmation that core logic of barebones reflection invocations over a serialized array still works
- RpcTransport transport = new LoopbackTransport();
RpcSerializer serializer = new GsonSerializer();
// this should return 30
@@ -93,7 +88,6 @@ public void testSerializedReflectionWithPrimitiveArrayAndNullAndNull() throws Cl
assertEquals("hellonullworld", result);
}
-
public String concatStrings(String[] strings) {
StringBuilder sb = new StringBuilder();
for (String s : strings) {
@@ -102,6 +96,7 @@ public String concatStrings(String[] strings) {
return sb.toString();
}
+
public int worstCaseScenarioTest(int a, int... addBeforeMultiplying) {
int result = 0;
for (int i : addBeforeMultiplying) {
diff --git a/meteor-core/src/test/java/dev/pixelib/meteor/core/MeteorLoopbackTest.java b/meteor-core/src/test/java/dev/pixelib/meteor/core/MeteorLoopbackTest.java
index 3f288a1..5c8a38a 100644
--- a/meteor-core/src/test/java/dev/pixelib/meteor/core/MeteorLoopbackTest.java
+++ b/meteor-core/src/test/java/dev/pixelib/meteor/core/MeteorLoopbackTest.java
@@ -9,10 +9,10 @@
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
-public class MeteorLoopbackTest {
+class MeteorLoopbackTest {
@Test
- public void testLoopbackFunctionality() {
+ void testLoopbackFunctionality() {
Meteor meteor = new Meteor(new LoopbackTransport());
// register a procedure
diff --git a/meteor-core/src/test/java/dev/pixelib/meteor/core/invocations/PendingInvocationTest.java b/meteor-core/src/test/java/dev/pixelib/meteor/core/invocations/PendingInvocationTest.java
index 8cc21dc..3d0fe9f 100644
--- a/meteor-core/src/test/java/dev/pixelib/meteor/core/invocations/PendingInvocationTest.java
+++ b/meteor-core/src/test/java/dev/pixelib/meteor/core/invocations/PendingInvocationTest.java
@@ -13,32 +13,32 @@
import org.junit.jupiter.api.Timeout;
import java.util.Timer;
+import java.util.concurrent.CountDownLatch;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
-import static org.junit.jupiter.api.Assertions.assertEquals;
-import static org.junit.jupiter.api.Assertions.assertThrowsExactly;
+import static org.junit.jupiter.api.Assertions.*;
-public class PendingInvocationTest {
+class PendingInvocationTest {
// test thread pool
private static ThreadPoolExecutor threadPoolExecutor;
@BeforeAll
- public static void setUp() {
+ static void setUp() {
// create thread pool
threadPoolExecutor = new ThreadPoolExecutor(5, 5, 5, TimeUnit.SECONDS, new LinkedBlockingQueue<>());
}
@AfterAll
- public static void tearDown() {
+ static void tearDown() {
// shutdown thread pool
threadPoolExecutor.shutdown();
}
@Test
- public void testPendingInvocation() throws Throwable {
+ void testPendingInvocation() throws Throwable {
// base instance
OutgoingInvocationTracker outgoingInvocationTracker = new OutgoingInvocationTracker(new LoopbackTransport(), new GsonSerializer(), new RpcOptions(), new Timer());
@@ -48,13 +48,10 @@ public void testPendingInvocation() throws Throwable {
// complete invocation
threadPoolExecutor.execute(() -> {
- try {
- Thread.sleep(1000);
- } catch (InterruptedException e) {
- e.printStackTrace();
- }
- outgoingInvocationTracker.completeInvocation(
- new InvocationResponse(invocationDescriptor.getUniqueInvocationId(), testString)
+ assertDoesNotThrow(() -> {
+ new CountDownLatch(1).await(1, TimeUnit.SECONDS);
+ }, "Thread interrupted");
+ outgoingInvocationTracker.completeInvocation(new InvocationResponse(invocationDescriptor.getUniqueInvocationId(), testString)
);
});
@@ -64,7 +61,7 @@ public void testPendingInvocation() throws Throwable {
@Test
@Timeout(2) // seconds
- public void testTimeout() {
+ void testTimeout() {
RpcOptions options = new RpcOptions();
options.setTimeoutSeconds(1);
OutgoingInvocationTracker outgoingInvocationTracker = new OutgoingInvocationTracker(new LoopbackTransport(), new GsonSerializer(), options, new Timer());
diff --git a/meteor-core/src/test/java/dev/pixelib/meteor/core/trackers/IncomingInvocationTrackerTest.java b/meteor-core/src/test/java/dev/pixelib/meteor/core/trackers/IncomingInvocationTrackerTest.java
index 5ffee70..08beb3d 100644
--- a/meteor-core/src/test/java/dev/pixelib/meteor/core/trackers/IncomingInvocationTrackerTest.java
+++ b/meteor-core/src/test/java/dev/pixelib/meteor/core/trackers/IncomingInvocationTrackerTest.java
@@ -15,21 +15,19 @@ void testRegisterImplementation_thenSuccess() {
incomingInvocationTracker.registerImplementation(testMathFunctions, "test");
- assertEquals(incomingInvocationTracker.getImplementations().size(), 1);
- assertEquals(incomingInvocationTracker.getImplementations().get(MathFunctions.class).size(), 1);
+ assertEquals(1, incomingInvocationTracker.getImplementations().size());
+ assertEquals(1, incomingInvocationTracker.getImplementations().get(MathFunctions.class).size());
boolean matched = false;
for (ImplementationWrapper implementationWrapper : incomingInvocationTracker.getImplementations().get(MathFunctions.class)) {
- if (implementationWrapper.getImplementation() == testMathFunctions) {
- // inverted, because the namespace is nullable
- if ("test".equals(implementationWrapper.getNamespace())) {
+ if (implementationWrapper.getImplementation() == testMathFunctions && "test".equals(implementationWrapper.getNamespace())) {
if (matched) {
fail("Implementation registered twice");
return;
}
matched = true;
}
- }
+
}
assertTrue(matched, "Implementation not registered");
diff --git a/meteor-core/src/test/java/dev/pixelib/meteor/core/transport/packets/InvocationDescriptorTest.java b/meteor-core/src/test/java/dev/pixelib/meteor/core/transport/packets/InvocationDescriptorTest.java
index d200708..a3672e9 100644
--- a/meteor-core/src/test/java/dev/pixelib/meteor/core/transport/packets/InvocationDescriptorTest.java
+++ b/meteor-core/src/test/java/dev/pixelib/meteor/core/transport/packets/InvocationDescriptorTest.java
@@ -25,7 +25,7 @@ private void compareInstances(InvocationDescriptor a, InvocationDescriptor b) {
}
@Test
- public void testSerializationWithNamespace() throws ClassNotFoundException {
+ void testSerializationWithNamespace() throws ClassNotFoundException {
RpcSerializer defaultSerializer = new GsonSerializer();
InvocationDescriptor original = new InvocationDescriptor(
@@ -44,7 +44,7 @@ public void testSerializationWithNamespace() throws ClassNotFoundException {
}
@Test
- public void testSerializationWithoutNamespace() throws ClassNotFoundException {
+ void testSerializationWithoutNamespace() throws ClassNotFoundException {
RpcSerializer defaultSerializer = new GsonSerializer();
InvocationDescriptor original = new InvocationDescriptor(
diff --git a/meteor-core/src/test/java/dev/pixelib/meteor/core/utils/ArgumentMapperTest.java b/meteor-core/src/test/java/dev/pixelib/meteor/core/utils/ArgumentMapperTest.java
index 7dcd609..48f66ab 100644
--- a/meteor-core/src/test/java/dev/pixelib/meteor/core/utils/ArgumentMapperTest.java
+++ b/meteor-core/src/test/java/dev/pixelib/meteor/core/utils/ArgumentMapperTest.java
@@ -87,7 +87,6 @@ void testResolvePrimitive_nullValue() {
assertEquals("className cannot be null", exception.getMessage());
}
-
static class Example {
private void singleParamMethod(Integer integer) { }
private void multipleParamsMethod(Integer integer, String str, Double dd) { }
diff --git a/meteor-jedis/pom.xml b/meteor-jedis/pom.xml
index bace8f0..b6f0a0f 100644
--- a/meteor-jedis/pom.xml
+++ b/meteor-jedis/pom.xml
@@ -10,40 +10,34 @@
../pom.xml
- dev.pixelib.meteor.transport
meteor-jedis
-
- 17
- 17
- UTF-8
-
-
dev.pixelib.meteor
meteor-common
- ${revision}
+ ${project.version}
redis.clients
jedis
+ 7.5.2
- org.junit.jupiter
- junit-jupiter
+ com.github.fppt
+ jedis-mock
+ 1.1.15
test
- com.github.fppt
- jedis-mock
- 1.0.10
+ org.junit.jupiter
+ junit-jupiter
test
-
+
org.mockito
mockito-junit-jupiter
diff --git a/meteor-jedis/src/main/java/dev/pixelib/meteor/transport/redis/RedisPacketListener.java b/meteor-jedis/src/main/java/dev/pixelib/meteor/transport/redis/RedisPacketListener.java
index e45e437..c2f501a 100644
--- a/meteor-jedis/src/main/java/dev/pixelib/meteor/transport/redis/RedisPacketListener.java
+++ b/meteor-jedis/src/main/java/dev/pixelib/meteor/transport/redis/RedisPacketListener.java
@@ -1,16 +1,13 @@
package dev.pixelib.meteor.transport.redis;
-import dev.pixelib.meteor.base.interfaces.SubscriptionHandler;
import redis.clients.jedis.JedisPubSub;
-import java.util.Base64;
import java.util.Collection;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
-import java.util.function.Consumer;
import java.util.logging.Level;
import java.util.logging.Logger;
diff --git a/meteor-jedis/src/main/java/dev/pixelib/meteor/transport/redis/RedisSubscriptionThread.java b/meteor-jedis/src/main/java/dev/pixelib/meteor/transport/redis/RedisSubscriptionThread.java
index 4859cfb..3e52c9e 100644
--- a/meteor-jedis/src/main/java/dev/pixelib/meteor/transport/redis/RedisSubscriptionThread.java
+++ b/meteor-jedis/src/main/java/dev/pixelib/meteor/transport/redis/RedisSubscriptionThread.java
@@ -1,15 +1,13 @@
package dev.pixelib.meteor.transport.redis;
-import dev.pixelib.meteor.base.interfaces.SubscriptionHandler;
-import redis.clients.jedis.Jedis;
-import redis.clients.jedis.JedisPool;
+import redis.clients.jedis.Connection;
+import redis.clients.jedis.RedisClient;
import redis.clients.jedis.exceptions.JedisConnectionException;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
-import java.util.function.Consumer;
import java.util.logging.Level;
import java.util.logging.Logger;
@@ -18,7 +16,7 @@ public class RedisSubscriptionThread {
private final StringMessageBroker messageBroker;
private final Logger logger;
private final String defaultChannel;
- private final JedisPool jedisPool;
+ private final RedisClient redisClient;
private boolean isStopping = false;
private RedisPacketListener jedisPacketListener;
@@ -29,11 +27,11 @@ public class RedisSubscriptionThread {
return thread;
});
- public RedisSubscriptionThread(StringMessageBroker messageBroker, Logger logger, String channel, JedisPool jedisPool) {
+ public RedisSubscriptionThread(StringMessageBroker messageBroker, Logger logger, String channel, RedisClient redisClient) {
this.messageBroker = messageBroker;
this.logger = logger;
this.defaultChannel = channel;
- this.jedisPool = jedisPool;
+ this.redisClient = redisClient;
}
public CompletableFuture start() {
@@ -41,12 +39,11 @@ public CompletableFuture start() {
Runnable runnable = () -> {
while (!Thread.currentThread().isInterrupted()) {
- try (Jedis connection = jedisPool.getResource()) {
+ try (Connection connection = redisClient.getPool().getResource()) {
connection.ping();
logger.log(Level.FINE, "Redis connected!");
- //Start blocking
- connection.subscribe(jedisPacketListener, jedisPacketListener.getCustomSubscribedChannels().toArray(new String[]{}));
+ jedisPacketListener.proceed(connection, jedisPacketListener.getCustomSubscribedChannels().toArray(new String[]{}));
break;
} catch (JedisConnectionException e) {
if (isStopping) {
@@ -65,21 +62,20 @@ public CompletableFuture start() {
};
listenerThread.execute(runnable);
-
return isSubscribed();
-
}
public void stop() {
if (isStopping) return;
isStopping = true;
- jedisPacketListener.stop();
+ if (jedisPacketListener != null && jedisPacketListener.isSubscribed()) {
+ jedisPacketListener.stop();
+ }
listenerThread.shutdownNow();
}
public void subscribe(String channel, StringMessageBroker onReceive) {
jedisPacketListener.subscribe(channel, onReceive);
-
}
private CompletableFuture isSubscribed() {
@@ -98,7 +94,6 @@ private CompletableFuture isSubscribed() {
}
}
- // If it fails to subscribe within 5 attempts (5 seconds), throw an exception
throw new IllegalStateException("Failed to subscribe within the given timeframe");
});
}
diff --git a/meteor-jedis/src/main/java/dev/pixelib/meteor/transport/redis/RedisTransport.java b/meteor-jedis/src/main/java/dev/pixelib/meteor/transport/redis/RedisTransport.java
index 2a28289..ba6945c 100644
--- a/meteor-jedis/src/main/java/dev/pixelib/meteor/transport/redis/RedisTransport.java
+++ b/meteor-jedis/src/main/java/dev/pixelib/meteor/transport/redis/RedisTransport.java
@@ -3,38 +3,40 @@
import dev.pixelib.meteor.base.RpcTransport;
import dev.pixelib.meteor.base.enums.Direction;
import dev.pixelib.meteor.base.interfaces.SubscriptionHandler;
-import redis.clients.jedis.Jedis;
-import redis.clients.jedis.JedisPool;
+import redis.clients.jedis.RedisClient;
import java.io.IOException;
import java.util.Base64;
import java.util.Locale;
import java.util.UUID;
-import java.util.function.Consumer;
import java.util.logging.Logger;
public class RedisTransport implements RpcTransport {
+ private static final Base64.Decoder base64Decoder = Base64.getDecoder();
+ private static final Base64.Encoder base64Encoder = Base64.getEncoder();
+
private final Logger logger = Logger.getLogger(RedisTransport.class.getSimpleName());
- private final JedisPool jedisPool;
+ private final RedisClient redisClient;
private final String topic;
- private RedisSubscriptionThread redisSubscriptionThread;
private final UUID transportId = UUID.randomUUID();
+ private RedisSubscriptionThread redisSubscriptionThread;
private boolean ignoreSelf = true;
+ private boolean closed;
- public RedisTransport(JedisPool jedisPool, String topic) {
- this.jedisPool = jedisPool;
+ public RedisTransport(RedisClient redisClient, String topic) {
+ this.redisClient = redisClient;
this.topic = topic;
}
public RedisTransport(String url, String topic) {
- this.jedisPool = new JedisPool(url);
+ this.redisClient = RedisClient.create(url);
this.topic = topic;
}
public RedisTransport(String host, int port, String topic) {
- this.jedisPool = new JedisPool(host, port);
+ this.redisClient = RedisClient.create(host, port);
this.topic = topic;
}
@@ -45,27 +47,22 @@ public RedisTransport withIgnoreSelf(boolean ignoreSelf) {
@Override
public void send(Direction direction, byte[] bytes) {
- if (jedisPool.isClosed()) {
- throw new IllegalStateException("Jedis pool is closed");
- }
-
- try (Jedis connection = jedisPool.getResource()) {
- connection.publish(
- getTopicName(direction),
- transportId + Base64.getEncoder().encodeToString(bytes)
- );
+ if (closed) {
+ throw new IllegalStateException("RedisTransport is closed");
}
+ redisClient.publish(
+ getTopicName(direction),
+ transportId + base64Encoder.encodeToString(bytes)
+ );
}
@Override
public void subscribe(Direction direction, SubscriptionHandler onReceive) {
- if (jedisPool.isClosed()) {
- throw new IllegalStateException("Jedis pool is closed");
+ if (closed) {
+ throw new IllegalStateException("RedisTransport is closed");
}
-
- StringMessageBroker wrappedHandler = (message) -> {
+ StringMessageBroker wrappedHandler = message -> {
byte[] bytes = message.getBytes();
- // only split after UUID, which always has a length of 36
byte[] uuid = new byte[36];
System.arraycopy(bytes, 0, uuid, 0, 36);
@@ -77,14 +74,14 @@ public void subscribe(Direction direction, SubscriptionHandler onReceive) {
System.arraycopy(bytes, 36, data, 0, data.length);
try {
- return onReceive.onPacket(Base64.getDecoder().decode(data));
+ return onReceive.onPacket(base64Decoder.decode(data));
} catch (Exception e) {
throw new RuntimeException(e);
}
};
if (redisSubscriptionThread == null) {
- redisSubscriptionThread = new RedisSubscriptionThread(wrappedHandler, logger, getTopicName(direction), jedisPool);
+ redisSubscriptionThread = new RedisSubscriptionThread(wrappedHandler, logger, getTopicName(direction), redisClient);
redisSubscriptionThread.start().join();
} else {
redisSubscriptionThread.subscribe(getTopicName(direction), wrappedHandler);
@@ -97,10 +94,10 @@ public String getTopicName(Direction direction) {
@Override
public void close() throws IOException {
+ closed = true;
if (redisSubscriptionThread != null) {
redisSubscriptionThread.stop();
}
-
- jedisPool.close();
+ redisClient.close();
}
}
diff --git a/meteor-jedis/src/test/java/dev/pixelib/meteor/transport/redis/RedisPacketListenerTest.java b/meteor-jedis/src/test/java/dev/pixelib/meteor/transport/redis/RedisPacketListenerTest.java
index 0adb592..ee4e940 100644
--- a/meteor-jedis/src/test/java/dev/pixelib/meteor/transport/redis/RedisPacketListenerTest.java
+++ b/meteor-jedis/src/test/java/dev/pixelib/meteor/transport/redis/RedisPacketListenerTest.java
@@ -1,17 +1,16 @@
package dev.pixelib.meteor.transport.redis;
import com.github.fppt.jedismock.RedisServer;
-import dev.pixelib.meteor.base.interfaces.SubscriptionHandler;
-import org.junit.jupiter.api.Disabled;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
-import redis.clients.jedis.JedisPool;
+import redis.clients.jedis.Connection;
+import redis.clients.jedis.RedisClient;
-import java.util.Base64;
import java.util.Collection;
+import java.util.concurrent.CountDownLatch;
import java.util.logging.Logger;
import static org.junit.jupiter.api.Assertions.*;
@@ -32,10 +31,6 @@ void onMessage_withValidChannel() throws Exception {
redisPacketListener.onMessage(topic, expected);
verify(subscriptionHandler, times(1)).onRedisMessage(expected);
- verify(subscriptionHandler).onRedisMessage(argThat(argument -> {
- assertEquals(expected,argument);
- return true;
- }));
}
@Test
@@ -55,14 +50,10 @@ public boolean onRedisMessage(String message) throws Exception {
redisPacketListener.onMessage(topic, expected);
verify(handlerSub, times(1)).onRedisMessage(expected);
- verify(handlerSub).onRedisMessage(argThat(argument -> {
- assertEquals(expected,argument);
- return true;
- }));
}
@Test
- void onMessage_withUnKnownChannel() throws Exception {
+ void onMessage_withUnKnownChannel() {
String topic = "test";
String message = "message";
@@ -73,73 +64,77 @@ void onMessage_withUnKnownChannel() throws Exception {
@Test
@Timeout(10)
- @Disabled
- void subscribe_success() throws Exception{
+ void subscribe_success() throws Exception {
String topic = "test";
String newTopic = "newTopic";
RedisServer server = RedisServer.newRedisServer().start();
-
- JedisPool jedisPool = new JedisPool(server.getHost(), server.getBindPort());
-
-
- RedisPacketListener redisPacketListener = new RedisPacketListener(subscriptionHandler, topic, Logger.getAnonymousLogger());
-
- Thread runner = new Thread(() -> {
- jedisPool.getResource().subscribe(redisPacketListener, topic);
- });
-
- runner.start();
-
- while(!redisPacketListener.isSubscribed()) {
- Thread.sleep(20);
- }
-
- redisPacketListener.subscribe(newTopic, subscriptionHandler);
-
- assertTrue(redisPacketListener.getCustomSubscribedChannels().contains(newTopic));
-
- redisPacketListener.stop();
-
- while(redisPacketListener.isSubscribed()) {
- Thread.sleep(20);
+ try {
+ RedisClient redisClient = new RedisClient.Builder().hostAndPort(server.getHost(), server.getBindPort()).build();
+
+ CountDownLatch subscribedLatch = new CountDownLatch(1);
+ RedisPacketListener redisPacketListener = new RedisPacketListener(subscriptionHandler, topic, Logger.getAnonymousLogger()) {
+ @Override
+ public void onSubscribe(String channel, int subscribedChannels) {
+ super.onSubscribe(channel, subscribedChannels);
+ subscribedLatch.countDown();
+ }
+ };
+
+ Thread runner = new Thread(() -> {
+ Connection connection = redisClient.getPool().getResource();
+ redisPacketListener.proceed(connection, topic);
+ });
+ runner.start();
+
+ subscribedLatch.await();
+
+ redisPacketListener.subscribe(newTopic, subscriptionHandler);
+ assertTrue(redisPacketListener.getCustomSubscribedChannels().contains(newTopic));
+
+ redisPacketListener.stop();
+ runner.join();
+
+ redisClient.close();
+ } finally {
+ server.stop();
}
- jedisPool.close();
- server.stop();
-
-
- assertEquals(0, redisPacketListener.getSubscribedChannels());
}
@Test
@Timeout(10)
- @Disabled
- void stop_success() throws Exception{
+ void stop_success() throws Exception {
String topic = "test";
RedisServer server = RedisServer.newRedisServer().start();
-
- JedisPool jedisPool = new JedisPool(server.getHost(), server.getBindPort());
-
-
- RedisPacketListener redisPacketListener = new RedisPacketListener(subscriptionHandler, topic, Logger.getAnonymousLogger());
-
- Thread runner = new Thread(() -> {
- jedisPool.getResource().subscribe(redisPacketListener, topic);
- });
-
- runner.start();
-
- while(!redisPacketListener.isSubscribed()) {
- Thread.sleep(20);
+ try {
+ RedisClient redisClient = new RedisClient.Builder().hostAndPort(server.getHost(), server.getBindPort()).build();
+
+ CountDownLatch subscribedLatch = new CountDownLatch(1);
+ RedisPacketListener redisPacketListener = new RedisPacketListener(subscriptionHandler, topic, Logger.getAnonymousLogger()) {
+ @Override
+ public void onSubscribe(String channel, int subscribedChannels) {
+ super.onSubscribe(channel, subscribedChannels);
+ subscribedLatch.countDown();
+ }
+ };
+
+ Thread runner = new Thread(() -> {
+ Connection connection = redisClient.getPool().getResource();
+ redisPacketListener.proceed(connection, topic);
+ });
+ runner.start();
+
+ subscribedLatch.await();
+
+ redisPacketListener.stop();
+ runner.join();
+
+ redisClient.close();
+ } finally {
+ server.stop();
}
-
- redisPacketListener.stop();
- jedisPool.close();
- server.stop();
-
- assertEquals(0, redisPacketListener.getSubscribedChannels());
}
@Test
@@ -154,4 +149,4 @@ void getCustomSubscribedChannels_success() {
assertEquals(1, result.size());
assertTrue(result.contains(topic));
}
-}
\ No newline at end of file
+}
diff --git a/meteor-jedis/src/test/java/dev/pixelib/meteor/transport/redis/RedisSubscriptionThreadTest.java b/meteor-jedis/src/test/java/dev/pixelib/meteor/transport/redis/RedisSubscriptionThreadTest.java
index 788638e..de7f02e 100644
--- a/meteor-jedis/src/test/java/dev/pixelib/meteor/transport/redis/RedisSubscriptionThreadTest.java
+++ b/meteor-jedis/src/test/java/dev/pixelib/meteor/transport/redis/RedisSubscriptionThreadTest.java
@@ -3,10 +3,12 @@
import com.github.fppt.jedismock.RedisServer;
import com.github.fppt.jedismock.operations.server.MockExecutor;
import com.github.fppt.jedismock.server.ServiceOptions;
-import org.junit.jupiter.api.Disabled;
+import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
-import redis.clients.jedis.JedisPool;
+import org.junit.jupiter.api.Timeout;
+import redis.clients.jedis.RedisClient;
+import java.io.IOException;
import java.util.Collection;
import java.util.HashSet;
import java.util.concurrent.CompletionException;
@@ -16,30 +18,40 @@
class RedisSubscriptionThreadTest {
+ private RedisSubscriptionThread subThread;
+ private RedisServer server;
+ private RedisClient redisClient;
+
+ @AfterEach
+ void tearDown() throws IOException {
+ if (subThread != null) {
+ subThread.stop();
+ }
+ if (redisClient != null && !redisClient.getPool().isClosed()) {
+ redisClient.close();
+ }
+ if (server != null) {
+ server.stop();
+ }
+ }
+
@Test
- @Disabled
void start_success() throws Exception {
- RedisServer server = RedisServer.newRedisServer().start();
-
- JedisPool jedisPool = new JedisPool(server.getHost(), server.getBindPort());
- RedisSubscriptionThread subThread = new RedisSubscriptionThread(packet -> true, Logger.getAnonymousLogger(), "channel", jedisPool);
+ server = RedisServer.newRedisServer().start();
+ redisClient = new RedisClient.Builder().hostAndPort(server.getHost(), server.getBindPort()).build();
+ subThread = new RedisSubscriptionThread(packet -> true, Logger.getAnonymousLogger(), "channel", redisClient);
boolean result = subThread.start().join();
assertTrue(result);
- assertFalse(jedisPool.isClosed());
-
-
- subThread.stop();
- jedisPool.close();
- server.stop();
- assertTrue(jedisPool.isClosed());
+ assertFalse(redisClient.getPool().isClosed());
}
@Test
- void start_NoConnection() throws Exception {
- JedisPool jedisPool = new JedisPool("127.0.0.5", 2314);
- RedisSubscriptionThread subThread = new RedisSubscriptionThread(packet -> true, Logger.getAnonymousLogger(), "channel", jedisPool);
+ @Timeout(10)
+ void start_NoConnection() {
+ redisClient = new RedisClient.Builder().hostAndPort("127.0.0.5", 2314).build();
+ subThread = new RedisSubscriptionThread(packet -> true, Logger.getAnonymousLogger(), "channel", redisClient);
CompletionException exception = assertThrowsExactly(CompletionException.class, () -> {
subThread.start().join();
@@ -50,67 +62,50 @@ void start_NoConnection() throws Exception {
}
@Test
- @Disabled
- void start_reconnect() throws Exception {
- RedisServer server = RedisServer.newRedisServer().start();
-
- JedisPool jedisPool = new JedisPool(server.getHost(), server.getBindPort());
- RedisSubscriptionThread subThread = new RedisSubscriptionThread(packet -> true, Logger.getAnonymousLogger(), "channel", jedisPool);
-
- boolean result = subThread.start().join();
-
- assertTrue(result);
- assertFalse(jedisPool.isClosed());
-
- server.stop();
-
- Thread.sleep(100);
- server.start();
-
- subThread.stop();
- jedisPool.close();
- assertTrue(jedisPool.isClosed());
- }
-
- @Test
- @Disabled
+ @Timeout(10)
void subscribe_success() throws Exception {
Collection subscribedChannels = new HashSet<>();
- RedisServer server = RedisServer.newRedisServer()
+ server = RedisServer.newRedisServer()
.setOptions(ServiceOptions.withInterceptor((state, command, params) -> {
if ("subscribe".equals(command)) {
- subscribedChannels.add(params.get(0).toString());
+ subscribedChannels.add(params.getFirst().toString());
}
return MockExecutor.proceed(state, command, params);
}))
.start();
- JedisPool jedisPool = new JedisPool(server.getHost(), server.getBindPort());
- RedisSubscriptionThread subThread = new RedisSubscriptionThread(packet -> true, Logger.getAnonymousLogger(), "channel", jedisPool);
+ redisClient = new RedisClient.Builder().hostAndPort(server.getHost(), server.getBindPort()).build();
+ subThread = new RedisSubscriptionThread(packet -> true, Logger.getAnonymousLogger(), "channel", redisClient);
subThread.start().join();
- subThread.stop();
- jedisPool.close();
- server.stop();
- assertTrue(jedisPool.isClosed());
assertTrue(subscribedChannels.contains("channel"));
}
@Test
- @Disabled
- void stop_success() throws Exception {
- RedisServer server = RedisServer.newRedisServer().start();
-
- JedisPool jedisPool = new JedisPool(server.getHost(), server.getBindPort());
- RedisSubscriptionThread subThread = new RedisSubscriptionThread(packet -> true, Logger.getAnonymousLogger(), "channel", jedisPool);
+ @Timeout(10)
+ void stop_idempotent() throws Exception {
+ server = RedisServer.newRedisServer().start();
+ redisClient = new RedisClient.Builder().hostAndPort(server.getHost(), server.getBindPort()).build();
+ subThread = new RedisSubscriptionThread(packet -> true, Logger.getAnonymousLogger(), "channel", redisClient);
subThread.start().join();
+ subThread.stop();
+ assertDoesNotThrow(() -> subThread.stop());
+ }
+ @Test
+ @Timeout(10)
+ void stop_success() throws Exception {
+ server = RedisServer.newRedisServer().start();
+ redisClient = new RedisClient.Builder().hostAndPort(server.getHost(), server.getBindPort()).build();
+ subThread = new RedisSubscriptionThread(packet -> true, Logger.getAnonymousLogger(), "channel", redisClient);
+
+ subThread.start().join();
subThread.stop();
- jedisPool.close();
- server.stop();
- assertTrue(jedisPool.isClosed());
+ redisClient.close();
+
+ assertTrue(redisClient.getPool().isClosed());
}
-}
\ No newline at end of file
+}
diff --git a/meteor-jedis/src/test/java/dev/pixelib/meteor/transport/redis/RedisTransportTest.java b/meteor-jedis/src/test/java/dev/pixelib/meteor/transport/redis/RedisTransportTest.java
index 0187186..4c0185d 100644
--- a/meteor-jedis/src/test/java/dev/pixelib/meteor/transport/redis/RedisTransportTest.java
+++ b/meteor-jedis/src/test/java/dev/pixelib/meteor/transport/redis/RedisTransportTest.java
@@ -4,306 +4,186 @@
import com.github.fppt.jedismock.operations.server.MockExecutor;
import com.github.fppt.jedismock.server.ServiceOptions;
import dev.pixelib.meteor.base.enums.Direction;
-import org.junit.jupiter.api.Disabled;
+import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
-import org.opentest4j.AssertionFailedError;
-import redis.clients.jedis.JedisPool;
+import redis.clients.jedis.RedisClient;
import java.io.IOException;
-import java.util.*;
+import java.util.ArrayList;
+import java.util.Base64;
+import java.util.List;
import static org.junit.jupiter.api.Assertions.*;
class RedisTransportTest {
+ private RedisTransport transport;
+ private RedisServer server;
+
+ @AfterEach
+ void tearDown() throws IOException {
+ if (transport != null) {
+ transport.close();
+ }
+ if (server != null) {
+ server.stop();
+ }
+ }
+
@Test
- @Disabled
void send_validImplementation() throws IOException {
- String topic = "test";
- String channel = "test_implementation";
String message = "cool_message";
- List assertionErrors = new ArrayList<>();
+ List published = new ArrayList<>();
- RedisServer server = RedisServer.newRedisServer()
+ server = RedisServer.newRedisServer()
.setOptions(ServiceOptions.withInterceptor((state, command, params) -> {
- try {
- if ("publish".equals(command)) {
- // loop over all params and print them
- assertEquals(channel, params.get(0).toString(), "Channel name is not correct");
- assertEquals(message, new String(Base64.getDecoder().decode(params.get(1).data())), "Message is not correct");
- }
- } catch (AssertionFailedError e) {
- assertionErrors.add(e);
+ if ("publish".equals(command)) {
+ published.add(new String(params.getLast().data()));
}
return MockExecutor.proceed(state, command, params);
}))
.start();
- RedisTransport transport = new RedisTransport(server.getHost(), server.getBindPort(), topic);
+ transport = new RedisTransport(server.getHost(), server.getBindPort(), "test");
transport.send(Direction.IMPLEMENTATION, message.getBytes());
- transport.close();
- server.stop();
-
- if (!assertionErrors.isEmpty()) {
- throw assertionErrors.get(0);
- }
+ assertEquals(1, published.size());
+ String raw = published.getFirst();
+ String base64Data = raw.substring(36);
+ assertEquals(message, new String(Base64.getDecoder().decode(base64Data)));
}
@Test
- @Disabled
void send_validMethodProxy() throws IOException {
- String topic = "test";
- String channel = "test_method_proxy";
String message = "cool_message";
- List assertionErrors = new ArrayList<>();
+ List published = new ArrayList<>();
- RedisServer server = RedisServer.newRedisServer()
+ server = RedisServer.newRedisServer()
.setOptions(ServiceOptions.withInterceptor((state, command, params) -> {
- try {
- if ("publish".equals(command)) {
- // loop over all params and print them
- assertEquals(channel, params.get(0).toString(), "Channel name is not correct");
- assertEquals(message, new String(Base64.getDecoder().decode(params.get(1).data())), "Message is not correct");
- }
- } catch (AssertionFailedError e) {
- assertionErrors.add(e);
+ if ("publish".equals(command)) {
+ published.add(new String(params.getLast().data()));
}
return MockExecutor.proceed(state, command, params);
}))
.start();
- RedisTransport transport = new RedisTransport(server.getHost(), server.getBindPort(), topic);
+ transport = new RedisTransport(server.getHost(), server.getBindPort(), "test");
transport.send(Direction.METHOD_PROXY, message.getBytes());
- transport.close();
- server.stop();
-
- if (!assertionErrors.isEmpty()) {
- throw assertionErrors.get(0);
- }
+ assertEquals(1, published.size());
+ String raw = published.getFirst();
+ String base64Data = raw.substring(36);
+ assertEquals(message, new String(Base64.getDecoder().decode(base64Data)));
}
@Test
- @Disabled
- void subscribe_implementation() throws IOException, InterruptedException {
- String topic = "test";
- String channel = "test_implementation";
+ void subscribe_implementation() throws IOException {
+ List subscribedChannels = new ArrayList<>();
- List assertionErrors = new ArrayList<>();
-
- RedisServer server = RedisServer.newRedisServer()
+ server = RedisServer.newRedisServer()
.setOptions(ServiceOptions.withInterceptor((state, command, params) -> {
- try {
- if ("subscribe".equals(command)) {
- assertEquals(channel, params.get(0).toString(), "Channel name is not correct");
- }
- } catch (AssertionFailedError e) {
- assertionErrors.add(e);
+ if ("subscribe".equals(command)) {
+ subscribedChannels.add(params.getFirst().toString());
}
return MockExecutor.proceed(state, command, params);
}))
.start();
- RedisTransport transport = new RedisTransport(server.getHost(), server.getBindPort(), topic);
+ transport = new RedisTransport(server.getHost(), server.getBindPort(), "test");
transport.subscribe(Direction.IMPLEMENTATION, packet -> true);
- transport.close();
- server.stop();
-
- if (!assertionErrors.isEmpty()) {
- throw assertionErrors.get(0);
- }
+ assertEquals(1, subscribedChannels.size());
+ assertEquals("test_implementation", subscribedChannels.getFirst());
}
@Test
- @Disabled
void subscribe_methodProxy() throws IOException {
- String topic = "test";
- String channel = "test_method_proxy";
-
- List assertionErrors = new ArrayList<>();
+ List subscribedChannels = new ArrayList<>();
- RedisServer server = RedisServer.newRedisServer()
+ server = RedisServer.newRedisServer()
.setOptions(ServiceOptions.withInterceptor((state, command, params) -> {
- try {
- if ("subscribe".equals(command)) {
- assertEquals(channel, params.get(0).toString(), "Channel name is not correct");
- }
- } catch (AssertionFailedError e) {
- assertionErrors.add(e);
+ if ("subscribe".equals(command)) {
+ subscribedChannels.add(params.getFirst().toString());
}
return MockExecutor.proceed(state, command, params);
}))
.start();
- RedisTransport transport = new RedisTransport(server.getHost(), server.getBindPort(), topic);
+ transport = new RedisTransport(server.getHost(), server.getBindPort(), "test");
transport.subscribe(Direction.METHOD_PROXY, packet -> true);
- transport.close();
- server.stop();
-
- if (!assertionErrors.isEmpty()) {
- throw assertionErrors.get(0);
- }
+ assertEquals(1, subscribedChannels.size());
+ assertEquals("test_method_proxy", subscribedChannels.getFirst());
}
@Test
- @Disabled
void subscribe_secondSubscription() throws IOException {
- String topic = "test";
- String channelProxy = "test_method_proxy";
- String channelImpl = "test_implementation";
-
- Collection subscribedChannels = new HashSet<>();
-
- RedisServer server = RedisServer.newRedisServer()
- .setOptions(ServiceOptions.withInterceptor((state, command, params) -> {
- if ("subscribe".equals(command)) {
- subscribedChannels.add(params.get(0).toString());
- }
- return MockExecutor.proceed(state, command, params);
- }))
- .start();
-
- RedisTransport transport = new RedisTransport(server.getHost(), server.getBindPort(), topic);
+ server = RedisServer.newRedisServer().start();
+ transport = new RedisTransport(server.getHost(), server.getBindPort(), "test");
transport.subscribe(Direction.METHOD_PROXY, packet -> true);
transport.subscribe(Direction.IMPLEMENTATION, packet -> true);
-
- transport.close();
- server.stop();
-
- assertTrue(subscribedChannels.contains(channelProxy));
- assertTrue(subscribedChannels.contains(channelImpl));
}
@Test
- void getTopicName_withImplementationDirection() throws IOException {
- String topic = "test";
- String expected = "test_implementation";
-
- RedisServer server = RedisServer.newRedisServer().start();
- RedisTransport transport = new RedisTransport(server.getHost(), server.getBindPort(), topic);
-
+ void getTopicName_withImplementationDirection() {
+ transport = new RedisTransport("redis://localhost:6379", "test");
String resultTopicName = transport.getTopicName(Direction.IMPLEMENTATION);
- assertEquals(expected, resultTopicName, "Topic name is not correct");
-
-
- transport.close();
- server.stop();
+ assertEquals("test_implementation", resultTopicName);
}
@Test
- @Disabled
- void getTopicName_withMethodProxy() throws IOException {
- String topic = "test";
- String expected = "test_method_proxy";
-
- RedisServer server = RedisServer.newRedisServer().start();
- RedisTransport transport = new RedisTransport(server.getHost(), server.getBindPort(), topic);
-
+ void getTopicName_withMethodProxy() {
+ transport = new RedisTransport("redis://localhost:6379", "test");
String resultTopicName = transport.getTopicName(Direction.METHOD_PROXY);
- assertEquals(expected, resultTopicName, "Topic name is not correct");
-
-
- transport.close();
- server.stop();
+ assertEquals("test_method_proxy", resultTopicName);
}
-
-
@Test
- @Disabled
- void getTopic_nameWithNull() throws IOException {
- String topic = "test";
-
- RedisServer server = RedisServer.newRedisServer().start();
- RedisTransport transport = new RedisTransport(server.getHost(), server.getBindPort(), topic);
-
+ void getTopic_nameWithNull() {
+ transport = new RedisTransport("redis://localhost:6379", "test");
assertThrowsExactly(NullPointerException.class, () -> {
transport.getTopicName(null);
- }, "Method returned a topic name for a null direction");
-
-
- transport.close();
- server.stop();
+ });
}
@Test
- @Disabled
- void close_success() throws IOException {
+ void close_success() throws IOException {
String topic = "test";
-
- RedisServer server = RedisServer.newRedisServer().start();
-
- JedisPool jedisPool = new JedisPool(server.getHost(), server.getBindPort());
- RedisTransport transport = new RedisTransport(jedisPool, topic);
+ transport = new RedisTransport("localhost", 6379, topic);
transport.close();
- server.stop();
-
- assertThrowsExactly(IllegalStateException.class, () -> {
- transport.send(Direction.IMPLEMENTATION, "test".getBytes());
- }, "Method did not throw an exception when trying to send a message after closing");
- assertThrowsExactly(IllegalStateException.class, () -> {
- transport.subscribe(Direction.IMPLEMENTATION, (data) -> true);
- }, "Method did not throw an exception when trying to send a message after closing");
-
- assertTrue(jedisPool.isClosed());
+ assertThrowsExactly(IllegalStateException.class, () -> transport.send(Direction.IMPLEMENTATION, "test".getBytes()));
+ assertThrowsExactly(IllegalStateException.class, () -> transport.subscribe(Direction.IMPLEMENTATION, data -> true));
}
-
@Test
- @Disabled
- void close_whenAlreadyClosed() throws IOException {
- String topic = "test";
-
- RedisServer server = RedisServer.newRedisServer().start();
-
- JedisPool jedisPool = new JedisPool(server.getHost(), server.getBindPort());
- RedisTransport transport = new RedisTransport(jedisPool, topic);
+ void close_whenAlreadyClosed() throws IOException {
+ server = RedisServer.newRedisServer().start();
+ transport = new RedisTransport(server.getHost(), server.getBindPort(), "test");
+ transport.subscribe(Direction.IMPLEMENTATION, packet -> true);
transport.close();
- server.stop();
-
- assertTrue(jedisPool.isClosed());
+ assertDoesNotThrow(() -> transport.close());
}
@Test
- @Disabled
- void construct_withJedisPool() throws IOException {
- String topic = "test";
-
- RedisServer server = RedisServer.newRedisServer().start();
-
- JedisPool jedisPool = new JedisPool(server.getHost(), server.getBindPort());
- RedisTransport transport = new RedisTransport(jedisPool, topic);
+ void construct_withRedisClient() {
+ RedisClient client = RedisClient.create("localhost", 6379);
+ transport = new RedisTransport(client, "test");
- assertFalse(jedisPool.isClosed());
-
- transport.close();
- server.stop();
-
- assertTrue(jedisPool.isClosed());
+ assertNotNull(transport);
}
- @Test
- @Disabled
- void construct_withUrl() throws IOException {
- String topic = "test";
- RedisServer server = RedisServer.newRedisServer().start();
+ @Test
+ void construct_withUrl() {
+ transport = new RedisTransport("redis://localhost:6379", "test");
- RedisTransport transport = new RedisTransport("redis://" + server.getHost() + ":" + server.getBindPort(), topic);
assertNotNull(transport);
-
- transport.close();
- server.stop();
-
}
-}
\ No newline at end of file
+}
diff --git a/pom.xml b/pom.xml
index 057cf73..a2efaf0 100644
--- a/pom.xml
+++ b/pom.xml
@@ -13,6 +13,21 @@
Meteor is a lightweight, fast and easy to use Java RPC library.
https://github.com/pixelib/Meteor
+
+ meteor-common
+ meteor-core
+ examples
+ meteor-jedis
+ benchmarks
+
+
+
+ 0.0.1-localbuild
+ 21
+ 21
+ UTF-8
+
+
duckelekuuk
@@ -25,7 +40,7 @@
mindgamesnl
Mats
Mindgamesnl
- https://github.com/Duckelekuuk
+ https://github.com/Mindgamesnl
mats@toetmats.nl
@@ -36,21 +51,6 @@
https://github.com/pixelib/Meteor
-
- meteor-common
- meteor-core
- examples
- meteor-jedis
- benchmarks
-
-
-
- 1.0.2-localbuild
- 17
- 17
- UTF-8
-
-
GNU General Public License v3.0
@@ -70,101 +70,58 @@
-
+
+ com.google.code.gson
+ gson
+ 2.14.0
+
org.junit.jupiter
junit-jupiter
- 5.10.0
+ 6.1.1
+ test
org.mockito
mockito-junit-jupiter
- 4.8.1
+ 5.23.0
test
-
-
- redis.clients
- jedis
- 5.0.2
-
-
- io.netty
- netty-buffer
- 4.1.97.Final
-
-
- com.google.code.gson
- gson
- 2.10.1
-
+
+
+ org.projectlombok
+ lombok
+ 1.18.46
+ provided
+
+
-
-
-
- org.apache.maven.plugins
- maven-surefire-plugin
- 3.1.2
-
-
- org.apache.maven.plugins
- maven-compiler-plugin
- 3.11.0
-
-
- org.apache.maven.plugins
- maven-resources-plugin
- 3.3.1
-
-
- org.apache.maven.plugins
- maven-shade-plugin
- 3.5.0
-
-
-
- org.jacoco
- jacoco-maven-plugin
- 0.8.9
-
-
-
- org.apache.maven.plugins
- maven-javadoc-plugin
- 3.5.0
-
-
-
- org.apache.maven.plugins
- maven-source-plugin
- 3.3.0
-
-
-
- org.apache.maven.plugins
- maven-deploy-plugin
- 3.1.1
-
-
-
- org.codehaus.mojo
- flatten-maven-plugin
- 1.5.0
-
-
-
-
+
+ org.apache.maven.plugins
+ maven-compiler-plugin
+ 3.15.0
+
+
+
+ org.projectlombok
+ lombok
+ 1.18.46
+
+
+
+
org.jacoco
jacoco-maven-plugin
+ 0.8.15
@@ -183,6 +140,7 @@
maven-javadoc-plugin
+ 3.12.0
all,-missing
@@ -198,6 +156,7 @@
maven-source-plugin
+ 3.4.0
true
@@ -214,6 +173,7 @@
org.codehaus.mojo
flatten-maven-plugin
+ 1.7.3
true
resolveCiFriendliesOnly
@@ -244,10 +204,22 @@
release
+
+ org.sonatype.central
+ central-publishing-maven-plugin
+ 0.11.0
+ true
+
+ central
+ true
+ true
+
+
+ org.apache.maven.plugins
maven-gpg-plugin
- 3.1.0
+ 3.2.8
--pinentry-mode