Skip to content

Instantly share code, notes, and snippets.

@MeiyappanKannappa
Created October 17, 2023 10:19
Show Gist options
  • Select an option

  • Save MeiyappanKannappa/a984b07b58e82107c65ffeffafffa168 to your computer and use it in GitHub Desktop.

Select an option

Save MeiyappanKannappa/a984b07b58e82107c65ffeffafffa168 to your computer and use it in GitHub Desktop.
public Mono<Void> handle(WebSocketSession session) {
UriComponentsBuilder builder = UriComponentsBuilder.fromUri(session.getHandshakeInfo().getUri());
Map<String,String> queryParams = builder.build().getQueryParams().toSingleValueMap();
String uniqueId = queryParams.get("email");
log.info("Websocket connected for id "+uniqueId);
sink.asFlux().subscribe(new WebSocketSessionDataPublisher(uniqueId, session));
Flux<WebSocketMessage> messageFlux = session.receive().share();
Mono<Void> sendPing = session.send(
Flux.interval(Duration.ofSeconds(2),Duration.ofSeconds(2))
.map(aLong -> session.pingMessage(dataBufferFactory -> session.bufferFactory().allocateBuffer())));
Flux<WebSocketMessage> pong = messageFlux.filter(webSocketMessage -> webSocketMessage.getType()==WebSocketMessage.Type.PONG)
.doOnNext(webSocketMessage -> {
log.info("Recieved Pong from "+uniqueId+" "+webSocketMessage.getPayloadAsText());
sink.tryEmitNext(new WSMessage(uniqueId,"Recieved Pong from "+uniqueId+" "+webSocketMessage.getPayloadAsText()));
});
Flux<String> input = messageFlux.filter(webSocketMessage -> webSocketMessage.getType() == WebSocketMessage.Type.TEXT)
.map(WebSocketMessage::getPayloadAsText)
.doOnNext(s -> {
log.info("Recieved message from "+uniqueId);
Sinks.EmitResult emitResult = sink.tryEmitNext(new WSMessage(uniqueId, "This is message for " + uniqueId));
log.info("Emit result status "+emitResult.name()+" "+emitResult.isSuccess());
});
return Flux.merge(sendPing,pong,input).then();
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment