In a reactive WebSocketHandler, what exactly controls when the session closes, and how do close, errors, and cleanup propagate through the returned Mono<Void>?
answer
- Session open == returned Mono unterminated
- complete→normal close, error→error close, cancel on shutdown
- receive() completes on client close → drives your Mono
- Infinite send: merge receive().then() to drain + detect close
- doFinally = your afterConnectionClosed; session.close(status) proactive
basics
~20 sThe session stays open while the Mono<Void> you return is unterminated. When that Mono completes, the framework closes the session normally; when it errors, the session closes with an error status. The client disconnecting completes receive(), which typically completes your Mono.
solid answer
~40 sThe connection lifetime equals the lifetime of the Mono<Void> returned from handle(). The framework subscribes to it: normal completion → the framework sends a close frame (normal status); an error signal → close with an error/close status. Because receive() completes when the client sends a close frame or the transport ends, a pipeline built on receive() naturally terminates on client disconnect. If your outbound stream is infinite (e.g., Flux.interval) and you don't also depend on receive(), the session stays open until the client leaves or the transport drops — so you should merge receive().then() to react to client close and to drain/release inbound buffers. For cleanup, attach doFinally/doOnError/doOnCancel to the pipeline. You can also close proactively with session.close(CloseStatus). Cancellation (server shutdown) propagates as a cancel signal down the stream.
code
java · 25 linesimport org.springframework.web.reactive.socket.*;
import reactor.core.publisher.*;
import java.time.Duration;
public class PushHandler implements WebSocketHandler {
private final SessionRegistry registry;
public PushHandler(SessionRegistry registry) { this.registry = registry; }
@Override
public Mono<Void> handle(WebSocketSession session) {
registry.add(session);
Flux<WebSocketMessage> push = Flux.interval(Duration.ofSeconds(1))
.map(i -> session.textMessage("tick-" + i));
Mono<Void> out = session.send(push);
Mono<Void> drain = session.receive().then(); // drains inbound, completes on client close
// End the session as soon as EITHER the client leaves or the push errors.
return Mono.firstWithSignal(out, drain)
.doFinally(sig -> registry.remove(session)); // cleanup on complete/error/cancel
}
interface SessionRegistry { void add(WebSocketSession s); void remove(WebSocketSession s); }
}go deeper
Know completing the returned Mono closes the session and client disconnect completes receive().
Handle client-close by depending on receive() and add doFinally cleanup.
Reason about infinite-outbound traps, error→close-status mapping, and proactive session.close(status).
Design robust lifecycle: drain semantics, timeout/heartbeat strategy, cancellation on shutdown, registry cleanup.
## The core rule The framework **subscribes to the `Mono<Void>` you return** and keeps the WebSocket open until that Mono **terminates**: - **onComplete** → framework closes the session with a **normal** close status (1000). - **onError** → framework closes the session with an error/close status (e.g., 1011 / server error), and the error is logged. - **Cancel** (e.g., the server is shutting down or the connection is being torn down) → propagates as a `cancel` signal through your reactive chain. So "when does the session close?" == "when does my returned Mono terminate?" ## How client disconnect flows in When the client sends a WebSocket **close frame** or the TCP connection drops: - `session.receive()` **completes** (or errors on abnormal drop). - Any pipeline whose completion depends on `receive()` then completes → your returned Mono completes → framework finalizes close. This is why an echo handler `session.send(session.receive().map(...))` closes cleanly on client disconnect: `receive()` completes → `send()`'s source completes → `send()`'s `Mono<Void>` completes. ## The infinite-outbound trap If you return only an infinite outbound stream: ```java return session.send(Flux.interval(Duration.ofSeconds(1)).map(i -> session.textMessage("" + i))); ``` The `send` Mono never completes on its own. The session ends only when the client disconnects (the write fails / transport signals) or you `session.close()`. Two issues: 1. **Inbound buffers**: if the client also sends frames, you never subscribed to `receive()`, so those frames aren't drained; pooled `DataBuffer`s can accumulate. Best practice: merge `session.receive().then()` even for push-only handlers so inbound frames are drained and released, and so a client close terminates your pipeline promptly. 2. **Prompt close detection**: depending on the server, an infinite outbound may not notice client close until the next write. Depending on `receive()` completion makes close detection immediate. Robust push pattern: ```java Mono<Void> out = session.send(pushFlux); Mono<Void> drain = session.receive().then(); return Mono.firstWithSignal(out, drain); // end as soon as either ends ``` ## Errors - An error anywhere in your chain becomes the Mono's `onError`, closing the session with an error status. Guard against expected transient errors with `onErrorResume` if you want to send a graceful close message first. - Uncaught errors are logged by the framework. ## Proactive close `session.close()` / `session.close(CloseStatus)` returns a `Mono<Void>` that closes the session with a specific status (e.g., `CloseStatus.POLICY_VIOLATION`, `CloseStatus.GOING_AWAY`). Use it to reject or terminate: e.g., after auth failure or an idle timeout you compose `session.close(CloseStatus.POLICY_VIOLATION)` into the returned Mono. ## Cleanup hooks Because there is no `afterConnectionClosed` callback, use reactive operators for teardown: - `doFinally(signalType -> ...)` — runs on complete, error, or cancel; ideal for removing the session from a registry, decrementing counters, releasing resources. - `doOnError(...)`, `doOnCancel(...)` — signal-specific hooks. ```java return session.send(out) .and(session.receive().then()) .doFinally(sig -> registry.remove(session.getId())); ``` ## Idle / timeouts and heartbeats Raw reactive WebSockets don't auto-heartbeat like STOMP. You can add `timeout(Duration)` on the pipeline to close idle sessions, or emit periodic `session.pingMessage(...)` in the outbound stream and rely on client PONGs. Beware `timeout` errors close with an error status unless you `onErrorResume` to a graceful close. ## Summary mental model One stream in, one stream out, **one returned Mono** = the whole session. Its terminal signal is the close trigger; `receive()` completion is how a client disconnect reaches you; `doFinally` is your close callback.
- You return only an infinite Flux.interval to send(); the client sends data too. What problems arise?You never subscribed to receive(), so inbound frames aren't drained and their pooled DataBuffers can accumulate, and client-close may not be detected promptly. Merge session.receive().then() so inbound is drained and disconnect terminates the pipeline.
- There is no afterConnectionClosed callback — how do you run cleanup when the session ends?Attach doFinally(signal -> ...) to the returned pipeline; it runs on complete, error, and cancel. Use doOnError/doOnCancel for signal-specific handling.
- How do you close a session with a specific status code?Compose session.close(CloseStatus.POLICY_VIOLATION) (or GOING_AWAY, etc.) into the returned Mono; it emits a close frame with that status.
saying these in an interview costs you the question
- Believing the session closes on its own regardless of the returned Mono
- Returning an infinite send() without draining receive()
- Expecting an afterConnectionClosed callback in the reactive model
- Thinking an error in the pipeline is swallowed rather than closing the session
- Assuming raw reactive WebSockets auto-heartbeat like STOMP