ó
    °"³j�  ã            
      óÂ  • S r SSKJr  SSK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JrJrJrJ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S9r\" SSS9r \" SSS9r!\" SSS9r"\" SSS9r#\" SSS9r$\" SSS9r%\ " S S5      5       r&\" SS9 " S S\\\ 4   5      5       r'\" S\\"\!/\"4   \!\"4S9r(\" S\\'\\ 4   \"\!/\"4   \\ \!\"4S9r)\" S\)\\ \!\"4   \(\!\"4   -  \\ \!\"4S9r* S-S jr+S.S jr,S/S  jr-S0S! jr. " S" S#\5      r/\" S$\/SS%9r0S1S& jr1\ " S' S(\\#   5      5       r2\" SS9 " S) S*\\\ \!\"4   5      5       r3\ " S+ S,\\\ \4   5      5       r4g)2zíJoin operations and reducers for graph execution.

This module provides the core components for joining parallel execution paths
in a graph, including various reducer types that aggregate data from multiple
sources into a single output.
é    )ÚannotationsN)Úabstractmethod)ÚCallableÚIterableÚMapping)Ú	dataclass)ÚAnyÚGenericÚLiteralÚcastÚoverload)ÚProtocolÚSelfÚTypeAliasTypeÚTypeVar)ÚBaseNodeÚEndÚGraphRunContext)ÚForkIDÚ	ForkStackÚJoinIDÚStateTT)Úinfer_varianceÚDepsTÚInputTÚOutputTÚTÚKÚVc                  ó<   • \ rS rSr% SrS\S'   S\S'   SrS\S	'   S
rg)Ú	JoinStateé   zOThe state of a join during graph execution associated to a particular fork run.r	   Úcurrentr   Údownstream_fork_stackFÚboolÚcancelled_sibling_tasks© N)Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__Ú__annotations__r&   Ú__static_attributes__r'   ó    ÚP/home/mande/repo/quber/.venv/lib/python3.13/site-packages/pydantic_graph/join.pyr!   r!      s   ‡ áYàƒLØ$Ó$Ø$)Ð˜TÖ)r/   r!   F)Úinitc                  óv   • \ rS rSr% SrS\S'    S\S'    S\S'    SS	 jr\SS
 j5       r\SS j5       r	S r
Srg)ÚReducerContexté(   züContext information passed to reducer functions during graph execution.

The reducer context provides access to the current graph state and dependencies.

Type Parameters:
    StateT: The type of the graph state
    DepsT: The type of the dependencies
r   Ú_stater   Ú_depsr!   Ú_join_statec               ó(   • Xl         X l        X0l        g ©N)r5   r6   r7   )ÚselfÚstateÚdepsÚ
join_states       r0   Ú__init__ÚReducerContext.__init__:   s   € ØŒØŒ
Ø%Õr/   c                ó   • U R                   $ )zThe state of the graph run.)r5   ©r:   s    r0   r;   ÚReducerContext.state?   s   € ð �{‰{Ðr/   c                ó   • U R                   $ )zThe deps for the graph run.)r6   rA   s    r0   r<   ÚReducerContext.depsD   s   € ð �z‰zÐr/   c                ó&   • SU R                   l        g)zCancel all sibling tasks created from the same fork.

You can call this if you want your join to have early-stopping behavior.
TN)r7   r&   rA   s    r0   Úcancel_sibling_tasksÚ#ReducerContext.cancel_sibling_tasksI   s   € ð
 48ˆ×ÑÕ0r/   )r6   r7   r5   N)r;   r   r<   r   r=   r!   )Úreturnr   )rH   r   )r(   r)   r*   r+   r,   r-   r>   Úpropertyr;   r<   rF   r.   r'   r/   r0   r3   r3   (   sT   ‡ ñð ƒNØ"ØƒLØ4ØÓØ1ô&ð
 óó ðð óó ðõ8r/   r3   ÚPlainReducerFunction)Útype_paramsÚContextReducerFunctionÚReducerFunctionc                ó   • g)z8A reducer that discards all input data and returns None.Nr'   ©r#   Úinputss     r0   Úreduce_nullrQ   e   s   € àr/   c                ó(   • U R                  U5        U $ )z!A reducer that appends to a list.)ÚappendrO   s     r0   Úreduce_list_appendrT   j   ó   € à‡N�N�6ÔØ€Nr/   c                ó(   • U R                  U5        U $ )zA reducer that extends a list.)ÚextendrO   s     r0   Úreduce_list_extendrX   p   rU   r/   c                ó(   • U R                  U5        U $ )zA reducer that updates a dict.)ÚupdaterO   s     r0   Úreduce_dict_updater[   v   rU   r/   c                  ó,   • \ rS rSrSr\SS j5       rSrg)ÚSupportsSumé|   z5A protocol for a type that supports adding to itself.c               ó   • g r9   r'   )r:   Úothers     r0   Ú__add__ÚSupportsSum.__add__   s   € àr/   r'   N)r`   r   rH   r   )r(   r)   r*   r+   r,   r   ra   r.   r'   r/   r0   r]   r]   |   s   † Ù?àóó ór/   r]   ÚNumericT)Úboundr   c                ó
   • X-   $ )zA reducer that sums numbers.r'   rO   s     r0   Ú
