@@ -216,17 +216,14 @@ private HttpServletStreamableServerTransportProvider(McpJsonMapper jsonMapper, S
216216 * down
217217 */
218218 private Disposable startSessionSweeper (Duration sessionSweepInterval ) {
219- return Flux .interval (sessionSweepInterval , sessionSweepInterval , Schedulers .boundedElastic ()).doOnNext (tick -> {
220- // Each sweep runs in its own subscription, so that a sweep failing does
221- // not terminate the interval and disable sweeping altogether
222- Mono .fromRunnable (this ::sweepSessions )
223- .doOnError (e -> logger .error ("Session sweep failed" , e ))
224- .onErrorComplete ()
225- .subscribe ();
226- }).onErrorComplete (error -> {
227- logger .error ("Session sweeper error" , error );
228- return true ;
229- }).subscribe ();
219+ return Flux .interval (sessionSweepInterval , sessionSweepInterval , Schedulers .boundedElastic ())
220+ .concatMap (
221+ tick -> sweepSessions ().doOnError (e -> logger .error ("Session sweep failed" , e )).onErrorComplete ())
222+ .onErrorComplete (error -> {
223+ logger .error ("Session sweeper error" , error );
224+ return true ;
225+ })
226+ .subscribe ();
230227 }
231228
232229 /**
@@ -235,25 +232,23 @@ private Disposable startSessionSweeper(Duration sessionSweepInterval) {
235232 * if it received a request during the interval which just elapsed. Anything else is a
236233 * session whose client went away without deleting it: the protocol lets a client
237234 * reconnect to a session, so nothing else ever reclaims it.
235+ * @return a Mono completing once every evicted session has been closed
238236 */
239- private void sweepSessions () {
237+ private Mono < Void > sweepSessions () {
240238 if (this .isClosing ) {
241- return ;
239+ return Mono . empty () ;
242240 }
243241 Set <String > active = this .activeSessions .getAndSet (ConcurrentHashMap .newKeySet ());
244- this .sessions .values ().removeIf (session -> {
245- if (session .hasOpenStream () || active .contains (session .getId ())) {
246- return false ;
247- }
248- logger .debug ("Evicting idle session {}" , session .getId ());
249- try {
250- session .closeGracefully ().block ();
251- }
252- catch (Exception e ) {
253- logger .warn ("Failed to close idle session {}: {}" , session .getId (), e .getMessage ());
254- }
255- return true ;
256- });
242+ return Flux .fromIterable (this .sessions .values ())
243+ .filter (session -> !session .hasOpenStream () && !active .contains (session .getId ()))
244+ .filter (session -> this .sessions .remove (session .getId (), session ))
245+ .flatMap (session -> {
246+ logger .debug ("Evicting idle session {}" , session .getId ());
247+ return session .closeGracefully ()
248+ .doOnError (e -> logger .warn ("Failed to close idle session {}: {}" , session .getId (), e .getMessage ()))
249+ .onErrorComplete ();
250+ })
251+ .then ();
257252 }
258253
259254 /**
0 commit comments