Ë
    ïñþiJ*  ã                   ó  — d 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
 ddlZddlmZmZ ddlm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 ddl m!Z!m"Z"  ejF                  e$«      Z% G d„ d«      Z&y)aÚ  
SSE Server Transport Module

This module implements a Server-Sent Events (SSE) transport layer for MCP servers.

Example usage:
```
    # Create an SSE transport at an endpoint
    sse = SseServerTransport("/messages/")

    # Create Starlette routes for SSE and message handling
    routes = [
        Route("/sse", endpoint=handle_sse, methods=["GET"]),
        Mount("/messages/", app=sse.handle_post_message),
    ]

    # Define handler functions
    async def handle_sse(request):
        async with sse.connect_sse(
            request.scope, request.receive, request._send
        ) as streams:
            await app.run(
                streams[0], streams[1], app.create_initialization_options()
            )
        # Return empty response to avoid NoneType error
        return Response()

    # Create and run Starlette app
    starlette_app = Starlette(routes=routes)
    uvicorn.run(starlette_app, host="127.0.0.1", port=port)
```

Note: The handle_sse function must return a Response to avoid a "TypeError: 'NoneType'
object is not callable" error when client disconnects. The example above returns
an empty Response() after the SSE connection ends to fix this.

See SseServerTransport class documentation for more details.
é    N)Úasynccontextmanager)ÚAny)Úquote)ÚUUIDÚuuid4)ÚMemoryObjectReceiveStreamÚMemoryObjectSendStream)ÚValidationError)ÚEventSourceResponse)ÚRequest)ÚResponse)ÚReceiveÚScopeÚSend)ÚTransportSecurityMiddlewareÚTransportSecuritySettings)ÚServerMessageMetadataÚSessionMessagec                   ó¤   ‡ — e Zd ZU dZeed<   eeee	e
z     f   ed<   eed<   ddededz  ddfˆ fd	„Zed
ededefd„«       Zd
edededdfd„Zˆ xZS )ÚSseServerTransporta  
    SSE server transport for MCP. This class provides _two_ ASGI applications,
    suitable to be used with a framework like Starlette and a server like Hypercorn:

        1. connect_sse() is an ASGI application which receives incoming GET requests,
           and sets up a new SSE stream to send server messages to the client.
        2. handle_post_message() is an ASGI application which receives incoming POST
           requests, which should contain client messages that link to a
           previously-established SSE session.
    Ú	_endpointÚ_read_stream_writersÚ	_securityNÚendpointÚsecurity_settingsÚreturnc                 ó  •— t         ‰| �  «        d|v s|j                  d«      sd|v sd|v rt        d|› d�«      ‚|j                  d«      sd|z   }|| _        i | _        t        |«      | _        t        j                  d|› �«       y	)
aÇ  
        Creates a new SSE server transport, which will direct the client to POST
        messages to the relative path given.

        Args:
            endpoint: A relative path where messages should be posted
                    (e.g., "/messages/").
            security_settings: Optional security settings for DNS rebinding protection.

        Note:
            We use relative paths instead of full URLs for several reasons:
            1. Security: Prevents cross-origin requests by ensuring clients only connect
               to the same origin they established the SSE connection with
            2. Flexibility: The server can be mounted at any path without needing to
               know its full URL
            3. Portability: The same endpoint configuration works across different
               environments (development, staging, production)

        Raises:
            ValueError: If the endpoint is a full URL instead of a relative path
        z://z//ú?ú#zGiven endpoint: z] is not a relative path (e.g., '/messages/'), expecting a relative path (e.g., '/messages/').ú/z.SseServerTransport initialized with endpoint: N)
