Ë
    ÚŠ±j�£  ã                   óØ   — d Z ddlZddlZddlZddlZddlZddl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 ddlmZmZmZmZmZ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#  G d	„ d
«      Z$y)z“
Dirty Arbiter Process

Asyncio-based arbiter that manages the dirty worker pool and routes
requests from HTTP workers to available dirty workers.
é    N)Úutilé   )Úget_app_workers_attributeÚparse_dirty_app_spec)Ú
DirtyErrorÚDirtyNoWorkersAvailableErrorÚDirtyTimeoutErrorÚDirtyWorkerError)ÚDirtyProtocolÚmake_error_responseÚmake_responseÚSTASH_OP_PUTÚSTASH_OP_GETÚSTASH_OP_DELETEÚSTASH_OP_KEYSÚSTASH_OP_CLEARÚSTASH_OP_INFOÚSTASH_OP_ENSUREÚSTASH_OP_DELETE_TABLEÚSTASH_OP_TABLESÚSTASH_OP_EXISTSÚMANAGE_OP_ADDÚMANAGE_OP_REMOVE)ÚDirtyWorkerc                   óV  — e Zd ZdZdj	                  «       D � ���cg c]  }t        t        d|z  «      ‘Œ c}}}} ZdZd&d„Z	d„ Z
d„ Zd	„ Zd
„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Zd'd„Zd„ Zd„ Zd„ Zd„ Zd„ Zd„ Z d(d„Z!d„ Z"d „ Z#d!„ Z$d"„ Z%d#„ Z&d)d$„Z'd%„ Z(yc c}}}} w )*ÚDirtyArbitera4  
    Dirty arbiter that manages the dirty worker pool.

    The arbiter runs an asyncio event loop and handles:
    - Spawning and managing dirty worker processes
    - Accepting connections from HTTP workers
    - Routing requests to available dirty workers
    - Monitoring worker health via heartbeat
    z*HUP QUIT INT TERM TTIN TTOU USR1 USR2 CHLDzSIG%sé   Nc                 óN  — || _         || _        d| _        t        j                  «       | _        || _        t        j                  d¬«      | _	        |xs* t        j                  j                  | j                  d«      | _        i | _        i | _        i | _        i | _        i | _        d| _        d| _        d| _        | j                   j*                  | _        d| _        d| _        i | _        i | _        i | _        i | _        i | _        g | _        i | _        | jA                  «        y)zù
        Initialize the dirty arbiter.

        Args:
            cfg: Gunicorn config
            log: Logger
            socket_path: Path to the arbiter's Unix socket
            pidfile: Well-known PID file location for orphan detection
        Nzgunicorn-dirty-)Úprefixzarbiter.sockr   T)!ÚcfgÚlogÚpidÚosÚgetpidÚppidÚpidfileÚtempfileÚmkdtempÚtmpdirÚpathÚjoinÚsocket_pathÚworkersÚworker_socketsÚworker_connectionsÚworker_queuesÚworker_consumersÚ_worker_rr_indexÚ
worker_ageÚaliveÚdirty_workersÚnum_workersÚ_serverÚ_loopÚ_pending_requestsÚ	app_specsÚapp_worker_mapÚworker_app_mapÚ_app_rr_indicesÚ_pending_respawnsÚstash_tablesÚ_parse_app_specs)Úselfr    r!   r,   r&   s        ú[/var/www/io.vulcan-creative.com/venv/lib/python3.12/site-packages/gunicorn/dirty/arbiter.pyÚ__init__zDirtyArbiter.__init__B   s  € ð ˆŒØˆŒØˆŒÜ—I‘I“KˆŒ	ØˆŒô ×&Ñ&Ð.?Ô@ˆŒØ&ò 
¬"¯'©'¯,©,Ø�K‰K˜ó+
ˆÔð ˆŒØ ˆÔØ"$ˆÔØˆÔØ "ˆÔØ !ˆÔØˆŒØˆŒ
ØŸ8™8×1Ñ1ˆÔàˆŒØˆŒ
Ø!#ˆÔð ˆŒà ˆÔà ˆÔà!ˆÔà!#ˆÔð ˆÔð 	×ÑÕó    c                 ó,  — | j                   j                  D ]H  }t        |«      \  }}|€	 t        |«      }|||dœ| j                  |<   t        «       | j                  |<   ŒJ y# t        $ r'}| j
                  j                  d||«       Y d}~ŒXd}~ww xY w)a‹  
        Parse all app specifications from config.

        Populates self.app_specs with parsed information about each app,
        including the import path and worker count limits.

        Worker count priority:
        1. Config override (e.g., "module:Class:2") - highest priority
        2. Class attribute (e.g., workers = 2 on the class)
        3. None (all workers) - default
        Nz,Could not read workers attribute from %s: %s)Úimport_pathÚworker_countÚoriginal_spec)
r    Ú
dirty_appsr   r   Ú	Exceptionr!   Úwarningr:   Úsetr;   )rA   ÚspecrF   rG   Úes        rB   r@   zDirtyArbiter._parse_app_specsy   s£   € ð —H‘H×'Ñ'ò 	5ˆDÜ(<¸TÓ(BÑ%ˆK˜ð Ð#ðÜ#<¸[Ó#I�Lð  +Ø ,Ø!%ñ+ˆD�N‰N˜;Ñ'ô 03«uˆD×Ñ Ò,ñ)	5øô !ò à—H‘H×$Ñ$ØFØ# Q÷ñ ûðús   «A#Á#	BÁ,BÂBc                 óp   — d}| j                   j                  «       D ]  }|d   }|€Œt        ||«      }Œ |S )a  
        Calculate minimum number of workers required by app specs.

        Returns the maximum worker_count across all apps that have limits.
        Apps with worker_count=None don't impose a minimum.

        Returns:
            int: Minimum workers required (at least 1)
        r   rG   )r:   ÚvaluesÚmax)rA   Úmin_requiredrM   rG   s       rB   Ú_get_minimum_workersz!DirtyArbiter._get_minimum_workers›   sI   € ð ˆØ—N‘N×)Ñ)Ó+ò 	?ˆDØ Ñ/ˆLØÑ'Ü" <°Ó>‘ð	?ð ÐrD   c                 ó  — g }| j                   j                  «       D ]b  \  }}|d   }t        | j                  j	                  |t        «       «      «      }|€|j                  |«       ŒL||k  sŒR|j                  |«       Œd |S )a�  
        Determine which apps a new worker should load.

        Returns a list of import paths for apps that need more workers.
        Apps with workers=None (all workers) are always included.
        Apps with worker limits are included only if they haven't
        reached their limit yet.

        Returns:
            List of import paths to load, or empty list if no apps need workers
        rG   )r:   ÚitemsÚlenr;   ÚgetrL   Úappend)rA   Ú	app_pathsrF   rM   rG   Úcurrent_workerss         rB   Ú_get_apps_for_new_workerz%DirtyArbiter._get_apps_for_new_worker¬   s‡   € ð ˆ	à!%§¡×!5Ñ!5Ó!7ò 		.ÑˆK˜Ø Ñ/ˆLÜ! $×"5Ñ"5×"9Ñ"9¸+ÄsÃuÓ"MÓNˆOð Ð#Ø× Ñ  Õ-à  <Ó/Ø× Ñ  Õ-ð		.ð ÐrD   c                 óÈ   — t        |«      | j                  |<   |D ]E  }|| j                  vrt        «       | j                  |<   | j                  |   j	                  |«       ŒG y)a?  
        Register which apps a worker has loaded.

        Updates both app_worker_map and worker_app_map to track the
        bidirectional relationship between workers and apps.

        Args:
            worker_pid: The PID of the worker
            app_paths: List of app import paths loaded by this worker
        N)Úlistr<   r;   rL   Úadd©rA   Ú
