1+ import atexit
12import functools
23import json
34import logging
5+ import queue
46import ssl
7+ import threading
8+ import time
59import uuid
610from base64 import b64encode
711from functools import cached_property , wraps
1822
1923LOG = logging .getLogger ("testgen" )
2024
25+ _BATCH_SIZE = 50 # Mixpanel /track array cap
26+ _FLUSH_INTERVAL_SEC = 10 # Time-based flush
27+ _QUEUE_MAX_SIZE = 1000 # Memory cap before drop-on-overflow
28+ _DRAIN_TIMEOUT_SEC = 5 # Bounded shutdown drain
29+
30+ _SHUTDOWN = object () # sentinel enqueued by drain() to stop the worker
31+
2132
2233def safe_method (method ):
2334 @wraps (method )
@@ -33,6 +44,14 @@ def wrapped(*args, **kwargs):
3344
3445class MixpanelService (Singleton ):
3546
47+ def __init__ (self ) -> None :
48+ self ._queue : queue .Queue | None = None
49+ self ._worker : threading .Thread | None = None
50+ self ._worker_lock = threading .Lock ()
51+ self ._started = False
52+ self ._stopped = False
53+ atexit .register (self .drain )
54+
3655 @cached_property
3756 @with_database_session
3857 def instance_id (self ):
@@ -54,33 +73,106 @@ def _hash_value(self, value: bytes | str, digest_size: int = 8) -> str:
5473
5574 @safe_method
5675 def send_event (self , event_name , include_usage = False , ** properties ):
57- self ._track ( event_name , include_usage = include_usage , ** properties )
76+ self ._enqueue ( self . _build_event ( event_name , include_usage = include_usage , ** properties ) )
5877
5978 def send_feedback (self , ** properties ):
60- # User-submitted feedback is content the user explicitly chose to share
61- # so it is not gated by the TG_ANALYTICS opt-out.
79+ # User-submitted feedback is content the user explicitly chose to share,
80+ # so it is not gated by the TG_ANALYTICS opt-out. It is a foreground action
81+ # posted synchronously — never enqueued — so it never starts the worker.
6282 try :
63- self ._track ( " feedback" , ** properties )
83+ self .send_mp_request ( "track?ip=1" , self . _build_event ( " feedback" , ** properties ) )
6484 except Exception :
6585 LOG .exception ("Error sending feedback" )
6686
67- def _track (self , event_name , include_usage = False , ** properties ):
87+ def _build_event (self , event_name , include_usage = False , ** properties ) -> dict :
6888 properties .setdefault ("instance_id" , self .instance_id )
6989 properties .setdefault ("edition" , settings .DOCKER_HUB_REPOSITORY )
7090 properties .setdefault ("version" , settings .VERSION )
7191 properties .setdefault ("username" , session .auth .user_display if session .auth else None )
7292 properties .setdefault ("distinct_id" , self .get_distinct_id (properties ["username" ]))
93+ properties .setdefault ("time" , int (time .time ()))
7394 if include_usage :
7495 properties .update (self .get_usage ())
7596
76- track_payload = {
97+ return {
7798 "event" : event_name ,
7899 "properties" : {
79100 "token" : settings .MIXPANEL_TOKEN ,
80101 ** properties ,
81- }
102+ },
82103 }
83- self .send_mp_request ("track?ip=1" , track_payload )
104+
105+ def _ensure_worker (self ) -> None :
106+ if self ._started :
107+ return
108+ with self ._worker_lock :
109+ if self ._started :
110+ return
111+ self ._queue = queue .Queue (maxsize = _QUEUE_MAX_SIZE )
112+ self ._worker = threading .Thread (target = self ._worker_loop , name = "mixpanel-flush" , daemon = True )
113+ self ._worker .start ()
114+ self ._started = True
115+
116+ def _enqueue (self , event : dict ) -> None :
117+ if self ._stopped :
118+ LOG .warning ("analytics worker stopped; dropping event" )
119+ return
120+ self ._ensure_worker ()
121+ try :
122+ self ._queue .put_nowait (event )
123+ except queue .Full :
124+ LOG .warning ("analytics queue full; dropping event" )
125+
126+ def _worker_loop (self ) -> None :
127+ while True :
128+ batch , stop = self ._next_batch ()
129+ if batch :
130+ self ._flush (batch )
131+ if stop :
132+ return
133+
134+ def _next_batch (self ) -> tuple [list [dict ], bool ]:
135+ """Block up to _FLUSH_INTERVAL_SEC for the first event, then drain up to
136+ _BATCH_SIZE without blocking. Returns (batch, stop)."""
137+ try :
138+ first = self ._queue .get (timeout = _FLUSH_INTERVAL_SEC )
139+ except queue .Empty :
140+ return [], False
141+ if first is _SHUTDOWN :
142+ return [], True
143+ batch = [first ]
144+ while len (batch ) < _BATCH_SIZE :
145+ try :
146+ event = self ._queue .get_nowait ()
147+ except queue .Empty :
148+ break
149+ if event is _SHUTDOWN :
150+ return batch , True
151+ batch .append (event )
152+ return batch , False
153+
154+ def _flush (self , events : list [dict ]) -> None :
155+ for start in range (0 , len (events ), _BATCH_SIZE ):
156+ chunk = events [start :start + _BATCH_SIZE ]
157+ try :
158+ self .send_mp_request ("track?ip=1" , chunk )
159+ except Exception :
160+ LOG .exception ("Failed to flush analytics batch" )
161+
162+ def drain (self ) -> None :
163+ """Flush queued events and stop the worker. Idempotent; bounded by
164+ _DRAIN_TIMEOUT_SEC. Called by the atexit hook and the server lifespan."""
165+ if not self ._started or self ._stopped :
166+ return
167+ self ._stopped = True
168+ try :
169+ self ._queue .put_nowait (_SHUTDOWN )
170+ except queue .Full :
171+ LOG .warning ("analytics queue full at shutdown; in-flight events may be dropped" )
172+ if self ._worker is not None :
173+ self ._worker .join (timeout = _DRAIN_TIMEOUT_SEC )
174+ if self ._worker .is_alive ():
175+ LOG .warning ("analytics drain timed out after %ss" , _DRAIN_TIMEOUT_SEC )
84176
85177 def get_ssl_context (self ):
86178 ssl_context = ssl .create_default_context ()
0 commit comments