§
    …ßjDD  ã                   ó¬   — 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mZ ddl	m
Z
 ddl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  G d	„ d
¦  «        ZdS )a	  
Dirty Worker Process

Asyncio-based worker that loads dirty apps and handles requests
from the DirtyArbiter.

Threading Model
---------------
Each dirty worker runs an asyncio event loop in the main thread for:
- Handling connections from the arbiter
- Managing heartbeat updates
- Coordinating task execution

Actual app execution runs in a ThreadPoolExecutor (separate threads):
- The number of threads is controlled by ``dirty_threads`` config (default: 1)
- Each thread can execute one app action at a time
- The asyncio event loop is NOT blocked by task execution

State and Global Objects
------------------------
Apps can maintain persistent state because:

1. Apps are loaded ONCE when the worker starts (in ``load_apps()``)
2. The same app instances are reused for ALL requests
3. App state (instance variables, loaded models, etc.) persists

Example::

    class MLApp(DirtyApp):
        def init(self):
            self.model = load_heavy_model()  # Loaded once, reused
            self.cache = {}                   # Persistent cache

        def predict(self, data):
            return self.model.predict(data)  # Uses loaded model

Thread Safety:
- With ``dirty_threads=1`` (default): No concurrent access, thread-safe by design
- With ``dirty_threads > 1``: Multiple threads share the same app instances,
  apps MUST be thread-safe (use locks, thread-local storage, etc.)

Heartbeat and Liveness
----------------------
The worker sends heartbeat updates to prove it's alive:

1. A dedicated asyncio task (``_heartbeat_loop``) runs independently
2. It updates the heartbeat file every ``dirty_timeout / 2`` seconds
3. Since tasks run in executor threads, they do NOT block heartbeats
4. The arbiter kills workers that miss heartbeat updates

Timeout Control
---------------
Execution timeout is enforced at two levels:

1. **Worker level**: Each task execution has a timeout (``dirty_timeout``).
   If exceeded, the worker returns a timeout error but the thread may
   continue running (Python threads cannot be cancelled).

2. **Arbiter level**: The arbiter also enforces timeout when waiting
   for worker response. Workers that don't respond are killed via SIGABRT.

Note: Since Python threads cannot be forcibly cancelled, a truly stuck
operation will continue until the worker is killed by the arbiter.
é    N)Úutil)Ú	WorkerTmpé   )Úload_dirty_apps)ÚDirtyAppErrorÚDirtyAppNotFoundErrorÚDirtyTimeoutErrorÚDirtyWorkerError)ÚDirtyProtocolÚmake_responseÚmake_error_responseÚmake_chunk_messageÚmake_end_messagec                   ó´   — e Zd ZdZd„ 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„ ZdS )ÚDirtyWorkerzÐ
    Dirty worker process that loads dirty apps and handles requests.

    Each worker runs its own asyncio event loop and listens on a
    worker-specific Unix socket for requests from the DirtyArbiter.
    c                 ó>   — g | ]}t          t          d |z  ¦  «        ‘ŒS )zSIG%s)ÚgetattrÚsignal)Ú.0Úxs     úH/var/www/html/venv/lib/python3.11/site-packages/gunicorn/dirty/worker.pyú
<listcomp>zDirtyWorker.<listcomp>h   s2   € ð 6ð 6ð 6°�w•v˜w¨™{Ñ+Ô+ð 6ð 6ð 6ó    zABRT HUP QUIT INT TERM USR1c                 óò   — || _         d| _        || _        || _        || _        || _        || _        d| _        d| _        d| _	        t          |¦  «        | _        i | _        d| _        d| _        d| _        dS )a?  
        Initialize a dirty worker.

        Args:
            age: Worker age (for identifying workers)
            ppid: Parent process ID
            app_paths: List of dirty app import paths
            cfg: Gunicorn config
            log: Logger
            socket_path: Path to this worker's Unix socket
        z	[booting]FTN)ÚageÚpidÚppidÚ	app_pathsÚcfgÚlogÚsocket_pathÚbootedÚabortedÚaliver   ÚtmpÚappsÚ_serverÚ_loopÚ	_executor)Úselfr   r   r   r   r    r!   s          r   Ú__init__zDirtyWorker.__init__k   sx   € ð ˆŒØˆŒØˆŒ	Ø"ˆŒØˆŒØˆŒØ&ˆÔØˆŒØˆŒØˆŒ