worker_pidrY   Úapp_paths       rB   Ú_register_worker_appsz"DirtyArbiter._register_worker_appsÇ   s`   € ô +/¨y«/ˆ×Ñ˜JÑ'à!ò 	:ˆHØ˜t×2Ñ2Ñ2Ü03³�×#Ñ# HÑ-Ø×Ñ Ñ)×-Ñ-¨jÕ9ñ	:rD   c                 ó¤   — | j                   j                  |g «      }|D ]/  }|| j                  v sŒ| j                  |   j                  |«       Œ1 y)zº
        Unregister a worker's apps when it exits.

        Removes the worker from all tracking maps.

        Args:
            worker_pid: The PID of the worker to unregister
        N)r<   Úpopr;   Údiscardr_   s       rB   Ú_unregister_workerzDirtyArbiter._unregister_workerÙ   sV   € ð ×'Ñ'×+Ñ+¨J¸Ó;ˆ	ð "ò 	BˆHØ˜4×.Ñ.Ò.Ø×#Ñ# HÑ-×5Ñ5°jÕAñ	BrD   c                 ó  — t        j                  «       | _        | j                  j	                  d| j                  «       | j
                  rD	 t        | j
                  d«      5 }|j                  t        | j                  «      «       ddd«       | j                  t         j                  d<   | j                  j                  | «       | j                  «        t!        j"                  d«       	 t%        j&                  | j)                  «       «       | j-                  «        y# 1 sw Y   Œ›xY w# t        $ r&}| j                  j                  d|«       Y d}~ŒÈd}~ww xY w# t*        $ r Y ŒZw xY w# | j-                  «        w xY w)z&Run the dirty arbiter (blocking call).z Dirty arbiter starting (pid: %s)ÚwNzFailed to write PID file: %sÚGUNICORN_DIRTY_SOCKETzdirty-arbiter)r#   r$   r"   r!   Úinfor&   ÚopenÚwriteÚstrÚIOErrorrK   r,   Úenvironr    Úon_dirty_startingÚinit_signalsr   Ú_setproctitleÚasyncioÚrunÚ
