
    i6                     T   d Z ddlZddlZddlZddlZddlmZmZmZm	Z	 ddl
mZ ddlmZmZmZ ddlZddlmZmZmZ ddlmZ dd	lmZmZmZmZmZmZmZmZm Z  d
dl!m"Z"m#Z# d
dl$m%Z% erddl&m'Z( ddlm)Z) d
dl*m+Z+  ejX                  e-      Z.de/e0ef   de/e0ef   fdZ1 G d d      Z2y)z8Query class for handling bidirectional control protocol.    N)AsyncIterableAsyncIterator	AwaitableCallable)suppress)TYPE_CHECKINGAnyLiteral)CallToolRequestCallToolRequestParamsListToolsRequest   )ProcessError)	PermissionModePermissionResultAllowPermissionResultDenyPermissionUpdateSDKControlPermissionRequestSDKControlRequestSDKControlResponseSDKHookCallbackRequestToolPermissionContext   )
TaskHandlespawn_detached)	Transport)Server)
SessionKey)TranscriptMirrorBatcherhook_outputreturnc                 p    i }| j                         D ]   \  }}|dk(  r||d<   |dk(  r||d<   |||<   " |S )zConvert Python-safe field names to CLI-expected field names.

    The Python SDK uses `async_` and `continue_` to avoid keyword conflicts,
    but the CLI expects `async` and `continue`. This function performs the
    necessary conversion.
    async_async	continue_continue)items)r    	convertedkeyvalues       g/volume1/homes/robertsu/coba/app/.venv/lib/python3.12/site-packages/claude_agent_sdk/_internal/query.py_convert_hook_output_for_clir,   *   sT     I!'')
U(?!&IgK$)Ij!"IcN *     c                      e Zd ZdZ	 	 	 	 	 	 	 d8dededeeeee	f   e
geeez     f   dz  deeeeee	f      f   dz  deedf   dz  d	ed
eeeee	f   f   dz  dedz  dee   ed   z  dz  fdZd9dZdddeddfdZdeee	f   dz  fdZd:dZde	defdZdeddfdZd:dZdeddfdZ	 d;deee	f   dedeee	f   fdZded eee	f   deee	f   fd!Zdeee	f   fd"Zdeee	f   fd#Zd:d$Z d%e!ddfd&Z"d'edz  ddfd(Z#d)eddfd*Z$deddfd+Z%ded,eddfd-Z&d.eddfd/Z'd:d0Z(d1e)eee	f      ddfd2Z*de+eee	f      fd3Z,d:d4Z-d:d5Z.de+eee	f      fd6Z/deee	f   fd7Z0y)<QueryzHandles bidirectional control protocol on top of Transport.

    This class manages:
    - Control request/response routing
    - Hook callbacks
    - Tool permission callbacks
    - Message streaming
    - Initialization handshake
    N	transportis_streaming_modecan_use_toolhookssdk_mcp_servers	McpServerinitialize_timeoutagentsexclude_dynamic_sectionsskillsallc
                    || _         || _        || _        || _        |xs i | _        |xs i | _        || _        || _        |	| _        i | _	        i | _
        i | _        d| _        d| _        t        j                  t         t"        t$        f      d      \  | _        | _        d| _        t-               | _        i | _        d| _        d| _        d| _        t        j8                         | _        d| _        d| _        y)a0  Initialize Query with transport and callbacks.

        Args:
            transport: Low-level transport for I/O
            is_streaming_mode: Whether using streaming (bidirectional) mode
            can_use_tool: Optional callback for tool permission requests
            hooks: Optional hook configurations
            sdk_mcp_servers: Optional SDK MCP server instances
            initialize_timeout: Timeout in seconds for the initialize request
            agents: Optional agent definitions to send via initialize
            exclude_dynamic_sections: Optional preset-prompt flag to send via
                initialize (see ``SystemPromptPreset``)
            skills: Optional skill allowlist to send via initialize so the CLI
                can filter which skills are loaded into the system prompt
        r   d   )max_buffer_sizeNF) _initialize_timeoutr0   r1   r2   r3   r4   _agents_exclude_dynamic_sections_skillspending_control_responsespending_control_resultshook_callbacksnext_callback_id_request_counteranyiocreate_memory_object_streamdictstrr	   _message_send_message_receive
