ó
    ±"³j¬Ó  ã                  ó|  • S SK Jr  S SKJrJrJrJrJrJ	r	  S SK
JrJr  S SKJr  S SKJrJrJr  S SKJr  S SKJr  S SKJr  S S	KJrJrJrJrJr  S S
KrS SKJ r   S SK!J"r"  SSK#J$r$J%r%J&r'J(r(  SSK)J*r*  SSK+J,r,J-r-J.r.J/r/J0r0J1r1J2r2  SSK3J4r4J5r5J6r6  SSK7J8r8  SSK&J9r9J:r:  SSK;J<r<J=r=  SSK>J?r?  SSK@JArA  SSKBJCrCJDrD  \(       a  SSKEJFrF  SSKGJHrH  SrI\" SS9 " S S\\4\<4   5      5       rJ\" SS 9 " S! S"\\4\<4   5      5       rK " S# S$\\4\<4   5      rL\" SS%9 " S& S'\\<   5      5       rM        S*S( jrN      S+S) jrOg
),é    )Úannotations)ÚAsyncGeneratorÚAsyncIteratorÚ	AwaitableÚCallableÚIterableÚIterator)ÚAbstractAsyncContextManagerÚaclosing)Údeepcopy)Ú	dataclassÚfieldÚreplace)Údatetime)ÚDecimal)ÚTracebackType)ÚTYPE_CHECKINGÚAnyÚGenericÚcastÚoverloadN)ÚValidationError)ÚSelfé   )Ú_utilsÚ
exceptionsÚmessagesÚmodels)Úbest_effort_price)ÚOutputDataT_invÚOutputSchemaÚOutputValidatorÚOutputValidatorFuncÚTextOutputSchemaÚrun_image_process_hooksÚrun_output_with_hooks)Ú
AgentDepsTÚ
RunContextÚdispatch_event_stream)ÚSyncStreamBridge)ÚAgentStreamEventÚModelResponseStreamEvent)ÚOutputDataTÚ
ToolOutput)ÚToolManager)ÚDeferredToolRequests)ÚRunUsageÚUsageLimits)ÚAbstractCapability©ÚAgentRunResult)r-   r    r.   r#   ÚStreamedRunResultSyncT)Úkw_onlyc                  óŠ  • \ rS rSr% S\S'   S\S'   S\S'   S\S	'   S
\S'   S\S'   S\S'   S\S'   \" SSS9rS\S'   \" \SS9rS\S'   \" SSS9r	S\S'   \" SS9r
S\S'   \" SSS9rS\S '   \" \R                  SS!9rS"\S#'   \" S$ SS!9rS%\S&'   S' rS(S).SCS* jjrS(S).SDS+ jjrSS(S,.SES- jjrSFS. jrSFS/ jr\SGS0 j5       r\SHS1 j5       r\SHS2 j5       r\SIS3 j5       r\SJS4 j5       r\SKS5 j5       r\SLS6 j5       rSMS7 jrSNS8 jrSS9.     SOS: jjr SPS; jr!SS(S,.     SES< jjr"SQS= jr#SFS> jr$SRS? jr%SSS@ jr&STSA jr'SBr(g)UÚAgentStreamé3   úmodels.StreamedResponseÚ_raw_stream_responsezOutputSchema[OutputDataT]Ú_output_schemazmodels.ModelRequestParametersÚ_model_request_parametersz.list[OutputValidator[AgentDepsT, OutputDataT]]Ú_output_validatorszRunContext[AgentDepsT]Ú_run_ctxúUsageLimits | NoneÚ_usage_limitsúToolManager[AgentDepsT]Ú_tool_managerzAbstractCapability[AgentDepsT]Ú_root_capabilityNF)ÚdefaultÚreprz*Callable[[], dict[str, Any] | None] | NoneÚ_metadata_getterz$Callable[[], list[AgentStreamEvent]]Ú_event_stream_buffer_getter©rF   Úinitz&AsyncIterator[AgentStreamEvent] | NoneÚ_events_iterator©rK   r1   Ú_initial_run_ctx_usagezOutputDataT | NoneÚ_cached_output)Údefault_factoryrK   z
anyio.LockÚ_anext_lockc                 ó:   • [         [        R                     " 5       $ ©N)ÚsetÚanyioÚCancelScope© ó    ÚO/home/mande/repo/quber/.venv/lib/python3.13/site-packages/pydantic_ai/result.pyÚ<lambda>ÚAgentStream.<lambda>E   s   € ÌÌU×M^ÑM^ÒI_ÔIarX   zset[anyio.CancelScope]Ú_pull_scopesc                óL   • [        U R                  R                  5      U l        g rS   )r   r@   ÚusagerN   ©Úselfs    rY   Ú__post_init__ÚAgentStream.__post_init__G   s   € Ü&.¨t¯}©}×/BÑ/BÓ&CˆÕ#rX   çš™™™™™¹?©Údebounce_byc              óB  #   • U R                   b  [        U R                   5      7v •  gSnU R                  US9  Sh  v•N nU R                  R                  b!  U(       a  UR
                  UR
                  :X  a  MC  Un U R                  USS9I Sh  v•N 7v •  Mc   N^ N! [        [        R                  4 a     M„  f = f
 U R                  R                  bD  U R                  nU R                  U5      I Sh  v•N  U l         [        U R                   5      7v •  gg7f)z4Asynchronously stream the (validated) agent outputs.Nrd   T©Úallow_partial)rO   r   Ústream_responser<   Úfinal_result_eventÚpartsÚvalidate_response_outputr   r   Ú
ModelRetryÚresponse)r`   re   Úlast_responsern   s       rY   Ústream_outputÚAgentStream.stream_outputJ   s  é € à×ÑÑ*Ü˜4×.Ñ.Ó/Ó/Øà8<ˆØ"×2Ñ2¸{Ð2ÑK÷ 
	�(Ø×(Ñ(×;Ñ;ÑCÞ (§.¡.°M×4GÑ4GÓ"GáØ$ˆMðØ ×9Ñ9¸(ÐRVÐ9ÐW×WÕWñ
	ñ XøÜ#¤Z×%:Ñ%:Ð;ó Úðúð Lð ×$Ñ$×7Ñ7ÑCØ—}‘}ˆHð )-×(EÑ(EÀhÓ(O×"OÐ"OˆDÔÜ˜4×.Ñ.Ó/Ô/ð Dùsb   ‚8DºC¾B¿CÁ=DÂ B!ÂBÂB!ÂDÂCÂB!Â!B?Â;DÂ>B?Â?;DÃ:C=Ã;$Dc              ó¶  #   • U R                   nUR                  S:X  a/  UR                   H  nUR                  5       (       d  M  U7v •    O   [        R
                  " U R                  5       U5       ISh  v•N nU  Sh  v•N nU R                   7v •  M   N  N
 SSS5      ISh  v•N    O! , ISh  v•N  (       d  f       O= fU R                   7v •  g7f)u‰  Asynchronously stream the (unvalidated) model responses for the agent.

Yields `ModelResponse` snapshots â€” `state='incomplete'` while streaming is in flight,
followed by one final `state='complete'` snapshot (or `'interrupted'` if `cancel()` was
called). If the underlying response already has accumulated content when this is called,
a pre-stream yield surfaces it before iteration begins.
Ú
incompleteN)rn   Ústaterk   Úhas_contentr   Úgroup_by_temporalÚ_model_response_events)r`   re   ÚmsgÚpartÚ
group_iterÚ_itemss         rY   ri   ÚAgentStream.stream_responseg   s¥   é € ð �m‰mˆØ�9‰9˜Ó$ØŸ	œ	�Ø×#Ñ#×%Ó%Ø“IÙñ "ô
 ×+Ò+¨D×,GÑ,GÓ,IÈ;×WÔWÐ[eÙ *÷ $�fØ—m‘mÕ#ñ Xñ$ 
÷ X×W×W×W×WÐWúð �m‰mÔùsq   ‚?CÁ2CÁ7BÁ8CÁ;B/Á>BÂBÂBÂB/ÂCÂBÂB/ÂCÂ(B+Â)CÂ/CÂ5B8Â6CÃC©Údeltare   c          
    óò  #   • [        U R                  [        5      (       d  [        R                  " S5      e[        U R
                  [        5      (       a  U R
                  7v •  gU(       a   U R                  SUS9  Sh  v•N nU7v •  M  U R                  SUS9  Sh  v•N nU R                   H/  nUR                  U[        U R                  SS95      I Sh  v•N nM1     U7v •  MO   Nj
 g NN N
 g7f)uÈ  Stream the text result as an async iterable.

!!! note
    [`TextOutput`][pydantic_ai.output.TextOutput] functions are not applied â€” use
    [`stream_output()`][pydantic_ai.result.AgentStream.stream_output] instead.
    Result validators will NOT be called on the text result if `delta=True`.

Args:
    delta: if `True`, yield each chunk of text as it is received, if `False` (default), yield the full text
        up to the current point.
    debounce_by: by how much (if at all) to debounce/group the response chunks by. `None` means no debouncing.
        Debouncing is particularly important for long structured responses to reduce the overhead of
        performing validation as each token is received.
ú2stream_text() can only be used with text responsesNTr}   F©Úpartial_output)Ú
isinstancer=   r$   r   Ú	UserErrorrO   ÚstrÚ_stream_response_textr?   Úvalidater   r@   )r`   r~   re   ÚtextÚ	validators        rY   Ústream_textÚAgentStream.stream_text|   sÛ   é € ô ˜$×-Ñ-Ô/?×@Ñ@Ü×&Ò&Ð'[Ó\Ð\ô
 �d×)Ñ)¬3×/Ñ/Ø×%Ñ%Ó%ØæØ"×8Ñ8¸tÐQ\Ð8Ñ]÷ �dØ•
à"×8Ñ8¸uÐR]Ð8Ñ^÷ �dØ!%×!8Ô!8�IØ!*×!3Ñ!3°D¼'À$Ç-Á-Ð`dÑ:eÓ!f×f’Dñ "9à•
ñÑ]ñáfñ _ùs`   ‚A<C7Á>C/ÂC-ÂC/ÂC7ÂC5Â"C1Â#C5Â&7C7ÃC3ÃC7Ã-C/Ã/C7Ã1C5Ã3C7Ã5C7c              ƒ  óT   #   • U R                   R                  5       I Sh  v•N   g N7f)aw  Cancel local stream consumption and request provider shutdown.

Whether this stops remote generation or closes the underlying transport depends on the provider SDK.

This stops only the current model response; the run continues. To end the whole run,
use [`AgentRun.cancel()`][pydantic_ai.run.AgentRun.cancel] or
[`RunContext.cancel()`][pydantic_ai.tools.RunContext.cancel].
N)r<   Úcancelr_   s    rY   r�   ÚAgentStream.cancelž   s   é € ð ×'Ñ'×.Ñ.Ó0×0Ó0ùs   ‚( &¡(c              ƒ  ó,   #   • U   Sh  v•N nM   N
 g7f)z>Consume all remaining events from the stream, discarding them.NrW   ©r`   Ú_s     rY   ÚdrainÚAgentStream.drain©   s   é € á÷ 	�!Ùñ	‘tùs   ‚…‰Š��’c                ó.   • U R                   R                  $ )ú5Whether the stream has been cancelled via `cancel()`.)r<   Ú	cancelledr_   s    rY   r–   ÚAgentStream.cancelled®   ó   € ð ×(Ñ(×2Ñ2Ð2rX   c                ó`   • U R                   R                  c   eU R                   R                  $ ©ú(The unique identifier for the agent run.)r@   Úrun_idr_   s    rY   rœ   ÚAgentStream.run_id³   s*   € ð �}‰}×#Ñ#Ñ/Ð/Ð/Ø�}‰}×#Ñ#Ð#rX   c                ó`   • U R                   R                  c   eU R                   R                  $ ©ú?The unique identifier for the conversation this run belongs to.)r@   Úconversation_idr_   s    rY   r¡   ÚAgentStream.conversation_id¹   s*   € ð �}‰}×,Ñ,Ñ8Ð8Ð8Ø�}‰}×,Ñ,Ð,rX   c                óh   • U R                   b  U R                  5       $ U R                  R                  $ ©ú7Metadata associated with this agent run, if configured.)rH   r@   Úmetadatar_   s    rY   r¦   ÚAgentStream.metadata¿   s/   € ð × Ñ Ñ,Ø×(Ñ(Ó*Ð*Ø�}‰}×%Ñ%Ð%rX   c                ó6   • U R                   R                  5       $ )z&Get the current state of the response.)r<   Úgetr_   s    rY   rn   ÚAgentStream.responseÆ   s   € ð ×(Ñ(×,Ñ,Ó.Ð.rX   c                óÖ  • U R                   U R                  R                  -   nU R                  R                  R                  c¤  [	        U R                  R                  U R
                  R                  U R
                  R                  U R
                  R                  U R
                  R                  S9nUb0  UR                  =(       d    [        S5      UR                  -   Ul        U$ )úpReturn the usage of the whole run.

!!! note
    This won't return the full usage until the stream is finished.
)Ú
model_nameÚprovider_api_urlÚprovider_nameÚgenai_request_timestampr   )rN   r<   r^   Úcostr   rn   r­   Úprovider_urlr¯   Ú	timestampr   Útotal_price)r`   r^   Úprices      rY   r^   ÚAgentStream.usageË   s±   € ð ×+Ñ+¨d×.GÑ.G×.MÑ.MÑMˆØ×$Ñ$×*Ñ*×/Ñ/Ñ7Ü%Ø×)Ñ)×/Ñ/ØŸ=™=×3Ñ3Ø!%§¡×!;Ñ!;Ø"Ÿm™m×9Ñ9Ø(,¯©×(?Ñ(?ñˆEð Ñ Ø#Ÿj™j×6¬G°A«J¸%×:KÑ:KÑK�”
ØˆrX   c                ó.   • U R                   R                  $ ©ú"Get the timestamp of the response.)r<   r³   r_   s    rY   r³   ÚAgentStream.timestampâ   r˜   rX   c              ƒ  óê   #   • U R                   b  [        U R                   5      $ U   Sh  v•N nM   N
 U R                  U R                  5      I Sh  v•N  U l         [        U R                   5      $ 7f)z=Stream the whole response, validate the output and return it.N)rO   r   rl   rn   r�   s     rY   Ú
get_outputÚAgentStream.get_outputç   sh   é € à×ÑÑ*Ü˜D×/Ñ/Ó0Ð0ñ ÷ 	�!Ùñ	�tð %)×$AÑ$AÀ$Ç-Á-Ó$P×PÐPˆÔÜ˜×+Ñ+Ó,Ð,ùs+   ‚%A3§4«2¬4¯A3²4´A3ÁAÁA3c                óî   • U R                   c'  U R                  (       a  [        R                  " S5      eU R                  R
                  n[        [        U R                   5      Ub  UR                  4$ S4$ )aB  The validated output and the name of the output tool that produced it, if any.

Every path that consumes a stream to the end validates the response and caches the result,
so by then `_cached_output` holds it. Cancelling does not: it marks the stream complete
without a final response, and there is no output to settle on.
Nz±The stream was cancelled before it produced an output, so this run has no settled result. The messages recorded up to the interruption are still available from `all_messages()`.)	rO   r–   r   r„   r<   rj   r   r-   Ú	tool_name)r`   rj   s     rY   Ú_settled_outputÚAgentStream._settled_outputô   sx   € ð ×ÑÑ&¨4¯>¯>Ü×&Ò&ð$óð ð
 "×6Ñ6×IÑIÐä”˜d×1Ñ1Ó2Ø,>Ñ,JÐ×(Ñ(ð
ð 	
àPTð
ð 	
rX   rg   c             ƒ  ó`  ^#   • U R                   R                  nUc  [        R                  " S5      eUR                  m U R
                  R                  (       an  Tbk  [        U4S jUR                   5       S5      nUc  [        R                  " ST< 35      eU R                  R                  UU R
                  USS9I Sh  v•N $ [        UR                  U R                  5      =n(       aA  U R
                  R                  (       d  [        R                  " S5      e[        [        U5      $ U R
                  R                   (       a6  UR"                  (       a%  U R%                  UR"                  S   US	9I Sh  v•N $ U R
                  R&                  =n(       a±  S
nUR(                   HU  n[+        U[,        R.                  5      (       a  XxR0                  -  nM2  [+        U[,        R2                  5      (       d  MS  S
nMW     [5        U R6                  US9n	[9        UUU	U R:                  U R
                  USU R<                  S9I Sh  v•N $ [        R                  " S5      e GN¢ Në N! [>        [        R@                  4 a$  n
U(       d  [        R                  " S5      U
ee Sn
A
ff = f7f)ú%Validate a structured result message.Nz'Invalid response, unable to find outputc              3  óJ   >#   • U  H  oR                   T:X  d  M  Uv •  M     g 7frS   )r¿   )Ú.0ry   Úoutput_tool_names     €rY   Ú	<genexpr>Ú7AgentStream.validate_response_output.<locals>.<genexpr>  s   øé € Ð_Ò&8˜d¿N¹NÐN^Ñ<^—T‘TÒ&8ùs   ƒ#š	#z/Invalid response, unable to find tool call for F)Úschemarh   Úwrap_validation_errorsz¯A deferred tool call was present, but `DeferredToolRequests` is not among output types. To resolve this, add `DeferredToolRequests` to the list of output types for this agent.r   rg   Ú r�   )rˆ   Úrun_contextÚ
capabilityrÉ   rh   rÊ   Úoutput_validatorsz/Invalid response, unable to process text outputzZOutput validation failed during streaming, and retries are not supported in `run_stream()`)!r<   rj   r   ÚUnexpectedModelBehaviorr¿   r=   ÚtoolsetÚnextÚ
tool_callsrD   Úhandle_output_tool_callÚ_get_deferred_tool_requestsÚallows_deferred_toolsr„   r   r-   Úallows_imageÚimagesÚ_validate_image_outputÚtext_processorrk   rƒ   Ú	_messagesÚTextPartÚcontentÚNativeToolCallPartr   r@   r&   rE   r?   r   rm   )r`   Úmessagerh   rj   Ú	tool_callÚdeferred_tool_requestsrÙ   rˆ   ry   Úrun_ctxÚerÆ   s              @rY   rl   Ú$AgentStream.validate_response_output  s|  øé € ð "×6Ñ6×IÑIÐØÑ%Ü×4Ò4Ð5^Ó_Ð_à-×7Ñ7Ðð6	Ø×"Ñ"×*×*Ð/?Ñ/KÜ Ü_ g×&8Ò&8Ó_Øó�	ð Ñ$Ü$×<Ò<ØIÐJZÑI]Ð^óð ð "×/Ñ/×GÑGØØ×.Ñ.Ø"/Ø+0ð	 Hð ÷ ð ô ,GÀw×GYÑGYÐ[_×[mÑ[mÓ+nÐnÐ'ÕnØ×*Ñ*×@×@Ü$×.Ò.ð Jóð ô œKÐ)?Ó@Ð@Ø×$Ñ$×1×1°g·n·nØ!×8Ñ8¸¿¹ÈÑ9JÐZgÐ8Ðh×hÐhØ#'×#6Ñ#6×#EÑ#EÐE�ÕEØ�Ø#ŸMœM�DÜ! $¬	×(:Ñ(:×;Ñ;Ø§¡Ñ,šÜ# D¬)×*FÑ*F×GÓGð  "šñ *ô " $§-¡-ÀÑN�Ü2Ø"ØØ 'Ø#×4Ñ4Ø×.Ñ.Ø"/Ø+0Ø&*×&=Ñ&=ñ	÷ 	ð 	ô !×8Ò8ØEóð òIñ iñ	øô  ¤×!6Ñ!6Ð7ó 	Þ Ü ×8Ò8Øpóàðð ûð	üsŽ   ƒ<J.Á BI- ÃI&ÃI- ÃJ.Ã	A'I- Ä0J.Ä1AI- Å=I)Å>I- ÆJ.ÆA:I- È AI- ÉI+ÉI- ÉJ.ÉI- É)I- É+I- É-J+ÊJ&Ê&J+Ê+J.c             ƒ  óº   #   • [        U R                  US9n[        [        [	        UU R
                  UU R                  SU R                  S9I Sh  v•N 5      $  N7f)zARun process hooks (including output validators) for image output.r�   F)rÍ   rÌ   rÉ   rÊ   rÎ   N)r   r@   r   r-   r%   rE   r=   r?   )r`   Úimagerh   rá   s       rY   rØ   Ú"AgentStream._validate_image_outputI  s[   é € ä˜$Ÿ-™-¸ÑFˆÜÜÜ)ØØ×0Ñ0Ø#Ø×*Ñ*Ø',Ø"&×"9Ñ"9ñ÷ ó

ð 
	
ñùs   ‚AAÁA
Á	Ac              óh  ^ ^^#   • SU 4S jjmSUU4S jjn[        U" 5       5       ISh  v•N nU(       a  U  Sh  v•N nU7v •  M  / nU  Sh  v•N nUR                  U5        SR                  U5      7v •  M0   NQ NA
 O N1
 SSS5      ISh  v•N    g! , ISh  v•N  (       d  f       g= f7f)z1Stream the response as an async iterable of text.c                ó<  >#   • TR                   n [        U R                  5       HJ  u  p[        U[        R
                  5      (       d  M&  UR                  (       d  M9  UR                  U47v •  ML     S nT  S h  v•N n[        U[        R                  5      (       aw  [        UR                  [        R
                  5      (       aN  UR                  R                  (       a3  UR                  nUR                  R                  UR                  47v •  MŸ  [        U[        R                  5      (       ax  [        UR                  [        R                  5      (       aO  UR                  R                  (       a4  UR                  nUR                  R                  UR                  47v •  GM6  [        U[        R                  5      (       d  GMX  [        UR                  [        R                  5      (       d  GM„  Uc  GMŠ  SUR                  47v •  S nGM    GNœ
 g 7f)Nz

)rn   Ú	enumeraterk   rƒ   rÚ   rÛ   rÜ   ÚPartStartEventry   ÚindexÚPartDeltaEventr~   ÚTextPartDeltaÚcontent_deltarÝ   )rx   Úiry   Úlast_text_indexÚeventr`   s        €rY   Ú_stream_text_deltas_ungroupedÚHAgentStream._stream_response_text.<locals>._stream_text_deltas_ungrouped`  s_  øé € ð —-‘-ˆCÜ$ S§Y¡YÖ/‘�Ü˜d¤I×$6Ñ$6×7Ó7¸D¿L¿L¹LØŸ,™,¨˜/Õ)ñ 0ð +/ˆOÙ#÷ +�eä˜u¤i×&>Ñ&>×?Ñ?Ü" 5§:¡:¬y×/AÑ/A×BÑBØŸ
™
×*×*à&+§k¡k�OØŸ*™*×,Ñ,¨e¯k©kÐ9Õ9ä˜u¤i×&>Ñ&>×?Ñ?Ü" 5§;¡;´	×0GÑ0G×HÑHØŸ™×1×1à&+§k¡k�OØŸ+™+×3Ñ3°U·[±[Ð@Ö@ä˜u¤i×&>Ñ&>×?Ô?Ü" 5§:¡:¬y×/KÑ/K×LÔLØ'Ô3ð ! %§+¡+Ð-Ó-Ø&*“Oò-+™tùsI   ƒAHÁHÁHÁ7HÁ;HÁ<HÁ?EHÇ'HÇ;HÈHÈHÈHc            
    ó0  >#   • [         R                  " T" 5       T5       IS h  v•N n U   S h  v•N nSR                  U VVs/ s H  u  p#UPM	     snn5      7v •  M4   N: N1s  snnf 
 S S S 5      IS h  v•N    g ! , IS h  v•N  (       d  f       g = f7f)NrË   )r   rv   Újoin)rz   ÚitemsrÜ   r‘   rò   re   s       €€rY   Ú_stream_text_deltasÚ>AgentStream._stream_response_text.<locals>._stream_text_deltas‚  su   øé € Ü×/Ò/Ñ0MÓ0OÐQ\×]Ô]ÐakÙ#-÷ E˜%àŸ'™'¹UÔ"CºU©z¨w£7¹UÒ"CÓDÕDñ ^ñEùã"Cð $.÷ ^×]×]×]×]Ð]üsp   ƒ!B¤A¥B¨A<«A)¯A!°A)³A<ÁA#ÁA<ÁBÁ!A)Á#A<Á*BÁ5A8Á6BÁ<BÂBÂBÂBNrË   )ÚreturnzAsyncIterator[tuple[str, int]])rù   zAsyncGenerator[str, None])r   Úappendrõ   )r`   r~   re   r÷   Údeltas_iterrˆ   Údeltasrò   s   ` `    @rY   r†   Ú!AgentStream._stream_response_textX  s�   úé € ÷ 	+÷D	Eð 	Eô Ñ/Ó1×2Ô2°kÞÙ"-÷ ˜$Ø•Jð %'�Ù"-÷ *˜$Ø—M‘M $Ô'ØŸ'™' &›/Õ)ñ 3ñ¡+ñ* +÷ 3×2×2×2×2Ð2üs‰   …&B2«A=¬B2¯
B¹B½A?¾BÁBÁBÁBÁBÁ(BÁ=B2Á?BÂBÂBÂBÂB2ÂBÂB2ÂB/ÂB!ÂB/Â+B2c                óH  ^ • T R                   cz  [        T R                  T R                  U 4S j5      n[	        T R
                  R                  T R                  [        T R                  T R                  U5      5      S95      T l         T R                  T R                   5      $ )z}Stream [`AgentStreamEvent`][pydantic_ai.messages.AgentStreamEvent]s, interleaving events emitted into the run's event buffer.c                 óJ   >• T R                   T R                  R                  -   $ rS   )rN   r<   r^   r_   s   €rY   rZ   Ú'AgentStream.__aiter__.<locals>.<lambda>œ  s   ø€ ˜×3Ñ3°d×6OÑ6O×6UÑ6UÒUrX   )Ústream)rL   Ú#_get_usage_checking_stream_responser<   rB   ÚaiterrE   Úwrap_run_event_streamr@   r)   Ú_events_iterÚ_pull_shared)r`   Ú	base_iters   ` rY   Ú	__aiter__ÚAgentStream.__aiter__”  s’   ø€ à× Ñ Ñ(ô <Ø×)Ñ)Ø×"Ñ"ÜUóˆIô %*Ø×%Ñ%×;Ñ;Ø—M‘MÔ*?ÀÇÁÈt×O`ÑO`ÐajÓOkÓ*lð <ð ó%ˆDÔ!ð × Ñ  ×!6Ñ!6Ó7Ð7rX   c              ƒ  ó”  #   • U R                   nUb„  U R                   H  nUR                  5         M     [        R                  " SS9   U R
                   ISh  v•N   [        R                  " U5      I Sh  v•N   SSS5      ISh  v•N   SSS5        gg N< N  N! , ISh  v•N  (       d  f       N'= f! , (       d  f       g= f7f)aW  Close the event stream when a consumer walks away before exhausting it.

The event iterator owns the capability chain, which can otherwise stay suspended with
resources held, like a `ProcessEventStream` handler task parked on its receive stream.

Any in-flight shared pull is cancelled and drained before the iterator is closed. The
close is shielded because graph teardown can run inside an already-cancelled scope.

The closed iterator is kept in place rather than discarded, so a later `__aiter__()` ends
immediately instead of building a second chain (and a second handler) over a spent stream.
NT)Úshield)rL   r\   r�   rU   rV   rQ   r   Úaclose_if_supported)r`   Úevents_iteratorÚscopes      rY   Úaclose_eventsÚAgentStream.aclose_events¨  s“   é € ð ×/Ñ/ˆØÑ&Ø×*Ô*�Ø—‘–ñ +ä×"Ò"¨$Ó/Ø×+×+Ó+Ü ×4Ò4°_ÓE×EÐE÷ ,×+÷ 0Ð/ð 'ñ ,ÙE÷ ,×+×+Ð+ú÷ 0Õ/üsx   ‚ACÁ	B7ÁBÁB7ÁBÁ8BÁ9BÁ=B7ÂBÂ	B7Â
CÂB7ÂBÂB7ÂB4	Â#B&Â$B4	Â0B7Â7
CÃCc               ó¼  #   •  U R                    IS h  v•N   S n[        R                  " 5        nU R                  R	                  U5          [        U5      I S h  v•N n U R                  R                  U5         S S S 5        WR                  (       a   S S S 5      IS h  v•N   g Uc   eS S S 5      IS h  v•N   W7v •  MÁ   N° Nj! [         a7     U R                  R                  U5        S S S 5        S S S 5      IS h  v•N    g f = f! U R                  R                  U5        f = f! , (       d  f       Nµ= f Nš N†! , IS h  v•N  (       d  f       N›= f7frS   )	rQ   rU   rV   r\   ÚaddÚanextÚStopAsyncIterationÚdiscardÚcancel_called)r`   r  rñ   r  s       rY   r  ÚAgentStream._pull_shared¼  s  é € ð Ø×'×'Ó'Ø15�Ü×&Ò&Ô(¨EØ×%Ñ%×)Ñ)¨%Ô0ð9ð#Ü*/°Ó*@×$@™Eð ×)Ñ)×1Ñ1°%Õ8÷ )ð ×&×&Ø÷ (×'Ð'ð Ñ(Ð(Ð(÷ (×'ð ‹Kñ Ù'ñ %AøÜ1ó #Ø"à×)Ñ)×1Ñ1°%Ô8÷ )÷ (×'Ñ'ð#ûð ×)Ñ)×1Ñ1°%Õ8ú÷ )Õ(ú÷ (×'×'Ò'üsÙ   ‚E”C•E˜E°D-ÁC	ÁCÁC	Á"D-Á>EÂEÂ#D>Â$EÂ)EÂ.EÂ9E Â:EÃC	Ã	
D
ÃDÃD-Ã/EÃ7EÄDÄEÄ	D
Ä
DÄD*Ä*D-Ä-
D;	Ä7EÄ>EÅ EÅEÅEÅ	EÅEc               óÞ   #   • U   Sh  v•N n[        U[        R                  [        R                  -  [        R                  -  [        R
                  -  5      (       d  M]  U7v •  Md   N_
 g7f)zcIterate only the model response stream events, dropping events emitted into the run's event buffer.N)rƒ   rÚ   rê   rì   ÚPartEndEventÚFinalResultEvent)r`   rñ   s     rY   rw   Ú"AgentStream._model_response_eventsÑ  sb   é € á÷ 	�%ÜØÜ×(Ñ(Ü×*Ñ*ñ+ä×(Ñ(ñ)ô ×,Ñ,ñ-÷ó ð •ñ	™4ùs,   ‚A-…A+‰A)ŠA+�AA-Á"A-Á)A+Á+A-c               óî   #   •  U R                  5       =n(       a-  UR                  S5      7v •  U R                  5       =n(       a  M-   [        U5      I S h  v•N nU7v •  M_   N! [         a     g f = f7f)Nr   )rI   Úpopr  r  )r`   r  Úbufferrñ   s       rY   r  ÚAgentStream._events_iterÝ  sv   é € Øð !×<Ñ<Ó>Ð>�&Õ>Ø—j‘j “mÓ#ð !×<Ñ<Ó>Ð>�&×>ðÜ# IÓ.×.�ð ‹Kñ ñ /øÜ%ó Ùðüs<   ‚AA5Á	A% ÁA#ÁA% ÁA5Á#A% Á%
