ó
    Ú°›j%  ã                   óh   • S r SSKrSSKJrJrJr  SSKJrJr  SSK	J
r
   " S S5      r " S S	5      rg)
z,Module for running endpoints asynchronously.é    N)ÚAnyÚDictÚOptional)ÚFINAL_STATESÚis_completed)ÚClientSessionc                   ó˜   • \ rS rSrSrS\S\S\S\4S jrSS\S	\	\\
4   4S
 jjrS	\4S jrS rSS\S	\
4S jjrS	\
4S jrS	\4S jrSrg)ÚJobé   z5Class representing a job for an asynchronous endpointÚendpoint_idÚjob_idÚsessionÚheadersc                 óh   • SSK Jn  Xl        X l        X0l        XPl        X@l        SU l        SU l        g)zÏ
Initialize a Job instance.

Args:
    endpoint_id: The identifier for the endpoint.
    job_id: The identifier for the job.
    session: The aiohttp ClientSession.
    headers: Headers to use for requests.
r   )Úendpoint_url_baseN)Úrunpodr   r   r   r   r   Ú
job_statusÚ
job_output)Úselfr   r   r   r   r   s         Úc/home/mande/repo/quber/.venv/lib/python3.13/site-packages/runpod/endpoint/asyncio/asyncio_runner.pyÚ__init__ÚJob.__init__   s2   € õ	
ð 'ÔØŒØŒØ!2ÔØŒàˆŒØˆ�ó    ÚsourceÚreturnc              ƒ   óZ  #   • U R                    SU R                   SU SU R                   3nU R                  R	                  X R
                  S9I Sh  v•N nUR                  5       I Sh  v•N n[        US   5      (       a!  US   U l        UR	                  SS5      U l	        U$  NR N<7f)z~Returns the raw json of the status, raises an exception if invalid.

Args:
    source: The URL source path of the job status.
Ú/©r   NÚstatusÚoutput)
r   r   r   r   Úgetr   Újsonr   r   r   )r   r   Ú
status_urlÚ	job_states       r   Ú
_fetch_jobÚJob._fetch_job&   s¤   é € ð ×%Ñ%Ð& a¨×(8Ñ(8Ð'9¸¸6¸(À!ÀDÇKÁKÀ=ÐQð 	ð Ÿ,™,×*Ñ*¨:¿|¹|Ð*ÐL×Lˆ	Ø#Ÿ.™.Ó*×*ˆ	ä˜	 (Ñ+×,Ñ,Ø'¨Ñ1ˆDŒOØ'Ÿm™m¨H°dÓ;ˆDŒOàÐñ MÙ*ùs$   ‚AB+ÁB'ÁB+Á,B)Á-;B+Â)B+c              ƒ   óz   #   • U R                   b  U R                   $ U R                  5       I Sh  v•N nUS   $  N	7f)zAGets jobs' status

Returns:
    COMPLETED, FAILED or IN_PROGRESS
Nr   )r   r%   )r   r$   s     r   r   Ú
Job.status8   s:   é € ð �?‰?Ñ&Ø—?‘?Ð"àŸ/™/Ó+×+ˆ	Ø˜Ñ"Ð"ñ ,ùs   ‚-;¯9°
;c              ƒ   óò   #   • [        U R                  5       I S h  v•N 5      (       dG  [        R                  " S5      I S h  v•N   [        U R                  5       I S h  v•N 5      (       d  MF  g g  NU N0 N7f)Né   )r   r   ÚasyncioÚsleep)r   s    r   Ú_wait_for_completionÚJob._wait_for_completionD   sL   é € Ü T§[¡[£]×2×3Ñ3Ü—-’- Ó"×"Ð"ô  T§[¡[£]×2×3Õ3Ñ2Ù"ñ  3ùs9   ‚A7›A1œ&A7ÁA3ÁA7ÁA5Á A7Á/A7Á3A7Á5A7Útimeoutc              ƒ   óD  #   • U R                   b  U R                   $  [        R                  " U R                  5       U5      I Sh  v•N   U R                  5       I Sh  v•N nUR                  SS5      $  N.! [        R                   a  n[	        S5      UeSnAff = f N@7f)zpWaits for serverless API job to complete or fail

Returns:
    Output of job
Raises:
    KeyError if job Failed
NzJob timed out.r    )r   r+   Úwait_forr-   ÚTimeoutErrorr%   r!   )r   r/   ÚexcÚjob_datas       r   r    Ú
Job.outputH   s�   é € ð �?‰?Ñ&Ø—?‘?Ð"ð	:Ü×"Ò" 4×#<Ñ#<Ó#>ÀÓH×HÐHð Ÿ™Ó*×*ˆØ�|‰|˜H dÓ+Ð+ñ IøÜ×#Ñ#ó 	:ÜÐ/Ó0°cÐ9ûð	:úñ +ùsF   ‚B �(A6 ÁA4ÁA6 Á
B ÁBÁB Á4A6 Á6BÂ
BÂBÂB c                ó:  #   •  [         R                  " S5      I Sh  v•N   U R                  SS9I Sh  v•N nUS   [        ;  d  [	        UR                  S/ 5      5      S:”  a"  UR                  S/ 5       H  nUS   7v •  M     OUS   [        ;   a  gM“   Ny Nd7f)z>Returns a generator that yields the output of the job request.r*   NÚstream)r   r   r   r    )r+   r,   r%   r   Úlenr!   )r   Ústream_partialÚchunks      r   r7   Ú
Job.stream[   s�   é € àÜ—-’- Ó"×"Ð"Ø#'§?¡?¸( ?Ð#C×CˆNà˜xÑ(´Ó<Ü�~×)Ñ)¨(°BÓ7Ó8¸1Ó<à+×/Ñ/°¸"Ö=�EØ ™/Õ)ò >à Ñ)¬\Ó9Øñ Ù"ÙCùs    ‚B�BžB´BµA#BÂBc              ƒ   óD  #   • U R                    SU R                   SU R                   3nU R                  R	                  XR
                  S9 ISh  v•N nUR                  5       I Sh  v•N sSSS5      ISh  v•N   $  N- N N	! , ISh  v•N  (       d  f       g= f7f)z=Cancels current job

Returns:
    Output of cancel operation
r   z/cancel/r   N)r   r   r   r   Úpostr   r"   )r   Ú
cancel_urlÚresps      r   ÚcancelÚ
Job.canceli   sw   é € ð ×.Ñ.Ð/¨q°×1AÑ1AÐ0BÀ(È4Ï;É;È-ÐXˆ
Ø—<‘<×$Ñ$ Z¿¹Ð$×FÑFÈ$ØŸ™›×$÷ G×FÓFÙ$÷ G×F×FÐFüsZ   ‚AB ÁB ÁB ÁBÁ*BÁ+BÁ.B Á:BÁ;B ÂBÂB ÂBÂBÂBÂB )r   r   r   r   r   r   r   N)r   )r   )Ú__name__Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__Ústrr   Údictr   r   r   r%   r   r-   Úintr    r7   r@   Ú__static_attributes__© r   r   r
   r
      s~   † Ù?ð Cð °ð ¸}ð ÐW[ô ñ. sð ¸$¸sÀC¸x¹.õ ð$
#˜cô 
#ò#ñ, Cð ,°õ ,ð&˜cô ð%˜d÷ %r   r
   c                   óh   • \ rS rSrSr SS\S\S\\   4S jjrS\	S	\
4S
 jrS	\	4S jrS	\	4S jrSrg)ÚEndpointét   zClass for running endpointNr   r   Úapi_keyc                 óä   • SSK JnJn  Xl        X l        XPl        U SU R                   S3U l        U=(       d    UU l        U R                  c  [        S5      eSSU R                   3S	.U l        g)
zÃ
Initialize an async Endpoint instance.

Args:
    endpoint_id: The identifier for the endpoint.
    session: The aiohttp ClientSession.
    api_key: Optional API key for this endpoint instance.
r   )rO   r   r   z/runNz(API key must be provided or set globallyzapplication/jsonzBearer )zContent-TypeÚAuthorization)r   rO   r   r   r   Úendpoint_urlÚRuntimeErrorr   )r   r   r   rO   Úglobal_api_keyr   s         r   r   ÚEndpoint.__init__w   sw   € ÷	
ð
 'ÔØŒØ!2ÔØ0Ð1°°4×3CÑ3CÐ2DÀDÐIˆÔð ×0 .ˆŒà�<‰<ÑÜÐIÓJÐJð /Ø& t§|¡| nÐ5ñ
