
    )`i/                       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j%        e&          Z' G d d          Z(dS )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
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

    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.
    NFappMCPServer[Any, Any]event_storeEventStore | Nonejson_responsebool	statelesssecurity_settings TransportSecuritySettings | Noneretry_interval
int | Nonec                    || _         || _        || _        || _        || _        || _        t          j                    | _        i | _	        d | _
        t          j                    | _        d| _        d S )NF)r   r   r   r   r   r!   anyioLock_session_creation_lock_server_instances_task_group	_run_lock_has_started)selfr   r   r   r   r   r!   s          v/home/jaya/work/projects/VOICE-AGENT/VIET/agent-env/lib/python3.11/site-packages/mcp/server/streamable_http_manager.py__init__z%StreamableHTTPSessionManager.__init__<   so     &*"!2, ',jll#KM  !    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
                                         nX# t                              d           |j        	                                 d| _        | 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)r)   r*   RuntimeErrorr$   create_task_groupr(   loggerinfocancel_scopecancelr'   clear)r+   tgs     r,   runz StreamableHTTPSessionManager.runV   sz     & > 	% 	% 	% 	% 	% 	% 	% 	%  "Y   !%D	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% 	% *,, 	/ 	/ 	/ 	/ 	/ 	/ 	/!DKK@AAA/JKKK&&(((#' &,,.... JKKK&&(((#' &,,.....	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/ 	/s=   A  
A
A
*"EC&AE&AD;;E
EE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(   r2   r   _handle_stateless_request_handle_stateful_request)r+   r;   r<   r=   s       r,   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Er.   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_statusrG   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	Exceptionr4   	exception)rG   streamsread_streamwrite_streamhttp_transportr+   s       r,   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<)rG   rH   )r4   debugr   r   r   r$   TASK_STATUS_IGNOREDr(   startrB   	terminate)r+   r;   r<   r=   rT   rS   s   `    @r,   r@   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 &&(((((((((((r.   c                   K   t          ||          }|j                            t                    }|O| j        v rF j        |         }t
                              d           |                    |||           d{V  dS |t
                              d            j        4 d{V  t                      j
        }t          | 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/          t0          d                    }	t3          |	                    dd          t6          j        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
        Nz1Session already exists, handling request directlyzCreating new transport)rD   rE   r   r   r!   z'Created new transport with session ID: rF   rG   rH   r/   r>   c                  K                                    4 d {V }|\  }}|                                  	 j                            ||j                                        d           d {V  n># t
          $ r1}t                              dj         d| d           Y d }~nd }~ww xY wj        rEj        j	        v r7j
        s0t                              dj         d           j	        j        = nQ# j        rEj        j	        v r7j
        s0t                              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 )	NFrJ   zSession z
 crashed: T)exc_infozCleaning up crashed session z from active instances.)rK   rL   r   r:   rM   rN   r4   errorrD   r'   is_terminatedr5   )rG   rP   rQ   rR   erS   r+   s        r,   
run_serverzIStreamableHTTPSessionManager._handle_stateful_request.<locals>.run_server   s     -5577 Z Z Z Z Z Z Z74;1\#++---Z"&(,, + , $ F F H H*/	 #/ # #          )   "LL W>+H W WTU W W)- )         !/ =
Z$2$ATE[$[$[(6(D %\ !'%8'5'D%8 %8 %8!" !" !"
 %)$:>;X$Y !/ =
Z$2$ATE[$[$[(6(D %\ !'%8'5'D%8 %8 %8!" !" !"
 %)$:>;X$Y Y Y Y Y Y7Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z Z ZsN   E%;A54D 5
B0?'B+&D +B00D 3AE% AEE%%
E/2E/z2.0zserver-errorzSession not found)codemessage)jsonrpcidr\   T)by_aliasexclude_nonezapplication/json)contentstatus_code
media_type)rG   rH   r/   r>   )r	   headersgetr   r'   r4   rU   rB   r&   r   hexr   r   r   r   r!   rD   r5   r$   rV   r(   rW   r   r   r   r
   model_dump_jsonr   	NOT_FOUND)r+   r;   r<   r=   requestrequest_mcp_session_id	transportnew_session_idr_   error_responseresponserS   s   `          @r,   rA   z5StreamableHTTPSessionManager._handle_stateful_request   sR      %))!(!4!45J!K!K "-2HDLb2b2b./EFILLLMMM**5'4@@@@@@@@@F!)LL12222 3J 3J 3J 3J 3J 3J 3J 3J!&!>#1-1-? $ 0&*&<#'#6" " " &4@@@HV&~'DEVnVVWWW INHa Z Z Z Z Z Z Z Z 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3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3J 3Jp *!(/    N  &66SW6XX&0-  H
 (5'400000000000s   /CF
FF)NFFNN)r   r   r   r   r   r   r   r   r   r    r!   r"   )r/   r0   )r;   r   r<   r   r=   r   r/   r>   )__name__
__module____qualname____doc__r-   
contextlibasynccontextmanagerr:   rB   r@   rA    r.   r,   r   r      s         @ *.#>B%)" " " " "4 #&/ &/ &/ $#&/PF F F F2/) /) /) /)b`1 `1 `1 `1 `1 `1r.   r   ))rw   
__future__r   rx   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   	getLoggerrt   r4   r   rz   r.   r,   <module>r      s   5 5 " " " " " "      ) ) ) ) ) )                                & & & & & & ( ( ( ( ( ( 0 0 0 0 0 0 0 0 0 0 : : : : : :         
 D C C C C C > > > > > > > > > >		8	$	$K1 K1 K1 K1 K1 K1 K1 K1 K1 K1r.   