_read_taskset_child_tasks_inflight_requests_initialized_closed_initialization_resultEvent_first_result_event_last_error_result_text_transcript_mirror_batcher)
selfr0   r1   r2   r3   r4   r6   r7   r8   r9   s
             r+   __init__zQuery.__init__H   s   > $6 "!2([b
.4")A& BD&NP$=? ! ! 5:4U4UcN5
51D1 .2-0U9;!=A# $);;= 
 48$ KO'r-   r!   c                     || _         y)a  Attach a batcher that receives ``transcript_mirror`` frames.

        When set, the read loop peels ``transcript_mirror`` frames off stdout
        (they are not yielded to consumers), enqueues them on the batcher, and
        flushes before yielding each ``result`` message.
        N)rW   )rX   batchers     r+   set_transcript_mirror_batcherz#Query.set_transcript_mirror_batcher   s     +2'r-   r)   zSessionKey | Noneerrorc           	         dd||t        t        j                               |r|j                  dd      ndd}	 | j                  j                  |       y# t        $ r }t        j                  d|       Y d}~yd}~ww xY w)uo  Surface a :meth:`SessionStore.append` failure as a system message.

        Called from the batcher's ``on_error``; the dropped batch is not
        retried (at-most-once delivery), so this is the consumer's only signal.
        Non-blocking — if the message buffer is full the error is logged and
        dropped rather than back-pressuring the read loop.
        systemmirror_error
session_id )typesubtyper]   r)   uuidra   z/Dropping mirror_error message (buffer full): %sN)	rJ   re   uuid4getrK   send_nowait	Exceptionloggerwarning)rX   r)   r]   msges        r+   report_mirror_errorzQuery.report_mirror_error   su     %

%7:#'',3
	Q**3/ 	QNNLaPP	Qs   A 	A?A::A?c                 4  K   | j                   syi }| j                  r| j                  j                         D ]  \  }}|s	g ||<   |D ]  }g }|j                  dg       D ]F  }d| j                   }| xj                  dz  c_        || j
                  |<   |j                  |       H |j                  d      |d}|j                  d      |j                  d      |d<   ||   j                  |         d|r|ndd	}	| j                  r| j                  |	d
<   | j                  | j                  |	d<   t        | j                  t              r| j                  |	d<   | j                  |	| j                         d{   }
d| _        |
| _        |
S 7 w)zInitialize control protocol if in streaming mode.

        Returns:
            Initialize response with supported commands, or None if not streaming
        Nr3   hook_r   matcher)rq   hookCallbackIdstimeout
initialize)rd   r3   r7   excludeDynamicSectionsr9   )rs   T)r1   r3   r'   rg   rE   rD   appendr?   r@   
isinstancerA   list_send_control_requestr>   rQ   rS   )rX   hooks_configeventmatchersrq   callback_idscallbackcallback_idhook_matcher_configrequestresponses              r+   rt   zQuery.initialize   s     %% (*::#'::#3#3#5x*,L'#+')(/GR(@H,1$2G2G1H*IK 11Q61?GD//<(//<	 )A (/{{9'=/;?+ #;;y1==D[[=S/	:$U+223FG $, $6( $%1\t#
 << $GH))5040N0NG,- dllD) $GH 33T55 4 
 
 !&.#
s   >FE FFFc                 b   K   | j                   t        | j                               | _         yyw)z&Start reading messages from transport.N)rM   r   _read_messagesrX   s    r+   startzQuery.start   s*     ??",T-@-@-BCDO #s   -/coroc                     t        |      }| j                  j                  |       |j                  | j                  j                         |S )z5Spawn a child task that will be cancelled on close().)r   rO   addadd_done_callbackdiscard)rX   r   tasks      r+   
spawn_taskzQuery.spawn_task   s?    d#d#t00889r-   r   c                      |d    j                   j                  |            }| j                  <   dt        ddf fd}|j	                  |       y)z>Spawn a control request handler and track it for cancellation.
request_id_tr!   Nc                 >    j                   j                  d        y )N)rP   pop)r   req_idrX   s    r+   _donez3Query._spawn_control_request_handler.<locals>._done   s    ##''5r-   )r   _handle_control_requestrP   r   r   )rX   r   r   r   r   s   `   @r+   _spawn_control_request_handlerz$Query._spawn_control_request_handler   sY    &t;;GDE*.'	6j 	6T 	6 	u%r-   c                 
  K   	 | j                   j                         2 3 d{   }| j                  r nI|j                  d      }|dk(  r|j                  di       }|j                  d      }|| j                  v rk| j                  |   }|j                  d      dk(  r)t        |j                  dd            | j                  |<   n|| j                  |<   |j                          |d	k(  r |}| j                  s| j                  |       |d
