
    Ji0                       d dl mZ d dlZd dlZd dlmZmZ d dlmZm	Z	m
Z
mZmZmZmZmZmZmZ d dlmZmZ d dlmZ d dlmZmZ d dlmZ d d	lmZmZ d d
lm Z  d dl!m"Z"m#Z#m$Z$ d dl%m&Z&m'Z'm(Z(m)Z)m*Z*m+Z+m,Z, d dl-m.Z.m/Z/ d dl0m1Z1m2Z2m3Z3 d dl4m5Z5 d dl6m7Z8 ddl9m:Z: erd dl;m<Z< d dl=m>Z>  ede      Z?e@eAeBe   f   ZC ej                         ZE G d deF      ZG G d dee:         ZHy)    )annotationsN)Counterdefaultdict)
TYPE_CHECKINGAnyAsyncGenerator	AwaitableCallable	CoroutineGenericTypeTypeVarcast)active_instrument_tagsinstrument_tags)WorkflowRuntimeError)Event
StartEvent)WorkflowHandler)control_looprebuild_state_from_ticks)BrokerState)PluginWorkflowRuntimeas_snapshottable)AddCollectedEvent	AddWaiterDeleteCollectedEventDeleteWaiterStepWorkerContextStepWorkerStateContextVarWaitingForEvent)StepWorkerFunctionas_step_worker_function)TickAddEventTickCancelRunWorkflowTick)workflow_registry)_nanoid   )MODEL_T)Workflow)ContextT)boundc                      e Zd Zy)UnserializableKeyWarningN)__name__
__module____qualname__     j/var/www/html/BankruptcyAI-uat/bankruptcy-ai/venv/lib/python3.12/site-packages/workflows/runtime/broker.pyr1   r1   B   s    r6   r1   c                  h   e Zd ZU dZded<   ded<   ded<   ded	<   d
ed<   ded<   ded<   ded<   	 	 	 	 	 	 	 	 	 	 d#dZd$dZ	 	 	 	 d%	 	 	 	 	 	 	 	 	 	 	 d&dZd'dZe	d(d       Z
e	d)d       Ze	d*d       Zd+dZ	 d,	 	 	 	 	 	 	 d-dZd,d.dZd/dZ	 	 	 	 d0	 	 	 	 	 	 	 	 	 	 	 d1dZd2dZd3d Zd4d!Zd'd"Zy)5WorkflowBrokerz
    The workflow broker manages starting up and connecting a workflow handler, a runtime, and triggering the
    execution of the workflow. From there it manages communication between the workflow and the outside world.
    Context[MODEL_T]_contextr   _runtimer   _pluginbool_is_runningzWorkflowHandler | None_handlerr,   	_workflowzlist[asyncio.Task]_workersBrokerState | None_init_statec                t    || _         || _        || _        d| _        d | _        || _        g | _        d | _        y )NF)r;   r<   r=   r?   r@   rA   rB   rD   )selfworkflowcontextruntimeplugins        r7   __init__zWorkflowBroker.__init__V   s>       !r6   c                     t        j                  |       j                  j                         d fd}j	                  |       S )Nc                \    	 j                   j                         y # t        $ r Y y w xY wN)rB   remove