ˆ�r   Úendpoint_inputr   c              ƒ   óœ  #   • U R                   R                  U R                  U R                  SU0S9 ISh  v•N nUR	                  5       I Sh  v•N nSSS5      ISh  v•N   U R                  R                  5       nWS   US'   [        U R                  US   U R                   U5      $  Ns N] NO! , ISh  v•N  (       d  f       Nd= f7f)zz
Runs endpoint with specified input.

Args:
    endpoint_input: any dictionary with input

Returns:
    Newly created job
Úinput)r   r"   NÚidzX-Request-ID)r   r=   rR   r   r"   Úcopyr
   r   )r   rV   r?   Ú	json_respÚjob_headerss        r   ÚrunÚEndpoint.run–   s³   é € ð —<‘<×$Ñ$Ø×Ñ t§|¡|¸7ÀNÐ:Sð %÷ 
ñ 
àØ"Ÿi™i›k×)ˆI÷
÷ 
ð —l‘l×'Ñ'Ó)ˆØ&/°¡oˆ�NÑ#Ü�4×#Ñ# Y¨t¡_°d·l±lÀKÓPÐPñ
ñ *÷
÷ 
÷ 
ð 
üsW   ‚6C¸B,¹C¼B2ÁB.ÁB2ÁCÁ B0Á!ACÂ.B2Â0CÂ2C	Â8B;Â9C	ÃCc              ƒ   ó,  #   • U R                    SU R                   S3nU R                  R                  XR                  S9 ISh  v•N nUR                  5       I Sh  v•N sSSS5      ISh  v•N   $  N- N N	! , ISh  v•N  (       d  f       g= f7f)z<
Checks health of endpoint

Returns:
    Health of endpoint
r   z/healthr   N)r   r   r   r!   r   r"   )r   Ú
health_urlr?   s      r   ÚhealthÚEndpoint.healthª   so   é € ð ×.Ñ.Ð/¨q°×1AÑ1AÐ0BÀ'ÐJˆ
à—<‘<×#Ñ# J¿¹Ð#×EÑEÈØŸ™›×$÷ F×EÓEÙ$÷ F×E×EÐEüóZ   ‚ABÁA4ÁBÁ
A:ÁA6ÁA:Á"BÁ.A8Á/BÁ6A:Á8BÁ:BÂ BÂBÂBc              ƒ   ó,  #   • U R                    SU R                   S3nU R                  R                  XR                  S9 ISh  v•N nUR                  5       I Sh  v•N sSSS5      ISh  v•N   $  N- N N	! , ISh  v•N  (       d  f       g= f7f)z5
Purges queue of endpoint

Returns:
    Purge status
r   z/purge-queuer   N)r   r   r   r=   r   r"   )r   Ú	purge_urlr?   s      r   Úpurge_queueÚEndpoint.purge_queue¶   so   é € ð ×-Ñ-Ð.¨a°×0@Ñ0@Ð/AÀÐNˆ	à—<‘<×$Ñ$ Y¿¹Ð$×EÑEÈØŸ™›×$÷ F×EÓEÙ$÷ F×E×EÐEürc   )rO   r   rR   r   r   r   )N)rB   rC   rD   rE   rF   rG   r   r   r   rH   r
   r]   ra   rf   rJ   rK   r   r   rM   rM   t   sW   † Ù$ð +/ñ
 Cð 
°-ð 
Ø" 3™-õ
ð>Q¨ð Q°ô Qð(
%˜dô 
%ð
% 4÷ 
%r   rM   )rF   r+   Útypingr   r   r   Úrunpod.endpoint.helpersr   r   Úrunpod.http_clientr   r
   rM   rK   r   r   Ú<module>rk      s2   ðÙ 2ó ß &Ñ &ç >Ý ,÷e%ñ e%÷PL%ò L%r   