k(  rC|j                  d      }|r.| j                  j                  |d      }|r|j                          8|dk(  r0| j                  "| j                  j                  |d   |d          m|dk(  r| j                  "| j                  j                          d{    | j                  j                          |j                  d      rI|j                  d      xs g }	dj!                  |	      xs t#        |j                  dd            | _        n(d| _        n |dk(  r|j                  d      dk(  sd| _        | j&                  j)                  |       d{    a| j                  At+        j>                  d      5  | j                  j                          d{    ddd       | j                  j                          tA        t*        jB                        5  | j&                  jE                  ddi       ddd       | j&                  jG                          y7 7 7 6 # t+        j,                         $ r t.        j1                  d        t
        $ r}
t3        | j                  j5                               D ]3  \  }}|| j                  vs|
| j                  |<   |j                          5 t7        |
t8              r<| j$                  0d| j$                   }t.        j1                  d|
j:                         n#t#        |
      }t.        j=                  d|
        | j&                  j)                  d|d       d{  7   Y d}
~
d}
~
ww xY w7 # 1 sw Y   xY w# 1 sw Y   qxY w# | j                  Ot+        j>                  d      5  | j                  j                          d{  7   ddd       n# 1 sw Y   nxY w| j                  j                          tA        t*        jB                        5  | j&                  jE                  ddi       ddd       n# 1 sw Y   nxY w| j&                  jG                          w xY ww)z,Read messages from transport and route them.Nrc   control_responser   r   rd   r]   Unknown errorcontrol_requestcontrol_cancel_requesttranscript_mirrorfilePathentriesresultis_errorerrorsz; zunknown errorr_   session_state_changedzRead task cancelledz&Claude Code returned an error result: z<Replacing ProcessError (exit code %s) with result error textzFatal error in message reader: )rc   r]   T)shieldend)$r0   read_messagesrR   rg   rB   ri   rC   rN   r   rP   r   cancelrW   enqueueflushrU   joinrJ   rV   rK   sendrG   get_cancelled_exc_classrj   debugrx   r'   rw   r   	exit_coder]   CancelScoper   
WouldBlockrh   close)rX   messagemsg_typer   r   r{   r   	cancel_idinflightr   rm   
error_texts               r+   r   zQuery._read_messages   s    |	'!%!=!=!? H7g<<";;v. 11&{{:r:H!)l!;J!T%C%CC $ > >z J#<<	2g=GP (Wo FHD88D HPD88D		!22 29G<<;;GD!99 'L 9I #'#:#:#>#>y$#O#$OO-!44 66B77??#J/1C  x' 66B"==CCEEE,,002{{:.!(X!6!<"7;yy7H 8C#KK	?CM4 8<4(I.2II 48D0 ((--g666J ..:&&d399??AAA 4 $$((* %**+""..? ,$$&wH7h F( 7Q "@T ,,. 	LL./ 	R%)$*H*H*N*N*P%Q!
ET%A%AA?@D00<IIK &R !\*t/K/K/W<3346  RKK
 !V
>qcBC$$))7Z*PQQQ3	RB B 43 ,+ ..:&&d399??AAA 433 $$((* %**+""..? ,++$$&s  U9M
 MM MFM
 9M:B>M
 8M9M
 ?"U9!Q??Q< Q?;U9?R#U9 MM
 M
 M
 
5Q9?7Q47B1Q4(Q+)Q4.R 4Q99R <Q??R	U9RU9#U6<S)S
S) 	U6)S2.:U6(U	U6U"U66U9c                 &  K   |d   }|d   }|d   }	 i }|dk(  r|}|d   }| j                   st        d      t        d|j                  d      xs g D cg c]  }t	        j
                  |       c}|j                  d	      |j                  d
      |j                  d      |j                  d      |j                  d      |j                  d      |j                  d      	      }	| j                  |d   |d   |	       d{   }
