Class WebsocketEndpoint
java.lang.Object
io.fluxzero.testserver.websocket.WebsocketEndpoint
- Direct Known Subclasses:
ConsumerEndpoint, EventSourcingEndpoint, KeyValueEndPoint, ProducerEndpoint, SchedulingEndpoint, SearchEndpoint
-
Nested Class Summary
Nested ClassesModifier and TypeClassDescriptionprotected static classprotected static classprotected static class -
Field Summary
FieldsModifier and TypeFieldDescriptionprotected static Durationprotected static Durationprotected booleanprotected final AtomicBoolean -
Constructor Summary
ConstructorsModifierConstructorDescriptionprotectedprotectedWebsocketEndpoint(CommandIdempotencyStore commandIdempotencyStore) protectedWebsocketEndpoint(Executor requestExecutor) protectedWebsocketEndpoint(Executor requestExecutor, CommandIdempotencyStore commandIdempotencyStore) -
Method Summary
Modifier and TypeMethodDescriptionprotected voidabort(ServerWebsocketSession session, String reason) createTasks(io.fluxzero.common.api.RequestBatch<?> batch, ServerWebsocketSession session) protected io.fluxzero.common.api.JsonTypedeserializeRequest(ServerWebsocketSession session, byte[] bytes) protected voiddispatchRequest(ServerWebsocketSession session, io.fluxzero.common.api.JsonType request) protected voiddoSendResult(ServerWebsocketSession session, io.fluxzero.common.api.RequestResult result) protected Optional<WebsocketEndpoint.SessionBacklog> findAlternativeBacklog(ServerWebsocketSession closedSession) protected StringgetClientId(ServerWebsocketSession session) protected StringgetClientName(ServerWebsocketSession session) protected Stringprotected io.fluxzero.common.serialization.compression.CompressionAlgorithmprotected StringgetNamespace(ServerWebsocketSession session) protected StringgetRequestHeaders(ServerWebsocketSession session) protected Stringprotected voidhandleMessage(ServerWebsocketSession session, io.fluxzero.common.api.JsonType message) voidonClose(ServerWebsocketSession session, io.fluxzero.sdk.common.websocket.WebsocketCloseReason closeReason) voidonError(ServerWebsocketSession session, Throwable e) voidonMessage(byte[] bytes, ServerWebsocketSession session) voidonOpen(ServerWebsocketSession session) voidonPong(ByteBuffer message, ServerWebsocketSession session) protected voidregisterMetrics(io.fluxzero.common.api.JsonType event, ServerWebsocketSession session) protected WebsocketEndpoint.PingRegistrationschedulePing(ServerWebsocketSession session) protected voidsendPing(ServerWebsocketSession session) protected voidsendResultBatch(ServerWebsocketSession session, List<io.fluxzero.common.api.RequestResult> results) protected io.fluxzero.common.api.MetadatasessionMetadata(ServerWebsocketSession session) protected booleanshouldHandleIdempotently(io.fluxzero.common.api.Command command) protected voidshutDown()Close all sessions on the websocket after an optional delay.protected booleansubmitRequestTask(ServerWebsocketSession session, io.fluxzero.common.api.JsonType request, Runnable task)
-
Field Details
-
pingTimeout
-
pingDelay
-
shuttingDown
-
shutDown
protected volatile boolean shutDown
-
-
Constructor Details
-
WebsocketEndpoint
protected WebsocketEndpoint() -
WebsocketEndpoint
-
WebsocketEndpoint
-
WebsocketEndpoint
protected WebsocketEndpoint(Executor requestExecutor, CommandIdempotencyStore commandIdempotencyStore)
-
-
Method Details
-
onOpen
-
deserializeRequest
protected io.fluxzero.common.api.JsonType deserializeRequest(ServerWebsocketSession session, byte[] bytes) -
onMessage
-
dispatchRequest
protected void dispatchRequest(ServerWebsocketSession session, io.fluxzero.common.api.JsonType request) -
shouldHandleIdempotently
protected boolean shouldHandleIdempotently(io.fluxzero.common.api.Command command) -
submitRequestTask
protected boolean submitRequestTask(ServerWebsocketSession session, io.fluxzero.common.api.JsonType request, Runnable task) -
handleMessage
protected void handleMessage(ServerWebsocketSession session, io.fluxzero.common.api.JsonType message) -
doSendResult
protected void doSendResult(ServerWebsocketSession session, io.fluxzero.common.api.RequestResult result) -
createTasks
protected Stream<Runnable> createTasks(io.fluxzero.common.api.RequestBatch<?> batch, ServerWebsocketSession session) -
sendResultBatch
protected void sendResultBatch(ServerWebsocketSession session, List<io.fluxzero.common.api.RequestResult> results) -
findAlternativeBacklog
protected Optional<WebsocketEndpoint.SessionBacklog> findAlternativeBacklog(ServerWebsocketSession closedSession) -
schedulePing
-
sendPing
-
onPong
-
abort
-
onClose
public void onClose(ServerWebsocketSession session, io.fluxzero.sdk.common.websocket.WebsocketCloseReason closeReason) -
onError
-
shutDown
protected void shutDown()Close all sessions on the websocket after an optional delay. During the delay we don't handle new requests but will be able to send back results. -
getCompressionAlgorithm
protected io.fluxzero.common.serialization.compression.CompressionAlgorithm getCompressionAlgorithm(ServerWebsocketSession session) -
getRequestHeaders
-
getNamespace
-
getClientId
-
getClientName
-
getClientSdkVersion
-
getRuntimeVersion
-
getNegotiatedSessionId
-
registerMetrics
protected void registerMetrics(io.fluxzero.common.api.JsonType event, ServerWebsocketSession session) -
sessionMetadata
-