_run_asyncÚKeyboardInterruptÚ_cleanup_sync)rA   ÚfrN   s      rB   rt   zDirtyArbiter.runê   s(  € ä—9‘9“;ˆŒØ�‰�‰Ð8¸$¿(¹(ÔCð �<Š<ðDÜ˜$Ÿ,™,¨Ó,ð +°Ø—G‘GœC §¡›MÔ*÷+ð /3×.>Ñ.>Œ�
‰
Ð*Ñ+ð 	�‰×"Ñ" 4Ô(ð 	×ÑÔô 	×Ñ˜?Ô+ð	!Ü�K‰K˜Ÿ™Ó)Ô*ð ×ÑÕ ÷-+ð +ûäò DØ—‘× Ñ Ð!?À×CÑCûðDûô" !ò 	Ùð	ûð ×ÑÕ úsT   ÁD. Á#%D"ÂD. Ã.#E  Ä"D+Ä'D. Ä.	EÄ7EÅEÅ 	E,Å)E/ Å+E,Å,E/ Å/Fc                 óN  — | j                   D ]&  }t        j                  |t        j                  «       Œ( t        j                  t        j                  | j                  «       t        j                  t        j
                  | j                  «       t        j                  t        j                  | j                  «       t        j                  t        j                  | j                  «       t        j                  t        j                  | j                  «       t        j                  t        j                  | j                  «       t        j                  t        j                  | j                  «       t        j                  t        j                  | j                  «       y)zSet up signal handlers.N)ÚSIGNALSÚsignalÚSIG_DFLÚSIGTERMÚ_signal_handlerÚSIGQUITÚSIGINTÚSIGHUPÚSIGUSR1ÚSIGCHLDÚSIGTTINÚSIGTTOU)rA   Úsigs     rB   rq   zDirtyArbiter.init_signals
  sé   € à—<‘<ò 	/ˆCÜ�M‰M˜#œvŸ~™~Õ.ð	/ô 	�‰”f—n‘n d×&:Ñ&:Ô;Ü�‰”f—n‘n d×&:Ñ&:Ô;Ü�‰”f—m‘m T×%9Ñ%9Ô:Ü�‰”f—m‘m T×%9Ñ%9Ô:Ü�‰”f—n‘n d×&:Ñ&:Ô;Ü�‰”f—n‘n d×&:Ñ&:Ô;Ü�‰”f—n‘n d×&:Ñ&:Ô;Ü�‰”f—n‘n d×&:Ñ&:Õ;rD   c                 ó(  ‡ — |t         j                  k(  r+‰ j                  r‰ j                  j                  ˆ fd„«       y|t         j                  k(  r‰ j
                  j                  «        y|t         j                  k(  r+‰ j                  r‰ j                  j                  ˆ fd„«       y|t         j                  k(  rf‰ xj                  dz  c_	        ‰ j
                  j                  d‰ j                  «       ‰ j                  r‰ j                  j                  ˆ fd„«       y|t         j                  k(  r¢‰ j                  «       }‰ j                  |k  r‰ j
                  j                  d|«       y‰ xj                  dz  c_	        ‰ j
                  j                  d‰ j                  «       ‰ j                  r‰ j                  j                  ˆ fd	„«       yd
‰ _        ‰ j                  r&‰ j                  j                  ‰ j                  «       yy)zHandle signals.c                  óJ   •— t        j                  ‰ j                  «       «      S ©N)rs   Úcreate_taskÚ_handle_sigchld©rA   s   €rB   ú<lambda>z.DirtyArbiter._signal_handler.<locals>.<lambda>  s   ø€ œG×/Ñ/°×0DÑ0DÓ0FÓG€ rD   Nc                  óJ   •— t        j                  ‰ j                  «       «      S r‰   )rs   rŠ   ÚreloadrŒ   s   €rB   r�   z.DirtyArbiter._signal_handler.<locals>.<lambda>+  s   ø€ œG×/Ñ/°·±³Ó>€ rD   r   z'SIGTTIN: Increasing dirty workers to %sc                  óJ   •— t        j                  ‰ j                  «       «      S r‰   ©rs   rŠ   Úmanage_workersrŒ   s   €rB   r�   z.DirtyArbiter._signal_handler.<locals>.<lambda>6  ó   ø€ œG×/Ñ/°×0CÑ0CÓ0EÓF€ rD   zASIGTTOU: Cannot decrease below %s workers (required by app specs)z'SIGTTOU: Decreasing dirty workers to %sc                  óJ   •— t        j                  ‰ j                  «       «      S r‰   r‘   rŒ   s   €rB   r�   z.DirtyArbiter._signal_handler.<locals>.<lambda>I  r“   rD   F)r{   rƒ   r8   Úcall_soon_threadsafer‚   r!   Úreopen_filesr�   r„   r6   rj   r…   rS   rK   r4   Ú	_shutdown)rA   r†   ÚframeÚmin_workerss   `   rB   r~   zDirtyArbiter._signal_handler  s–  ø€ à”&—.‘.Ò à�zŠzØ—
‘
×/Ñ/ÛGôð à”&—.‘.Ò à�H‰H×!Ñ!Ô#Øà”&—-‘-Òà�zŠzØ—
‘
×/Ñ/Û>ôð à”&—.‘.Ò à×Ò Ñ!ÕØ�H‰H�M‰MÐCØ×*Ñ*ô,à�zŠzØ—
‘
×/Ñ/ÛFôð à”&—.‘.Ò à×3Ñ3Ó5ˆKØ×Ñ ;Ò.Ø—‘× Ñ ð.àôð
 Ø×Ò Ñ!ÕØ�H‰H�M‰MÐCØ×*Ñ*ô,à�zŠzØ—
‘
×/Ñ/ÛFôð ð ˆŒ
Ø�:Š:Ø�J‰J×+Ñ+¨D¯N©NÕ;ð rD   c                 óR   — | j                   r| j                   j                  «        yy)zInitiate async shutdown.N)r7   ÚcloserŒ   s    rB   r—   zDirtyArbiter._shutdownR  s   € à�<Š<Ø�L‰L×ÑÕ ð rD   c              ƒ   ó–  K  — t        j                  «       | _        t        j                  j                  | j                  «      rt        j                  | j                  «       t        j                  | j                  | j                  ¬«      ƒ d{  –—† | _
        t        j                  | j                  d«       | j                  j                  d| j                  «       | j                  «       ƒ d{  –—†  t        j                  | j!                  «       «      }	 | j                  4 ƒd{  –—†  | j                  j#                  «       ƒ d{  –—†  ddd«      ƒd{  –—†  |j)                  «        	 |ƒ d{  –—†  | j+                  «       ƒ d{  –—†  y7 �Œ7 Œ¦7 Œo7 ŒO7 ŒA# 1 ƒd{  –—†7  sw Y   ŒQxY w# t         j$                  t&        f$ r Y Œow xY w7 Œ\# t         j$                  $ r Y Œow xY w7 Œ_# |j)                  «        	 |ƒ d{  –—†7   n# t         j$                  $ r Y nw xY w| j+                  «       ƒ d{  –—†7   w xY w­w)z/Main async loop - start server, manage workers.)r*   Ni€  zDirty arbiter listening on %s)rs   Úget_running_loopr8   r#   r*   Úexistsr,   ÚunlinkÚstart_unix_serverÚhandle_clientr7   Úchmodr!   rj   r’   rŠ   Ú_worker_monitorÚserve_foreverÚCancelledErrorÚRuntimeErrorÚcancelÚstop)rA   Úmonitor_tasks     rB   ru   zDirtyArbiter._run_asyncW  sÜ  è ø€ ä×-Ñ-Ó/ˆŒ
ô �7‰7�>‰>˜$×*Ñ*Ô+Ü�I‰I�d×&Ñ&Ô'ô %×6Ñ6Ø×ÑØ×!Ñ!ô
÷ 
ˆŒô 	�‰�×!Ñ! 5Ô)à�‰�‰Ð5°t×7GÑ7GÔHð ×!Ñ!Ó#×#Ð#ô ×*Ñ*¨4×+?Ñ+?Ó+AÓBˆð	Ø—|‘|÷ 3ñ 3Ø—l‘l×0Ñ0Ó2×2Ð2÷3÷ 3ð ×ÑÔ!ðØ"×"Ð"ð —)‘)“+×Ñð;
ùð 	$øð3øØ2øð3ø÷ 3÷ 3ñ 3ûä×&Ñ&¬Ð5ò 	áð	úð #ùÜ×)Ñ)ò Ùðúð ùð ×ÑÔ!ðØ"×"Ò"øÜ×)Ñ)ò Ùðúð —)‘)“+×Òüs-  ‚BI	ÂFÂA"I	Ã5FÃ6'I	ÄF9 Ä.FÄ/F9 Ä2F$ÅF ÅF$ÅF9 Å F"Å!F9 Å%I	Å6G Å;GÅ<G Æ I	ÆG3ÆI	ÆI	ÆF9 Æ F$Æ"F9 Æ$F6Æ*F-Æ+F6Æ2F9 Æ9GÇG5 ÇGÇG5 ÇG ÇG0Ç-I	Ç/G0Ç0I	Ç5IÈHÈHÈHÈIÈH)È&IÈ(H)È)IÈ?IÉ IÉI	c              ƒ   óŽ  K  — | j                   r¯t        j                  d«      ƒ d{  –—†  t        j                  «       | j
                  k7  r3| j                  j                  d«       d| _         | j                  «        y| j                  «       ƒ d{  –—†  | j                  «       ƒ d{  –—†  | j                   rŒ®yy7 Œ—7 Œ-7 Œ­w)z1Periodically check worker health and manage pool.g      ð?Nz+Parent changed, shutting down dirty arbiterF)r4   rs   Úsleepr#   Úgetppidr%   r!   rK   r—   Úmurder_workersr’   rŒ   s    rB   r£   zDirtyArbiter._worker_monitor  s”   è ø€ à�jŠjÜ—-‘- Ó$×$Ð$ô �z‰z‹|˜tŸy™yÒ(Ø—‘× Ñ Ð!NÔOØ"�”
Ø—‘Ô Øà×%Ñ%Ó'×'Ð'Ø×%Ñ%Ó'×'Ð'ð �j�jØ$øð (øØ'ús:   ‚%C§B?¨A+CÂCÂCÂ+CÂ,CÂ=CÃCÃCc              ƒ   óz   K  — | j                  «        | j                  r| j                  «       ƒ d{  –—†  yy7 Œ­w)z#Handle SIGCHLD - reap dead workers.N)Úreap_workersr4   r’   rŒ   s    rB   r‹   zDirtyArbiter._handle_sigchldŽ  s3   è ø€ à×ÑÔà�:Š:Ø×%Ñ%Ó'×'Ñ'ð Ø'ús   ‚0;²9³;c              ƒ   ó¶  K  — | j                   j                  d«       	 | j                  rà	 t        j                  |«      ƒ d{  –—† }|j                  d«      }|t        j                  k(  r| j                  ||«      ƒ d{  –—†  nv|t        j                  k(  r| j                  ||«      ƒ d{  –—†  nH|t        j                  k(  r| j                  ||«      ƒ d{  –—†  n| j                  ||«      ƒ d{  –—†  | j                  rŒà|j#                  «        	 |j%                  «       ƒ d{  –—†  y7 Œð# t
        j                  $ r Y ŒAw xY w7 ŒÍ7 Œ¡7 Œu7 Œ\# t        $ r&}| j                   j!                  d|«       Y d}~Œwd}~ww xY w7 ŒZ# t        $ r Y yw xY w# |j#                  «        	 |j%                  «       ƒ d{  –—†7   w # t        $ r Y w w xY wxY w­w)a
  
        Handle a connection from an HTTP worker.

        Routes requests to available dirty workers and returns responses.
        Supports both regular responses and streaming (chunk-based) responses.
        Also handles stash (shared state) operations.
        z&New client connection from HTTP workerNÚtypezClient connection error: %s)r!   Údebugr4   r   Úread_message_asyncrs   ÚIncompleteReadErrorrW   ÚMSG_TYPE_STASHÚhandle_stash_requestÚMSG_TYPE_STATUSÚhandle_status_requestÚMSG_TYPE_MANAGEÚhandle_manage_requestÚroute_requestrJ   Úerrorr›   Úwait_closed)rA   ÚreaderÚwriterÚmessageÚmsg_typerN   s         rB   r¡   zDirtyArbiter.handle_client•  s¦  è ø€ ð 	�‰�‰Ð?Ô@ð	Ø—*’*ðÜ$1×$DÑ$DÀVÓ$L×L�Gð #Ÿ;™; vÓ.�ð œ}×;Ñ;Ò;Ø×3Ñ3°G¸VÓD×DÑDà¤×!>Ñ!>Ò>Ø×4Ñ4°W¸fÓE×EÑEà¤×!>Ñ!>Ò>Ø×4Ñ4°W¸fÓE×EÑEð ×,Ñ,¨W°fÓ=×=Ð=ð' —*“*ð. �L‰LŒNðØ×(Ñ(Ó*×*Ñ*ð/ MùÜ×2Ñ2ò Ùðúð Eøð Føð Føð >ùÜò 	=Ø�H‰H�N‰NÐ8¸!×<Ñ<ûð	=úð
 +ùÜò Ùðûð �L‰LŒNðØ×(Ñ(Ó*×*Ò*øÜò Ùðÿs  ‚GŸE ¬D7 ÁD5ÁD7 Á	9E ÂEÂ-E Â0EÂ1-E ÃEÃE Ã9EÃ:E ÄGÄF Ä/F
Ä0F Ä4GÄ5D7 Ä7EÅ
E ÅEÅE ÅE ÅE ÅE Å	FÅ!FÅ=F ÆFÆF Æ
F Æ	FÆGÆFÆGÆGÆ-GÇ GÇGÇGÇ	GÇGÇGÇGÇGc              ƒ   ó  K  — |j                  dd«      }|j                  d«      }| j                  |«      ƒ d{  –—† }|€h| j                  st        d«      }n%|r| j                  rt        |«      }nt        d«      }t        ||«      }t        j                  ||«      ƒ d{  –—†  y|| j                  vr| j                  |«      ƒ d{  –—†  | j                  |   }t        j                  «       j                  «       }	|j                  |||	f«      ƒ d{  –—†  	 |	ƒ d{  –—†  y7 Œî7 Œ‡7 Œa7 Œ7 Œ# t        $ rC}
t        |t!        d|
› �|¬«      «      }t        j                  ||«      ƒ d{  –—†7   Y d}
~
yd}
~
ww xY w­w)aä  
        Route a request to an available dirty worker via queue.

        Each worker has a dedicated queue and consumer task. Requests are
        submitted to the queue and processed sequentially by the consumer.

        For streaming responses, messages (chunks) are forwarded directly
        to the client_writer as they arrive from the worker.

        Args:
            request: Request message dict
            client_writer: StreamWriter to send responses to client
        ÚidÚunknownra   NzNo dirty workers availablezRequest failed: ©Ú	worker_id)rW   Ú_get_available_workerr-   r   r:   r   r   r   Úwrite_message_asyncr0   Ú_start_worker_consumerrs   r�   Úcreate_futureÚputrJ   r
   )rA   ÚrequestÚclient_writerÚ
request_idra   r`   r¼   ÚresponseÚqueueÚfuturerN   s              rB   r»   zDirtyArbiter.route_request½  sp  è ø€ ð —[‘[  yÓ1ˆ
Ø—;‘;˜zÓ*ˆð  ×5Ñ5°hÓ?×?ˆ
ØÐà—<’<Ü"Ð#?Ó@‘Ù˜dŸnšnä4°XÓ>‘ä"Ð#?Ó@�Ü*¨:°uÓ=ˆHÜ×3Ñ3°MÀ8ÓL×LÐLØð ˜T×/Ñ/Ñ/Ø×-Ñ-¨jÓ9×9Ð9à×"Ñ" :Ñ.ˆÜ×)Ñ)Ó+×9Ñ9Ó;ˆð �i‰i˜ -°Ð8Ó9×9Ð9ð	MØ�L‰Lð5 @øð Møð
 :øð 	:øð ùÜò 	MÜ*ØÜ Ð#3°A°3Ð!7À:ÔNóˆHô  ×3Ñ3°MÀ8ÓL×LÖLûð	Müs�   ‚8FºD)»A(FÂ#D+Â$'FÃD-ÃAFÄD/ÄFÄD3 Ä#D1Ä$D3 Ä(FÄ+FÄ-FÄ/FÄ1D3 Ä3	E?Ä<3E:Å/E2Å0E:Å5FÅ:E?Å?Fc              ƒ   ó¸   ‡ ‡‡K  — t        j                  «       Š‰‰ j                  ‰<   ˆˆ ˆfd„}t        j                   |«       «      }|‰ j                  ‰<   y­w)z3Start a consumer task for a worker's request queue.c               “   óê  •K  — ‰j                   ry	 ‰j                  «       ƒ d {  –—† \  } }}	 ‰j                  ‰| |«      ƒ d {  –—†  |j                  «       s|j	                  d «       ‰j                  «        	 ‰j                   rŒxy y 7 Œe7 ŒG# t
        $ r+}|j                  «       s|j                  |«       Y d }~ŒSd }~ww xY w# ‰j                  «        w xY w# t        j                  $ r Y y w xY w­wr‰   )
