Xj^ ddlZddlmZddlmZddlmZddlmZddl m Z m Z m Z m Z ddlZddlZddlmZmZddlmZdd lmZdd lmZdd lmZmZmZdd lmZdd lm Z m!Z!m"Z"m#Z#m$Z$m%Z%m&Z&m'Z'm(Z(m)Z)m*Z*m+Z+m,Z,m-Z-m.Z.m/Z/m0Z0e de$e/Z1e de%e0Z2e de#e.Z3e de$e/Z4e deZ5e de#e.Z6e7e8zZ9Gdde Z:Gdde e4e2fZ;Gdde e1e3e2e4e6fZ>'99r?c.|jjSr3)rN cancel_calledrTs r6razRequestResponder.cancelleds!!///r?r3)r1rGr1N)r9r:r;r< RequestIdr!Metar(rrrrQrStype BaseExceptionrr[r&rrcrhpropertyboolrjrar4r?r6rArA7s  8-1$((4/!   OPRUUV* 2C}%,C%C$ C  C kI&=$&  :4::0400r?rAc 6eZdZUdZeeeeezfe d<e e d<eee e e ffe d<eeefe d<ede d< d1d eeezd eed ee d eed edzddf dZdeddfdZdefdZdeedzdedzdedzdedzfdZ d2dedee dedzde!dedzde f dZ" d1de#dedzddfdZ$d ed!e e%zddfd"Z&d3d#Z'd$edefd%Z(d&eddfd'Z)d(e e e fddfd)Z*deddfd*Z+ d4d+e,e zd,e-d-e-dzd&e,dzddf d.Z.d/e e e fezezddfd0Z/y)5 BaseSessiona Implements an MCP "session" on top of read/write streams, including features like request/response linking, notifications, and progress. This class is an async context manager that automatically starts processing messages when entered. _response_streams _request_id _in_flight_progress_callbacksr_response_routersN read_stream write_streamreceive_request_typereceive_notification_typeread_timeout_secondsr1c||_||_i|_d|_||_||_||_i|_i|_g|_ t|_ y)Nr) _read_stream _write_streamrvrw_receive_request_type_receive_notification_type_session_read_timeout_secondsrxryrzr _exit_stack)r5r{r|r}r~rs r6rQzBaseSession.__init__s^()!#%9"*C'-A*#% !#)+r?routerc:|jj|y)a Register a response router to handle responses for non-standard requests. Response routers are checked in order before falling back to the default response stream mechanism. This is used by TaskResultHandler to route responses for queued task requests back to their resolvers. WARNING: This is an experimental API that may change without notice. Args: router: A ResponseRouter implementation N)rzappend)r5rs r6add_response_routerzBaseSession.add_response_routers %%f-r?cKtj|_|jjd{|jj |j |S7+wr3)rLcreate_task_group _task_group __aenter__ start_soon _receive_looprTs r6rzBaseSession.__aenter__sS 224))+++ ##D$6$67  ,s7A'A%,A'rUrVrWcK|jjd{|jjj |jj |||d{S7M7wr3)racloser cancel_scoperh __aexit__r\s r6rzBaseSession.__aexit__sd %%''' %%,,.%%//'6JJJ ( Ks"A2A.AA2)A0*A20A2rD result_typerequest_read_timeout_secondsmetadataprogress_callbackc K|j}|dz|_tjttzd\}}||j |<|j ddd} |2d| vri| d<d| dvri| dd<|| ddd<||j|< tdd |d | } |jjtt| | d{d} ||j} n&|j|jj} tj| 5|j!d{} dddt3 trt%| j4|j7| j8|j j;|d|jj;|d|j=d{|j=d{S77#1swYxYw#t"$rJt%t't(j*j,d |j.j0d | dwxYw77n#|j j;|d|jj;|d|j=d{7|j=d{7wxYww)a> Sends a request and wait for a response. Raises an McpError if the response contains an error. If a request read timeout is provided, it will take precedence over the session read timeout. Do not use this method to emit notifications! Use send_notification() instead. Tjsonby_aliasmode exclude_noneNparams_meta progressToken2.0)jsonrpcidr0rz(Timed out while waiting for response to z . Waited z seconds.rfr0r4)rwrLcreate_memory_object_streamrrrv model_dumpryrrsendrr total_secondsr fail_afterreceive TimeoutErrorrrhttpxcodesREQUEST_TIMEOUT __class__r9 isinstanceerrormodel_validateresultpopr) r5rDrrrrrBresponse_streamresponse_stream_reader request_datajsonrpc_requesttimeoutresponse_or_errors r6 send_requestzBaseSession.send_requests %% %>272S2STcfrTr2stu2v//-<z*))4fSW)X  (|+)+ X&l84424 X&w/?IL "7 +O <3DD $ $Z 0( 2,O $$)).P_A`ks*tu u uG+76DDF33?<<JJL %%g.O.D.L.L.N(N%O+\:06677"112C2J2JK  " " & &z4 8  $ $ ( (T :!((* * *(//1 1 1C v)OOO "[[88F&0099:)&iy2   ( + 1  " " & &z4 8  $ $ ( (T :!((* * *(//1 1 1sBKAI!G6?I!H-G;G9G;H?I! A KIK0I1K6I!9G;;HHAII!KK!A K -J0.K K K  K notificationrelated_request_idc Ktd ddi|jddd}tt||r t |nd}|j j |d{y7w) zk Emits a notification, which is a one-way message that does not expect a response. rrTrr)rNrr4)rrrrrrr)r5rrjsonrpc_notificationsession_messages r6send_notificationzBaseSession.send_notification<su 3  %%t&t%T )"#78Ug*>PQmq   %%o666sA"A,$A*%A,rBr]c rKt|trGtd||}tt |}|j j |d{ytd||jddd}tt |}|j j |d{y7^7w)Nrrrrr0Trr)rrr) rrrrrrrrr)r5rBr] jsonrpc_errorrjsonrpc_responses r6rbzBaseSession._send_responseQs h *(:XVM,^M5RSO$$))/: : :.**DvTX*Y   -^DT5UVO$$))/: : : ; ;s%AB7B3AB7-B5.B75B7c Kj4d{j4d{ j23d{}t|trj |d{3t|j j tr jj|j j jddd}t|j j j|j jr |j jjnd|fd|j}|j |j"<j%|d{|j&sj |d{et|j j t:r j<j|j j jddd}t|j t>rU|j jj@}|j vrj |jCd{nt|j tDr|j jjF} | jHvr|jH| } | |j jjJ|j jjL|j jj d{jQ|d{j |d{3jS|d{N7u7c7P7*7!7#t$r}t)j*d|t)j,d|j j t/d|j j jt1t2d d  }t5t7| }jj9|d{7Yd}~(d}~wwxYw779#t$r!}t)jNd|Yd}~\d}~wwxYw7P7:#t$r:}t)j*d|d|j j Yd}~d}~wwxYw7h6nW#tTjV$rt)j,dYn-t$r"}t)jXd|Yd}~nd}~wwxYwt[j\j_D]e\} } t1t`d} | j9t/d| | d{7| jcd{7X#t$rYcwxYwj\jen#t[j\j_D]e\} } t1t`d} | j9t/d| | d{7| jcd{7X#t$rYcwxYwj\jewxYwdddd{7n#1d{7swYnxYwdddd{7y#1d{7swYyxYww)NTrrcPjj|jdSr3)rxrrB)rr5s r6z+BaseSession._receive_loop..tsdoo6I6I!,,X\6]r?)rBrCrDrErFrHzFailed to validate request: z Message that failed validation: rzInvalid request parametersrerrz)Progress callback raised an exception: %sz!Failed to validate notification: z. Message was: zRead stream closed by clientz%Unhandled exception in receive loop: zConnection closedr)3rrr Exception_handle_incomingr0rootrrrrrArrmetarrxrB_received_requestrKloggingwarningdebugrrrrrrrrr requestIdrhr rryr.r/r_received_notification_handle_responserLClosedResourceError exceptionlistrvitemsrrclear)r5r0validated_request respondereerror_responserr cancelled_idprogress_tokencallbackrstreamrs` r6rzBaseSession._receive_loop_s   j /j /   j /j /f /%)%6%6M=M='!'95"33G<<<#GOO$8$8.I"K040J0J0Y0Y ' 4 4 ? ?TZim ? n1-)9+2??+?+?+B+B#4#9#9#@#@.?-C-C-J-J-O-O%)(9(,,]181A1A )IENDOOI,@,@A"&"8"8"CCC#,#7#7&*&;&;I&F F F$$GOO$8$8:MN"+/+J+J+Y+Y ' 4 4 ? ?TZim ? n,L *,*;*;=RS/;/@/@/G/G/Q/Q #/4??#B*.//,*G*N*N*P$P$P$.l.?.?AU#V5A5F5F5M5M5[5[N(69Q9Q'Q373K3KN3[ ).2:0<0A0A0H0H0Q0Q0<0A0A0H0H0N0N0<0A0A0H0H0P0P3.-.-.'+&A&A,&O O O&*&;&;L&I I I#33G<< <= O!!$I!"MNN O#'t'='='C'C'E"FJB%+9W,>UW,:Z( W!V$ "W:V= ;WZ( W Z( W Z(,;Z( Y( Y  Y( !Y$"Y( 'Z( Y4 1Z3Y4 4ZZ( [!Z$"[(Z: .Z1/Z: 6[= [$[  [$[![ [![$ response_idct|tr t|S|S#t$rt j d|dY|SwxYw)a  Normalize a response ID to match how request IDs are stored. Since the client always sends integer IDs, we normalize string IDs to integers when possible. This matches the TypeScript SDK approach: https://github.com/modelcontextprotocol/typescript-sdk/blob/a606fb17909ea454e83aab14c73f14ea45c04448/src/shared/protocol.ts#L861 Args: response_id: The response ID from the incoming message. Returns: The normalized ID (int if possible, otherwise original value). z Response ID z/ cannot be normalized to match pending requests)rr>int ValueErrorrr)r5rs r6_normalize_request_idz!BaseSession._normalize_request_idsV k3 ' o;'' o,{o=l mn os "AAr0cJK|jj}t|ttzsy|j |j }t|tr0|jD] }|j||js yn5|jxsi}|jD]}|j||sy|jj|d}|r|j|d{y|jt!d|d{y7+7w)z Handle an incoming response or error message. Checks response routers first (e.g., for task-related responses), then falls back to the normal response stream mechanism. Nz.Received response with an unknown request ID: )r0rrrrrrrz route_errorrrroute_responservrrrrZ)r5r0rrr response_datars r6rzBaseSession._handle_responses ## $, >? 009  dL )00 %%k4::>  -1KK,=2M00 ((mD  ''++K> ++d# # #'' 7efmen5o(pq q q $ qs0BD#5D#>5D#3D4%D#D!D#!D#rc Kyw)z Can be overridden by subclasses to handle a request without needing to listen on the message stream. If the request is responded to within this method, it will not be forwarded on to the message stream. Nr4)r5rs r6rzBaseSession._received_requestr8c Kyw)z Can be overridden by subclasses to handle a notification without needing to listen on the message stream. Nr4)r5rs r6rz"BaseSession._received_notificationrr8rr.r/c Kyw)zh Sends a progress notification for a request that is currently being processed. Nr4)r5rr.r/r0s r6send_progress_notificationz&BaseSession.send_progress_notificationrr8reqc Kyw)zCA generic handler for incoming messages. Overwritten by subclasses.Nr4)r5rs r6rzBaseSession._handle_incoming#s r8r3)NNNrm)NN)0r9r:r;r<dictrnr rr__annotations__rrAr(r&r-rr rrrpr+rrQrrrrrqrrsrr%r)rrr'rrrbrrrrrr>r=rrr4r?r6rurusI'=oP\>\']]^^Y 0+1M NNOOi455,--26,.~ /IJ,-^<,#?3 , $((<#= ,($., ,* .. .T .$ K}%, K% K$ K  K":>$(04 J2J2.)J2'0$&6 J2 " J2 '- J2 J2^047'7&,7  7* ;y ;KR[D[ ;`d ;k/Zy*%rn%r%rN 1A/S^B^1_ dh  9M RV #"  c     t|  t     o{: ;>R RU^ ^   r?ru)=rcollections.abcr contextlibrdatetimertypesrtypingrrr r rLranyio.streams.memoryr r pydanticr typing_extensionsrmcp.shared.exceptionsrmcp.shared.messagerrrmcp.shared.response_routerr mcp.typesrrrrrrrrrrrrr r!r"r#r$r%r&r'r(r)r+r>rrnr-rArur4r?r6rs$%22 R"*UU5(~}mD m\<@ /1CEWX+]MJ);57IK]^ #I (h0w ;<h0VF    F r?