From ac773d97e944a350483ed9d2f18f7176d2412697 Mon Sep 17 00:00:00 2001 From: rstoyanchev Date: Tue, 29 Apr 2025 16:43:27 +0300 Subject: [PATCH] Polishing contribution Closes gh-34828 --- .../web/reactive/socket/WebSocketHandler.java | 99 +++++++++---------- 1 file changed, 45 insertions(+), 54 deletions(-) diff --git a/spring-webflux/src/main/java/org/springframework/web/reactive/socket/WebSocketHandler.java b/spring-webflux/src/main/java/org/springframework/web/reactive/socket/WebSocketHandler.java index ab6c8aa2ec..ead8608fe8 100644 --- a/spring-webflux/src/main/java/org/springframework/web/reactive/socket/WebSocketHandler.java +++ b/spring-webflux/src/main/java/org/springframework/web/reactive/socket/WebSocketHandler.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2021 the original author or authors. + * Copyright 2002-2025 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -23,78 +23,69 @@ import org.reactivestreams.Publisher; import reactor.core.publisher.Mono; /** - * Handler for a WebSocket session. - * - *

A server {@code WebSocketHandler} is mapped to requests with + * Handler for a WebSocket messages. You can use it as follows: + *

* - *

Use {@link WebSocketSession#receive() session.receive()} to compose on - * the inbound message stream, and {@link WebSocketSession#send(Publisher) - * session.send(publisher)} for the outbound message stream. Below is an - * example, combined flow to process inbound and to send outbound messages: + *

{@link WebSocketSession#receive() session.receive()} handles inbound + * messages, while {@link WebSocketSession#send(Publisher) session.send} + * sends outbound messages. Below is an example of handling inbound messages + * and responding to every message: * *

- * class ExampleHandler implements WebSocketHandler {
+ *	class ExampleHandler implements WebSocketHandler {
  *
- * 	@Override
- * 	public Mono<Void> handle(WebSocketSession session) {
- *
- * 		Flux<WebSocketMessage> output = session.receive()
- * 			.doOnNext(message -> {
- * 				// This is for side effects such as
- * 				// - Logging incoming messages
- * 				// - Updating some metrics or counters
- * 				// - Performing access checks or validations (non-blocking)
- * 				System.out.println("Got message: " + message.getPayloadAsText());
- * 			})
- * 			.concatMap(message -> {
- * 				// This is where you handle the actual processing of the incoming message. It
- * 				// might involve:
- * 				// - Parsing the message content (e.g., JSON parsing)
- * 				// - Invoking a reactive service (e.g., database, HTTP call, etc.)
- * 				// - Returning a transformed value, typically a Mono<String> or Mono<SomeType>
- * 				// if you're mapping to another data format
- * 				return Mono.just(message.getPayloadAsText());
- * 			})
- * 			.map(value -> {
- * 				// This is where you produce one or more responses for the message
- * 				return session.textMessage("Echo " + value));
- * 			});
- *
- * 		return session.send(output);
- * 	}
- * }
+ *		@Override
+ *		public Mono<Void> handle(WebSocketSession session) {
+ *			Flux<WebSocketMessage> output = session.receive()
+ * 				.doOnNext(message -> {
+ * 					// Imperative calls without a return value:
+ * 					// perform access checks, log, validate, update metrics.
+ * 					// ...
+ * 				})
+ * 				.concatMap(message -> {
+ * 					// Async, non-blocking calls:
+ * 					// parse messages, call a database, make remote calls.
+ * 					// Return the same message, or a transformed value
+ * 					// ...
+ * 				});
+ * 			return session.send(output);
+ *		}
+ *	}
  * 
* *

If processing inbound and sending outbound messages are independent * streams, they can be joined together with the "zip" operator: * *

- * class ExampleHandler implements WebSocketHandler {
+ *	class ExampleHandler implements WebSocketHandler {
  *
- * 	@Override
- * 	public Mono<Void> handle(WebSocketSession session) {
+ *		@Override
+ *		public Mono<Void> handle(WebSocketSession session) {
  *
- * 		Mono<Void> input = session.receive()
- *			.doOnNext(message -> {
- * 				// ...
- * 			})
- * 			.concatMap(message -> {
- * 				// ...
- * 			})
- * 			.then();
+ *			Mono<Void> input = session.receive()
+ *				.doOnNext(message -> {
+ * 					// ...
+ * 				})
+ * 				.concatMap(message -> {
+ * 					// ...
+ * 				})
+ * 				.then();
  *
- *		Flux<String> source = ... ;
- * 		Mono<Void> output = session.send(source.map(session::textMessage));
+ *			Flux<String> source = ... ;
+ * 			Mono<Void> output = session.send(source.map(session::textMessage));
  *
- * 		return Mono.zip(input, output).then();
- * 	}
- * }
+ * 			return Mono.zip(input, output).then();
+ *		}
+ *	}
  * 
* *

A {@code WebSocketHandler} must compose the inbound and outbound streams