Ë
    èlÖi'V  ã                   óŽ  — 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mZmZm	Z	m
Z
mZ d dlmZ d dlmZ d dlmZmZmZmZmZmZmZ d dlmZ d dlmZmZ  e j:                  «         e j<                  e«      Z e jC                  e jD                  «       dZ#d	Z$d
dhZ%ddhZ& G d„ d«      Z' G d„ de«      Z( G d„ de(«      Z) G d„ de*«      Z+y)é    N)ÚOptionalÚUnionÚTupleÚDictÚcastÚCallable)ÚQueue)ÚSelfDescribingJson)ÚPayloadDictÚPayloadDictListÚHttpProtocolÚMethodÚSuccessCallbackÚFailureCallbackÚEmitterProtocol)Úone_of)Ú
EventStoreÚInMemoryEventStoreé
   zAiglu:com.snowplowanalytics.snowplow/payload_data/jsonschema/1-0-4ÚhttpÚhttpsÚgetÚpostc                   ó2   — e Zd ZU eed<   eed<   dedefd„Zy)Ú	Requesterr   r   c                 ó8   — t        | d|«       t        | d|«       y )Nr   r   )Úsetattr)Úselfr   r   s      ú^/var/www/html/strategist-ai/venv_dbt/lib/python3.12/site-packages/snowplow_tracker/emitters.pyÚ__init__zRequester.__init__9   s   € ô 	��f˜dÔ#Ü��e˜SÕ!ó    N)Ú__name__Ú
__module__Ú__qualname__r   Ú__annotations__r    © r!   r   r   r   5   s   … Ø
ƒNØ	ƒMð"˜Xð "¨Hô "r!   r   c                   óü  — e Zd ZdZddddddddddi ddfdededee   d	ed
ee   dee	   dee
   dee   deeeeeef   f      dedee   deeef   dee   deej$                     ddfd„Ze	 	 	 d/dededee   d	edef
d„«       Zdeddfd„Zdefd„Zd0d„Zdedefd„Zdedefd„Zd0d„Zededefd „«       Zd!eddfd"„Zd#eddfd$„Z d#eddfd%„Z!d0d&„Z"ed'eddfd(„«       Z#dedefd)„Z$d0d*„Z%d0d+„Z&d0d,„Z'd0d-„Z(d0d.„Z)y)1ÚEmitterzl
    Synchronously send Snowplow events to a Snowplow collector
    Supports both GET and POST requests
    r   Nr   é<   ÚendpointÚprotocolÚportÚmethodÚ
batch_sizeÚ
on_successÚ
on_failureÚ
byte_limitÚrequest_timeoutÚmax_retry_delay_secondsÚbuffer_capacityÚcustom_retry_codesÚevent_storeÚsessionÚreturnc                 ó
  — t        |t        «       t        |t        «       t        j	                  ||||«      | _        || _        |€$|€t        t        ¬«      }nt        |t        ¬«      }|| _	        |€|dk(  rt        }nd}|�||kD  r|}|| _        || _        |€dnd| _        |	| _        || _        || _        t#        j$                  «       | _        t)        | d¬«      | _        t)        | d	¬«      | _        |
| _        d| _        || _        t        j5                  d
| j
                  z   «       |€/t7        t8        j:                  t8        j<                  ¬«      | _        yt7        |j:                  |j<                  ¬«      | _        y)a§  
        :param endpoint:    The collector URL. If protocol is not set in endpoint it will automatically set to "https://" - this is done automatically.
        :type  endpoint:    string
        :param protocol:    The protocol to use - http or https. Defaults to https.
        :type  protocol:    protocol
        :param port:        The collector port to connect to
        :type  port:        int | None
        :param method:      The HTTP request method. Defaults to post.
        :type  method:      method
        :param batch_size:  The maximum number of queued events before the buffer is flushed. Default is 10.
        :type  batch_size:  int | None
        :param on_success:  Callback executed after every HTTP request in a flush has status code 200
                            Gets passed one argument, an array of dictionaries corresponding to the sent events' payloads
        :type  on_success:  function | None
        :param on_failure:  Callback executed if at least one HTTP request in a flush has status code other than 200
                            Gets passed two arguments:
                            1) The number of events which were successfully sent
                            2) An array of dictionaries corresponding to the unsent events' payloads
        :type  on_failure:  function | None
        :param byte_limit:  The size event list after reaching which queued events will be flushed
        :type  byte_limit:  int | None
        :param request_timeout: Timeout for the HTTP requests. Can be set either as single float value which
                                 applies to both "connect" AND "read" timeout, or as tuple with two float values
                                 which specify the "connect" and "read" timeouts separately
        :type request_timeout:  float | tuple | None
        :param max_retry_delay_seconds:     Set the maximum time between attempts to send failed events to the collector. Default 60 seconds
        :type max_retry_delay_seconds:      int
        :param buffer_capacity: The maximum capacity of the event buffer.
                                When the buffer is full new events are lost.
        :type buffer_capacity: int
        :param  custom_retry_codes: Set custom retry rules for HTTP status codes received in emit responses from the Collector.
                                    By default, retry will not occur for status codes 400, 401, 403, 410 or 422. This can be overridden here.
                                    Note that 2xx codes will never retry as they are considered successful.
        :type   custom_retry_codes: dict
        :param  event_store:    Stores the event buffer and buffer capacity. Default is an InMemoryEventStore object with buffer_capacity of 10,000 events.
        :type   event_store:    EventStore | None
        :param  session:    Persist parameters across requests by using a session object
        :type   session:    requests.Session | None
        N)Úlogger)r4   r:   r   é   r   T)ÚemitterÚ	repeatingFz"Emitter initialized with endpoint )r   r   ) r   Ú	PROTOCOLSÚMETHODSr(   Úas_collector_urir*   r-   r   r:   r6   ÚDEFAULT_MAX_LENGTHr.   r1   Úbytes_queuedr2   r/   r0   Ú	threadingÚRLockÚlockÚ
FlushTimerÚtimerÚretry_timerr3   Úretry_delayr5   Úinfor   Úrequestsr   r   Úrequest_method)r   r*   r+   r,   r-   r.   r/   r0   r1   r2   r3   r4   r5   r6   r7   s                  r   r    zEmitter.__init__G   sP  € ôp 	ˆxœÔ#Üˆv”wÔä×0Ñ0°¸8ÀTÈ6ÓRˆŒàˆŒàÐØÐ&Ü0¼Ô?‘ä0Ø$3¼Fô�ð 'ˆÔàÐØ˜ÒÜ/‘
à�
àÐ&¨:¸Ò+GØ(ˆJà$ˆŒØ$ˆŒØ$.Ð$6™D¸AˆÔØ.ˆÔà$ˆŒØ$ˆŒä—O‘OÓ%ˆŒ	ä¨¸Ô=ˆŒ
Ü%¨d¸eÔDˆÔà'>ˆÔ$Ø./ˆÔà"4ˆÔÜ�‰Ð8¸4¿=¹=ÑHÔIàˆ?Ü"+´·±ÄHÇLÁLÔ"QˆDÕä"+°·±À7Ç;Á;Ô"OˆDÕr!   c                 ó>  — t        | «      dk  rt        d«      ‚| j                  d«      } | j                  d«      d   t        v r)| j                  d«      }t        t        |d   «      }|d   } |dk(  rd}nd}|€|dz   | z   |z   S |dz   | z   d	z   t        |«      z   |z   S )
aª  
        :param endpoint:  The raw endpoint provided by the user
        :type  endpoint:  string
        :param protocol:  The protocol to use - http or https
        :type  protocol:  protocol
        :param port:      The collector port to connect to
        :type  port:      int | None
        :param method:    Either `get` or `post` HTTP method
        :type  method:    method
        :rtype:           string
        r;   zNo endpoint provided.ú/z://r   r   z/iz#/com.snowplowanalytics.snowplow/tp2ú:)ÚlenÚ
ValueErrorÚrstripÚsplitr>   r   r   Ústr)r*   r+   r,   r-   Úendpoint_arrÚpaths         r   r@   zEmitter.as_collector_uri±   sµ   € ô$ ˆx‹=˜1ÒÜÐ4Ó5Ð5à—?‘? 3Ó'ˆà�>‰>˜%Ó  Ñ#¤yÑ0Ø#Ÿ>™>¨%Ó0ˆLÜœL¨,°q©/Ó:ˆHØ# A‘ˆHà�UŠ?Ø‰Dà8ˆDØˆ<Ø˜eÑ# hÑ.°Ñ5Ð5à˜eÑ# hÑ.°Ñ4´s¸4³yÑ@À4ÑGÐGr!   Úpayloadc                 ó¸  — | j                   5  | j                  �'| xj                  t        t        |«      «      z  c_        | j                  dk(  r7| j
                  j                  |D �ci c]  }|t        ||   «      “Œ c}«       n| j
                  j                  |«       | j                  «       r| j                  «        ddd«       yc c}w # 1 sw Y   yxY w)zØ
        Adds an event to the buffer.
        If the maximum size has been reached, flushes the buffer.

        :param payload:   The name-value pairs for the event
        :type  payload:   dict(string:\*)
        Nr   )	rE   rB   rP   rT   r-   r6   Ú	add_eventÚreached_limitÚflush)r   rW   Úkeys      r   ÚinputzEmitter.inputÖ   s³   € ð �Y‰Yñ 
	Ø× Ñ Ð,Ø×!Ò!¤S¬¨W«Ó%6Ñ6Õ!à�{‰{˜fÒ$Ø× Ñ ×*Ñ*ÈgÖ+VÀs¨C´°W¸S±\Ó1BÑ,BÒ+VÕWà× Ñ ×*Ñ*¨7Ô3à×!Ñ!Ô#Ø—
‘
”÷
	ð 
	ùò
 ,W÷
	ð 
	ús   �ACÁ)C
Á?ACÃCÃCc                 óô   — | j                   €'| j                  j                  «       | j                  k\  S | j                  xs d| j                   k\  xs' | j                  j                  «       | j                  k\  S )zW
        Checks if event-size or bytes limit are reached

        :rtype: bool
        r   )r1   r6   Úsizer.   rB   ©r   s    r   rZ   zEmitter.reached_limitê   so   € ð �?‰?Ð"Ø×#Ñ#×(Ñ(Ó*¨d¯o©oÑ=Ð=ð ×!Ñ!Ò& QØ—‘ñ!ò Oà$(×$4Ñ$4×$9Ñ$9Ó$;¸t¿¹Ñ$NðOr!   c                 ó
  — | j                   5  | j                  j                  «       r
	 ddd«       y| j                  j	                  «       }| j                  |«       | j                  �d| _        ddd«       y# 1 sw Y   yxY w)zB
        Sends all events in the buffer to the collector.
        Nr   )rE   rH   Ú	is_activer6   Úget_events_batchÚsend_eventsrB   )r   rd   s     r   r[   zEmitter.flush÷   sw   € ð �Y‰Yñ 	&Ø×Ñ×)Ñ)Ô+Ø÷	&ð 	&ð ×*Ñ*×;Ñ;Ó=ˆKØ×Ñ˜[Ô)Ø× Ñ Ð,Ø$%�Ô!÷	&÷ 	&ñ 	&ús   �A9²>A9Á9BÚdatac                 ód  — t         j                  d| j                  z  «       t         j                  d|z  «       	 | j                  j                  | j                  |ddi| j                  ¬«      }|j                  S # t        j                  $ r}t         j                  |«       Y d}~yd}~ww xY w)zZ
        :param data:  The array of JSONs to be sent
        :type  data:  string
        zSending POST request to %s...úPayload: %szContent-Typezapplication/json; charset=utf-8)re   ÚheadersÚtimeoutNéÿÿÿÿ)r:   rJ   r*   ÚdebugrL   r   r2   rK   ÚRequestExceptionÚwarningÚstatus_code)r   re   ÚrÚes       r   Ú	http_postzEmitter.http_post  sš   € ô
 	�‰Ð3°d·m±mÑCÔDÜ�‰�] TÑ)Ô*ð		Ø×#Ñ#×(Ñ(Ø—‘ØØ'Ð)JÐKØ×,Ñ,ð	 )ó ˆAð �}‰}Ðøô	 ×(Ñ(ò 	Ü�N‰N˜1ÔÜûð	ús   ¼5A= Á=B/ÂB*Â*B/c                 ó^  — t         j                  d| j                  z  «       t         j                  d|z  «       	 | j                  j                  | j                  || j                  ¬«      }|j                  S # t        j                  $ r}t         j                  |«       Y d}~yd}~ww xY w)z`
        :param payload:  The event properties
        :type  payload:  dict(string:\*)
        zSending GET request to %s...rg   )Úparamsri   Nrj   )r:   rJ   r*   rk   rL   r   r2   rK   rl   rm   rn   )r   rW   ro   rp   s       r   Úhttp_getzEmitter.http_get  s�   € ô
 	�‰Ð2°T·]±]ÑBÔCÜ�‰�] WÑ,Ô-ð	Ø×#Ñ#×'Ñ'Ø—‘ g°t×7KÑ7Kð (ó ˆAð �}‰}Ðøô	 ×(Ñ(ò 	Ü�N‰N˜1ÔÜûð	ús   ¼2A: Á:B,ÂB'Â'B,c                 óx   — t         j                  d«       | j                  «        t         j                  d«       y)z€
        Calls the flush method of the base Emitter class.
        This is guaranteed to be blocking, not asynchronous.
        zStarting synchronous flush...zFinished synchronous flushN)r:   rk   r[   rJ   r`   s    r   Ú
sync_flushzEmitter.sync_flush(  s'   € ô
 	�‰Ð4Ô5Ø�
‰
ŒÜ�‰Ð0Õ1r!   rn   c                 ó"   — d| cxk  xr dk  S c S )zz
        :param status_code:  HTTP status code
        :type  status_code:  int
        :rtype:              bool
        éÈ   i,  r&   )rn   s    r   Úis_good_status_codezEmitter.is_good_status_code1  s   € ð �kÖ' CÑ'Ð'Ñ'Ð'r!   Úevtsc                 ó˜  — t        |«      dkD  �r¦t        j                  dt        |«      z  «       t        j	                  |«       g }g }| j
                  dk(  rRt        t        |«      j                  «       }| j                  |«      }t        j                  |«      }|r||z  }nQ||z  }nK| j
                  dk(  r<|D ]7  }| j                  |«      }t        j                  |«      }|r||gz  }Œ2||gz  }Œ9 | j                  �t        |«      dkD  r| j                  |«       | j                  �)t        |«      dkD  r| j                  t        |«      |«       | j                  «      r"| j                  «        | j!                  |«       y| j"                  j%                  |d«       | j'                  «        yt        j                  d«       y)zd
        :param evts: Array of events to be sent
        :type  evts: list(dict(string:\*))
        r   zAttempting to send %s eventsr   r   NFz$Skipping flush since buffer is empty)rP   r:   rJ   r(   Úattach_sent_timestampr-   r
   ÚPAYLOAD_DATA_SCHEMAÚ	to_stringrq   ry   rt   r/   r0   Ú_should_retryÚ_set_retry_delayÚ_retry_failed_eventsr6   ÚcleanupÚ_reset_retry_delay)r   rz   Úsuccess_eventsÚfailure_eventsre   rn   Úrequest_succeededÚevts           r   rd   zEmitter.send_events:  sŽ  € ô
 ˆt‹9�q‹=Ü�K‰KÐ6¼¸T»ÑBÔCä×)Ñ)¨$Ô/ØˆNØˆNà�{‰{˜fÒ$Ü)Ô*=¸tÓD×NÑNÓP�Ø"Ÿn™n¨TÓ2�Ü$+×$?Ñ$?ÀÓ$LÐ!Ù$Ø" dÑ*‘Nà" dÑ*‘Nà—‘ Ò%Øò 0�CØ"&§-¡-°Ó"4�KÜ(/×(CÑ(CÀKÓ(PÐ%á(Ø&¨3¨%Ñ/™à&¨3¨%Ñ/™ð0ð �‰Ð*¬s°>Ó/BÀQÒ/FØ—‘ Ô/Ø�‰Ð*¬s°>Ó/BÀQÒ/FØ—‘¤ NÓ 3°^ÔDà×!Ñ! +Ô.Ø×%Ñ%Ô'Ø×)Ñ)¨.Õ9à× Ñ ×(Ñ(¨¸Ô?Ø×'Ñ'Õ)ä�K‰KÐ>Õ?r!   ri   c                 ó<   — | j                   j                  |¬«       y)z�
        Set an interval at which failed events will be retried

        :param timeout:   interval in seconds
        :type  timeout:   int | float
        ©ri   N)rH   Ústart©r   ri   s     r   Ú_set_retry_timerzEmitter._set_retry_timerg  s   € ð 	×Ñ×Ñ wÐÕ/r!   c                 ó<   — | j                   j                  |¬«       y)z™
        Set an interval at which the buffer will be flushed
        :param timeout:   interval in seconds
        :type  timeout:   int | float
        r‰   N)rG   rŠ   r‹   s     r   Úset_flush_timerzEmitter.set_flush_timerp  s   € ð 	�
‰
×Ñ ÐÕ)r!   c                 ó8   — | j                   j                  «        y)z0
        Abort automatic async flushing
        N)rG   Úcancelr`   s    r   Úcancel_flush_timerzEmitter.cancel_flush_timerx  s   € ð 	�
‰
×ÑÕr!   Úeventsc                 ó:   — dt         ddfd„}| D ]
  } ||«       Œ y)zÝ
        Attach (by mutating in-place) current timestamp in milliseconds
        as `stm` param

        :param events: Array of events to be sent
        :type  events: list(dict(string:\*))
        :rtype: None
        rp   r8   Nc           	      óx   — | j                  dt        t        t        j                  «       «      dz  «      i«       y )NÚstmiè  )ÚupdaterT   ÚintÚtime)rp   s    r   r–   z-Emitter.attach_sent_timestamp.<locals>.update‰  s(   € Ø�H‰H�eœS¤¤T§Y¡Y£[Ó!1°DÑ!8Ó9Ð:Õ;r!   )r   )r’   r–   Úevents      r   r|   zEmitter.attach_sent_timestamp~  s-   € ð	<”kð 	< dó 	<ð ò 	ˆEÙ�5�Mñ	r!   c                 óŒ   — t         j                  |«      ry|| j                  j                  «       v r| j                  |   S |dvS )z 
        Checks if a request should be retried

        :param  status_code: Response status code
        :type   status_code: int
        :rtype: bool
        F)i�  i‘  i“  iš  i¦  )r(   ry   r5   Úkeys)r   rn   s     r   r   zEmitter._should_retry�  sI   € ô ×&Ñ& {Ô3Øà˜$×1Ñ1×6Ñ6Ó8Ñ8Ø×*Ñ*¨;Ñ7Ð7àÐ";Ð;Ð;r!   c                 ó‚   — t        j                   «       }t        | j                  dz  |z   | j                  «      | _        y)z5
        Sets a delay to retry failed events
        é   N)ÚrandomÚminrI   r3   )r   Úrandom_noises     r   r€   zEmitter._set_retry_delayŸ  s7   € ô —}‘}“ˆÜØ×Ñ˜qÑ  <Ñ/°×1MÑ1Mó
ˆÕr!   c                 ó   — d| _         y)z)
        Resets retry delay to 0
        r   N)rI   r`   s    r   rƒ   zEmitter._reset_retry_delay¨  s   € ð ˆÕr!   c                 ór   — | j                   j                  |d«       | j                  | j                  «       y)z‹
        Adds failed events back to the buffer to retry

        :param  failed_events: List of failed events
        :type   List
        TN)r6   r‚   rŒ   rI   )r   Úfailed_eventss     r   r�   zEmitter._retry_failed_events®  s.   € ð 	×Ñ× Ñ  °Ô5Ø×Ñ˜d×.Ñ.Õ/r!   c                 ó8   — | j                   j                  «        y)z'
        Cancels a retry timer
        N)rH   r�   r`   s    r   Ú_cancel_retry_timerzEmitter._cancel_retry_timer¸  s   € ð 	×Ñ×ÑÕ!r!   c                  ó   — y ©Nr&   r`   s    r   Úasync_flushzEmitter.async_flush¿  s   € Ør!   )r   Nr   ©r8   N)*r"   r#   r$   Ú__doc__rT   r   r   r—   r   r   r   r   Úfloatr   r   Úboolr   rK   ÚSessionr    Ústaticmethodr@   r   r]   rZ   r[   rq   rt   rv   ry   r   rd   rŒ   rŽ   r‘   r|   r   r€   rƒ   r�   r¥   r¨   r&   r!   r   r(   r(   A   s�  „ ñð ")Ø"ØØ$(Ø04Ø04Ø$(ØGKØ')Ø)-Ø.0Ø,0Ø.2ñhPàðhPð ðhPð �s‰mð	hPð
 ðhPð ˜S‘MðhPð ˜_Ñ-ðhPð ˜_Ñ-ðhPð ˜S‘MðhPð " %¨¨u°U¸E°\Ñ/BÐ(BÑ"CÑDðhPð "%ðhPð " #™ðhPð !  d ™OðhPð ˜jÑ)ðhPð ˜(×*Ñ*Ñ+ðhPð  
ó!hPðT ð ")Ø"Øñ	"HØð"Hàð"Hð �s‰mð"Hð ð	"Hð
 
ò"Hó ð"HðH˜[ð ¨Tó ð(O˜tó Oó
&ð˜cð  có ð( ð °ó ó"2ð ð(¨ð (°ò (ó ð(ð+@ ð +@°Dó +@ðZ0¨ð 0°$ó 0ð* uð *°ó *óð ð oð ¸$ò ó ðð <¨ð <°ó <ó 
óó0ó"ôr!   r(   c            !       ó  ‡ — e Zd ZdZdddddddddddi ddfdeded	ee   d
edee   dee	   dee
   dedee   deeeeeef   f      dedee   deeef   dee   deej$                     ddf ˆ fd„Zdd„Zdd„Zdd„Zˆ xZS )ÚAsyncEmitterz;
    Uses threads to send HTTP requests asynchronously
    r   Nr   r;   r)   r*   r+   r,   r-   r.   r/   r0   Úthread_countr1   r2   r3   r4   r5   r6   r7   r8   c                 óô   •— t         t        | �  ||||||||	|
|||||¬«       t        «       | _        t        |«      D ]9  }t        j                  | j                  ¬«      }d|_	        |j                  «        Œ; y)aè  
        :param endpoint:    The collector URL. If protocol is not set in endpoint it will automatically set to "https://" - this is done automatically.
        :type  endpoint:    string
        :param protocol:    The protocol to use - http or https. Defaults to http.
        :type  protocol:    protocol
        :param port:        The collector port to connect to
        :type  port:        int | None
        :param method:      The HTTP request method
        :type  method:      method
        :param batch_size: The maximum number of queued events before the buffer is flushed. Default is 10.
        :type  batch_size: int | None
        :param on_success:  Callback executed after every HTTP request in a flush has status code 200
                            Gets passed one argument, an array of dictionaries corresponding to the sent events' payloads
        :type  on_success:  function | None
        :param on_failure:  Callback executed if at least one HTTP request in a flush has status code other than 200
                            Gets passed two arguments:
                            1) The number of events which were successfully sent
                            2) An array of dictionaries corresponding to the unsent events' payloads
        :type  on_failure:  function | None
        :param thread_count: Number of worker threads to use for HTTP requests
        :type  thread_count: int
        :param byte_limit:  The size event list after reaching which queued events will be flushed
        :type  byte_limit:  int | None
        :param max_retry_delay_seconds:     Set the maximum time between attempts to send failed events to the collector. Default 60 seconds
        :type max_retry_delay_seconds:      int
        :param buffer_capacity: The maximum capacity of the event buffer.
                                When the buffer is full new events are lost.
        :type buffer_capacity: int
        :param  event_store:    Stores the event buffer and buffer capacity. Default is an InMemoryEventStore object with buffer_capacity of 10,000 events.
        :type   event_store:    EventStore
        :param  session:    Persist parameters across requests by using a session object
        :type   session:    requests.Session | None
        )r*   r+   r,   r-   r.   r/   r0   r1   r2   r3   r4   r5   r6   r7   )ÚtargetTN)Úsuperr°   r    r	   ÚqueueÚrangerC   ÚThreadÚconsumeÚdaemonrŠ   )r   r*   r+   r,   r-   r.   r/   r0   r±   r1   r2   r3   r4   r5   r6   r7   ÚiÚtÚ	__class__s                     €r   r    zAsyncEmitter.__init__È  s‡   ø€ ôf 	Œl˜DÑ*ØØØØØ!Ø!Ø!Ø!Ø+Ø$;Ø+Ø1Ø#Øð 	+ô 	
ô  "›GˆŒ
Ü�|Ó$ò 	ˆAÜ× Ñ ¨¯©Ô5ˆAØˆAŒHØ�G‰G�Iñ	r!   c                 ó–   — 	 | j                  «        | j                  j                  «        | j                  j	                  «       dk  ry ŒI)Nr;   )r[   rµ   Újoinr6   r_   r`   s    r   rv   zAsyncEmitter.sync_flush  s;   € ØØ�J‰JŒLØ�J‰J�O‰OÔØ×Ñ×$Ñ$Ó&¨Ò*Øð	 r!   c                 óÒ   — | j                   5  | j                  j                  | j                  j	                  «       «       | j
                  �d| _        ddd«       y# 1 sw Y   yxY w)z‡
        Removes all dead threads, then creates a new thread which
        executes the flush method of the base Emitter class
        Nr   )rE   rµ   Úputr6   rc   rB   r`   s    r   r[   zAsyncEmitter.flush  sS   € ð
 �Y‰Yñ 	&Ø�J‰J�N‰N˜4×+Ñ+×<Ñ<Ó>Ô?Ø× Ñ Ð,Ø$%�Ô!÷	&÷ 	&ñ 	&ús   �AAÁA&c                 ó�   — 	 | j                   j                  «       }| j                  |«       | j                   j                  «        ŒFr§   )rµ   r   rd   Ú	task_done)r   rz   s     r   r¸   zAsyncEmitter.consume"  s8   € ØØ—:‘:—>‘>Ó#ˆDØ×Ñ˜TÔ"Ø�J‰J× Ñ Ô"ð r!   r©   )r"   r#   r$   rª   rT   r   r   r—   r   r   r   r   r«   r   r   r¬   r   rK   r­   r    rv   r[   r¸   Ú__classcell__)r¼   s   @r   r°   r°   Ã  sB  ø„ ñð "(Ø"ØØ$(Ø04Ø04ØØ$(ØGKØ')Ø)-Ø.0Ø,0Ø.2ñ!GàðGð ðGð �s‰mð	Gð
 ðGð ˜S‘MðGð ˜_Ñ-ðGð ˜_Ñ-ðGð ðGð ˜S‘MðGð " %¨¨u°U¸E°\Ñ/BÐ(BÑ"CÑDðGð "%ðGð " #™ðGð !  d ™OðGð ˜jÑ)ðGð  ˜(×*Ñ*Ñ+ð!Gð" 
õ#GóRó&÷#r!   r°   c                   ód   — e Zd ZdZdedefd„Zdedefd„Zdd	„Z	defd
„Z
deddfd„Zdeddfd„Zy)rF   zO
    Internal class used by the Emitter to schedule flush calls for later.
    r<   r=   c                 ó`   — || _         || _        d | _        t        j                  «       | _        y r§   )r<   r=   rG   rC   rD   rE   )r   r<   r=   s      r   r    zFlushTimer.__init__.  s%   € ØˆŒØ"ˆŒØ04ˆŒ
Ü—O‘OÓ%ˆ�	r!   ri   r8   c                 ó˜   — | j                   5  | j                  �
	 d d d «       y| j                  |¬«       	 d d d «       y# 1 sw Y   y xY w)NFr‰   T)rE   rG   Ú_schedule_timerr‹   s     r   rŠ   zFlushTimer.start4  sK   € Ø�Y‰Yñ 	Ø�z‰zÐ%Ø÷	ð 	ð ×$Ñ$¨WÐ$Ô5Ø÷	÷ 	ñ 	ús   �A ¤A Á A	Nc                 ó    — | j                   5  | j                  �!| j                  j                  «        d | _        d d d «       y # 1 sw Y   y xY wr§   )rE   rG   r�   r`   s    r   r�   zFlushTimer.cancel<  s?   € Ø�Y‰Yñ 	"Ø�z‰zÐ%Ø—
‘
×!Ñ!Ô#Ø!�”
÷	"÷ 	"ñ 	"ús   �.AÁAc                 ób   — | j                   5  | j                  d ucd d d «       S # 1 sw Y   y xY wr§   )rE   rG   r`   s    r   rb   zFlushTimer.is_activeB  s*   € Ø�Y‰Yñ 	*Ø—:‘: TÐ)÷	*÷ 	*ò 	*ús   �%¥.c                 óÄ   — | j                   5  | j                  r| j                  |«       nd | _        d d d «       | j                  j                  «        y # 1 sw Y   Œ$xY wr§   )rE   r=   rÇ   rG   r<   r[   r‹   s     r   Ú_firezFlushTimer._fireF  sL   € Ø�Y‰Yñ 	"Ø�~Š~Ø×$Ñ$ WÕ-à!�”
÷		"ð 	�‰×ÑÕ÷	"ð 	"ús   �&AÁAc                 ó¨   — t        j                  || j                  |g«      | _        d| j                  _        | j                  j                  «        y )NT)rC   ÚTimerrË   rG   r¹   rŠ   r‹   s     r   rÇ   zFlushTimer._schedule_timerO  s8   € Ü—_‘_ W¨d¯j©j¸7¸)ÓDˆŒ
Ø ˆ�
‰
ÔØ�
‰
×ÑÕr!   r©   )r"   r#   r$   rª   r(   r¬   r    r«   rŠ   r�   rb   rË   rÇ   r&   r!   r   rF   rF   )  sd   „ ñð& ð &°Dó &ð˜Uð  tó ó"ð*˜4ó *ð˜Uð  tó ð uð °ô r!   rF   ),Úloggingr˜   rC   rK   rž   Útypingr   r   r   r   r   r   rµ   r	   Ú%snowplow_tracker.self_describing_jsonr
   Úsnowplow_tracker.typingr   r   r   r   r   r   r   Úsnowplow_tracker.contractsr   Úsnowplow_tracker.event_storer   r   ÚbasicConfigÚ	getLoggerr"   r:   ÚsetLevelÚINFOrA   r}   r>   r?   r   r(   r°   ÚobjectrF   r&   r!   r   ú<module>rÙ      sÂ   ðó$ Û Û Û Û ß ?× ?Ý å D÷÷ ñ õ .ß Gð €× Ñ Ô Ø	ˆ×	Ñ	˜8Ó	$€Ø ‡��—‘Ô àÐ àGð ð �WÐ€	Ø�&ˆ/€÷	"ñ 	"ôˆoô ôDc#�7ô c#ôL)�õ )r!   