A2Á/A5Á1A2Á2A5)rO   rL   rN   ©re   úfloat | Nonerù   zAsyncIterator[OutputDataT]©re   r!  rù   z&AsyncIterator[_messages.ModelResponse]©r~   Úboolre   r!  rù   zAsyncIterator[str]©rù   ÚNone©rù   r$  ©rù   r…   ©rù   zdict[str, Any] | None©rù   ú_messages.ModelResponse©rù   r1   ©rù   r   ©rù   r-   )rù   ztuple[OutputDataT, str | None]©rÞ   r+  rh   r$  rù   r-   )rå   z_messages.BinaryImagerh   r$  rù   r-   )rù   úAsyncIterator[AgentStreamEvent])r  r0  rù   r0  )rù   ú'AsyncIterator[ModelResponseStreamEvent])r  r1  rù   r0  ))Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__annotations__r   rH   ÚlistrI   rL   rN   rO   rU   ÚLockrQ   r\   ra   rp   ri   rŠ   r�   r’   Úpropertyr–   rœ   r¡   r¦   rn   r^   r³   r¼   rÀ   rl   rØ   r†   r  r  r  rw   r  Ú__static_attributes__rW   rX   rY   r9   r9   3   sô  ‡ à1Ó1Ø-Ó-Ø<Ó<ØFÓFØ$Ó$Ø%Ó%Ø*Ó*Ø4Ó4ÙCHÐQUÐ\aÑCbÐÐ@ÓbÙHMÐVZÐafÑHgÐÐ!EÓgá?DÈTÐX]Ñ?^ÐÐ<Ó^Ù',°%Ñ'8Ð˜HÓ8Ù).°tÀ%Ñ)H€NÐ&ÓHá#°E·J±JÀUÑK€K�ÓKÙ+0ÑAaÐhmÑ+n€LÐ(ÓnòDð BE÷ 0ð: DG÷ ð* 27ÐTW÷  ôD	1ôð
 ó3ó ð3ð ó$ó ð$ð
 ó-ó ð-ð
 ó&ó ð&ð ó/ó ð/ð óó ðð, ó3ó ð3ô-ô