r4   rW   Ú_execute_on_workerÚdoneÚ
set_resultrJ   Úset_exceptionÚ	task_doners   r¥   )rÌ   rÍ   rÑ   rN   rÐ   rA   r`   s       €€€rB   Úconsumerz5DirtyArbiter._start_worker_consumer.<locals>.consumerö  sÐ   øè ø€ Ø—*’*ðØ;@¿9¹9»;×5FÑ2�G˜]¨Fð
*Ø"×5Ñ5Ø&¨°ó÷ ð ð  &Ÿ{™{œ}Ø"×-Ñ-¨dÔ3ð
 Ÿ™Õ)ð —*•*à5Føðùô
 %ò 4Ø%Ÿ{™{œ}Ø"×0Ñ0°Ô3ÿøð4ûð Ÿ™Õ)ûÜ×-Ñ-ò Ùðüs…   ƒC3‘C ¤B
¥C ®B ÁBÁ%B Á*C Á:C3ÂC3Â
C ÂB Â	CÂ!B=Â8C Â=CÃC ÃCÃC ÃC0Ã-C3Ã/C0Ã0C3N)rs   ÚQueuer0   rŠ   r1   )rA   r`   rÙ   ÚtaskrÐ   s   ``  @rB   rÉ   z#DirtyArbiter._start_worker_consumerñ  sK   úè ø€ ä—‘“ˆØ).ˆ×Ñ˜:Ñ&ö	ô$ ×"Ñ"¡8£:Ó.ˆØ,0ˆ×Ñ˜jÒ)ùs   …AAc              ƒ   óÚ  K  — |j                  dd«      }	 | j                  |«      ƒ d{  –—† \  }}t        j                  ||«      ƒ d{  –—†  	 	 t	        j
                  t        j                  |«      | j                  j                  ¬«      ƒ d{  –—† }|j                  d«      }	|	t        j                  k(  rt        j                  ||«      ƒ d{  –—†  Œ‹|	t        j                  k(  rt        j                  ||«      ƒ d{  –—†  y|	t        j                  t        j                   fv rt        j                  ||«      ƒ d{  –—†  y| j"                  j%                  d|	«       �Œ7 �ŒB7 �Œ$7 ŒÞ# t        j                  $ r] | j                  |«       t        |t        d| j                  j                  «      «      }t        j                  ||«      ƒ d{  –—†7   Y yw xY w7 �Œ7 Œâ7 Œ£# t&        $ rq}
| j"                  j)                  d||
«       | j                  |«       t        |t+        d	|
› �|¬
«      «      }t        j                  ||«      ƒ d{  –—†7   Y d}
~
yd}
~
ww xY w­w)a  
        Execute request on a specific worker (called by consumer).

        Handles both regular responses and streaming (chunk-based) responses.
        For streaming, chunk and end messages are forwarded directly to the
        client_writer as they arrive from the worker.
        rÃ   rÄ   N)ÚtimeoutzWorker timeoutr±   z$Unknown message type from worker: %sz Error executing on worker %s: %szWorker communication failed: rÅ   )rW   Ú_get_worker_connectionr   rÈ   rs   Úwait_forr³   r    Údirty_timeoutÚTimeoutErrorÚ_close_worker_connectionr   r	   ÚMSG_TYPE_CHUNKÚMSG_TYPE_ENDÚMSG_TYPE_RESPONSEÚMSG_TYPE_ERRORr!   rK   rJ   r¼   r
   )rA   r`   rÌ   rÍ   rÎ   r¾   r¿   rÀ   rÏ   rÁ   rN   s              rB   rÔ   zDirtyArbiter._execute_on_worker  s5  è ø€ ð —[‘[  yÓ1ˆ
ð4	MØ#'×#>Ñ#>¸zÓ#J×J‰NˆF�FÜ×3Ñ3°F¸GÓD×DÐDð ðÜ$+×$4Ñ$4Ü%×8Ñ8¸Ó@Ø $§¡× 6Ñ 6ô%÷ �Gð  #Ÿ;™; vÓ.�ð œ}×;Ñ;Ò;Ü'×;Ñ;¸MÈ7ÓS×SÐSØð œ}×9Ñ9Ò9Ü'×;Ñ;¸MÈ7ÓS×SÐSØð ¤× ?Ñ ?Ü -× <Ñ <ð >ñ >ä'×;Ñ;¸MÈ7ÓS×SÐSØð —‘× Ñ Ð!GÈÔRñK ð	 KùØDùð
ùô ×+Ñ+ò 
ð ×1Ñ1°*Ô=Ü2Ø"Ü)Ð*:¸D¿H¹H×<RÑ<RÓSó �Hô (×;Ñ;¸MÈ8ÓT×TÑTÙð
úð  Tùð
 Tøð Tùô ò 	MØ�H‰H�N‰NÐ=¸zÈ1ÔMØ×)Ñ)¨*Ô5Ü*ØÜ Ð#@ÀÀÐ!DØ+5ô7óˆHô
  ×3Ñ3°MÀ8ÓL×LÖLûð	Müsã   ‚I+–G. ªE,« G. ÁE/ÁG. ÁAE4 ÂE2ÂE4 Â=G. ÃG'Ã1G. ÄG*ÄG. ÄI+Ä;G. ÅG,Å	G. ÅI+ÅG. Å/G. Å2E4 Å4A'G$ÇGÇG$Ç!G. Ç"I+Ç#G$Ç$G. Ç*G. Ç,G. Ç.	I(Ç7A!I#ÉIÉI#ÉI+É#I(É(I+c              ƒ   óº  K  — |r4| j                   r(|| j                  v rt        | j                  |   «      }n$yt        | j                  j	                  «       «      }|sy|rG| j                   r;| j
                  j                  |d«      }|dz   t        |«      z  | j
                  |<   n"| j                  }|dz   t        |«      z  | _        ||t        |«      z     S ­w)añ  
        Get an available worker PID using round-robin selection.

        If app_path is provided, only returns workers that have loaded
        that specific app. Uses per-app round-robin to ensure fair
        distribution among eligible workers.

        Args:
            app_path: Optional import path of the target app. If None,
                     returns any worker using global round-robin.

        Returns:
            Worker PID or None if no eligible workers are available.
        Nr   r   )	r:   r;   r]   r-   Úkeysr=   rW   rV   r2   )rA   ra   Úeligible_pidsÚidxs       rB   rÇ   z"DirtyArbiter._get_available_workerK  sÑ   è ø€ ñ  ˜Ÿšð ˜4×.Ñ.Ñ.Ü $ T×%8Ñ%8¸Ñ%BÓ C‘ð ô ! §¡×!2Ñ!2Ó!4Ó5ˆMáØñ ˜ŸšØ×&Ñ&×*Ñ*¨8°QÓ7ˆCØ.1°A©g¼¸]Ó9KÑ-KˆD× Ñ  Ò*à×'Ñ'ˆCØ%(¨1¡W´°MÓ0BÑ$BˆDÔ!à˜S¤3 }Ó#5Ñ5Ñ6Ð6ùs   ‚CCc              ƒ   óÄ  K  — || j                   v r| j                   |   S | j                  j                  |«      }|st        d|› �«      ‚t	        d«      D ]@  }t
        j                  j                  |«      r n-t        j                  d«      ƒ d{  –—†  ŒB t        d|› �«      ‚t        j                  |«      ƒ d{  –—† \  }}||f| j                   |<   ||fS 7 ŒI7 Œ­w)z%Get or create connection to a worker.zNo socket for worker é2   çš™™™™™¹?NzWorker socket not ready: )r/   r.   rW   r   Úranger#   r*   rž   rs   r«   Úopen_unix_connection)rA   r`   r,   Ú_r¾   r¿   s         rB   rÞ   z#DirtyArbiter._get_worker_connectionu  sä   è ø€ à˜×0Ñ0Ñ0Ø×*Ñ*¨:Ñ6Ð6à×)Ñ)×-Ñ-¨jÓ9ˆÙÜÐ4°Z°LÐAÓBÐBô �r“ò 	HˆAÜ�w‰w�~‰~˜kÔ*ÙÜ—-‘- Ó$×$Ñ$ð	Hô
 Ð8¸¸ÐFÓGÐGä&×;Ñ;¸KÓH×H‰ˆ�Ø/5°vÐ.>ˆ×Ñ 
