Fixed notification handling

This commit is contained in:
2026-02-23 09:19:05 +01:00
parent b80ce6ea01
commit 563ea08a3f
3 changed files with 18 additions and 12 deletions
@@ -83,15 +83,13 @@ public final class Session {
this.listener.methods.put(method.getName(), wrapper);
}
public <Req, Res, Err> void onNotification(Method<Req, Res, Err> method,
NotificationHandler<Req, Res, Err> handler) {
NotificationHandler<JsonNode, JsonNode, JsonNode> wrapper = (id, value) -> {
return CompletableFuture.supplyAsync(() -> {
public <Req, Res, Err> void onNotification(Method<Req, Res, Err> method, NotificationHandler<Req> handler) {
NotificationHandler<JsonNode> wrapper = (value) -> {
CompletableFuture.runAsync(() -> {
try {
return handler.handle(id, listener.mapper.treeToValue(value, method.getReqClass())).get()
.intoJSON(listener.mapper);
handler.handle(listener.mapper.treeToValue(value, method.getReqClass()));
} 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 {
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();
private final ConcurrentHashMap<Integer, CompletableFuture<JsonNode>> pending = new ConcurrentHashMap<>();
@@ -58,6 +58,16 @@ 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);
@@ -1,8 +1,6 @@
package dev.selimaj.session.types;
import java.util.concurrent.CompletableFuture;
@FunctionalInterface
public interface NotificationHandler<Req, Res, Err> {
CompletableFuture<SessionResult<Res, Err>> handle(int id, Req data);
public interface NotificationHandler<Req> {
void handle(Req data);
}