ÚsuperÚ__init__Ú
startswithÚ
ValueErrorr   r   r   r   ÚloggerÚdebug)Úselfr   r   Ú	__class__s      €úN/root/aria/mcps/aria-brain/venv/lib/python3.12/site-packages/mcp/server/sse.pyr"   zSseServerTransport.__init__P   sŸ   ø€ ô. 	‰ÑÔð �HÑ × 3Ñ 3°DÔ 9¸SÀH¹_ÐPSÐW_ÑP_ÜØ" 8 *ð -Bð Bóð ð ×"Ñ" 3Ô'Ø˜X‘~ˆHà!ˆŒØ$&ˆÔ!Ü4Ð5FÓGˆŒÜ�‰ÐEÀhÀZÐPÕQó    ÚscopeÚreceiveÚsendc                óJ  ‡‡‡‡‡‡‡K  — |d   dk7  r t         j                  d«       t        d«      ‚t        ||«      }| j                  j                  |d¬«      ƒ d {  –—† }|r ||||«      ƒ d {  –—†  t        d«      ‚t         j                  d«       t        j                  d	«      \  Š}t        j                  d	«      \  }Št        «       Š‰| j                  ‰<   t         j                  d
‰› �«       |j                  dd«      }|j                  d«      | j                  z   }	t        |	«      › d‰j                  › �Št        j                  t         t"        t$        f      d	«      \  ŠŠˆˆˆfd„Št        j&                  «       4 ƒd {  –—† }
dt(        dt*        dt,        fˆˆˆˆˆfd„}t         j                  d«       |
j/                  ||||«       t         j                  d«       ||f­–— d d d «      ƒd {  –—†  y 7 �Œ¦7 �Œ•7 Œ|7 Œ# 1 ƒd {  –—†7  sw Y   y xY w­w)NÚtypeÚhttpz%connect_sse received non-HTTP requestz)connect_sse can only handle HTTP requestsF©Úis_postzRequest validation failedzSetting up SSE connectionr   zCreated new session with ID: Ú	root_pathÚ r    z?session_id=c            
   “   ó2  •K  — t         j                  d«       ‰4 ƒd {  –—†  ‰4 ƒd {  –—†  ‰j                  d‰dœ«      ƒ d {  –—†  t         j                  d‰› �«       ‰2 3 d {  –—† } t         j                  d| › �«       ‰j                  d| j                  j	                  dd¬«      dœ«      ƒ d {  –—†  ŒY7 Œž7 Œ•7 Œ{7 ŒZ7 Œ6 d d d «      ƒd {  –—†7   n# 1 ƒd {  –—†7  sw Y   nxY wd d d «      ƒd {  –—†7   y # 1 ƒd {  –—†7  sw Y   y xY w­w)	NzStarting SSE writerr   )ÚeventÚdatazSent endpoint event: zSending message via SSE: ÚmessageT)Úby_aliasÚexclude_none)r%   r&   r-   r8   Úmodel_dump_json)Úsession_messageÚclient_post_uri_dataÚsse_stream_writerÚwrite_stream_readers    €€€r)   Ú
sse_writerz2SseServerTransport.connect_sse.<locals>.sse_writer¥   s  øè ø€ Ü�L‰LÐ.Ô/Ø(÷ ñ Ð*=÷ ñ Ø'×,Ñ,°zÐK_Ñ-`Óa×aÐaÜ—‘Ð4Ð5IÐ4JÐKÔLà-@÷ ð ˜/Ü—L‘LÐ#<¸_Ð<MÐ!NÔOØ+×0Ñ0à%.Ø$3×$;Ñ$;×$KÑ$KÐUYÐhlÐ$KÓ$mñó÷ ñ ðøð øØaøðøðøð .A÷	÷ ÷ ÷ ó ú÷ ÷ ÷ ÷ ó üsÌ   ƒDŸB> D£DªC «D®CÁCÁCÁ%CÁ)C
Á*CÁ-ACÂ8C
Â9CÂ>DÃ DÃCÃCÃCÃCÃ	DÃCÃDÃC-	Ã!C$Ã"C-	Ã)DÃ0DÃ;C>Ã<DÄDÄDÄ	DÄDr+   r,   r-   c              “   óä   •K  —  t        ‰‰¬«      | ||«      ƒ d{  –—†  ‰j                  «       ƒ d{  –—†  ‰j                  «       ƒ d{  –—†  t        j                  d‰› �«       y7 ŒM7 Œ77 Œ!­w)zð
                The EventSourceResponse returning signals a client close / disconnect.
                In this case we close our side of the streams to signal the client that
                the connection has been closed.
                )ÚcontentÚdata_sender_callableNzClient session disconnected )r   ÚacloseÚloggingr&   )r+   r,   r-   Úread_stream_writerÚ
session_idÚsse_stream_readerr@   r?   s      €€€€€r)   Úresponse_wrapperz8SseServerTransport.connect_sse.<locals>.response_wrapper¶   sx   øè ø€ ð fÔ)Ð2CÐZdÔeØ˜7 Dó÷ ð ð )×/Ñ/Ó1×1Ð1Ø)×0Ñ0Ó2×2Ð2Ü—‘Ð <¸Z¸LÐIÕJðøð 2øØ2ús1   ƒA0œA*�A0´A,µA0ÁA.ÁA0Á,A0Á.A0zStarting SSE response taskzYielding read and write streams)r%   Úerrorr$   r   r   Úvalidate_requestr&   ÚanyioÚcreate_memory_object_streamr   r   ÚgetÚrstripr   r   ÚhexÚdictÚstrr   Úcreate_task_groupr   r   r   Ú
start_soon)r'   r+   r,   r-   ÚrequestÚerror_responseÚread_streamÚwrite_streamr3   Úfull_message_path_for_clientÚtgrI   r=   rF   rG   rH   r>   r@   r?   s               @@@@@@@r)   Úconnect_ssezSseServerTransport.connect_ssey   s  þè ø€ à�‰=˜FÒ"Ü�L‰LÐ@ÔAÜÐHÓIÐIô ˜% Ó)ˆØ#Ÿ~™~×>Ñ>¸wÐPUÐ>ÓV×VˆÙÙ  ¨°Ó6×6Ð6ÜÐ8Ó9Ð9ä�‰Ð0Ô1ô +0×*KÑ*KÈAÓ*NÑ'Ð˜KÜ,1×,MÑ,MÈaÓ,PÑ)ˆÐ)ä“Wˆ
Ø0Bˆ×!Ñ! *Ñ-Ü�‰Ð4°Z°LÐAÔBð —I‘I˜k¨2Ó.ˆ	ð (1×'7Ñ'7¸Ó'<¸t¿~¹~Ñ'MÐ$ô #(Ð(DÓ"EÐ!FÀlÐS]×SaÑSaÐRbÐcÐä/4×/PÑ/PÔQUÔVYÔ[^ÐV^ÑQ_Ñ/`ÐabÓ/cÑ,ÐÐ,ö	ô ×*Ñ*Ó,÷ 	.ð 	.°ðK¬eð K¼gð KÌT÷ Kñ Kô �L‰LÐ5Ô6Ø�M‰MÐ*¨E°7¸DÔAä�L‰LÐ:Ô;Ø Ð-Ó-÷'	.÷ 	.ð 	.ðg Wùà6ùðb	.øð 	.ø÷ 	.÷ 	.ñ 	.üsn   ‰AH#ÁHÁH#Á2HÁ3DH#ÆH
ÆH#ÆA"HÇ3H#Ç>HÇ?H#ÈH#È
H#ÈH#ÈH ÈHÈH ÈH#c              ƒ   ót  K  — t         j                  d«       t        ||«      }| j                  j	                  |d¬«      ƒ d {  –—† }|r ||||«      ƒ d {  –—† S |j
                  j                  d«      }|€4t         j                  d«       t        dd¬«      } ||||«      ƒ d {  –—† S 	 t        |¬	«      }t         j                  d