ð( JOñ@Ø.ð@ØBFð@à	õ@ôD
ð   %À#ñ:*Øð:*Ø3?ð:*à	õ:*ôx8ô(Fô(ô*
÷rX   r9   FrM   c                  óP  • \ rS rSr% SrS\S'   S\S'   SrS\S	'   SrS
\S'   SrS\S'   \	" SSS9r
S\S'    \          S/S j5       r\        S0S j5       r   S1           S2S jjr\S3S j5       rSS.S4S jjrSS.S5S jjrSS.S4S jjrSS.S5S jjrSS.S6S jjrSSS.S7S jjrSS.S8S  jjrS9S! jr\S:S" j5       r\S;S# j5       r\S<S$ j5       r\S=S% j5       r\S>S& j5       r\S>S' j5       rSS(.     S?S) jjrS@S* jrSASBS+ jjrSCS, jr \SDS- j5       r!S.r"g)EÚStreamedRunResultií  zFResult of a streamed run that returns structured data via a tool call.úlist[_messages.ModelMessage]Ú_all_messagesÚintÚ_new_message_indexNú+AgentStream[AgentDepsT, OutputDataT] | NoneÚ_stream_responseú$Callable[[], Awaitable[None]] | NoneÚ_on_completeú"AgentRunResult[OutputDataT] | NoneÚ_run_resultFrJ   r$  Úis_completec                ó   • g rS   rW   )r`   Úall_messagesÚnew_message_indexri   Úon_completes        rY   Ú__init__ÚStreamedRunResult.__init__  ó   € ð rX   c               ó   • g rS   rW   )r`   rI  rJ  Ú
run_results       rY   rL  rM    rN  rX   c                óN   • Xl         X l        X0l        X@l        XPl        S U l        g rS   )r>  r@  rB  rD  rF  Ú_traceparent_value)r`   rI  rJ  ri   rK  rP  s         rY   rL  rM    s/   € ð *ÔØ"3Ôà /ÔØ'ÔØ%ÔØ.2ˆÔð	rX   c                ó¢  • SSK Jn  SSKJn  U R                  =nb  U$ U R
                  b›  U R                  (       d  [        R                  " S5      eU R
                  R                  5       u  pEU" UUU" U R                  U R                  U R                  U R                  U R                  S9U R                  U R                   S9$ [#        S5      e)uŒ  This run as an [`AgentRunResult`][pydantic_ai.run.AgentRunResult], once the stream has finished.

A `StreamedRunResult` reads its values off the stream that is producing them, so it lives only as
long as that stream does. `AgentRunResult` is the settled form of the same run â€” the same output,
messages, usage and IDs â€” and it has a
[serialized shape](../message-history.md#storing-complete-run-results) you can store and load. Reach
for this to hand the run to code that outlives the stream.

Raises:
    UserError: If the stream hasn't finished, or was cancelled before producing an output,
        so the run has no settled output to settle on.
r   )ÚGraphAgentStater4   z“The run is still streaming, so it has no settled result yet. Await `stream_output()`, `stream_text()`, `stream_response()` or `get_output()` first.)Úmessage_historyr^   rœ   r¡   r¦   )ÚoutputÚ_output_tool_nameÚ_stater@  rR  ú)No stream response or run result provided)Ú_agent_graphrT  Úrunr5   rF  rB  rG  r   r„   rÀ   r>  r^   rœ   r¡   r¦   r@  rR  Ú
ValueError)r`   rT  r5   rP  rV  rÆ   s         rY   ÚresultÚStreamedRunResult.result*  sÌ   € õ 	2Ý'à×*Ñ*Ð*ˆJÑ7ØÐØ×"Ñ"Ñ.Ø×#×#Ü ×*Ò*ðTóð ð (,×'<Ñ'<×'LÑ'LÓ'NÑ$ˆFÙ!ØØ"2Ù&Ø$(×$6Ñ$6ØŸ*™*ØŸ;™;Ø$(×$8Ñ$8Ø!Ÿ]™]ñð $(×#:Ñ#:Ø#'×#:Ñ#:ñð ô ÐHÓIÐIrX   ©Úoutput_tool_return_contentc               ó6   • Ub  [        S5      eU R                  $ )až  Return the history of _messages.

