ó
    ±"³jµ  ã                  ó  • S r SSKJr  SSKJrJr  SSKJr  SSKJ	r	J
r
  SSK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  \(       a  SSKJr  SSKJr  SSKJr  SS jr         SS jr! " S S\\   5      r"      SS jr#g)zTAuto-injected capability that drains the pending message queue at appropriate times.é    )Úannotations)ÚTYPE_CHECKINGÚAny)ÚModelRequestNode)ÚPendingMessageÚPendingMessageQueue)Úfill_run_metadata)ÚAbstractCapabilityÚCapabilityOrdering)Ú	UserError)ÚEnqueuedMessagesEventÚModelMessageÚModelRequest)Ú
RunContext)ÚEnd)Ú_agent_graph)ÚModelRequestContext)ÚFinalResultc                óL   • U R                   n[        U[        5      (       d   eU$ ©N)Úpending_messagesÚ
isinstancer   )ÚctxÚqueues     Úg/home/mande/repo/quber/.venv/lib/python3.13/site-packages/pydantic_ai/capabilities/_pending_messages.pyÚ_queuer      s(   € ð × Ñ €EÜ�eÔ0×1Ñ1Ð1Ð1Ø€Ló    c               óf   • / nU R                    H  n[        XAUS9  UR                  U5        M      U$ )aÒ  Stamp a pending message's messages' `timestamp` / `run_id` / `conversation_id` where unset.

Each [`PendingMessage`][pydantic_ai._enqueue.PendingMessage] carries one or more built
[`ModelMessage`][pydantic_ai.messages.ModelMessage]s (assembled at enqueue time by
[`PendingMessage.from_content`][pydantic_ai._enqueue.PendingMessage.from_content]); this only
fills in framework-tracked metadata that the producer left unset, so producer-supplied values
are preserved.
)Úrun_idÚconversation_id)Úmessagesr	   Úappend)ÚpendingÚfallback_run_idÚfallback_conversation_idr!   Úmessages        r   Ú_stamped_messagesr'      s7   € ð $&€HØ×#Ô#ˆÜ˜'ÐKcÒdØ�‰˜Ö ñ $ð €Or   c                  óL   • \ rS rSrSrSS jr\S	S j5       r      S
S jrSr	g)ÚPendingMessageDrainCapabilityé3   uM  Drains the pending message queue at appropriate times.

- `'asap'` messages drain at the earliest opportunity: into the next
  [`ModelRequest`][pydantic_ai.messages.ModelRequest] via `before_model_request`,
  or â€” if the agent would otherwise terminate â€” redirected through a new
  `ModelRequestNode` at the end of the run.
- `'when_idle'` messages drain only when the agent would otherwise terminate
  and no `'asap'` messages remain, after any `'asap'` redirect.

