ó
    Ú°›jK7  ã                   óÆ   • S r SSKrSSKrSSKrSSKrSSKJrJrJr  SSK	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  \" 5       rS
 rS\S\4S jr " S S5      rg)zh
runpod | serverless | rp_scale.py
Provides the functionality for scaling the runpod serverless worker.
é    N)ÚAnyÚDictÚSeté   )ÚAsyncClientSessionÚClientSessionÚTooManyRequestsé   )Ú_job_stop_urlÚget_jobÚget_stop_signalsÚ
handle_job)ÚRunPodLoggerÚ_reset_batch_idÚ_set_batch_id)ÚJobsProgressÚIS_LOCAL_TESTc                 ób   • [         R                  " XU5      n[        R                  SU 35        g )NzUncaught exception | )Ú	tracebackÚformat_exceptionÚlogÚerror)Úexc_typeÚ	exc_valueÚexc_tracebackÚexcs       Ú_/home/mande/repo/quber/.venv/lib/python3.13/site-packages/runpod/serverless/modules/rp_scale.pyÚ_handle_uncaught_exceptionr      s(   € Ü
×
$Ò
$ X¸-Ó
H€CÜ‡I�IÐ% c UÐ+Õ,ó    Úcurrent_concurrencyÚreturnc                 ó   • U $ )zÓ
Default concurrency modifier.

This function returns the current concurrency without any modification.

Args:
    current_concurrency (int): The current concurrency.

Returns:
    int: The current concurrency.
© )r    s    r   Ú_default_concurrency_modifierr$      s
   € ð Ðr   c                   ó°   • \ rS rSrSrS\\\4   4S jrS r	S r
S rS rS	 rS
 rS\4S jrS\4S jrS\4S jrS\4S jrS\S\4S jrS\S\4S jrSrg)Ú	JobScaleré(   zV
Job Scaler. This class is responsible for scaling the number of concurrent requests.
Úconfigc                 ó  • [         R                  " 5       U l        SU l        Xl        [        5       U l        0 U l        [        U l	        SU l
        [         R                  " U R                  S9U l        [        U l        [        U l        SU l        [$        U l        UR)                  S5      =n(       a  X l        [*        (       d  g U R                  R)                  S5      =n(       a  X0l        U R                  R)                  S5      =n(       a  X@l        U R                  R)                  S5      =n(       a  XPl        U R                  R)                  S5      =n(       a  X`l	        U R                  R)                  S	5      =n(       a  Xpl
        g g )
Nr
   éZ   ©ÚmaxsizeÚconcurrency_modifierÚjobs_fetcherÚjobs_fetcher_timeoutÚjobs_handlerÚstop_signals_fetcherÚstop_signals_fetcher_timeout)ÚasyncioÚEventÚ_shutdown_eventr    r(   r   Újob_progressÚ
jobs_tasksr   r1   r2   ÚQueueÚ
jobs_queuer$   r-   r   r.   r/   r   r0   Úgetr   )Úselfr(   r-   r.   r/   r0   r1   r2   s           r   Ú__init__ÚJobScaler.__init__-   s<  € Ü&Ÿ}š}›ˆÔØ#$ˆÔ ØŒÜ(›NˆÔð 46ˆŒä$4ˆÔ!Ø,.ˆÔ)ä!Ÿ-š-°×0HÑ0HÑIˆŒä$AˆÔ!Ü#ˆÔØ$&ˆÔ!Ü&ˆÔà#)§:¡:Ð.DÓ#EÐEÐÕEØ(<Ô%çŠ}ààŸ;™;Ÿ?™?¨>Ó:Ð:ˆ<Õ:Ø ,Ôà#'§;¡;§?¡?Ð3IÓ#JÐJÐÕJØ(<Ô%àŸ;™;Ÿ?™?¨>Ó:Ð:ˆ<Õ:Ø ,Ôà#'§;¡;§?¡?Ð3IÓ#JÐJÐÕJØ(<Ô%à+/¯;©;¯?©?Ð;YÓ+ZÐZÐ'ÕZØ0LÕ-ð [r   c              ƒ   ó®  #   • U R                  U R                  5      U l        U R                  (       a%  U R                  U R                  R                  :X  a  g U R	                  5       S:”  a   [
        R                  " S5      I S h  v•N   M4  [
        R                  " U R                  S9U l        [        R                  SU R                   35        g  NL7f)Nr   r
   r+   z.JobScaler.set_scale | New concurrency set to: )
