Ë
    ïñþi‰7  ã                  ó  — d Z ddlmZ ddlZddlZddlmZ ddlmZ ddl	m
Z
 ddlmZ ddlZddlmZ dd	lmZ dd
lmZ ddlmZmZmZ ddlmZ ddlmZmZmZ ddlm Z  ddl!m"Z"m#Z#m$Z$  ejJ                  e&«      Z' G d„ d«      Z(y)z/StreamableHTTP Session Manager for MCP servers.é    )ÚannotationsN)ÚAsyncIterator)Ú
HTTPStatus)ÚAny)Úuuid4)Ú
TaskStatus)ÚRequest)ÚResponse)ÚReceiveÚScopeÚSend)ÚServer)ÚMCP_SESSION_ID_HEADERÚ
EventStoreÚStreamableHTTPServerTransport)ÚTransportSecuritySettings)ÚINVALID_REQUESTÚ	ErrorDataÚJSONRPCErrorc                  ó®   — e Zd ZdZ	 	 	 	 	 	 d	 	 	 	 	 	 	 	 	 	 	 	 	 d	d„Zej                  d
d„«       Z	 	 	 	 	 	 	 	 dd„Z	 	 	 	 	 	 	 	 dd„Z		 	 	 	 	 	 	 	 dd„Z
y)ÚStreamableHTTPSessionManagera†  
    Manages StreamableHTTP sessions with optional resumability via event store.

    This class abstracts away the complexity of session management, event storage,
    and request handling for StreamableHTTP transports. It handles:

    1. Session tracking for clients
    2. Resumability via an optional event store
    3. Connection management and lifecycle
    4. Request handling and transport setup
    5. Idle session cleanup via optional timeout

    Important: Only one StreamableHTTPSessionManager instance should be created
    per application. The instance cannot be reused after its run() context has
    completed. If you need to restart the manager, create a new instance.

    Args:
        app: The MCP server instance
        event_store: Optional event store for resumability support. If provided, enables resumable connections
            where clients can reconnect and receive missed events. If None, sessions are still tracked but not
            resumable.
        json_response: Whether to use JSON responses instead of SSE streams
        stateless: If True, creates a completely fresh transport for each request with no session tracking or
            state persistence between requests.
        security_settings: Optional transport security settings.
        retry_interval: Retry interval in milliseconds to suggest to clients in SSE retry field. Used for SSE
            polling behavior.
        session_idle_timeout: Optional idle timeout in seconds for stateful sessions. If set, sessions that
            receive no HTTP requests for this duration will be automatically terminated and removed. When
            retry_interval is also configured, ensure the idle timeout comfortably exceeds the retry interval to
            avoid reaping sessions during normal SSE polling gaps. Default is None (no timeout). A value of 1800
            (30 minutes) is recommended for most deployments.
    Nc                ó6  — |�|dk  rt        d«      ‚|r|�t        d«      ‚|| _        || _        || _        || _        || _        || _        || _        t        j                  «       | _        i | _        d | _        t        j                  «       | _        d| _        y )Nr   z9session_idle_timeout must be a positive number of secondsz7session_idle_timeout is not supported in stateless modeF)Ú
