From 563ea08a3fcfe9fd5f1320f279a0d68b38b30280 Mon Sep 17 00:00:00 2001 From: Klesti Selimaj Date: Mon, 23 Feb 2026 09:19:05 +0100 Subject: [PATCH] Fixed notification handling --- src/main/java/dev/selimaj/session/Session.java | 12 +++++------- .../java/dev/selimaj/session/SessionListener.java | 12 +++++++++++- .../selimaj/session/types/NotificationHandler.java | 6 ++---- 3 files changed, 18 insertions(+), 12 deletions(-) diff --git a/src/main/java/dev/selimaj/session/Session.java b/src/main/java/dev/selimaj/session/Session.java index 83c9ee9..2434612 100644 --- a/src/main/java/dev/selimaj/session/Session.java +++ b/src/main/java/dev/selimaj/session/Session.java @@ -83,15 +83,13 @@ public final class Session { this.listener.methods.put(method.getName(), wrapper); } - public void onNotification(Method method, - NotificationHandler handler) { - NotificationHandler wrapper = (id, value) -> { - return CompletableFuture.supplyAsync(() -> { + public void onNotification(Method method, NotificationHandler handler) { + NotificationHandler 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(); } }); }; diff --git a/src/main/java/dev/selimaj/session/SessionListener.java b/src/main/java/dev/selimaj/session/SessionListener.java index b7efdbe..2ab8aee 100644 --- a/src/main/java/dev/selimaj/session/SessionListener.java +++ b/src/main/java/dev/selimaj/session/SessionListener.java @@ -17,7 +17,7 @@ import dev.selimaj.session.types.NotificationHandler; public class SessionListener implements WebSocket.Listener { final ConcurrentHashMap> methods = new ConcurrentHashMap<>(); - final ConcurrentHashMap> notificationHandlers = new ConcurrentHashMap<>(); + final ConcurrentHashMap> notificationHandlers = new ConcurrentHashMap<>(); final ObjectMapper mapper = new ObjectMapper(); private final ConcurrentHashMap> 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 handler = notificationHandlers.get(r.method()); + + if (handler != null) { + try { + handler.handle(r.data()); + } catch (Exception e) { + e.printStackTrace(); + } + } } ws.request(1); diff --git a/src/main/java/dev/selimaj/session/types/NotificationHandler.java b/src/main/java/dev/selimaj/session/types/NotificationHandler.java index a7dced9..3ce7cb5 100644 --- a/src/main/java/dev/selimaj/session/types/NotificationHandler.java +++ b/src/main/java/dev/selimaj/session/types/NotificationHandler.java @@ -1,8 +1,6 @@ package dev.selimaj.session.types; -import java.util.concurrent.CompletableFuture; - @FunctionalInterface -public interface NotificationHandler { - CompletableFuture> handle(int id, Req data); +public interface NotificationHandler { + void handle(Req data); }