Skip to content

Commit f4fe475

Browse files
feat(traces): span export queue with retries
Adds SpanExporter, the in-memory queue ended spans wait in until they are batched and sent, separate from the events queue. A timer flushes every flush_interval and a full batch flushes at once; only one flush runs at a time. A full queue drops the incoming span, never a queued parent. Failures back off exponentially with jitter, floored by Retry-After (clamped to 30 s, extended by a later deadline but never shortened), and automatic sends pause while backing off. A batch refused across 8 backoff windows is dropped, a 413 halves the batch and ramps back, and other 4xx drop it. flush(timeout) always sends the first batch, so a serverless handler with no budget left still ships spans. Not reachable from the client. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TkZAsCciW4PV8ZdcCHmAbA
1 parent 8abb393 commit f4fe475

3 files changed

Lines changed: 1390 additions & 2 deletions

File tree

‎posthog/test/tracing/helpers.py‎

Lines changed: 51 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
"""Shared fakes for the tracing pipeline tests."""
1+
"""Shared fakes for the tracing pipeline and export tests."""
22

33
import threading
44
import time
@@ -8,13 +8,18 @@
88

99
import pytest
1010

11+
from posthog.tracing import _export as export_module
1112
from posthog.tracing._config import resolve_traces_config
1213
from posthog.tracing._drops import DropLog
14+
from posthog.tracing._export import SpanExporter
1315
from posthog.tracing._pipeline import PostHogTraces
16+
from posthog.tracing._transport import SendOutcome
1417

1518
TRACE_ID = "4bf92f3577b34da6a3ce929d0e0e4736"
1619
SPAN_ID = "00f067aa0ba902b7"
1720

21+
RealTimer = threading.Timer
22+
1823

1924
class FakeTimer:
2025
"""Records the delay it was armed with; fires only when a test says so."""
@@ -39,6 +44,21 @@ def fire(self):
3944
self.fn()
4045

4146

47+
class FakeSender:
48+
def __init__(self, *outcomes):
49+
self.outcomes = list(outcomes)
50+
self.payloads: list = []
51+
52+
def __call__(self, client, payload):
53+
self.payloads.append(payload)
54+
if len(self.outcomes) > 1:
55+
return self.outcomes.pop(0)
56+
return self.outcomes[0] if self.outcomes else SendOutcome("ok")
57+
58+
def batches(self):
59+
return [p["resourceSpans"][0]["scopeSpans"][0]["spans"] for p in self.payloads]
60+
61+
4262
class RecordingExporter:
4363
"""Stands in for the export queue: keeps every record it is handed."""
4464

@@ -70,6 +90,13 @@ def fake_timers():
7090
yield FakeTimer
7191

7292

93+
@pytest.fixture(autouse=True)
94+
def no_jitter():
95+
# Backoff delays are asserted exactly; TestJitter covers the spread.
96+
with mock.patch.object(export_module, "_draw_jitter", return_value=1.0):
97+
yield
98+
99+
73100
@pytest.fixture
74101
def clock():
75102
state = {"now": 1000.0}
@@ -90,5 +117,27 @@ def make(client=None, context=None, **config):
90117
return pipeline, exporter, active
91118

92119

120+
def make_traces(sender=None, client=None, context=None, **config):
121+
"""A pipeline over a real exporter whose sender is ``sender``."""
122+
config.setdefault("flush_interval", 5)
123+
client = client or SimpleNamespace(disabled=False, send=True)
124+
sender = sender or FakeSender(SendOutcome("ok"))
125+
active: ContextVar = ContextVar("active", default=None)
126+
resolved = resolve_traces_config(config)
127+
drops = DropLog(resolved.flush_interval)
128+
pipeline = PostHogTraces(
129+
client,
130+
resolved,
131+
lambda: context or {},
132+
active,
133+
SpanExporter(client, resolved, drops, send=sender),
134+
drops,
135+
)
136+
return pipeline, sender, active
137+
138+
93139
def queued(pipeline):
94-
return pipeline._exporter.records
140+
exporter = pipeline._exporter
141+
if isinstance(exporter, RecordingExporter):
142+
return exporter.records
143+
return exporter._queue

0 commit comments

Comments
 (0)