Ñ+Ø�vˆ~Ðð %øð Iús$   ‚BC ÂCÂ,C Â?CÃ C ÃC c                 ó~   — || j                   v r/| j                   j                  |«      \  }}|j                  «        yy)zClose connection to a worker.N)r/   rd   r›   )rA   r`   Ú_readerr¿   s       rB   râ   z%DirtyArbiter._close_worker_connectionŠ  s8   € à˜×0Ñ0Ñ0Ø"×5Ñ5×9Ñ9¸*ÓE‰OˆG�VØ�L‰L�Nð 1rD   c              ƒ   óª  K  — |j                  dd«      }t        j                  «       }g }| j                  j	                  «       D ]f  \  }}	 |j
                  j                  «       }t        ||z
  d«      }	|j                  ||j                  t        |dg «      t        |dd«      |	dœ«       Œh |j                  d	„ ¬
«       | j                  |t!        |«      | j"                  r#t%        | j"                  j'                  «       «      ng dœ}
t)        ||
«      }t+        j,                  ||«      ƒ d{  –—†  y# t        t        t        f$ r d}	Y ŒØw xY w7 Œ!­w)zô
        Handle a status query request.

        Returns information about the dirty arbiter and its workers.

        Args:
            message: Status request message
            client_writer: StreamWriter to send response to client
        rÃ   rÄ   é   NrY   ÚbootedF)r"   ÚageÚappsrõ   Úlast_heartbeatc                 ó   — | d   S )Nrö   © )rh   s    rB   r�   z4DirtyArbiter.handle_status_request.<locals>.<lambda>±  s
   € ¨¨%©€ rD   ©Úkey)Úarbiter_pidr-   rG   r÷   )rW   ÚtimeÚ	monotonicr-   rU   ÚtmpÚlast_updateÚroundÚOSErrorÚ
ValueErrorÚAttributeErrorrX   rö   ÚgetattrÚsortr"   rV   r:   r]   rè   r   r   rÈ   )rA   rÀ   rÍ   rÎ   ÚnowÚworkers_infor"   Úworkerr  rø   ÚresultrÏ   s               rB   r¸   z"DirtyArbiter.handle_status_request”  s>  è ø€ ð —[‘[  yÓ1ˆ
Ü�n‰nÓˆàˆØŸ<™<×-Ñ-Ó/ò 	‰KˆC�ð&Ø$Ÿj™j×4Ñ4Ó6�Ü!& s¨[Ñ'8¸!Ó!<�ð ×ÑØØ—z‘zÜ ¨°RÓ8Ü! &¨(°EÓ:Ø"0ñ!õ ð	ð 	×ÑÑ0ÐÔ1ð  Ÿ8™8Ø#Ü Ó-Ø37·>²>”D˜Ÿ™×,Ñ,Ó.Ô/Àrñ	
ˆô ! ¨VÓ4ˆÜ×/Ñ/°¸xÓH×HÑHøô+ œZ¬Ð8ò &Ø!%’ð&úð* 	Iús7   ‚A	EÁ)D5Á5B:EÄ/EÄ0EÄ5EÅEÅEÅEc              ƒ   óT  ‡ K  — |j                  dd«      }|j                  d«      }t        dt        |j                  dd«      «      «      }	 |t        k(  r±d}t	        |«      D ]K  }‰ j                  «       }|�‰ xj                  dz  c_        |dz  }t        j                  d«      ƒ d{  –—†  ŒM |dk(  r)d	d
|ddt        ‰ j                  «      ‰ j                  dœ}�n]d	d
||t        ‰ j                  «      ‰ j                  dœ}�n5|t        k(  ró‰ j                  «       }	d}
t	        |«      D ]¬  }‰ j                  |	k  r n›t        ‰ j                  «      dk  r n�‰ xj                  dz  c_        t        ‰ j                  j                  «       ˆ fd„¬«      }‰ j                  |t         j"                  «       |
dz  }
t        j                  d«      ƒ d{  –—†  Œ® d	d||
t        ‰ j                  «      ‰ j                  dœ}n9t%        d|› �«      }t'        ||«      }t)        j*                  ||«      ƒ d{  –—†  y‰ j,                  j/                  d|t        k(  rd
nd||j                  d|j                  dd«      «      «       t1        ||«      }t)        j*                  ||«      ƒ d{  –—†  y7 �Œ7 ŒÝ7 Œ~7 Œ# t2        $ rc}‰ j,                  j5                  d|«       t'        |t%        t7        |«      «      «      }t)        j*                  ||«      ƒ d{  –—†7   Y d}~yd}~ww xY w­w)zý
        Handle a worker management request.

        Supports adding or removing dirty workers via protocol messages.

        Args:
            message: Manage request message
            client_writer: StreamWriter to send response to client
        rÃ   rÄ   Úopr   Úcountr   Nrí   Tr^   z)All apps have reached their worker limits)ÚsuccessÚ	operationÚ	requestedÚspawnedÚreasonÚtotal_workersÚtarget_workers)r  r  r  r  r  r  c                 ó6   •— ‰j                   |    j                  S r‰   ©r-   rö   ©ÚprA   s    €rB   r�   z4DirtyArbiter.handle_manage_request.<locals>.<lambda>ú  s   ø€ °4·<±<À±?×3FÑ3F€ rD   rû   Úremove)r  r  r  Úremovedr  r  zUnknown manage operation: z6Worker management: %s %d workers (spawned/removed: %d)r  r  zManage operation error: %s)rW   rQ   Úintr   rî   Úspawn_workerr6   rs   r«   rV   r-   r   rS   Úminrè   Úkill_workerr{   r}   r   r   r   rÈ   r!   rj   r   rJ   r¼   rm   )rA   rÀ   rÍ   rÎ   r  r  r  rð   r  r™   r  Ú
oldest_pidr¼   rÏ   rN   s   `              rB   rº   z"DirtyArbiter.handle_manage_request½  sê  øè ø€ ð —[‘[  yÓ1ˆ
Ø�[‰[˜ÓˆÜ�A”s˜7Ÿ;™; w°Ó2Ó3Ó4ˆðN	MØ”]Ò"à�Ü˜u›ò -�AØ!×.Ñ.Ó0�FØÐ)Ø×(Ò(¨AÑ-Õ(Ø 1™˜Ü!Ÿ-™-¨Ó,×,Ñ,ð-ð ˜a’<à#'Ø%*Ø%*Ø#$Ø"MÜ),¨T¯\©\Ó):Ø*.×*:Ñ*:ñ’Fð $(Ø%*Ø%*Ø#*Ü),¨T¯\©\Ó):Ø*.×*:Ñ*:ñ’Fð Ô'Ò'à"×7Ñ7Ó9�Ø�ä˜u›ò -�AØ×'Ñ'¨;Ò6ÙÜ˜4Ÿ<™<Ó(¨AÒ-Ùà×$Ò$¨Ñ)Õ$ô "% T§\¡\×%6Ñ%6Ó%8Û)Fô"H�Jà×$Ñ$ Z´·±Ô@Ø˜q‘L�GÜ!Ÿ-™-¨Ó,×,Ñ,ð-ð   $Ø!)Ø!&Ø&Ü%(¨¯©Ó%6Ø&*×&6Ñ&6ñ‘ô #Ð%?À¸tÐ#DÓE�Ü.¨z¸5ÓA�Ü#×7Ñ7¸ÀxÓP×PÐPØà�H‰H�M‰MÐRØ#%¬Ò#6™%¸HØØ Ÿ*™* Y°·
±
¸9ÀaÓ0HÓIôKô
 % Z°Ó8ˆHÜ×3Ñ3°MÀ8ÓL×LÑLðA -ùðR -øð Qøð Mùäò 	MØ�H‰H�N‰NÐ7¸Ô;Ü*¨:´zÄ#ÀaÃ&Ó7IÓJˆHÜ×3Ñ3°MÀ8ÓL×LÖLûð	Müs•   ƒA	L(ÁAJ9 Â*J0Â+D*J9 ÇJ3ÇA J9 È6J5È7J9 È;L(È<A.J9 Ê*J7Ê+J9 Ê/L(Ê0J9 Ê3J9 Ê5J9 Ê7J9 Ê9	L%ËAL ÌLÌL ÌL(Ì L%Ì%L(c           	   ƒ   ód  K  — |j                  dd«      }|j                  d«      }|j                  dd«      }|j                  d«      }|j                  d«      }|j                  d«      }	 d	}	|t        k(  r3|| j                  vri | j                  |<   || j                  |   |<   d
}	�nX|t        k(  r?|| j                  vrddi}	�n;|| j                  |   vrddi}	�n$| j                  |   |   }	�n|t        k(  r7|| j                  v r%|| j                  |   v r| j                  |   |= d
}	�nÔd}	�nÐ|t
        k(  rl|| j                  vrg }	�nµt        | j                  |   j                  «       «      }