t        |
t              rWd|
j                  |
j                  n|d}|
j                  }|
j                  D cg c]  }|j                          c}|d<   nPt        |
t              r-d|
j                  d}|
j                  r$|
j                  |d<   nt        dt        |
             |dk(  rp|}|d   }| j                   j                  |      }|st        d|        ||j                  d      |j                  d	      ddi       d{   }t#        |      }n|dk(  rt|j                  d      }|j                  d      }|r|st        d       t        |t$              sJ t        |t&              sJ | j)                  ||       d{   }d!|i}nt        d"|       d#d$||d%d&}| j*                  j-                  t/        j0                  |      d'z          d{    yc c}w 7 c c}w 7 7 i7 # t3        j4                         $ r  t        $ rV}d#d(|t%        |      d)d&}| j*                  j-                  t/        j0                  |      d'z          d{  7   Y d}~yd}~ww xY ww)*z)Handle incoming control request from CLI.r   r   rd   r2   inputz#canUseTool callback is not providedNpermission_suggestionstool_use_idagent_idblocked_pathdecision_reasontitledisplay_namedescription)	signalsuggestionsr   r   r   r   r   r   r   	tool_nameallow)behaviorupdatedInputupdatedPermissionsdeny)r   r   	interruptzkTool permission callback must return PermissionResult (PermissionResultAllow or PermissionResultDeny), got hook_callbackr   zNo hook callback found for ID: r   mcp_messageserver_namer   z.Missing server_name or message for MCP requestmcp_responsez%Unsupported control request subtype: r   success)rd   r   r   )rc   r   
r]   )rd   r   r]   )r2   ri   r   rg   r   	from_dictrw   r   updated_inputupdated_permissionsto_dictr   r   r   	TypeErrorrc   rD   r,   rJ   rI   _handle_sdk_mcp_requestr0   writejsondumpsrG   r   )rX   r   r   request_datard   response_datapermission_requestoriginal_inputscontextr   
permissionhook_callback_requestr   r~   r    r   r   r   success_responserm   error_responses                         r+   r   zQuery._handle_control_requestw  s
    \*
y)y)v	J,.M.(BN"!3G!<((#$IJJ/ /223KLRPRR! S  )2215R! !3 6 6} E/33J?!3!7!7!G$6$:$:;L$M,009!3!7!7!G 2 6 6} E" "&!2!2&{3&w/"  h(=>$+  (55A %22!/%M  33? /7.J.J?.J
 '..0.J?&:;  *>?17HDTDT$UM))5=5G5Gk2# F  GK  LT  GU  FV  W  O+@L%3MB..22;?#&Ek]$STT$, $$W- $$]3t$%  !=[ IM)*..}=*..y9"+#$TUU "+s333!+t444%)%A%A&   "0 >  "Gy QRR +(", -4 ..&&tzz2B'Cd'JKKKu!"?*& $ L,,. 	  
	J +&", V2N ..&&tzz.'AD'HIII
	Js   NAL L
2BL L
AL L*C L *L+A?L *L+AL ?L L NL L L L  N8AN	>N?N	N	NNrs   c                   K   | j                   st        d      | xj                  dz  c_        d| j                   dt        j                  d      j                          }t        j                         }|| j                  |<   d||d}| j                  j                  t        j                  |      dz          d	{    	 t        j                  |      5  |j                          d	{    d	d	d	       | j                  j!                  |      }| j                  j!                  |d	       t#        |t              r||j%                  d
i       }t#        |t&              r|S i S 7 7 }# 1 sw Y   |xY w# t(        $ r[}| j                  j!                  |d	       | j                  j!                  |d	       t        d|j%                  d             |d	}~ww xY ww)zSend control request to CLI and wait for response.

        Args:
            request: The control request to send
            timeout: Timeout in seconds to wait for response (default 60s)
        z'Control requests require streaming moder   req__   r   )rc   r   r   r   Nr   zControl request timeout: rd   )r1   ri   rF   osurandomhexrG   rT   rB   r0   r   r   r   
