From bbc005c7572a0583d2c1034f1c6cf1a319b28616 Mon Sep 17 00:00:00 2001 From: sadcenter Date: Sun, 22 May 2022 14:22:20 +0200 Subject: [PATCH 1/3] hermes --- .idea/compiler.xml | 3 ++- .idea/encodings.xml | 2 ++ messenger-hermes/pom.xml | 19 +++++++++++++++++++ pom.xml | 1 + 4 files changed, 24 insertions(+), 1 deletion(-) create mode 100644 messenger-hermes/pom.xml diff --git a/.idea/compiler.xml b/.idea/compiler.xml index 293159c..e316254 100644 --- a/.idea/compiler.xml +++ b/.idea/compiler.xml @@ -15,11 +15,12 @@ - + + diff --git a/.idea/encodings.xml b/.idea/encodings.xml index 774fbaa..596b248 100644 --- a/.idea/encodings.xml +++ b/.idea/encodings.xml @@ -10,6 +10,8 @@ + + diff --git a/messenger-hermes/pom.xml b/messenger-hermes/pom.xml new file mode 100644 index 0000000..2b59996 --- /dev/null +++ b/messenger-hermes/pom.xml @@ -0,0 +1,19 @@ + + + + messenger + io.github.sadcenter + 2.1.6 + + 4.0.0 + + messenger-hermes + + + 17 + 17 + + + \ No newline at end of file diff --git a/pom.xml b/pom.xml index 5903aff..79533bb 100644 --- a/pom.xml +++ b/pom.xml @@ -15,6 +15,7 @@ messenger-nats messenger-nats-luckperms messenger-redis + messenger-hermes From eeb8c9b1da0ef7fdf8793cf353f85f6d92b3a9be Mon Sep 17 00:00:00 2001 From: sadcenter Date: Sun, 29 May 2022 02:50:02 +0200 Subject: [PATCH 2/3] get me out of this hell --- messenger-benchmarks/pom.xml | 19 +++ messenger-codec-fst/pom.xml | 74 --------- messenger-codec-msgpack/pom.xml | 81 ---------- messenger-codecs/pom.xml | 59 ++++++++ .../messenger/codecs/binary}/FstCodec.java | 3 +- .../codecs/binary}/MessagePackCodec.java | 11 +- .../messenger/codecs/json/JacksonCodec.java | 74 +++++++++ .../src/test/java/CodecsTest.java | 78 ++++++++++ messenger-core/pom.xml | 47 +----- .../sadcenter/messenger/channel/Channel.java | 6 +- .../sadcenter/messenger/packet/Packet.java | 5 +- .../messenger/packet/PacketRequest.java | 16 +- messenger-hermes/pom.xml | 26 +++- .../messenger/hermes/HermesMessenger.java | 140 ++++++++++++++++++ messenger-nats-luckperms/pom.xml | 48 ++---- messenger-nats/pom.xml | 106 ++++--------- .../messenger/nats/NatsMessenger.java | 8 +- .../nats/tests/NatsMessengerTest.java | 105 ++++++------- messenger-redis/pom.xml | 103 +++++-------- .../messenger/redis/RedisMessenger.java | 84 ++++++----- .../messenger/redis/codec/RedisByteCodec.java | 4 - .../redis/codec/StringByteArrayCodec.java | 1 - .../redis/request/RedisRequestHandler.java | 31 ---- .../redis/tests/RedisMessengerTest.java | 59 +++----- pom.xml | 54 ++++--- 25 files changed, 643 insertions(+), 599 deletions(-) create mode 100644 messenger-benchmarks/pom.xml delete mode 100644 messenger-codec-fst/pom.xml delete mode 100644 messenger-codec-msgpack/pom.xml create mode 100644 messenger-codecs/pom.xml rename {messenger-codec-fst/src/main/java/io/github/sadcenter/messenger/codec => messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary}/FstCodec.java (83%) rename {messenger-codec-msgpack/src/main/java/io/github/sadcenter/messenger/codec => messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary}/MessagePackCodec.java (88%) create mode 100644 messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/json/JacksonCodec.java create mode 100644 messenger-codecs/src/test/java/CodecsTest.java create mode 100644 messenger-hermes/src/main/java/io/github/sadcenter/messenger/hermes/HermesMessenger.java delete mode 100644 messenger-redis/src/main/java/io/github/sadcenter/messenger/redis/request/RedisRequestHandler.java diff --git a/messenger-benchmarks/pom.xml b/messenger-benchmarks/pom.xml new file mode 100644 index 0000000..9d65d22 --- /dev/null +++ b/messenger-benchmarks/pom.xml @@ -0,0 +1,19 @@ + + + + messenger + io.github.sadcenter + 3.0.0 + + 4.0.0 + + messenger-benchmarks + + + 17 + 17 + + + \ No newline at end of file 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 83% 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..f05a212 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,5 +1,6 @@ -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") 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 88% 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..b3a1387 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,4 +1,4 @@ -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; @@ -7,12 +7,13 @@ 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 +44,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 +59,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..ea34f09 --- /dev/null +++ b/messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/json/JacksonCodec.java @@ -0,0 +1,74 @@ +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); + } + + @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/CodecsTest.java b/messenger-codecs/src/test/java/CodecsTest.java new file mode 100644 index 0000000..8393044 --- /dev/null +++ b/messenger-codecs/src/test/java/CodecsTest.java @@ -0,0 +1,78 @@ +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-hermes/pom.xml b/messenger-hermes/pom.xml index 2b59996..d60fad6 100644 --- a/messenger-hermes/pom.xml +++ b/messenger-hermes/pom.xml @@ -5,15 +5,31 @@ messenger io.github.sadcenter - 2.1.6 + 3.0.0 4.0.0 messenger-hermes - - 17 - 17 - + + + io.github.sadcenter + messenger-core + 3.0.0 + compile + + + pl.allegro.tech.hermes + hermes-client + 1.13.0 + compile + + + pl.allegro.tech.hermes + hermes-mock + 1.13.0 + test + + \ No newline at end of file diff --git a/messenger-hermes/src/main/java/io/github/sadcenter/messenger/hermes/HermesMessenger.java b/messenger-hermes/src/main/java/io/github/sadcenter/messenger/hermes/HermesMessenger.java new file mode 100644 index 0000000..f999123 --- /dev/null +++ b/messenger-hermes/src/main/java/io/github/sadcenter/messenger/hermes/HermesMessenger.java @@ -0,0 +1,140 @@ +package io.github.sadcenter.messenger.hermes; + +import io.github.sadcenter.messenger.Messenger; +import io.github.sadcenter.messenger.channel.Channel; +import io.github.sadcenter.messenger.codec.Codec; +import io.github.sadcenter.messenger.handler.PacketHandler; +import io.github.sadcenter.messenger.handler.PacketHandlerAbstract; +import io.github.sadcenter.messenger.packet.Packet; +import io.github.sadcenter.messenger.packet.PacketRequest; +import java.util.Set; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeoutException; +import org.jetbrains.annotations.NotNull; +import pl.allegro.tech.hermes.client.HermesClient; +import pl.allegro.tech.hermes.client.HermesClientBuilder; +import pl.allegro.tech.hermes.client.HermesMessage; + +@SuppressWarnings("unchecked") +public class HermesMessenger implements Messenger { + + private final HermesClient hermesClient; + private Codec codec; + + public HermesMessenger(@NotNull HermesClient hermesClient) { + this.hermesClient = hermesClient; + } + + public HermesMessenger(@NotNull HermesClientBuilder builder) { + this.hermesClient = builder.build(); + } + + @Override + public void initialize() { + // + } + + @Override + public void close() { + try { + this.hermesClient.close(50L, 1000L); + } catch (InterruptedException | TimeoutException e) { + e.printStackTrace(); + } + } + + @Override + public void publish(@NotNull String channel, @NotNull Packet packet) { + this.hermesClient.publish(HermesMessage.hermesMessage(channel, this.codec.encode(packet)).build()); + } + + @Override + public void reply(@NotNull Channel channel, @NotNull PacketRequest packet) { + this.publish(channel.getReplyTo(), packet); + } + + @Override + public CompletableFuture request(@NotNull String channel, + @NotNull T packet) { + return null; + } + + @Override + public void listen(@NotNull String channel, @NotNull PacketHandler handler, + @NotNull Class type) { + + } + + @Override + public void listen(@NotNull Set channels, + @NotNull PacketHandler handler, @NotNull Class type) { + + } + + @Override + public void listen(@NotNull PacketHandler handler, @NotNull Class type, + @NotNull String... channels) { + + } + + @Override + public void listen(@NotNull PacketHandlerAbstract abstractHandler) { + + } + + @Override + public void listen(@NotNull PacketHandlerAbstract... abstractHandlers) { + + } + + @Override + public void set(@NotNull String id, @NotNull String key, @NotNull String value) { + + } + + @Override + public void set(@NotNull String id, @NotNull String key, @NotNull T value) { + + } + + @Override + public CompletableFuture get(@NotNull String id, @NotNull String key) { + return null; + } + + @Override + public CompletableFuture get(@NotNull String id, @NotNull String key, + @NotNull Class type) { + return null; + } + + @Override + public CompletableFuture getString(@NotNull String id, @NotNull String key) { + return null; + } + + @Override + public void delete(@NotNull String id, @NotNull String key) { + + } + + @Override + public void purge(@NotNull String id, @NotNull String key) { + + } + + @Override + public void from(@NotNull String from) { + + } + + @Override + public @NotNull T codec() { + return (T) this.codec; + } + + @Override + public void codec(@NotNull T codec) { + this.codec = codec; + } +} 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 79533bb..7b775e9 100644 --- a/pom.xml +++ b/pom.xml @@ -7,19 +7,20 @@ 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-hermes + messenger-codecs + messenger-benchmarks UTF-8 + 11 @@ -40,40 +41,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 + @@ -84,6 +93,7 @@ + shitzuu-me-repo-releases From a1fcaebd91734f9fa56384e0f75debc1557914ab Mon Sep 17 00:00:00 2001 From: sadcenter Date: Tue, 5 Jul 2022 23:13:49 +0200 Subject: [PATCH 3/3] code clean-up and adding benchmarks --- .idea/compiler.xml | 17 +-- .idea/encodings.xml | 9 ++ .idea/misc.xml | 2 + messenger-benchmarks/pom.xml | 72 ++++++++- .../messenger/BenchmarkApplication.java | 11 ++ .../benchmarks/CodecBenchmarkTest.java | 40 +++++ .../benchmarks/MessengerBenchmarkTest.java | 5 + .../benchmarks/example/ExamplePacket.java | 43 ++++++ .../messenger/codecs/binary/FstCodec.java | 10 +- .../codecs/binary/MessagePackCodec.java | 1 + .../messenger/codecs/json/JacksonCodec.java | 4 + .../messenger/codecs}/CodecsTest.java | 2 + messenger-hermes/pom.xml | 35 ----- .../messenger/hermes/HermesMessenger.java | 140 ------------------ pom.xml | 1 - 15 files changed, 200 insertions(+), 192 deletions(-) create mode 100644 messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/BenchmarkApplication.java create mode 100644 messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/benchmarks/CodecBenchmarkTest.java create mode 100644 messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/benchmarks/MessengerBenchmarkTest.java create mode 100644 messenger-benchmarks/src/main/java/io/github/sadcenter/messenger/benchmarks/example/ExamplePacket.java rename messenger-codecs/src/test/java/{ => io/github/sadcenter/messenger/codecs}/CodecsTest.java (97%) delete mode 100644 messenger-hermes/pom.xml delete mode 100644 messenger-hermes/src/main/java/io/github/sadcenter/messenger/hermes/HermesMessenger.java diff --git a/.idea/compiler.xml b/.idea/compiler.xml index e316254..6bfebb7 100644 --- a/.idea/compiler.xml +++ b/.idea/compiler.xml @@ -6,25 +6,20 @@ - - - - - - - - - - + - + + + + + \ No newline at end of file diff --git a/.idea/encodings.xml b/.idea/encodings.xml index 596b248..4fd1572 100644 --- a/.idea/encodings.xml +++ b/.idea/encodings.xml @@ -2,18 +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 index 9d65d22..dcbb223 100644 --- a/messenger-benchmarks/pom.xml +++ b/messenger-benchmarks/pom.xml @@ -11,9 +11,73 @@ messenger-benchmarks - - 17 - 17 - + + + 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-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary/FstCodec.java b/messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary/FstCodec.java index f05a212..776f1a3 100644 --- a/messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary/FstCodec.java +++ b/messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary/FstCodec.java @@ -6,7 +6,15 @@ @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-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary/MessagePackCodec.java b/messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary/MessagePackCodec.java index b3a1387..d190abe 100644 --- a/messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary/MessagePackCodec.java +++ b/messenger-codecs/src/main/java/io/github/sadcenter/messenger/codecs/binary/MessagePackCodec.java @@ -3,6 +3,7 @@ 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; 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 index ea34f09..530e9b8 100644 --- 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 @@ -39,6 +39,10 @@ public JacksonCodec() { this.objectMapper.registerSubtypes(PacketRequest.class); } + public JacksonCodec(ObjectMapper objectMapper) { + this.objectMapper = objectMapper; + } + @Override public byte[] encode(Object data) { try { diff --git a/messenger-codecs/src/test/java/CodecsTest.java b/messenger-codecs/src/test/java/io/github/sadcenter/messenger/codecs/CodecsTest.java similarity index 97% rename from messenger-codecs/src/test/java/CodecsTest.java rename to messenger-codecs/src/test/java/io/github/sadcenter/messenger/codecs/CodecsTest.java index 8393044..ddda936 100644 --- a/messenger-codecs/src/test/java/CodecsTest.java +++ b/messenger-codecs/src/test/java/io/github/sadcenter/messenger/codecs/CodecsTest.java @@ -1,3 +1,5 @@ +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; diff --git a/messenger-hermes/pom.xml b/messenger-hermes/pom.xml deleted file mode 100644 index d60fad6..0000000 --- a/messenger-hermes/pom.xml +++ /dev/null @@ -1,35 +0,0 @@ - - - - messenger - io.github.sadcenter - 3.0.0 - - 4.0.0 - - messenger-hermes - - - - io.github.sadcenter - messenger-core - 3.0.0 - compile - - - pl.allegro.tech.hermes - hermes-client - 1.13.0 - compile - - - pl.allegro.tech.hermes - hermes-mock - 1.13.0 - test - - - - \ No newline at end of file diff --git a/messenger-hermes/src/main/java/io/github/sadcenter/messenger/hermes/HermesMessenger.java b/messenger-hermes/src/main/java/io/github/sadcenter/messenger/hermes/HermesMessenger.java deleted file mode 100644 index f999123..0000000 --- a/messenger-hermes/src/main/java/io/github/sadcenter/messenger/hermes/HermesMessenger.java +++ /dev/null @@ -1,140 +0,0 @@ -package io.github.sadcenter.messenger.hermes; - -import io.github.sadcenter.messenger.Messenger; -import io.github.sadcenter.messenger.channel.Channel; -import io.github.sadcenter.messenger.codec.Codec; -import io.github.sadcenter.messenger.handler.PacketHandler; -import io.github.sadcenter.messenger.handler.PacketHandlerAbstract; -import io.github.sadcenter.messenger.packet.Packet; -import io.github.sadcenter.messenger.packet.PacketRequest; -import java.util.Set; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.TimeoutException; -import org.jetbrains.annotations.NotNull; -import pl.allegro.tech.hermes.client.HermesClient; -import pl.allegro.tech.hermes.client.HermesClientBuilder; -import pl.allegro.tech.hermes.client.HermesMessage; - -@SuppressWarnings("unchecked") -public class HermesMessenger implements Messenger { - - private final HermesClient hermesClient; - private Codec codec; - - public HermesMessenger(@NotNull HermesClient hermesClient) { - this.hermesClient = hermesClient; - } - - public HermesMessenger(@NotNull HermesClientBuilder builder) { - this.hermesClient = builder.build(); - } - - @Override - public void initialize() { - // - } - - @Override - public void close() { - try { - this.hermesClient.close(50L, 1000L); - } catch (InterruptedException | TimeoutException e) { - e.printStackTrace(); - } - } - - @Override - public void publish(@NotNull String channel, @NotNull Packet packet) { - this.hermesClient.publish(HermesMessage.hermesMessage(channel, this.codec.encode(packet)).build()); - } - - @Override - public void reply(@NotNull Channel channel, @NotNull PacketRequest packet) { - this.publish(channel.getReplyTo(), packet); - } - - @Override - public CompletableFuture request(@NotNull String channel, - @NotNull T packet) { - return null; - } - - @Override - public void listen(@NotNull String channel, @NotNull PacketHandler handler, - @NotNull Class type) { - - } - - @Override - public void listen(@NotNull Set channels, - @NotNull PacketHandler handler, @NotNull Class type) { - - } - - @Override - public void listen(@NotNull PacketHandler handler, @NotNull Class type, - @NotNull String... channels) { - - } - - @Override - public void listen(@NotNull PacketHandlerAbstract abstractHandler) { - - } - - @Override - public void listen(@NotNull PacketHandlerAbstract... abstractHandlers) { - - } - - @Override - public void set(@NotNull String id, @NotNull String key, @NotNull String value) { - - } - - @Override - public void set(@NotNull String id, @NotNull String key, @NotNull T value) { - - } - - @Override - public CompletableFuture get(@NotNull String id, @NotNull String key) { - return null; - } - - @Override - public CompletableFuture get(@NotNull String id, @NotNull String key, - @NotNull Class type) { - return null; - } - - @Override - public CompletableFuture getString(@NotNull String id, @NotNull String key) { - return null; - } - - @Override - public void delete(@NotNull String id, @NotNull String key) { - - } - - @Override - public void purge(@NotNull String id, @NotNull String key) { - - } - - @Override - public void from(@NotNull String from) { - - } - - @Override - public @NotNull T codec() { - return (T) this.codec; - } - - @Override - public void codec(@NotNull T codec) { - this.codec = codec; - } -} diff --git a/pom.xml b/pom.xml index 7b775e9..3ba6548 100644 --- a/pom.xml +++ b/pom.xml @@ -13,7 +13,6 @@ messenger-nats messenger-nats-luckperms messenger-redis - messenger-hermes messenger-codecs messenger-benchmarks