Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
563ea08a3f |
@@ -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);
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user