Skip to content

Commit d6920ae

Browse files
author
ci bot
committed
Merge branch 'batched-telemetry' into 'enterprise'
MCP & UI: async batched Mixpanel telemetry (TG-1076) See merge request dkinternal/testgen/dataops-testgen!571
2 parents e130f3c + 73f00a4 commit d6920ae

7 files changed

Lines changed: 539 additions & 114 deletions

File tree

testgen/common/mixpanel_service.py

Lines changed: 100 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,11 @@
1+
import atexit
12
import functools
23
import json
34
import logging
5+
import queue
46
import ssl
7+
import threading
8+
import time
59
import uuid
610
from base64 import b64encode
711
from functools import cached_property, wraps
@@ -18,6 +22,13 @@
1822

1923
LOG = 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

2233
def safe_method(method):
2334
@wraps(method)
@@ -33,6 +44,14 @@ def wrapped(*args, **kwargs):
3344

3445
class 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()

testgen/mcp/permissions.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -80,6 +80,11 @@ def set_mcp_token(token: str | None) -> None:
8080
_mcp_token.set(token)
8181

8282

83+
def get_mcp_username() -> str | None:
84+
"""Return the authenticated username for the current MCP request (or None)."""
85+
return _mcp_username.get()
86+
87+
8388
def get_authorized_mcp_user() -> User:
8489
"""Get the authenticated and authorized User for the current MCP request.
8590

0 commit comments

Comments
 (0)