5 Commits
Author SHA1 Message Date
selimaj-dev 563ea08a3f Fixed notification handling 2026-02-23 09:19:05 +01:00
selimaj-dev b80ce6ea01 Notification handler 2026-02-23 06:28:12 +01:00
selimaj-dev b236de082a Fixed ping bug 2026-02-23 05:54:56 +01:00
selimaj-dev 1ba5cd8159 HOWWWW 2026-02-20 02:22:58 +01:00
selimaj-dev 34d86a23ab Fixed build artifacts 2026-02-20 02:06:05 +01:00
7 changed files with 122 additions and 48 deletions
+1 -1
View File
@@ -33,7 +33,7 @@ repositories {
#### Add the dependency
```groovy
dependencies {
implementation 'com.github.selimaj-dev:session-java:v0.1.3'
implementation 'com.github.selimaj-dev:session-java:0.1.3'
}
```
+9 -1
View File
@@ -4,7 +4,7 @@ plugins {
}
group = "dev.selimaj.session"
version = "0.1.2"
version = "0.1.3"
java {
toolchain {
@@ -34,3 +34,11 @@ dependencies {
test {
useJUnitPlatform()
}
publishing {
publications {
mavenJava(MavenPublication) {
from components.java
}
}
}
+21 -5
View File
@@ -7,6 +7,7 @@ import dev.selimaj.session.types.SessionResult;
import dev.selimaj.session.types.Message;
import dev.selimaj.session.types.Method;
import dev.selimaj.session.types.MethodHandler;
import dev.selimaj.session.types.NotificationHandler;
import java.net.URI;
import java.net.http.HttpClient;
@@ -63,16 +64,17 @@ public final class Session {
listener.notify(ws, method, data);
}
<Req, Res, Err> CompletableFuture<Res> request(
public <Req, Res, Err> CompletableFuture<Res> request(
Method<Req, Res, Err> method, Req req) throws Exception {
return listener.request(ws, method.getName(), req, method.getResClass());
}
<Req, Res, Err> void onRequest(Method<Req, Res, Err> method,
MethodHandler<Req> handler) {
MethodHandler<JsonNode> wrapper = (id, value) -> {
public <Req, Res, Err> void onRequest(Method<Req, Res, Err> method,
MethodHandler<Req, Res, Err> handler) {
MethodHandler<JsonNode, JsonNode, JsonNode> wrapper = (id, value) -> {
try {
return handler.handle(id, listener.mapper.treeToValue(value, method.getReqClass()));
return handler.handle(id, listener.mapper.treeToValue(value, method.getReqClass()))
.intoJSON(listener.mapper);
} catch (Exception e) {
return SessionResult.error(TextNode.valueOf(e.getMessage()));
}
@@ -81,6 +83,20 @@ public final class Session {
this.listener.methods.put(method.getName(), wrapper);
}
public <Req, Res, Err> void onNotification(Method<Req, Res, Err> method, NotificationHandler<Req> handler) {
NotificationHandler<JsonNode> wrapper = (value) -> {
CompletableFuture.runAsync(() -> {
try {
handler.handle(listener.mapper.treeToValue(value, method.getReqClass()));
} catch (Exception e) {
e.printStackTrace();
}
});
};
this.listener.notificationHandlers.put(method.getName(), wrapper);
}
public void close() throws Exception {
ws.sendClose(0, "Session closed");
}
@@ -3,6 +3,7 @@ package dev.selimaj.session;
import java.net.http.WebSocket;
import java.nio.ByteBuffer;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
import java.util.concurrent.CompletionStage;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicInteger;
@@ -12,9 +13,11 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import dev.selimaj.session.types.Message;
import dev.selimaj.session.types.MethodHandler;
import dev.selimaj.session.types.NotificationHandler;
public class SessionListener implements WebSocket.Listener {
final ConcurrentHashMap<String, MethodHandler<JsonNode>> methods = new ConcurrentHashMap<>();
final ConcurrentHashMap<String, MethodHandler<JsonNode, JsonNode, JsonNode>> methods = new ConcurrentHashMap<>();
final ConcurrentHashMap<String, NotificationHandler<JsonNode>> notificationHandlers = new ConcurrentHashMap<>();
final ObjectMapper mapper = new ObjectMapper();
private final ConcurrentHashMap<Integer, CompletableFuture<JsonNode>> pending = new ConcurrentHashMap<>();
@@ -32,20 +35,19 @@ public class SessionListener implements WebSocket.Listener {
}
if (msg instanceof Message.Request r) {
MethodHandler<JsonNode> handler = methods.get(r.method());
MethodHandler<JsonNode, JsonNode, JsonNode> handler = methods.get(r.method());
if (handler != null) {
handler.handle(r.id(), r.data())
.thenAccept(res -> {
try {
var res = handler.handle(r.id(), r.data());
if (res.isError()) {
respondError(ws, r.id(), res.value());
respondError(ws, r.id(), res.getJsonNode(mapper));
} else {
respond(ws, r.id(), res.value());
respond(ws, r.id(), res.getJsonNode(mapper));
}
} catch (Exception ignored) {
}
});
}
} else if (msg instanceof Message.Response r) {
CompletableFuture<JsonNode> fut = pending.remove(r.id());
@@ -56,16 +58,45 @@ public class SessionListener implements WebSocket.Listener {
if (fut != null)
fut.completeExceptionally(
new RuntimeException(r.error().toString()));
} else if (msg instanceof Message.Notification r) {
NotificationHandler<JsonNode> handler = notificationHandlers.get(r.method());
if (handler != null) {
try {
handler.handle(r.data());
} catch (Exception e) {
e.printStackTrace();
}
}
}
ws.request(1);
return null;
}
@Override
public void onOpen(WebSocket ws) {
ws.request(1);
}
@Override
public CompletionStage<?> onPing(WebSocket ws, ByteBuffer message) {
ws.request(1);
return ws.sendPong(message);
}
@Override
public CompletionStage<?> onPong(WebSocket ws, ByteBuffer message) {
ws.request(1);
return null;
}
@Override
public CompletionStage<?> onClose(WebSocket ws, int statusCode, String reason) {
ws.request(1);
return null;
}
void send(WebSocket ws, Message msg) throws Exception {
ws.sendText(mapper.writeValueAsString(msg), true);
}
@@ -82,27 +113,35 @@ public class SessionListener implements WebSocket.Listener {
send(ws, new Message.Notification(method, mapper.valueToTree(data)));
}
<Req, Res> CompletableFuture<Res> request(
public <Req, Res> CompletableFuture<Res> request(
WebSocket ws,
String method,
Req data,
Class<Res> resType) throws Exception {
Class<Res> resType) {
int id = this.id.incrementAndGet();
CompletableFuture<JsonNode> fut = new CompletableFuture<>();
pending.put(id, fut);
return CompletableFuture
.runAsync(() -> {
try {
send(ws, new Message.Request(
id,
method,
mapper.valueToTree(data)));
return fut.thenApply(json -> {
} catch (Exception e) {
throw new CompletionException(e);
}
})
.thenCompose(v -> {
CompletableFuture<JsonNode> fut = new CompletableFuture<>();
pending.put(id, fut);
return fut;
})
.thenApply(json -> {
try {
return mapper.treeToValue(json, resType);
} catch (Exception e) {
throw new RuntimeException(e);
throw new CompletionException(e);
}
});
}
@@ -1,8 +1,6 @@
package dev.selimaj.session.types;
import java.util.concurrent.CompletableFuture;
@FunctionalInterface
public interface MethodHandler<Req> {
CompletableFuture<SessionResult> handle(int id, Req data);
public interface MethodHandler<Req, Res, Err> {
SessionResult<Res, Err> handle(int id, Req data);
}
@@ -0,0 +1,6 @@
package dev.selimaj.session.types;
@FunctionalInterface
public interface NotificationHandler<Req> {
void handle(Req data);
}
@@ -1,15 +1,22 @@
package dev.selimaj.session.types;
import java.util.concurrent.CompletableFuture;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
public record SessionResult(boolean isError, JsonNode value) {
public static CompletableFuture<SessionResult> ok(JsonNode value) {
return CompletableFuture.completedFuture(new SessionResult(false, value));
public record SessionResult<T, E>(boolean isError, Object value) {
public static <T, E> SessionResult<T, E> ok(T value) {
return new SessionResult<>(false, value);
}
public static CompletableFuture<SessionResult> error(JsonNode value) {
return CompletableFuture.completedFuture(new SessionResult(false, value));
public static <T, E> SessionResult<T, E> error(E err) {
return new SessionResult<>(false, err);
}
public SessionResult<JsonNode, JsonNode> intoJSON(ObjectMapper mapper) {
return new SessionResult<>(this.isError, mapper.valueToTree(this.value));
}
public JsonNode getJsonNode(ObjectMapper mapper) {
return mapper.valueToTree(this.value);
}
}