Skip to content

Commit 267f783

Browse files
committed
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.
1 parent ecf3d66 commit 267f783

2 files changed

Lines changed: 100 additions & 32 deletions

File tree

posthog/test/tracing/test_transport.py

Lines changed: 63 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,8 @@
22
import json
33
import threading
44
import time
5-
from datetime import datetime, timezone
5+
from datetime import datetime, timedelta, timezone
6+
from email.utils import format_datetime
67
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
78
from types import SimpleNamespace
89
from unittest import mock
@@ -149,11 +150,10 @@ def test_reads_retry_after_delta_seconds(self):
149150
assert outcome == SendOutcome("retry-later", 120.0)
150151

151152
def test_reads_retry_after_http_date(self):
152-
outcome, _ = send(
153-
session=mock_session(503, {"Retry-After": "Wed, 21 Oct 2099 07:28:00 GMT"})
154-
)
153+
when = format_datetime(datetime.now(timezone.utc) + timedelta(seconds=120))
154+
outcome, _ = send(session=mock_session(503, {"Retry-After": when}))
155155
assert outcome.kind == "retry-later"
156-
assert outcome.retry_after is not None and outcome.retry_after > 0
156+
assert outcome.retry_after is not None and 100 < outcome.retry_after <= 120
157157

158158

159159
NOW = datetime(2026, 9, 10, 12, 0, 0, tzinfo=timezone.utc)
@@ -195,13 +195,25 @@ def test_ignores_an_unparseable_retry_after(self):
195195
outcome, _ = send(session=mock_session(429, {"Retry-After": "10 minutes"}))
196196
assert outcome == SendOutcome("retry-later", None)
197197

198-
def test_survives_a_throwing_headers_object(self):
199-
response = mock.Mock(status_code=503)
200-
response.headers.get.side_effect = RuntimeError("no headers")
201-
session = mock.Mock()
202-
session.post.return_value = response
203-
outcome, _ = send(session=session)
204-
assert outcome == SendOutcome("retry-later", None)
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
205217

206218

207219
class _ChunkedHandler(BaseHTTPRequestHandler):
@@ -230,16 +242,19 @@ def log_message(self, *args):
230242
def local_server():
231243
servers = []
232244

233-
def start(status):
234-
server = ThreadingHTTPServer(("127.0.0.1", 0), _ChunkedHandler)
245+
def start(status, body=None):
246+
handler = _ChunkedHandler if body is None else _SizedHandler
247+
server = ThreadingHTTPServer(("127.0.0.1", 0), handler)
235248
server.status = status
249+
server.body = body
250+
server.connections = 0
236251
server.stop = threading.Event()
237252
thread = threading.Thread(target=server.serve_forever, daemon=True)
238253
thread.start()
239254
servers.append((server, thread))
240255
return "http://127.0.0.1:{}".format(server.server_port)
241256

242-
yield start
257+
yield servers, start
243258
for server, thread in servers:
244259
server.stop.set()
245260
server.shutdown()
@@ -248,18 +263,48 @@ def start(status):
248263

249264

250265
class TestResponseBody:
251-
def test_closes_the_response_without_reading_the_body(self):
266+
def test_closes_the_response_without_reading_an_unsized_body(self):
252267
_, session = send()
253268
assert session.post.call_args[1]["stream"] is True
254269
assert session.post.return_value.close.called
255270

256271
def test_does_not_wait_for_a_dripping_error_body(self, local_server):
257-
client = fake_client(host=local_server(503), timeout=0.5)
272+
_, start = local_server
273+
client = fake_client(host=start(503), timeout=0.5)
258274
started = time.monotonic()
259275
outcome = send_traces_batch(client, PAYLOAD)
260276
assert outcome == SendOutcome("retry-later", None)
261277
assert time.monotonic() - started < 2
262278

263279
def test_a_completed_response_is_still_ok(self, local_server):
264-
client = fake_client(host=local_server(200), timeout=0.5)
280+
_, start = local_server
281+
client = fake_client(host=start(200), timeout=0.5)
265282
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

posthog/tracing/_transport.py