reduce_sumrf   ‡   s   € àÑÐr/   c                  ó"   • \ rS rSrSrSS jrSrg)ÚReduceFirstValueéŒ   zRA reducer that returns the first value it encounters, and cancels all other tasks.c                ó&   • UR                  5         U$ )zThe reducer function.)rF   )r:   Úctxr#   rP   s       r0   Ú__call__ÚReduceFirstValue.__call__�   s   € à× Ñ Ô"Øˆr/   r'   N)rk   zReducerContext[object, object]r#   r   rP   r   rH   r   )r(   r)   r*   r+   r,   rl   r.   r'   r/   r0   rh   rh   Œ   s
   † á\÷r/   rh   c                  óÖ   • \ rS rSr% SrS\S'   S\S'   S\S'   S	\S
'   S\S'   SSS.         SS jjr\S 5       r\S 5       r	SS jr
\SSS jj5       r\SS j5       rSSS jjrSrg)ÚJoiné–   að  A join operation that synchronizes and aggregates parallel execution paths.

A join defines how to combine outputs from multiple parallel execution paths
using a [`ReducerFunction`][pydantic_graph.join.ReducerFunction]. It specifies which fork
it joins (if any) and manages the initialization of reducers.

Type Parameters:
    StateT: The type of the graph state
    DepsT: The type of the dependencies
    InputT: The type of input data to join
    OutputT: The type of the final joined output
r   Úidú/ReducerFunction[StateT, DepsT, InputT, OutputT]Ú_reducerúCallable[[], OutputT]Ú_initial_factoryúForkID | NoneÚparent_fork_idzLiteral['closest', 'farthest']Úpreferred_parent_forkNÚfarthest)rw   rx   c               ó@   • Xl         X l        X0l        X@l        XPl        g r9   )rq   rs   ru   rw   rx   )r:   rq   ÚreducerÚinitial_factoryrw   rx   s         r0   r>   ÚJoin.__init__«   s    € ð ŒØŒØ /ÔØ,ÔØ%:Õ"r/   c                ó   • U R                   $ r9   )rs   rA   s    r0   r{   ÚJoin.reducerº   s   € à�}‰}Ðr/   c                ó   • U R                   $ r9   )ru   rA   s    r0   r|   ÚJoin.initial_factory¾   s   € à×$Ñ$Ð$r/   c                ó>  • [        [        R                  " U R                  5      R                  5      nUS:X  a-  [        [        [        [        4   U R                  5      " X#5      $ [        [        [        [        [        [        4   U R                  5      " XU5      $ )Né   )ÚlenÚinspectÚ	signaturer{   Ú
parametersr   rJ   r   r   rL   r   r   )r:   rk   r#   rP   Ún_parameterss        r0   ÚreduceÚJoin.reduceÂ   su   € Üœ7×,Ò,¨T¯\©\Ó:×EÑEÓFˆØ˜1ÓÜÔ,¬V´W¨_Ñ=¸t¿|¹|ÔLÈWÓ]Ð]äÔ.¬v´u¼fÄgÐ/MÑNÐPT×P\ÑP\Ô]Ð^aÐlrÓsÐsr/   c                ó   • g r9   r'   ©r:   rP   s     r0   Úas_nodeÚJoin.as_nodeÉ   s   € ØGJr/   c                ó   • g r9   r'   rŒ   s     r0   r�   rŽ   Ì   s   € ØBEr/   c                ó   • [        X5      $ )zÅCreate a join node with bound inputs.

Args:
    inputs: The input data to bind to this join, or None

Returns:
    A [`JoinNode`][pydantic_graph.join.JoinNode] with this join and the bound inputs
)ÚJoinNoderŒ   s     r0   r�   rŽ   Ï   s   € ô ˜Ó%Ð%r/   )ru   rs   rq   rw   rx   )
rq   r   r{   rr   r|   rt   rw   rv   rx   zLiteral['farthest', 'closest'])rk   zReducerContext[StateT, DepsT]r#   r   rP   r   rH   r   r9   )rP   ÚNonerH   úJoinNode[StateT, DepsT])rP   r   rH   r“   )rP   zInputT | NonerH   r“   )r(   r)   r*   r+   r,   r-   r>   rI   r{   r|   r‰   r   r�   r.   r'   r/   r0   ro   ro   –   s¼   ‡ ñð 	ƒJØ=Ó=Ø+Ó+Ø!Ó!Ø9Ó9ð )-Ø@Jñ;ð ð;ð Að	;ð
 /ð;ð &ð;ð  >õ;ð ñó ðð ñ%ó ð%ôtð ÝJó ØJàÛEó ØE÷	&ñ 	&r/   ro   c                  ó<   • \ rS rSr% SrS\S'    S\S'    S