Args:
    output_tool_return_content: The return content of the tool call to set in the last message.
        This provides a convenient way to modify the content of the output tool call if you want to continue
        the conversation and want to set the response to the output tool call. If `None`, the last message will
        not be modified.

Returns:
    List of messages.
zISetting output tool return content is not supported for this result type.)ÚNotImplementedErrorr>  ©r`   r`  s     rY   rI  ÚStreamedRunResult.all_messagesT  s"   € ð &Ñ1Ü%Ð&qÓrÐrØ×!Ñ!Ð!rX   c               óZ   • [         R                  R                  U R                  US95      $ )aý  Return all messages from [`all_messages`][pydantic_ai.result.StreamedRunResult.all_messages] as JSON bytes.

Args:
    output_tool_return_content: The return content of the tool call to set in the last message.
        This provides a convenient way to modify the content of the output tool call if you want to continue
        the conversation and want to set the response to the output tool call. If `None`, the last message will
        not be modified.

Returns:
    JSON bytes representing the messages.
r_  )rÚ   ÚModelMessagesTypeAdapterÚ	dump_jsonrI  rc  s     rY   Úall_messages_jsonÚ#StreamedRunResult.all_messages_jsone  ó/   € ô ×1Ñ1×;Ñ;Ø×ÑÐ9SÐÐTó
ð 	
rX   c               ó:   • U R                  US9U R                  S $ )á  Return the messages produced during this run.

Messages provided via `message_history` and messages from older runs are excluded.

Args:
    output_tool_return_content: The return content of the tool call to set in the last message.
        This provides a convenient way to modify the content of the output tool call if you want to continue
        the conversation and want to set the response to the output tool call. If `None`, the last message will
        not be modified.

Returns:
    List of new messages.
r_  N)rI  r@  rc  s     rY   Únew_messagesÚStreamedRunResult.new_messagesu  s(   € ð × Ñ Ð<VÐ ÐWÐX\×XoÑXoÐXqÐrÐrrX   c               óZ   • [         R                  R                  U R                  US95      $ )a  Return new messages from [`new_messages`][pydantic_ai.result.StreamedRunResult.new_messages] as JSON bytes.

Args:
    output_tool_return_content: The return content of the tool call to set in the last message.
        This provides a convenient way to modify the content of the output tool call if you want to continue
        the conversation and want to set the response to the output tool call. If `None`, the last message will
        not be modified.

Returns:
    JSON bytes representing the new messages.
r_  )rÚ   rf  rg  rm  rc  s     rY   Únew_messages_jsonÚ#StreamedRunResult.new_messages_json…  rj  rX   rc   rd   c              ó^  #   • U R                   b2  U R                   R                  7v •  U R                  5       I Sh  v•N   gU R                  b)  U R                  R	                  US9  Sh  v•N nU7v •  M  [        S5      e NF N
 U R                  U R
                  5      I Sh  v•N    g7f)a  Stream the output as an async iterable.