ValueErrorÚRuntimeErrorÚappÚevent_storeÚjson_responseÚ	statelessÚsecurity_settingsÚretry_intervalÚsession_idle_timeoutÚanyioÚLockÚ_session_creation_lockÚ_server_instancesÚ_task_groupÚ	_run_lockÚ_has_started)Úselfr   r   r   r   r   r    r!   s           úb/root/aria/mcps/aria-brain/venv/lib/python3.12/site-packages/mcp/server/streamable_http_manager.pyÚ__init__z%StreamableHTTPSessionManager.__init__A   s    € ð  Ð+Ð0DÈÒ0IÜÐXÓYÐYÙÐ-Ð9ÜÐXÓYÐYàˆŒØ&ˆÔØ*ˆÔØ"ˆŒØ!2ˆÔØ,ˆÔØ$8ˆÔ!ô ',§j¡j£lˆÔ#ØKMˆÔð  ˆÔäŸ™›ˆŒØ!ˆÕó    c               óÞ  K  — | j                   4 ƒd{  –—†  | j                  rt        d«      ‚d| _        ddd«      ƒd{  –—†  t        j                  «       4 ƒd{  –—† }|| _        t        j                  d«       	 d­–— t        j                  d«       |j                  j                  «        d| _        | j                  j                  «        ddd«      ƒd{  –—†  y7 ŒÒ7 Œ¦# 1 ƒd{  –—†7  sw Y   Œ¶xY w7 Œ # t        j                  d«       |j                  j                  «        d| _        | j                  j                  «        w xY w7 Œu# 1 ƒd{  –—†7  sw Y   yxY w­w)aw  
        Run the session manager with proper lifecycle management.

        This creates and manages the task group for all session operations.

        Important: This method can only be called once per instance. The same
        StreamableHTTPSessionManager instance cannot be reused after this
        context manager exits. Create a new instance if you need to restart.

        Use this in the lifespan context manager of your Starlette app:

        @contextlib.asynccontextmanager
        async def lifespan(app: Starlette) -> AsyncIterator[None]:
            async with session_manager.run():
                yield
        NzyStreamableHTTPSessionManager .run() can only be called once per instance. Create a new instance if you need to run again.Tz&StreamableHTTP session manager startedz,StreamableHTTP session manager shutting down)r'   r(   r   r"   Úcreate_task_groupr&   ÚloggerÚinfoÚcancel_scopeÚcancelr%   Úclear)r)   Útgs     r*   Úrunz StreamableHTTPSessionManager.runb   s7  è ø€ ð& —>‘>÷ 	%ñ 	%Ø× Ò Ü"ðYóð ð !%ˆDÔ÷	%÷ 	%ô ×*Ñ*Ó,÷ 	/ð 	/°à!ˆDÔÜ�K‰KÐ@ÔAð/Üä—‘ÐJÔKà—‘×&Ñ&Ô(Ø#'�Ô à×&Ñ&×,Ñ,Ô.÷	/÷ 	/ð 	/ð	%øð 	%ø÷ 	%÷ 	%ñ 	%úð	/ùô —‘ÐJÔKà—‘×&Ñ&Ô(Ø#'�Ô à×&Ñ&×,Ñ,Õ.úð	/ø÷ 	/÷ 	/ñ 	/üsŸ   ‚E-“C&”E-—C*¶E-ÁC(ÁE-ÁC?ÁE-Á"EÂ DÂAEÃE-Ã EÃ!E-Ã(E-Ã*C<Ã0C3Ã1C<Ã8E-ÄAEÅEÅE-ÅE*ÅE!ÅE*Å&E-c              ƒ  óÈ   K  — | j                   €t        d«      ‚| j                  r| j                  |||«      ƒ d{  –—†  y| j	                  |||«      ƒ d{  –—†  y7 Œ!7 Œ­w)a  
        Process ASGI request with proper session handling and transport setup.

        Dispatches to the appropriate handler based on stateless mode.

        Args:
            scope: ASGI scope
            receive: ASGI receive function
            send: ASGI send function
        Nz6Task group is not initialized. Make sure to use run().)r&   r   r   Ú_handle_stateless_requestÚ_handle_stateful_request)r)   ÚscopeÚreceiveÚsends       r*   Úhandle_requestz+StreamableHTTPSessionManager.handle_request‹   sc   è ø€ ð  ×ÑÐ#ÜÐWÓXÐXð �>Š>Ø×0Ñ0°¸ÀÓF×FÑFà×/Ñ/°°wÀÓE×EÑEð GøàEús!   ‚:A"¼A½A"ÁA ÁA"Á A"c              ƒ  ó„  ‡ ‡K  — t         j                  d«       t        d‰ j                  d‰ j                  ¬«      Št
        j                  dœdˆˆ fd„}‰ j                  €J ‚‰ j                  j                  |«      ƒ d{  –—†  ‰j                  |||«      ƒ d{  –—†  ‰j                  «       ƒ d{  –—†  y7 Œ87 Œ7 Œ	­w)zÝ
        Process request in stateless mode - creating a new transport for each request.

        Args:
            scope: ASGI scope
            receive: ASGI receive function
            send: ASGI send function
        z7Stateless mode: Creating new transport for this requestN)Úmcp_session_idÚis_json_response_enabledr   r   ©Útask_statusc              “  óˆ  •K  — ‰j                  «       4 ƒd {  –—† }|\  }}| j                  «        	 ‰j                  j                  ||‰j                  j	                  «       d¬«      ƒ d {  –—†  d d d «      ƒd {  –—†  y 7 Œj7 Œ# t
        $ r t        j                  d«       Y Œ5w xY w7 Œ-# 1 ƒd {  –—†7  sw Y   y xY w­w)NT©r   zStateless session crashed)ÚconnectÚstartedr   r5   Úcreate_initialization_optionsÚ	Exceptionr/   Ú	exception)rA   ÚstreamsÚread_streamÚwrite_streamÚhttp_transportr)   s       €€r*   Úrun_stateless_serverzTStreamableHTTPSessionManager._handle_stateless_request.<locals>.run_stateless_server¼   sÈ   øè ø€ Ø%×-Ñ-Ó/÷ Bð B°7Ø,3Ñ)�˜\Ø×#Ñ#Ô%ðBØŸ(™(Ÿ,™,Ø#Ø$ØŸ™×>Ñ>Ó@Ø"&ð	 'ó ÷ ð ÷	B÷ Bñ Bøðùô !ò BÜ×$Ñ$Ð%@ÖAðBúðBø÷ B÷ Bñ Büss   ƒC˜B™CœB-³:BÁ-BÁ.BÁ2CÁ=B+Á>CÂBÂB(Â%B-Â'B(Â(B-Â+CÂ-B?Â3B6Â4B?Â;C)rA   úTaskStatus[None])r/   Údebugr   r   r   r"   ÚTASK_STATUS_IGNOREDr&   Ústartr<   Ú	terminate)r)   r9   r:   r;   rM   rL   s   `    @r*   r7   z6StreamableHTTPSessionManager._handle_stateless_request¤   sÀ   ùè ø€ ô 	�‰ÐNÔOä6ØØ%)×%7Ñ%7ØØ"×4Ñ4ô	
ˆô KP×JcÑJc÷ 	Bð 	Bð ×ÑÐ+Ð+Ð+à×Ñ×$Ñ$Ð%9Ó:×:Ð:ð ×+Ñ+¨E°7¸DÓA×AÐAð ×&Ñ&Ó(×(Ñ(ð 	;øð 	Bøð 	)ús6   „A=C ÂB:ÂC ÂB<ÂC Â4B>Â5C Â<C Â>C c              ƒ  óî  ‡ ‡K  — t        ||«      }|j                  j                  t        «      }|�–|‰ j                  v rˆ‰ j                  |   }t
        j                  d«       |j                  �<‰ j                  �0t        j                  «       ‰ j                  z   |j                  _        |j                  |||«      ƒ d{  –—†  y|�€*t
        j                  d«       ‰ j                  4 ƒd{  –—†  t        «       j                  }t!        |‰ j"                  ‰ j$                  ‰ j&                  ‰ j(                  ¬«      Š‰j*                  €J ‚‰‰ j                  ‰j*                  <   t
        j-                  d|› �«       t        j.                  dœdˆˆ fd„}‰ j0                  €J ‚‰ j0                  j3                  |«      ƒ d{  –—†  ‰j                  |||«      ƒ d{  –—†  ddd«      ƒd{  –—†  yt5        dd	t7        t8        d
¬«      ¬«      }	t;        |	j=                  dd¬«      t>        j@                  d¬«      }
 |
|||«      ƒ d{  –—†  y7 �Œ�7 �Œe7 Œ“7 Œz7 Œl# 1 ƒd{  –—†7  sw Y   yxY w7 Œ&­w)zÝ
        Process request in stateful mode - maintaining session state between requests.

        Args:
            scope: ASGI scope
            receive: ASGI receive function
            send: ASGI send function
        Nz1Session already exists, handling request directlyzCreating new transport)r>   r?   r   r   r    z'Created new transport with session ID: r@   c              “  ó&  •K  — ‰j                  «       4 ƒd {  –—† }|\  }}| j                  «        	 t        j                  «       }‰j                  �-t        j
                  «       ‰j                  z   |_        |‰_        |5  ‰j                  j                  ||‰j                  j                  «       d¬«      ƒ d {  –—†  d d d «       |j                  ro‰j                  €J ‚t        j                  d‰j                  › d�«       ‰j                  j!                  ‰j                  d «       ‰j#                  «       ƒ d {  –—†  ‰j                  r_‰j                  ‰j                  v rG‰j(                  s;t        j                  d‰j                  › d�«       ‰j                  ‰j                  = 	 d d d «      ƒd {  –—†  y 7 �Œ©7 �Œ# 1 sw Y   �ŒxY w7 Œ“# t$        $ r& t        j'                  d‰j                  › d�«       Y Œ¿w xY w# ‰j                  ra‰j                  ‰j                  v rH‰j(                  s;t        j                  d‰j                  › d�«       ‰j                  ‰j                  = w w w w xY w7 Œ¾# 1 ƒd {  –—†7  sw Y   y xY w­w)NFrC   zSession z idle timeoutz crashedzCleaning up crashed session z from active instances.)rD   rE   r"   ÚCancelScoper!   Úcurrent_timeÚdeadlineÚ
idle_scoper   r5   rF   Úcancelled_caughtr>   r/   r0   r%   ÚpoprR   rG   rH   Úis_terminated)rA   rI   rJ   rK   rX   rL   r)   s        €€r*   Ú
run_serverzIStreamableHTTPSessionManager._handle_stateful_request.<locals>.run_server  sy  øè ø€ Ø-×5Ñ5Ó7÷ 'Zð 'Z¸7Ø4;Ñ1˜ \Ø#×+Ñ+Ô-ð$Zô
 */×):Ñ):Ó)<˜JØ#×8Ñ8ÐDÜ6;×6HÑ6HÓ6JÈT×MfÑMfÑ6f 