fail_afterwaitrC   r   rw   rg   rI   TimeoutError)	rX   r   rs   r   r{   r   r   r   rm   s	            r+   ry   zQuery._send_control_request  s     %%EFF 	"D112!BJJqM4E4E4G3HI
 5:&&z2 &$
 nn""4::o#>#EFFF	Y!!'*jjl"" + 1155jAF**..z4@&),"JJz26M$.}d$C=KK 	G
 # +*  	Y**..z4@((,,Z>7I8N7OPQWXX	Ysn   B=G!?E* G!E: E..E,/E.3A4E: 'G!(E: )G!,E..E73E: :	GAGGG!r   r   c           
        K   || j                   vrd|j                  d      dd| dddS | j                   |   }|j                  d      }|j                  d	i       }	 |d
k(  r6d|j                  d      ddi i|j                  |j                  xs ddddS |dk(  r+t	        |      }|j
                  j                  t              }|rG ||       d{   }g }	|j                  j                  D ]  }
|
j                  |
j                  |
j                  r<t        |
j                  d      r|
j                  j                         n|
j                  ni d}|
j                  r|
j                  j                  d      |d<   |
j                  r|
j                  |d<   |	j                  |        d|j                  d      d|	idS |dk(  r:t        |t!        |j                  d      |j                  di                   }|j
                  j                  t              }|r ||       d{   }g }|j                  j"                  D ]k  }t%        |dd      }|d k(  r |j                  d t%        |d d!      d"       6|d#k(  r,|j                  d#t%        |d$d!      t%        |d%d!      d&       g|d'k(  rg }t%        |dd      }t%        |d(d      }t%        |d)d      }|r|j                  |       |r|j                  t'        |             |r|j                  |       |j                  d |rd*j)                  |      nd+d"       |d,k(  rRt%        |d,d      }|r,t        |d       r |j                  d |j*                  d"       ?t,        j/                  d-       Vt,        j/                  d.|       n d/|i}t        |j                  d0      r|j                  j0                  rd|d0<   d|j                  d      |dS |d1k(  rdi d2S d|j                  d      dd3| dddS 7 X7 # t2        $ r+}d|j                  d      d4t'        |      ddcY d}~S d}~ww xY ww)5a  Handle an MCP request for an SDK server.

        This acts as a bridge between JSONRPC messages from the CLI
        and the in-process MCP server. Ideally the MCP SDK would provide
        a method to handle raw JSONRPC, but for now we route manually.

        Args:
            server_name: Name of the SDK MCP server
            message: The JSONRPC message

        Returns:
            The response message
        z2.0idizServer 'z' not found)coder   )jsonrpcr   r]   methodparamsrt   z
