Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b80ce6ea01 | ||
|
|
b236de082a | ||
|
|
1ba5cd8159 |
@@ -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,22 @@ 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(() -> {
|
||||
try {
|
||||
return handler.handle(id, listener.mapper.treeToValue(value, method.getReqClass())).get()
|
||||
.intoJSON(listener.mapper);
|
||||
} catch (Exception e) {
|
||||
return SessionResult.error(TextNode.valueOf(e.getMessage()));
|
||||
}
|
||||
});
|
||||
};
|
||||
|
||||
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, JsonNode, 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 {
|
||||
if (res.isError()) {
|
||||
respondError(ws, r.id(), res.value());
|
||||
} else {
|
||||
respond(ws, r.id(), res.value());
|
||||
}
|
||||
} catch (Exception ignored) {
|
||||
}
|
||||
});
|
||||
try {
|
||||
var res = handler.handle(r.id(), r.data());
|
||||
|
||||
if (res.isError()) {
|
||||
respondError(ws, r.id(), res.getJsonNode(mapper));
|
||||
} else {
|
||||
respond(ws, r.id(), res.getJsonNode(mapper));
|
||||
}
|
||||
} catch (Exception ignored) {
|
||||
}
|
||||
}
|
||||
} else if (msg instanceof Message.Response r) {
|
||||
CompletableFuture<JsonNode> fut = pending.remove(r.id());
|
||||
@@ -58,14 +60,33 @@ public class SessionListener implements WebSocket.Listener {
|
||||
new RuntimeException(r.error().toString()));
|
||||
}
|
||||
|
||||
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,28 +103,36 @@ 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);
|
||||
|
||||
send(ws, new Message.Request(
|
||||
id,
|
||||
method,
|
||||
mapper.valueToTree(data)));
|
||||
|
||||
return fut.thenApply(json -> {
|
||||
try {
|
||||
return mapper.treeToValue(json, resType);
|
||||
} catch (Exception e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
});
|
||||
return CompletableFuture
|
||||
.runAsync(() -> {
|
||||
try {
|
||||
send(ws, new Message.Request(
|
||||
id,
|
||||
method,
|
||||
mapper.valueToTree(data)));
|
||||
} 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 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,8 @@
|
||||
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);
|
||||
}
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user