ó
    ±"³j.Œ  ã                  ó¨  • S r SSKJr  SSKrSSKrSSKJrJrJrJ	r	  SSK
JrJrJr  SSKJrJrJr  SSKJrJr  SSK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"J#r#J$r$J%r%  S	SK&J'r'  S	SK(J)r)  SSKJ*r*J+r+  / SQr,Sr-Sr.Sr/ Sr0 \S   r1 \Rd                  " S\Rf                  S9r4S S jr5S!S jr6S"S jr7S#S jr8S$S jr9\ " S S\+5      5       r:g)%uÝ  Model-continuation primitives shared by the agent graph.

A *continuation* happens when a model returns a `ModelResponse` with
`state == 'suspended'` (Anthropic `pause_turn`, OpenAI background mode, â€¦): the
graph re-issues the request with the suspended response echoed back, and the
provider resumes the same logical turn. This module owns the provider-agnostic
glue for stitching those segments back into a single response/stream:

- [`merge_responses`][pydantic_ai.models._continuation.merge_responses] folds a
  continuation response into the one it continues.
- [`merge_mode`][pydantic_ai.models._continuation.merge_mode] reports whether a
  continuation *replaces* or *accumulates*, so the streamed composite can reindex
  parts consistently with the merge.
- [`_ContinuationStreamedResponse`][pydantic_ai.models._continuation._ContinuationStreamedResponse]
  drives the streamed loop, presenting every segment as one continuous stream.

This module is deliberately decoupled from `_agent_graph`: it imports only from
`models`, `messages`, `usage`, `exceptions`, and the stdlib. Pluggable timing is
injected as `sleep_func` so the loop stays free of `now_utc()`/RNG and replays
deterministically under durable executors (e.g. Temporal).
é    )ÚannotationsN)ÚAsyncGeneratorÚAsyncIteratorÚ	AwaitableÚCallable)ÚAbstractContextManagerÚnullcontextÚsuppress)Ú	dataclassÚfieldÚreplace)ÚdatetimeÚtimezone)ÚAnyÚLiteralé   )Ú_utils)Ú
RunContext)ÚUnexpectedModelBehavior)ÚFinalResultEventÚModelMessageÚModelResponseÚModelResponseStateÚModelResponseStreamEventÚPartDeltaEventÚPartEndEventÚPartStartEvent)ÚModelSettings)ÚRequestUsageé   )ÚModelÚStreamedResponse)ÚMAX_BACKGROUND_POLLSÚMAX_GENERATION_CONTINUATIONSÚ	MergeModeÚcancel_suspended_jobÚ
merge_modeÚmerge_responsesÚ_ContinuationStreamedResponseÚ__pydantic_ai__Úreplace_previous_responseé
   iè  )úreplace-same-idúreplace-newÚ