|› �«       | j                  j                  |«      }	|	s7t         j                  d|› �«       t        dd¬«      } ||||«      ƒ d {  –—† S |j                  «       ƒ d {  –—† }
t         j                  d|
› �«       	 t        j                  j                  |
«      }t         j                  d|› �«       t'        |¬«      }t)        ||¬«      }t         j                  d|› �«       t        dd¬«      } ||||«      ƒ d {  –—†  |	j%                  |«      ƒ d {  –—†  y 7 �Œµ7 �Œ¤7 �ŒV# t        $ r; t         j                  d|› �«       t        dd¬«      } ||||«      ƒ d {  –—†7  cY S w xY w7 �Œ'7 �Œ# t         $ rY}t         j#                  d«       t        dd¬«      } ||||«      ƒ d {  –—†7   |	j%                  |«      ƒ d {  –—†7   Y d }~y d }~ww xY w7 ŒÙ7 ŒÂ­w)NzHandling POST messageTr1   rG   z#Received request without session_idzsession_id is requiredi�  )Ústatus_code)rP   zParsed session ID: zReceived invalid session ID: zInvalid session IDzCould not find session for ID: zCould not find sessioni”  zReceived JSON: zValidated client message: zFailed to parse messagezCould not parse message)Úrequest_context)Úmetadataz#Sending session message to writer: ÚAcceptedéÊ   )r%   r&   r   r   rK   Úquery_paramsrN   Úwarningr   r   r$   r   ÚbodyÚtypesÚJSONRPCMessageÚmodel_validate_jsonr
   Ú	exceptionr-   r   r   )r'   r+   r,   r-   rU   rV   Úsession_id_paramÚresponserG   Úwriterrd   r8   Úerrr_   r<   s                  r)   Úhandle_post_messagez&SseServerTransport.handle_post_messageÉ   s|  è ø€ Ü�‰Ð,Ô-Ü˜% Ó)ˆð  $Ÿ~™~×>Ñ>¸wÐPTÐ>ÓU×UˆÙÙ'¨¨w¸Ó=×=Ð=à"×/Ñ/×3Ñ3°LÓAÐØÐ#Ü�N‰NÐ@ÔAÜÐ 8ÀcÔJˆHÙ! %¨°$Ó7×7Ð7ð	8ÜÐ"2Ô3ˆJÜ�L‰LÐ.¨z¨lÐ;Ô<ð ×*Ñ*×.Ñ.¨zÓ:ˆÙÜ�N‰NÐ<¸Z¸LÐIÔJÜÐ 8ÀcÔJˆHÙ! %¨°$Ó7×7Ð7à—\‘\“^×#ˆÜ�‰� t fÐ-Ô.ð	Ü×*Ñ*×>Ñ>¸tÓDˆGÜ�L‰LÐ5°g°YÐ?Ô@ô )¸ÔAˆÜ(¨¸8ÔDˆÜ�‰Ð:¸?Ð:KÐLÔMÜ˜J°CÔ8ˆÙ�u˜g tÓ,×,Ð,Ø�k‰k˜/Ó*×*Ñ*ðW Vùà=ùð 8úô
 ò 	8Ü�N‰NÐ:Ð;KÐ:LÐMÔNÜÐ 4À#ÔFˆHÙ! %¨°$Ó7×7Ð7Ò7ð	8úð 8ùà#úô ò 	Ü×ÑÐ6Ô7ÜÐ 9ÀsÔKˆHÙ˜5 '¨4Ó0×0Ñ0Ø—+‘+˜cÓ"×"Ñ"Üûð	úð 	-øØ*úsç   ‚AJ8ÁG9ÁJ8ÁG<ÁAJ8Â)G?Â*J8Â/$H ÃAJ8Ä"I	Ä#J8Ä:IÄ;J8Å7I ÆAJ8ÇJ4ÇJ8Ç3J6Ç4J8Ç<J8Ç?J8È;IÈ=I È>IÉJ8ÉIÉJ8ÉJ8É	J1É/J,ÊJ