Ý˜S‘>”>ˆŒØˆŒ	ØˆŒØˆŒ
ØˆŒˆˆr   c                 ó   — d| j         › d�S )Nz<DirtyWorker ú>)r   ©r*   s    r   Ú__str__zDirtyWorker.__str__‡   s   € Ø*˜tœxÐ*Ð*Ð*Ð*r   c                 ó8   — | j                              ¦   «          dS )zUpdate heartbeat timestamp.N)r%   Únotifyr.   s    r   r1   zDirtyWorker.notifyŠ   s   € àŒ�ŠÑÔÐÐÐr   c                 ó†  — | j         j        r3| j         j                             ¦   «         D ]\  }}|t          j        |<   Œt          j        | j         j        | j         j        | j         j	        ¬¦  «         t          j
        ¦   «          t          j        | j                             ¦   «         ¦  «         | j                             ¦   «          |                      ¦   «          |                      ¦   «          t          j        ¦   «         | _        | j                              | ¦  «         d| _        |                      ¦   «          dS )zÂ
        Initialize the worker process after fork.

        This is called in the child process after fork. It sets up
        the environment, loads apps, and starts the main run loop.
        )Ú
initgroupsTN)r   ÚenvÚitemsÚosÚenvironr   Úset_owner_processÚuidÚgidr3   ÚseedÚclose_on_execr%   Úfilenor    Úinit_signalsÚ	load_appsÚgetpidr   Údirty_worker_initr"   Úrun)r*   ÚkÚvs      r   Úinit_processzDirtyWorker.init_processŽ   s  € ð Œ8Œ<ð 	"Øœœ×*Ò*Ñ,Ô,ð "ð "‘��1Ø !•”
˜1‘�åÔ˜tœxœ|¨T¬X¬\Ø*.¬(Ô*=ð	?ñ 	?ô 	?ð 	?õ 	Œ	‰Œˆõ 	Ô˜4œ8Ÿ?š?Ñ,Ô,Ñ-Ô-Ð-ØŒ×ÒÑ Ô Ð ð 	×ÒÑÔÐð 	�ŠÑÔÐõ ”9‘;”;ˆŒØŒ×"Ò" 4Ñ(Ô(Ð(ð ˆŒØ�Š‰
Œ
ˆ
ˆ
ˆ
r   c                 óÀ  — | 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        ¦  «         dS )zSet up signal handlers.N)	ÚSIGNALSr   ÚSIG_DFLÚSIGTERMÚ_signal_handlerÚSIGQUITÚSIGINTÚSIGABRTÚSIGUSR1)r*   Úsigs     r   r>   zDirtyWorker.init_signals²   s«   € ð ”<ð 	/ð 	/ˆCÝŒM˜#�vœ~Ñ.Ô.Ð.Ð.õ 	Œ•f”n dÔ&:Ñ;Ô;Ð;ÝŒ•f”n dÔ&:Ñ;Ô;Ð;ÝŒ•f”m TÔ%9Ñ:Ô:Ð:õ 	Œ•f”n dÔ&:Ñ;Ô;Ð;õ 	Œ•f”n dÔ&:Ñ;Ô;Ð;Ð;Ð;r   c                 óº   — |t           j        k    r| j                             ¦   «          dS d| _        | j        r!| j                             | j        ¦  «         dS dS )z(Handle signals by setting alive = False.NF)r   rN   r    Úreopen_filesr$   r(   Úcall_soon_threadsafeÚ	_shutdown)r*   rO   Úframes      r   rJ   zDirtyWorker._signal_handlerÃ   sa   € à•&”.Ò Ð ØŒH×!Ò!Ñ#Ô#Ð#ØˆFàˆŒ
ØŒ:ð 	<ØŒJ×+Ò+¨D¬NÑ;Ô;Ð;Ð;Ð;ð	<ð 	<r   c                 óJ   — | j         r| j                              ¦   «          dS dS )zInitiate async shutdown.N)r'   Úcloser.   s    r   rS   zDirtyWorker._shutdownÍ   s0   € àŒ<ð 	!ØŒL×ÒÑ Ô Ð Ð Ð ð	!ð 	!r   c                 óÈ  — 	 t          | j        ¦  «        | _        | j                             ¦   «         D ]\  }}| j                             d|¦  «         	 |                     ¦   «          | j                             d|¦  «         ŒQ# t          $ r"}| j         	                    d||¦  «         ‚ d}~ww xY wdS # t          $ r!}| j         	                    d|¦  «         ‚ d}~ww xY w)zLoad all configured dirty apps.zLoaded dirty app: %szInitialized dirty app: %sz%Failed to initialize dirty app %s: %sNzFailed to load dirty apps: %s)
r   r   r&   r5   r    ÚdebugÚinitÚinfoÚ	ExceptionÚerror©r*   ÚpathÚappÚes       r   r?   zDirtyWorker.load_appsÒ   s  € ð	Ý'¨¬Ñ7Ô7ˆDŒIØ!œYŸ_š_Ñ.Ô.ð ð ‘	��cØ”—’Ð5°tÑ<Ô<Ð<ðØ—H’H‘J”J�JØ”H—M’MÐ"=¸tÑDÔDÐDÐDøÝ ð ð ð Ø”H—N’NÐ#JØ#'¨ñ,ô ,ð ,àøøøøðøøøðð øõ ð 	ð 	ð 	ØŒH�NŠNÐ:¸AÑ>Ô>Ð>Øøøøøð	øøøs<   ‚AB6 Á/BÂB6 Â
B1ÂB,Â,B1Â1B6 Â6
C!Ã CÃC!c                 ó  — ddl m} | j        j        } ||d| j        › d�¬¦  «        | _        | j                             d|¦  «         	 t          j	        ¦   «         | _
        t          j        | j
        ¦  «         | j
                             |                      ¦   «         ¦  «         n2# t          $ r%}| j                             d|¦  «         Y d}~nd}~ww xY w|                      ¦   «          dS # |                      ¦   «          w xY w)	z Run the main asyncio event loop.r   )ÚThreadPoolExecutorzdirty-worker-ú-)Úmax_workersÚthread_name_prefixz#Created thread pool with %d threadszWorker error: %sN)Úconcurrent.futuresrb   r   Údirty_threadsr   r)   r    rX   ÚasyncioÚnew_event_loopr(   Úset_event_loopÚrun_until_completeÚ