accumulate)Útzc                óò   • U R                   n[        R                  " U5      (       d  gUR                  [        5      n[        [        R                  " U5      =(       a    UR                  [        5      5      $ )zUWhether `response` carries the `FallbackModel` `replace_previous_response` directive.F)Úmetadatar   Úis_str_dictÚgetÚ_PYDANTIC_AI_METADATA_KEYÚboolÚ_REPLACE_PREVIOUS_RESPONSE_KEY)Úresponser2   Ú	namespaces      Ú]/home/mande/repo/quber/.venv/lib/python3.13/site-packages/pydantic_ai/models/_continuation.pyÚ_has_replace_markerr;   x   sU   € à× Ñ €HÜ×Ò˜h×'Ñ'ØØ—‘Ô6Ó7€IÜ”×"Ò" 9Ó-×_°)·-±-Ô@^Ó2_Ó`Ð`ó    c                ó\  • [         R                  " U 5      (       d  U $ U R                  [        5      n[         R                  " U5      (       a
  [        U;   d  U $ UR                  5        VVs0 s H  u  p#U[        :w  d  M  X#_M     nnn0 U En U(       a
  X[        '   U $ U [        	 U $ s  snnf )a  Return `metadata` without the transient `replace_previous_response` marker (other keys intact).

Copies before mutating so the caller's dicts (including the shared `__pydantic_ai__` namespace,
which also holds the `FallbackModel` continuation pin) aren't touched.
)r   r3   r4   r5   r7   Úitems)r2   r9   ÚkÚvs       r:   Ú_strip_replace_markerrA   �   s¥   € ô ×Ò˜h×'Ñ'ØˆØ—‘Ô6Ó7€IÜ×Ò˜y×)Ñ)Ô.LÐPYÓ.YØˆØ"+§/¡/Ô"3Ô[Ò"3™$˜!°qÔ<ZÑ7Z“�’Ñ"3€IÑ[Ø�(ˆ|€HÞØ.7Ô*Ñ+ð €Oð Ô.Ð/Ø€Oùó \s   Á-B(ÂB(c                óø   • [        U5      (       a  gU R                  (       a  U R                  UR                  :X  a  gU R                  (       a,  UR                  (       a  U R                  UR                  :w  a  gg)u,  Classify how `new` folds into `existing` â€” see [`MergeMode`][pydantic_ai.models._continuation.MergeMode].

The single decision path shared by [`merge_responses`][pydantic_ai.models._continuation.merge_responses],
the continuation-count ceilings, and the streamed composite's part-index reindexing.
r.   r-   r/   )r;   Úprovider_response_idÚ
model_name)ÚexistingÚnews     r:   r'   r'   •   sW   € ô ˜3×ÑØØ×$×$¨×)FÑ)FÈ#×JbÑJbÓ)bØ Ø××˜sŸ~Ÿ~°(×2EÑ2EÈÏÉÓ2WØØr<   c                ót  • [        X5      nUS:X  a  UnOUS:X  a!  [        XR                  UR                  -   S9nOX[        U/ U R                  QUR                  QU R                  UR                  -   UR                  =(       d    U R                  S9nU R
                  (       a+  [        U0 U R
                  EUR
                  =(       d    0 ES9nU R                  (       a+  [        U0 U R                  EUR                  =(       d    0 ES9n[        UR                  5      nXCR                  La	  [        X4S9nU$ )ud  Merge a continuation response into the one it continues.

On any `'replace-*'` mode (same `provider_response_id`, a model change, or a `FallbackModel`
`replace_previous_response` directive), replace the content with the new response. Fresh-generation
replacements retain the prior response's billed usage; same-job polling replaces its cumulative usage
snapshot. Otherwise accumulate parts and usage, and use other fields from the new response.

Either way, `provider_details` and `metadata` accumulate across the turn's segments (latest-wins)
so turn-scoped data a later segment omits isn't lost â€” see below.
r-   r.   )Úusage)ÚpartsrH   rC   )Úprovider_details)r2   )r'   r   rH   rI   rC   rJ   r2   rA   )rE   rF   ÚmodeÚmergedÚstrippeds        r:   r(   r(   ¦   s  € ô �hÓ$€DØÐ Ó Ø‰Ø	�Ó	Ü˜§N¡N°S·Y±YÑ$>Ñ?‰ô
 ØØ/�H—N‘NÐ/ S§Y¡YÐ/Ø—.‘. 3§9¡9Ñ,Ø!$×!9Ñ!9×!Z¸X×=ZÑ=Zñ	
ˆð × × Ü˜Ð2r°X×5NÑ5NÐ2rÐSY×SjÑSj×SpÐnpÐ2rÑsˆØ××Ü˜Ð*Z¨X×->Ñ->Ð*ZÀ6Ç?Á?×CXÐVXÐ*ZÑ[ˆô % V§_¡_Ó5€HØ—‘Ò&Ü˜Ñ3ˆØ€Mr<   c              ƒ  óH  #   • [         R                  " U R                  U5      5      n [         R                  " U5      I Sh  v•N   g N! [         R                   a6    [        [        5         UI Sh  v•N    SSS5        e ! , (       d  f       e = f[         a     gf = f7f)u=  Best-effort teardown of a server-side suspended/background job that survives cancellation.

When the trigger is a workflow/task cancellation (e.g. Temporal), awaiting the (activity-wrapped)
cancel from inside an already-cancelled scope would raise `CancelledError` before the cancel runs,
silently leaking the job. Shield the cancel so it completes before the cancellation propagates;
Temporal's workflow loop respects `asyncio.shield`. Any error from the cancel itself is swallowed â€”
a failing teardown must not replace the error (or cancellation) that aborted the run.
N)ÚasyncioÚensure_futureÚcancel_suspended_responseÚshieldÚCancelledErrorr
   Ú	Exception)Úmodelr8   Újobs      r:   r&   r&   Ø   s{   é € ô ×
Ò
 × ?Ñ ?ÀÓ IÓ
J€CðÜ�nŠn˜SÓ!×!Ó!øÜ×!Ñ!ó ä”iÕ Ø�I‰I÷ !à÷ !Ô àúÜó Ùðüs^   ‚&B"©A
 ÁAÁA
 ÁB"ÁA
 Á
#BÁ-BÁ3A6Á4BÁ9	BÂ
B	ÂBÂB"ÂBÂB"c                  ó  ^ • \ rS rSr% SrS\S'   S\S'   S\S'   S	\S
'   S\S'   S\S'   S\S'   S\S'   SrS\S'   \rS\S'   \	r
S\S'   \" SSS9rS\S'   \" SSS9rS\S'   \" SSS9rS\S'   \" SSS9rS\S '   \" SSS9rS!\S"'   S4S# jrS5U 4S$ jjr    S6S% jr          S7S& jrS8S' jr\S9S( j5       rS:S) jrS;S* jr\S<S+ j5       rS=S, jrS>S- jrS>S. jr\S?S/ j5       r\S@S0 j5       r \S@S1 j5       r!\SAS2 j5       r"S3r#U =r$$ )Br)   éí   aC  A [`StreamedResponse`][pydantic_ai.models.StreamedResponse] that stitches continuation segments into one stream.

Each segment is an ordinary `model.request_stream(...)` sub-stream. Their events
are re-emitted as a single continuous stream, with part indices offset so parts
from accumulated segments (Anthropic `pause_turn`) don't collide, while replaced
segments (OpenAI background `retrieve`) keep reusing the same index space.

`get()` returns the live merged snapshot at any point; `usage` sums/replaces in
lockstep with the merge so the graph accounts for it exactly once.
r!   rU   zModelSettings | NoneÚmodel_settingszlist[ModelMessage]Úbase_messageszRunContext[Any] | NoneÚrun_contextÚintÚmax_generation_continuationsz"Callable[[float], Awaitable[None]]Ú
sleep_funczCallable[[RequestUsage], None]Úcheck_usagezCallable[[ModelResponse], None]Úfinalize_responseNúModelResponse | NoneÚinitial_suspended_responseÚmax_background_pollsz)Callable[[], AbstractContextManager[Any]]Úsegment_contextF)ÚdefaultÚinitÚ_merged_responsezStreamedResponse | NoneÚ_current_subr6   Ú_stoppedÚ	_detachedz5AsyncGenerator[ModelResponseStreamEvent, None] | NoneÚ_segment_iteratorc                óž   • U R                   c5  U R                  5       U l        U R                  U R                  5      U l         U R                   $ )aý  Stream every segment as one continuous event stream.

This intentionally bypasses the base `StreamedResponse.__aiter__`'s `iterator_with_final_event`
/ `iterator_with_part_end` wrappers: each sub-stream is already wrapped by them, so it emits
fully-formed `PartStart`/`PartDelta`/`PartEnd` and `FinalResultEvent`s. The composite only
applies reindexing + final-result capture (inside `_get_event_iterator`) and the cancel-guard
(reproducing the base `_finished`/`_cancelled` transitions) on top.

One minor semantic gap: `PartEndEvent.next_part_kind` is `None` at each sub-stream boundary
(a segment can't see the next segment's first part), whereas a single-segment stream would
populate it. This is acceptable because parts never merge across segment boundaries.
)Ú_event_iteratorÚ_get_event_iteratorrk   Ú_iterator_with_cancel_guard©Úselfs    r:   Ú	__aiter__Ú'_ContinuationStreamedResponse.__aiter__  sF   € ð ×ÑÑ'Ø%)×%=Ñ%=Ó%?ˆDÔ"Ø#'×#CÑ#CÀD×DZÑDZÓ#[ˆDÔ Ø×#Ñ#Ð#r<   c                ó|   >• [         TU ]  5       nU R                  b  / UQU R                  R                  5       Q7nU$ )uE  Cancel-teardown errors to suppress, extended with the in-flight sub-stream's own.

The cancel-guard tears the current segment down via its `close_stream()`, so the transport
error it raises is whatever that sub's transport produces â€” httpx for most providers, but
botocore (Bedrock) or grpc (xAI) for others, which each report their own types via an
override. Consult the in-flight sub (`_current_sub` is still set at exception time) on top of
the httpx default, so a non-httpx teardown error is suppressed into a clean `'interrupted'`
stop rather than escaping to the consumer.
)ÚsuperÚget_stream_cancel_errorsrh   )rq   ÚerrorsÚ	__class__s     €r:   rv   Ú6_ContinuationStreamedResponse.get_stream_cancel_errors,  sA   ø€ ô ‘Ñ1Ó3ˆØ×ÑÑ(ØM�vÐM × 1Ñ 1× JÑ JÓ LÑMˆFØˆr<   c               ó  #   •  U  S h  v•N nU R                   c  [        R                  " 5       U l         U7v •  M7   N2
 U R                  (       d  SU l        g g ! U R                  5        a    U R                  (       d  e  g f = f7f)NT)Ú_first_chunk_monotonicÚtimeÚperf_counterÚ
_cancelledÚ	_finishedrv   Ú	cancelled)rq   ÚiteratorÚevents      r:   ro   Ú9_ContinuationStreamedResponse._iterator_with_cancel_guard;  s|   é € ð	&Ù'÷ �eØ×.Ñ.Ñ6ô 37×2CÒ2CÓ2E�DÔ/Ø•ñ˜xð —?—?Ø!%�•ð #øð	 ×,Ñ,Ó.ó 	Ø—>—>Øñ "ð	üsB   ‚B„A †?Š=‹?Ž/A ½?¿A Á BÁ&BÂ BÂBÂBc                óø   • UR                   nUS:X  a5  US-  nX@R                  :”  a  [        SU< SU R                   S35      e X44$ US-  nX0R                  :”  a  [        SU< SU R                   S35      eX44$ )u  Count a suspended re-issue against its ceiling (same-id poll vs everything else), raising if exceeded.

