ó
    °"³jÂ1  ã                  ó  • % S r SSKJr  SSKrSSKrSSKrSSKrSSKJ	r	  SSK
Jr  SSKJr  SSKJrJr  \(       a  SS	KJr  S
r " S S5      r " S S5      r\R,                   " S S5      5       r\" SSS9rS\S'   \SS j5       rSS jrg)u  Run-scoped cancellation controller for first-party run cancellation.

First-party cancellation (`AgentRun.cancel()`, `RunContext.cancel()`) is implemented by
cancelling the asyncio task that drives the run: that wakes whatever the run is blocked on (a
model stream, tool tasks, a suspended-job poll) and reuses the exact same teardown machinery as
external cancellation â€” streams are closed, in-flight tool tasks are cancelled and drained,
suspended server-side jobs are best-effort cancelled, and completed work is recorded to message
history. At the outer edge of [`Agent.iter()`][pydantic_ai.agent.Agent.iter], after teardown, the
resulting `CancelledError` is translated back into
[`RunCancelled`][pydantic_ai.exceptions.RunCancelled] â€” but only if the cancellation was ours:

- The controller counts every `Task.cancel()` it issues. On catching `CancelledError`, the
  outer edge consumes exactly that many cancellations via `Task.uncancel()` (mirroring what
  `asyncio.timeout()` does for its own cancellation).
- If `Task.cancelling()` is still positive afterwards, an *external* cancellation raced in; it
  wins, and the `CancelledError` keeps propagating as itself.

On Python 3.10, `Task.cancelling()`/`Task.uncancel()` don't exist, so the race cannot be
disambiguated: a requested first-party cancellation is translated to `RunCancelled` even if an
external cancellation arrived at the same time (documented degraded behavior).

The controller is runtime-only state: it holds a live task reference and is never serialized.
é    )ÚannotationsN)Ú	Generator)Úcontextmanager)Ú
ContextVar)ÚTYPE_CHECKINGÚAnyé   )ÚAgentRun)ÚCancellationTokenÚ
RunBindingÚRunCancellationÚprovide_run_bindingÚtake_run_bindingc                  óT   • \ rS rSrSrS
S jr\SS j5       rS
S jrSS jr	SS jr
Srg	)r   é*   a  A thread-safe handle for cancelling one or more agent runs.

A token is permanently cancelled after [`cancel`][pydantic_ai.CancellationToken.cancel] is
called. The same token may be passed to multiple concurrent runs, in which case all of them
are cancelled.
c                ód   • SU l         [        5       U l        [        R                  " 5       U l        g ©NF)Ú
_cancelledÚsetÚ_registrationsÚ	threadingÚLockÚ_lock©Úselfs    ÚP/home/mande/repo/quber/.venv/lib/python3.13/site-packages/pydantic_ai/_cancel.pyÚ__init__ÚCancellationToken.__init__2   s!   € ØˆŒÜ47³EˆÔÜ—^’^Ó%ˆ�
ó    c                óh   • U R                      U R                  sSSS5        $ ! , (       d  f       g= f)z(Whether cancellation has been requested.N)r   r   r   s    r   Ú	cancelledÚCancellationToken.cancelled7   ó   € ð �Z‹ZØ—?‘?÷ �Z�Zúó   �#£
1c                óð   • U R                      U R                  (       a
   SSS5        gSU l        [        U R                  5      nSSS5        W H  nUR	                  5         M     g! , (       d  f       N(= f)zpCancel every live run registered with this token.

This method is idempotent and may be called from any thread.
NT)r   r   Útupler   Úcancel)r   ÚregistrationsÚcancellations      r   r'   ÚCancellationToken.cancel=   sZ   € ð
 �Z‹ZØ��Ø÷ ˆZð #ˆDŒOÜ! $×"5Ñ"5Ó6ˆM÷	 ó *ˆLØ×ÑÖ!ò *÷ �Zús   �A'©A'Á'
A5c                óâ   • U R                      U R                  (       a  SnOU R                  R                  U5        SnS S S 5        W(       a  UR	                  5         g g ! , (       d  f       N'= f)NTF)r   r   r   Úaddr'   )r   r)   Úshould_cancels      r   Ú	_registerÚCancellationToken._registerM   sQ   € Ø�Z‹ZØ��Ø $‘à×#Ñ#×'Ñ'¨Ô5Ø %�÷ ö Ø×ÑÕ!ð ÷ �Zús   �2A Á 
A.c                ó†   • U R                      U R                  R                  U5        S S S 5        g ! , (       d  f       g = f©N)r   r   Údiscard)r   r)   s     r   Ú_unregisterÚCancellationToken._unregisterW   s'   € Ø�Z‹ZØ×Ñ×'Ñ'¨Ô5÷ �ZŽZús	   �2²
A )r   r   r   N©ÚreturnÚNone©r6   Úbool)r)   r   r6   r7   )Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__r   Úpropertyr!   r'   r.   r3   Ú__static_attributes__© r   r   r   r   *   s/   † ñô&ð
 ó#ó ð#ô
