Ë
    8ø»j7  ã                  ó  — d dl m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	 d dl
mZmZmZ d dlmZ erd dlmZ d d	lmZmZ d d
lmZ d dlmZ  G d„ d«      Z ed«       G d„ de«      «       Z ed«       G d„ de«      «       Zy)é    )ÚannotationsN)ÚTask)ÚUTCÚdatetime)ÚThread)ÚTYPE_CHECKINGÚClassVarÚSelf)Ú
docs_group)ÚTracebackType)Ú	LogClientÚLogClientAsync)ÚHttpResponse)ÚTimeoutc                  ó^   — e Zd ZU dZdZ	 dZded<   	 ddœdd„Zdd	„Zdd
œdd„Z	e
dd„«       Zy)ÚStreamedLogBasez>Base class for streaming and buffering chunked Actor run logs.FÚ
no_timeoutzClassVar[Timeout]Ú_stream_timeoutT©Ú
from_startc               óê   — | j                   rd|_        || _        t        t           «       | _        t        j                  d«      | _        |rd | _        y t        j                  t        ¬«      | _        y )NTs5   (?:\n|^)(\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}\.\d{3}Z))Útz)Ú_force_propagateÚ	propagateÚ
_to_loggerÚlistÚbytesÚ_stream_bufferÚreÚcompileÚ_split_markerr   Únowr   Ú_relevancy_time_limit)ÚselfÚ	to_loggerr   s      úQ/var/www/html/GAP/venv/lib/python3.12/site-packages/apify_client/_streamed_log.pyÚ__init__zStreamedLogBase.__init__#   sU   € Ø× Ò Ø"&ˆIÔØ#ˆŒÜ"¤5™k›mˆÔÜŸZ™ZÐ(aÓbˆÔÙ>H°dˆÕ"ÌhÏlÉlÔ^aÔNbˆÕ"ó    c                ó¤   — |}| j                   j                  |«       t        j                  | j                  |«      r| j                  d¬«       y y )NF©Úinclude_last_part)r   Úappendr   Úfindallr!   Ú_log_buffer_content)r$   ÚdataÚ	new_chunks      r&   Ú_process_new_dataz!StreamedLogBase._process_new_data+   sE   € Øˆ	Ø×Ñ×"Ñ" 9Ô-Ü�:‰:�d×(Ñ(¨)Ô4à×$Ñ$°uÐ$Õ=ð 5r(   r*   c               ó4  — t        j                  | j                  dj                  | j                  «      «      dd }|r|ddd…   }|ddd…   }g | _        n|ddd…   }|ddd…   }|dd | _        t        ||d¬«      D ]—  \  }}|j                  d	«      }|j                  d	«      }| j                  r%t        j                  |«      }	|	| j                  k  rŒY||z   }
| j                  j                  | j                  |
«      |
j                  «       ¬
«       Œ™ y)a  Merge the whole buffer and split it into parts based on the marker.

        Log the messages created from the split parts and remove them from buffer.
        The last part could be incomplete, and so it can be left unprocessed in the buffer until later.
        r(   é   Nr   é   éþÿÿÿF)Ústrictzutf-8)ÚlevelÚmsg)r   Úsplitr!   Újoinr   ÚzipÚdecoder#   r   Úfromisoformatr   ÚlogÚ_guess_log_level_from_messageÚstrip)r$   r+   Ú	all_partsÚmessage_markersÚmessage_contentsÚmarkerÚcontentÚdecoded_markerÚdecoded_contentÚlog_timeÚmessages              r&   r.   z#StreamedLogBase._log_buffer_content2   s&  € ô —H‘H˜T×/Ñ/°·±¸$×:MÑ:MÓ1NÓOÐPQÐPRÐSˆ	ÙØ'¨¨¨1¨™oˆOØ(¨¨¨A¨™ÐØ"$ˆDÕà'¨¨"¨Q¨Ñ/ˆOØ(¨¨2¨a¨Ñ0Ðà"+¨B¨C .ˆDÔä" ?Ð4DÈUÔSò 		h‰OˆF�GØ#Ÿ]™]¨7Ó3ˆNØ%Ÿn™n¨WÓ5ˆOØ×)Ò)Ü#×1Ñ1°.ÓA�Ø˜d×8Ñ8Ò8àØ$ Ñ6ˆGØ�O‰O×Ñ d×&HÑ&HÈÓ&QÐW^×WdÑWdÓWfÐÕgñ		hr(   c                ór   — d}t        j                  «       }|D ]  }|| v sŒ||   c S  t         j                  S )z%Guess the log level from the message.)ÚCRITICALÚFATALÚERRORÚWARNÚWARNINGÚINFOÚDEBUGÚNOTSET)ÚloggingÚgetLevelNamesMappingrP   )rI   Úknown_levelsÚlevel_names_to_levelsr7   s       r&   r?   z-StreamedLogBase._guess_log_level_from_messageN   sG   € ð dˆÜ '× <Ñ <Ó >ÐØ!ò 	4ˆEØ˜ÒØ,¨UÑ3Ò3ð	4ô �|‰|Ðr(   N)r%   úlogging.Loggerr   ÚboolÚreturnÚNone)r/   r   rY   rZ   )r+   rX   rY   rZ   )rI   ÚstrrY   Úint)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   r   Ú__annotations__r'   r1   r.   Ústaticmethodr?   © r(   r&   r   r      sQ   … ÙHàÐØ_à)5€OÐ&Ó5ðð IMõ có>ð @Eõ hð8 ò	ó ñ	r(   r   ÚOtherc                  ót   ‡ — e Zd ZU dZdZded<   	 ddœdˆ fd„Zdd„Zdd	„Zdd
