
    Ji3                        d dl mZ d dlmZ d dlmZmZmZmZ d dl	Z	d dl
m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 d dlmZmZ dd	Z G d
 d      Z	 d	 	 	 	 	 ddZy)    )annotations)asynccontextmanager)AnyAsyncGeneratorAsyncIteratoroverloadN)Context)Event
StartEvent)CancelHandlerResponseHandlerDataHandlersListResponseHealthResponseSendEventResponseStatusWorkflowsListResponse)EventEnvelopeEventEnvelopeWithMetadatac                   	 | j                          y# t        j                  $ r}d|j                  j                  cxk  rdk  rn  |j                  j
                  dd }|j                  j                  }|j                  j                  }|j                  j                  }t        j                  | d|j                  j                   d| d| d| 	|j                  |j                        | d}~ww xY w)	zw
    Raise an HTTPStatusError with the first 200 characters of the response body
    for 400 and 500 level errors.
    i  iX  N    z for z. Response: )requestresponse)
raise_for_statushttpxHTTPStatusErrorr   status_codetextr   methodurlreason_phrase)r   ebody_previewr   r    r   s         i/var/www/html/BankruptcyAI-uat/bankruptcy-ai/venv/lib/python3.12/site-packages/workflows/client/client.py_raise_for_status_with_bodyr%   !   s    
!!#   !**((.3. 	 ::??4C0LYY%%F))--C**00K''-q!9!9 :%xq\ZfYgh		 	
 	s    C8CC33C8c                  $   e Zd Zedd       Ze	 	 dd       Zddd	 	 	 ddZedd       ZddZddZ	 	 	 d	 	 	 	 	 	 	 	 	 dd	Z		 	 	 d	 	 	 	 	 	 	 	 	 dd
Z
	 	 d	 	 	 	 	 	 	 ddZ	 d	 	 	 	 	 	 	 ddZddZ	 	 d	 	 	 	 	 ddZddZ	 d 	 	 	 	 	 d!dZy)"WorkflowClientc                    y N )selfhttpx_clients     r$   __init__zWorkflowClient.__init__7   s    <?    c                    y r)   r*   )r+   base_urls     r$   r-   zWorkflowClient.__init__9   s    
 r.   N)r,   r0   c               \    ||t        d      ||t        d      || _        || _        y )Nz0Either httpx_client or base_url must be providedz5Only one of httpx_client or base_url must be provided)
ValueErrorr,   r0   )r+   r,   r0   s      r$   r-   zWorkflowClient.__init__@   sA     H$4OPP#(<TUU( r.   c                  K   | j                   r| j                    y t        j                  | j                  xs d      4 d {   }| d d d       d {    y 7 7 # 1 d {  7  sw Y   y xY ww)N )r0   )r,   r   AsyncClientr0   )r+   clients     r$   _get_clientzWorkflowClient._get_clientM   sh     ###(($--2E2F  &      sH   AA;A"A;A&A;A$A;$A;&A8,A/-A84A;c                "  K   | j                         4 d{   }|j                  d       d{   }t        |       t        j                  |j                               cddd      d{    S 7 \7 E7 	# 1 d{  7  sw Y   yxY ww)z
        Check whether the workflow server is helathy or not

        Returns:
            HealthResponse: health response from the workflow
        Nz/health)r7   getr%   r   model_validatejsonr+   r6   r   s      r$   
is_healthyzWorkflowClient.is_healthyU   s      ##% 	B 	B#ZZ	22H'1!00A	B 	B 	B2	B 	B 	B 	BT   BA4BA:A61A:"B.A8/B6A:8B:B BBBc                "  K   | j                         4 d{   }|j                  d       d{   }t        |       t        j                  |j                               cddd      d{    S 7 \7 E7 	# 1 d{  7  sw Y   yxY ww)z
        List workflows

        Returns:
            WorkflowsListResponse: List of workflow names available through the server.
        Nz
/workflows)r7   r9   r%   r   r:   r;   r<   s      r$   list_workflowszWorkflowClient.list_workflowsa   s      ##% 	I 	I#ZZ55H'1(77H	I 	I 	I5	I 	I 	I 	Ir>   c                4  K   |	 t        |d      }t        |t              r	 |j                         }|xs d|xs i d}|r||d<   | j                         4 d{   }|j                  d	| d
|       d{   }t        |       t        j                  |j                               cddd      d{    S # t        $ r}t        d|       d}~ww xY w# t        $ r}t        d|       d}~ww xY w7 7 7 G# 1 d{  7  sw Y   yxY ww)al  
        Run the workflow and wait until completion.

        Args:
            start_event (Union[StartEvent, dict[str, Any], None]): start event class or dictionary representation (optional, defaults to None and get passed as an empty dictionary if not provided).
            context: Context or serialized representation of it (optional, defaults to None if not provided)
            handler_id (Optional[str]): Workflow handler identifier to continue from a previous completed run.

        Returns:
            HandlerData: Data representing the handler running the workflow (including result and metadata)
        NT)bare4Impossible to serialize the start event because of: 0Impossible to serialize the context because of: r4   start_eventcontext
handler_id/workflows/z/runr;   )_serialize_event	Exceptionr2   
isinstancer	   to_dictr7   postr%   r   r:   r;   	r+   workflow_namerH   rF   rG   r"   request_bodyr6   r   s	            r$   run_workflowzWorkflowClient.run_workflowo   sA    $ ".{F
 gw'Y!//+ ',"}"
 )3L&##% 	? 	?#[[m_D1 )  H (1--hmmo>	? 	? 	?   J1#N   Y #STUSV!WXXY	?	? 	? 	? 	?s   DB? DC (DC=D D;C?<1D-D9D:D?	CCCD	C:'C55C::D?DDD	D
DDc                R  K   |	 t        |      }t        |t              r	 |j                         }|xs t        t                     |xs i d}|r||d<   | j                         4 d{   }|j                  d| d|       d{   }t        |       t        j                  |j                               cddd      d{    S # t        $ r}t        d|       d}~ww xY w# t        $ r}t        d|       d}~ww xY w7 7 7 G# 1 d{  7  sw Y   yxY ww)	aE  
        Run the workflow in the background.

        Args:
            start_event (Union[StartEvent, dict[str, Any], None]): start event class or dictionary representation (optional, defaults to None and get passed as an empty dictionary if not provided).
            context: Context or serialized representation of it (optional, defaults to None if not provided)
            handler_id (Optional[str]): Workflow handler identifier to continue from a previous completed run.

        Returns:
            HandlerData: data representing the handler running the workflow.
        NrC   rD   rE   rH   rI   z/run-nowaitrJ   )rK   rL   r2   rM   r	   rN   r   r7   rO   r%   r   r:   r;   rP   s	            r$   run_workflow_nowaitz"WorkflowClient.run_workflow_nowait   sG    $ ".{;
 gw'Y!//+ 'H*::<*H}"(
 )3L&##% 	? 	?#[[m_K8| )  H (1--hmmo>	? 	? 	?   J1#N   Y #STUSV!WXXY	?	? 	? 	? 	?s   D'C D'C- 9D'+D,D'/D