ValueError)_rF   tasks    r7   _remove_taskz2WorkflowBroker._execute_task.<locals>._remove_taskj   s,    $$T* s    	++)rQ   asyncio.Task[Any]returnNone)asynciocreate_taskrB   appendadd_done_callback)rF   cororS   rR   s   `  @r7   _execute_taskzWorkflowBroker._execute_taskf   s?    ""4(T"	 	|,r6   Nc                   	  j                   t        d       _        d	 fd}t               } j	                   ||t        j                                     }t         j                  ||      		 _         	S )z0Start the workflow run. Can only be called once.z?this WorkflowBroker already run or running. Cannot start again.c                  K   t        d| i|      5  d_        t        j                  d       d {             d {    	 xs t	        j
                        }	 d }i }j                         j                         D ]   \  }}t        |d|      }t        |      ||<   " t        j                  j                  t        |      }t        j                  | j                  j                   |j"                         	 |j%                  ||        d {   }	t        j&                  |        j)                  |	       |r!j-                         sj/                  |                d {    d_        	 d d d        y 7 [7 M7 u# t        j&                  |        w xY w# t*        $ r}
|
}Y d }
~
wd }
~
ww xY w7 O#          d {  7   d_        w xY w# 1 sw Y   y xY ww)Nrun_idTr   __func__)r_   rG   rJ   rH   stepsF)r   r?   rW   sleepr   from_workflow
_get_stepsitemsgetattrr$   r(   get_registered_workflowr=   r   register_runr<   r;   ra   workflow_function
delete_run_set_stop_event	Exceptiondoneset_exception)r_   tags
init_stateexception_raisedstep_workersname	step_funcunbound
registeredworkflow_resulteafter_completebefore_startpreviousresultrF   start_eventrG   s              r7   _run_workflowz+WorkflowBroker.start.<locals>._run_workflow   s     (F!;d!;< 5- $( mmA&&&+&.((.-!)!P[-F-Fx-PJ"-+/(FH/7/B/B/D/J/J/L ROD) '.iY&OG1H1QL.	R &7%N%N$dllL,&

 *66#)%-#'==$(MM","2"2A4>4P4P + * &5 /O .88@..? (%{{}"001AB%1,...',D$k5- 5- '(8/ .88@$ -+,(- / &1,...',D$k5- 5-s   G4 G(FG(FG(
G
$B#F0FFF#&F0	#G
,G(8G9G(	G4G(G(FF--F00	G9G ;G
 GG
G(
G%G
G%%G((G1-G4)ro   )ctxr_   run_task)r_   strro   zdict[str, Any]rU   rV   )	r@   r   rD   nanoidr\   r   getr   r;   )
rF   rG   r{   r}   rz   ry   r~   r_   r   r|   s
   ``````   @r7   startzWorkflowBroker.startu   s     ==$&Q  $6	- 6	-r  %%&'='A'A'CD
 !

 r6   c                h    | j                  | j                  j                  t                            y rN   )r\   r<   
send_eventr&   rF   s    r7   
cancel_runzWorkflowBroker.cancel_run   s!    4==33MODEr6   c                    | j                   S rN   )r?   r   s    r7   
is_runningzWorkflowBroker.is_running   s    r6   c                    | j                   }| j                  xs t        j                  | j                        }t        ||      }|S rN   )	_tick_logrD   r   rc   rA   r   )rF   ticksstate	new_states       r7   _statezWorkflowBroker._state   s<      MK$=$=dnn$M,UE:	r6   c                f    t        | j                        }|t        d      |j                         S )NzPlugin is not snapshottable)r   r<   r   replay)rF   snapshottables     r7   r   zWorkflowBroker._tick_log   s1    (7 &'DEE##%%r6   c                   K   | j                   j                  j                         D cg c]'  }| j                   j                  |   j                  r|) c}S c c}w wrN   )r   workerskeysin_progress)rF   steps     r7   running_stepszWorkflowBroker.running_steps   sS      ++002
{{""4(44 
 	
 
s   'A,AAc           	        | j                  d      }|xs d}|j                  j                  j                  |g       }t	        |      t	        |D cg c]  }t        |       c}      z
  }|t	        t        |      g      k7  r>t        |      |v r0|j                  j                  j                  t        ||             y g }t        t              }	||gz   D ]  }|	t        |         j                  |       ! |D ]%  }
|j                  |	|
   j                  d             ' |j                  j                  j                  t        |             |S c c}w )Ncollect_eventsfndefault)event_ideventr   )r   )_get_step_ctxr   collected_eventsr   r   typereturnsreturn_valuesrY   r   r   listpopr   )rF   evexpected	buffer_idstep_ctxr   rx   remaining_event_typestotalby_typee_types              r7   r   zWorkflowBroker.collect_events   sG    %%)9%:*	#>>::>>y"M ' 1G./T!W/5
 !
 !GT"XJ$77Bx00  ..55%yC d#!RD( 	'ADG##A&	'  	1FLL,,Q/0	1 	&&--.BI.VW' 0s   E
c                |   ||| j                   j                         vrt        d| d      | j                   j                         |   }|j                  }t	        |      |j
                  vrt        d| dt	        |             | j                  | j                  j                  t        ||                   y )NzStep z does not existz does not accept event of type )r   	step_name)
rA   rd   r   _step_configr   accepted_eventsr\   r<   r   r%   )rF   messager   rt   step_configs        r7   r   zWorkflowBroker.send_event  s    4>>4466*U4&+HII 113D9I#00KG}K$?$??*D6!@gP  	MM$$\4%PQ	
r6   c                b    	 t        j                         S # t        $ r t        | d      w xY w)Nz/ may only be called from within a step function)r!   r   LookupErrorr   )rF   r   s     r7   r   zWorkflowBroker._get_step_ctx  s=    	,0022 	&$EF 	s    .c           	       K   | j                  d      }|j                  j                  }|xs i }| j                  |      }t	        |      }	xs d| d|	 t        fd|D        d       }
|
|
j                  t        t        ||||            |j                  j                  j                  t                     t        t        |
j                        S w)Nwait_for_eventr   waiter_rQ   c              3  B   K   | ]  }|j                   k(  s|  y wrN   	waiter_id).0wr   s     r7   	<genexpr>z0WorkflowBroker.wait_for_event.<locals>.<genexpr>7  s     PQq{{i7OqPs   )r   requirementstimeout
event_typewaiter_eventr   )r   r   collected_waiters_get_full_pathr   nextresolved_eventr"   r   r   r   rY   r   r   r.   )rF   r   r   r   r   r   r   r   	event_strrequirements_strwaiters      `       r7   r   zWorkflowBroker.wait_for_event%  s      %%)9%:$NN<<#)r ''
3	|,I79+Q7G6H!I	P"3PRVW>V22:!'!-#)!-  **11,2ST60011s   CCc                8    |j                    d|j                   S )N.)r3   r2   )rF   ev_types     r7   r   zWorkflowBroker._get_full_pathF  s!    $$%Qw'7'7&899r6   c                6    | j                   j                         S )z8The internal queue used for streaming events to callers.)r<   stream_published_eventsr   s    r7   r   z&WorkflowBroker.stream_published_eventsI  s    }}4466r6   c                ^    |+| j                  | j                  j                  |             y y rN   )r\   r<   write_to_event_stream)rF   r   s     r7   write_event_to_streamz$WorkflowBroker.write_event_to_streamN  s)    >t}}BB2FG r6   c                $  K   | j                   j                  t                      d{    | j                  D ]  }|j	                           | j                  j                          | j                   j                          d{    y7 b7 w)zCancels the running workflow loop

        Cancels all outstanding workers, waits for them to finish, and marks the
        broker as not running. Queues and state remain available so callers can
        inspect or drain leftover events.
        N)r<   r   r&   rB   cancelclearclose)rF   workers     r7   shutdownzWorkflowBroker.shutdownR  so      mm&&}777mm 	FMMO	mm!!###	 	8 	$s"   'BBABBBB)