|r.|
D �cg c]#  }t        j                  t        |«      |«      r|‘Œ% }
}|
}	�n[|t        k(  r/|| j                  v r| j                  |   j                  «        d
}	�n#|t        k(  r0|| j                  vrddi}	�nt        | j                  |   «      |dœ}	nê|t        k(  r || j                  vri | j                  |<   d
}	nÁ|t        k(  r!|| j                  v r| j                  |= d
}	nšd}	n—|t         k(  r$t        | j                  j                  «       «      }	nj|t"        k(  r(|| j                  vrd}	nP|€d
}	nK|| j                  |   v }	n9t%        d|› �«      }t'        ||«      }t)        j*                  ||«      ƒ d	{  –—†  y	t-        |	t.        «      r{d|	v rw|	d   }|dk(  rt%        d|› �«      }n(|dk(  rt%        d|› �«      }nt%        t        |	«      «      }d|j1                  «       j3                  dd«      › d�|_        t'        ||«      }nt7        ||	«      }t)        j*                  ||«      ƒ d	{  –—†  y	c c}w 7 ŒÀ7 Œ# t8        $ rc}| j:                  j=                  d|«       t'        |t%        t        |«      «      «      }t)        j*                  ||«      ƒ d	{  –—†7   Y d	}~y	d	}~ww xY w­w)a!  
        Handle a stash operation directly in the arbiter.

        All stash tables are stored in arbiter memory for simplicity
        and fast access.

        Args:
            message: Stash operation message
            client_writer: StreamWriter to send response to client
        rÃ   rÄ   r  ÚtableÚ rü   ÚvalueÚpatternNTr¼   Úkey_not_foundFÚtable_not_found)Úsizer"  zUnknown stash operation: zTable not found: zKey not found: ÚStashrð   ÚErrorzStash operation error: %s)rW   r   r?   r   r   r   r]   rè   Úfnmatchrm   r   Úclearr   rV   r   r   r   r   r   r   r   rÈ   Ú