D1D<D'D	D'	C*C%%C**D'-	D	6DD		D'DD'D$DD$ D'c           	      K   |rdnd}d| }| j                         4 d{   }	 |j                  d|d||dddid	      4 d{   }|j                  d
k(  rt        d      |j                  dk(  r"	 ddd      d{    ddd      d{    yt	        |       |j                         2 3 d{   }|j                         st        j                  |      }	|	 57 7 7 i7 [7 86 ddd      d{  7   n# 1 d{  7  sw Y   nxY wnI# t        j                  $ r t        d|       t        j                  $ r}
t        d|
       d}
~
ww xY wddd      d{  7   y# 1 d{  7  sw Y   yxY ww)a  
        Stream events as they are produced by the workflow.

        Args:
            handler_id (str): ID of the handler running the workflow
            include_internal_events (bool): Include internal workflow events. Defaults to False.
            lock_timeout (float): Timeout (in seconds) for acquiring the lock to iterate over the events.

        Returns:
            AsyncGenerator[EventEnvelopeWithMetadata, None]: Generator for the events that are streamed as instances of `EventEnvelopeWithMetadata`.
        truefalse/events/NGET)sseinclude_internalacquire_timeout
Connectionz
keep-alive)paramsheaderstimeouti  zHandler not found   z(Timeout waiting for events from handler z#Failed to connect to event stream: )r7   streamr   r2   r%   aiter_linesstripr   model_validate_jsonr   TimeoutExceptionTimeoutErrorRequestErrorConnectionError)r+   rH   include_internal_eventslock_timeout
incl_interr    r6   r   lineeventr"   s              r$   get_workflow_eventsz"WorkflowClient.get_workflow_events   s    "  7VG
%##%  	Q  	QQ!==&,6+7
 *<8  ) 
 ( ( ++s2()<==!--4!( ( 	Q  	Q  	Q( 09&.&:&:&< ( (d::<$=$Q$QRV$WE"'K3 	Q( ( 	Q,(&<)( ( ( ( (2 )) ">zlK  %% Q%(KA3&OPPQ? 	Q  	Q  	Q  	Q  	Qs	   FC(FE6 DC*	D+D7DC,DFC.FD3C27C0
8C2;DD(F*D,D.F0C22D3D>D?DD	DD	DE63E!EE!!E6$F/E20F6F<E?=FFc                  K   	 t        |      }d|i}|r||d<   | j                         4 d{   }|j	                  d| |       d{   }t        |       t        j                  |j                               cddd      d{    S # t        $ r}t        d|       d}~ww xY w7 7 d7 (# 1 d{  7  sw Y   yxY ww)a  
        Send an event to the workflow.

        Args:
            handler_id (str): ID of the handler of the running workflow to send the event to
            event (Event | dict[str, Any] | str): Event to send, represented as an Event object, a dictionary or a serialized string.
            step (Optional[str]): Step to send the event to (optional, defaults to None)

        Returns:
            SendEventResponse: Confirmation of the send operation
        z,Error while serializing the provided event: Nro   steprY   rJ   )	rK   rL   r2   r7   rO   r%   r   r:   r;   )	r+   rH   ro   rr   serialized_eventr"   rR   r6   r   s	            r$   
send_eventzWorkflowClient.send_event   s     "	Q/?/F )01A'B#'L ##% 	E 	E#[[8J<)@|[TTH'1$33HMMOD		E 	E 	E  	QKA3OPP	Q
	ET	E 	E 	E 	Esx   C
B C
B/C
B5B11B5>C

B3C
	B,B''B,,C
1B53C
5C;B><CC
c                @   K   | j                  |       d{   S 7 w)z6
        Deprecated. Use get_handler instead.
        N)get_handler)r+   rH   s     r$   
get_resultzWorkflowClient.get_result  s      %%j1111s   c                ,  K   | j                         4 d{   }|j                  d||d       d{   }t        |       t        j                  |j                               cddd      d{    S 7 a7 E7 	# 1 d{  7  sw Y   yxY ww)aq  
        Get all the workflow handlers.
        Args:
            status (list[Status] | None): List of statuses (e.g. "running", "completed", etc. ) to filter by. Defaults to None.
            workflow_name (list[str] | None): List of workflow names to filter by. Defaults to None.
        Returns:
            HandlersListResponse: List of workflow handlers.
        Nz	/handlers)statusrQ   r_   )r7   r9   r%   r   r:   r;   )r+   ry   rQ   r6   r   s        r$   get_handlerszWorkflowClient.get_handlers#  s      ##% 
	H 
	H#ZZ$%2 (  H (1'66x}}G
	H 
	H 
	H
	H 
	H 
	H 
	HsT   BA9BA?A;1A?'B3A=4B;A?=B?BBBBc                (  K   | j                         4 d{   }|j                  d|        d{   }t        |       t        j                  |j                               cddd      d{    S 7 _7 E7 	# 1 d{  7  sw Y   yxY ww)z
        Get a single workflow handler by identifier.

        Args:
            handler_id (str): ID of the handler associated with the workflow run

        Returns:
            HandlerData: Handler metadata persisted by the server.
        N
/handlers/)r7   r9   r%   r   r:   r;   )r+   rH   r6   r   s       r$   rv   zWorkflowClient.get_handler<  s}      ##% 	? 	?#ZZ*ZL(ABBH'1--hmmo>		? 	? 	?B	? 	? 	? 	?sT   BA7BA=A91A=%B1A;2B9A=;B=BBBBc                :  K   | j                         4 d{   }|j                  d| dd|rdndi       d{   }t        |       t        j                  |j                               cddd      d{    S 7 h7 E7 	# 1 d{  7  sw Y   yxY ww)a   
        Stop and cancel a workflow run.

        Args:
            handler_id (str): ID of the handler associated with the workflow run
            purge (bool): Whether or not to delete the run also from the persistent storage. Defaults to false
        Nr}   z/cancelpurgerW   rX   rz   )r7   rO   r%   r   r:   r;   )r+   rH   r   r6   r   s        r$   cancel_handlerzWorkflowClient.cancel_handlerL  s      ##% 	I 	I#[[ZL05g> )  H (1(77H	I 	I 	I	I 	I 	I 	IsT   BB B!BB1B.B:B;BBBBBBB)r,   zhttpx.AsyncClient)r0   str)r,   zhttpx.AsyncClient | Noner0   
str | None)returnz AsyncIterator[httpx.AsyncClient])r   r   )r   r   )NNN)
rQ   r   rH   r   rF   z"StartEvent | dict[str, Any] | NonerG   zContext | dict[str, Any] | Noner   r   )F   )rH   r   rk   boolrl   floatr   z/AsyncGenerator[EventEnvelopeWithMetadata, None]r)   )rH   r   ro   Event | dict[str, Any]rr   r   r   r   )rH   r   r   r   )NN)ry   zlist[Status] | NonerQ   zlist[str] | Noner   r   F)rH   r   r   r   r   r   )__name__
__module____qualname__r   r-   r   r7   r=   r@   rS   rU   rp   rt   rw   r{   rv   r   r*   r.   r$   r'   r'   6   s   ? ?   26#	! /! 	!  
BI" "&:>37+?+? +? 8	+?
 1+? 
+?` "&:>37+?+? +? 8	+?
 1+? 
+?` ).	4Q4Q "&4Q 	4Q
 
94Qt  	EE &E 	E
 
E<2 '+*.H#H (H 
	H2?" .3II&*I	Ir.   r'   c                    t        | t              r| S |r| j                         S t        j                  |       j                         S )N)ro   )rM   dict
model_dumpr   
from_event)ro   rB   s     r$   rK   rK   `  sI     %  	 %%E2==?r.   )r   zhttpx.Responser   Noner   )ro   r   rB   r   r   zdict[str, Any])
__future__r   
contextlibr   typingr   r   r   r   r   	workflowsr	   workflows.eventsr
   r   workflows.protocolr   r   r   r   r   r   r   &workflows.protocol.serializable_eventsr   r   r%   r'   rK   r*   r.   r$   <module>r      sl    # *    .  *gI gIV	 16	!	)-		r.   