Created
October 17, 2023 09:02
-
-
Save MeiyappanKannappa/7c2928fa3745f195798f471ca7545143 to your computer and use it in GitHub Desktop.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| @Configuration | |
| @EnableWebFlux | |
| @Slf4j | |
| public class SimpleWebSocketHandler implements WebSocketHandler { | |
| @Autowired | |
| Sinks.Many<Object> sink; | |
| @Bean | |
| public SimpleUrlHandlerMapping simpleUrlHandlerMapping() { | |
| Map<String, WebSocketHandler> urlMap = new HashMap<>(); | |
| urlMap.put("/ws", this::handle); | |
| SimpleUrlHandlerMapping handlerMapping = new SimpleUrlHandlerMapping(); | |
| handlerMapping.setUrlMap(urlMap); | |
| handlerMapping.setOrder(1); | |
| return handlerMapping; | |
| } | |
| 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(); | |
| 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(input).then(); | |
| } | |
| private class WebSocketSessionDataPublisher implements Consumer<Object> { | |
| String uniqueId; | |
| WebSocketSession session; | |
| public WebSocketSessionDataPublisher(String uniqueId, WebSocketSession session){ | |
| this.uniqueId = uniqueId; | |
| this.session = session; | |
| } | |
| @Override | |
| public void accept(Object o) { | |
| WSMessage message = (WSMessage) o; | |
| log.info("Message to be sent for id "+message.getUniqueId()); | |
| if(this.uniqueId.equalsIgnoreCase(message.getUniqueId())){ | |
| log.info("[Before Send]Message to be sent for id "+message.getUniqueId()); | |
| session.send(Mono.just(session.textMessage("Filtered Message for "+uniqueId))).subscribe(); | |
| } | |
| } | |
| } | |
| @Data | |
| @AllArgsConstructor | |
| static class WSMessage{ | |
| private String uniqueId; | |
| private Object message; | |
| } | |
| } |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment