§
    ÷žyj£  ã                  óà   — U d Z 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m	Z	m
Z
  ej        e¦  «        ZdZdZ G d„ d¦  «        Zdad	ed
<    ej        ¦   «         Zdd„Zdd„Zddd„ZeZg d¢ZdS )uù  Monitoring emitter: fire-and-forget queue + background dispatcher.

The emitter is the single seam between producers (gateway status hooks, the
diagnostic log handler) and consumers (the OTLP streamers). Its contract is
the hot-path invariant:

    ``emit()`` MUST return in O(microseconds), MUST NOT block on disk/network,
    and MUST NEVER raise into the caller. A monitoring failure is logged
    locally and dropped â€” it can never affect the gateway or a session.

Mechanism:
  * ``emit(event)`` does a non-blocking ``queue.put_nowait`` wrapped in a bare
    except. On a full queue it drops the *oldest* event and counts the drop.
  * A daemon thread drains the queue and fans each batch out to subscribers
    (the OTLP metric/span/log streamers). Each subscriber is fail-isolated â€”
    a slow or raising subscriber never affects the hot path or its peers.

Nothing is persisted here. Monitoring is an egress path, not a local store;
if no subscriber is attached, events simply age out of the ring buffer.
é    )ÚannotationsN)ÚAnyÚDictÚOptionali'  é   c                  ój   — e Zd ZdZddœdd„Zdd„Zdd„Zdd„Zdd„Zdd„Z	dd„Z
ddd„Zdd„Zdd„ZdS )ÚMonitoringEmitterz?Owns the queue, the dispatcher thread, and the subscriber list.T©Úenabledr   ÚboolÚreturnÚNonec               óø   — || _         t          j        t          ¬¦  «        | _        d| _        d| _        t          j        ¦   «         | _	        d| _
        t          j        ¦   «         | _        d | _        g | _        d S )N)Úmaxsizer   F)Ú_enabledÚqueueÚQueueÚ
_MAX_QUEUEÚ_qÚ_droppedÚ_dispatchedÚ	threadingÚEventÚ_stopÚ_startedÚLockÚ_lockÚ_threadÚ_subscribers)Úselfr   s     ú>/home/ragecks/.hermes/hermes-agent/agent/monitoring/emitter.pyÚ__init__zMonitoringEmitter.__init__'   sh   € ØˆŒÝ16´ÅZÐ1PÑ1PÔ1PˆŒØˆŒØˆÔÝ”_Ñ&Ô&ˆŒ
ØˆŒÝ”^Ñ%Ô%ˆŒ
Ø37ˆŒð #%ˆÔÐÐó    Úeventr   c                ó´  — | j         sdS 	 t          |d¦  «        r|                     ¦   «         nt          |¦  «        }|                     dt          j        ¦   «         ¦  «         |                      ¦   «          	 | j         	                    |¦  «         dS # t          j        $ r… 	 | j                             ¦   «          | j                             ¦   «          | xj        dz  c_        | j         	                    |¦  «         n # t          $ r | xj        dz  c_        Y nw xY wY dS Y dS w xY w# t          $ r  t                                dd¬¦  «         Y dS w xY w)z€Enqueue an event. Never blocks, never raises.

        ``event`` may be a dataclass with ``to_dict()`` or a plain dict.
        NÚto_dictÚts_nsé   zmonitoring emit failedT©Úexc_info)r   Úhasattrr&   ÚdictÚ
setdefaultÚtimeÚtime_nsÚ_ensure_startedr   Ú
put_nowaitr   ÚFullÚ
get_nowaitÚ	task_doner   Ú	ExceptionÚloggerÚdebug)r    r$   Úpayloads      r!   ÚemitzMonitoringEmitter.emit5   sƒ  € ð
 Œ}ð 	ØˆFð	BÝ)0°¸	Ñ)BÔ)BÐS�e—m’m‘o”o�oÍÈUÉÌˆGØ×Ò˜w­¬©¬Ñ7Ô7Ð7Ø× Ò Ñ"Ô"Ð"ð
'Ø”×"Ò" 7Ñ+Ô+Ð+Ð+Ð+øÝ”:ð 'ð 'ð 'ð'Ø”G×&Ò&Ñ(Ô(Ð(Ø”G×%Ò%Ñ'Ô'Ð'Ø�M”M QÑ&�M”MØ”G×&Ò& wÑ/Ô/Ð/Ð/øÝ ð 'ð 'ð 'Ø�M”M QÑ&�M”M�M�Mð'øøøð 0Ð/Ð/à!�M�Mð'øøøøõ ð 	Bð 	Bð 	BÝ�LŠLÐ1¸DˆLÑAÔAÐAÐAÐAÐAð	Bøøøs[   ‹A.D- Á:B ÂD*Â&ADÄD*ÄD ÄD*ÄD Ä D*Ä#D- Ä&D- Ä)D*Ä*D- Ä-&EÅEc                ó  — | j         rd S | j        5  | j         r	 d d d ¦  «         d S t          j        | j        dd¬¦  «        | _        | j                             ¦   «          d| _         d d d ¦  «         d S # 1 swxY w Y   d S )Nzhermes-monitoring-dispatchT©ÚtargetÚnameÚdaemon)r   r   r   ÚThreadÚ_runr   Ústart©r    s    r!   r0   z!MonitoringEmitter._ensure_startedO   sô   € ØŒ=ð 	ØˆFØŒZð 	!ð 	!ØŒ}ð Øð	!ð 	!ð 	!ñ 	!ô 	!ð 	!ð 	!ð 	!õ %Ô+Ø”yÐ'CÈDðñ ô ˆDŒLð ŒL×ÒÑ Ô Ð Ø ˆDŒMð	!ð 	!ð 	!ñ 	!ô 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!ð 	!øøøð 	!ð 	!ð 	!ð 	!ð 	!ð 	!s   ‘	A5§AA5Á5A9Á<A9c                ór  — | j                              ¦   «         �s	 | j                             d¬¦  «        }n# t          j        $ r Y ŒHw xY w|g}t          |¦  «        t          k     r[	 |                     | j         	                    ¦   «         ¦  «         n# t          j        $ r Y nw xY wt          |¦  «        t          k     °[	 |  
                    |¦  «         |D ]}| j                             ¦   «          Œn## |D ]}| j                             ¦   «          Œw xY w| j                              ¦   «         �¯d S d S )Ng      à?©Útimeout)r   Úis_setr   Úgetr   ÚEmptyÚlenÚ_DRAIN_BATCHÚappendr3   Ú	_dispatchr4   )r    ÚfirstÚbatchÚ_s       r!   r@   zMonitoringEmitter._run[   sg  € Ø”*×#Ò#Ñ%Ô%ñ 	(ðØœŸš¨C˜Ñ0Ô0��øÝ”;ð ð ð Ø�ðøøøà�GˆEÝ�e‘*”*�|Ò+Ð+ðØ—L’L ¤×!3Ò!3Ñ!5Ô!5Ñ6Ô6Ð6Ð6øÝ”{ð ð ð Ø�Eðøøøõ �e‘*”*�|Ò+Ð+ð
(Ø—’˜uÑ%Ô%Ð%àð (ð (�AØ”G×%Ò%Ñ'Ô'Ð'Ð'ð(ø˜ð (ð (�AØ”G×%Ò%Ñ'Ô'Ð'Ð'ð(øøøð ”*×#Ò#Ñ%Ô%ñ 	(ð 	(ð 	(ð 	(ð 	(s-   œ8 ¸A
Á	A
Á),B ÂB(Â'B(ÃC8 Ã8 Dc                óÞ   — t          | j        ¦  «        D ]:}	  ||¦  «         Œ# t          $ r t                               dd¬¦  «         Y Œ7w xY w| xj        t          |¦  «        z  c_        d S )Nzmonitoring subscriber failedTr)   )Úlistr   r5   r6   r7   r   rI   )r    rN   Úsubs      r!   rL   zMonitoringEmitter._dispatchm   s�   € å˜Ô)Ñ*Ô*ð 	Lð 	LˆCðLØ��E‘
”
�
�
øÝð Lð Lð LÝ—’Ð;Àd�ÑKÔKÐKÐKÐKðLøøøàÐÔ�C ™JœJÑ&ÐÔÐÐs   ˜$¤&AÁAc                óZ   — || j         vr| j                              |¦  «         d| _        dS )z?Register a live batch subscriber (callable(batch: list[dict])).TN)r   rK   r   ©r    Úcallbacks     r!   Ú	subscribezMonitoringEmitter.subscribev   s2   € à˜4Ô,Ð,Ð,ØÔ×$Ò$ XÑ.Ô.Ð.ØˆŒˆˆr#   c                ó~   — 	 | j                              |¦  «         n# t          $ r Y nw xY w| j         s	d| _        d S d S )NF)r   ÚremoveÚ
ValueErrorr   rT   s     r!   ÚunsubscribezMonitoringEmitter.unsubscribe|   sa   € ð	ØÔ×$Ò$ XÑ.Ô.Ð.Ð.øÝð 	ð 	ð 	ØˆDð	øøøàÔ ð 	"Ø!ˆDŒMˆMˆMð	"ð 	"s   ‚ �
*©*ç       @rE   Úfloatc                óÐ   ‡ ‡— |dk    rdS t          j        ¦   «         Šd
ˆˆ fd„}t          j        |dd¬¦  «        }|                     ¦   «          ‰                     |¬	¦  «         dS )zCWait boundedly for queued and in-flight batches to finish dispatch.r   Nr   r   c                 ób   •— ‰j                              ¦   «          ‰                      ¦   «          d S ©N)r   ÚjoinÚset)Úfinishedr    s   €€r!   Ú_wait_for_completionz5MonitoringEmitter.flush.<locals>._wait_for_completionŒ   s#   ø€ ØŒG�LŠL‰NŒNˆNØ�LŠL‰NŒNˆNˆNˆNr#   zhermes-monitoring-flushTr;   rD   ©r   r   )r   r   r?   rA   Úwait)r    rE   rc   Úwaiterrb   s   `   @r!   ÚflushzMonitoringEmitter.flush…   s�   øø€ à�aŠ<ˆ<ØˆFå”?Ñ$Ô$ˆð	ð 	ð 	ð 	ð 	ð 	ð 	õ Ô!Ø'Ø*Øð
ñ 
ô 
ˆð
 	�Š‰ŒˆØ�Š˜gˆÑ&Ô&Ð&Ð&Ð&r#   úDict[str, int]c                óv   — | j                              ¦   «         | j        | j        t	          | j        ¦  «        dœS )N)ÚqueuedÚ
dispatchedÚdroppedÚsubscribers)r   Úqsizer   r   rI   r   rB   s    r!   ÚstatszMonitoringEmitter.stats˜   s7   € à”g—m’m‘o”oØÔ*Ø”}Ý˜tÔ0Ñ1Ô1ð	
ð 
ð 	
r#   c                óŠ   — | j                              ¦   «          | j        �| j                             d¬¦  «         d| _        d S )Nr[   rD   F)r   ra   r   r`   r   rB   s    r!   ÚclosezMonitoringEmitter.close    s@   € ØŒ
�ŠÑÔÐØŒ<Ð#ØŒL×Ò cÐÑ*Ô*Ð*ØˆŒˆˆr#   N)r   r   r   r   ©r$   r   r   r   rd   )r[   )rE   r\   r   r   )r   rh   )Ú__name__Ú
__module__Ú__qualname__Ú__doc__r"   r9   r0   r@   rL   rV   rZ   rg   ro   rq   © r#   r!   r	   r	   $   sø   € € € € € ØIÐIà*.ð %ð %ð %ð %ð %ð %ðBð Bð Bð Bð4
!ð 
!ð 
!ð 
!ð(ð (ð (ð (ð$'ð 'ð 'ð 'ðð ð ð ð"ð "ð "ð "ð'ð 'ð 'ð 'ð 'ð&
ð 
ð 
ð 
ðð ð ð ð ð r#   r	   úOptional[MonitoringEmitter]Ú_EMITTERr   c                 ó˜   — t           �t           S t          5  t           €t          d¬¦  «        a ddd¦  «         n# 1 swxY w Y   t           S )z+Return the process-wide monitoring emitter.NFr
   )ry   Ú_EMITTER_LOCKr	   rw   r#   r!   Úget_emitterr|   ¬   s‰   € õ ÐÝˆÝ	ð 8ð 8ÝÐõ )°Ð7Ñ7Ô7ˆHð	8ð 8ð 8ñ 8ô 8ð 8ð 8ð 8ð 8ð 8ð 8øøøð 8ð 8ð 8ð 8õ
 €Os   –:º>Á>r$   r   r   c                óH   — t          ¦   «                              | ¦  «         dS )z1Module-level convenience: emit via the singleton.N)r|   r9   )r$   s    r!   r9   r9   ¹   s    € å�M„M×Ò�uÑÔÐÐÐr#   Úemitterc                óÀ   — t           5  t          �4| t          ur+	 t                               ¦   «          n# t          $ r Y nw xY w| addd¦  «         dS # 1 swxY w Y   dS )z Swap the singleton (tests only).N)r{   ry   rq   r5   )r~   s    r!   Úreset_emitter_for_testsr€   ¾   s»   € õ 
ð ð ÝÐ Gµ8Ð$;Ð$;ðÝ—’Ñ Ô Ð Ð øÝð ð ð Ø�ðøøøàˆðð ð ñ ô ð ð ð ð ð ð ð øøøð ð ð ð ð ð s0   ˆAš4³A´
A¾AÁ AÁAÁAÁA)r	   ÚTelemetryEmitterr|   r9   r€   )r   r	   rr   r_   )r~   rx   r   r   )rv   Ú
__future__r   Úloggingr   r   r.   Útypingr   r   r   Ú	getLoggerrs   r6   r   rJ   r	   ry   Ú__annotations__r   r{   r|   r9   r€   r�   Ú__all__rw   r#   r!   ú<module>rˆ      s:  ððð ð ð* #Ð "Ð "Ð "Ð "Ð "à €€€Ø €€€Ø Ð Ð Ð Ø €€€Ø &Ð &Ð &Ð &Ð &Ð &Ð &Ð &Ð &Ð &à	ˆÔ	˜8Ñ	$Ô	$€à€
Ø€ð@ð @ð @ð @ð @ñ @ô @ð @ðH )-€Ð ,Ð ,Ð ,Ñ ,Ø�	”Ñ Ô €ð
ð 
ð 
ð 
ðð ð ð ð
	ð 	ð 	ð 	ð 	ð %Ð ðð ð €€€r#   