„Z		 	 	 	 	 	 	 	 dd„Z
dd„Zˆ xZS )ÚStreamedLogaÍ  Streams Actor run log output to a Python logger in a background thread.

    The log stream is consumed in a background thread and each log message is forwarded to the provided logger with
    an appropriate log level inferred from the message content.

    Can be used as a context manager, which automatically starts and stops the streaming thread. Alternatively,
    call `start` and `stop` manually. Obtain an instance via `RunClient.get_streamed_log`.
    é   zClassVar[float]Ú_stop_timeout_sTr   c               ó`   •— t         ‰| �  ||¬«       || _        d| _        d| _        d| _        y)aì  Initialize `StreamedLog`.

        Args:
            log_client: The log client used to stream raw log data from the Actor run.
            to_logger: The logger to which the log messages will be forwarded.
            from_start: If `True`, all logs from the start of the Actor run will be streamed. If `False`, only newly
                arrived logs will be streamed. This can be useful for long-running Actors in stand-by mode where only
                recent logs are relevant.
        ©r%   r   NF)Úsuperr'   Ú_log_clientÚ_streaming_threadÚ_log_streamÚ_stop_logging©r$   Ú
log_clientr%   r   Ú	__class__s       €r&   r'   zStreamedLog.__init__m   s7   ø€ ô 	‰Ñ 9¸ÐÔDØ%ˆÔØ04ˆÔØ04ˆÔØ"ˆÕr(   c                ó
  — | j                   r%| j                   j                  «       rt        d«      ‚d| _        t	        j
                  | j                  d¬«      | _         | j                   j                  «        | j                   S )z{Start the streaming thread.

        The caller is responsible for cleanup by calling the `stop` method when done.
        zStreaming thread already activeFT)ÚtargetÚdaemon)rm   Úis_aliveÚRuntimeErrorro   Ú	threadingr   Ú_stream_logÚstart©r$   s    r&   rz   zStreamedLog.start}   sl   € ð
 ×!Ò! d×&<Ñ&<×&EÑ&EÔ&GÜÐ@ÓAÐAØ"ˆÔä!*×!1Ñ!1¸×9IÑ9IÐRVÔ!WˆÔØ×Ñ×$Ñ$Ô&Ø×%Ñ%Ð%r(   c                ó˜  — | j                   st        d«      ‚d| _        | j                  }|�	 |j	                  «        | j                   j                  | j                  ¬«       | j                   j                  «       r| j                  j                  d«       yd| _         y# t
        $ r | j                  j                  d«       Y ŒŠw xY w)a”  Signal the streaming thread to stop logging and wait up to `_stop_timeout_s` for it to finish.

        A thread that outlives the wait is a daemon with `_stop_logging` set, so it exits after at most one more chunk,
        and only then does its buffered tail reach the logger. Its handle is kept while it is alive, so `start` cannot
        revive it beside a second thread on the same buffer.
        zStreaming thread is not activeTNzClosing the log stream failed:)ÚtimeoutzMLog streaming thread outlived the stop timeout; it ends after the next chunk.)rm   rw   ro   rn   ÚcloseÚ	Exceptionr   Ú	exceptionr:   rh   rv   Údebug)r$   Ú
log_streams     r&   ÚstopzStreamedLog.stopŠ   s¶   € ð ×%Ò%ÜÐ?Ó@Ð@Ø!ˆÔà×%Ñ%ˆ
ØÐ!ðLØ× Ñ Ô"ð 	×Ñ×#Ñ#¨D×,@Ñ,@Ð#ÔAØ×!Ñ!×*Ñ*Ô,à�O‰O×!Ñ!Ð"qÕrà%)ˆDÕ"øô ò Là—‘×)Ñ)Ð*JÖKðLús   ®B" Â"$C	ÃC	c                ó&   — | j                  «        | S )zdStart the streaming thread within the context. Exiting the context will finish the streaming thread.©rz   r{   s    r&   Ú	__enter__zStreamedLog.__enter__£   s   € à�
‰
ŒØˆr(   c                ó$   — | j                  «        y)zStop the streaming thread.N©rƒ   ©r$   Úexc_typeÚexc_valÚexc_tbs       r&   Ú__exit__zStreamedLog.__exit__¨   s   € ð 	�	‰	�r(   c                ó$  — 	 | j                   j                  d| j                  ¬«      5 }|s
	 d d d «       y || _        	 | j                  r$	 d | _        	 | j                  d¬«       d d d «       y |j                  «       D ]!  }| j                  |«       | j                  sŒ! n d | _        	 | j                  d¬«       d d d «       y # t        $ r | j                  j                  d«       Y Œ†w xY w# t        $ r | j                  j                  d«       Y ŒYw xY w# d | _        	 | j                  d¬«       w # t        $ r | j                  j                  d«       Y w w xY wxY w# 1 sw Y   y xY w# t        $ r˜}| j                  r!| j                  j                  d|«       Y d }~y | j                   j                  j                  |«      r| j                  j                  d«       n | j                  j                  d«       Y d }~y Y d }~y d }~ww xY w)NT©Úrawr}   r*   ú0Log redirection stopped due to unexpected error:z6Log streaming stopped while `stop` was in progress: %rú8Log streaming stopped: the log stream request timed out.)rl   Ústreamr   rn   ro   r.   r   r   r€   Ú
iter_bytesr1   r�   Ú_http_clientÚis_timeout_errorÚwarning©r$   r‚   r/   Úexcs       r&   ry   zStreamedLog._stream_log®   s  € ð!	^Ø×!Ñ!×(Ñ(¨T¸4×;OÑ;OÐ(ÓPð fÐT^Ù!Ø÷fð fð $.�Ô ðfà×)Ò)Øð (,�DÔ$ðfà×0Ñ0À4Ð0ÔH÷#fð fð !+× 5Ñ 5Ó 7ò "˜Ø×.Ñ.¨tÔ4Ø×-Ó-Ù!ð"ð
 (,�DÔ$ðfà×0Ñ0À4Ð0ÔH÷#fð føô$ %ò fð Ÿ™×1Ñ1Ð2dÖeðfûœ9ò fð Ÿ™×1Ñ1Ð2dÖeðfûð	 (,�DÔ$ðfà×0Ñ0À4Ð0ÕHøÜ$ò fð Ÿ™×1Ñ1Ð2dÖeðfý÷%fð fûô, ò 
	^Ø×!Ò!à—‘×%Ñ%Ð&^Ð`cÔdÜØ×Ñ×,Ñ,×=Ñ=¸cÔBà—‘×'Ñ'Ð(bÕcð —‘×)Ñ)Ð*\×]Ñ]ô dûð
	^úsÔ   ‚'E. ©E"­E. ¶E"¾DÁE"ÁCÁ%E. Á.0DÂDÂ"E"Â*C/Â<E. Ã$C,Ã)E"Ã+C,Ã,E"Ã/$DÄE"ÄDÄE"ÄEÄ"D5Ä4EÄ5$E	ÅEÅE	ÅEÅE"Å"E+Å'E. Å+E. Å.	HÅ7(H
Æ$AH
È
H)rq   r   r%   rW   r   rX   rY   rZ   )rY   r   ©rY   rZ   ©rY   r
   ©rŠ   ztype[BaseException] | Noner‹   zBaseException | NonerŒ   zTracebackType | NonerY   rZ   )r]   r^   r_   r`   rh   ra   r'   rz   rƒ   r†   r�   ry   Ú__classcell__©rr   s   @r&   rf   rf   [   s_   ø… ñð ()€O�_Ó(ðð `d÷ #ó &ó*ó2ð