2024-11-05toolsz1.0.0)nameversion)protocolVersioncapabilities
serverInfo)r   r   r   z
tools/list)r   N
model_dump)r   r   inputSchemaT)exclude_noneannotations_metaz
tools/callr   	arguments)r   r   )r   r   rc   textrb   )rc   r   imagedatamimeType)rc   r   r  resource_linkurir   r   zResource linkresourcez>Binary embedded resource cannot be converted to text, skippingz4Unsupported content type %r in tool result, skippingcontentisErrorznotifications/initialized)r   r   zMethod 'i)r4   rg   r   r   r   request_handlersrootr   r   r   hasattrr   r   metarv   r   r   r  getattrrJ   r   r   rj   rk   r  ri   )rX   r   r   serverr   r   r   handlerr   
tools_datatool	tool_datacall_requestr  item	item_typepartsr   r  descr  r   rm   s                          r+   r   zQuery._handle_sdk_mcp_request$  s      d222 kk$'"!)+kB  %%k2X&Xr*O	 %  %!++d++7#R) %+KK'-~~'@'	  <'*&9 11556FG#*7#33F!#J & 1 1$(II+/+;+;  $// $+4+;+;\#J !% 0 0 ; ; =%)%5%5 "$
5	  ++7;7G7G7R7R-1 8S 8Im4  9915Ig.")))4% !2( $)%kk$/#*J"7  <'.!0#ZZ/6::kSU;V  !1155oF#*<#88F G & 3 3$+D&$$?	$.#NN)/vr9R S ''1#NN,3,3D&",E07j"0M!" '/9$&E#*4#>D")$t"<C#*4#ED# %T 2" %SX 6# %T 2#NN,2', -1IIe,<)8	!" '*4'.tZ'FH'GHf,E '/V W &$d!" #NN V )U !4^ &/$8Mv{{I66;;;N;N37i0 $)%kk$/"/  66#(B77 !kk$'"(xx{5ST Q 4J 9R  	 kk$'"(SV< 	s   AQ/ :P8 Q/?P8 P2C0P8 Q/A*P8 6P57GP8 Q/	P8 Q/P8 1Q/2P8 5P8 8	Q, Q'!Q,"Q/'Q,,Q/c                 D   K   | j                  ddi       d{   S 7 w)z)Get current MCP server connection status.rd   
mcp_statusNry   r   s    r+   get_mcp_statuszQuery.get_mcp_status  s"     //L0IJJJJ     c                 D   K   | j                  ddi       d{   S 7 w)z<Get a breakdown of current context window usage by category.rd   get_context_usageNr  r   s    r+   r  zQuery.get_context_usage  s#     //<O0PQQQQr  c                 F   K   | j                  ddi       d{    y7 w)zSend interrupt control request.rd   r   Nr  r   s    r+   r   zQuery.interrupt  s     (()[)ABBBs   !!modec                 H   K   | j                  d|d       d{    y7 w)zChange permission mode.set_permission_mode)rd   r  Nr  )rX   r  s     r+   r   zQuery.set_permission_mode  s)     ((0
 	
 	
   " "modelc                 H   K   | j                  d|d       d{    y7 w)zChange the AI model.	set_model)rd   r"  Nr  )rX   r"  s     r+   r$  zQuery.set_model  s)     ((&
 	
 	
r!  user_message_idc                 H   K   | j                  d|d       d{    y7 w)zRewind tracked files to their state at a specific user message.

        Requires file checkpointing to be enabled via the `enable_file_checkpointing` option.

        Args:
            user_message_id: UUID of the user message to rewind to
        rewind_files)rd   r%  Nr  )rX   r%  s     r+   r'  zQuery.rewind_files  s+      (()#2
 	
 	
r!  c                 H   K   | j                  d|d       d{    y7 w)zReconnect a disconnected or failed MCP server.

        Args:
            server_name: The name of the MCP server to reconnect
        mcp_reconnect)rd   
serverNameNr  )rX   r   s     r+   reconnect_mcp_serverzQuery.reconnect_mcp_server   s+      ((*)
 	
 	
r!  enabledc                 J   K   | j                  d||d       d{    y7 w)zEnable or disable an MCP server.

        Args:
            server_name: The name of the MCP server to toggle
            enabled: Whether the server should be enabled
        
mcp_toggle)rd   r*  r,  Nr  )rX   r   r,  s      r+   toggle_mcp_serverzQuery.toggle_mcp_server  s.      ((')"
 	
 	
s   #!#task_idc                 H   K   | j                  d|d       d{    y7 w)zkStop a running task.

        Args:
            task_id: The task ID from task_notification events
        	stop_task)rd   r0  Nr  )rX   r0  s     r+   r2  zQuery.stop_task  s+      ((&"
 	
 	
r!  c                 P  K   | j                   s| j                  rdt        j                  dt	        | j                          dt        | j                         d       | j                  j                          d{    | j                  j                          d{    y7 '7 w)a  Wait for the first result (if needed) then close stdin.

        If SDK MCP servers or hooks require bidirectional communication,
        keeps stdin open until the first result arrives. The control protocol
        requires stdin to remain open for the entire conversation, so no
        timeout is applied. The event is guaranteed to fire: either when the
        result message arrives, or in _read_messages' finally block if the
        process exits early.
        z?Waiting for first result before closing stdin (sdk_mcp_servers=z, has_hooks=)N)
r4   r3   rj   r   lenboolrU   r   r0   	end_inputr   s    r+   wait_for_result_and_end_inputz#Query.wait_for_result_and_end_input)  s      4::LL$$'(<(<$=#> ?!$**-.a1
 **//111nn&&((( 2(s$   A8B&:B";!B&B$B&$B&streamc                 P  K   	 |2 3 d{   }| j                   r n:| j                  j                  t        j                  |      dz          d{    Q| j                          d{    y7 e7  6 7 # t        $ r"}t        j                  d|        Y d}~yd}~ww xY ww)zStream input messages to transport.

        If SDK MCP servers or hooks are present, waits for the first result
        before closing stdin to allow bidirectional control protocol communication.
        Nr   zError streaming input: )	rR   r0   r   r   r   r8  ri   rj   r   )rX   r9  r   rm   s       r+   stream_inputzQuery.stream_input=  s     	8!' Gg<<nn**4::g+>+EFFF44666G G "(
 7 	8LL21#677	8ss   B&A8 A4A0A4AA8 A2A8 *A6+A8 /B&0A42A8 4A8 6A8 8	B#BB&B##B&c                   K   | j                   2 3 d{   }|j                  d      dk(  r y|j                  d      dk(  rt        |j                  dd            | T7 O6 yw)z,Receive SDK messages (not control messages).Nrc   r   r]   r   )rL   rg   ri   rX   r   s     r+   receive_messageszQuery.receive_messagesM  s_     !22 	'{{6"e+V$/G_ EFFM	2s&   A'A%A#A%AA'#A%%A'c                    K   d| _         | j                  "| j                  j                          d{    t        | j                        D ]  }|j                           | j                  V| j                  j                         s<| j                  j                          | j                  j                          d{    d| _        | j                  j                          | j                  j                          d{    y7 7 J7 	w)zClose the query and transport.TN)rR   rW   r   rx   rO   r   rM   doner   rK   r0   )rX   r   s     r+   r   zQuery.closeX  s      **61177999**+DKKM ,??&t/C/C/EOO""$//&&((( 	  "nn""$$$! :
 ) 	%s5   1DDBD?D
 ADDD