isinstanceÚdictÚtitleÚreplaceÚ
error_typer   rJ   r!   r¼   )rA   rÀ   rÍ   rÎ   r  r"  rü   r$  r%  r  Úall_keysÚkr¼   rÏ   r1  rN   s                   rB   r¶   z!DirtyArbiter.handle_stash_request  s/  è ø€ ð —[‘[  yÓ1ˆ
Ø�[‰[˜ÓˆØ—‘˜G RÓ(ˆØ�k‰k˜%Ó ˆØ—‘˜GÓ$ˆØ—+‘+˜iÓ(ˆða	MØˆFà”\Ò!à × 1Ñ 1Ñ1Ø/1�D×%Ñ% eÑ,Ø05�×!Ñ! %Ñ(¨Ñ-Ø’à”|Ò#Ø × 1Ñ 1Ñ1Ø% Ð7’FØ × 1Ñ 1°%Ñ 8Ñ8Ø% Ð7’Fà!×.Ñ.¨uÑ5°cÑ:’Fà”Ò&Ø˜D×-Ñ-Ñ-°#¸×9JÑ9JÈ5Ñ9QÑ2QØ×)Ñ)¨%Ñ0°Ð5Ø!’Fà"’Fà”}Ò$Ø × 1Ñ 1Ñ1Ø’Fä# D×$5Ñ$5°eÑ$<×$AÑ$AÓ$CÓD�HÙØ/7ö $I¨!Ü'.§¡´s¸1³v¸wÔ'Gò %&ð $I˜ð $Ià%’Fà”~Ò%Ø˜D×-Ñ-Ñ-Ø×%Ñ% eÑ,×2Ñ2Ô4Ø’à”}Ò$Ø × 1Ñ 1Ñ1Ø%Ð'8Ð9’Fô !$ D×$5Ñ$5°eÑ$<Ó =Ø!&ñ‘Fð
 ”Ò&Ø × 1Ñ 1Ñ1Ø/1�D×%Ñ% eÑ,Ø‘àÔ,Ò,Ø˜D×-Ñ-Ñ-Ø×)Ñ)¨%Ð0Ø!‘Fà"‘Fà”Ò&Ü˜d×/Ñ/×4Ñ4Ó6Ó7‘à”Ò&Ø × 1Ñ 1Ñ1Ø"‘FØ�[Ø!‘Fà  D×$5Ñ$5°eÑ$<Ð<‘Fô #Ð%>¸r¸dÐ#CÓD�Ü.¨z¸5ÓA�Ü#×7Ñ7¸ÀxÓP×PÐPØô ˜&¤$Ô'¨G°vÑ,=Ø# G™_�
ØÐ!2Ò2Ü&Ð):¸5¸'Ð'BÓC‘EØ ?Ò2Ü&¨¸¸Ð'>Ó?‘Eä&¤s¨6£{Ó3�EØ%*¨:×+;Ñ+;Ó+=×+EÑ+EÀcÈ2Ó+NÐ*OÈuÐ#U�Ô Ü.¨z¸5ÓA‘ä(¨°VÓ<�ä×3Ñ3°MÀ8ÓL×LÑLùò{$IðX Qøð" Mùäò 	MØ�H‰H�N‰NÐ6¸Ô:Ü*¨:´zÄ#ÀaÃ&Ó7IÓJˆHÜ×3Ñ3°MÀ8ÓL×LÖLûð	Müs|   ‚A)P0Á,DO Å9(N8Æ!EO Ë<N=Ë=O ÌP0ÌB0O Î2N?Î3O Î7P0Î8O Î?O Ï	P-Ï
AP(ÐP ÐP(Ð#P0Ð(P-Ð-P0c              ƒ   óh  ‡ K  — ‰ j                   sy‰ j                  }‰ j                   rmt        ‰ j                  «      |k  rU‰ j	                  «       }|€nBt        j                  d«      ƒ d{  –—†  ‰ j                   rt        ‰ j                  «      |k  rŒUt        ‰ j                  «      |kD  rt        ‰ j                  j                  «       ˆ fd„¬«      }‰ j                  |t        j                  «       t        j                  d«      ƒ d{  –—†  t        ‰ j                  «      |kD  rŒ~yy7 ŒÁ7 Œ!­w)z%Maintain the number of dirty workers.Nrí   c                 ó6   •— ‰j                   |    j                  S r‰   r  r  s    €rB   r�   z-DirtyArbiter.manage_workers.<locals>.<lambda>£  s   ø€ ¨4¯<©<¸©?×+>Ñ+>€ rD   rû   )r4   r6   rV   r-   r  rs   r«   r  rè   r  r{   r}   )rA   r6   r  r   s   `   rB   r’   zDirtyArbiter.manage_workers�  sê   øè ø€ à�zŠzØà×&Ñ&ˆð �jŠjœS §¡Ó.°Ò<Ø×&Ñ&Ó(ˆFØˆ~àÜ—-‘- Ó$×$Ð$ð �jŠjœS §¡Ó.°Ó<ô �$—,‘,Ó +Ò-ä˜TŸ\™\×.Ñ.Ó0Û!>ô@ˆJà×Ñ˜Z¬¯©Ô8Ü—-‘- Ó$×$Ð$ô �$—,‘,Ó +Õ-ð %øð %ús1   ƒA)D2Á,D.Á-(D2ÂA8D2ÄD0ÄD2Ä,D2Ä0D2c                 óF  — | j                   r| j                   j                  d«      }n6|r$t        | j                  j	                  «       «      }n| j                  «       }|s| j                  j                  d«       y| xj                  dz  c_        t        j                  j                  | j                  d| j                  › d�«      }t        | j                  | j                  || j                  | j                  |¬«      }t        j                   «       }|dk7  rr||_        || j"                  |<   || j$                  |<   | j'                  ||«       | j                  j)                  | |«       | j                  j+                  d||«       |S t        j,                  «       |_        	 t/        j0                  d	| j                  j2                  › d
�«       |j5                  «        t        j6                  d«       y# t8        $ r7}t        j6                  |j:                  �|j:                  nd«       Y d}~yd}~wt<        $ r^ | j                  j?                  d«       |j@                  st        j6                  | jB                  «       t        j6                  d«       Y yw xY w)a  
        Spawn a new dirty worker.

        Worker app assignment follows these priorities:
        1. If there are pending respawns (from dead workers), use those apps
        2. Otherwise, determine apps for a new worker based on allocation
        3. If force_all_apps=True, spawn with all apps regardless of limits

        Args:
            force_all_apps: If True, spawn worker with all apps ignoring limits

        Returns:
            Worker PID in parent process, or None if no apps need workers
        r   z)No apps need more workers, skipping spawnNr   zworker-z.sock)rö   r%   rY   r    r!   r,   z,Spawned dirty worker (pid: %s) with apps: %szdirty-worker [ú]z!Exception in dirty worker process)"r>   rd   r]   r:   rè   r[   r!   r²   r3   r#   r*   r+   r)   r   r"   r    Úforkr-   r.   rb   Údirty_post_forkrj   r$   r   rr   Ú	proc_nameÚinit_processÚ_exitÚ
SystemExitÚcoderJ   Ú	exceptionrõ   ÚWORKER_BOOT_ERROR)rA   Úforce_all_appsrY   r,   r
  r"   rN   s          rB   r  zDirtyArbiter.spawn_worker§  sÿ  € ð  ×!Ò!Ø×.Ñ.×2Ñ2°1Ó5‰IÙä˜TŸ^™^×0Ñ0Ó2Ó3‰Ið ×5Ñ5Ó7ˆIáØ�H‰H�N‰NÐFÔGØà�Š˜1Ñ�Ü—g‘g—l‘lØ�K‰K˜7 4§?¡?Ð"3°5Ð9ó
ˆô Ø—‘Ø—‘ØØ—‘Ø—‘Ø#ô
ˆô �g‰g‹iˆØ�!Š8àˆFŒJØ &ˆD�L‰L˜ÑØ'2ˆD×Ñ Ñ$ð ×&Ñ& s¨IÔ6à�H‰H×$Ñ$ T¨6Ô2Ø�H‰H�M‰MÐHØ˜yô*àˆJô —Y‘Y“[ˆŒ
ð
	Ü×Ñ °·±×0BÑ0BÐ/CÀ1ÐEÔFØ×ÑÔ!Ü�H‰H�Q�KøÜò 	:Ü�H‰H˜qŸv™vÐ1�Q—V’V°q×9Ñ9ûÜò 	Ø�H‰H×ÑÐBÔCØ—=’=Ü—‘˜×/Ñ/Ô0Ü�H‰H�QŽKð		ús    Æ(AG; Ç;	J È-H6È6A'J ÊJ c                 óÄ   — 	 t        j                  ||«       y# t        $ r=}|j                  t        j                  k(  r| j                  |«       Y d}~yY d}~yd}~ww xY w)zKill a worker by PID.N)r#   Úkillr  ÚerrnoÚESRCHÚ_cleanup_worker)rA   r"   r†   rN   s       rB   r  zDirtyArbiter.kill_workerï  sK   € ð	*Ü�G‰G�C˜ÕøÜò 	*Ø�w‰wœ%Ÿ+™+Ò%Ø×$Ñ$ S×)Ñ)ô &ûð	*ús   ‚ ™	A¢.AÁAc                 ó²  — | j                  |«       || j                  v r*| j                  |   j                  «        | j                  |= | j                  j	                  |d«       || j
                  v r5t        | j
                  |   «      }|r| j                  j                  |«       | j                  |«       | j                  j	                  |d«      }|r| j                  j                  | |«       | j                  j	                  |d«      }|r7t        j                  j!                  |«      r	 t        j"                  |«       yyy# t$        $ r Y yw xY w)z¤
        Clean up after a worker exits.

        Saves the dead worker's app list to pending respawns so the
        replacement worker gets the same apps.
        N)râ   r1   r§   r0   rd   r<   r]   r>   rX   rf   r-   r    Údirty_worker_exitr.   r#   r*   rž   rŸ   r  )rA   r"   Ú	dead_appsr
  r,   s        rB   rF  zDirtyArbiter._cleanup_worker÷  s.  € ð 	×%Ñ% cÔ*ð �$×'Ñ'Ñ'Ø×!Ñ! #Ñ&×-Ñ-Ô/Ø×%Ñ% cÐ*ð 	×Ñ×Ñ˜s DÔ)ð �$×%Ñ%Ñ%Ü˜T×0Ñ0°Ñ5Ó6ˆIÙØ×&Ñ&×-Ñ-¨iÔ8ð 	×Ñ Ô$à—‘×!Ñ! # tÓ,ˆÙØ�H‰H×&Ñ& t¨VÔ4Ø×)Ñ)×-Ñ-¨c°4Ó8ˆÙœ2Ÿ7™7Ÿ>™>¨+Ô6ðÜ—	‘	˜+Õ&ð 7ˆ;øô ò Ùðús   Ä2E
 Å
	EÅEc              ƒ   ó,  K  — | j                   j                  syt        | j                  j	                  «       «      D ]¾  \  }}	 t        j                  «       |j                  j                  «       z
  | j                   j                  k  rŒN	 |j                  sD| j                  j                  d|«       d|_        | j                  |t        j                   «       ŒŸ| j                  |t        j"                  «       ŒÀ y# t        t        f$ r Y ŒÓw xY w­w)z!Kill workers that have timed out.NzDIRTY WORKER TIMEOUT (pid:%s)T)r    rà   r]   r-   rU   rþ   rÿ   r   r  r  r  Úabortedr!   Úcriticalr  r{   ÚSIGABRTÚSIGKILL)rA   r"   r
  s      rB   r­   zDirtyArbiter.murder_workers  sÙ   è ø€ à�x‰x×%Ò%Øä §¡× 2Ñ 2Ó 4Ó5ò 	6‰KˆC�ðÜ—>‘>Ó# f§j¡j×&<Ñ&<Ó&>Ñ>À$Ç(Á(×BXÑBXÒXØð Yð
 —>’>Ø—‘×!Ñ!Ð"AÀ3ÔGØ!%�”Ø× Ñ  ¤f§n¡nÕ5à× Ñ  ¤f§n¡nÕ5ñ	6øô œZÐ(ò Ùðüs,   ‚ADÁAC?Â