ÊJ,Ê!J$Ê"J,Ê'J8Ê,J1Ê1J8Ê6J8)N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__rR   Ú__annotations__rQ   r   r	   r   Ú	Exceptionr   r   r"   r   r   r   r   r[   rm   Ú__classcell__)r(   s   @r)   r   r   @   s¨   ø… ñ	ð ƒNØ˜tÐ%;¸NÈYÑ<VÑ%WÐWÑXÓXØ*Ó*ñ'R ð 'RÐ9RÐUYÑ9Yð 'RÐeiõ 'RðR ðM. uð M.°wð M.Àdò M.ó ðM.ð^0+¨uð 0+¸wð 0+Èdð 0+ÐW[÷ 0+r*   r   )'rq   rE   Ú
contextlibr   Útypingr   Úurllib.parser   Úuuidr   r   rL   Úanyio.streams.memoryr   r	   Úpydanticr
   Ússe_starletter   Ústarlette.requestsr   Ústarlette.responsesr   Ústarlette.typesr   r   r   Ú	mcp.typesre   Úmcp.server.transport_securityr   r   Úmcp.shared.messager   r   Ú	getLoggerrn   r%   r   © r*   r)   ú<module>r„      s`   ðñ%óN Ý *Ý Ý ß ã ß RÝ $Ý -Ý &Ý (ß 0Ñ 0å ÷÷ Eà	ˆ×	Ñ	˜8Ó	$€÷y+ò y+r*   