_run_asyncr[   r\   Ú_cleanup)r*   rb   Únum_threadsr`   s       r   rB   zDirtyWorker.runã   s  € ð 	:Ð9Ð9Ð9Ð9Ð9ð ”hÔ,ˆØ+Ð+Ø#Ø:¨t¬xÐ:Ð:Ð:ð
ñ 
ô 
ˆŒð 	Œ�ŠÐ<¸kÑJÔJÐJð	Ý Ô/Ñ1Ô1ˆDŒJÝÔ" 4¤:Ñ.Ô.Ð.ØŒJ×)Ò)¨$¯/ª/Ñ*;Ô*;Ñ<Ô<Ð<Ð<øÝð 	2ð 	2ð 	2ØŒH�NŠNÐ-¨qÑ1Ô1Ð1Ð1Ð1Ð1Ð1Ð1øøøøð	2øøøð �MŠM‰OŒOˆOˆOˆOøˆD�MŠM‰OŒOˆOˆOøøøs1   Á
AB( Â'C0 Â(
CÂ2CÃC0 ÃCÃC0 Ã0Dc              ƒ   óH  K  — t           j                             | j        ¦  «        rt          j        | j        ¦  «         t          j        | j        | j        ¬¦  «        ƒ d{V —†| _        t          j	        | j        d¦  «         | j
                             d| j        | j        ¦  «         t          j        |                      ¦   «         ¦  «        }	 | j        4 ƒd{V —† | j                             ¦   «         ƒ d{V —† ddd¦  «        ƒd{V —† n# 1 ƒd{V —†swxY w Y   n# t
          j        $ r Y nw xY w|                     ¦   «          	 |ƒ d{V —† dS # t
          j        $ r Y dS w xY w# |                     ¦   «          	 |ƒ d{V —† w # t
          j        $ r Y w w xY wxY w)z6Main async loop - start server and handle connections.)r^   Ni€  zDirty worker %s listening on %s)r6   r^   Úexistsr!   Úunlinkrh   Ústart_unix_serverÚhandle_connectionr'   Úchmodr    rZ   r   Úcreate_taskÚ_heartbeat_loopÚserve_foreverÚCancelledErrorÚcancel)r*   Úheartbeat_tasks     r   rl   zDirtyWorker._run_asyncù   s�  è è € õ Œ7�>Š>˜$Ô*Ñ+Ô+ð 	(ÝŒI�dÔ&Ñ'Ô'Ð'õ %Ô6ØÔ"ØÔ!ð
ñ 
ô 
ð 
ð 
ð 
ð 
ð 
ð 
ˆŒõ 	Œ�Ô! 5Ñ)Ô)Ð)àŒ�ŠÐ7Ø”h Ô 0ñ	2ô 	2ð 	2õ !Ô,¨T×-AÒ-AÑ-CÔ-CÑDÔDˆð
	Ø”|ð 3ð 3ð 3ð 3ð 3ð 3ð 3ð 3Ø”l×0Ò0Ñ2Ô2Ð2Ð2Ð2Ð2Ð2Ð2Ð2ð3ð 3ð 3ñ 3ô 3ð 3ð 3ð 3ð 3ð 3ð 3ð 3ð 3ð 3ð 3ð 3ð 3ð 3ð 3ð 3ð 3ð 3ð 3øøøð 3ð 3ð 3ð 3øøåÔ%ð 	ð 	ð 	ØˆDð	øøøð ×!Ò!Ñ#Ô#Ð#ðØ$Ð$Ð$Ð$Ð$Ð$Ð$Ð$Ð$Ð$øÝÔ)ð ð ð Ø��ðøøøøð ×!Ò!Ñ#Ô#Ð#ðØ$Ð$Ð$Ð$Ð$Ð$Ð$Ð$Ð$øÝÔ)ð ð ð Ø�ðøøøøøøs„   ÃD# Ã DÃ?D# Ä
DÄD# ÄDÄD# Ä"E- Ä#D5Ä2E- Ä4D5Ä5E- ÅE ÅE*Å)E*Å-F!ÆFÆF!ÆFÆF!ÆFÆF!c              ƒ   ó    K  — | j         rD|                      ¦   «          t          j        | j        j        dz  ¦  «        ƒ d{V —† | j         °BdS dS )zPeriodically update heartbeat.g       @N)r$   r1   rh   Úsleepr   Údirty_timeoutr.   s    r   rv   zDirtyWorker._heartbeat_loop  si   è è € àŒjð 	>Ø�KŠK‰MŒMˆMÝ”- ¤Ô 6¸Ñ <Ñ=Ô=Ð=Ð=Ð=Ð=Ð=Ð=Ð=ð Œjð 	>ð 	>ð 	>ð 	>ð 	>r   c              ƒ   ód  K  — | j                              d¦  «         	 | j        rT	 t          j        |¦  «        ƒ d{V —†}n# t
          j        $ r Y n'w xY w|                      ||¦  «        ƒ d{V —† | j        °Tn2# t          $ r%}| j          	                    d|¦  «         Y d}~nd}~ww xY w| 
                    ¦   «          	 |                     ¦   «         ƒ d{V —† dS # t          $ r Y dS w xY w# | 
                    ¦   «          	 |                     ¦   «         ƒ d{V —† w # t          $ r Y w w xY wxY w)zl
        Handle a connection from the arbiter.

        Each connection can send multiple requests.
        zNew connection from arbiterNzConnection error: %s)r    rX   r$   r   Úread_message_asyncrh   ÚIncompleteReadErrorÚhandle_requestr[   r\   rV   Úwait_closed)r*   ÚreaderÚwriterÚmessager`   s        r   rs   zDirtyWorker.handle_connection   sµ  è è € ð 	Œ�ŠÐ4Ñ5Ô5Ð5ð	Ø”*ð ;ðÝ$1Ô$DÀVÑ$LÔ$LÐLÐLÐLÐLÐLÐL�G�GøÝÔ2ð ð ð à�Eðøøøð
 ×)Ò)¨'°6Ñ:Ô:Ð:Ð:Ð:Ð:Ð:Ð:Ð:ð ”*ð ;øøõ ð 	6ð 	6ð 	6ØŒH�NŠNÐ1°1Ñ5Ô5Ð5Ð5Ð5Ð5Ð5Ð5øøøøð	6øøøð �LŠL‰NŒNˆNðØ×(Ò(Ñ*Ô*Ð*Ð*Ð*Ð*Ð*Ð*Ð*Ð*Ð*øÝð ð ð Ø��ðøøøøð �LŠL‰NŒNˆNðØ×(Ò(Ñ*Ô*Ð*Ð*Ð*Ð*Ð*Ð*Ð*Ð*øÝð ð ð Ø�ðøøøøøøsˆ   žA: ¦A Á A: ÁAÁA: ÁAÁ&A: Á9C. Á:
B)ÂB$ÂC. Â$B)Â)C. ÃC Ã
C+Ã*C+Ã.D/ÄDÄD/Ä
D,Ä)D/Ä+D,Ä,D/c           
   ƒ   óŽ  K  — |                      dt          t          j        ¦   «         ¦  «        ¦  «        }|                      d¦  «        }|t          j        k    r=t          |t          d|› �¦  «        ¦  «        }t	          j        ||¦  «        ƒ d{V —† dS |                      d¦  «        }|                      d¦  «        }|                      dg ¦  «        }|                      di ¦  «        }	|  	                    ¦   «          	 |  
                    ||||	¦  «        ƒ d{V —†}
t          j        |
¦  «        r|                      ||
|¦  «        ƒ d{V —† dS t          j        |
¦  «        r|                      ||
|¦  «        ƒ d{V —† dS t!          ||
¦  «        }t	          j        ||¦  «        ƒ d{V —† dS # t"          $ r…}t%          j        ¦   «         }| j                             d	||||¦  «         t          |t-          t          |¦  «        |||¬
¦  «        ¦  «        }t	          j        ||¦  «        ƒ d{V —† Y d}~dS d}~ww xY w)ai  
        Handle a single request message.

        Supports both regular (non-streaming) and streaming responses.
        For streaming, detects if the result is a generator and sends
        chunk messages followed by an end message.

        Args:
            message: Request dict from protocol
            writer: StreamWriter for sending responses
        ÚidÚtypezUnknown message type: NÚapp_pathÚactionÚargsÚkwargszError executing %s.%s: %s
