
    Rmj<                       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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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j'        e(          Z) G d d          Z*dS )z/StreamableHTTP Session Manager for MCP servers.    )annotationsN)AsyncIterator)Any)uuid4)
TaskStatus)Request)Response)ReceiveScopeSend)AuthenticatedUserAuthorizationContextauthorization_context)Server)MCP_SESSION_ID_HEADER
EventStoreStreamableHTTPServerTransport)TransportSecuritySettings)INVALID_REQUEST	ErrorDataJSONRPCErrorc                  b    e Zd ZdZ	 	 	 	 	 	 dd dZej        d!d            Zd"dZd"dZ	d"dZ
dS )#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.
    NFappMCPServer[Any, Any]event_storeEventStore | Nonejson_responsebool	statelesssecurity_settings TransportSecuritySettings | Noneretry_interval
int | Nonesession_idle_timeoutfloat | Nonec                T   ||dk    rt          d          |r|t          d          || _        || _        || _        || _        || _        || _        || _        t          j
                    | _        i | _        i | _        d | _        t          j
                    | _        d| _        d S )Nr   z9session_idle_timeout must be a positive number of secondsz7session_idle_timeout is not supported in stateless modeF)
ValueErrorRuntimeErrorr   r   r   r    r!   r#   r%   anyioLock_session_creation_lock_server_instances_session_owners_task_group	_run_lock_has_started)selfr   r   r   r    r!   r#   r%   s           j/home/thesage/.hermes/hermes-agent/venv/lib/python3.11/site-packages/mcp/server/streamable_http_manager.py__init__z%StreamableHTTPSessionManager.__init__A   s      +0D0I0IXYYY 	Z-9XYYY&*"!2,$8! ',jll#KM AC  !    returnAsyncIterator[None]c               .  K   | j         4 d{V  | j        rt          d          d| _        ddd          d{V  n# 1 d{V swxY w Y   t          j                    4 d{V }|| _        t                              d           	 dW V  t                              d           |j        	                                 d| _        | j
                                         | j                                         nq# t                              d           |j        	                                 d| _        | j
                                         | j                                         w xY w	 ddd          d{V  dS # 1 d{V swxY w Y   dS )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)r0   r1   r)   r*   create_task_groupr/   loggerinfocancel_scopecancelr-   clearr.   )r2   tgs     r3   runz StreamableHTTPSessionManager.rune   s     & > 	% 	% 	% 	% 	% 	% 	% 	%  "Y   !%D	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% *,, 	- 	- 	- 	- 	- 	- 	-!DKK@AAA	-JKKK&&(((#' &,,...$**,,,, JKKK&&(((#' &,,...$**,,,,,	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	- 	-s=   A  
A
A
*"FC?A-F?A.E--F
FFscoper   receiver
   sendr   Nonec                   K   | j         t          d          | j        r|                     |||           d{V  dS |                     |||           d{V  dS )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)r2   rA   rB   rC   s       r3   handle_requestz+StreamableHTTPSessionManager.handle_request   s        #WXXX > 	F00FFFFFFFFFFF//wEEEEEEEEEEEr5   c                d   K   t                               d           t          d j        d j                  t
          j        dd fd} j        J  j                            |           d{V  	                    |||           d{V  
                                 d{V  dS )	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_statusrM   TaskStatus[None]c                  K                                    4 d {V }|\  }}|                                  	 j                            ||j                                        d           d {V  n*# t
          $ r t                              d           Y nw xY wd d d           d {V  d S # 1 d {V swxY w Y   d S )NTr    zStateless session crashed)connectstartedr   r@   create_initialization_options	Exceptionr:   	exception)rM   streamsread_streamwrite_streamhttp_transportr2   s       r3   run_stateless_serverzTStreamableHTTPSessionManager._handle_stateless_request.<locals>.run_stateless_server   s     %--// B B B B B B B7,3)\##%%%B(,,#$>>@@"&	 '           ! B B B$$%@AAAAABB B B B B B B B B B B B B B B B B B B B B B B B B B B B B Bs4   B2;A54B25$BB2BB22
B<?B<)rM   rN   )r:   debugr   r   r!   r*   TASK_STATUS_IGNOREDr/   startrH   	terminate)r2   rA   rB   rC   rZ   rY   s   `    @r3   rF   z6StreamableHTTPSessionManager._handle_stateless_request   s      	NOOO6%)%7"4	
 
 
 KPJc 	B 	B 	B 	B 	B 	B 	B 	B 	B +++$$%9::::::::: ++E7DAAAAAAAAA &&(((((((((((r5   c                L   K   t          ||          }|j                            t                    }|                    d          }t	          |t
                    rt          |          nd}|&| j        v r j        |         }| j                            |          k    rt          
                    d|dd                    t          ddt          t          d          	          }	t          |	                    d
d
          dd          }
 |
|||           d{V  dS t                              d           |j        , j        %t'          j                     j        z   |j        _        |                    |||           d{V  dS |)t                              d            j        4 d{V  t1                      j        }t5          | j         j         j         j                  j        J || j        j        <    j        j        <   t                               d|            t&          j!        dd fd} j"        J  j"        #                    |           d{V                      |||           d{V  ddd          d{V  dS # 1 d{V swxY w Y   dS t          ddt          t          d          	          }	t          |	                    d
d
          dd          }
 |
|||           d{V  dS )z
        Process request in stateful mode - maintaining session state between requests.

        Args:
            scope: ASGI scope
            receive: ASGI receive function
            send: ASGI send function
        userNz\Rejecting request for session %s: credential does not match the one that created the session@   z2.0zserver-errorzSession not found)codemessage)jsonrpciderrorT)by_aliasexclude_nonei  zapplication/json)status_code
media_typez1Session already exists, handling request directlyzCreating new transport)rJ   rK   r   r!   r#   z'Created new transport with session ID: rL   rM   rN   r6   rD   c                .  K                                    4 d {V }|\  }}|                                  	 t          j                    }j        't          j                    j        z   |_        |_        |5  j        	                    ||j        
                                d           d {V  d d d            n# 1 swxY w Y   |j        rj        J t                              dj         d           j                            j        d            j                            j        d                                             d {V  n3# t&          $ r& t                              dj         d           Y nw xY wj        rej        j        v rWj        sPt                              dj         d           j        j        = j                            j        d            nt# j        rfj        j        v rYj        sSt                              dj         d           j        j        = j                            j        d            w w w w xY wd d d           d {V  d S # 1 d {V swxY w Y   d S )NFrP   zSession z idle timeoutz crashedzCleaning up crashed session z from active instances.)rQ   rR   r*   CancelScoper%   current_timedeadline
idle_scoper   r@   rS   cancelled_caughtrJ   r:   r;   r-   popr.   r^   rT   rU   is_terminated)rM   rV   rW   rX   ro   rY   r2   s        r3   
run_serverzIStreamableHTTPSessionManager._handle_stateful_request.<locals>.run_server  s     -5577 )^ )^ )^ )^ )^ )^ )^74;1\#++---&^
 */):)<)<J#8D6;6H6J6JTMf6f
 3<F 9!+ " "&*hll$/$0$(H$J$J$L$L.3	 '3 '" '" !" !" !" !" !" !" !"" " " " " " " " " " " " " " "  *: A'5'D'P'P'P &,c~7T,c,c,c d d d $ 6 : :>;XZ^ _ _ _ $ 4 8 89VX\ ] ] ]&4&>&>&@&@ @ @ @ @ @ @ @( a a a",,-_8U-_-_-_`````a !/ =^$2$ATE[$[$[(6(D %\ !'%8'5'D%8 %8 %8!" !" !"
 %)$:>;X$Y $ 4 8 89VX\ ] ] ] !/ =^$2$ATE[$[$[(6(D %\ !'%8'5'D%8 %8 %8!" !" !"
 %)$:>;X$Y $ 4 8 89VX\ ] ] ] ]^$[$[A)^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^ )^sn   JAE<<C8EC	EC	BEG=-F
G=FG=A-J=A1I..J
JJ)rM   rN   r6   rD   )$r   headersgetr   
isinstancer   r   r-   r.   r:   warningr   r   r   r	   model_dump_jsonr[   ro   r%   r*   rm   rn   rH   r,   r   hexr   r   r   r!   r#   rJ   r;   r\   r/   r]   )r2   rA   rB   rC   requestrequest_mcp_session_idr`   	requestor	transportbodyresponsenew_session_idrs   rY   s   `            @r3   rG   z5StreamableHTTPSessionManager._handle_stateful_request   s      %))!(!4!45J!K!Kyy  3=dDU3V3V`)$///\`	 "-2HDLb2b2b./EFID0445KLLLL r*3B3/   $!nI?dw<x<x<x   $(($T(JJ #1  
 hugt444444444LLLMMM#/D4M4Y050B0D0DtG`0`	$-**5'4@@@@@@@@@F!)LL12222 CJ CJ CJ CJ CJ CJ CJ CJ!&!>#1-1-? $ 0&*&<#'#6" " " &4@@@(JSD()FGHV&~'DEVnVVWWW INHa *^ *^ *^ *^ *^ *^ *^ *^ *^Z '333&,,Z888888888 %33E7DIIIIIIIIIGCJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJ CJL  .	`s8t8t8t  D  $$d$FFTWdv  H (5'400000000000s   CJ11
J;>J;)NFFNNN)r   r   r   r   r   r   r    r   r!   r"   r#   r$   r%   r&   )r6   r7   )rA   r   rB   r
   rC   r   r6   rD   )__name__
__module____qualname____doc__r4   
contextlibasynccontextmanagerr@   rH   rF   rG    r5   r3   r   r      s           J *.#>B%)-1"" "" "" "" ""H #'- '- '- $#'-RF F F F2/) /) /) /)b~1 ~1 ~1 ~1 ~1 ~1r5   r   )+r   
__future__r   r   loggingcollections.abcr   typingr   uuidr   r*   	anyio.abcr   starlette.requestsr   starlette.responsesr	   starlette.typesr
   r   r   &mcp.server.auth.middleware.bearer_authr   r   r   mcp.server.lowlevel.serverr   	MCPServermcp.server.streamable_httpr   r   r   mcp.server.transport_securityr   	mcp.typesr   r   r   	getLoggerr   r:   r   r   r5   r3   <module>r      s   5 5 " " " " " "      ) ) ) ) ) )                          & & & & & & ( ( ( ( ( ( 0 0 0 0 0 0 0 0 0 0 q q q q q q q q q q : : : : : :         
 D C C C C C > > > > > > > > > >		8	$	$y1 y1 y1 y1 y1 y1 y1 y1 y1 y1r5   