The pydantic validator for structured data will be called in
[partial mode](https://docs.pydantic.dev/dev/concepts/experimental/#partial-validation)
on each iteration.

Args:
    debounce_by: by how much (if at all) to debounce/group the output chunks by. `None` means no debouncing.
        Debouncing is particularly important for long structured outputs to reduce the overhead of
        performing validation as each token is received.

Returns:
    An async iterable of the response data.
Nrd   rY  )rF  rV  Ú_marked_completedrB  rp   rn   r\  )r`   re   rV  s      rY   rp   ÚStreamedRunResult.stream_output•  sš   é € ð ×ÑÑ'Ø×"Ñ"×)Ñ)Ó)Ø×(Ñ(Ó*×*Ñ*Ø×"Ñ"Ñ.Ø $× 5Ñ 5× CÑ CÐP[Ð CÑ \÷ �fØ•ô ÐHÓIÐIñ +ñÐ \à×(Ñ(¨¯©Ó7×7Ò7ùs?   ‚:B-¼B½+B-Á(BÁ,BÁ-BÁ0B-ÂBÂB-Â&B)Â'B-r}   c              óÜ  #   • U R                   bq  [        U R                   R                  [        5      (       d  [        R
                  " S5      eU R                   R                  7v •  U R                  5       I Sh  v•N   gU R                  b)  U R                  R                  XS9  Sh  v•N nU7v •  M  [        S5      e NF N
 U R                  U R                  5      I Sh  v•N    g7f)uÎ  Stream the text result as an async iterable.

!!! note
    [`TextOutput`][pydantic_ai.output.TextOutput] functions are not applied â€” use
    [`stream_output()`][pydantic_ai.result.StreamedRunResult.stream_output] instead.
    Result validators will NOT be called on the text result if `delta=True`.

Args:
    delta: if `True`, yield each chunk of text as it is received, if `False` (default), yield the full text
        up to the current point.
    debounce_by: by how much (if at all) to debounce/group the response chunks by. `None` means no debouncing.
        Debouncing is particularly important for long structured responses to reduce the overhead of
        performing validation as each token is received.
Nr€   r}   rY  )rF  rƒ   rV  r…   r   r„   rs  rB  rŠ   rn   r\  )r`   r~   re   rˆ   s       rY   rŠ   ÚStreamedRunResult.stream_text®  sÆ   é € ð ×ÑÑ'ô ˜d×.Ñ.×5Ñ5´s×;Ñ;Ü ×*Ò*Ð+_Ó`Ð`Ø×"Ñ"×)Ñ)Ó)Ø×(Ñ(Ó*×*Ñ*Ø×"Ñ"Ñ.Ø"×3Ñ3×?Ñ?ÀeÐ?Ñe÷ �dØ•
ô ÐHÓIÐIñ +ñÐeà×(Ñ(¨¯©Ó7×7Ò7ùsB   ‚A9C,Á;CÁ<+C,Â'CÂ+CÂ,CÂ/C,ÃCÃC,Ã%C(Ã&C,c              óH  #   • U R                   b(  U R                  7v •  U R                  5       I Sh  v•N   gU R                  b-  SnU R                  R	                  US9  Sh  v•N nU7v •  UnM  [        S5      e NJ N
 Uc   eU R                  U5      I Sh  v•N    g7f)aŸ  Stream the response as an async iterable of `ModelResponse` snapshots.

Each yielded `ModelResponse` is the current state of the response: `response.state` is
`'incomplete'` while streaming is in flight and `'complete'` (or `'interrupted'` if
[`cancel()`][pydantic_ai.result.StreamedRunResult.cancel] was called) on the final yield.

Args:
    debounce_by: by how much (if at all) to debounce/group the response chunks by. `None` means no debouncing.
        Debouncing is particularly important for long structured responses to reduce the overhead of
        performing validation as each token is received.

Returns:
    An async iterable of `ModelResponse` snapshots.
Nrd   rY  )rF  rn   rs  rB  ri   r\  )r`   re   Úlast_msgrx   s       rY   ri   Ú!StreamedRunResult.stream_responseÌ  s¥   é € ð ×ÑÑ'Ø—-‘-ÓØ×(Ñ(Ó*×*Ñ*Ø×"Ñ"Ñ.Ø7;ˆHØ!×2Ñ2×BÑBÈ{ÐBÑ[÷ �cØ“	Ø’ô ÐHÓIÐIñ +ñÐ[ð Ñ'Ð'Ð'Ø×(Ñ(¨Ó2×2Ò2ùs?   ‚0B"²A=³-B"Á BÁ$A?Á%BÁ(B"Á?BÂB"ÂBÂB"c              ƒ  óN  #   • U R                   b0  U R                   R                  nU R                  5       I Sh  v•N   U$ U R                  bG  U R                  R	                  5       I Sh  v•N nU R                  U R
                  5      I Sh  v•N   U$ [        S5      e Ne N6 N7f)ú2Stream the whole response, validate and return it.NrY  )rF  rV  rs  rB  r¼   rn   r\  )r`   rV  s     rY   r¼   ÚStreamedRunResult.get_outputë  s’   é € à×ÑÑ'Ø×%Ñ%×,Ñ,ˆFØ×(Ñ(Ó*×*Ð*ØˆMØ×"Ñ"Ñ.Ø×0Ñ0×;Ñ;Ó=×=ˆFØ×(Ñ(¨¯©Ó7×7Ð7ØˆMäÐHÓIÐIñ +ñ >Ù7ùs3   ‚7B%¹Bº0B%Á*B!Á+"B%ÂB#ÂB%Â!B%Â#B%c                ó¤   • U R                   b  U R                   R                  $ U R                  b  U R                  R                  $ [        S5      e)ú)Return the current state of the response.rY  )rF  rn   rB  r\  r_   s    rY   rn   ÚStreamedRunResult.responseø  sL   € ð ×ÑÑ'Ø×#Ñ#×,Ñ,Ð,Ø×"Ñ"Ñ.Ø×(Ñ(×1Ñ1Ð1äÐHÓIÐIrX   c                ó�   • U R                   b  U R                   R                  $ U R                  b  U R                  R                  $ g)r¥   N)rF  r¦   rB  r_   s    rY   r¦   ÚStreamedRunResult.metadata  sC   € ð ×ÑÑ'Ø×#Ñ#×,Ñ,Ð,Ø×"Ñ"Ñ.Ø×(Ñ(×1Ñ1Ð1àrX   c                ó¤   • U R                   b  U R                   R                  $ U R                  b  U R                  R                  $ [        S5      e)r¬   rY  )rF  r^   rB  r\  r_   s    rY   r^   ÚStreamedRunResult.usage  sL   € ð ×ÑÑ'Ø×#Ñ#×)Ñ)Ð)Ø×"Ñ"Ñ.Ø×(Ñ(×.Ñ.Ð.äÐHÓIÐIrX   c                ó¤   • U R                   b  U R                   R                  $ U R                  b  U R                  R                  $ [        S5      e)r¹   rY  )rF  r³   rB  r\  r_   s    rY   r³   ÚStreamedRunResult.timestamp  sL   € ð ×ÑÑ'Ø×#Ñ#×-Ñ-Ð-Ø×"Ñ"Ñ.Ø×(Ñ(×2Ñ2Ð2äÐHÓIÐIrX   c                ó¤   • U R                   b  U R                   R                  $ U R                  b  U R                  R                  $ [        S5      e)r›   rY  )rF  rœ   rB  r\  r_   s    rY   rœ   ÚStreamedRunResult.run_id$  sL   € ð ×ÑÑ'Ø×#Ñ#×*Ñ*Ð*Ø×"Ñ"Ñ.Ø×(Ñ(×/Ñ/Ð/äÐHÓIÐIrX   c                ó¤   • U R                   b  U R                   R                  $ U R                  b  U R                  R                  $ [        S5      e)r    rY  )rF  r¡   rB  r\  r_   s    rY   r¡   Ú!StreamedRunResult.conversation_id.  sL   € ð ×ÑÑ'Ø×#Ñ#×3Ñ3Ð3Ø×"Ñ"Ñ.Ø×(Ñ(×8Ñ8Ð8äÐHÓIÐIrX   rg   c             ƒ  óÆ   #   • U R                   b  U R                   R                  $ U R                  b!  U R                  R                  XS9I Sh  v•N $ [	        S5      e N7f)rÃ   Nrg   rY  )rF  rV  rB  rl   r\  ©r`   rÞ   rh   s      rY   rl   Ú*StreamedRunResult.validate_response_output8  s`   é € ð ×ÑÑ'Ø×#Ñ#×*Ñ*Ð*Ø×"Ñ"Ñ.Ø×.Ñ.×GÑGÈÐGÐm×mÐmäÐHÓIÐIñ nùs   ‚AA!ÁAÁA!c                óÈ   • U R                   (       a6  U R                   R                  Ul        U R                   R                  Ul        U R                  R	                  U5        g)zYAppend a model response to the message history with the correct run and conversation IDs.N)rB  rœ   r¡   r>  rú   )r`   rÞ   s     rY   Ú_record_responseÚ"StreamedRunResult._record_responseC  sF   € à× × Ø!×2Ñ2×9Ñ9ˆGŒNØ&*×&;Ñ&;×&KÑ&KˆGÔ#Ø×Ñ×!Ñ! 'Õ*rX   c              ƒ  óÚ   #   • SSK Jn  U R                  (       a  g SU l        U" 5       U l        Ub  U R	                  U5        U R
                  b  U R                  5       I S h  v•N   g g  N7f)Nr   )Úcurrent_otel_traceparentT)Ú_instrumentationr‘  rG  rR  rŽ  rD  )r`   rÞ   r‘  s      rY   rs  Ú#StreamedRunResult._marked_completedJ  sa   é € Ý>à××ØØˆÔÙ":Ó"<ˆÔØÑØ×!Ñ! 'Ô*Ø×ÑÑ(Ø×#Ñ#Ó%×%Ñ%ð )Ù%ùs   ‚A A+Á"A)Á#A+c              ƒ  óØ   #   • U R                   bW  U R                   R                  5       I Sh  v•N   U R                  (       d#  SU l        U R                  U R                  5        ggg N:7f)aà  Cancel local stream consumption and request provider shutdown.

Whether this stops remote generation or closes the underlying transport depends on the provider SDK.

The interrupted response state is recorded in the message history so that
`all_messages()` includes it.

This stops only the current model response; the run continues. To end the whole run,
use [`AgentRun.cancel()`][pydantic_ai.run.AgentRun.cancel] or
[`RunContext.cancel()`][pydantic_ai.tools.RunContext.cancel].
NT)rB  r�   rG  rŽ  rn   r_   s    rY   r�   ÚStreamedRunResult.cancelV  s]   é € ð × Ñ Ñ,Ø×'Ñ'×.Ñ.Ó0×0Ð0ð ×#×#Ø#'�Ô Ø×%Ñ% d§m¡mÕ4ð $ð -Ù0ùs   ‚+A*­A(®;A*c                óJ   • U R                   b  U R                   R                  $ g)r•   F)rB  r–   r_   s    rY   r–   ÚStreamedRunResult.cancelledk  s%   € ð × Ñ Ñ,Ø×(Ñ(×2Ñ2Ð2àrX   )r>  r@  rD  rF  rB  rR  rG  )
rI  r=  rJ  r?  ri   rA  rK  rC  rù   r&  )rI  r=  rJ  r?  rP  úAgentRunResult[OutputDataT]rù   r&  )NNN)rI  r=  rJ  r?  ri   rA  rK  rC  rP  rE  rù   r&  )rù   r˜  ©r`  ú
str | Nonerù   r=  ©r`  rš  rù   Úbytesr   r#  r"  r.  r*  r)  r,  r-  r(  r/  )rÞ   r+  rù   r&  rS   )rÞ   z_messages.ModelResponse | Nonerù   r&  r%  r'  )#r2  r3  r4  r5  Ú__doc__r6  rB  rD  rF  r   rG  r   rL  r9  r]  rI  rh  rm  rp  rp   rŠ   ri   r¼   rn   r¦   r^   r³   rœ   r¡   rl   rŽ  rs  r�   r–   r:  rW   rX   rY   r<  r<  í  sX  ‡ áPà/Ó/ØÓàDHÐÐAÓHØ9=€LÐ6Ó=à6:€KÐ3Ó:á e°%Ñ8€K�Ó8ðð ðà2ðð ðð Eð	ð
 :ðð 
óó ðð ðà2ðð ðð
 0ðð 
óó ðð HLØ<@Ø9=ðà2ðð ðð Eð	ð
 :ðð 7ðð 
õð* ó'Jó ð'JðR HL÷ "ð" MQ÷ 
ð  HL÷ sð  MQ÷ 
ð  BE÷ Jð2 27ÐTW÷ Jð< DG÷ Jô>Jð óJó ðJð óó ðð óJó ðJð óJó ðJð óJó ðJð óJó ðJð JOñ	JØ.ð	JØBFð	Jà	õ	Jô+ö
&ô5ð* óó órX   r<  c                  óz  • \ rS rSr% SrS\S'   S S jrS!S jr        S"S jrSS	.S#S
 jjr	SS	.S$S jjr
SS	.S#S jjrSS	.S$S jjrSS.S%S jjrSSS.S&S jjrSS.S'S jjrS(S jr\S)S j5       r\S*S j5       r\S+S j5       r\S,S j5       r\S,S j5       r\S-S j5       rSS.S.S jjr\S/S j5       rSrg)0r6   it  a8  Synchronous wrapper for [`StreamedRunResult`][pydantic_ai.result.StreamedRunResult] that only exposes sync methods.

All of the run's async work happens on the caller's event loop. Context-manager and iterator
lifecycles remain in stable tasks, so cancel scopes entered and exited by the agent graph never
straddle tasks and OpenTelemetry spans stay correctly nested. The wrapper must be used and closed
on the thread where it was created.

This is a synchronous context manager; the underlying stream is cleaned up on exit:

```python
from pydantic_ai import Agent