Only a `'replace-same-id'` re-suspension â€” a passive re-poll of one long-running background job â€”
gets the generous `max_background_polls` ceiling. A model-change or `FallbackModel`-directed
replace is *fresh* generation, not the same job, so it counts against the strict `max_generation_continuations`
cap alongside accumulate re-suspensions â€” otherwise a chain of fresh suspensions is the exact
runaway the strict cap guards against. The first re-issue has no prior merge to classify
(`last_mode is None`) so it counts as strict, harmless since both ceilings allow at least one.
See `MAX_BACKGROUND_POLLS`. Returns the updated `(accumulate_count, replace_count)`.
r-   r    zModel response for job z1 remained suspended after polling the maximum of z timeszModel response z( was suspended more than the maximum of )rC   rc   r   r]   )rq   r8   Ú	last_modeÚaccumulate_countÚreplace_countÚjob_ids         r:   Ú_count_continuationÚ1_ContinuationStreamedResponse._count_continuationR  s·   € ð ×.Ñ.ˆØÐ)Ó)Ø˜QÑˆMØ×8Ñ8Ó8Ü-Ø-¨f©ZÐ7hØ×0Ñ0Ð1°ð9óð ð 9ð  Ð.Ð.ð  Ñ!ÐØ×"CÑ"CÓCÜ-Ø% f¡ZÐ/WØ×8Ñ8Ð9¸ðAóð ð  Ð.Ð.r<   c               ó  #   • SnSnS nU R                   nSn  U R                  (       d  U R                  (       a  GO_Uc  U R                  nO”UR                  S:X  aƒ  U R                  XCX5      u  pU R                  R                  U5      =n(       a<  U R                  U5      I S h  v•N   U R                  (       d  U R                  (       a  OÍ/ U R                  QUPnOO»X@l	        S nU R                  5          U R                  R                  X`R                  U R                  U R                  5       IS h  v•N n	X�l        U	  S h  v•N n
[!        U
["        5      (       a  X l        U
7v •  M+  Uc  U R'                  XIU5      nU R)                  X¨5      7v •  MV  X@l	         U R                  b*  U R-                  U R                  R+                  5       5        g g  GN4 N¤ N•
 S S S 5      IS h  v•N    O! , IS h  v•N  (       d  f       O= fS S S 5        O! , (       d  f       O= fU=(       d    SnW	R+                  5       nUc$  UR                  S:X  a  U R-                  U5        UnO8U R-                  U5        U R-                  U5        [/        XK5      n[1        XK5      nXÀl	        S U l        UR2                  U l        U R7                  UR2                  5        UnGMÈ  ! [8         a    e [:         aV    UbQ  UR                  S:X  aA  U R                  (       d0  U R                  (       d  [=        U R                  U5      I S h  v•N    e f = f! U R                  b*  U R-                  U R                  R+                  5       5        f f = f7f)Nr   Ú	suspended)rb   r~   ri   rZ   Ústater‰   rU   Úcontinuation_delayr^   rg   rd   Úrequest_streamrY   Úmodel_request_parametersr[   rh   Ú
isinstancer   Úfinal_result_eventÚ_segment_offsetÚ_reindexr4   r`   r'   r(   rH   Ú_usager_   ÚGeneratorExitÚBaseExceptionr&   )rq   r†   r‡   r…   r8   Úlast_segment_offsetÚmessagesÚdelayÚsegment_offsetÚsubr‚   Úsub_responserL   s                r:   rn   Ú1_ContinuationStreamedResponse._get_event_iteratorp  s"  é € ð ÐØˆð '+ˆ	Ø×2Ñ2ˆð  ÐðV	@ØØ—?—? d§m§mÙàÑ#Ø#×1Ñ1‘HØ—^‘^ {Ó2Ø6:×6NÑ6NØ Ð-=ó7Ñ3Ð$ð !%§
¡
× =Ñ =¸hÓ GÐG�uÕGØ"Ÿo™o¨eÓ4×4Ð4ð  Ÿ?Ÿ?¨d¯m¯mØ!Ø> ×!3Ñ!3Ð>°XÐ>‘Hàð
 )1Ô%ð .2�Ø×)Ñ)Õ+Ø#Ÿz™z×8Ñ8Ø ×"5Ñ"5°t×7TÑ7TÐVZ×VfÑVf÷ ô  àØ,/Ô)Ù+.÷ G %Ü)¨%Ô1A×BÑBØ:?Ô 7Ø&+£Ù (Ø-Ñ5Ø15×1EÑ1EÀhÐUhÓ1i Ø"&§-¡-°Ó"FÕFð6 %-Õ!ð" × Ñ Ñ,Ø×&Ñ& t×'8Ñ'8×'<Ñ'<Ó'>Õ?ð -òS 5ñ$ ñG¨3÷	 ÷  ÷  ÷  ÷  ð  ú÷ ,×+Ö+úð '5×&9¸Ð#ð  #Ÿw™w›y�ØÑ#Ø#×)Ñ)¨[Ó8Ø×.Ñ.¨|Ô<Ø)‘Fð ×*Ñ*¨8Ô4Ø×*Ñ*¨<Ô8ô !+¨8Ó B�IÜ,¨XÓD�Fà(.Ô%Ø$(�Ô!Ø$Ÿl™l�”Ø× Ñ  §¡Ô.Ø!�òC øôH ó 	ð Üó 	ð Ñ#¨¯©¸+Ó(EÈtÏÏÐbf×bo×boÜ*¨4¯:©:°xÓ@×@Ñ@Øð	ûð × Ñ Ñ,Ø×&Ñ& t×'8Ñ'8×'<Ñ'<Ó'>Õ?ð -üsÐ   ‚N
˜BK! Â'GÂ(AK! Ã9A HÄ9GÄ:HÄ=	G5ÅG"Å
G 
ÅG"ÅAG5ÆK! Æ#8N
ÇK! ÇHÇ G"Ç"G5Ç#HÇ.G1Ç/HÇ5H	Ç;G>Ç<H	ÈHÈ	K! È
H&È"B?K! Ë!A"MÍMÍMÍM Í9NÎN
c                ó„   • U c  g[        XR                  5       5      nUS:X  a  [        U R                  5      $ US:X  a  gU$ )u5  Index at which the current segment's parts begin in the stitched stream.

Shares [`merge_mode`][pydantic_ai.models._continuation.merge_mode]'s decision so reindexing
matches the eventual merge:

- `'accumulate'` appends after all prior parts (offset = number of prior parts).
- `'replace-same-id'` (a background job re-polled under the same `provider_response_id`) re-emits
  the *same* parts in the *same* index space, so it reuses the replaced segment's offset.
- `'replace-new'` (a model change, or a `FallbackModel` `replace_previous_response` directive)
  supersedes the whole prior response â€” `merge_responses` keeps only the new parts, indexed from
  0 â€” so its events must start at offset 0 too, or the live event indices would drift past the
  final response's (e.g. after one or more accumulated segments).
r   r/   r.   )r'   r4   ÚlenrI   )r8   rœ   r˜   rK   s       r:   r“   Ú-_ContinuationStreamedResponse._segment_offsetÜ  sE   € ð ÑØÜ˜(§G¡G£IÓ.ˆØ�<ÓÜ�x—~‘~Ó&Ð&Ø�=Ó ØØ"Ð"r<   c                ó€   • U(       a6  [        U[        [        [        45      (       a  [	        XR
                  U-   S9$ U$ )N)Úindex)r‘   r   r   r   r   r£   )rq   r‚   Úoffsets      r:   r”   Ú&_ContinuationStreamedResponse._reindexô  s1   € Þ”j ¬¼ÌÐ(V×WÑWÜ˜5¯©°fÑ(<Ñ=Ð=Øˆr<   c                ó|   • U R                   nU R                  =nb   UR                  5       nUc  U$ [        X5      $ U$ )z@The merged response so far, folding in any in-flight sub-stream.)rg   rh   r4   r(   )rq   rL   rœ   r�   s       r:   Ú	_snapshotÚ'_ContinuationStreamedResponse._snapshotù  sB   € à×&Ñ&ˆØ×$Ñ$Ð$ˆCÑ1ØŸ7™7›9ˆLØ#)¡>�<Ð\´ÀvÓ7\Ð\Øˆr<   c                óX   • U R                  5       nUb  UR                  $ U R                  $ )u˜  Live usage across all segments so far, including the in-flight sub-stream.

The composite's `_usage` is only refreshed when a segment completes, so â€” unlike a
plain segment, whose model updates `_usage` live during iteration â€” reading it mid
segment would omit the in-flight sub's usage. Fold in the current sub's live snapshot
so consumers (e.g. `AgentStream.usage`) see the running total at any point.
)r§   rH   r•   )rq   Úsnapshots     r:   rH   Ú#_ContinuationStreamedResponse.usage  s(   € ð —>‘>Ó#ˆØ!)Ñ!5ˆx�~‰~ÐF¸4¿;¹;ÐFr<   c                ó  • U R                  5       nU R                  (       a  SnO=U R                  (       a  SnO)U R                  (       a  Ub  UR                  S:X  a  SnOSnUc  [        / U R                  US9$ [        XS9$ )ud  Build the live merged [`ModelResponse`][pydantic_ai.messages.ModelResponse] across all segments so far.

The composite normally resolves the whole `suspended â†’ â€¦ â†’ complete` chain, so mid-run it's
`'complete'` once the loop exits, `'interrupted'` if cancelled, and `'incomplete'` while a
segment is still in flight. The one case it *does* surface `'suspended'` is a **detach**: the
consumer stopped iterating and the stream was torn down via `aclose()` â€” not `cancel()` â€” while
the current/last segment is itself a still-pending suspended job. The server-side job survives
(detach doesn't cancel it), so recording `'suspended'` makes the run resumable later, matching
the non-streaming path where a persisted suspended response can be resumed. A real `cancel()`
(`_cancelled`) also cancels the server-side job, so it stays `'interrupted'` and non-resumable.
ÚcompleteÚinterruptedrŒ   Ú
incomplete)rI   rD   r�   )r�   )r§   r   r~   rj   r�   r   rD   r   )rq   rª   r�   s      r:   r4   Ú!_ContinuationStreamedResponse.get  so   € ð —>‘>Ó#ˆð �>�>Ø‰EØ�_�_à!‰EØ�^�^ Ñ 4¸¿¹È;Ó9Và‰Eà ˆEàÑÜ  r°d·o±oÈUÑSÐSÜ�xÑ-Ð-r<   c              ƒ  ó>  #   • SU l          U R                  b"  U R                  R                  5       I Sh  v•N   [        U R                  U R                  5       5      I Sh  v•N   g N1 N! [        U R                  U R                  5       5      I Sh  v•N    f = f7f)zOStop the continuation loop and cancel any server-side suspended/background job.TN)ri   rh   Úclose_streamr&   rU   r4   rp   s    r:   r²   Ú*_ContinuationStreamedResponse.close_stream+  sv   é € àˆŒð	?Ø× Ñ Ñ,Ø×'Ñ'×4Ñ4Ó6×6Ð6ô
 ' t§z¡z°4·8±8³:Ó>×>Ñ>ñ 7ñ
 ?øÔ& t§z¡z°4·8±8³:Ó>×>Ò>üsI   ‚B‹*A+ µA'¶A+ º'BÁ!A)Á"BÁ'A+ Á)BÁ+(BÂBÂBÂBc              ƒ  óæ   #   • SU l         U R                  b$   U R                  R                  5       I Sh  v•N   gg N! [         a&  n[        R
                  " U5      (       d  e  SnAgSnAff = f7f)uö  Tear down the in-flight sub-stream (its HTTP connection) without cancelling any job.

Each segment's `model.request_stream(...)` context manager lives inside the stitching
async generator, so â€” unlike an ordinary `StreamedResponse` whose connection the agent
graph closed via `async with request_stream(...)` â€” an early break from a consumer
doesn't reliably propagate `aclose()` to it. Closing the generator here runs that
context manager's teardown (mirroring the pre-stitching behavior), and is safe to call
once the consumer has stopped iterating (including after a normal, fully-drained stream,
where it's a no-op). Cancellation of a server-side job stays on the `close_stream()`
path, driven by `AgentStream.cancel()`.

Closes the inner segment generator directly: the outer cancel-guard's
`async for â€¦ in iterator` does not forward `aclose()` to it, so closing the cancel-guard
alone would leave each segment's `request_stream(...)` context manager (and its
connection) open until garbage collection.
TN)rj   rk   ÚacloseÚRuntimeErrorr   Ú"is_async_generator_already_running)rq   Úexcs     r:   rµ   Ú$_ContinuationStreamedResponse.aclose7  sj   é € ð* ˆŒØ×!Ñ!Ñ-ðØ×,Ñ,×3Ñ3Ó5×5Ñ5ð .á5øÜó 
ô ×@Ò@À×EÑEØô Fûð
üs7   ‚A1˜> µ<¶> ºA1¼> ¾
A.ÁA)Á$A1Á)A.Á.A1c                óÆ   • U R                   b  U R                   R                  $ U R                  b1  U R                  R                  (       a  U R                  R                  $ g)NÚ )rh   rD   rg   rp   s    r:   rD   Ú(_ContinuationStreamedResponse.model_name\  sO   € à×ÑÑ(Ø×$Ñ$×/Ñ/Ð/Ø× Ñ Ñ,°×1FÑ1F×1Q×1QØ×(Ñ(×3Ñ3Ð3Ør<   c                ó’   • U R                   b  U R                   R                  $ U R                  b  U R                  R                  $ S $ ©N)rh   Úprovider_namerg   rp   s    r:   r¿   Ú+_ContinuationStreamedResponse.provider_named  sC   € à×ÑÑ(Ø×$Ñ$×2Ñ2Ð2Ø6:×6KÑ6KÑ6Wˆt×$Ñ$×2Ñ2ÐaÐ]aÐar<   c                ó’   • U R                   b  U R                   R                  $ U R                  b  U R                  R                  $ S $ r¾   )rh   Úprovider_urlrg   rp   s    r:   rÂ   Ú*_ContinuationStreamedResponse.provider_urlj  sC   € à×ÑÑ(Ø×$Ñ$×1Ñ1Ð1Ø59×5JÑ5JÑ5Vˆt×$Ñ$×1Ñ1Ð`Ð\`Ð`r<   c                óš   • U R                   b  U R                   R                  $ U R                  b  U R                  R                  $ [        $ r¾   )rh   Ú	timestamprg   Ú_FALLBACK_TIMESTAMPrp   s    r:   rÅ   Ú'_ContinuationStreamedResponse.timestampp  sC   € à×ÑÑ(Ø×$Ñ$×.Ñ.Ð.Ø26×2GÑ2GÑ2Sˆt×$Ñ$×.Ñ.ÐlÔYlÐlr<   )
rh   rj   rm   r   r{   rg   rk   ri   r•   r’   )Úreturnú'AsyncIterator[ModelResponseStreamEvent])rÈ   ztuple[type[BaseException], ...])r�   rÉ   rÈ   rÉ   )
r8   r   r…   zMergeMode | Noner†   r\   r‡   r\   rÈ   ztuple[int, int])rÈ   z.AsyncGenerator[ModelResponseStreamEvent, None])r8   ra   rœ   r"   r˜   r\   rÈ   r\   )r‚   r   r¤   r\   rÈ   r   )rÈ   ra   )rÈ   r   )rÈ   r   )rÈ   ÚNone)rÈ   Ústr)rÈ   z
str | None)rÈ   r   )%Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__Ú__annotations__rb   r#   rc   r	   rd   r   rg   rh   ri   rj   rk   rr   rv   ro   r‰   rn   Ústaticmethodr“   r”   r§   ÚpropertyrH   r4   r²   rµ   rD   r¿   rÂ   rÅ   Ú__static_attributes__Ú__classcell__)rx   s   @r:   r)   r)   í   s±  ø‡ ñ	ð ƒLØ(Ó(Ø%Ó%Ø'Ó'Ø"%Ó%Ø2Ó2Ø/Ó/Ø6Ó6Ø7;ÐÐ 4Ó;ð !5Ð˜#Ó4ð
 BM€OÐ>ÓLá-2¸4ÀeÑ-LÐÐ*ÓLÙ,1¸$ÀUÑ,K€LÐ)ÓKÙ 5¨uÑ5€HˆdÓ5ñ  E°Ñ6€IˆtÓ6ñ PUÐ]aÐhmÑOnÐÐLÓnô$÷$ð&Ø?ð&à	0ô&ð./Ø%ð/Ø2Bð/ØVYð/Øjmð/à	ô/ô<j@ðX ó#ó ð#ô.ô
ð ó	Gó ð	Gô.ô<
?ô#ðJ óó ðð óbó ðbð
 óaó ðað
 ómó ömr<   r)   )r8   r   rÈ   r6   )r2   r   rÈ   zdict[str, Any] | None)rE   r   rF   r   rÈ   r%   )rE   r   rF   r   rÈ   r   )rU   r!   r8   r   rÈ   rÊ   );rÐ   Ú
__future__r   rO   r|   Úcollections.abcr   r   r   r   Ú
contextlibr   r	   r
   Údataclassesr   r   r   r   r   Útypingr   r   r»   r   Ú_run_contextr   Ú
exceptionsr   r™   r   r   r   r   r   r   r   r   Úsettingsr   rH   r   r!   r"   Ú__all__r5   r7   r$   r#   r%   ÚfromtimestampÚutcrÆ   r;   rA   r'   r(   r&   r)   © r<   r:   Ú<module>râ      så   ðñõ, #ã Û ß NÓ Nß DÑ Dß 1Ñ 1ß 'ß å Ý %Ý 0÷	÷ 	ó 	õ %Ý  ß %ò€ð" .Ð Ø!<Ð à!Ð ðð Ð ðð ÐBÑC€	ðð" ×,Ò,¨Q°8·<±<Ñ@Ð ôaôô(ô"/ôdð* ôFmÐ$4ó Fmó ñFmr<   