%s)r‰   rŠ   Ú	traceback)ÚgetÚstrÚuuidÚuuid4r   ÚMSG_TYPE_REQUESTr   r
   Úwrite_message_asyncr1   ÚexecuteÚinspectÚisgeneratorÚ_stream_sync_generatorÚ
isasyncgenÚ_stream_async_generatorr   r[   r�   Ú
format_excr    r\   r   )r*   r…   r„   Ú
request_idÚmsg_typeÚresponser‰   rŠ   r‹   rŒ   Úresultr`   Útbs                r   r�   zDirtyWorker.handle_request;  s§  è è € ð —[’[ ¥s­4¬:©<¬<Ñ'8Ô'8Ñ9Ô9ˆ
Ø—;’;˜vÑ&Ô&ˆà•}Ô5Ò5Ð5Ý*ØÝ Ð!D¸(Ð!DÐ!DÑEÔEñô ˆHõ  Ô3°F¸HÑEÔEÐEÐEÐEÐEÐEÐEÐEØˆFà—;’;˜zÑ*Ô*ˆØ—’˜XÑ&Ô&ˆØ�{Š{˜6 2Ñ&Ô&ˆØ—’˜X rÑ*Ô*ˆð 	�Š‰Œˆð	FØŸ<š<¨°&¸$ÀÑGÔGÐGÐGÐGÐGÐGÐGˆFõ Ô" 6Ñ*Ô*ð JØ×1Ò1°*¸fÀfÑMÔMÐMÐMÐMÐMÐMÐMÐMÐMÐMÝÔ# FÑ+Ô+ð JØ×2Ò2°:¸vÀvÑNÔNÐNÐNÐNÐNÐNÐNÐNÐNÐNõ )¨°VÑ<Ô<�Ý#Ô7¸ÀÑIÔIÐIÐIÐIÐIÐIÐIÐIÐIÐIøÝð 		Fð 		Fð 		FÝÔ%Ñ'Ô'ˆBØŒH�NŠNÐ:Ø# V¨Q°ñ4ô 4ð 4å*ØÝ�c !™fœf¨xÀØ(*ð,ñ ,ô ,ñô ˆHõ
  Ô3°F¸HÑEÔEÐEÐEÐEÐEÐEÐEÐEÐEÐEÐEÐEÐEÐEøøøøð		Føøøs&   ÄAF5 Å1F5 Æ+F5 Æ5
IÆ?A:H?È?Ic           	   ƒ   óî  ‡‡
K  — t          ¦   «         Š
ˆ
ˆfd„}	 t          j        ¦   «         }	 |                     | j        |¦  «        ƒ d{V —†}|‰
u rn>t          j        |t          ||¦  «        ¦  «        ƒ d{V —† |                      ¦   «          Œdt          j        |t          |¦  «        ¦  «        ƒ d{V —† n�# t          $ r€}t          j        ¦   «         }| j                             d||¦  «         t          |t!          t#          |¦  «        |¬¦  «        ¦  «        }	t          j        ||	¦  «        ƒ d{V —† Y d}~nd}~ww xY w‰                     ¦   «          dS # ‰                     ¦   «          w xY w)zá
        Stream chunks from a synchronous generator.

        Args:
            request_id: Request ID for the messages
            gen: Sync generator to iterate
            writer: StreamWriter for sending messages
        c                  óH   •— 	 t          ‰¦  «        S # t          $ r ‰ cY S w xY w©N)ÚnextÚStopIteration)Ú
_EXHAUSTEDÚgens   €€r   Ú	_get_nextz5DirtyWorker._stream_sync_generator.<locals>._get_next~  s;   ø€ ð"Ý˜C‘y”yÐ øÝ ð "ð "ð "Ø!Ð!Ð!Ð!ð"øøøs   ƒ ’! !TNúError during streaming: %s
%s©r�   )Úobjectrh   Úget_running_loopÚrun_in_executorr)   r   r“   r   r1   r   r[   r�   rš   r    r\   r   r   r�   rV   )r*   r›   r¦   r„   r§   ÚloopÚchunkr`   rŸ   r�   r¥   s     `       @r   r—   z"DirtyWorker._stream_sync_generatorq  sî  øøè è € õ ‘X”Xˆ
ð	"ð 	"ð 	"ð 	"ð 	"ð 	"ð	ÝÔ+Ñ-Ô-ˆDð
à"×2Ò2°4´>À9ÑMÔMÐMÐMÐMÐMÐMÐM�Ø˜JÐ&Ð&Øå#Ô7ØÕ.¨z¸5ÑAÔAñô ð ð ð ð ð ð ð ð —’‘”�ð
õ  Ô3ØÕ(¨Ñ4Ô4ñô ð ð ð ð ð ð ð ð øõ ð 	Fð 	Fð 	FåÔ%Ñ'Ô'ˆBØŒH�NŠNÐ;¸QÀÑCÔCÐCÝ*ØÝ�c !™fœf°Ð3Ñ3Ô3ñô ˆHõ  Ô3°F¸HÑEÔEÐEÐEÐEÐEÐEÐEÐEÐEÐEÐEÐEÐEøøøøð	Føøøð �IŠI‰KŒKˆKˆKˆKøˆC�IŠI‰KŒKˆKˆKøøøs1   šB B; Â:E Â;
EÃA6E Ä;E Å EÅE ÅE4c           	   ƒ   óz  K  — 	 |2 3 d{V —†}t          j        |t          ||¦  «        ¦  «        ƒ d{V —† |                      ¦   «          ŒE6 t          j        |t	          |¦  «        ¦  «        ƒ d{V —† n�# t
          $ r€}t          j        ¦   «         }| j         	                    d||¦  «         t          |t          t          |¦  «        |¬¦  «        ¦  «        }t          j        ||¦  «        ƒ d{V —† Y d}~nd}~ww xY w|                     ¦   «         ƒ d{V —† dS # |                     ¦   «         ƒ d{V —† w xY w)zä
        Stream chunks from an asynchronous generator.

        Args:
            request_id: Request ID for the messages
            gen: Async generator to iterate
            writer: StreamWriter for sending messages
        Nr¨   r©   )r   r“   r   r1   r   r[   r�   rš   r    r\   r   r   r�   Úaclose)r*   r›   r¦   r„   r®   r`   rŸ   r�   s           r   r™   z#DirtyWorker._stream_async_generator¡  sÌ  è è € ð	Ø"ð ð ð ð ð ð ð �eå#Ô7ØÕ.¨z¸5ÑAÔAñô ð ð ð ð ð ð ð ð —’‘”��ð  #õ  Ô3ØÕ(¨Ñ4Ô4ñô ð ð ð ð ð ð ð ð øõ ð 	Fð 	Fð 	FåÔ%Ñ'Ô'ˆBØŒH�NŠNÐ;¸QÀÑCÔCÐCÝ*ØÝ�c !™fœf°Ð3Ñ3Ô3ñô ˆHõ  Ô3°F¸HÑEÔEÐEÐEÐEÐEÐEÐEÐEÐEÐEÐEÐEÐEøøøøð	Føøøð —*’*‘,”,ÐÐÐÐÐÐÐÐÐø�#—*’*‘,”,ÐÐÐÐÐÐÐÐøøøs;   „A5 †AŒA(A5 Á4D Á5