Ô 3Ø<F Ô 9à!+ñ "Ø&*§h¡h§l¡lØ$/Ø$0Ø$(§H¡H×$JÑ$JÓ$LØ.3ð	 '3ó '"÷ !"ð !"÷"ð  *×:Ò:Ø'5×'DÑ'DÐ'PÐ PÐ'PÜ &§¡¨h°~×7TÑ7TÐ6UÐUbÐ,cÔ dØ $× 6Ñ 6× :Ñ :¸>×;XÑ;XÐZ^Ô _Ø&4×&>Ñ&>Ó&@× @Ð @ð
 !/× =Ò =Ø$2×$AÑ$AÀT×E[ÑE[Ñ$[Ø(6×(DÒ(Dä &§¡Ø$BØ'5×'DÑ'DÐ&Eð F8ð%8ô!"ð
 %)×$:Ñ$:¸>×;XÑ;XÑ$Y÷O'Z÷ 'Zñ 'Zùð!"ù÷"ñ "úð !AùÜ(ò aÜ"×,Ñ,¨x¸×8UÑ8UÐ7VÐV^Ð-_Ö`ðaûð !/× =Ò =Ø$2×$AÑ$AÀT×E[ÑE[Ñ$[Ø(6×(DÒ(Dä &§¡Ø$BØ'5×'DÑ'DÐ&Eð F8ð%8ô!"ð
 %)×$:Ñ$:¸>×;XÑ;XÑ$Yð )Eð %\ð !>úð='Zø÷ 'Z÷ 'Zñ 'Züs­   ƒJ˜G™JœI<³AGÂ;GÂ=GÂ>GÃA>GÅ GÅGÅA+I<Æ0JÆ;I:Æ<JÇGÇG	Ç	GÇ,HÈHÈHÈHÈA/I7É7I<É:JÉ<JÊJÊJÊ
Jz2.0zserver-errorzSession not found)ÚcodeÚmessage)ÚjsonrpcÚidÚerrorT)Úby_aliasÚexclude_nonezapplication/json)ÚcontentÚstatus_codeÚ
media_type)rA   rN   ÚreturnÚNone)!r	   ÚheadersÚgetr   r%   r/   rO   rX   r!   r"   rV   rW   r<   r$   r   Úhexr   r   r   r   r    r>   r0   rP   r&   rQ   r   r   r   r
   Úmodel_dump_jsonr   Ú	NOT_FOUND)r)   r9   r:   r;   ÚrequestÚrequest_mcp_session_idÚ	transportÚnew_session_idr\   Úerror_responseÚresponserL   s   `          @r*   r8   z5StreamableHTTPSessionManager._handle_stateful_requestÕ   sj  ùè ø€ ô ˜% Ó)ˆØ!(§¡×!4Ñ!4Ô5JÓ!KÐð "Ð-Ð2HÈD×LbÑLbÑ2bØ×.Ñ.Ð/EÑFˆIÜ�L‰LÐLÔMà×#Ñ#Ð/°D×4MÑ4MÐ4YÜ05×0BÑ0BÓ0DÀt×G`ÑG`Ñ0`�	×$Ñ$Ô-Ø×*Ñ*¨5°'¸4Ó@×@Ð@Øà!Ñ)ä�L‰LÐ1Ô2Ø×2Ñ2÷ ?Jñ ?JÜ!&£§¡�Ü!>Ø#1Ø-1×-?Ñ-?Ø $× 0Ñ 0Ø&*×&<Ñ&<Ø#'×#6Ñ#6ô"�ð &×4Ñ4Ð@Ð@Ð@ØHV�×&Ñ& ~×'DÑ'DÑEÜ—‘ÐEÀnÐEUÐVÔWô IN×HaÑHa÷ (Zð (ZðV ×'Ñ'Ð3Ð3Ð3à×&Ñ&×,Ñ,¨ZÓ8×8Ð8ð %×3Ñ3°E¸7ÀDÓI×IÐI÷?J÷ ?Jð ?JôH *ØØ!ÜÜ(Ø/ôôˆNô  Ø&×6Ñ6ÀÐSWÐ6ÓXÜ&×0Ñ0Ø-ôˆHñ
 ˜5 '¨4Ó0×0Ñ0ðo Aùð?Jùðx 9øð Jøð?Jø÷ ?J÷ ?Jñ ?Júðb 1úsŒ   „B>I5ÃIÃ-I5Ã0IÃ1I5Ã4CIÇIÇIÇIÇ IÇ$I5Ç/IÇ0AI5ÉI3ÉI5ÉI5ÉIÉIÉI5ÉI0É$I'É%I0É,I5)NFFNNN)r   zMCPServer[Any, Any]r   zEventStore | Noner   Úboolr   rt   r   z TransportSecuritySettings | Noner    z