r-   r    r9   r,   Úcurrent_occupancyr3   Úsleepr8   r   Údebug©r;   s    r   Ú	set_scaleÚJobScaler.set_scaleW   s�   é € Ø#'×#<Ñ#<¸T×=UÑ=UÓ#VˆÔ à�?�? × 8Ñ 8¸D¿O¹O×<SÑ<SÓ Sàà×$Ñ$Ó&¨Ó*ä—-’- Ó"×"Ð"Ùä!Ÿ-š-°×0HÑ0HÑIˆŒÜ�	‰	Ø<¸T×=UÑ=UÐ<VÐWõ	
ñ	 #ùs   ‚BCÂCÂACc                 ór  • [         [        l         [        R                  " [        R                  U R
                  5        [        R                  " [        R                  U R
                  5        [        R                  " U R                  5       5        g! [         a    [        R                  S5         NFf = f)zº
This is required for the worker to be able to shut down gracefully
when the user sends a SIGTERM or SIGINT signal. This is typically
the case when the worker is running in a container.
z5Signal handling is only supported in the main thread.N)r   ÚsysÚ
excepthookÚsignalÚSIGTERMÚhandle_shutdownÚSIGINTÚ
ValueErrorr   Úwarnr3   ÚrunrB   s    r   ÚstartÚJobScaler.starth   su   € ô 4ŒŒð	Nä�MŠMœ&Ÿ.™.¨$×*>Ñ*>Ô?Ü�MŠMœ&Ÿ-™-¨×)=Ñ)=Ô>ô 	�Š�D—H‘H“JÕøô ó 	NÜ�H‰HÐLÖMð	Nús   ‘AB ÂB6Â5B6c                 óV   • [         R                  SU S35        U R                  5         g)a[  
Called when the worker is signalled to shut down.

This function is called when the worker receives a signal to shut down, such as
SIGTERM or SIGINT. It sets the shutdown event, which will cause the worker to
exit its main loop and shut down gracefully.

Args:
    signum: The signal number that was received.
    frame: The current stack frame.
zReceived shutdown signal: Ú.N)r   rA   Úkill_worker)r;   ÚsignumÚframes      r   rJ   ÚJobScaler.handle_shutdown{   s&   € ô 	�	‰	Ð.¨v¨h°aÐ8Ô9Ø×ÑÕr   c              ƒ   ó®  #   • [        5        IS h  v•N n[        R                  " U R                  U5      5      n[        R                  " U R	                  U5      5      n[        R                  " U R                  U5      5      nX#U/n[        R                  " U6 I S h  v•N   S S S 5      IS h  v•N   g  N£ N N	! , IS h  v•N  (       d  f       g = f7f)N)r   r3   Úcreate_taskÚget_jobsÚrun_jobsÚmonitor_stop_signalsÚgather)r;   ÚsessionÚjobtake_taskÚjobrun_taskÚjobstop_taskÚtaskss         r   rN   ÚJobScaler.runŠ   s™   é € ä%×'Ô'¨7ä"×.Ò.¨t¯}©}¸WÓ/EÓFˆLÜ!×-Ò-¨d¯m©m¸GÓ.DÓEˆKÜ"×.Ò.¨t×/HÑ/HÈÓ/QÓRˆLà!°Ð=ˆEô —.’. %Ð(×(Ð(÷ (×'Ò'ñ )÷ (×'×'Ð'üsW   ‚C‘B5’C•B
B;ÂB7Â B;Â$CÂ/B9Â0CÂ7B;Â9CÂ;CÃCÃCÃCc                 ó@   • U R                   R                  5       (       + $ )z,
Return whether the worker is alive or not.
)r5   Úis_setrB   s    r   Úis_aliveÚJobScaler.is_alive—   s   € ð ×'Ñ'×.Ñ.Ó0Ô0Ð0r   c                 ób   • [         R                  S5        U R                  R                  5         g)z
Whether to kill the worker.
zKill worker.N)r   rA   r5   ÚsetrB   s    r   rS   ÚJobScaler.kill_worker�   s"   € ô 	�	‰	�.Ô!Ø×Ñ× Ñ Õ"r   r!   c                 óÂ   • U R                   R                  5       nU R                  R                  5       n[        R                  SU R                   SU SU 35        X!-   $ )Nz JobScaler.status | concurrency: z	; queue: z; progress: )r9   Úqsizer6   Úget_job_countr   rA   r    )r;   Úcurrent_queue_countÚcurrent_progress_counts      r   r?   ÚJobScaler.current_occupancy¤   sn   € Ø"Ÿo™o×3Ñ3Ó5ÐØ!%×!2Ñ!2×!@Ñ!@Ó!BÐä�	‰	Ø.¨t×/GÑ/GÐ.HÈ	ÐReÐQfÐfrð  tJð  sKð  Lô	
ð &Ñ;Ð;r   r]   c           	   ƒ   ó:  #   • U R                  5       (       GaÁ  U R                  5       I Sh  v•N   U R                  U R                  5       -
  nUS::  a5  [        R                  S5        [        R                  " S5      I Sh  v•N   M†   [        R                  S5        [        R                  " U R                  X5      U R                  S9I Sh  v•N nU(       d7  [        R                  S5         [        R                  " S5      I Sh  v•N   GM  U HZ  nU R                  R                  U5      I Sh  v•N   U R                  R                  U5        [        R                  SUS	   5        M\     [        R                  S
U R                  R!                  5        35        [        R                  " S5      I Sh  v•N   U R                  5       (       a  GMÀ  gg GN¯ GN\ GN NØ N¯! ["         a7    [        R                  S5        [        R                  " S5      I Sh  v•N     N…[        R$                   a    [        R                  S5        e [        R&                   a    [        R                  S5         NØ[(         a#  n[        R                  SU S35         SnANÿSnAf[*         aB  n[        R-                  S[/        U5      R0                   S[3        U5       35         SnAGNISnAff = f GN8! [        R                  " S5      I Sh  v•N    f = f7f)z§
Retrieve multiple jobs from the server in batches using blocking requests.

Runs the block in an infinite loop while the worker is alive.

Adds jobs to the JobsQueue
Nr   z2JobScaler.get_jobs | Queue is full. Retrying soon.r
   z.JobScaler.get_jobs | Starting job acquisition.©Útimeoutz&JobScaler.get_jobs | No jobs acquired.z
Job QueuedÚidzJobs in queue: z?JobScaler.get_jobs | Too many requests. Debounce for 5 seconds.é   z+JobScaler.get_jobs | Request was cancelled.z9JobScaler.get_jobs | Job acquisition timed out. Retrying.z'JobScaler.get_jobs | Unexpected error: rR   z!Failed to get job. | Error Type: ú | Error Message: )re   rC   r    r?   r   rA   r3   r@   Úwait_forr.   r/   r9   Úputr6   ÚaddÚinfork   r	   ÚCancelledErrorÚTimeoutErrorÚ	TypeErrorÚ	Exceptionr   ÚtypeÚ__name__Ústr)r;   r]   Újobs_neededÚacquired_jobsÚjobr   s         r   rY   ÚJobScaler.get_jobs­   s[  é € ð �m‰m�oŠoØ—.‘.Ó"×"Ð"à×2Ñ2°T×5KÑ5KÓ5MÑMˆKØ˜aÓÜ—	‘	ÐNÔOÜ—m’m AÓ&×&Ð&Ùð&'Ü—	‘	ÐJÔKô '.×&6Ò&6Ø×%Ñ% gÓ;Ø ×5Ñ5ñ'÷ !�ö
 %Ü—I‘IÐFÔGØô6 —m’m AÓ&×&Ó&ó3 )�CØŸ/™/×-Ñ-¨cÓ2×2Ð2Ø×%Ñ%×)Ñ)¨#Ô.Ü—I‘I˜l¨C°©IÖ6ñ )ô
 —‘˜?¨4¯?©?×+@Ñ+@Ó+BÐ*CÐDÔEô( —m’m AÓ&×&Ð&ð_ �m‰m�oŽoÚ"ò
 'ò!ñD 'ñ1 3øô #ó 'Ü—	‘	ØUôô —m’m AÓ&×&Ó&Ü×)Ñ)ó Ü—	‘	ÐGÔHØÜ×'Ñ'ó WÜ—	‘	ÐUÖVÜó NÜ—	‘	ÐCÀEÀ7È!ÐL×MÑMûÜó Ü—	‘	Ø7¼¸U»×8LÑ8LÐ7MÐM_Ô`cÐdiÓ`jÐ_kÐl÷ò ûðúò 'ø”g—m’m AÓ&×&Ò&üs÷   ‚*L¬G­ALÂGÂLÂ
AG( ÃG!Ã G( Ã2LÄG$ÄLÄ#G( Ä6G&Ä7A,G( Æ#LÆ<K4Æ=LÇLÇLÇ!G( Ç$LÇ&G( Ç(8K1È H#È!K1È&K7 È(AK1É9K7 É;	K1ÊJ"ÊK7 Ê"K1Ê/7K,Ë&K7 Ë,K1Ë1K7 Ë4LË7LÌLÌLÌLc              ƒ   óÖ  #   • [        5       nSnU R                  5       (       d   U R                  R                  5       (       Gd¤  [	        U5      U R
                  :  aÂ  U R                  R                  5       (       d£  U R                  R                  5       I Sh  v•N n[        R                  " U R                  X5      5      nUR                  U5        XPR                  US   '   [	        U5      U R
                  :  a!  U R                  R                  5       (       d  M£  U(       aj  [	        U5      nXc:w  a  [        R                  SU 35        Un[        R                  " U[        R                  SS9I Sh  v•N u  pxUR!                  U5        O[        R"                  " S5      I Sh  v•N   U R                  5       (       a  GM‚  U R                  R                  5       (       d  GM¤  [        R$                  " USS06I Sh  v•N n	U	 HS  n
['        U
[(        5      (       d  M  ['        U
[        R*                  5      (       a  M;  [        R-                  S	U
 35        MU     g GNÇ Në N» Ne7f)
zœ
Retrieve jobs from the jobs queue and process them concurrently.

Runs the block in an infinite loop while the worker is alive or jobs queue is not empty.
r   Nrs   zJobs in progress: gš™™™™™¹?)Úreturn_whenrr   Úreturn_exceptionsTz8JobScaler.run_jobs | Task failed during shutdown drain: )rh   re   r9   ÚemptyÚlenr    r:   r3   rX   r   rx   r7   r   ry   ÚwaitÚFIRST_COMPLETEDÚdifference_updater@   r\   Ú
isinstancer}   rz   r   )r;   r]   ra   Úlast_task_countrƒ   ÚtaskÚcurrent_task_countÚdoneÚpendingÚresultsÚresults              r   rZ   ÚJobScaler.run_jobsæ   sÄ  é € ô $'£5ˆàˆØ�m‰m�o‰o T§_¡_×%:Ñ%:×%<Ò%<ä�e“*˜t×7Ñ7Ó7ÀÇÁ×@UÑ@U×@WÑ@WØ ŸO™O×/Ñ/Ó1×1�ä×*Ò*¨4¯?©?¸7Ó+HÓI�Ø—	‘	˜$”Ø-1—‘  D¡	Ñ*ô �e“*˜t×7Ñ7Ó7ÀÇÁ×@UÑ@U×@WÓ@Wö Ü%(¨£ZÐ"Ø%Ó8Ü—H‘HÐ1Ð2DÐ1EÐFÔGØ&8�Oä&-§l¢lØ¤w×'>Ñ'>Èñ'÷ !‘�ð
 ×'Ñ'¨Õ-ô —m’m CÓ(×(Ð(ð3 �m‰m�oŒo T§_¡_×%:Ñ%:×%<Ô%<ô>  Ÿš¨ÐFÀÑF×FˆÛˆFÜ˜&¤)×,Ó,´ZÀÌ×H^ÑH^×5_Ó5_Ü—	‘	ÐTÐU[ÐT\Ð]Ö^ò ò; 2ñ!ñ )ñ Gùsh   ‚BI)ÂI ÂBI)ÄAI)Å7I#Å81I)Æ)I%Æ*I)ÇI)Ç(I)ÈI'ÈI)È"I)ÉI)É#I)É%I)É'I)c           	   ƒ   ó6  #   • U R                   [        L a!  [        5       c  [        R	                  S5        gU R                  5       (       a´   [        R                  " U R                  U5      U R                  S9I Sh  v•N nU H  nU R                  U5      I Sh  v•N   M     U(       d  [        R                  " S5      I Sh  v•N   [        R                  " S
5      I Sh  v•N   U R                  5       (       a  M³  gg N‚ Nf N?! [         a"    [        R                  " S5      I Sh  v•N     Nh[        R                   a    [        R                  S5        e [        R                   a    [        R                  S5         N»[         aa  n[        R!                  S[#        U5      R$                   S	['        U5       35        [        R                  " S5      I Sh  v•N     SnAGN SnAff = f GN! [        R                  " S
5      I Sh  v•N    f = f7f)a0  
Long-polls the dedicated stop channel and stops signalled jobs.

Runs in an infinite loop while the worker is alive. The Runpod server
signals a request to be stopped (for example when it is cancelled or
times out) and this loop stops just that in-progress job, leaving the
worker's other jobs running.
NzƒJobScaler.monitor_stop_signals | Stop channel could not be derived from the job-take URL; per-job stop is disabled for this worker.rq   r
   rt   z7JobScaler.monitor_stop_signals | Request was cancelled.z?JobScaler.monitor_stop_signals | Stop poll timed out. Retrying.z-JobScaler.monitor_stop_signals | Error Type: ru   r   )r1   r   r   r   rM   re   r3   rv   r2   Ústop_jobr@   r	   rz   rA   r{   r}   r   r~   r   r€   )r;   r]   Újob_idsÚjob_idr   s        r   r[   ÚJobScaler.monitor_stop_signals  s˜  é € ð ×$Ñ$Ô(8Ò8¼]»_Ñ=TÜ�H‰HðSôð à�m‰m�o‰oð'ô !(× 0Ò 0Ø×-Ñ-¨gÓ6Ø ×=Ñ=ñ!÷ �ó &�FØŸ-™-¨Ó/×/Ò/ñ &ö ô "Ÿ-š-¨Ó*×*Ð*ô —m’m AÓ&×&Ð&ð9 �m‰m�o�oññ
 0ñ
 +øÜ"ó 'Ü—m’m AÓ&×&Ó&Ü×)Ñ)ó Ü—	‘	ÐSÔTØÜ×'Ñ'ó ]Ü—	‘	Ð[Ö\Üó 'Ü—	‘	ØCÄDÈÃK×DXÑDXÐCYÐYkÔloÐpuÓlvÐkwÐxôô —m’m AÓ&×&×&ûð	'úò 'ø”g—m’m AÓ&×&Ò&üsÈ   ‚A
HÁ1D Á>DÁ?D ÂDÂ(D ÃDÃD Ã
HÃ#G2Ã$HÃ?HÄD ÄD ÄD Ä#G/Ä*D-Ä+G/Ä0G5 Ä2AG/ÆG5 Æ	G/ÆAG*ÇG!ÇG*Ç$G5 Ç*G/Ç/G5 Ç2HÇ5HÈHÈHÈHr™   c              ƒ   óÈ   #   • U R                   R                  U5      nUc  [        R                  SU S35        g[        R	                  SU5        UR                  5         g7f)zÊ
Stop a single in-progress job by cancelling its running task.

Args:
    job_id: The id of the job to stop.

Returns:
    True if a matching in-progress job was found and stopped,
    False otherwise.
z,JobScaler.stop_job | No in-progress job for rR   FzStopping job.T)r7   r:   r   rA   ry   Úcancel)r;   r™   r�   s      r   r—   ÚJobScaler.stop_jobA  sS   é € ð �‰×"Ñ" 6Ó*ˆØ‰<Ü�I‰IÐDÀVÀHÈAÐNÔOØä�‰� &Ô)Ø�‰ŒØùs   ‚A A"rƒ   c              ƒ   óÄ  #   • [        UR                  S5      5      n [        R                  SUS   5        U R	                  XR
                  U5      I Sh  v•N   U R
                  R                  SS5      (       a  U R                  5         U R                  R                  5         U R                  R                  U5        U R                   R#                  US   S5        [        R                  S	US   5        [%        U5        g N®! [        R                   a    [        R                  SUS   5        e [         a"  n[        R                  SU 3US   5        e SnAff = f! U R                  R                  5         U R                  R                  U5        U R                   R#                  US   S5        [        R                  S	US   5        [%        U5        f = f7f)
zQ
Process an individual job. This function is run concurrently for multiple jobs.
ÚbatchIdzHandling Jobrs   NÚrefresh_workerFzJob stopped.zError handling job: zFinished Job)r   r:   r   rA   r0   r(   rS   r3   rz   ry   r}   r   r9   Ú	task_doner6   Úremover7   Úpopr   )r;   r]   rƒ   Úbatch_id_tokenÚerrs        r   r   ÚJobScaler.handle_jobU  sx  é € ô ' s§w¡w¨yÓ'9Ó:ˆð	,Ü�I‰I�n c¨$¡iÔ0à×#Ñ# G¯[©[¸#Ó>×>Ð>à�{‰{�‰Ð/°×7Ñ7Ø× Ñ Ô"ð �O‰O×%Ñ%Ô'ð ×Ñ×$Ñ$ SÔ)Ø�O‰O×Ñ  D¡	¨4Ô0ä�I‰I�n c¨$¡iÔ0Ü˜NÕ+ñ- ?øô
 ×%Ñ%ó 	Ü�H‰H�^ S¨¡YÔ/Øäó 	Ü�I‰IÐ,¨S¨EÐ2°C¸±IÔ>Øûð	ûð �O‰O×%Ñ%Ô'ð ×Ñ×$Ñ$ SÔ)Ø�O‰O×Ñ  D¡	¨4Ô0ä�I‰I�n c¨$¡iÔ0Ü˜NÕ+üsH   ‚G ž8D ÁDÁ5D ÂA9G ÄD Ä7E Ä>EÅE Å E# Å#A:GÇG )r5   r-   r(   r    r6   r.   r/   r0   r9   r7   r1   r2   N)r   Ú
__module__Ú__qualname__Ú__firstlineno__Ú__doc__r   r€   r   r<   rC   rO   rJ   rN   re   rS   Úintr?   r   rY   rZ   r[   Úboolr—   Údictr   Ú__static_attributes__r#   r   r   r&   r&   (   s—   † ñð(M˜t C¨ H™~ô (MòT
ò" ò&ò)ò1ò#ð< 3ô <ð7' mô 7'ðr+_ mô +_ðZ,'°-ô ,'ð\ Sð ¨Tô ð(,¨ð ,¸D÷ ,r   r&   )rª   r3   rH   rF   r   Útypingr   r   r   Úhttp_clientr   r   r	   Úrp_jobr   r   r   r   Ú	rp_loggerr   r   r   Úworker_stater   r   r   r   r«   r$   r&   r#   r   r   Ú<module>r´      s^   ðñó
 Û Û 
Û ß !Ñ !ç MÑ Mß HÓ Hß CÑ Cß 5áƒn€ò-ð
°sð ¸sô ÷K,ò K,r   