rG   r,   rH   r:   rI   r   rJ   r   rU   rV   )r[   zCoroutine[Any, Any, Any]rU   rT   )NNNN)rG   r,   r{   rC   r}   zStartEvent | Nonerz   $Callable[[], Awaitable[None]] | Nonery   r   rU   r   )rU   rV   )rU   r>   )rU   r   )rU   zlist[WorkflowTick])rU   z	list[str]rN   )r   r   r   zlist[Type[Event]]r   
str | NonerU   zlist[Event] | None)r   r   r   r   rU   rV   )r   r   rU   r    )NNNi  )r   zType[T]r   Event | Noner   r   r   zdict[str, Any] | Noner   zfloat | NonerU   r.   )r   zType[Event]rU   r   )rU   zAsyncGenerator[Event, None])r   r   rU   rV   )r2   r3   r4   __doc____annotations__rK   r\   r   r   propertyr   r   r   r   r   r   r   r   r   r   r   r   r5   r6   r7   r9   r9   F   s   
 O$$  ##   "  !	 
   
  $ (,)-=A?CUU %U '	U
 ;U =U 
UpF       & &
 OS#4AK	@
" &* $.2 $22 #2 	2
 ,2 2 
2B:7
H$r6   r9   )I
__future__r   rW   loggingcollectionsr   r   typingr   r   r   r	   r
   r   r   r   r   r   &llama_index_instrumentation.dispatcherr   r   workflows.errorsr   workflows.eventsr   r   workflows.handlerr   workflows.runtime.control_loopr   r   &workflows.runtime.types.internal_stater   workflows.runtime.types.pluginr   r   r   workflows.runtime.types.resultsr   r   r   r   r    r!   r"   %workflows.runtime.types.step_functionr#   r$   workflows.runtime.types.ticksr%   r&   r'   #workflows.runtime.workflow_registryr(   workflows.utilsr)   r   context.state_storer+   	workflowsr,   workflows.context.contextr-   r.   dictr   r   EventBuffer	getLoggerloggerWarningr1   r9   r5   r6   r7   <module>r      s    #   ,   2 . Q > T T   T S A - )"1 Cu3U#$					w 	W$WW% W$r6   