Skip to content

Commit d7d43c8

Browse files
feat(traces): span batch transport (#952)
* feat(traces): span batch transport Adds the HTTP transport: one gzipped OTLP JSON POST per batch to {host}/i/v1/traces with bearer auth, classified as ok (2xx), too large (413, or a body over the 10 MiB hosted ingestion limit, refused without a request), retriable (408, 429, 5xx, transport errors, carrying any Retry-After as delta-seconds or HTTP-date) or fatal (other 4xx). Not reachable from the client. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01TkZAsCciW4PV8ZdcCHmAbA * fix(traces): classify the export response without reading its body A requests timeout bounds read inactivity, not the whole response, so a proxy answering 503 with a body that keeps dripping held the exporter's single flight open indefinitely. Stream the response, classify it from status and headers, and close it unread. * fix(traces): drain a small response body so the pooled connection survives Closing an unread streamed response tears the connection down, so every batch paid a new TCP and TLS handshake. A body whose Content-Length is at most 64 KiB is read before close, and a fatal status now logs it, so a bad key shows the server's error rather than a bare status. The outcome kind is a Literal, gzip runs at level 6, and the dead Retry-After guard and its mock-only test are gone. --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
1 parent d369f6e commit d7d43c8

3 files changed

Lines changed: 473 additions & 0 deletions

File tree

Lines changed: 310 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,310 @@
1+
import gzip
2+
import json
3+
import threading
4+
import time
5+
from datetime import datetime, timedelta, timezone
6+
from email.utils import format_datetime
7+
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
8+
from types import SimpleNamespace
9+
from unittest import mock
10+
11+
import pytest
12+
import requests
13+
14+
from posthog.tracing._transport import (
15+
OTLP_MAX_BODY_BYTES,
16+
SendOutcome,
17+
parse_retry_after,
18+
send_traces_batch,
19+
)
20+
from posthog.version import VERSION
21+
22+
PAYLOAD = {"resourceSpans": [{"scopeSpans": [{"spans": [{"name": "x"}]}]}]}
23+
24+
25+
def fake_client(**overrides):
26+
base = dict(
27+
disabled=False,
28+
send=True,
29+
host="https://us.example.com/",
30+
api_key="phc_test_key",
31+
timeout=7,
32+
)
33+
base.update(overrides)
34+
return SimpleNamespace(**base)
35+
36+
37+
def mock_session(status_code=200, headers=None):
38+
session = mock.Mock()
39+
session.post.return_value = mock.Mock(
40+
status_code=status_code, headers=headers or {}
41+
)
42+
return session
43+
44+
45+
def send(client=None, payload=PAYLOAD, session=None):
46+
session = session or mock_session()
47+
with mock.patch("posthog.tracing._transport._get_session", return_value=session):
48+
outcome = send_traces_batch(client or fake_client(), payload)
49+
return outcome, session
50+
51+
52+
class TestRequestShape:
53+
def test_posts_to_the_traces_endpoint_with_bearer_auth(self):
54+
outcome, session = send()
55+
assert outcome == SendOutcome("ok")
56+
args, kwargs = session.post.call_args
57+
assert args[0] == "https://us.example.com/i/v1/traces"
58+
assert kwargs["headers"]["Authorization"] == "Bearer phc_test_key"
59+
assert kwargs["headers"]["User-Agent"] == "posthog-python/" + VERSION
60+
assert kwargs["timeout"] == 7
61+
62+
def test_does_not_put_the_project_key_in_the_query_string(self):
63+
_, session = send()
64+
assert "token=" not in session.post.call_args[0][0]
65+
66+
def test_gzips_the_json_body_and_says_so(self):
67+
_, session = send()
68+
kwargs = session.post.call_args[1]
69+
assert kwargs["headers"]["Content-Encoding"] == "gzip"
70+
assert kwargs["headers"]["Content-Type"] == "application/json"
71+
assert json.loads(gzip.decompress(kwargs["data"])) == PAYLOAD
72+
73+
def test_falls_back_to_the_default_timeout(self):
74+
_, session = send(fake_client(timeout=None))
75+
assert session.post.call_args[1]["timeout"] == 15
76+
77+
78+
class TestGates:
79+
def test_disabled_client_is_fatal_without_a_request(self):
80+
outcome, session = send(fake_client(disabled=True))
81+
assert outcome.kind == "fatal"
82+
assert not session.post.called
83+
84+
def test_send_false_is_ok_without_a_request(self):
85+
outcome, session = send(fake_client(send=False))
86+
assert outcome.kind == "ok"
87+
assert not session.post.called
88+
89+
def test_an_oversized_body_is_too_large_without_a_request(self):
90+
payload = {"resourceSpans": [{"blob": "x" * (OTLP_MAX_BODY_BYTES + 1)}]}
91+
outcome, session = send(payload=payload)
92+
assert outcome.kind == "too-large"
93+
assert outcome.measured_locally
94+
assert not session.post.called
95+
96+
def test_a_413_is_not_marked_as_measured_locally(self):
97+
outcome, _ = send(session=mock_session(413))
98+
assert outcome.kind == "too-large"
99+
assert not outcome.measured_locally
100+
101+
def test_the_limit_is_what_hosted_ingestion_accepts(self):
102+
assert OTLP_MAX_BODY_BYTES == 10 * 1024 * 1024
103+
104+
def test_a_body_exactly_at_the_limit_is_sent(self):
105+
# {"s":""} is 8 bytes of JSON around the string.
106+
outcome, session = send(payload={"s": "x" * (OTLP_MAX_BODY_BYTES - 8)})
107+
assert outcome.kind == "ok"
108+
assert session.post.called
109+
110+
def test_a_body_one_byte_over_the_limit_is_not_sent(self):
111+
outcome, session = send(payload={"s": "x" * (OTLP_MAX_BODY_BYTES - 7)})
112+
assert outcome.kind == "too-large"
113+
assert not session.post.called
114+
115+
def test_measures_the_uncompressed_body(self):
116+
# Compresses to a few kilobytes.
117+
outcome, session = send(payload={"s": "\u2603" * OTLP_MAX_BODY_BYTES})
118+
assert outcome.kind == "too-large"
119+
assert not session.post.called
120+
121+
122+
class TestOutcomes:
123+
@pytest.mark.parametrize(
124+
"status,kind",
125+
[
126+
(200, "ok"),
127+
(204, "ok"),
128+
(413, "too-large"),
129+
(408, "retry-later"),
130+
(429, "retry-later"),
131+
(500, "retry-later"),
132+
(503, "retry-later"),
133+
(400, "fatal"),
134+
(401, "fatal"),
135+
(404, "fatal"),
136+
],
137+
)
138+
def test_maps_status_codes(self, status, kind):
139+
outcome, _ = send(session=mock_session(status))
140+
assert outcome.kind == kind
141+
142+
def test_a_transport_error_is_retriable(self):
143+
session = mock.Mock()
144+
session.post.side_effect = requests.exceptions.ConnectionError("down")
145+
outcome, _ = send(session=session)
146+
assert outcome == SendOutcome("retry-later")
147+
148+
def test_reads_retry_after_delta_seconds(self):
149+
outcome, _ = send(session=mock_session(429, {"Retry-After": "120"}))
150+
assert outcome == SendOutcome("retry-later", 120.0)
151+
152+
def test_reads_retry_after_http_date(self):
153+
when = format_datetime(datetime.now(timezone.utc) + timedelta(seconds=120))
154+
outcome, _ = send(session=mock_session(503, {"Retry-After": when}))
155+
assert outcome.kind == "retry-later"
156+
assert outcome.retry_after is not None and 100 < outcome.retry_after <= 120
157+
158+
159+
NOW = datetime(2026, 9, 10, 12, 0, 0, tzinfo=timezone.utc)
160+
161+
162+
class TestParseRetryAfter:
163+
@pytest.mark.parametrize(
164+
"value,expected",
165+
[
166+
("120", 120.0),
167+
(" 30 ", 30.0),
168+
("60, 120", 60.0),
169+
("Thu, 10 Sep 2026 12:00:30 GMT", 30.0),
170+
],
171+
)
172+
def test_reads_both_wire_forms(self, value, expected):
173+
assert parse_retry_after(value, NOW) == expected
174+
175+
@pytest.mark.parametrize(
176+
"value",
177+
[
178+
None,
179+
"",
180+
"0",
181+
"-5",
182+
"+5",
183+
"5.5",
184+
"1e3",
185+
"10 minutes",
186+
"Wed, 21 Oct 2015 07:28:00 GMT",
187+
"Thu, 10 Sep 2026 12:00:00 GMT",
188+
42,
189+
],
190+
)
191+
def test_treats_anything_else_as_absent(self, value):
192+
assert parse_retry_after(value, NOW) is None
193+
194+
def test_ignores_an_unparseable_retry_after(self):
195+
outcome, _ = send(session=mock_session(429, {"Retry-After": "10 minutes"}))
196+
assert outcome == SendOutcome("retry-later", None)
197+
198+
199+
class _SizedHandler(BaseHTTPRequestHandler):
200+
protocol_version = "HTTP/1.1"
201+
202+
def setup(self):
203+
super().setup()
204+
self.server.connections += 1
205+
206+
def do_POST(self):
207+
self.rfile.read(int(self.headers.get("Content-Length", 0)))
208+
body = self.server.body
209+
self.send_response(self.server.status)
210+
self.send_header("Content-Type", "application/json")
211+
self.send_header("Content-Length", str(len(body)))
212+
self.end_headers()
213+
self.wfile.write(body)
214+
215+
def log_message(self, *args):
216+
pass
217+
218+
219+
class _ChunkedHandler(BaseHTTPRequestHandler):
220+
def do_POST(self):
221+
self.rfile.read(int(self.headers.get("Content-Length", 0)))
222+
self.send_response(self.server.status)
223+
self.send_header("Transfer-Encoding", "chunked")
224+
self.end_headers()
225+
if self.server.status < 300:
226+
self.wfile.write(b"0\r\n\r\n")
227+
return
228+
# An error body that drips a chunk every 10 ms and never finishes.
229+
while not self.server.stop.is_set():
230+
try:
231+
self.wfile.write(b"1\r\nx\r\n")
232+
self.wfile.flush()
233+
except OSError:
234+
return
235+
time.sleep(0.01)
236+
237+
def log_message(self, *args):
238+
pass
239+
240+
241+
@pytest.fixture
242+
def local_server():
243+
servers = []
244+
245+
def start(status, body=None):
246+
handler = _ChunkedHandler if body is None else _SizedHandler
247+
server = ThreadingHTTPServer(("127.0.0.1", 0), handler)
248+
server.status = status
249+
server.body = body
250+
server.connections = 0
251+
server.stop = threading.Event()
252+
thread = threading.Thread(target=server.serve_forever, daemon=True)
253+
thread.start()
254+
servers.append((server, thread))
255+
return "http://127.0.0.1:{}".format(server.server_port)
256+
257+
yield servers, start
258+
for server, thread in servers:
259+
server.stop.set()
260+
server.shutdown()
261+
server.server_close()
262+
thread.join(2)
263+
264+
265+
class TestResponseBody:
266+
def test_closes_the_response_without_reading_an_unsized_body(self):
267+
_, session = send()
268+
assert session.post.call_args[1]["stream"] is True
269+
assert session.post.return_value.close.called
270+
271+
def test_does_not_wait_for_a_dripping_error_body(self, local_server):
272+
_, start = local_server
273+
client = fake_client(host=start(503), timeout=0.5)
274+
started = time.monotonic()
275+
outcome = send_traces_batch(client, PAYLOAD)
276+
assert outcome == SendOutcome("retry-later", None)
277+
assert time.monotonic() - started < 2
278+
279+
def test_a_completed_response_is_still_ok(self, local_server):
280+
_, start = local_server
281+
client = fake_client(host=start(200), timeout=0.5)
282+
assert send_traces_batch(client, PAYLOAD) == SendOutcome("ok")
283+
284+
def test_drains_a_small_sized_body_so_the_connection_is_reused(self, local_server):
285+
servers, start = local_server
286+
client = fake_client(host=start(200, body=b"{}"), timeout=2)
287+
for _ in range(5):
288+
assert send_traces_batch(client, PAYLOAD) == SendOutcome("ok")
289+
assert servers[0][0].connections == 1
290+
291+
def test_leaves_a_large_sized_body_unread(self):
292+
response = mock.Mock(
293+
status_code=200, headers={"Content-Length": str(64 * 1024 + 1)}
294+
)
295+
type(response).text = mock.PropertyMock(side_effect=AssertionError("read"))
296+
session = mock.Mock()
297+
session.post.return_value = response
298+
assert send(session=session)[0] == SendOutcome("ok")
299+
assert response.close.called
300+
301+
def test_logs_the_server_error_body_on_a_fatal_status(self, caplog):
302+
body = '{"error": "invalid api key"}'
303+
response = mock.Mock(
304+
status_code=401, headers={"Content-Length": str(len(body))}, text=body
305+
)
306+
session = mock.Mock()
307+
session.post.return_value = response
308+
with caplog.at_level("ERROR", logger="posthog"):
309+
assert send(session=session)[0] == SendOutcome("fatal")
310+
assert 'HTTP 401: {"error": "invalid api key"}' in caplog.text

0 commit comments

Comments
 (0)