Ø2ðØ=QðØ[oðà	ó÷"^r(   rf   c                  ób   ‡ — e Zd ZdZddœd
ˆ fd„Zdd„Zdd„Zdd„Z	 	 	 	 	 	 	 	 dd„Zdd	„Z	ˆ xZ
S )ÚStreamedLogAsyncaÛ  Streams Actor run log output to a Python logger in an asyncio task.

    The log stream is consumed in a background asyncio task and each log message is forwarded to the provided logger
    with an appropriate log level inferred from the message content.

    Can be used as an async context manager, which automatically starts and cancels the streaming task. Alternatively,
    call `start` and `stop` manually. Obtain an instance via `RunClientAsync.get_streamed_log`.
    Tr   c               óD   •— t         ‰| �  ||¬«       || _        d| _        y)a÷  Initialize `StreamedLogAsync`.

        Args:
            log_client: The async log client used to stream raw log data from the Actor run.
            to_logger: The logger to which the log messages will be forwarded.
            from_start: If `True`, all logs from the start of the Actor run will be streamed. If `False`, only newly
                arrived logs will be streamed. This can be useful for long-running Actors in stand-by mode where only
                recent logs are relevant.
        rj   N)rk   r'   rl   Ú_streaming_taskrp   s       €r&   r'   zStreamedLogAsync.__init__Þ   s'   ø€ ô 	‰Ñ 9¸ÐÔDØ%ˆÔØ,0ˆÕr(   c                óÌ   — | j                   r%| j                   j                  «       st        d«      ‚t        j                  | j                  «       «      | _         | j                   S )zyStart the streaming task.

        The caller is responsible for cleanup by calling the `stop` method when done.
        zStreaming task already active)r¢   Údonerw   ÚasyncioÚcreate_taskry   r{   s    r&   rz   zStreamedLogAsync.startì   sR   € ð
 ×Ò¨×(<Ñ(<×(AÑ(AÔ(CÜÐ>Ó?Ð?Ü&×2Ñ2°4×3CÑ3CÓ3EÓFˆÔØ×#Ñ#Ð#r(   c              ƒ  óô   K  — | j                   st        d«      ‚| j                   j                  «        	 | j                   ƒ d{  –—†  d| _         y7 Œ# t        j                  $ r Y Œw xY w# d| _         w xY w­w)zStop the streaming task.zStreaming task is not activeN)r¢   rw   Úcancelr¥   ÚCancelledErrorr{   s    r&   rƒ   zStreamedLogAsync.stopö   sr   è ø€ à×#Ò#ÜÐ=Ó>Ð>à×Ñ×#Ñ#Ô%ð	(Ø×&Ñ&×&Ð&ð $(ˆDÕ ð	 'ùÜ×%Ñ%ò 	Ùð	ûð $(ˆDÕ üsF   ‚2A8µA ÁAÁA Á	A8ÁA ÁA)Á&A, Á(A)Á)A, Á,	A5Á5A8c              ƒ  ó.   K  — | j                  «        | S ­w)z`Start the streaming task within the context. Exiting the context will cancel the streaming task.r…   r{   s    r&   Ú
__aenter__zStreamedLogAsync.__aenter__  s   è ø€ à�
‰
ŒØˆùs   ‚c              ƒ  ó@   K  — | j                  «       ƒ d{  –—†  y7 Œ­w)zCancel the streaming task.Nrˆ   r‰   s       r&   Ú	__aexit__zStreamedLogAsync.__aexit__  s   è ø€ ð �i‰i‹k×Òús   ‚–—c              ƒ  ó8  K  — 	 | j                   j                  d| j                  ¬«      4 ƒd {  –—† }|s	 d d d «      ƒd {  –—†  y 	 |j                  «       2 3 d {  –—† }| j	                  |«       Œ7 ŒD7 Œ37 Œ6 	 	 | j                  d¬«       nl# t        $ r | j                  j                  d«       Y nFw xY w# 	 | j                  d¬«       w # t        $ r | j                  j                  d«       Y w w xY wxY wd d d «      ƒd {  –—†7   y # 1 ƒd {  –—†7  sw Y   y xY w# t        $ rk}| j                   j                  j                  |«      r| j                  j                  d«       n | j                  j                  d«       Y d }~y Y d }~y d }~ww xY w­w)NTr�   r*   r‘   r’   )rl   r“   r   Úaiter_bytesr1   r.   r   r   r€   r•   r–   r—   r˜   s       r&   ry   zStreamedLogAsync._stream_log  s‹  è ø€ ð	^Ø×'Ñ'×.Ñ.°4À×AUÑAUÐ.ÓV÷ fð fÐZdÙ!Ø÷f÷ fð fð
fØ&0×&<Ñ&<Ó&>÷ 5ð 5˜dØ×.Ñ.¨tÕ4ðføð føð5øÑ&>ðfà×0Ñ0À4Ð0ÕHøÜ$ò fð Ÿ™×1Ñ1Ð2dÖeðfûðfà×0Ñ0À4Ð0ÕHøÜ$ò fð Ÿ™×1Ñ1Ð2dÖeðfý÷f÷ f÷ f÷ fó fûô ò 	^Ø×Ñ×,Ñ,×=Ñ=¸cÔBà—‘×'Ñ'Ð(bÕcð —‘×)Ñ)Ð*\×]Ñ]ô dûð	^üs  ‚F„+D# ¯A4°D# ³D·D# ÁA6ÁD# ÁFÁ	B:ÁA:ÁA8ÁA:Á!B:Á4D# Á6D# Á8A:Á:B:Á=BÂDÂ$B7Â4DÂ6B7Â7DÂ:C9Â<CÃC9Ã$C6	Ã3C9Ã5C6	Ã6C9Ã9DÃ<D# ÄD
ÄD# ÄFÄD ÄDÄD ÄD# ÄFÄ D# Ä#	FÄ,AFÆ
FÆFÆF)rq   r   r%   rW   r   rX   rY   rZ   )rY   r   rš   r›   rœ   )r]   r^   r_   r`   r'   rz   rƒ   r«   r­   ry   r�   rž   s   @r&   r    r    Ó   sN   ø„ ñð ei÷ 1ó$ó(óð
Ø2ðØ=QðØ[oðà	ó÷^r(   r    )Ú
__future__r   r¥   rS   r   rx   r   r   r   r   Útypingr   r	   r
   Úapify_client._docsr   Útypesr   Úapify_client._resource_clientsr   r   Úapify_client.http_clientsr   Úapify_client.typesr   r   rf   r    rc   r(   r&   ú<module>r·      s�   ðÝ "ã Û Û 	Û Ý ß "Ý ß 0Ñ 0å )áÝ#çHÝ6Ý*÷Bñ BñJ ˆGÓôt^�/ó t^ó ðt^ñn ˆGÓôP^�ó P^ó ñP^r(   