C?Á?A6C:Ã5D Ã:C?Ã?D ÄD:c           	   ƒ   óÆ  ‡‡‡‡K  — || j         vrt          |¦  «        ‚| j         |         Š| j        j        dk    r| j        j        nd}t	          j        ¦   «         }	 t	          j        |                     | j        ˆˆˆˆfd„¦  «        |¬¦  «        ƒ d{V —†}|S # t          j	        $ r6 | j
                             d|‰|¦  «         t          d|› d‰› d�|¬¦  «        ‚w xY w)	a„  
        Execute an action on a dirty app.

        The action runs in a thread pool executor to avoid blocking the
        asyncio event loop. Execution timeout is enforced using
        ``dirty_timeout`` config.

        Args:
            app_path: Import path of the dirty app
            action: Action name to execute
            args: Positional arguments
            kwargs: Keyword arguments

        Returns:
            Result from the app action

        Raises:
            DirtyAppNotFoundError: If app is not loaded
            DirtyTimeoutError: If execution exceeds timeout
            DirtyAppError: If execution fails
        r   Nc                  ó   •—  ‰‰ g‰¢R i ‰¤ŽS r¢   © )rŠ   r_   r‹   rŒ   s   €€€€r   ú<lambda>z%DirtyWorker.execute.<locals>.<lambda>æ  s!   ø€ ˜C˜C Ð8¨Ð8Ð8Ð8°Ð8Ð8€ r   )Útimeoutz%Execution timeout for %s.%s after %dszExecution of ú.z
 timed out)r&   r   r   r}   rh   r«   Úwait_forr¬   r)   ÚTimeoutErrorr    Úwarningr	   )	r*   r‰   rŠ   r‹   rŒ   rµ   r­   rž   r_   s	     ```   @r   r”   zDirtyWorker.executeÂ  sH  øøøøè è € ð, ˜4œ9Ð$Ð$Ý'¨Ñ1Ô1Ð1àŒi˜Ô!ˆØ,0¬HÔ,BÀQÒ,FÐ,F�$”(Ô(Ð(ÈDˆõ Ô'Ñ)Ô)ˆð	Ý"Ô+Ø×$Ò$Ø”NØ8Ð8Ð8Ð8Ð8Ð8Ð8ñô ð  ðñ ô ð ð ð ð ð ð ˆFð ˆMøÝÔ#ð 		ð 		ð 		àŒH×ÒØ7Ø˜& 'ñô ð õ $Ø= Ð=Ð=¨6Ð=Ð=Ð=Øðñ ô ð ð		øøøs   Á<B ÂAC c                 ó’  — | j         r#| j                              dd¬¦  «         d| _         | j                             ¦   «         D ]h\  }}	 |                     ¦   «          | j                             d|¦  «         Œ6# t          $ r&}| j                             d||¦  «         Y d}~Œad}~ww xY w	 | j	                             ¦   «          n# t          $ r Y nw xY w	 t          j                             | j        ¦  «        rt          j        | j        ¦  «         n# t          $ r Y nw xY w| j                             d| j        ¦  «         dS )zClean up resources on shutdown.FT)ÚwaitÚcancel_futuresNzClosed dirty app: %szError closing dirty app %s: %szDirty worker %s exiting)r)   Úshutdownr&   r5   rV   r    rX   r[   r\   r%   r6   r^   rp   r!   rq   rZ   r   r]   s       r   rm   zDirtyWorker._cleanupö  sx  € ð Œ>ð 	"ØŒN×#Ò#¨¸tÐ#ÑDÔDÐDØ!ˆDŒNð œŸšÑ*Ô*ð 	Jð 	J‰IˆD�#ðJØ—	’	‘”�Ø”—’Ð5°tÑ<Ô<Ð<Ð<øÝð Jð Jð JØ”—’Ð?ÀÀqÑIÔIÐIÐIÐIÐIÐIÐIøøøøðJøøøð	ØŒH�NŠNÑÔÐÐøÝð 	ð 	ð 	ØˆDð	øøøð	ÝŒw�~Š~˜dÔ.Ñ/Ô/ð ,Ý”	˜$Ô*Ñ+Ô+Ð+øøÝð 	ð 	ð 	ØˆDð	øøøð 	Œ�ŠÐ/°´Ñ:Ô:Ð:Ð:Ð:s<   Á
/A:Á:
B*ÂB%Â%B*Â.C Ã
CÃCÃ=D Ä
D$Ä#D$N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__ÚsplitrG   r+   r/   r1   rE   r>   rJ   rS   r?   rB   rl   rv   rs   r�   r—   r™   r”   rm   r³   r   r   r   r   `   sP  € € € € € ðð ð6ð 6Ø,×2Ò2Ñ4Ô4ð6ñ 6ô 6€Gðð ð ð8+ð +ð +ðð ð ð"ð "ð "ðH<ð <ð <ð"<ð <ð <ð!ð !ð !ð
ð ð ð"ð ð ð,ð ð ðB>ð >ð >ðð ð ð64Fð 4Fð 4Fðl.ð .ð .ð`ð ð ðB2ð 2ð 2ðh;ð ;ð ;ð ;ð ;r   r   )rÁ   rh   r•   r6   r   r�   r�   Úgunicornr   Úgunicorn.workers.workertmpr   r_   r   Úerrorsr   r   r	   r
   Úprotocolr   r   r   r   r   r   r³   r   r   ú<module>rÇ      s?  ðð
?ð ?ðB €€€Ø €€€Ø 	€	€	€	Ø €€€Ø Ð Ð Ð Ø €€€à Ð Ð Ð Ð Ð Ø 0Ð 0Ð 0Ð 0Ð 0Ð 0à  Ð  Ð  Ð  Ð  Ð  ðð ð ð ð ð ð ð ð ð ð ð ðð ð ð ð ð ð ð ð ð ð ð ð ð ðr;ð r;ð r;ð r;ð r;ñ r;ô r;ð r;ð r;ð r;r   