Lines changed: 37 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@
88
from dataclasses import dataclass
99
from datetime import datetime, timezone
1010
from email.utils import parsedate_to_datetime
11-
from typing import Any, Optional
11+
from typing import Any, Literal, Optional
1212

1313
import requests
1414

@@ -26,6 +26,15 @@
2626

2727
_RETRIABLE_STATUSES = frozenset({408, 429})
2828

29+
# Level 9 costs about twice the CPU of 6 for about 1% fewer bytes.
30+
_GZIP_LEVEL = 6
31+
32+
# A body this small is read so the pooled connection survives the response;
33+
# an unread stream is torn down on close, and every batch would then pay a
34+
# new TCP and TLS handshake.
35+
_MAX_DRAINED_BODY_BYTES = 64 * 1024
36+
_LOGGED_BODY_CHARS = 512
37+
2938
_DELTA_SECONDS_RE = re.compile(r"^\d+$")
3039
_NUMERIC_RE = re.compile(r"^[+-]?[\d.]+$")
3140

@@ -34,7 +43,7 @@
3443
class SendOutcome:
3544
"""How one export attempt went: ``ok``, ``retry-later``, ``too-large`` or ``fatal``."""
3645

37-
kind: str
46+
kind: Literal["ok", "retry-later", "too-large", "fatal"]
3847
retry_after: Optional[float] = None
3948
# Too large by the SDK's own measure, so no request was spent (too-large only).
4049
measured_locally: bool = False
@@ -100,7 +109,7 @@ def send_traces_batch(client: Any, payload: dict) -> SendOutcome:
100109
try:
101110
response = _get_session().post(
102111
url,
103-
data=gzip.compress(serialized),
112+
data=gzip.compress(serialized, compresslevel=_GZIP_LEVEL),
104113
headers={
105114
"Content-Type": "application/json",
106115
"Content-Encoding": "gzip",
@@ -113,26 +122,40 @@ def send_traces_batch(client: Any, payload: dict) -> SendOutcome:
113122
except requests.exceptions.RequestException as e:
114123
log.debug("Span batch request failed: %s", e)
115124
return SendOutcome("retry-later")
116-
# Status and headers alone classify the response, so the body is never
117-
# read: the timeout bounds read inactivity, and a body that keeps dripping
118-
# would otherwise hold the exporter's single flight open indefinitely.
125+
# Status and headers classify the response. A body of unknown or large
126+
# size is left unread: the timeout bounds read inactivity, and a body that
127+
# keeps dripping would otherwise hold the exporter's single flight open.
119128
try:
120-
return _classify(response)
129+
return _classify(response, _read_small_body(response))
121130
finally:
122131
response.close()
123132

124133

125-
def _classify(response: requests.Response) -> SendOutcome:
134+
def _read_small_body(response: requests.Response) -> Optional[str]:
135+
try:
136+
length = int(response.headers.get("Content-Length", ""))
137+
except (TypeError, ValueError):
138+
return None
139+
if not 0 <= length <= _MAX_DRAINED_BODY_BYTES:
140+
return None
141+
try:
142+
return response.text
143+
except requests.exceptions.RequestException:
144+
return None
145+
146+
147+
def _classify(response: requests.Response, body: Optional[str]) -> SendOutcome:
126148
status = response.status_code
127149
if status < 300:
128150
return OK
129151
if status == 413:
130152
return TOO_LARGE
131153
if status >= 500 or status in _RETRIABLE_STATUSES:
132-
try:
133-
retry_after = parse_retry_after(response.headers.get("Retry-After"))
134-
except Exception:
135-
retry_after = None
136-
return SendOutcome("retry-later", retry_after)
137-
log.error("Failed to send span batch: HTTP %s", status)
154+
return SendOutcome(
155+
"retry-later", parse_retry_after(response.headers.get("Retry-After"))
156+
)
157+
detail = body.strip()[:_LOGGED_BODY_CHARS] if body else ""
158+
log.error(
159+
"Failed to send span batch: HTTP %s%s", status, ": " + detail if detail else ""
160+
)
138161
return FATAL

0 commit comments

Comments
 (0)