agent = Agent('openai:gpt-5.2')

def main():
    with agent.run_stream_sync('What is the capital of the UK?') as response:
        print(response.get_output())
        #> The capital of the UK is London.
```

Using it without a `with` block also works for backwards compatibility. Garbage collection requests
best-effort cleanup on the owner loop, but it cannot drive a stopped owner loop from another thread
or while another loop is running. A `with` block should be used whenever deterministic cleanup matters.
z*StreamedRunResult[AgentDepsT, OutputDataT]Ú_streamed_run_resultc                ó˜   • [        U[        5      (       a  [        S5      e[        USS9U l        U R                  R
                  U l        g )Nzª`StreamedRunResultSync` now takes the `run_stream()` context manager rather than an already-entered `StreamedRunResult`; use `agent.run_stream_sync(...)` to construct it.z`run_stream`)Úasync_alternative)rƒ   r<  Ú	TypeErrorr*   Ú_bridger  rŸ  )r`   Úrun_stream_cms     rY   rL  ÚStreamedRunResultSync.__init__�  sG   € Ü�mÔ%6×7Ñ7ô ðióð ô (¨ÈÑXˆŒØ$(§L¡L×$7Ñ$7ˆÕ!rX   c                ó   • U $ rS   rW   r_   s    rY   Ú	__enter__ÚStreamedRunResultSync.__enter__œ  s   € ØˆrX   c                ó>   • U R                   R                  XU45        g rS   )r£  Úshutdown)r`   Úexc_typeÚexc_valÚexc_tbs       rY   Ú__exit__ÚStreamedRunResultSync.__exit__Ÿ  s   € ð 	�‰×Ñ˜x°&Ð9Õ:rX   Nr_  c               ó4   • U R                   R                  US9$ )a�  Return the history of messages.

Args:
    output_tool_return_content: The return content of the tool call to set in the last message.
        This provides a convenient way to modify the content of the output tool call if you want to continue
        the conversation and want to set the response to the output tool call. If `None`, the last message will
        not be modified.

Returns:
    List of messages.
r_  )rŸ  rI  rc  s     rY   rI  Ú"StreamedRunResultSync.all_messages§  s   € ð ×(Ñ(×5Ñ5ÐQkÐ5ÐlÐlrX   c               ó4   • U R                   R                  US9$ )a  Return all messages from [`all_messages`][pydantic_ai.result.StreamedRunResultSync.all_messages] as JSON bytes.

Args:
    output_tool_return_content: The return content of the tool call to set in the last message.
        This provides a convenient way to modify the content of the output tool call if you want to continue
        the conversation and want to set the response to the output tool call. If `None`, the last message will
        not be modified.

Returns:
    JSON bytes representing the messages.
r_  )rŸ  rh  rc  s     rY   rh  Ú'StreamedRunResultSync.all_messages_jsonµ  ó   € ð ×(Ñ(×:Ñ:ÐVpÐ:ÐqÐqrX   c               ó4   • U R                   R                  US9$ )rl  r_  )rŸ  rm  rc  s     rY   rm  Ú"StreamedRunResultSync.new_messagesÃ  s   € ð ×(Ñ(×5Ñ5ÐQkÐ5ÐlÐlrX   c               ó4   • U R                   R                  US9$ )a  Return new messages from [`new_messages`][pydantic_ai.result.StreamedRunResultSync.new_messages] as JSON bytes.

Args:
    output_tool_return_content: The return content of the tool call to set in the last message.
        This provides a convenient way to modify the content of the output tool call if you want to continue
        the conversation and want to set the response to the output tool call. If `None`, the last message will
        not be modified.

Returns:
    JSON bytes representing the new messages.
r_  )rŸ  rp  rc  s     rY   rp  Ú'StreamedRunResultSync.new_messages_jsonÓ  r´  rX   rc   rd   c               ó^   ^^• U R                   mU R                  R                  UU4S j5      $ )a  Stream the output as an iterable.

The pydantic validator for structured data will be called in
[partial mode](https://docs.pydantic.dev/dev/concepts/experimental/#partial-validation)
on each iteration.

Args:
    debounce_by: by how much (if at all) to debounce/group the output chunks by. `None` means no debouncing.
        Debouncing is particularly important for long structured outputs to reduce the overhead of
        performing validation as each token is received.

Returns:
    An iterable of the response data.
c                 ó"   >• TR                  T S9$ ©Nrd   )rp   ©re   r]  s   €€rY   rZ   Ú5StreamedRunResultSync.stream_output.<locals>.<lambda>ñ  s   ø€ °×0DÑ0DÐQ\Ð0DÑ0]rX   ©rŸ  r£  Ústream_sync©r`   re   r]  s    `@rY   rp   Ú#StreamedRunResultSync.stream_outputá  s&   ù€ ð ×*Ñ*ˆØ�|‰|×'Ñ'Õ(]Ó^Ð^rX   Fr}   c               ób   ^^^• U R                   mU R                  R                  UUU4S j5      $ )uÌ  Stream the text result as an iterable.

!!! note
    [`TextOutput`][pydantic_ai.output.TextOutput] functions are not applied â€” use
    [`stream_output()`][pydantic_ai.result.StreamedRunResultSync.stream_output] instead.
    Result validators will NOT be called on the text result if `delta=True`.

Args:
    delta: if `True`, yield each chunk of text as it is received, if `False` (default), yield the full text
        up to the current point.
    debounce_by: by how much (if at all) to debounce/group the response chunks by. `None` means no debouncing.
        Debouncing is particularly important for long structured responses to reduce the overhead of
        performing validation as each token is received.
c                 ó$   >• TR                  TT S9$ )Nr}   )rŠ   )re   r~   r]  s   €€€rY   rZ   Ú3StreamedRunResultSync.stream_text.<locals>.<lambda>  s   ø€ °×0BÑ0BÈÐ\gÐ0BÑ0hrX   r¾  )r`   r~   re   r]  s    ``@rY   rŠ   Ú!StreamedRunResultSync.stream_textó  s&   ú€ ð ×*Ñ*ˆØ�|‰|×'Ñ'Ö(hÓiÐirX   c               ó^   ^^• U R                   mU R                  R                  UU4S j5      $ )a6  Stream the response as an iterable of `ModelResponse` snapshots.

Each yielded `ModelResponse` is the current state of the response: `response.state` is
`'incomplete'` while streaming is in flight and `'complete'` on the final yield.