A5DÃ?DÄDÄDÄDc                 ó\  — 	 	 t        j                  dt         j                  «      \  }}|syd}t        j                  |«      rt        j                  |«      }nGt        j
                  |«      r2t        j                  |«      }| j                  j                  d||«       || j                  k(  r| j                  j                  d|«       | j                  |«       | j                  j                  d|«       Œ÷# t        $ r(}|j                  t        j                  k7  r‚ Y d}~yd}~ww xY w)zReap dead worker processes.éÿÿÿÿNz)Dirty worker (pid:%s) killed by signal %sz$Dirty worker failed to boot (pid:%s)zDirty worker exited (pid:%s))r#   ÚwaitpidÚWNOHANGÚ	WIFEXITEDÚWEXITSTATUSÚWIFSIGNALEDÚWTERMSIGr!   rK   r@  r¼   rF  rj   r  rD  ÚECHILD)rA   ÚwpidÚstatusÚexitcoder†   rN   s         rB   r¯   zDirtyArbiter.reap_workers.  sî   € ð	ØÜ!Ÿz™z¨"¬b¯j©jÓ9‘��fÙØà�Ü—<‘< Ô'Ü!Ÿ~™~¨fÓ5‘HÜ—^‘^ FÔ+ÜŸ+™+ fÓ-�CØ—H‘H×$Ñ$Ð%PØ%)¨3ô0ð ˜t×5Ñ5Ò5Ø—H‘H—N‘NÐ#IÈ4ÔPà×$Ñ$ TÔ*Ø—‘—‘Ð<¸dÔCð# øô$ ò 	Ø�w‰wœ%Ÿ,™,Ò&Øô 'ûð	ús   ‚*C: ­CC: Ã:	D+ÄD&Ä&D+c              ƒ   óª  K  — | j                   j                  d«       t        | j                  j                  «      D ]/  }| j                  «        t        j                  d«      ƒ d{  –—†  Œ1 t        | j                  j                  «       «      }|| j                  j                  d D ]"  }| j                  |t        j                  «       Œ$ y7 Œh­w)z!Reload workers (SIGHUP handling).zReloading dirty workersrí   N)r!   rj   rî   r    r5   r  rs   r«   r]   r-   rè   r  r{   r}   )rA   rð   Úold_workersr"   s       rB   r�   zDirtyArbiter.reloadG  s¨   è ø€ à�‰�‰Ð/Ô0ô �t—x‘x×-Ñ-Ó.ò 	%ˆAØ×ÑÔÜ—-‘- Ó$×$Ñ$ð	%ô
 ˜4Ÿ<™<×,Ñ,Ó.Ó/ˆØ˜tŸx™x×5Ñ5Ð6Ð7ò 	2ˆCØ×Ñ˜S¤&§.¡.Õ1ñ	2ð	 %ús   ‚A&CÁ(CÁ)A)Cc              ƒ   ó  K  — | j                   j                  «       D ]  }|j                  «        Œ |rt        j                  nt        j
                  }t        j                  «       | j                  j                  z   }t        | j                  j                  «       «      D ]  }| j                  ||«       Œ | j                  rht        j                  «       |k  rQ| j                  «        t        j                  d«      ƒ d{  –—†  | j                  rt        j                  «       |k  rŒQt        | j                  j                  «       «      D ]"  }| j                  |t        j                   «       Œ$ | j                  «        y7 Œ�­w)zStop all workers.rí   N)r1   rP   r§   r{   r}   r   rþ   r    Údirty_graceful_timeoutr]   r-   rè   r  r¯   rs   r«   rN  )rA   ÚgracefulrÛ   r†   Úlimitr"   s         rB   r¨   zDirtyArbiter.stopU  s!  è ø€ ð ×)Ñ)×0Ñ0Ó2ò 	ˆDØ�K‰K�Mð	ñ !)Œf�nŠn¬f¯n©nˆÜ—	‘	“˜dŸh™h×=Ñ=Ñ=ˆô ˜Ÿ™×)Ñ)Ó+Ó,ò 	'ˆCØ×Ñ˜S #Õ&ð	'ð �lŠlœtŸy™y›{¨UÒ2Ø×ÑÔÜ—-‘- Ó$×$Ð$ð �lŠlœtŸy™y›{¨UÓ2ô
 ˜Ÿ™×)Ñ)Ó+Ó,ò 	2ˆCØ×Ñ˜S¤&§.¡.Õ1ð	2à×ÑÕð %ús   ‚DFÄFÄ'FÄ-AFc                 óè  — | j                   rIt        j                  j                  | j                   «      r 	 t        j                  | j                   «       t        j                  j                  | j                  «      r 	 t        j                  | j                  «       	 t        j                  | j                  «      D ]?  }t        j                  t        j                  j                  | j                  |«      «       ŒA t        j                  | j                  «       | j                  j                  d| j                  «       y# t
        $ r Y Œüw xY w# t
        $ r Y ŒÂw xY w# t
        $ r Y ŒPw xY w)zSynchronous cleanup on exit.zDirty arbiter exiting (pid: %s)N)r&   r#   r*   rž   rŸ   r  r,   Úlistdirr)   r+   Úrmdirr!   rj   r"   )rA   rx   s     rB   rw   zDirtyArbiter._cleanup_syncl  s  € ð �<Š<œBŸG™GŸN™N¨4¯<©<Ô8ðÜ—	‘	˜$Ÿ,™,Ô'ô
 �7‰7�>‰>˜$×*Ñ*Ô+ðÜ—	‘	˜$×*Ñ*Ô+ð
	Ü—Z‘Z §¡Ó,ò 8�Ü—	‘	œ"Ÿ'™'Ÿ,™, t§{¡{°AÓ6Õ7ð8ä�H‰H�T—[‘[Ô!ð 	�‰�‰Ð7¸¿¹ÕBøô% ò Ùðûô ò Ùðûô ò 	Ùð	ús6   ·E Â E Â B E% Å	EÅEÅ	E"Å!E"Å%	E1Å0E1)NNr‰   )F)T))Ú__name__Ú
__module__Ú__qualname__Ú__doc__Úsplitr  r{   rz   r@  rC   r@   rS   r[   rb   rf   rt   rq   r~   r—   ru   r£   r‹   r¡   r»   rÉ   rÔ   rÇ   rÞ   râ   r¸   rº   r¶   r’   r  r  rF  r­   r¯   r�   r¨   rw   )Ú.0Úxr  r{   s   0000rB   r   r   1   s   „ ñð <×AÑAÓC÷Eñ E°Œw”v˜w¨™{Õ+õ E€Gð Ðó5 òn 5òDò"ò6:ò$Bò"!ò@<ò8<òt!ò
&òP(ò(ò&òP2Mòh1ò4>Mó@(7òTò*ò'IòR\Mò|sMòj%ó.FòP*ò"òH6ò&ò22óó.Cùõa!Es   œB#
r   )%rg  rs   rD  r+  r#   r{   r'   rþ   Úgunicornr   Úappr   r   Úerrorsr   r   r	   r
   Úprotocolr   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r
  r   r   rú   rD   rB   ú<module>ro     s]   ðñ
ó Û Û Û 	Û Û Û å ç @÷ó ÷÷ ÷ ÷ ñ õ"  ÷SCò SCrD   