int | Noner!   zfloat | None)rg   zAsyncIterator[None])r9   r   r:   r   r;   r   rg   rh   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r+   Ú
contextlibÚasynccontextmanagerr5   r<   r7   r8   © r,   r*   r   r      s  „ ñ ðJ *.Ø#ØØ>BØ%)Ø-1ð"à ð"ð 'ð"ð ð	"ð
 ð"ð <ð"ð #ð"ð +ó"ðB ×#Ñ#ò&/ó $ð&/ðPFàðFð ðFð ð	Fð
 
óFð2/)àð/)ð ð/)ð ð	/)ð
 
ó/)ðbo1àðo1ð ðo1ð ð	o1ð
 
ôo1r,   r   ))rx   Ú
__future__r   ry   ÚloggingÚcollections.abcr   Úhttpr   Útypingr   Úuuidr   r"   Ú	anyio.abcr   Ústarlette.requestsr	   Ústarlette.responsesr
   Ústarlette.typesr   r   r   Úmcp.server.lowlevel.serverr   Ú	MCPServerÚmcp.server.streamable_httpr   r   r   Úmcp.server.transport_securityr   Ú	mcp.typesr   r   r   Ú	getLoggerru   r/   r   r{   r,   r*   ú<module>rŒ      sf   ðÙ 5å "ã Û Ý )Ý Ý Ý ã Ý  Ý &Ý (ß 0Ñ 0å :÷ñ õ
 Dß >Ñ >à	ˆ×	Ñ	˜8Ó	$€÷f1ò f1r,   