diff --git a/posthog/test/tracing/test_transport.py b/posthog/test/tracing/test_transport.py new file mode 100644 index 000000000..46d0916c4 --- /dev/null +++ b/posthog/test/tracing/test_transport.py @@ -0,0 +1,265 @@ +import gzip +import json +import threading +import time +from datetime import datetime, timezone +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from types import SimpleNamespace +from unittest import mock + +import pytest +import requests + +from posthog.tracing._transport import ( + OTLP_MAX_BODY_BYTES, + SendOutcome, + parse_retry_after, + send_traces_batch, +) +from posthog.version import VERSION + +PAYLOAD = {"resourceSpans": [{"scopeSpans": [{"spans": [{"name": "x"}]}]}]} + + +def fake_client(**overrides): + base = dict( + disabled=False, + send=True, + host="https://us.example.com/", + api_key="phc_test_key", + timeout=7, + ) + base.update(overrides) + return SimpleNamespace(**base) + + +def mock_session(status_code=200, headers=None): + session = mock.Mock() + session.post.return_value = mock.Mock( + status_code=status_code, headers=headers or {} + ) + return session + + +def send(client=None, payload=PAYLOAD, session=None): + session = session or mock_session() + with mock.patch("posthog.tracing._transport._get_session", return_value=session): + outcome = send_traces_batch(client or fake_client(), payload) + return outcome, session + + +class TestRequestShape: + def test_posts_to_the_traces_endpoint_with_bearer_auth(self): + outcome, session = send() + assert outcome == SendOutcome("ok") + args, kwargs = session.post.call_args + assert args[0] == "https://us.example.com/i/v1/traces" + assert kwargs["headers"]["Authorization"] == "Bearer phc_test_key" + assert kwargs["headers"]["User-Agent"] == "posthog-python/" + VERSION + assert kwargs["timeout"] == 7 + + def test_does_not_put_the_project_key_in_the_query_string(self): + _, session = send() + assert "token=" not in session.post.call_args[0][0] + + def test_gzips_the_json_body_and_says_so(self): + _, session = send() + kwargs = session.post.call_args[1] + assert kwargs["headers"]["Content-Encoding"] == "gzip" + assert kwargs["headers"]["Content-Type"] == "application/json" + assert json.loads(gzip.decompress(kwargs["data"])) == PAYLOAD + + def test_falls_back_to_the_default_timeout(self): + _, session = send(fake_client(timeout=None)) + assert session.post.call_args[1]["timeout"] == 15 + + +class TestGates: + def test_disabled_client_is_fatal_without_a_request(self): + outcome, session = send(fake_client(disabled=True)) + assert outcome.kind == "fatal" + assert not session.post.called + + def test_send_false_is_ok_without_a_request(self): + outcome, session = send(fake_client(send=False)) + assert outcome.kind == "ok" + assert not session.post.called + + def test_an_oversized_body_is_too_large_without_a_request(self): + payload = {"resourceSpans": [{"blob": "x" * (OTLP_MAX_BODY_BYTES + 1)}]} + outcome, session = send(payload=payload) + assert outcome.kind == "too-large" + assert outcome.measured_locally + assert not session.post.called + + def test_a_413_is_not_marked_as_measured_locally(self): + outcome, _ = send(session=mock_session(413)) + assert outcome.kind == "too-large" + assert not outcome.measured_locally + + def test_the_limit_is_what_hosted_ingestion_accepts(self): + assert OTLP_MAX_BODY_BYTES == 10 * 1024 * 1024 + + def test_a_body_exactly_at_the_limit_is_sent(self): + # {"s":""} is 8 bytes of JSON around the string. + outcome, session = send(payload={"s": "x" * (OTLP_MAX_BODY_BYTES - 8)}) + assert outcome.kind == "ok" + assert session.post.called + + def test_a_body_one_byte_over_the_limit_is_not_sent(self): + outcome, session = send(payload={"s": "x" * (OTLP_MAX_BODY_BYTES - 7)}) + assert outcome.kind == "too-large" + assert not session.post.called + + def test_measures_the_uncompressed_body(self): + # Compresses to a few kilobytes. + outcome, session = send(payload={"s": "\u2603" * OTLP_MAX_BODY_BYTES}) + assert outcome.kind == "too-large" + assert not session.post.called + + +class TestOutcomes: + @pytest.mark.parametrize( + "status,kind", + [ + (200, "ok"), + (204, "ok"), + (413, "too-large"), + (408, "retry-later"), + (429, "retry-later"), + (500, "retry-later"), + (503, "retry-later"), + (400, "fatal"), + (401, "fatal"), + (404, "fatal"), + ], + ) + def test_maps_status_codes(self, status, kind): + outcome, _ = send(session=mock_session(status)) + assert outcome.kind == kind + + def test_a_transport_error_is_retriable(self): + session = mock.Mock() + session.post.side_effect = requests.exceptions.ConnectionError("down") + outcome, _ = send(session=session) + assert outcome == SendOutcome("retry-later") + + def test_reads_retry_after_delta_seconds(self): + outcome, _ = send(session=mock_session(429, {"Retry-After": "120"})) + assert outcome == SendOutcome("retry-later", 120.0) + + def test_reads_retry_after_http_date(self): + outcome, _ = send( + session=mock_session(503, {"Retry-After": "Wed, 21 Oct 2099 07:28:00 GMT"}) + ) + assert outcome.kind == "retry-later" + assert outcome.retry_after is not None and outcome.retry_after > 0 + + +NOW = datetime(2026, 9, 10, 12, 0, 0, tzinfo=timezone.utc) + + +class TestParseRetryAfter: + @pytest.mark.parametrize( + "value,expected", + [ + ("120", 120.0), + (" 30 ", 30.0), + ("60, 120", 60.0), + ("Thu, 10 Sep 2026 12:00:30 GMT", 30.0), + ], + ) + def test_reads_both_wire_forms(self, value, expected): + assert parse_retry_after(value, NOW) == expected + + @pytest.mark.parametrize( + "value", + [ + None, + "", + "0", + "-5", + "+5", + "5.5", + "1e3", + "10 minutes", + "Wed, 21 Oct 2015 07:28:00 GMT", + "Thu, 10 Sep 2026 12:00:00 GMT", + 42, + ], + ) + def test_treats_anything_else_as_absent(self, value): + assert parse_retry_after(value, NOW) is None + + def test_ignores_an_unparseable_retry_after(self): + outcome, _ = send(session=mock_session(429, {"Retry-After": "10 minutes"})) + assert outcome == SendOutcome("retry-later", None) + + def test_survives_a_throwing_headers_object(self): + response = mock.Mock(status_code=503) + response.headers.get.side_effect = RuntimeError("no headers") + session = mock.Mock() + session.post.return_value = response + outcome, _ = send(session=session) + assert outcome == SendOutcome("retry-later", None) + + +class _ChunkedHandler(BaseHTTPRequestHandler): + def do_POST(self): + self.rfile.read(int(self.headers.get("Content-Length", 0))) + self.send_response(self.server.status) + self.send_header("Transfer-Encoding", "chunked") + self.end_headers() + if self.server.status < 300: + self.wfile.write(b"0\r\n\r\n") + return + # An error body that drips a chunk every 10 ms and never finishes. + while not self.server.stop.is_set(): + try: + self.wfile.write(b"1\r\nx\r\n") + self.wfile.flush() + except OSError: + return + time.sleep(0.01) + + def log_message(self, *args): + pass + + +@pytest.fixture +def local_server(): + servers = [] + + def start(status): + server = ThreadingHTTPServer(("127.0.0.1", 0), _ChunkedHandler) + server.status = status + server.stop = threading.Event() + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + servers.append((server, thread)) + return "http://127.0.0.1:{}".format(server.server_port) + + yield start + for server, thread in servers: + server.stop.set() + server.shutdown() + server.server_close() + thread.join(2) + + +class TestResponseBody: + def test_closes_the_response_without_reading_the_body(self): + _, session = send() + assert session.post.call_args[1]["stream"] is True + assert session.post.return_value.close.called + + def test_does_not_wait_for_a_dripping_error_body(self, local_server): + client = fake_client(host=local_server(503), timeout=0.5) + started = time.monotonic() + outcome = send_traces_batch(client, PAYLOAD) + assert outcome == SendOutcome("retry-later", None) + assert time.monotonic() - started < 2 + + def test_a_completed_response_is_still_ok(self, local_server): + client = fake_client(host=local_server(200), timeout=0.5) + assert send_traces_batch(client, PAYLOAD) == SendOutcome("ok") diff --git a/posthog/tracing/_transport.py b/posthog/tracing/_transport.py new file mode 100644 index 000000000..2bf4c9e28 --- /dev/null +++ b/posthog/tracing/_transport.py @@ -0,0 +1,138 @@ +"""HTTP transport for span batches: one POST to ``/i/v1/traces`` per batch.""" + +import gzip +import json +import logging +import math +import re +from dataclasses import dataclass +from datetime import datetime, timezone +from email.utils import parsedate_to_datetime +from typing import Any, Optional + +import requests + +from ..request import USER_AGENT, _get_session +from ..utils import remove_trailing_slash + +log = logging.getLogger("posthog") + +TRACES_PATH = "/i/v1/traces" + +# PostHog's hosted ingestion accepts 10 MiB (measured decompressed); a larger +# body can only come back 413, so it is refused without a request. A proxy or +# self-hosted deployment enforcing less is covered by the 413 path. +OTLP_MAX_BODY_BYTES = 10 * 1024 * 1024 + +_RETRIABLE_STATUSES = frozenset({408, 429}) + +_DELTA_SECONDS_RE = re.compile(r"^\d+$") +_NUMERIC_RE = re.compile(r"^[+-]?[\d.]+$") + + +@dataclass(frozen=True) +class SendOutcome: + """How one export attempt went: ``ok``, ``retry-later``, ``too-large`` or ``fatal``.""" + + kind: str + retry_after: Optional[float] = None + # Too large by the SDK's own measure, so no request was spent (too-large only). + measured_locally: bool = False + + +OK = SendOutcome("ok") +TOO_LARGE = SendOutcome("too-large") +TOO_LARGE_LOCALLY = SendOutcome("too-large", measured_locally=True) +FATAL = SendOutcome("fatal") + + +def parse_retry_after(value: Any, now: Optional[datetime] = None) -> Optional[float]: + """``Retry-After`` as seconds from now; ``None`` when absent, malformed or not in the future. + + Accepts delta-seconds or an HTTP-date. A repeated header arrives joined as + ``"60, 120"``; the first value is the outermost hop's. + """ + if not isinstance(value, str) or not value.strip(): + return None + raw = value.strip() + if re.match(r"^\d+\s*,", raw): + raw = raw.split(",", 1)[0].strip() + if _DELTA_SECONDS_RE.match(raw): + seconds = float(raw) + elif _NUMERIC_RE.match(raw): + return None + else: + try: + when = parsedate_to_datetime(raw) + except (TypeError, ValueError, IndexError): + return None + if when.tzinfo is None: + when = when.replace(tzinfo=timezone.utc) + seconds = (when - (now or datetime.now(timezone.utc))).total_seconds() + if not math.isfinite(seconds) or seconds <= 0: + return None + return seconds + + +def send_traces_batch(client: Any, payload: dict) -> SendOutcome: + """POST one OTLP batch with bearer auth and gzip, classifying the response. + + 2xx is ok; 413 is too large; 408, 429, 5xx and transport errors are + retriable; any other status is fatal. A 2xx is not proof of ingestion: an + unknown but well-formed key is accepted and the spans dropped downstream. + """ + if getattr(client, "disabled", False): + return FATAL + if not getattr(client, "send", True): + return OK + + serialized = json.dumps(payload, separators=(",", ":")).encode("utf-8") + if len(serialized) > OTLP_MAX_BODY_BYTES: + log.warning( + "Span batch is %s bytes, over the %s byte ingestion limit; not sending it", + len(serialized), + OTLP_MAX_BODY_BYTES, + ) + return TOO_LARGE_LOCALLY + + url = remove_trailing_slash(client.host) + TRACES_PATH + timeout = getattr(client, "timeout", 15) or 15 + try: + response = _get_session().post( + url, + data=gzip.compress(serialized), + headers={ + "Content-Type": "application/json", + "Content-Encoding": "gzip", + "Authorization": "Bearer {}".format(client.api_key), + "User-Agent": USER_AGENT, + }, + timeout=timeout, + stream=True, + ) + except requests.exceptions.RequestException as e: + log.debug("Span batch request failed: %s", e) + return SendOutcome("retry-later") + # Status and headers alone classify the response, so the body is never + # read: the timeout bounds read inactivity, and a body that keeps dripping + # would otherwise hold the exporter's single flight open indefinitely. + try: + return _classify(response) + finally: + response.close() + + +def _classify(response: requests.Response) -> SendOutcome: + status = response.status_code + if status < 300: + return OK + if status == 413: + return TOO_LARGE + if status >= 500 or status in _RETRIABLE_STATUSES: + try: + retry_after = parse_retry_after(response.headers.get("Retry-After")) + except Exception: + retry_after = None + return SendOutcome("retry-later", retry_after) + log.error("Failed to send span batch: HTTP %s", status) + return FATAL diff --git a/typings/requests/__init__.pyi b/typings/requests/__init__.pyi index 91034b4f4..75a3fa48c 100644 --- a/typings/requests/__init__.pyi +++ b/typings/requests/__init__.pyi @@ -8,6 +8,7 @@ class Response: text: str headers: dict[str, str] def json(self) -> Any: ... + def close(self) -> None: ... class Session: def mount(self, prefix: str, adapter: adapters.HTTPAdapter) -> None: ... @@ -19,6 +20,7 @@ class Session: data: str | bytes, headers: dict[str, str], timeout: int, + stream: bool = ..., ) -> Response: ... def get( self,