This capability is always auto-injected and placed outermost via
[`CapabilityOrdering`][pydantic_ai.capabilities.abstract.CapabilityOrdering].
Final draining happens at the graph-advancement seam after the full capability
hook chain, so user ordering constraints cannot make the queue close before a
hook has had a chance to redirect an apparent [`End`][pydantic_graph.End].
c                ó   • [        SS9$ )NÚ	outermost)Úposition)r   )Úselfs    r   Úget_orderingÚ*PendingMessageDrainCapability.get_orderingD   s   € Ü!¨;Ñ7Ð7r   c                ó   • g r   © )Úclss    r   Úget_serialization_nameÚ4PendingMessageDrainCapability.get_serialization_nameG   s   € àr   c           	   ƒ  óT  #   • [        U5      R                  S5      nU Hƒ  n[        XAR                  UR                  S9nUR
                  R                  U5        UR
                  R                  U5        UR                  [        UR                  [        U5      S95        M…     U$ 7f)u¶  Drain `'asap'` messages into the upcoming model request.

Each drained request is appended to both `request_context.messages` (so the model
sees it this step) and `ctx.messages` (so it persists in the agent's message
history). Stamps `timestamp`/`run_id`/`conversation_id` if the producer didn't â€”
`ModelRequestNode.run()` only stamps `self.request` (the current node's request),
and capabilities downstream of us might append more messages, so we can't rely on
that fixup.

Emits one [`EnqueuedMessagesEvent`][pydantic_ai.messages.EnqueuedMessagesEvent] per drained
[`enqueue`][pydantic_ai.tools.RunContext.enqueue] call, in enqueue order, describing the
messages exactly as delivered here.
Úasap©r$   r%   ©Ú
enqueue_idr!   )r   Úpop_priorityr'   r   r    r!   ÚextendÚ_emit_eventr   r:   Útuple)r.   r   Úrequest_contextÚdrainedr#   r!   s         r   Úbefore_model_requestÚ2PendingMessageDrainCapability.before_model_requestK   s�   é € ô$ ˜“+×*Ñ*¨6Ó2ˆÛˆGÜ(Ø¯©Èc×NaÑNañˆHð ×$Ñ$×+Ñ+¨HÔ5Ø�L‰L×Ñ Ô)Ø�O‰OÔ1¸W×=OÑ=OÔZ_Ð`hÓZiÑjÖkñ ð Ðùs   ‚B&B(r2   N)Úreturnr   )rC   ú
str | None)r   úRunContext[Any]r?   r   rC   r   )
Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__r/   Úclassmethodr4   rA   Ú__static_attributes__r2   r   r   r)   r)   3   sA   † ñô 8ð óó ððàðð -ðð 
÷	r   r)   c           
     ó˜  • [        U[        5      (       d  U$ [        U 5      R                  5       u  p#U(       d	  U(       d  U$ / UQUQnU Vs/ s H#  nU[	        XPR
                  U R                  S94PM%     nnU VVV	s/ s H  u  pxU  H  o™PM     M     n
nnn	U
S   n[        U[        5      (       d"  [        S[        U5      R                   S35      eU H[  u  pXU H$  n	X›Ld  M	  U R                  R                  U	5        M&     U R                  [        UR                  [!        U5      S95        M]     [#        US9$ s  snf s  sn	nnf )a‚  Drain pending messages after all capability hooks if the agent would terminate.

Drain `'asap'` messages first (anything that arrived after the most recent
`before_model_request` and would otherwise be lost), then `'when_idle'` messages.
Each priority is appended independently so the history keeps the priority split
visible (matches pi-mono's separate steering / follow-up turns). On the wire,
`_clean_message_history` re-merges adjacent requests with compatible instructions,
so the model still sees one turn.

The last resulting request becomes the redirect
[`ModelRequestNode`][pydantic_ai._agent_graph.ModelRequestNode]'s request; any
earlier ones are appended to `ctx.messages` so they appear in history before the
redirect. Emits one [`EnqueuedMessagesEvent`][pydantic_ai.messages.EnqueuedMessagesEvent]
per drained [`enqueue`][pydantic_ai.tools.RunContext.enqueue] call, in enqueue order.
r8   éÿÿÿÿz|Enqueued content must end with a `ModelRequest` so the agent has a request to respond to, but the last queued message is a `z`.r9   )Úrequest)r   r   r   Údrain_at_endr'   r   r    r   r   ÚtyperF   r!   r"   r=   r   r:   r>   r   )r   ÚresultÚleftover_asapÚ	when_idler@   r#   ÚstampedÚ_r   r&   r!   Úfinals               r   Údrain_pending_messages_at_endrX   h   sX  € ô& �fœc×"Ñ"Øˆô  & c›{×7Ñ7Ó9Ñ€MÞ¦Øˆà*�Ð* 	Ð*€Gñ óò
 ˆGð Ü˜g·z±zÐ\_×\oÑ\oÑpó	
ñ ð ð ñ 4;Õ[²7Ñ/˜AÔJZ¸w’ÑJZ‘±7€HÒ[ð �R‰L€EÜ�eœ\×*Ñ*Üð1Ü15°e³×1EÑ1EÐ0FÀbðJó
ð 	
ó &-Ñ!ˆÛ'ˆGØÔ#Ø—‘×#Ñ# GÖ,ñ (ð 	�‰Ü!¨W×-?Ñ-?Ì%ÐP`ÓJaÑbö	
ñ	 &-ô  EÑ*Ð*ùò9ùô \s   Á*E Á?EN)r   rE   rC   r   )r#   r   r$   rD   r%   rD   rC   zlist[ModelMessage])r   rE   rR   ú8_agent_graph.AgentNode[Any, Any] | End[FinalResult[Any]]rC   rY   )$rJ   Ú
__future__r   Útypingr   r   Úpydantic_ai._agent_graphr   Úpydantic_ai._enqueuer   r   Úpydantic_ai._utilsr	   Ú!pydantic_ai.capabilities.abstractr
   r   Úpydantic_ai.exceptionsr   Úpydantic_ai.messagesr   r   r   Úpydantic_ai.toolsr   Úpydantic_graphr   Úpydantic_air   Úpydantic_ai.modelsr   Úpydantic_ai.resultr   r   r'   r)   rX   r2   r   r   Ú<module>rg      s—   ðÙ Zå "ç %å 5ß DÝ 0ß TÝ ,ß RÑ RÝ (Ý æÝ(Ý6Ý.ôðØðð  ðð )ð	ð
 ôô*2Ð$6°sÑ$;ô 2ðj;+Ø	ð;+àDð;+ð >õ;+r   