Args:
    debounce_by: by how much (if at all) to debounce/group the response chunks by. `None` means no debouncing.
        Debouncing is particularly important for long structured responses to reduce the overhead of
        performing validation as each token is received.

Returns:
    An iterable of `ModelResponse` snapshots.
c                 ó"   >• TR                  T S9$ r»  )ri   r¼  s   €€rY   rZ   Ú7StreamedRunResultSync.stream_response.<locals>.<lambda>  s   ø€ °×0FÑ0FÐS^Ð0FÑ0_rX   r¾  rÀ  s    `@rY   ri   Ú%StreamedRunResultSync.stream_response  s&   ù€ ð ×*Ñ*ˆØ�|‰|×'Ñ'Õ(_Ó`Ð`rX   c                ó`   • U R                   R                  U R                  R                  5      $ )r{  )r£  ÚcallrŸ  r¼   r_   s    rY   r¼   Ú StreamedRunResultSync.get_output  s#   € à�|‰|× Ñ  ×!:Ñ!:×!EÑ!EÓFÐFrX   c                ó.   • U R                   R                  $ )r~  )rŸ  rn   r_   s    rY   rn   ÚStreamedRunResultSync.response  ó   € ð ×(Ñ(×1Ñ1Ð1rX   c                ó.   • U R                   R                  $ )r¬   )rŸ  r^   r_   s    rY   r^   ÚStreamedRunResultSync.usage  s   € ð ×(Ñ(×.Ñ.Ð.rX   c                ó.   • U R                   R                  $ r¸   )rŸ  r³   r_   s    rY   r³   ÚStreamedRunResultSync.timestamp(  r˜   rX   c                ó.   • U R                   R                  $ rš   )rŸ  rœ   r_   s    rY   rœ   ÚStreamedRunResultSync.run_id-  s   € ð ×(Ñ(×/Ñ/Ð/rX   c                ó.   • U R                   R                  $ rŸ   )rŸ  r¡   r_   s    rY   r¡   Ú%StreamedRunResultSync.conversation_id2  s   € ð ×(Ñ(×8Ñ8Ð8rX   c                ó.   • U R                   R                  $ r¤   )rŸ  r¦   r_   s    rY   r¦   ÚStreamedRunResultSync.metadata7  rÏ  rX   rg   c               óJ   ^ ^^• T R                   R                  UUU 4S j5      $ )rÃ   c                 ó8   >• TR                   R                  TT S9$ )Nrg   )rŸ  rl   )rh   rÞ   r`   s   €€€rY   rZ   Ú@StreamedRunResultSync.validate_response_output.<locals>.<lambda>?  s   ø€ �D×-Ñ-×FÑFÀwÐ^kÐFÑlrX   )r£  rË  r‹  s   ```rY   rl   Ú.StreamedRunResultSync.validate_response_output<  s   ú€ à�|‰|× Ñ Þló
ð 	
rX   c                ó.   • U R                   R                  $ )a  Whether the stream has all been received.

This is set to `True` when one of
[`stream_output`][pydantic_ai.result.StreamedRunResultSync.stream_output],
[`stream_text`][pydantic_ai.result.StreamedRunResultSync.stream_text],
[`stream_response`][pydantic_ai.result.StreamedRunResultSync.stream_response] or
[`get_output`][pydantic_ai.result.StreamedRunResultSync.get_output] completes.
)rŸ  rG  r_   s    rY   rG  Ú!StreamedRunResultSync.is_completeB  s   € ð ×(Ñ(×4Ñ4Ð4rX   )r£  rŸ  )r¤  zGAbstractAsyncContextManager[StreamedRunResult[AgentDepsT, OutputDataT]]rù   r&  )rù   r   )r«  ztype[BaseException] | Noner¬  zBaseException | Noner­  zTracebackType | Nonerù   r&  r™  r›  )re   r!  rù   zIterator[OutputDataT])r~   r$  re   r!  rù   zIterator[str])re   r!  rù   z!Iterator[_messages.ModelResponse]r.  r*  r,  r-  r(  r)  r/  r'  )r2  r3  r4  r5  r�  r6  rL  r§  r®  rI  rh  rm  rp  rp   rŠ   ri   r¼   r9  rn   r^   r³   rœ   r¡   r¦   rl   rG  r:  rW   rX   rY   r6   r6   t  s;  ‡ ñð2 EÓDô
8ôð;à,ð;ð &ð;ð %ð	;ð
 
ô;ð HL÷ mð MQ÷ rð HL÷ mð  MQ÷ rð <?÷ _ð$ ,1Èc÷ jð$ >A÷ aô"Gð ó2ó ð2ð ó/ó ð/ð ó3ó ð3ð ó0ó ð0ð ó9ó ð9ð ó2ó ð2ð ch÷ 
ð ó	5ó ó	5rX   r6   )rG   c                  ó^   • \ rS rSr% SrS\S'    SrS\S'    SrS\S'    \R                  r
S	rg)
ÚFinalResultiO  zNMarker class storing the final output of an agent run and associated metadata.r-   rV  Nrš  r¿   Útool_call_idrW   )r2  r3  r4  r5  r�  r6  r¿   râ  r   Údataclasses_no_defaults_reprÚ__repr__r:  rW   rX   rY   rá  rá  O  s3   ‡ áXàÓØ à €IˆzÓ Øbà#€L�*Ó#Øwà×2Ñ2ƒHrX   rá  c                ól   ^ ^^• Tb$  TR                  5       (       a  UUU 4S jnU" 5       $ [        T 5      $ )Nc                ó®   >#   • T  S h  v•N n TR                  T" 5       5        TR                  TR                  R                  5        U 7v •  MK   NF
 g 7frS   )Úcheck_tokensÚcheck_per_request_input_tokensr^   Úinput_tokens)ÚitemÚ	get_usageÚlimitsri   s    €€€rY   Ú_usage_checking_iteratorÚE_get_usage_checking_stream_response.<locals>._usage_checking_iteratorf  sF   øé € Ù-÷ �dØ×#Ñ#¡I£KÔ0Ø×5Ñ5°o×6KÑ6K×6XÑ6XÔYØ•
ñ™oùs&   ƒA†AŠA‹AŽAAÁAÁA)Úhas_token_limitsr  )ri   rì  rë  rí  s   ``` rY   r  r  _  s4   ú€ ð
 Ñ˜f×5Ñ5×7Ñ7÷	ñ (Ó)Ð)ä�_Ó%Ð%rX   c                ó  • / n/ nU  Hi  nUR                  UR                  5      nUc  M#  UR                  S:X  a  UR                  U5        MF  UR                  S:X  d  MX  UR                  U5        Mk     U(       d  U(       d  g[	        X2S9$ )zBGet the deferred tool requests from the model response tool calls.NÚ
unapprovedÚexternal)ÚcallsÚ	approvals)Úget_tool_defr¿   Úkindrú   r0   )rÒ   Útool_managerrô  ró  rß   Útool_defs         rY   rÔ   rÔ   q  s{   € ð /1€IØ*,€Eãˆ	Ø×,Ñ,¨Y×-@Ñ-@ÓAˆØÓØ�}‰} Ó,Ø× Ñ  Ö+Ø—‘ *Õ,Ø—‘˜YÖ'ñ  ö žØä eÑAÐArX   )ri   r;   rì  rA   rë  zCallable[[], RunUsage]rù   r1  )rÒ   z Iterable[_messages.ToolCallPart]r÷  rC   rù   zDeferredToolRequests | None)PÚ
__future__r   Ú_annotationsÚcollections.abcr   r   r   r   r   r	   Ú
contextlibr
   r   Úcopyr   Údataclassesr   r   r   r   Údecimalr   Útypesr   Útypingr   r   r   r   r   rU   Úpydanticr   Útyping_extensionsr   rË   r   r   r   rÚ   r   Ú_genai_pricesr   Ú_outputr    r!   r"   r#   r$   r%   r&   Ú_run_contextr'   r(   r)   Ú_sync_streamr*   r+   r,   rV  r-   r.   r÷  r/   Útoolsr0   r^   r1   r2   Úcapabilities.abstractr3   r[  r5   Ú__all__r9   r<  r6   rá  r  rÔ   rW   rX   rY   Ú<module>r     s]  ðÝ 2ç b× bß <Ý ß 1Ñ 1Ý Ý Ý ß >Õ >ã Ý $Ý "ç ?Ó ?Ý ,÷÷ ñ ÷ HÑ GÝ *ß @÷õ &Ý 'ß (æÝ9Ý#ð€ñ �4Ñôv�'˜* kÐ1Ñ2ó vó ðvñr �ÑôC˜ 
¨KÐ 7Ñ8ó Có ðCôLX5˜G J°Ð$;Ñ<ô X5ñv �Ñô3�'˜+Ñ&ó 3ó ð3ð&Ø,ð&àð&ð &ð&ð -ô	&ð$BØ0ðBØ@WðBà õBrX   