Xj%$UdZddlZddlZddlZddlmZmZddlmZm Z m Z ddl m Z ddl mZddlmZddlmZdd lmZmZddlZdd lmZmZdd lmZdd lmZdd lmZddl m!Z!ddl"m#Z#m$Z$m%Z%ddl&m'Z'm(Z(ddl)m*Z*m+Z+ddl,m-Z-ddl.m/Z/m0Z0m1Z1m2Z2m3Z3m4Z4m5Z5m6Z6m7Z7m8Z8m9Z9ejte;ZdZ?dZ@dZAdZBdZCeeDd<ejdZFeGZHeGZIeJeGefZKeGddZLe eLge dfZMGdd eZNGd!d"ZOy)#z StreamableHTTP Server Transport Module This module implements an HTTP transport layer with Streamable HTTP. The transport handles bidirectional communication using HTTP requests and responses, with streaming support for long-running operations. N)ABCabstractmethod)AsyncGenerator AwaitableCallable)asynccontextmanager) dataclass)partial) HTTPStatus)AnyFinal)MemoryObjectReceiveStreamMemoryObjectSendStream)ValidationError)EventSourceResponse)Request)Response)ReceiveScopeSend)TransportSecurityMiddlewareTransportSecuritySettings)ServerMessageMetadataSessionMessage)SUPPORTED_PROTOCOL_VERSIONS) DEFAULT_NEGOTIATED_VERSIONINTERNAL_ERRORINVALID_PARAMSINVALID_REQUEST PARSE_ERROR ErrorData JSONRPCErrorJSONRPCMessageJSONRPCRequestJSONRPCResponse RequestIdzmcp-session-idzmcp-protocol-versionz last-event-idzapplication/jsonztext/event-stream _GET_streamREQUEST_STREAM_BUFFER_SIZEz^[\x21-\x7E]+$c0eZdZUdZeed<dZedzed<y) EventMessagezM A JSONRPCMessage with an optional event ID for stream resumability. messageNevent_id)__name__ __module__ __qualname____doc__r#__annotations__r-strn/mnt/ssd/data/Dropbox/adrian/vault-secondbrain/venv/lib/python3.12/site-packages/mcp/server/streamable_http.pyr+r+PsHcDjr5r+cXeZdZdZedededzdefdZedede dedzfd Z y) EventStorez? Interface for resumability support via event storage. stream_idr,Nreturnc Kyw)a Stores an event for later retrieval. Args: stream_id: ID of the stream the event belongs to message: The JSON-RPC message to store, or None for priming events Returns: The generated event ID for the stored event Nr4)selfr9r,s r6 store_eventzEventStore.store_eventbs   last_event_id send_callbackc Kyw)a2 Replays events that occurred after the specified event ID. Args: last_event_id: The ID of the last event the client received send_callback: A callback function to send events to the client Returns: The stream ID of the replayed events Nr4)r<r?r@s r6replay_events_afterzEventStore.replay_events_afterps r>) r.r/r0r1rStreamIdr#EventIdr= EventCallbackrBr4r5r6r8r8]sj  8  nt>S  X_      %  D   r5r8c ^eZdZUdZdZeeezdzed<dZ e eezdzed<dZ eedzed<dZ e edzed<e ed< d:dedzd ed edzd edzd edzd df dZed efdZded dfdZd;dZdedededed ef dZdeded edzfdZdedeede ededzd df dZ e!dfdede"ded e#eefdzd e$f d!Z%e"jLdfd"edzde"d e#eefdzd e$fd#Z'ded edzfd$Z(d%ed efd&Z)ded dfd'Z*d(e+d)e,d*e-d dfd+Z.ded e/eeffd,Z0ded efd-Z1ded(e+d*e-d efd.Z2d(e+ded)e,d*e-d df d/Z3ded*e-d dfd0Z4ded*e-d dfd1Z5d;d2Z6ded*e-d dfd3Z7ded*e-d efd4Z8ded*e-d efd5Z9ded*e-d efd6Z:d7eded*e-d dfd8Z;ey)<StreamableHTTPServerTransportz HTTP server transport with event streaming support for MCP. Handles JSON-RPC messages in HTTP POST requests with SSE streaming. Supports optional JSON responses and session management. N_read_stream_writer _read_stream _write_stream_write_stream_reader _securitymcp_session_idis_json_response_enabled event_storesecurity_settingsretry_intervalr:c| tj|s td||_||_||_t ||_||_i|_ i|_ d|_ d|_ y)am Initialize a new StreamableHTTP server transport. Args: mcp_session_id: Optional session identifier for this connection. Must contain only visible ASCII characters (0x21-0x7E). is_json_response_enabled: If True, return JSON responses for requests instead of SSE streams. Default is False. event_store: Event store for resumability support. If provided, resumability will be enabled, allowing clients to reconnect and resume messages. security_settings: Optional security settings for DNS rebinding protection. retry_interval: Retry interval in milliseconds to suggest to clients in SSE retry field. When set, the server will send a retry field in SSE priming events to control client reconnection timing for polling behavior. Only used when event_store is provided. Raises: ValueError: If the session ID contains invalid characters. NzASession ID must only contain visible ASCII characters (0x21-0x7E)F) SESSION_ID_PATTERN fullmatch ValueErrorrMrN _event_storerrL_retry_interval_request_streams_sse_stream_writers _terminated idle_scope)r<rMrNrOrPrQs r6__init__z&StreamableHTTPServerTransport.__init__sy8  %.@.J.J>.Z`a a,(@%'45FG-  WY  48r5c|jS)z7Check if this transport has been explicitly terminated.)rZr<s r6 is_terminatedz+StreamableHTTPServerTransport.is_terminatedsr5 request_idc|jj|d}|r|j||jvr?|jj|\}}|j|jyy)a Close SSE connection for a specific request without terminating the stream. This method closes the HTTP connection for the specified request, triggering client reconnection. Events continue to be stored in the event store and will be replayed when the client reconnects with Last-Event-ID. Use this to implement polling behavior during long-running operations - client will reconnect after the retry interval specified in the priming event. Args: request_id: The request ID whose SSE stream should be closed. Note: This is a no-op if there is no active stream for the request ID. Requires event_store to be configured for events to be stored during the disconnect. N)rYpopcloserX)r<r`writer send_streamreceive_streams r6close_sse_streamz.StreamableHTTPServerTransport.close_sse_streamsp$))--j$?  LLN .. .*.*?*?*C*CJ*O 'K      " /r5c.|jty)aFClose the standalone GET SSE stream, triggering client reconnection. This method closes the HTTP connection for the standalone GET stream used for unsolicited server-to-client notifications. The client SHOULD reconnect with Last-Event-ID to resume receiving notifications. Use this to implement polling behavior for the notification stream - client will reconnect after the retry interval specified in the priming event. Note: This is a no-op if there is no active standalone SSE stream. Requires event_store to be configured for events to be stored during the disconnect. Currently, client reconnection for standalone GET streams is NOT implemented - this is a known gap (see test_standalone_get_stream_reconnection). N)rgGET_STREAM_KEYr^s r6close_standalone_sse_streamz9StreamableHTTPServerTransport.close_standalone_sse_streams" n-r5r,requestprotocol_versioncjr!|dk\rdfd }dfd }t|||}n t|}t||S)aJCreate a session message with metadata including close_sse_stream callback. The close_sse_stream callbacks are only provided when the client supports resumability (protocol version >= 2025-11-25). Old clients can't resume if the stream is closed early because they didn't receive a priming event. 2025-11-25c0KjywN)rg)r`r<sr6close_stream_callbackzTStreamableHTTPServerTransport._create_session_message..close_stream_callbacks%%j1sc.Kjywrp)rjr^sr6 close_standalone_stream_callbackz_StreamableHTTPServerTransport._create_session_message..close_standalone_stream_callback s002s)request_contextrgrjrtmetadatar:N)rVrr)r<r,rkr`rlrqrsrws` ` r6_create_session_messagez5StreamableHTTPServerTransport._create_session_messagesN   !1\!A 2 3- '!6,LH -WEHg99r5r9cK|jsy|dkry|jj|dd{}|dd}|j|j|d<|S7&w)aoStore the priming cursor for `stream_id` and return its SSE wire form. Called before the request is dispatched so the priming row precedes anything `message_router` can store for this stream. Returns `None` when no event store is configured or the client predates 2025-11-25 (older clients cannot parse the empty-data event). Nrn)iddataretry)rVr=rW)r<r9rlpriming_event_id priming_events r6_mint_priming_eventz1StreamableHTTPServerTransport._mint_priming_eventsm   l *!%!2!2!>!>y$!OO)92"F    +%)%9%9M' " Ps3AA'Asse_stream_writerrequest_stream_readerrcK |4d{|4d{||j|d{|23d{}|j|j|d{t|jjt t zs]dddd{dddd{tjd|jj|d|j|d{y777776w7k#1d{7swY{xYw7r#1d{7swYxYw#tj$rtjdYt$rtjdYwxYw7#tjd|jj|d|j|d{7wxYww)zFForward `_request_streams[request_id]` onto the SSE wire for one POST.Nz'SSE stream closed by close_sse_stream()zError in SSE writerzClosing SSE writer)send_create_event_data isinstancer,rootr%r"anyioClosedResourceErrorloggerdebug Exception exceptionrYrb_clean_up_memory_streams)r<r`rrr event_messages r6_run_sse_writerz-StreamableHTTPServerTransport._run_sse_writer(s <(  *?   ,+00???+@-+001H1H1WXXX!-"7"7"<"\]      LL- .  $ $ ( (T :// ; ; ;  ?X,A        (( D LLB C 4   2 3 4 < LL- .  $ $ ( (T :// ; ; ;s2GD:DD:D%DD%DDDD D D $D#D $/DD D% D !D%% D:0D#1D:5AG:F;GD:D%DD D D D%D DD D%#D:%D7+D. ,D73D::(F"F $FF FF G AGGGG error_message status_code error_codeheaderscdti}|r|j||jr|j|t<t ddt ||}t |jdd||S) z6Create an error response with a simple string message. Content-Typez2.0z server-error)coder,)jsonrpcr|errorTby_alias exclude_nonerr)CONTENT_TYPE_JSONupdaterMMCP_SESSION_ID_HEADERr"r!rmodel_dump_json)r<rrrrresponse_headerserror_responses r6_create_error_responsez4StreamableHTTPServerTransport._create_error_responseAs+,=>   # #G ,   6:6I6I 2 3&%   * *Dt * L#$  r5response_messagecdti}|r|j||jr|j|t<t |r|j ddnd||S)z,Create a JSON response from a JSONRPCMessagerTrNr)rrrMrrr)r<rrrrs r6_create_json_responsez3StreamableHTTPServerTransport._create_json_response`sh+,=>   # #G ,   6:6I6I 2 3Rb  , ,d , Nhl#$  r5c@|jjtS)z,Extract the session ID from request headers.)rgetr)r<rks r6_get_session_idz-StreamableHTTPServerTransport._get_session_idts""#899r5rc|d|jjddd}|jr|j|d<|S)z2Create event data dictionary from an EventMessage.r,Tr)eventr}r|)r,rr-)r<r event_datas r6rz0StreamableHTTPServerTransport._create_event_dataxsI!))994VZ9[  ! !,55Jt r5cK||jvrn |j|djd{|j|djd{|jj |dyy7J7$#t$rtj dYBwxYw#|jj |dwxYww)z/Clean up memory streams for a given request ID.rNz4Error closing memory streams - may already be closed)rXacloserrrrb)r<r`s r6rz6StreamableHTTPServerTransport._clean_up_memory_streamss .. . <++J7:AACCC++J7:AACCC %%))*d; /DC U ST U %%))*d;sVC #BB'BBB"C BBB%"B($B%%B((CC scopereceivercKt||}|jdk(}|jj||d{}|r||||d{y|jr3|j dt j}||||d{y|jdk(r|j||||d{y|jdk(r|j||d{y|jdk(r|j||d{y|j||d{y7777z7R7*7w)z6Application entry point that handles all HTTP requestsPOST)is_postNz&Not Found: Session has been terminatedGETDELETE) rmethodrLvalidate_requestrZrr NOT_FOUND_handle_post_request_handle_get_request_handle_delete_request_handle_unsupported_request)r<rrrrkrrresponses r6handle_requestz,StreamableHTTPServerTransport.handle_requests6%)..F*#~~>>wPW>XX  6 6 6    228$$H5'40 0 0  >>V #++E7GTJ J J ^^u $**7D9 9 9 ^^x '--gt< < <227DA A A+Y 6 1 K 9 < As{ED6+E=D8>)E'D:()ED<E,D>-E4E6E8E:E<E>Ec|jjdd}|jdDcgc]}|j}}t d|D}t d|D}||fScc}w)z6Check if the request accepts the required media types.acceptr{,c3FK|]}|jtywrp) startswithr.0 media_types r6 zFStreamableHTTPServerTransport._check_accept_headers..s_Jz,,->?_!c3FK|]}|jtywrp)rCONTENT_TYPE_SSErs r6rzFStreamableHTTPServerTransport._check_accept_headers..s]*j++,<=]rrrsplitstripany)r<rk accept_headerr accept_typeshas_jsonhas_sses r6_check_accept_headersz3StreamableHTTPServerTransport._check_accept_headersss++Hb9 =J=P=PQT=UVz ((*V V_R^__]P\]]  WsA1c|jjdd}|jddjdDcgc]}|j}}t d|DScc}w)z2Check if the request has the correct Content-Type.z content-typer{;rrc3.K|] }|tk(ywrp)r)rparts r6rzDStreamableHTTPServerTransport._check_content_type..sL4,,Lsr)r<rk content_typercontent_type_partss r6_check_content_typez1StreamableHTTPServerTransport._check_content_typesf**>2> 7C7I7I#7Nq7Q7W7WX[7\]tdjjl]]L9KLLL^sA-cVK|j|\}}|jr@|s=|jdtj}|||j |d{yy|r|s=|jdtj}|||j |d{yy7H7w)zEValidate Accept header based on response mode. Returns True if valid.z3Not Acceptable: Client must accept application/jsonNFzNNot Acceptable: Client must accept both application/json and text/event-streamT)rrNrr NOT_ACCEPTABLEr)r<rkrrrrrs r6_validate_accept_headerz5StreamableHTTPServerTransport._validate_accept_headers 66w?'  ( (66I--ugoot<<<w22`))H5'//48 8 8= 9s%AB)B%AB)B'B)'B)c K|j}| td |j|||d{sy|j|s3|j dt j }||||d{y|jd{} tj|} tj|} t%| j&t(xr| j&j*dk(} | ra|j,rp|j/|} | r]| |j,k7rN|j dt j0}||||d{y|j3||d{syt%| j&t(se|j5dt j6}||||d{t9|} t;| | }|j=|d{y| rI| j&j>r3t| j&j>jAd tBn#|jDjAtFtB}t| j&jH}|jJrvtMjNtPtR|jT|<|jT|d }t9|} t;| | }|j=|d{ d}|23d{}t%|jVj&tXtZzr|jV}n7t\j_d |jVj&j*y|r$|j5|}||||d{nGt\jad |j dt jb}||||d{|jk|d{y|jm||d{}tMjNtnd\}}||jp|<tMjNtPtR|jT|<|jT|d }ddtrd|j,rtt|j,ini}tw|ty|jz|||||} tMj|4d{}|j|||||j| |||}|j=|d{dddd{yy777#tj$rN} |j dt| t jt}||||d{7Yd} ~ yd} ~ wwxYw#t $rN} |j dt| t jt"}||||d{7Yd} ~ yd} ~ wwxYw777]7.7$7677E#td$rQt\jgd|j dt jbth}||||d{7YwxYw7#|jk|d{7wxYw777~7q#1d{7swYxYw#td$rdt\jgd|jd{7|jd{7|jk|d{7YywxYw#td$rz}t\jgd|j dt jbth}||||d{7|j=te|d{7Yd}~yd}~wwxYww)z2Handle POST requests containing JSON-RPC messages.NBNo read stream writer available. Ensure connect() is called first.z=Unsupported Media Type: Content-Type must be application/jsonz Parse error: zValidation error: initialize(Not Found: Invalid or expired session IDrurvprotocolVersionrz received: z1No response message received before stream closedz.Error processing request: No response receivedzError processing JSON responsezError processing requestrno-cache, no-transform keep-alivez Cache-Control Connectionrcontentdata_sender_callablerzSSE response errorzError handling POST request)BrHrUrrrr UNSUPPORTED_MEDIA_TYPEbodyjsonloadsJSONDecodeErrorr3 BAD_REQUESTr r#model_validaterrrrr$rrMrr_validate_request_headersrACCEPTEDrrrparamsrrrMCP_PROTOCOL_VERSION_HEADERr|rNrcreate_memory_object_streamr+r)rXr,r%r"rrrINTERNAL_SERVER_ERRORrrrrrSSEEventrYrrrr rcreate_task_group start_soonryr)r<rrkrrrdrr raw_messageer,is_initialization_requestrequest_session_idrwsession_messagerlr`rrrrrsse_stream_readerrtgerrs r6rz2StreamableHTTPServerTransport._handle_post_requests*)) >ab bv 55gudKKK++G466S55ugt444!'D "jj. (77 D7<<8`W\\=P=PT`=` &)&&)-)=)=g)F&*.@DDWDW.W#'#>#>F&00$'ugt<<<99'4HHHgllN;55''ugt4441I"08"Lkk/222-1D1DGLL''++,=?YZ[__(()DF`a  W\\__-J,,494U4UVb4c.5%%j1)-(=(=j(I!(L%0I"08"Lkk/222%D(,$0E[[m%m&;&;&@&@/T`B`a/N#O&ugt<<< %XY#'#>#>L&<<$'ugt<<<77 CCC '+&>&>zK[&\ \ 7<7X7XYa7bcd7e4!#47H((4494U4UVb4c.5%%j1)-(=(=j(I!(L%&>".$4HLGZGZ-t/B/BC`b  /-)0,,j:KMbdq*$  D$668;;B hwE*.*F*FwPWYceu*v$kk/::: ;;;AL5('' 66s1vh7OQ[QgQgituugt444 # 66(Q1**" ugt444 4=H5 3.3[/D== 9$$%EF#::2"88& H #5'48889D$77 CCC !]4;; ;;;; !D$$%9:+22444+2244477 CCC D     : ;22-00H 5'40 0 0++in- - -  s_\>U \>_>\>7U8\><_=\>U\>U,V8B\>X\>_\>.X/\>3_4A\>;X<1\>-X.\>2_3D\>:X;\>X-X$X! X$ BX-X'AX-%X*&X-*\>>Z ?\>_\>Z,C\>[3Z/4[7=Z84Z25Z89 [Z5[ _ \>\>\>V5'>V0%V(&V0+\>/_0V55\>8 X>X ?XX \> _ X\>\>\>\>\>!X$$X-'X-*X--AZ=Z>ZZ ZZ \> Z)"Z%#Z))\>/[2Z85[8[ >[ ?[ [1\;?\\;\\;2\53\;8\>9_:\;;\>> _A^<^"^<1^42^<7_<__cd Kj}| tdj|\}}|sGjdtj }||j |j|d{yj||d{sy|jjtx}rj|||d{yddtd}jrj|t<t j"vrGjdtj$}||j |j|d{yt'j(t*d\ } fd } t-| | | } ||j |j|d{y7;7#77h7#t.$rht0j3d  j5d{7| j5d{7j7t d{7YywxYww) a Handle GET request to establish SSE. This allows the server to communicate to the client without the client first sending data via HTTP POST. The server can send JSON-RPC requests and notifications on this stream. Nrz4Not Acceptable: Client must accept text/event-streamrrrz4Conflict: Only one SSE stream is allowed per sessionrcK tjttjt <jt d}4d{|4d{|23d{}j |}j|d{47E7<737 6dddd{7n#1d{7swYnxYwdddd{7n#1d{7swYnxYwn$#t$rtjdYnwxYwtjdjt d{7y#tjdjt d{7wxYww)NrzError in standalone SSE writerzClosing standalone SSE writer) rrr+r)rXrirrrrrrr)standalone_stream_readerrrr<rs r6standalone_sse_writerzPStreamableHTTPServerTransport._handle_get_request..standalone_sse_writersi D9>8Y8YZf8g.9%%n5,0+@+@+PQR+S(, A A.F A A/GAAm&*%<%<]%K /44Z@@@ A AAA0H A A A A A A A A A A C  !AB C <=33NCCC <=33NCCCs FA C0BC0CBC!B3$B (B )B ,&B3B B3C0CB B3 B3! C,B/-C3C 9B<:C C C0CC0C, C# !C,(C0/E0DEDE-FEF.E=6E97E==Frz Error in standalone SSE response)rHrUrrr rrrrrrLAST_EVENT_ID_HEADER_replay_eventsrrMrrirXCONFLICTrrrrrrrrr) r<rkrrd_rrr?rrrrs ` @r6rz1StreamableHTTPServerTransport._handle_get_requests)) >ab b//8 722F))H7=='//4@ @ @ 33GTBBB $OO//0DE E= E%%mWdC C C 6&,    -1-@-@G) * T22 222F##H7=='//4@ @ @ 05/P/PQY/Z[\/],, D6'%!6   @7=='//4@ @ @G AC D$ AR A @   ? @#**, , ,#**, , ,//? ? ?  @sA1H05F06H0F3:><>>cK|jsy|j|}|sG|jdtj}||j |j |d{y||jk7rG|jdtj}||j |j |d{yy7\7w)z'Validate the session ID in the request.TzBad Request: Missing session IDNFr)rMrrr rrrr)r<rkrrrs r6rz/StreamableHTTPServerTransport._validate_sessionCs"""11':"221&&H7=='//4@ @ @ !4!4 422:$$H7=='//4@ @ @ A As%A"C$C%AC:C;CCc:K|jjt}|t}|tvrfdj t}|j d|dd|ztj}||j|j|d{yy7w)z4Validate the protocol version header in the request.Nz, z+Bad Request: Unsupported protocol version: z. zSupported versions: FT) rrrrrjoinrr rrr)r<rkrrlsupported_versionsrs r6rz8StreamableHTTPServerTransport._validate_protocol_version`s#??../JK  #9  #> >!%+F!G 22=>N=OrR();(<=>&&H 7=='//4@ @ @ AsBBBBr?c Kjsy ddtd}jrj|t<|jj t t tjtd\ } fd}t|||} ||j|j|d{ j#d{|j#d{y75#t$rtj!dYSwxYw7C7-# j#d{7|j#d{7wxYw#t$rdtj!d j%d t&j(t*}||j|j|d{7YywxYww) z Replays events that would have been sent after the specified event ID. Only used when resumability is enabled. Nrrrrc\K 4d{dtddf fd }j|d{}|r| jvr j|< j |d{}| j |d{t jtt j|< j|d}|4d{|23d{} j|} j |d{4dddd{y77777T7K7#6dddd{7n#1d{7swYnxYw jj|d j|d{7# jj|d j|d{7wxYw7#1d{7swYyxYw#t j$rtjdYyt$rtj!dYywxYww)Nrr:cfKj|}j|d{y7wrp)rr)rrr<rs r6 send_eventzWStreamableHTTPServerTransport._replay_events..replay_sender..send_events+)-)@)@)OJ"3"8"8"DDDs &1/1rz.Replay SSE stream closed by close_sse_stream()zError in replay sender)r+rBrXrYrrrrr)rrbrrrrrr) r%r9r msg_readerrrrOr?replay_protocol_versionr<rs r6 replay_senderzCStreamableHTTPServerTransport._replay_events..replay_senders:*?0$O$OELETE +6*I*I-Yc*d$d %$:O:O)OOFW 8 8 C 7;6N6NyZq6r0r #0#<*;*@*@*O$O$ODICdCdeqCr$>D" 5 5i @.2-B-B9-Ma-P ,6!Q!Q?I%Q%Qm595L5L]5[ .?.D.DZ.P(P(PC$O$O$O%e1s$O!Q%Q)Q@J!Q!Q!Q!Q!Q !% 8 8 < ?sbH,GDG%G D"G  $F .D$/F  D& A F D(F ED.D* D."&ED, E GGGH,G"G $F &F (F *D.,E.E/ F :D=;F E E E F 0G F G  1G>G ?GG G GG GGH,G(H)H, H)&H,(H))H,rzError in replay responsezError replaying events)rVrrMrrrrrrrrrrrrrrrrr rr) r<r?rkrrrr(rrOr'rs `` @@@r6r z,StreamableHTTPServerTransport._replay_eventsvs ''  S A!9* 0G ""151D1D-.'.oo&9&9:UWq&r #493T3TU]3^_`3a 0 0+ ?+ ?\+)%2H  1w}}gootDDD(..000'..000 E =  !;< =10(..000'..000 A   5 622(00H 7=='//4@ @ @ AsGBE!C5=C3>C5EDE-D.E2G3C55DDDDEEE1D42E E  EEA$G8F;9G>GGGcnKtjttzd\}}tjtd\}|_|__|_tj4d{}fd}|j| ||ftjjD]}j|d{jj |jd{|jd{jd{|jd{dddd{y777a7K757#t$r"}t j#d|Yd}~Cd}~wwxYw#tjjD]}j|d{7jj |jd{7|jd{7jd{7|jd{7w#t$r"}t j#d|Yd}~wd}~wwxYwxYw74#1d{7swYyxYww)zContext manager that provides read and write streams for a connection. Yields: Tuple of (read_stream, write_stream) for bidirectional communication rNcfK 23d{}|j}d}t|jttzr"t |jj }|}n[|jOt|jtr5|jjt |jj}||nt}d}jr?jj||d{}tjd|d||jvr6 j|dj!t#||d{Ptjd|dk7g77&#t$j&t$j(f$r jj+|dYwxYw6y#t$j($r;j,rtjdYytj/dYyt0$rtj/dYywxYww) NzStored z from rzRequest stream z not found for message. Still processing message as the client might reconnect and replay.zRead stream closed by clientz3Unexpected closure of read stream in message routerzError in message router)r,rrr%r"r3r|rwrrelated_request_idrirVr=rrrXrr+rBrokenResourceErrorrrbrZrr)rr,target_request_id response_idrequest_stream_idr-r<write_stream_readers r6message_routerz=StreamableHTTPServerTransport.connect..message_routers6@1D..o"1"9"9,0)%gllOl4RS*-glloo*>K1<-,44@ * / 8 8 5!!0 8 8 K K W03O4L4L4_4_0`-ARA^,=dr) $(,,-1->->-J-JK\^e-f'fH"LL78*FCTBU)VW,0E0EES&*&;&;?@sH1GF?E2F?C$G3E54-G".E9E7E9G2F?5G7E99?F<8G;F<<GH14H.5H17H. H1H.+H1-H..H1r)rrrrrHrIrKrJrrrrXrrrrrr) r<read_stream_writer read_stream write_streamrr1r9rr0s ` @r6connectz%StreamableHTTPServerTransport.connects$+0*K*KN]fLf*ghi*j'K,1,M,Mn,]^_,`) )$6 '$7!)**,N @N @7 @t MM. ) @!<//!%d&;&;&@&@&B!CCI77 BBBC%%++-@,33555%,,...-44666&--///WN @N @N @FC 6.6/ @LL#:1#!>??@"&d&;&;&@&@&B!CCI77 BBBC%%++-@,33555%,,...-44666&--/// @LL#:1#!>??@YN @N @N @N @sDA:J5>E/?J5J F)":J E1 J >E;E3E;)E5*E;E7E;E9E; J5)J*J51J 3E;5E;7E;9E;; F&F!J !F&&J );J$G' %!JI,H I,3H6 4I, I I,%I( &I,+J, J 5J JJ JJ J5 J2&J) 'J2.J5)FNNNrx)?r.r/r0r1rHrrrr2rIrrJrKrr3boolr8rintr\propertyr_r&rgrjr#rryrCrrr+rrr dictrrrrrrrrrrrtuplerrrrrrrrrrrr rrr5r4r5r6rGrGsJVZ/0JKdRYQUL+NY,FG$NUCGM).9D@GMQ3NCdJQ** */)->B%) .9d .9#'.9 $& .9 5t; .9 d .9 .9` t  #9##8.&::: :  :  :@8sW_bfWf$<<2(;< 9F <  $ <  <:*)-       c3h$&    D#---)-  (4/   c3h$&    (:w:3::    <  2 4    i@i@r5rG)Pr1rloggingreabcrrcollections.abcrrr contextlibr dataclassesr functoolsr httpr typingr r ranyio.streams.memoryrrpydanticr sse_starletterstarlette.requestsrstarlette.responsesrstarlette.typesrrrmcp.server.transport_securityrrmcp.shared.messagerrmcp.shared.versionr mcp.typesrrrrr r!r"r#r$r%r& getLoggerr.rrrr rrrir)r2compilerSr3rCrDr9rr+rEr8rGr4r5r6rPs)  #??*! R$-&(00E:       8 $)4&'& %'E& RZZ 12   S>     ,489 # # L{@{@r5