diff --git a/.idea/compiler.xml b/.idea/compiler.xml
index 293159c..6bfebb7 100644
--- a/.idea/compiler.xml
+++ b/.idea/compiler.xml
@@ -6,24 +6,20 @@
-
-
-
-
-
-
-
-
-
-
-
+
+
+
+
+
+
+
\ No newline at end of file
diff --git a/.idea/encodings.xml b/.idea/encodings.xml
index 774fbaa..4fd1572 100644
--- a/.idea/encodings.xml
+++ b/.idea/encodings.xml
@@ -2,16 +2,27 @@
+
+
+
+
+
+
+
+
+
+
+
diff --git a/.idea/misc.xml b/.idea/misc.xml
index cc819e8..fe2e3e3 100644
--- a/.idea/misc.xml
+++ b/.idea/misc.xml
@@ -10,6 +10,8 @@
diff --git a/messenger-benchmarks/pom.xml b/messenger-benchmarks/pom.xml
new file mode 100644
index 0000000..dcbb223
--- /dev/null
+++ b/messenger-benchmarks/pom.xml
@@ -0,0 +1,83 @@
+
+
+
+ messenger
+ io.github.sadcenter
+ 3.0.0
+
+ 4.0.0
+
+ messenger-benchmarks
+
+
+
+ org.openjdk.jmh
+ jmh-core
+ 1.35
+ compile
+
+
+ io.github.sadcenter
+ messenger-codecs
+ 3.0.0
+ compile
+
+
+ org.openjdk.jmh
+ jmh-generator-annprocess
+ 1.35
+ compile
+
+
+ io.github.sadcenter
+ messenger-core
+ 3.0.0
+ compile
+
+
+ de.ruedigermoeller
+ fst
+ 3.0.3
+ compile
+
+
+ org.msgpack
+ jackson-dataformat-msgpack
+ 0.9.1
+ compile
+
+
+ com.fasterxml.jackson.core
+ jackson-databind
+ 2.13.3
+ compile
+
+
+ com.fasterxml.jackson.core
+ jackson-core
+ 2.13.3
+ compile
+
+
+
+
+
+
+ org.apache.maven.plugins
+ maven-surefire-plugin
+
+ --add-opens=java.base/java.lang=ALL-UNNAMED
+ --add-opens=java.base/java.math=ALL-UNNAMED
+ --add-opens=java.base/java.util=ALL-UNNAMED
+ --add-opens=java.base/java.util.concurrent=ALL-UNNAMED
+ --add-opens=java.base/java.net=ALL-UNNAMED
+ --add-opens=java.base/java.text=ALL-UNNAMED
+ --add-opens=java.sql/java.sql=ALL-UNNAMED
+
+
+
+
+
+
\ No newline at end of file
diff --git a/messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/BenchmarkApplication.java b/messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/BenchmarkApplication.java
new file mode 100644
index 0000000..5c7e5e6
--- /dev/null
+++ b/messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/BenchmarkApplication.java
@@ -0,0 +1,11 @@
+package io.github.sadcenter.messenger;
+
+import java.io.IOException;
+import org.openjdk.jmh.Main;
+
+public class BenchmarkApplication {
+
+ public static void main(String[] args) throws IOException {
+ Main.main(args);
+ }
+}
diff --git a/messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/benchmarks/CodecBenchmarkTest.java b/messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/benchmarks/CodecBenchmarkTest.java
new file mode 100644
index 0000000..b8bd387
--- /dev/null
+++ b/messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/benchmarks/CodecBenchmarkTest.java
@@ -0,0 +1,40 @@
+package io.github.sadcenter.messenger.benchmarks;
+
+import io.github.sadcenter.messenger.benchmarks.example.ExamplePacket;
+import io.github.sadcenter.messenger.codec.Codec;
+import io.github.sadcenter.messenger.codecs.binary.FstCodec;
+import io.github.sadcenter.messenger.codecs.binary.MessagePackCodec;
+import io.github.sadcenter.messenger.codecs.json.JacksonCodec;
+import java.lang.ref.PhantomReference;
+import java.util.List;
+import org.openjdk.jmh.annotations.Benchmark;
+import org.openjdk.jmh.annotations.BenchmarkMode;
+import org.openjdk.jmh.annotations.Fork;
+import org.openjdk.jmh.annotations.Mode;
+import org.openjdk.jmh.annotations.Scope;
+import org.openjdk.jmh.annotations.State;
+
+@State(Scope.Benchmark)
+public class CodecBenchmarkTest {
+
+ private final ExamplePacket examplePacket = new ExamplePacket("Kacper", "Krzychala", 14);
+
+ private final JacksonCodec jacksonCodec = new JacksonCodec();
+ private final MessagePackCodec msgPackCodec = new MessagePackCodec();
+
+
+ @Benchmark
+ @BenchmarkMode(Mode.Throughput)
+ @Fork(value = 1)
+ public void benchJacksonSerialization() {
+ jacksonCodec.encode(examplePacket);
+ }
+
+ @Benchmark
+ @BenchmarkMode(Mode.Throughput)
+ @Fork(value = 1)
+ public void benchMessagePackSerialization() {
+ msgPackCodec.encode(examplePacket);
+ }
+
+}
diff --git a/messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/benchmarks/MessengerBenchmarkTest.java b/messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/benchmarks/MessengerBenchmarkTest.java
new file mode 100644
index 0000000..2f03036
--- /dev/null
+++ b/messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/benchmarks/MessengerBenchmarkTest.java
@@ -0,0 +1,5 @@
+package io.github.sadcenter.messenger.benchmarks;
+
+public class MessengerBenchmarkTest {
+
+}
diff --git a/messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/benchmarks/example/ExamplePacket.java b/messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/benchmarks/example/ExamplePacket.java
new file mode 100644
index 0000000..f9878c0
--- /dev/null
+++ b/messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/benchmarks/example/ExamplePacket.java
@@ -0,0 +1,43 @@
+package io.github.sadcenter.messenger.benchmarks.example;
+
+import io.github.sadcenter.messenger.packet.Packet;
+import java.util.Objects;
+
+public class ExamplePacket extends Packet {
+
+ private final String name, surname;
+ private final int age;
+
+ public ExamplePacket(String name, String surname, int age) {
+ this.name = name;
+ this.surname = surname;
+ this.age = age;
+ }
+
+ @Override
+ public String toString() {
+ return "ExamplePacket{" +
+ "name='" + name + '\'' +
+ ", surname='" + surname + '\'' +
+ ", age=" + age +
+ '}';
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) {
+ return true;
+ }
+ if (o == null || getClass() != o.getClass()) {
+ return false;
+ }
+ ExamplePacket that = (ExamplePacket) o;
+ return age == that.age && Objects.equals(name, that.name) && Objects.equals(
+ surname, that.surname);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(name, surname, age);
+ }
+}
diff --git a/messenger-codec-fst/pom.xml b/messenger-codec-fst/pom.xml
deleted file mode 100644
index 36aff9a..0000000
--- a/messenger-codec-fst/pom.xml
+++ /dev/null
@@ -1,74 +0,0 @@
-
-
-
- messenger
- io.github.sadcenter
- 2.1.6
-
- 4.0.0
-
- messenger-codec-fst
-
-
- UTF-8
-
-
-
-
- de.ruedigermoeller
- fst
- 2.56
- compile
-
-
- io.github.sadcenter
- messenger-core
- 2.1.6
- provided
-
-
-
-
-
-
- org.apache.maven.plugins
- maven-compiler-plugin
- 3.10.1
-
- 1.8
- 1.8
-
-
-
- org.apache.maven.plugins
- maven-source-plugin
- 3.2.1
-
-
- org.apache.maven.plugins
- maven-shade-plugin
- 3.2.4
-
-
- package
-
- shade
-
-
-
-
- ${project.build.directory}/dependency-reduced-pom.xml
-
-
-
-
-
- src/main/resources
- true
-
-
-
-
-
\ No newline at end of file
diff --git a/messenger-codec-msgpack/pom.xml b/messenger-codec-msgpack/pom.xml
deleted file mode 100644
index 3009cc2..0000000
--- a/messenger-codec-msgpack/pom.xml
+++ /dev/null
@@ -1,81 +0,0 @@
-
-
-
- messenger
- io.github.sadcenter
- 2.1.6
-
- 4.0.0
-
- messenger-codec-msgpack
-
-
- UTF-8
-
-
-
-
- org.msgpack
- jackson-dataformat-msgpack
- 0.9.0
- compile
-
-
- com.fasterxml.jackson.core
- jackson-databind
- 2.13.2.2
- provided
-
-
- com.fasterxml.jackson.core
- jackson-core
- 2.13.2
- provided
-
-
- io.github.sadcenter
- messenger-core
- 2.1.6
- provided
-
-
-
-
-
-
- org.apache.maven.plugins
- maven-compiler-plugin
- 3.10.1
-
- 11
- 11
-
-
-
- org.apache.maven.plugins
- maven-shade-plugin
- 3.2.4
-
-
- package
-
- shade
-
-
-
-
- ${project.build.directory}/dependency-reduced-pom.xml
-
-
-
-
-
- src/main/resources
- true
-
-
-
-
-
\ No newline at end of file
diff --git a/messenger-codecs/pom.xml b/messenger-codecs/pom.xml
new file mode 100644
index 0000000..a071c41
--- /dev/null
+++ b/messenger-codecs/pom.xml
@@ -0,0 +1,59 @@
+
+
+
+ messenger
+ io.github.sadcenter
+ 3.0.0
+
+ 4.0.0
+
+ messenger-codecs
+
+
+
+ de.ruedigermoeller
+ fst
+ 3.0.3
+ provided
+
+
+ org.msgpack
+ jackson-dataformat-msgpack
+ 0.9.1
+ provided
+
+
+ com.fasterxml.jackson.core
+ jackson-databind
+ 2.13.3
+ provided
+
+
+ com.fasterxml.jackson.core
+ jackson-core
+ 2.13.3
+ provided
+
+
+ io.github.sadcenter
+ messenger-core
+ 3.0.0
+ provided
+
+
+ org.junit.jupiter
+ junit-jupiter-engine
+ 5.8.2
+ test
+
+
+ org.junit.jupiter
+ junit-jupiter-api
+ 5.8.2
+ test
+
+
+
+
\ No newline at end of file
diff --git a/messenger-codec-fst/src/main/java/io/github/sadcenter/messenger/codec/FstCodec.java b/messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary/FstCodec.java
similarity index 55%
rename from messenger-codec-fst/src/main/java/io/github/sadcenter/messenger/codec/FstCodec.java
rename to messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary/FstCodec.java
index ea4a8d6..776f1a3 100644
--- a/messenger-codec-fst/src/main/java/io/github/sadcenter/messenger/codec/FstCodec.java
+++ b/messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary/FstCodec.java
@@ -1,11 +1,20 @@
-package io.github.sadcenter.messenger.codec;
+package io.github.sadcenter.messenger.codecs.binary;
+import io.github.sadcenter.messenger.codec.Codec;
import org.nustaq.serialization.FSTConfiguration;
@SuppressWarnings("unchecked")
public class FstCodec implements Codec {
- private final FSTConfiguration fstConfiguration = FSTConfiguration.createDefaultConfiguration();
+ private final FSTConfiguration fstConfiguration;
+
+ public FstCodec() {
+ this.fstConfiguration = FSTConfiguration.createDefaultConfiguration();
+ }
+
+ public FstCodec(FSTConfiguration configuration) {
+ this.fstConfiguration = configuration;
+ }
@Override
public byte[] encode(Object data) {
diff --git a/messenger-codec-msgpack/src/main/java/io/github/sadcenter/messenger/codec/MessagePackCodec.java b/messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary/MessagePackCodec.java
similarity index 86%
rename from messenger-codec-msgpack/src/main/java/io/github/sadcenter/messenger/codec/MessagePackCodec.java
rename to messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary/MessagePackCodec.java
index d9f1288..d190abe 100644
--- a/messenger-codec-msgpack/src/main/java/io/github/sadcenter/messenger/codec/MessagePackCodec.java
+++ b/messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary/MessagePackCodec.java
@@ -1,18 +1,20 @@
-package io.github.sadcenter.messenger.codec;
+package io.github.sadcenter.messenger.codecs.binary;
import com.fasterxml.jackson.annotation.JsonAutoDetect.Visibility;
import com.fasterxml.jackson.annotation.JsonInclude.Include;
import com.fasterxml.jackson.core.JsonGenerator;
+import com.fasterxml.jackson.core.JsonpCharacterEscapes;
import com.fasterxml.jackson.databind.DeserializationFeature;
import com.fasterxml.jackson.databind.MapperFeature;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.SerializationFeature;
+import io.github.sadcenter.messenger.codec.Codec;
import io.github.sadcenter.messenger.packet.Packet;
import io.github.sadcenter.messenger.packet.PacketRequest;
import java.util.Arrays;
import org.msgpack.jackson.dataformat.MessagePackFactory;
-@SuppressWarnings({"deprecation", "unchecked"})
+@SuppressWarnings({"unchecked"})
public class MessagePackCodec implements Codec {
private final ObjectMapper objectMapper;
@@ -43,7 +45,7 @@ public MessagePackCodec(ObjectMapper objectMapper) {
}
@Override
- public T decode(byte[] bytes, Class clazz) {
+ public R decode(byte[] bytes, Class clazz) {
if (bytes == null) {
return null;
}
@@ -58,8 +60,8 @@ public T decode(byte[] bytes, Class clazz) {
}
@Override
- public T decode(byte[] bytes) {
- return (T) this.decode(bytes, Packet.class);
+ public R decode(byte[] bytes) {
+ return (R) this.decode(bytes, Packet.class);
}
@Override
diff --git a/messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/json/JacksonCodec.java b/messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/json/JacksonCodec.java
new file mode 100644
index 0000000..530e9b8
--- /dev/null
+++ b/messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/json/JacksonCodec.java
@@ -0,0 +1,78 @@
+package io.github.sadcenter.messenger.codecs.json;
+
+import com.fasterxml.jackson.annotation.JsonAutoDetect.Visibility;
+import com.fasterxml.jackson.annotation.JsonInclude.Include;
+import com.fasterxml.jackson.core.JsonGenerator;
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.DeserializationFeature;
+import com.fasterxml.jackson.databind.MapperFeature;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.databind.SerializationFeature;
+import io.github.sadcenter.messenger.codec.Codec;
+import io.github.sadcenter.messenger.packet.Packet;
+import io.github.sadcenter.messenger.packet.PacketRequest;
+import java.nio.charset.StandardCharsets;
+
+@SuppressWarnings("unchecked")
+public class JacksonCodec implements Codec {
+
+ private final ObjectMapper objectMapper;
+
+ public JacksonCodec() {
+ this.objectMapper = new ObjectMapper();
+
+ this.objectMapper.setSerializationInclusion(Include.NON_NULL);
+ this.objectMapper.setVisibility(
+ this.objectMapper
+ .getSerializationConfig()
+ .getDefaultVisibilityChecker()
+ .withFieldVisibility(Visibility.ANY)
+ .withGetterVisibility(Visibility.NONE)
+ .withSetterVisibility(Visibility.NONE)
+ .withCreatorVisibility(Visibility.NONE));
+ this.objectMapper
+ .disable(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES)
+ .enable(JsonGenerator.Feature.WRITE_BIGDECIMAL_AS_PLAIN)
+ .disable(SerializationFeature.FAIL_ON_EMPTY_BEANS)
+ .enable(MapperFeature.SORT_PROPERTIES_ALPHABETICALLY);
+ this.objectMapper.registerSubtypes(Packet.class);
+ this.objectMapper.registerSubtypes(PacketRequest.class);
+ }
+
+ public JacksonCodec(ObjectMapper objectMapper) {
+ this.objectMapper = objectMapper;
+ }
+
+ @Override
+ public byte[] encode(Object data) {
+ try {
+ return this.objectMapper.writeValueAsString(data).getBytes(StandardCharsets.UTF_8);
+ } catch (JsonProcessingException e) {
+ e.printStackTrace();
+ }
+
+ return new byte[0];
+ }
+
+ @Override
+ public R decode(byte[] data, Class type) {
+ try {
+ return this.objectMapper.readValue(new String(data), type);
+ } catch (JsonProcessingException e) {
+ e.printStackTrace();
+ }
+
+ return null;
+ }
+
+ @Override
+ public R decode(byte[] data) {
+ try {
+ return (R) this.objectMapper.readValue(new String(data), Packet.class);
+ } catch (JsonProcessingException e) {
+ e.printStackTrace();
+ }
+
+ return null;
+ }
+}
diff --git a/messenger-codecs/src/test/java/io/github/sadcenter/messenger/codecs/CodecsTest.java b/messenger-codecs/src/test/java/io/github/sadcenter/messenger/codecs/CodecsTest.java
new file mode 100644
index 0000000..ddda936
--- /dev/null
+++ b/messenger-codecs/src/test/java/io/github/sadcenter/messenger/codecs/CodecsTest.java
@@ -0,0 +1,80 @@
+package io.github.sadcenter.messenger.codecs;
+
+import com.fasterxml.jackson.annotation.JsonCreator;
+import com.fasterxml.jackson.annotation.JsonProperty;
+import io.github.sadcenter.messenger.codecs.binary.FstCodec;
+import io.github.sadcenter.messenger.codecs.binary.MessagePackCodec;
+import io.github.sadcenter.messenger.codecs.json.JacksonCodec;
+import io.github.sadcenter.messenger.packet.Packet;
+import java.util.Objects;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+class CodecsTest {
+
+ private static final String NOT_MATCH_MESSAGE = "The results does not match expected object.";
+ private static final TestPacket EXPECTED = new TestPacket("Kacper", "Krzychala");
+
+ @Test
+ void testFst() {
+ FstCodec codec = new FstCodec();
+
+ byte[] encode = codec.encode(EXPECTED);
+ assertEquals(EXPECTED, codec.decode(encode), NOT_MATCH_MESSAGE);
+ }
+
+ @Test
+ void testJackson() {
+ JacksonCodec codec = new JacksonCodec();
+
+ byte[] encode = codec.encode(EXPECTED);
+ assertEquals(EXPECTED, codec.decode(encode), NOT_MATCH_MESSAGE);
+ }
+
+ @Test
+ void testMsgPack() {
+ MessagePackCodec codec = new MessagePackCodec();
+
+ byte[] encode = codec.encode(EXPECTED);
+ assertEquals(EXPECTED, codec.decode(encode), NOT_MATCH_MESSAGE);
+ }
+
+ private static final class TestPacket extends Packet {
+
+ private final String name, surname;
+
+ @JsonCreator
+ private TestPacket(@JsonProperty("name") String name, @JsonProperty("surname") String surname) {
+ this.name = name;
+ this.surname = surname;
+ }
+
+ public String getName() {
+ return name;
+ }
+
+ public String getSurname() {
+ return surname;
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) {
+ return true;
+ }
+ if (o == null || getClass() != o.getClass()) {
+ return false;
+ }
+ TestPacket that = (TestPacket) o;
+ return Objects.equals(name, that.name) && Objects.equals(surname,
+ that.surname);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(name, surname);
+ }
+ }
+
+}
diff --git a/messenger-core/pom.xml b/messenger-core/pom.xml
index 84423cf..6541d9c 100644
--- a/messenger-core/pom.xml
+++ b/messenger-core/pom.xml
@@ -8,7 +8,7 @@
messenger
io.github.sadcenter
- 2.1.6
+ 3.0.0
@@ -19,56 +19,15 @@
com.fasterxml.jackson.core
jackson-databind
- 2.13.2.2
+ 2.13.3
compile
com.fasterxml.jackson.core
jackson-core
- 2.13.2
+ 2.13.3
compile
-
-
-
- org.apache.maven.plugins
- maven-compiler-plugin
- 3.10.1
-
- 11
- 11
-
-
-
- org.apache.maven.plugins
- maven-source-plugin
- 3.2.1
-
-
- org.apache.maven.plugins
- maven-shade-plugin
- 3.2.4
-
-
- package
-
- shade
-
-
-
-
- ${project.build.directory}/dependency-reduced-pom.xml
-
-
-
-
-
- src/main/resources
- true
-
-
-
-
\ No newline at end of file
diff --git a/messenger-core/src/main/java/io/github/sadcenter/messenger/channel/Channel.java b/messenger-core/src/main/java/io/github/sadcenter/messenger/channel/Channel.java
index c5e63ac..a9893b0 100644
--- a/messenger-core/src/main/java/io/github/sadcenter/messenger/channel/Channel.java
+++ b/messenger-core/src/main/java/io/github/sadcenter/messenger/channel/Channel.java
@@ -19,21 +19,21 @@ public Channel(String to, String from, String replyTo) {
* @return Channel where packet was sent
*/
public String getChannel() {
- return channel;
+ return this.channel;
}
/**
* @return Channel where packet was sent from.
*/
public String getFrom() {
- return from;
+ return this.from;
}
/**
* @return Channel where reply should be sent.
*/
public String getReplyTo() {
- return replyTo;
+ return this.replyTo;
}
@Contract(value = " -> new", pure = true)
diff --git a/messenger-core/src/main/java/io/github/sadcenter/messenger/packet/Packet.java b/messenger-core/src/main/java/io/github/sadcenter/messenger/packet/Packet.java
index 3925f0e..0aa01f4 100644
--- a/messenger-core/src/main/java/io/github/sadcenter/messenger/packet/Packet.java
+++ b/messenger-core/src/main/java/io/github/sadcenter/messenger/packet/Packet.java
@@ -6,6 +6,9 @@
@JsonTypeInfo(use = JsonTypeInfo.Id.CLASS)
public class Packet implements Serializable {
- //Where packet is sent from
+ /**
+ * Where packet is sent from.
+ * May be set to "unknown" if server didn't set up this future.
+ */
public String from;
}
diff --git a/messenger-core/src/main/java/io/github/sadcenter/messenger/packet/PacketRequest.java b/messenger-core/src/main/java/io/github/sadcenter/messenger/packet/PacketRequest.java
index 6abe7b1..054a5b8 100644
--- a/messenger-core/src/main/java/io/github/sadcenter/messenger/packet/PacketRequest.java
+++ b/messenger-core/src/main/java/io/github/sadcenter/messenger/packet/PacketRequest.java
@@ -7,10 +7,8 @@
@JsonTypeInfo(use = JsonTypeInfo.Id.CLASS)
public class PacketRequest extends Packet {
+ //TODO: make error more customizable
private String errorMessage;
- private final Map attributes = new HashMap<>();
- // bruh XDD
- // TODO delete that ^^ (this is a temp fix)
public boolean isExceptionally() {
return this.errorMessage != null;
@@ -23,16 +21,4 @@ public String getErrorMessage() {
public void setErrorMessage(String errorMessage) {
this.errorMessage = errorMessage;
}
-
- public void setAttribute(String attribute, String value) {
- this.attributes.put(attribute, value);
- }
-
- public String getValue(String attribute) {
- return this.attributes.get(attribute);
- }
-
- public void deleteAttribute(String attribute) {
- this.attributes.remove(attribute);
- }
}
diff --git a/messenger-nats-luckperms/pom.xml b/messenger-nats-luckperms/pom.xml
index 922511f..5942018 100644
--- a/messenger-nats-luckperms/pom.xml
+++ b/messenger-nats-luckperms/pom.xml
@@ -5,22 +5,24 @@
messenger
io.github.sadcenter
- 2.1.6
+ 3.0.0
4.0.0
messenger-nats-luckperms
-
- UTF-8
-
-
+
+ io.github.sadcenter
+ messenger-core
+ 3.0.0
+ provided
+
io.github.sadcenter
messenger-nats
- 2.1.6
- compile
+ 3.0.0
+ provided
net.luckperms
@@ -31,38 +33,6 @@
-
-
- org.apache.maven.plugins
- maven-compiler-plugin
- 3.10.1
-
- 11
- 11
-
-
-
- org.apache.maven.plugins
- maven-source-plugin
- 3.2.1
-
-
- org.apache.maven.plugins
- maven-shade-plugin
- 3.2.4
-
-
- package
-
- shade
-
-
-
-
- ${project.build.directory}/dependency-reduced-pom.xml
-
-
-
src/main/resources
diff --git a/messenger-nats/pom.xml b/messenger-nats/pom.xml
index 67833b1..be2588a 100644
--- a/messenger-nats/pom.xml
+++ b/messenger-nats/pom.xml
@@ -5,22 +5,18 @@
messenger
io.github.sadcenter
- 2.1.6
+ 3.0.0
4.0.0
messenger-nats
-
- UTF-8
-
-
io.github.sadcenter
messenger-core
- 2.1.6
- compile
+ 3.0.0
+ provided
io.nats
@@ -28,17 +24,10 @@
2.14.0
compile
-
-
- io.github.sadcenter
- messenger-codec-msgpack
- 2.1.6
- test
-
io.github.sadcenter
- messenger-codec-fst
- 2.1.6
+ messenger-codecs
+ 3.0.0
test
@@ -65,67 +54,30 @@
1.1.2
test
+
+ de.ruedigermoeller
+ fst
+ 3.0.3
+ test
+
+
+ org.msgpack
+ jackson-dataformat-msgpack
+ 0.9.1
+ test
+
+
+ com.fasterxml.jackson.core
+ jackson-databind
+ 2.13.3
+ test
+
+
+ com.fasterxml.jackson.core
+ jackson-core
+ 2.13.3
+ test
+
-
-
-
- org.apache.maven.plugins
- maven-compiler-plugin
- 3.10.1
-
- 11
- 11
-
-
-
- org.apache.maven.plugins
- maven-source-plugin
- 3.2.1
-
-
- maven-surefire-plugin
- 2.19
-
-
- org.junit.platform
- junit-platform-surefire-provider
- 1.0.0
-
-
-
- --add-opens=java.base/java.lang=ALL-UNNAMED
- --add-opens=java.base/java.math=ALL-UNNAMED
- --add-opens=java.base/java.util=ALL-UNNAMED
- --add-opens=java.base/java.util.concurrent=ALL-UNNAMED
- --add-opens=java.base/java.net=ALL-UNNAMED
- --add-opens=java.base/java.text=ALL-UNNAMED
- --add-opens=java.sql/java.sql=ALL-UNNAMED
-
-
-
- org.apache.maven.plugins
- maven-shade-plugin
- 3.2.4
-
-
- package
-
- shade
-
-
-
-
- ${project.build.directory}/dependency-reduced-pom.xml
-
-
-
-
-
- src/main/resources
- true
-
-
-
-
\ No newline at end of file
diff --git a/messenger-nats/src/main/java/io/github/sadcenter/messenger/nats/NatsMessenger.java b/messenger-nats/src/main/java/io/github/sadcenter/messenger/nats/NatsMessenger.java
index 71e2a8a..e4a8fe9 100644
--- a/messenger-nats/src/main/java/io/github/sadcenter/messenger/nats/NatsMessenger.java
+++ b/messenger-nats/src/main/java/io/github/sadcenter/messenger/nats/NatsMessenger.java
@@ -16,7 +16,6 @@
import io.nats.client.api.KeyValueConfiguration;
import io.nats.client.api.StorageType;
import java.io.IOException;
-import java.time.Duration;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashSet;
@@ -151,7 +150,6 @@ private void createIfNotExists(String bucket)
keyValueManagement.create(
new KeyValueConfiguration.Builder()
.name(bucket)
- .ttl(Duration.ofMinutes(1))
.storageType(StorageType.Memory)
.build());
}
@@ -176,7 +174,8 @@ public void set(@NotNull String id, @NotNull String key, @NotNull T value) {
() -> {
try {
this.createIfNotExists(id);
- this.connection.keyValue(id).put(key, this.codec.encode(value));
+ // this.connection.keyValue(id).put(key, this.codec.encode(value));
+ // System.out.println(this.connection.keyValue(id).get(key).getValueAsString());
} catch (IOException | JetStreamApiException | InterruptedException e) {
e.printStackTrace();
}
@@ -265,7 +264,8 @@ public void codec(@NotNull T codec) {
this.codec = codec;
}
+ @Deprecated
public Connection getConnection() {
- return connection;
+ return this.connection;
}
}
diff --git a/messenger-nats/src/test/java/io/github/sadcenter/messenger/nats/tests/NatsMessengerTest.java b/messenger-nats/src/test/java/io/github/sadcenter/messenger/nats/tests/NatsMessengerTest.java
index dc427ac..06f45fd 100644
--- a/messenger-nats/src/test/java/io/github/sadcenter/messenger/nats/tests/NatsMessengerTest.java
+++ b/messenger-nats/src/test/java/io/github/sadcenter/messenger/nats/tests/NatsMessengerTest.java
@@ -5,17 +5,20 @@
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import com.fasterxml.jackson.annotation.JsonCreator;
+import com.fasterxml.jackson.annotation.JsonProperty;
import io.github.sadcenter.messenger.Messenger;
import io.github.sadcenter.messenger.channel.Channel;
-import io.github.sadcenter.messenger.codec.FstCodec;
-import io.github.sadcenter.messenger.codec.MessagePackCodec;
+import io.github.sadcenter.messenger.codec.Codec;
+import io.github.sadcenter.messenger.codecs.binary.FstCodec;
+import io.github.sadcenter.messenger.codecs.binary.MessagePackCodec;
+import io.github.sadcenter.messenger.codecs.json.JacksonCodec;
import io.github.sadcenter.messenger.handler.PacketHandler;
import io.github.sadcenter.messenger.nats.NatsMessenger;
-import io.github.sadcenter.messenger.packet.Packet;
import io.github.sadcenter.messenger.packet.PacketRequest;
-import io.nats.client.Options;
-import java.io.IOException;
+import io.nats.client.Options.Builder;
import java.time.Duration;
+import java.util.List;
import java.util.Objects;
import java.util.concurrent.CompletableFuture;
import np.com.madanpokharel.embed.nats.EmbeddedNatsConfig;
@@ -23,12 +26,21 @@
import np.com.madanpokharel.embed.nats.NatsServerConfig;
import np.com.madanpokharel.embed.nats.NatsVersion;
import np.com.madanpokharel.embed.nats.ServerType;
+import org.jetbrains.annotations.NotNull;
+import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.RepeatedTest;
import org.junit.jupiter.api.RepetitionInfo;
class NatsMessengerTest {
+ private static final List CODECS =
+ List.of(new JacksonCodec(), new FstCodec(), new MessagePackCodec());
+ private static final int CODEC_SIZE = 3;
+ private static Messenger MESSENGER;
+ private static EmbeddedNatsServer SERVER;
+
@BeforeAll
static void initAll() throws Exception {
EmbeddedNatsConfig config =
@@ -43,24 +55,33 @@ static void initAll() throws Exception {
.build())
.build();
- EmbeddedNatsServer natsServer = new EmbeddedNatsServer(config);
- natsServer.startServer();
- assertTrue(natsServer.isServerRunning(), "Server is not running");
+ SERVER = new EmbeddedNatsServer(config);
+ SERVER.startServer();
+ assertTrue(SERVER.isServerRunning(), "Server is not running (?)");
+ Thread.sleep(1L);
+ MESSENGER =
+ new NatsMessenger(
+ new Builder().server("nats://127.0.0.1:7656").maxReconnects(1000).build());
}
- @RepeatedTest(2)
- void testPublish(RepetitionInfo repetitionInfo) throws Exception {
- Messenger messenger =
- new NatsMessenger(
- new Options.Builder().server("nats://127.0.0.1:7656").build(),
- repetitionInfo.getCurrentRepetition() % 2 == 1
- ? new MessagePackCodec()
- : new FstCodec());
+ @AfterAll
+ static void close() throws InterruptedException {
+ Thread.sleep(1L);
+ MESSENGER.close();
+ SERVER.stopServer();
+ }
- messenger.listen(new TestHandler(messenger), TestPacket.class, "messenger:test");
+ @BeforeEach
+ void initCodec(RepetitionInfo repetitionInfo) {
+ MESSENGER.codec(CODECS.get(repetitionInfo.getCurrentRepetition() - 1));
+ }
+
+ @RepeatedTest(CODEC_SIZE)
+ void testPublish() {
+ MESSENGER.listen(new TestHandler(MESSENGER), TestPacket.class, "messenger:test");
TestPacket sentPacket = new TestPacket("Kacper", "Krzychała");
- CompletableFuture request = messenger.request("messenger:test", sentPacket);
+ CompletableFuture request = MESSENGER.request("messenger:test", sentPacket);
await()
.atMost(Duration.ofSeconds(2))
.until(
@@ -68,35 +89,25 @@ void testPublish(RepetitionInfo repetitionInfo) throws Exception {
request.isDone() && !request.isCancelled() && !request.isCompletedExceptionally());
request.thenAccept(
receivedPacket -> assertEquals(sentPacket, receivedPacket, "Packets are not the same"));
-
- messenger.close();
}
- @RepeatedTest(2)
- void testKeyValue(RepetitionInfo repetitionInfo) throws IOException, InterruptedException {
- Messenger messenger =
- new NatsMessenger(
- new Options.Builder().server("nats://127.0.0.1:7656").build(),
- repetitionInfo.getCurrentRepetition() % 2 == 1
- ? new MessagePackCodec()
- : new FstCodec());
+ @RepeatedTest(CODEC_SIZE)
+ void testKeyValue() {
TestPacket testPacket = new TestPacket("Kacper", "Krzychala");
- messenger.set("messenger", "test", testPacket);
- CompletableFuture future = messenger
- .get("messenger", "test");
- future
- .whenComplete(
- (packet, throwable) -> {
- assertNull(throwable, throwable.getMessage());
- assertEquals(testPacket, packet, "Packets are not the same");
- });
+ MESSENGER.set("messenger", "test", testPacket);
+ CompletableFuture future = MESSENGER.get("messenger", "test");
+ future.whenComplete(
+ (packet, throwable) -> {
+ assertNull(throwable);
+ assertEquals(testPacket, packet, "Packets are not the same");
+ });
}
private static final class TestHandler implements PacketHandler {
private final Messenger messenger;
- private TestHandler(Messenger messenger) {
+ private TestHandler(@NotNull Messenger messenger) {
this.messenger = messenger;
}
@@ -113,24 +124,14 @@ private static final class TestPacket extends PacketRequest {
private final String name;
private final String surname;
- public TestPacket(String name, String surname) {
+ @JsonCreator
+ public TestPacket(
+ @NotNull @JsonProperty("name") String name,
+ @NotNull @JsonProperty("surname") String surname) {
this.name = name;
this.surname = surname;
}
- // Jackson requires an empty constructor
- public TestPacket() {
- this(null, null);
- }
-
- public String getName() {
- return name;
- }
-
- public String getSurname() {
- return surname;
- }
-
@Override
public boolean equals(Object o) {
if (this == o) {
diff --git a/messenger-redis/pom.xml b/messenger-redis/pom.xml
index a4605de..82c7ac6 100644
--- a/messenger-redis/pom.xml
+++ b/messenger-redis/pom.xml
@@ -5,7 +5,7 @@
messenger
io.github.sadcenter
- 2.1.6
+ 3.0.0
4.0.0
@@ -31,26 +31,19 @@
com.github.ben-manes.caffeine
caffeine
- 3.0.6
+ 3.1.0
compile
io.github.sadcenter
messenger-core
- 2.1.6
- compile
-
-
-
- io.github.sadcenter
- messenger-codec-msgpack
- 2.1.6
- test
+ 3.0.0
+ provided
io.github.sadcenter
- messenger-codec-fst
- 2.1.6
+ messenger-codecs
+ 3.0.0
test
@@ -74,64 +67,42 @@
com.github.fppt
jedis-mock
- 1.0.1
+ 1.0.2
+ test
+
+
+ de.ruedigermoeller
+ fst
+ 3.0.3
+ test
+
+
+ org.msgpack
+ jackson-dataformat-msgpack
+ 0.9.1
+ test
+
+
+ com.google.code.gson
+ gson
+ 2.9.0
+ test
+
+
+ com.fasterxml.jackson.core
+ jackson-databind
+ 2.13.3
+ test
+
+
+ com.fasterxml.jackson.core
+ jackson-core
+ 2.13.3
test
-
-
- org.apache.maven.plugins
- maven-compiler-plugin
- 3.10.1
-
- 11
- 11
-
-
-
- org.apache.maven.plugins
- maven-source-plugin
- 3.2.1
-
-
- org.apache.maven.plugins
- maven-shade-plugin
- 3.2.4
-
-
- package
-
- shade
-
-
-
-
- ${project.build.directory}/dependency-reduced-pom.xml
-
-
-
- maven-surefire-plugin
- 2.19
-
-
- org.junit.platform
- junit-platform-surefire-provider
- 1.0.0
-
-
-
- --add-opens=java.base/java.lang=ALL-UNNAMED
- --add-opens=java.base/java.math=ALL-UNNAMED
- --add-opens=java.base/java.util=ALL-UNNAMED
- --add-opens=java.base/java.util.concurrent=ALL-UNNAMED
- --add-opens=java.base/java.net=ALL-UNNAMED
- --add-opens=java.base/java.text=ALL-UNNAMED
- --add-opens=java.sql/java.sql=ALL-UNNAMED
-
-
-
src/main/resources
diff --git a/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/RedisMessenger.java b/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/RedisMessenger.java
index ee06cef..fe74ca2 100644
--- a/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/RedisMessenger.java
+++ b/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/RedisMessenger.java
@@ -14,7 +14,6 @@
import io.github.sadcenter.messenger.redis.codec.RedisByteCodec;
import io.github.sadcenter.messenger.redis.codec.StringByteArrayCodec;
import io.github.sadcenter.messenger.redis.packet.RedisPacketHandler;
-import io.github.sadcenter.messenger.redis.request.RedisRequestHandler;
import io.lettuce.core.RedisClient;
import io.lettuce.core.RedisURI;
import io.lettuce.core.api.StatefulRedisConnection;
@@ -27,6 +26,8 @@
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.Executors;
+import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import org.apache.commons.pool2.impl.GenericObjectPool;
@@ -36,20 +37,8 @@
@SuppressWarnings("all")
public class RedisMessenger implements Messenger {
- private final Cache> requestCache =
- Caffeine.newBuilder()
- .expireAfterWrite(30, TimeUnit.SECONDS)
- .removalListener(
- (RemovalListener>)
- (key, callbackFuture, removalCause) -> {
- if (callbackFuture != null && removalCause == RemovalCause.EXPIRED) {
- callbackFuture.completeExceptionally(new TimeoutException());
- }
- })
- .build();
- private final Cache idCache = Caffeine.newBuilder()
- .expireAfterWrite(60, TimeUnit.SECONDS)
- .build();
+ private static final ScheduledExecutorService WORKER = Executors.newSingleThreadScheduledExecutor();
+
private final Set subscribedChannels = new HashSet<>();
private final RedisClient redisClient;
private final GenericObjectPool> connectionPool;
@@ -60,29 +49,27 @@ public class RedisMessenger implements Messenger {
private String from;
public RedisMessenger(@NotNull RedisURI redisURI, @NotNull Codec codec) {
- this.codec(codec);
redisURI.setDatabase(0);
+
+ this.codec(codec);
this.redisClient = RedisClient.create(redisURI);
GenericObjectPoolConfig> poolConfig =
new GenericObjectPoolConfig<>();
- poolConfig.setMinIdle(50); //TODO tweak to optimal
- poolConfig.setMaxIdle(200); //TODO tweak to optimal
- poolConfig.setMaxTotal(200); //TODO tweak to optimal
+ poolConfig.setMinIdle(50); // TODO tweak to optimal
+ poolConfig.setMaxIdle(200); // TODO tweak to optimal
+ poolConfig.setMaxTotal(200); // TODO tweak to optimal
RedisByteCodec redisCodec = new RedisByteCodec<>(this.codec);
- this.subConnection = redisClient.connectPubSub(redisCodec, redisURI);
+ this.subConnection = this.redisClient.connectPubSub(redisCodec, redisURI);
this.connectionPool =
ConnectionPoolSupport.createGenericObjectPool(
- () -> redisClient.connect(redisCodec, redisURI), poolConfig);
- this.keyConnection = redisClient.connect(new StringByteArrayCodec());
-
- initialize();
+ () -> this.redisClient.connect(redisCodec, redisURI), poolConfig);
+ this.keyConnection = this.redisClient.connect(new StringByteArrayCodec());
}
@Override
public void initialize() {
- this.listen(new RedisRequestHandler(this.requestCache), PacketRequest.class, "messenger:temps");
}
@Override
@@ -94,7 +81,7 @@ public void close() {
public void publish(@NotNull String channel, @NotNull Packet packet) {
packet.from = this.from;
- try (StatefulRedisConnection connection = connectionPool.borrowObject()) {
+ try (StatefulRedisConnection connection = this.connectionPool.borrowObject()) {
connection.async().publish(channel, packet);
} catch (Exception e) {
e.printStackTrace();
@@ -110,28 +97,51 @@ public void reply(@NotNull Channel channel, @NotNull PacketRequest packet) {
public CompletableFuture request(
@NotNull String channel, @NotNull PacketRequest packet) {
CompletableFuture future = new CompletableFuture<>();
- try (StatefulRedisConnection connection = connectionPool.borrowObject()) {
- UUID uuid = UUID.randomUUID();
- packet.setAttribute("uuid", uuid.toString());
+ try (StatefulRedisConnection connection = this.connectionPool.borrowObject()) {
+ String id = UUID.randomUUID().toString();
+ listenToRequest(id, packet, future);
connection.async().publish(channel, packet);
- this.requestCache.put(uuid, future);
} catch (Exception e) {
e.printStackTrace();
future.completeExceptionally(e);
}
-
return future;
}
+ private void listenToRequest(@NotNull String uniqueId, @NotNull PacketRequest packet, @NotNull CompletableFuture future) {
+ this.subConnection.addListener(new RedisPacketHandler(
+ new PacketHandler() {
+ @Override
+ public void handle(Channel channel, PacketRequest packet) {
+ future.complete(packet);
+ }
+ }, Set.of(uniqueId), packet.getClass()));
+
+ if (!this.subscribedChannels.contains(uniqueId)) {
+ this.subConnection.async().subscribe(uniqueId);
+ this.subscribedChannels.add(uniqueId);
+ }
+
+ WORKER.schedule(() -> {
+ if(!future.isDone()) {
+ if (this.subscribedChannels.contains(uniqueId)) {
+ this.subConnection.async().unsubscribe(uniqueId);
+ this.subscribedChannels.remove(uniqueId);
+ }
+ }
+ }, 30L, TimeUnit.SECONDS);
+
+ }
+
@Override
public void listen(
@NotNull Set channels, @NotNull PacketHandler handler, @NotNull Class type) {
- subConnection.addListener(new RedisPacketHandler<>(handler, channels, type));
+ this.subConnection.addListener(new RedisPacketHandler<>(handler, channels, type));
for (String channel : channels) {
- if (!subscribedChannels.contains(channel)) {
- subConnection.sync().subscribe(channel);
- subscribedChannels.add(channel);
+ if (!this.subscribedChannels.contains(channel)) {
+ this.subConnection.sync().subscribe(channel);
+ this.subscribedChannels.add(channel);
}
}
}
@@ -162,7 +172,7 @@ public void listen(PacketHandlerAbstract>... abstractHandlers) {
@Override
public void set(@NotNull String id, @NotNull String key, @NotNull String value) {
- keyConnection.async().set(id + "_" + key, value.getBytes(StandardCharsets.UTF_8));
+ this.keyConnection.async().set(id + "_" + key, value.getBytes(StandardCharsets.UTF_8));
}
@Override
@@ -225,6 +235,6 @@ public void codec(@NotNull T codec) {
}
public String from() {
- return from;
+ return this.from;
}
}
diff --git a/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/codec/RedisByteCodec.java b/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/codec/RedisByteCodec.java
index e0d4a18..f6f330f 100644
--- a/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/codec/RedisByteCodec.java
+++ b/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/codec/RedisByteCodec.java
@@ -36,8 +36,4 @@ public ByteBuffer encodeValue(V value) {
byte[] encode = this.codec.encode(value);
return ByteBuffer.wrap(encode);
}
-
- public Codec getCodec() {
- return codec;
- }
}
diff --git a/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/codec/StringByteArrayCodec.java b/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/codec/StringByteArrayCodec.java
index ad7adad..ad5fa91 100644
--- a/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/codec/StringByteArrayCodec.java
+++ b/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/codec/StringByteArrayCodec.java
@@ -15,7 +15,6 @@ public String decodeKey(ByteBuffer bytes) {
public byte[] decodeValue(ByteBuffer bytes) {
byte[] array = new byte[bytes.remaining()];
bytes.get(array);
-
return array;
}
diff --git a/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/request/RedisRequestHandler.java b/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/request/RedisRequestHandler.java
deleted file mode 100644
index b91e968..0000000
--- a/messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/request/RedisRequestHandler.java
+++ /dev/null
@@ -1,31 +0,0 @@
-package io.github.sadcenter.messenger.redis.request;
-
-import com.github.benmanes.caffeine.cache.Cache;
-import io.github.sadcenter.messenger.channel.Channel;
-import io.github.sadcenter.messenger.handler.PacketHandler;
-import io.github.sadcenter.messenger.packet.PacketRequest;
-import io.lettuce.core.pubsub.RedisPubSubAdapter;
-import java.util.UUID;
-import java.util.concurrent.CompletableFuture;
-
-public class RedisRequestHandler implements PacketHandler {
-
- private final Cache> requestCache;
-
- public RedisRequestHandler(Cache> requestCache) {
- this.requestCache = requestCache;
- }
-
- @Override
- public void handle(Channel channel, PacketRequest packet) {
- UUID uuid = UUID.fromString(packet.getValue("uuid"));
- CompletableFuture requestFuture =
- this.requestCache.getIfPresent(uuid);
-
- if (requestFuture == null) {
- return;
- }
-
- requestFuture.complete(packet);
- }
-}
diff --git a/messenger-redis/src/test/java/io/github/sadcenter/messenger/redis/tests/RedisMessengerTest.java b/messenger-redis/src/test/java/io/github/sadcenter/messenger/redis/tests/RedisMessengerTest.java
index b3a0b94..e0355be 100644
--- a/messenger-redis/src/test/java/io/github/sadcenter/messenger/redis/tests/RedisMessengerTest.java
+++ b/messenger-redis/src/test/java/io/github/sadcenter/messenger/redis/tests/RedisMessengerTest.java
@@ -4,17 +4,21 @@
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
+import com.fasterxml.jackson.annotation.JsonCreator;
+import com.fasterxml.jackson.annotation.JsonProperty;
import com.github.fppt.jedismock.RedisServer;
import io.github.sadcenter.messenger.Messenger;
import io.github.sadcenter.messenger.channel.Channel;
-import io.github.sadcenter.messenger.codec.FstCodec;
-import io.github.sadcenter.messenger.codec.MessagePackCodec;
+import io.github.sadcenter.messenger.codec.Codec;
+import io.github.sadcenter.messenger.codecs.binary.FstCodec;
+import io.github.sadcenter.messenger.codecs.binary.MessagePackCodec;
+import io.github.sadcenter.messenger.codecs.json.JacksonCodec;
import io.github.sadcenter.messenger.handler.PacketHandler;
import io.github.sadcenter.messenger.packet.PacketRequest;
import io.github.sadcenter.messenger.redis.RedisMessenger;
import io.lettuce.core.RedisURI;
-import java.io.IOException;
import java.time.Duration;
+import java.util.List;
import java.util.Objects;
import java.util.concurrent.CompletableFuture;
import org.junit.jupiter.api.BeforeAll;
@@ -25,29 +29,27 @@ class RedisMessengerTest {
private static String HOST;
private static int PORT;
+ private static final List CODECS =
+ List.of(new FstCodec(), new MessagePackCodec());
+ private static final int CODEC_SIZE = 2;
@BeforeAll
static void initAll() throws Exception {
- RedisServer server = RedisServer
- .newRedisServer()
- .start();
+ RedisServer server = RedisServer.newRedisServer().start();
HOST = server.getHost();
PORT = server.getBindPort();
}
- @RepeatedTest(2)
- void testPublish(RepetitionInfo repetitionInfo) throws Exception {
+ @RepeatedTest(CODEC_SIZE)
+ void testPublish(RepetitionInfo repetitionInfo) {
Messenger messenger =
new RedisMessenger(
- RedisURI.create(HOST, PORT),
- repetitionInfo.getCurrentRepetition() % 2 == 1
- ? new MessagePackCodec()
- : new FstCodec());
-
+ RedisURI.create(HOST, PORT), CODECS.get(repetitionInfo.getCurrentRepetition() - 1));
+ messenger.initialize();
messenger.listen(new TestHandler(messenger), TestPacket.class, "messenger:test");
- TestPacket sentPacket = new TestPacket("Kacper", "Krzychała");
+ TestPacket sentPacket = new TestPacket("Kacper", "Krzychała");
CompletableFuture request = messenger.request("messenger:test", sentPacket);
await()
.atMost(Duration.ofSeconds(3))
@@ -56,23 +58,18 @@ void testPublish(RepetitionInfo repetitionInfo) throws Exception {
request.isDone() && !request.isCancelled() && !request.isCompletedExceptionally());
request.thenAccept(
receivedPacket -> assertEquals(sentPacket, receivedPacket, "Packets are not the same"));
-
messenger.close();
}
- @RepeatedTest(2)
- void testKeyValue(RepetitionInfo repetitionInfo) throws IOException, InterruptedException {
+ @RepeatedTest(CODEC_SIZE)
+ void testKeyValue(RepetitionInfo repetitionInfo) {
Messenger messenger =
new RedisMessenger(
- RedisURI.create(HOST, PORT),
- repetitionInfo.getCurrentRepetition() % 2 == 1
- ? new MessagePackCodec()
- : new FstCodec());
+ RedisURI.create(HOST, PORT), CODECS.get(repetitionInfo.getCurrentRepetition() - 1));
TestPacket testPacket = new TestPacket("Kacper", "Krzychala");
messenger.set("messenger", "test", testPacket);
- CompletableFuture future = messenger
- .get("messenger", "test");
+ CompletableFuture future = messenger.get("messenger", "test");
future.whenComplete(
(packet, throwable) -> {
assertNull(throwable);
@@ -101,24 +98,12 @@ private static final class TestPacket extends PacketRequest {
private final String name;
private final String surname;
- public TestPacket(String name, String surname) {
+ @JsonCreator
+ public TestPacket(@JsonProperty("name") String name, @JsonProperty("surname") String surname) {
this.name = name;
this.surname = surname;
}
- // Jackson requires an empty constructor
- public TestPacket() {
- this(null, null);
- }
-
- public String getName() {
- return name;
- }
-
- public String getSurname() {
- return surname;
- }
-
@Override
public boolean equals(Object o) {
if (this == o) {
diff --git a/pom.xml b/pom.xml
index 5903aff..3ba6548 100644
--- a/pom.xml
+++ b/pom.xml
@@ -7,18 +7,19 @@
io.github.sadcenter
messenger
pom
- 2.1.6
+ 3.0.0
messenger-core
- messenger-codec-fst
- messenger-codec-msgpack
messenger-nats
messenger-nats-luckperms
messenger-redis
+ messenger-codecs
+ messenger-benchmarks
UTF-8
+ 11
@@ -39,40 +40,48 @@
11
11
-
-
- org.jetbrains
- annotations
- 23.0.0
-
-
org.apache.maven.plugins
maven-source-plugin
3.2.1
-
-
- attach-sources
-
- jar
-
-
-
org.apache.maven.plugins
- maven-javadoc-plugin
- 3.3.2
+ maven-shade-plugin
+ 3.2.4
- attach-javadocs
+ package
- jar
+ shade
+
+ ${project.build.directory}/dependency-reduced-pom.xml
+
+
+
+ maven-surefire-plugin
+ 2.19
+
+
+ org.junit.platform
+ junit-platform-surefire-provider
+ 1.0.0
+
+
+
+ --add-opens=java.base/java.lang=ALL-UNNAMED
+ --add-opens=java.base/java.math=ALL-UNNAMED
+ --add-opens=java.base/java.util=ALL-UNNAMED
+ --add-opens=java.base/java.util.concurrent=ALL-UNNAMED
+ --add-opens=java.base/java.net=ALL-UNNAMED
+ --add-opens=java.base/java.text=ALL-UNNAMED
+ --add-opens=java.sql/java.sql=ALL-UNNAMED
+
@@ -83,6 +92,7 @@
+
shitzuu-me-repo-releases