4646StreamWriter = ContextSendStream [SessionMessageOrError ]
4747StreamReader = ContextReceiveStream [SessionMessage ]
4848
49+
50+ async def _send_or_ignore_closed (read_stream_writer : StreamWriter , message : SessionMessageOrError ) -> bool :
51+ try :
52+ await read_stream_writer .send (message )
53+ except (anyio .BrokenResourceError , anyio .ClosedResourceError ):
54+ logger .debug ("Read stream closed before Streamable HTTP message could be delivered" , exc_info = True )
55+ return False
56+ return True
57+
58+
4959MCP_SESSION_ID = "mcp-session-id"
5060LAST_EVENT_ID = "last-event-id"
5161
@@ -179,17 +189,17 @@ async def _handle_sse_event(
179189 # Otherwise, return False to continue listening
180190 return isinstance (message , JSONRPCResponse | JSONRPCError )
181191
182- # Forwarding to a closed read stream lands here when the caller cancels mid-SSE
183- # (BrokenResourceError, not a parse failure); coverage is timing-dependent in the
184- # streaming story's modern HTTP cancellation leg.
192+ except ( anyio . BrokenResourceError , anyio . ClosedResourceError ):
193+ logger . debug ( "Read stream closed while forwarding SSE message" , exc_info = True )
194+ return True
185195 except Exception as exc : # pragma: lax no cover
186196 logger .exception ("Error parsing SSE message" )
187197 if original_request_id is not None :
188198 error_data = ErrorData (code = PARSE_ERROR , message = f"Failed to parse SSE message: { exc } " )
189199 error_msg = SessionMessage (JSONRPCError (jsonrpc = "2.0" , id = original_request_id , error = error_data ))
190- await read_stream_writer . send ( error_msg )
200+ await _send_or_ignore_closed ( read_stream_writer , error_msg )
191201 return True
192- await read_stream_writer . send ( exc )
202+ await _send_or_ignore_closed ( read_stream_writer , exc )
193203 return False
194204 else : # pragma: no cover
195205 logger .warning (f"Unknown SSE event: { sse .event } " )
@@ -443,7 +453,7 @@ async def _handle_sse_response(
443453 if last_event_id is None :
444454 error_data = ErrorData (code = CONNECTION_CLOSED , message = "SSE stream disconnected before response completed" )
445455 error_msg = SessionMessage (JSONRPCError (jsonrpc = "2.0" , id = original_request_id , error = error_data ))
446- await ctx .read_stream_writer . send ( error_msg )
456+ await _send_or_ignore_closed ( ctx .read_stream_writer , error_msg )
447457 return
448458
449459 logger .info ("SSE stream disconnected, reconnecting..." )
@@ -467,7 +477,7 @@ async def _handle_reconnection(
467477 data = {"last_event_id" : last_event_id },
468478 )
469479 error_msg = SessionMessage (JSONRPCError (jsonrpc = "2.0" , id = original_request_id , error = error_data ))
470- await ctx .read_stream_writer . send ( error_msg )
480+ await _send_or_ignore_closed ( ctx .read_stream_writer , error_msg )
471481 logger .debug (f"Max reconnection attempts ({ MAX_RECONNECTION_ATTEMPTS } ) exceeded" )
472482 return
473483
0 commit comments