S jrSrg	)r‘   éÛ   a‘  A `BaseNode` that represents a builder join with bound inputs.

`JoinNode` lets a [`BaseNode`][pydantic_graph.BaseNode] subclass hand off to a builder
[`Join`][pydantic_graph.join.Join] by wrapping the join together with the value it should
receive as `inputs`. It is not meant to be run directly; returning a `JoinNode` from a
`BaseNode.run` method tells the graph builder which join to invoke next.
zJoin[StateT, DepsT, Any, Any]Újoinr	   rP   c              ƒ  ó    #   • [        S5      e7f)zÑAttempt to run the join node.

Args:
    ctx: The graph execution context

Returns:
    The result of step execution

Raises:
    NotImplementedError: Always raised as StepNode is not meant to be run directly
z�`JoinNode` is not meant to be run directly, it is meant to be returned from a `BaseNode` subclass to indicate a transition to a builder join.)ÚNotImplementedError)r:   rk   s     r0   ÚrunÚJoinNode.runë   s   é € ô "ð \ó
ð 	
ùs   ‚r'   N)rk   zGraphRunContext[StateT, DepsT]rH   z'BaseNode[StateT, DepsT, Any] | End[Any])r(   r)   r*   r+   r,   r-   r™   r.   r'   r/   r0   r‘   r‘   Û   s   ‡ ñð (Ó'ØàƒKØ(÷
r/   r‘   )r#   r’   rP   r	   rH   r’   )r#   úlist[T]rP   r   rH   r›   )r#   r›   rP   zIterable[T]rH   r›   )r#   ú
dict[K, V]rP   zMapping[K, V]rH   rœ   )r#   rc   rP   rc   rH   rc   )5r,   Ú
__future__r   r…   Úabcr   Úcollections.abcr   r   r   Údataclassesr   Útypingr	   r
   r   r   r   Útyping_extensionsr   r   r   r   Úpydantic_graphr   r   r   Úpydantic_graph.id_typesr   r   r   r   r   r   r   r   r   r   r!   r3   rJ   rL   rM   rQ   rT   rX   r[   r]   rc   rf   rh   ro   r‘   r'   r/   r0   Ú<module>r¥      s   ðñõ #ã Ý ß 7Ñ 7Ý !ß 8Õ 8ç DÓ Dç 9Ñ 9ß =Ñ =á	�¨$Ñ	/€Ù�¨Ñ-€Ù	�¨$Ñ	/€Ù
�)¨DÑ
1€ÙˆC Ñ%€ÙˆC Ñ%€ÙˆC Ñ%€ð ÷*ð *ó ð*ñ �Ñô%8�W˜V U˜]Ñ+ó %8ó ð%8ñP %ØØˆg�vÐ Ð'Ñ(Ø˜Ð!ñÐ ñ
 'ØØˆn˜V U˜]Ñ+¨W°fÐ=¸wÐFÑGØ˜ ¨Ð0ñÐ ñ
  ØØ˜6 5¨&°'Ð9Ñ:Ð=QÐRXÐZaÐRaÑ=bÑbØ˜ ¨Ð0ñ€ð
ô
ô
ôôô�(ô ñ �: [ÀÑF€ôð
 ô�w˜q‘zó ó ðñ �ÑôA&ˆ7�6˜5 &¨'Ð1Ñ2ó A&ó ðA&ðH ô
ˆx˜  sÐ*Ñ+ó 
ó ñ
r/   