"ô "÷6r   r   c                  óž   • \ rS rSrSrSS jr\SS j5       r\SS j5       rSSS jjr	SS jr
SS	 jrSS
 jrSS jrSS jrSS jrSS jrSrg)r   é\   aM  Tracks first-party cancellation of a single agent run.

One instance per run, shared (by reference) between the run's public handles and its
internals. The task driving the run binds itself with [`bind`][pydantic_ai._cancel.RunCancellation.bind]
at each step boundary, so `cancel()` always cancels the task currently doing the work.
c                óŒ   • S U l         S U l        0 U l        SU l        SU l        [
        R                  " 5       U l        / U l        g r   )	Ú_ownerÚ_loopÚ_issuedÚ
_requestedÚ	_finishedr   ÚRLockr   Ú_tokensr   s    r   r   ÚRunCancellation.__init__d   s:   € Ø37ˆŒØ7;ˆŒ
Ø8:ˆŒØˆŒØˆŒÜ—_’_Ó&ˆŒ
Ø02ˆ�r   c                óh   • U R                      U R                  sSSS5        $ ! , (       d  f       g= f)zVWhether a first-party cancellation has been requested. Sticky for the life of the run.N)r   rH   r   s    r   Úcancel_requestedÚ RunCancellation.cancel_requestedm   r#   r$   c                óz   • U R                      [        U R                  5      sSSS5        $ ! , (       d  f       g= f)zXWhether a [`CancellationToken`][pydantic_ai.CancellationToken] was attached to this run.N)r   r9   rK   r   s    r   Ú	has_tokenÚRunCancellation.has_tokens   s!   € ð �Z‹ZÜ˜Ÿ™Ó%÷ �Z�Zús   �,¬
:Nc                óR  • Uc   [         R                  " 5       nUc  gU R                     Xl        UR                  5       U l        [        R                  S:¼  ac  XR                  ;   aT  [        U R                  U   UR                  5       5      U R                  U'   U R                  U   S:X  a  U R                  U	 U R                  (       a1  U R                  (       d   XR                  ;  a  U R                  U5        SSS5        g! [         a     gf = f! , (       d  f       g= f)a.  Bind the task that is currently driving the run.

Called at run start and at each step boundary, so manual `AgentRun.next()` driving from
a different task than the one that started the run still gets cancelled correctly.
If a cancellation was requested before any task was bound (e.g. `cancel()` on a
lazily-started run) or was issued to a previous driving task, it is (re-)delivered to
this one. A caller that catches and uncancels the controller's own cancellation takes
over its bookkeeping; the issued count is re-synchronized at the next step boundary.
N©é   é   r   )ÚasyncioÚcurrent_taskÚRuntimeErrorr   rE   Úget_looprF   ÚsysÚversion_inforG   ÚminÚ
cancellingrH   rI   Ú_issue©r   Útasks     r   ÚbindÚRunCancellation.bindy   sá   € ð ‰<ðÜ×+Ò+Ó-�ð ‰<ØØ�Z‹ZØŒKØŸ™›ˆDŒJÜ×Ñ 7Ó*¨t·|±|Ó/Cô
 &)¨¯©°dÑ);¸T¿_¹_Ó=NÓ%O�—‘˜TÑ"Ø—<‘< Ñ%¨Ó*ØŸ™ TÐ*Ø�� t§~§~¸$ÇlÁlÓ:Rð —‘˜DÔ!÷ ˆZøô	  ó Ùðú÷ �Zús   …D ªCDÄ