DDc                 8    | j                   j                          y)aW  Close the receive side of the message stream.

        Call once the consumer has finished iterating ``receive_messages()``.
        ``close()`` leaves this open so a still-draining consumer can read
        buffered messages; the consumer is responsible for closing it to
        avoid a ``ResourceWarning`` from anyio's ``__del__``.
        N)rL   r   r   s    r+   close_receive_streamzQuery.close_receive_streamp  s     	##%r-   c                 "    | j                         S )z#Return async iterator for messages.)r>  r   s    r+   	__aiter__zQuery.__aiter__{  s    $$&&r-   c                 V   K   | j                         2 3 d{   }|c S 7 6 t        w)zGet next message.N)r>  StopAsyncIterationr=  s     r+   	__anext__zQuery.__anext__  s-     !224 	'N	4  s   )" ")"))NNN      N@NNN)r[   r   r!   N)r!   N)rH  )1__name__
__module____qualname____doc__r   r6  r   rJ   rI   r	   r   r   r   r   rx   floatr
   rY   r\   rn   rt   r   r   r   r   r   r   r   ry   r   r  r  r   r   r   r$  r'  r+  r/  r2  r8  r   r;  r   r>  r   rB  rD  rG   r-   r+   r/   r/   =   s9   $ 8<9=$(370448DODO  DO $sCx."78+.BBCE
 		DO Cd38n--.5DO c;./$6DO "DO S$sCx.()D0DO #'+DO S	GEN*T1DOL2Q': Q3 Q4 Q*2$sCx.4"7 2hD
s z 	&6G 	&D 	&~'@|J5F |J4 |J~ 9=-YCH~-Y05-Y	c3h-Y^mm)-c3hm	c3hm^Kd38n KRc3h RC
n 
 

S4Z 
D 

# 
$ 

c 
d 

3 
 
$ 

s 
t 
)(8tCH~)F 84 8 	d38n(E 	%0&'=c3h8 '!c3h !r-   r/   )3rL  r   loggingr   re   collections.abcr   r   r   r   
contextlibr   typingr   r	   r
   rG   	mcp.typesr   r   r   _errorsr   typesr   r   r   r   r   r   r   r   r   _task_compatr   r   r0   r   
mcp.serverr   r5   r   transcript_mirror_batcherr   	getLoggerrI  rj   rI   rJ   r,   r/   rN  r-   r+   <module>rZ     s    >   	  M M  . .   #
 
 
 5  ."B			8	$d38n c3h &F! F!r-   