1 Commits
Author SHA1 Message Date
selimaj-dev 563ea08a3f Fixed notification handling 2026-02-23 09:19:05 +01:00
3 changed files with 18 additions and 12 deletions
@@ -83,15 +83,13 @@ public final class Session {
this.listener.methods.put(method.getName(), wrapper); this.listener.methods.put(method.getName(), wrapper);
} }
public <Req, Res, Err> void onNotification(Method<Req, Res, Err> method, public <Req, Res, Err> void onNotification(Method<Req, Res, Err> method, NotificationHandler<Req> handler) {
NotificationHandler<Req, Res, Err> handler) { NotificationHandler<JsonNode> wrapper = (value) -> {
NotificationHandler<JsonNode, JsonNode, JsonNode> wrapper = (id, value) -> { CompletableFuture.runAsync(() -> {
return CompletableFuture.supplyAsync(() -> {
try { try {
return handler.handle(id, listener.mapper.treeToValue(value, method.getReqClass())).get() handler.handle(listener.mapper.treeToValue(value, method.getReqClass()));
.intoJSON(listener.mapper);
} catch (Exception e) { } catch (Exception e) {
return SessionResult.error(TextNode.valueOf(e.getMessage())); e.printStackTrace();
} }
}); });
}; };
@@ -17,7 +17,7 @@ import dev.selimaj.session.types.NotificationHandler;
public class SessionListener implements WebSocket.Listener { public class SessionListener implements WebSocket.Listener {
final ConcurrentHashMap<String, MethodHandler<JsonNode, JsonNode, JsonNode>> methods = new ConcurrentHashMap<>(); final ConcurrentHashMap<String, MethodHandler<JsonNode, JsonNode, JsonNode>> methods = new ConcurrentHashMap<>();
final ConcurrentHashMap<String, NotificationHandler<JsonNode, JsonNode, JsonNode>> notificationHandlers = new ConcurrentHashMap<>(); final ConcurrentHashMap<String, NotificationHandler<JsonNode>> notificationHandlers = new ConcurrentHashMap<>();
final ObjectMapper mapper = new ObjectMapper(); final ObjectMapper mapper = new ObjectMapper();
private final ConcurrentHashMap<Integer, CompletableFuture<JsonNode>> pending = new ConcurrentHashMap<>(); private final ConcurrentHashMap<Integer, CompletableFuture<JsonNode>> pending = new ConcurrentHashMap<>();
@@ -58,6 +58,16 @@ public class SessionListener implements WebSocket.Listener {
if (fut != null) if (fut != null)
fut.completeExceptionally( fut.completeExceptionally(
new RuntimeException(r.error().toString())); 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); ws.request(1);
@@ -1,8 +1,6 @@
package dev.selimaj.session.types; package dev.selimaj.session.types;
import java.util.concurrent.CompletableFuture;
@FunctionalInterface @FunctionalInterface
public interface NotificationHandler<Req, Res, Err> { public interface NotificationHandler<Req> {
CompletableFuture<SessionResult<Res, Err>> handle(int id, Req data); void handle(Req data);
} }