DÄDÄ
D&c                óä  • U R                      U R                  (       d  U R                  (       a
   SSS5        gSU l        U R                  nU R                  nUb  Ub  UR                  5       (       a
   SSS5        g SSS5         [        R                  " 5       nUWL a  U R                  5         gUR                  U R                  5        g! , (       d  f       NV= f! [         a    Sn NQf = f)zaRequest cancellation of the run from any thread.

Idempotent; a no-op once the run has finished.
NT)r   rI   rH   rE   rF   ÚdonerW   Úget_running_looprY   Ú_deliverÚcall_soon_threadsafe)r   ÚownerÚloopÚrunning_loops       r   r'   ÚRunCancellation.cancelš   s¹   € ð
 �Z‹ZØ�~�~ §§Ø÷ ˆZð #ˆDŒOØ—K‘KˆEØ—:‘:ˆDØ‰} ¡°·
±
·±Ø÷ ˆZð 1=÷ ð	 Ü"×3Ò3Ó5ˆLð ˜4ÒØ�M‰M�Oà×%Ñ% d§m¡mÕ4÷! �Zûô ó 	 ØŠLð	 ús"   �$Cº;CÂC  Ã
CÃ C/Ã.C/c                ó  • U R                      U R                  nU R                  (       d'  Ub$  UR                  5       (       d  XR                  ;   a
   S S S 5        g U R                  U5        S S S 5        g ! , (       d  f       g = fr1   )r   rE   rI   re   rG   r_   )r   ri   s     r   rg   ÚRunCancellation._deliver±   sR   € Ø�Z‹ZØ—K‘KˆEØ�~�~ ¡°%·*±*·,±,À%Ï<É<ÓBWØ÷ ˆZð �K‰K˜Ô÷	 �ZŽZús   �AA6ÁA6Á6
Bc                ó|   • U R                   R                  US5      S-   U R                   U'   UR                  5         g )Nr   r	   )rG   Úgetr'   r`   s     r   r_   ÚRunCancellation._issue¸   s/   € Ø!Ÿ\™\×-Ñ-¨d°AÓ6¸Ñ:ˆ�‰�TÑØ�‰�r   c                ó¨   • U R                      U R                  R                  U5        SSS5        UR                  U 5        g! , (       d  f       N = f)zCRegister this run with a cancellation token until the run finishes.N)r   rK   Úappendr.   )r   Útokens     r   Úattach_tokenÚRunCancellation.attach_token¼   s4   € à�Z‹ZØ�L‰L×Ñ Ô&÷ à�‰˜Õ÷ �Zús   �AÁ
Ac                óð   • U R                      SU l        [        U R                  5      nU R                  R	                  5         SSS5        W H  nUR                  U 5        M     g! , (       d  f       N)= f)z?Mark the run as finished: later `cancel()` calls become no-ops.TN)r   rI   r&   rK   Úclearr3   )r   Útokensrt   s      r   ÚfinishÚRunCancellation.finishÂ   sU   € à�Z‹ZØ!ˆDŒNÜ˜4Ÿ<™<Ó(ˆFØ�L‰L×ÑÔ ÷ ó ˆEØ×Ñ˜dÖ#ò ÷	 �Zús   �7A'Á'
A5c                ó˜  • U R                   (       d  g[        R                  S:  a  g [        R                  " 5       nUc  gU R                  R                  US5      nUS:”  aE  UR                  5       S:”  a1  UR                  5         US-  nUS:”  a  UR                  5       S:”  a  M1  UR                  5       S:H  $ ! [
         a     gf = f)uð  Resolve a caught `CancelledError` at the run's outer edge: is it ours to translate?

Consumes only the cancellations this controller issued to the calling task via
`Task.uncancel()`. Returns `True` if the cancellation was first-party and no external
cancellation is still pending (translate to `RunCancelled`); `False` if it must keep
propagating as `CancelledError`.

Must be called on the task the cancellation was delivered to.

Unlike `asyncio.timeout()`, which arbitrates against a baseline count captured at scope
entry, this check is baseline-free: a cancellation count already pending when the run
started makes a first-party cancel resolve as external. That is deliberate â€” the
conservative direction is "external wins".

One residual window escapes that guarantee: because attribution counts cancellations
rather than tracking their identity, if user code catches a first-party cancellation and
calls `Task.uncancel()` itself, then an external `Task.cancel()` arrives before the next
`bind()` with a matching count, `bind()`'s clamp keeps the stale issuance and this check
consumes the external cancellation as first-party. Reaching it requires user code to
uncancel a cancellation it was handed; a robust fix needs issuance-identity tracking (#7240).
FrT   Tr   r	   )
rH   r[   r\   rW   rX   rY   rG   Úpopr^   Úuncancel)r   ra   Úcounts      r   ÚresolveÚRunCancellation.resolveË   sº   € ð, ��ØÜ×Ñ˜gÓ%ð ð	Ü×'Ò'Ó)ˆDð ‰<ØØ—‘× Ñ   qÓ)ˆØ�a‹i˜DŸO™OÓ-°Ó1Ø�M‰MŒOØ�Q‰JˆEð �a‹i˜DŸO™OÓ-°Õ1ð �‰Ó  AÑ%Ð%øô ó 	Ùð	ús   ©B< Â<
C	ÃC	c                óD  • [         R                  S:¼  ar  U R                  R                  5        HT  u  pUR	                  5       (       a  M  [        U5       H)  nUR                  5       S:”  d  M  UR                  5         M+     MV     U R                  R                  5         g)a-  Release controller-issued cancellations that were never resolved.

This includes cancellations swallowed by user code or issued to a superseded driving
task. Releasing them prevents contamination of the tasks' outer cancellation bookkeeping,
such as `asyncio.timeout()` and AnyIO cancellation scopes.
rT   r   N)	r[   r\   rG   Úitemsre   Úranger^   r~   rx   )r   ra   r   Ú_s       r   Úrelease_issuedÚRunCancellation.release_issuedô   sn   € ô ×Ñ˜wÓ&Ø#Ÿ|™|×1Ñ1Ö3‘�Ø—y‘y—{“{Ü" 5ž\˜ØŸ?™?Ó,¨qÕ0Ø ŸM™MžOó *ñ  4ð
 	�‰×ÑÕr   )rI   rG   r   rF   rE   rH   rK   r5   r8   r1   )ra   zasyncio.Task[object] | Noner6   r7   )ra   zasyncio.Task[object]r6   r7   )rt   r   r6   r7   )r:   r;   r<   r=   r>   r   r?   rN   rQ   rb   r'   rg   r_   ru   rz   r€   r†   r@   rA   r   r   r   r   \   s^   † ñô3ð ó#ó ð#ð
 ó&ó ð&ö
"ôB5ô.ôôô$ô'&÷Rr   r   c                  óR   • \ rS rSr% Sr\R                  " \S9rS\	S'   Sr
S\	S'   S	rg)
r   i  zêBridge an `AgentRunEvents` handle to the run it starts.

The handle exists before its lazy background run, so it owns the cancellation controller.
`Agent.iter()` later attaches the live run state while retaining that same controller.
)Údefault_factoryr   r)   NzAgentRun[Any, Any] | NoneÚ	agent_runrA   )r:   r;   r<   r=   r>   ÚdataclassesÚfieldr   r)   Ú__annotations__rŠ   r@   rA   r   r   r   r     s)   ‡ ñð %0×$5Ò$5ÀoÑ$V€L�/ÓVØ+/€IÐ(Ö/r   r   zpydantic_ai.run_binding)ÚdefaultzContextVar[RunBinding | None]Ú_current_run_bindingc              #  óž   #   • [         R                  U 5      n Sv •  [         R                  U5        g! [         R                  U5        f = f7f)zGSet the binding for runs started in this context, resetting it on exit.N)r�   r   Úreset)Úbindingrt   s     r   r   r     s<   é € ô !×$Ñ$ WÓ-€Eð*Ûä×"Ñ" 5Õ)øÔ×"Ñ" 5Õ)üs   ‚A™3 �A³A
Á
Ac                 ó^   • [         R                  5       n U b  [         R                  S5        U $ )z‡Consume and return the pending binding at most once.

Consuming prevents nested agent runs from inheriting the outer handle's binding.
N)r�   rp   r   )r’   s    r   r   r     s+   € ô
 #×&Ñ&Ó(€GØÑÜ× Ñ  Ô&Ø€Nr   )r’   r   r6   zGenerator[None])r6   zRunBinding | None)r>   Ú
__future__r   Ú_annotationsrW   r‹   r[   r   Úcollections.abcr   Ú
contextlibr   Úcontextvarsr   Útypingr   r   Úrunr
   Ú__all__r   r   Ú	dataclassr   r�   r�   r   r   rA   r   r   Ú<module>r�      sš   ðòõ0 3ã Û Û 
Û Ý %Ý %Ý "ß %æÝà
k€÷/6ñ /6÷deñ eðP ×Ñ÷0ð 0ó ð0ñ 7AÐAZÐdhÑ6iÐ Ð3Ó ið ó*ó ð*õr   