From f97bf27e9a54ef071fb5ec036904a7afd4aa3b52 Mon Sep 17 00:00:00 2001 From: Eli Reisman Date: Tue, 6 Oct 2026 15:09:55 -0700 Subject: [PATCH 1/3] feat: AI lane config and size-aware batching --- AGENTS.md | 2 +- posthog/capture_compression.py | 22 +++- posthog/client.py | 136 ++++++++++++++++++++--- posthog/consumer.py | 39 +++++-- posthog/test/test_ai_capture_lane.py | 158 ++++++++++++++++++++++++--- posthog/test/test_consumer.py | 21 ++++ references/public_api_snapshot.txt | 7 +- 7 files changed, 341 insertions(+), 44 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 6b7a8d74..4d67e9ef 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -23,7 +23,7 @@ Follow [Public API changes](./CONTRIBUTING.md#public-api-changes). As an agent, Before changing capture configuration, serialization, routing, or retries, read the relevant implementation and tests. -Capture v1 is the only capture protocol (`capture` posts to `/i/v1/analytics/events`, `capture_ai` to `/i/v1/ai/events`); strictly typed v1 options and `$set`/`$set_once` relocation; compression (gzip, zlib-wrapped deflate, optional zstd, default none); partial-only per-event retries with stable identity; accumulated drop reporting even on 2xx; terminal v1 `429`; `Retry-After` as a minimum bounded by the shared 30s ceiling; and inline blocking retries with `sync_mode=True`. +Capture v1 is the only capture protocol (`capture` posts to `/i/v1/analytics/events`, `capture_ai` to `/i/v1/ai/events`); strictly typed v1 options and `$set`/`$set_once` relocation; compression (gzip, zlib-wrapped deflate, optional zstd, default none), set per lane by `capture_compression` and `capture_ai_compression`; partial-only per-event retries with stable identity; accumulated drop reporting even on 2xx; terminal v1 `429`; `Retry-After` as a minimum bounded by the shared 30s ceiling; and inline blocking retries with `sync_mode=True`. ## Mirror and build safety diff --git a/posthog/capture_compression.py b/posthog/capture_compression.py index 52a1d4c2..257c80c8 100644 --- a/posthog/capture_compression.py +++ b/posthog/capture_compression.py @@ -54,6 +54,7 @@ def _zstd_available() -> bool: def _coerce_explicit( value: Union[CaptureCompression, str], + name: str = "capture_compression", ) -> CaptureCompression: """Normalize an explicitly-supplied compression to a ``CaptureCompression``. @@ -68,7 +69,7 @@ def _coerce_explicit( if resolved is not None: return resolved raise ValueError( - f"invalid capture_compression {value!r}; expected a CaptureCompression " + f"invalid {name} {value!r}; expected a CaptureCompression " f"or one of {sorted(_ALIASES)}" ) @@ -122,3 +123,22 @@ def _resolve_capture_compression( ) return fallback return env_resolved + + +def _resolve_capture_ai_compression( + capture_ai_compression: Optional[Union[CaptureCompression, str]] = None, +) -> CaptureCompression: + """Resolve the AI lane's request-body compression. + + Explicit argument only, defaulting to ``NONE``. ``POSTHOG_CAPTURE_COMPRESSION`` + does not apply, so changing analytics compression never changes AI uploads. + """ + if capture_ai_compression is None: + return CaptureCompression.NONE + resolved = _coerce_explicit(capture_ai_compression, "capture_ai_compression") + if resolved is CaptureCompression.ZSTD and not _zstd_available(): + raise ValueError( + "capture_ai_compression 'zstd' requires the zstandard package; " + "install posthog[zstd]" + ) + return resolved diff --git a/posthog/client.py b/posthog/client.py index a5e23fe2..856692a7 100644 --- a/posthog/client.py +++ b/posthog/client.py @@ -2,6 +2,7 @@ import hashlib as _hashlib import inspect import json +from contextlib import contextmanager import logging import os import sys @@ -29,6 +30,7 @@ from posthog.tracing._span import inert_span as _inert_span from posthog.capture_compression import ( CaptureCompression, + _resolve_capture_ai_compression, _resolve_capture_compression, ) from posthog.capture_send import ( @@ -91,6 +93,7 @@ from posthog.request import ( USER_AGENT as _USER_AGENT, APIError, + DatetimeSerializer as _DatetimeSerializer, QuotaLimitError, RequestsConnectionError, RequestsTimeout, @@ -227,6 +230,23 @@ def _get_atexit_deadline() -> float: return _atexit_deadline +def _positive_config_value( + name: str, value, *, integer: bool = False, maximum: Optional[int] = None +): + """Return ``value`` if it is positive and no larger than ``maximum``. + + Bad lane config is a programming error, so it raises instead of falling + back to a default. + """ + allowed = (int,) if integer else (int, float) + if isinstance(value, bool) or not isinstance(value, allowed) or value <= 0: + kind = "integer" if integer else "number" + raise ValueError(f"{name} must be a positive {kind}, got {value!r}") + if maximum is not None and value > maximum: + raise ValueError(f"{name} must be at most {maximum}, got {value!r}") + return value + + def get_identity_state(passed) -> tuple[str, bool]: """Returns the distinct id to use, and whether this is a personless event or not""" stringified = stringify_id(passed) @@ -708,6 +728,10 @@ def __init__( exception_autocapture_refill_rate=ExceptionCapture.DEFAULT_REFILL_RATE, exception_autocapture_refill_interval_seconds=ExceptionCapture.DEFAULT_REFILL_INTERVAL_SECONDS, capture_compression: Optional[Union[CaptureCompression, str]] = None, + capture_ai_compression: Optional[Union[CaptureCompression, str]] = None, + capture_ai_max_queue_size: int = 1000, + capture_ai_timeout: float = 30, + capture_ai_max_event_bytes: int = AI_MAX_MSG_SIZE, secret_key=None, metrics: Optional[dict] = None, enable_full_ai_capture=False, @@ -726,7 +750,8 @@ def __init__( the corresponding ingestion host. debug: Enable verbose SDK logging and re-raise errors from public API methods. - max_queue_size: Maximum number of events buffered before upload. + max_queue_size: Maximum number of analytics events buffered before + upload. AI events use ``capture_ai_max_queue_size``. send: If False, queueing succeeds but events are not sent. on_error: Optional callback ``(error, batch)`` invoked when an upload fails: by background consumers, or on the calling thread in @@ -843,6 +868,21 @@ def __init__( strings ``"gzip"``/``"deflate"``). When omitted, the ``POSTHOG_CAPTURE_COMPRESSION`` env var is consulted, then no compression. + capture_ai_compression: Request-body compression for + ``capture_ai()`` uploads, set independently of + ``capture_compression``. Defaults to no compression, and the + env var does not apply. ``CaptureCompression.ZSTD`` suits large + AI payloads. + capture_ai_max_queue_size: Maximum number of AI events buffered + before upload. Defaults to 1000, lower than ``max_queue_size`` + because AI events are much larger. + capture_ai_timeout: Seconds allowed for one AI upload request. + Defaults to 30, longer than ``timeout`` because AI batches are + much larger. + capture_ai_max_event_bytes: Largest serialized AI event the SDK + sends; a larger one is dropped with an error log. Defaults to + the AI endpoint's ceiling plus envelope headroom, and may only + be lowered. Examples: ```python @@ -946,6 +986,21 @@ def __init__( self._library_version = VERSION self._sdk_info = f"{self._library_id}/{self._library_version}" self.capture_compression = _resolve_capture_compression(capture_compression) + self.capture_ai_compression = _resolve_capture_ai_compression( + capture_ai_compression + ) + capture_ai_max_queue_size = _positive_config_value( + "capture_ai_max_queue_size", capture_ai_max_queue_size, integer=True + ) + capture_ai_timeout = _positive_config_value( + "capture_ai_timeout", capture_ai_timeout + ) + capture_ai_max_event_bytes = _positive_config_value( + "capture_ai_max_event_bytes", + capture_ai_max_event_bytes, + integer=True, + maximum=AI_MAX_MSG_SIZE, + ) self.super_properties = super_properties # Release id from POSTHOG_RELEASE_ID, attached to every event. Resolved # here so the env var is read once per client. @@ -1055,34 +1110,36 @@ def __init__( api_key=self.api_key, host=self.host, on_error=on_error, - max_queue_size=max_queue_size, thread_count=thread, send=send, flush_at=flush_at, flush_interval=flush_interval, max_retries=self.max_retries, - timeout=timeout, historical_migration=historical_migration, sdk_info=self._sdk_info, ) self._analytics_lane = _Lane( name="analytics", **lane_defaults, + max_queue_size=max_queue_size, + timeout=timeout, endpoint=_CAPTURE_V1_PATH, max_msg_size=MAX_MSG_SIZE, capture_compression=self.capture_compression, eager_start=not sync_mode, ) # The AI lane posts to its own endpoint so multi-MB AI events stay off - # the analytics endpoint's smaller caps. It sends uncompressed. Lazy - # start, so the many clients that never emit AI events pay for no extra - # threads. + # the analytics endpoint's smaller caps, with its own queue, timeout, + # size guard and compression. Lazy start, so the many clients that never + # emit AI events pay for no extra threads. self._ai_lane = _Lane( name="ai", **lane_defaults, + max_queue_size=capture_ai_max_queue_size, + timeout=capture_ai_timeout, endpoint=_CAPTURE_AI_V1_PATH, - max_msg_size=AI_MAX_MSG_SIZE, - capture_compression=CaptureCompression.NONE, + max_msg_size=capture_ai_max_event_bytes, + capture_compression=self.capture_ai_compression, eager_start=False, ) self._lanes = [self._analytics_lane, self._ai_lane] @@ -2435,6 +2492,23 @@ def _enqueue(self, msg, disable_geoip, lane=None, property_allowlist=None): if self.sync_mode: self.log.debug("enqueued with blocking %s.", msg["event"]) + try: + event_size = len(json.dumps(msg, cls=_DatetimeSerializer).encode()) + except Exception: + self.log.error("Unable to serialize event for sizing, dropping.") + return None + if event_size > lane.max_msg_size: + # Log only name and size: AI events may carry unredacted + # multimodal payloads that must not leak into logs. + self.log.error( + "Event %s (%d bytes) exceeds the %dKiB limit for %s, dropping.", + msg["event"], + event_size, + lane.max_msg_size // 1024, + lane.endpoint, + ) + return None + def send_sync() -> None: # Sync mode bypasses the lane's queue but keeps its wire config, # so AI events still post to the AI endpoint. @@ -2443,7 +2517,7 @@ def send_sync() -> None: self.host, [msg], compression=lane.capture_compression, - timeout=self.timeout, + timeout=lane.timeout, max_retries=self.max_retries, historical_migration=self.historical_migration, sdk_info=self._sdk_info, @@ -2681,13 +2755,14 @@ def flush(self, timeout_seconds: Optional[float] = 10) -> None: # Spans drain with events: serverless handlers call flush(), not # shutdown(), and leaving spans on their own timer would lose them. span_flush = self._start_span_flush(timeout_seconds) - if timeout_seconds is None: - for lane in self._lanes: - lane.flush(None) - else: - deadline = time.monotonic() + timeout_seconds - for lane in self._lanes: - lane.flush(max(0.0, deadline - time.monotonic())) + with self._drain_lanes_together(): + if timeout_seconds is None: + for lane in self._lanes: + lane.flush(None) + else: + deadline = time.monotonic() + timeout_seconds + for lane in self._lanes: + lane.flush(max(0.0, deadline - time.monotonic())) if span_flush is not None: # The last span request is bounded only by the request # timeout, so the wait is not. @@ -2824,7 +2899,36 @@ def _run_lifecycle_cleanup( self.log.exception(log_message) errors.append(error) + @contextmanager + def _drain_lanes_together(self): + """Signal every lane to drain before waiting on any of them. + + Lanes then drain in parallel under one budget. Otherwise a lane keeps + batching on its normal cadence while the client waits on the lane before it. + """ + signals: list[_DrainSignal] = [] + try: + for lane in self._lanes: + signal = lane._drain_signal + signal.request() + signals.append(signal) + yield + finally: + for signal in signals: + signal.complete() + def _flush_or_discard_queues(self, errors: list[Exception]) -> None: + try: + with self._drain_lanes_together(): + self._flush_or_discard_each_lane(errors) + return + except Exception as error: + self.log.exception("Failed to signal lane drains during lifecycle cleanup") + errors.append(error) + # Each lane's flush signals its own drain, so lanes still drain one by one. + self._flush_or_discard_each_lane(errors) + + def _flush_or_discard_each_lane(self, errors: list[Exception]) -> None: for lane in self._lanes: try: if any(consumer.is_alive() for consumer in lane.consumers): diff --git a/posthog/consumer.py b/posthog/consumer.py index 22944064..a0779a3e 100644 --- a/posthog/consumer.py +++ b/posthog/consumer.py @@ -21,13 +21,21 @@ MAX_MSG_SIZE = 900 * 1024 # 900KiB per event -# AI events carry LLM inputs/outputs and post to a dedicated endpoint whose -# pipeline accepts larger messages than analytics ingestion, so the AI lane -# grants a higher per-event ceiling. `next()` appends an item before checking -# BATCH_SIZE_LIMIT, so worst-case request body is BATCH_SIZE_LIMIT + -# AI_MAX_MSG_SIZE (~13MiB) — keep that sum under the 20MiB server body cap. -AI_MAX_MSG_SIZE = 8 * 1024 * 1024 # 8MiB per event - +# The AI endpoint's per-event ceiling. The endpoint applies it to the +# serialized properties alone. +AI_MAX_PROPERTIES_SIZE = 8 * 1024 * 1024 +# The local guard measures the whole serialized event, so it adds headroom for +# the rest of the event. Without it, an event whose properties sit at the +# ceiling is refused here although the endpoint accepts it. The guard stays +# coarse: it skips a doomed multi-megabyte upload, it does not reproduce the +# endpoint's check. +AI_ENVELOPE_HEADROOM = 64 * 1024 +# The AI lane's per-event guard, and the upper bound for +# `capture_ai_max_event_bytes`. +AI_MAX_MSG_SIZE = AI_MAX_PROPERTIES_SIZE + AI_ENVELOPE_HEADROOM + +# A batch closes before it appends an event that would take it past this, so +# a request carries at most this much event data, or one larger event alone. # The maximum request body size is currently 20MiB, let's be conservative # in case we want to lower it in the future. BATCH_SIZE_LIMIT = 5 * 1024 * 1024 @@ -265,6 +273,11 @@ def next(self): queue.task_done() pending_items -= 1 continue + if items and total_size + item_size > BATCH_SIZE_LIMIT: + self._return_to_queue_head(item) + pending_items -= 1 + self.log.debug("hit batch size limit (size: %d)", total_size) + break items.append(item) total_size += item_size if total_size >= BATCH_SIZE_LIMIT: @@ -284,6 +297,18 @@ def next(self): return items + def _return_to_queue_head(self, item) -> None: + """Put a dequeued event back at the head of the queue for the next batch. + + The event stays counted in ``unfinished_tasks``, because it was never + marked done. Keeping it in the queue, not in the consumer, means a + stop, a discard or a fork accounts for it like any other queued event. + """ + queue = self.queue + with queue.not_empty: + queue.queue.appendleft(item) + queue.not_empty.notify() + def request(self, batch): """Upload the batch to this consumer's `endpoint` with the capture v1 partial-retry submitter.""" diff --git a/posthog/test/test_ai_capture_lane.py b/posthog/test/test_ai_capture_lane.py index 568f00fe..458d412a 100644 --- a/posthog/test/test_ai_capture_lane.py +++ b/posthog/test/test_ai_capture_lane.py @@ -1,14 +1,17 @@ +import os import threading import unittest import uuid from unittest import mock +from parameterized import parameterized + import posthog from posthog.ai.utils import _capture_ai_event, finalize_ai_content, with_privacy_mode -from posthog.capture_compression import CaptureCompression +from posthog.capture_compression import CAPTURE_COMPRESSION_ENV_VAR, CaptureCompression from posthog.client import Client -from posthog.consumer import AI_MAX_MSG_SIZE, MAX_MSG_SIZE +from posthog.consumer import AI_MAX_MSG_SIZE, AI_MAX_PROPERTIES_SIZE, MAX_MSG_SIZE from posthog.capture_send import _CAPTURE_AI_V1_PATH, _CAPTURE_V1_PATH from posthog.version import VERSION from posthog.test.capture_helpers import patch_capture_send, sent_batch @@ -174,14 +177,45 @@ def test_ai_lane_accepts_multi_megabyte_events(self): batch = consumer.next() self.assertEqual([e["event"] for e in batch], ["$ai_generation"]) - def test_ai_lane_drops_events_over_its_cap(self): - client = self._client() + @parameterized.expand( + [ + ("properties_at_endpoint_ceiling", {}, AI_MAX_PROPERTIES_SIZE, True), + ("over_guard", {}, AI_MAX_MSG_SIZE, False), + ( + "over_lowered_cap", + {"capture_ai_max_event_bytes": 1024 * 1024}, + 2 * 1024 * 1024, + False, + ), + ] + ) + def test_ai_lane_size_guard(self, _name, config, payload_bytes, accepted): + client = Client(TEST_API_KEY, send=False, flush_interval=0.05, **config) client._ai_lane.start() consumer = client._ai_lane.consumers[0] - client._ai_lane.queue.put(self._sized_event("$ai_generation", AI_MAX_MSG_SIZE)) - self.assertEqual(consumer.next(), []) + client._ai_lane.queue.put(self._sized_event("$ai_generation", payload_bytes)) + self.assertEqual( + [e["event"] for e in consumer.next()], + ["$ai_generation"] if accepted else [], + ) self.assertTrue(client._ai_lane.queue.empty()) + def test_sync_mode_ai_event_over_cap_is_not_sent(self): + client = Client( + TEST_API_KEY, sync_mode=True, capture_ai_max_event_bytes=1024 * 1024 + ) + with patch_capture_send("client") as mock_send: + with self.assertLogs("posthog", level="ERROR") as logs: + result = client.capture_ai( + "$ai_generation", + distinct_id="d", + properties={"p": "x" * (2 * 1024 * 1024)}, + ) + + self.assertIsNone(result) + mock_send.assert_not_called() + self.assertIn("exceeds the 1024KiB limit", "\n".join(logs.output)) + def test_analytics_lane_rejects_events_over_900kib(self): client = self._client() consumer = client.consumers[0] @@ -191,18 +225,74 @@ def test_analytics_lane_rejects_events_over_900kib(self): class TestAiLaneWireConfig(unittest.TestCase): - """The AI lane posts to the AI endpoint uncompressed, whatever the - analytics `capture_compression`.""" + """The AI lane has its own endpoint, compression, timeout, queue and size + guard, independent of the analytics lane's settings.""" - def test_ai_lane_consumers_use_ai_endpoint_without_compression(self): - client = Client(TEST_API_KEY, send=False, capture_compression="gzip", thread=2) + @parameterized.expand( + [ + ("defaults", {}, CaptureCompression.NONE, 30, 1000, AI_MAX_MSG_SIZE), + ( + "configured", + { + "capture_ai_compression": "zstd", + "capture_ai_timeout": 45, + "capture_ai_max_queue_size": 50, + "capture_ai_max_event_bytes": 1024 * 1024, + }, + CaptureCompression.ZSTD, + 45, + 50, + 1024 * 1024, + ), + ] + ) + def test_ai_lane_consumers_use_ai_config( + self, _name, config, compression, timeout, queue_size, max_event_bytes + ): + with mock.patch.dict(os.environ, {CAPTURE_COMPRESSION_ENV_VAR: "deflate"}): + client = Client( + TEST_API_KEY, + send=False, + capture_compression="gzip", + timeout=9, + max_queue_size=77, + thread=2, + **config, + ) client._ai_lane.start() + self.assertEqual(client._ai_lane.queue.maxsize, queue_size) self.assertEqual(len(client._ai_lane.consumers), 2) for consumer in client._ai_lane.consumers: self.assertIs(consumer.queue, client._ai_lane.queue) self.assertEqual(consumer.endpoint, _CAPTURE_AI_V1_PATH) - self.assertEqual(consumer.max_msg_size, AI_MAX_MSG_SIZE) - self.assertEqual(consumer.capture_compression, CaptureCompression.NONE) + self.assertEqual(consumer.max_msg_size, max_event_bytes) + self.assertEqual(consumer.capture_compression, compression) + self.assertEqual(consumer.timeout, timeout) + self.assertEqual(client.queue.maxsize, 77) + for consumer in client._analytics_lane.consumers: + self.assertEqual(consumer.capture_compression, CaptureCompression.GZIP) + self.assertEqual(consumer.timeout, 9) + client.join() + + @parameterized.expand( + [ + ( + "event_bytes_over_ceiling", + "capture_ai_max_event_bytes", + AI_MAX_MSG_SIZE + 1, + ), + ("event_bytes_zero", "capture_ai_max_event_bytes", 0), + ("event_bytes_float", "capture_ai_max_event_bytes", 1024.5), + ("event_bytes_bool", "capture_ai_max_event_bytes", True), + ("queue_size_negative", "capture_ai_max_queue_size", -1), + ("queue_size_string", "capture_ai_max_queue_size", "100"), + ("timeout_zero", "capture_ai_timeout", 0), + ("compression_unknown", "capture_ai_compression", "br"), + ] + ) + def test_invalid_ai_config_raises(self, _name, field, value): + with self.assertRaisesRegex(ValueError, field): + Client(TEST_API_KEY, send=False, **{field: value}) def test_async_lanes_keep_separate_path_and_compression(self): client = Client(TEST_API_KEY, capture_compression="gzip", flush_interval=0.05) @@ -228,23 +318,57 @@ def test_async_lanes_keep_separate_path_and_compression(self): ) client.join() - def test_sync_lanes_keep_separate_path_and_compression(self): - client = Client(TEST_API_KEY, sync_mode=True, capture_compression="gzip") + def test_sync_lanes_keep_separate_path_compression_and_timeout(self): + client = Client( + TEST_API_KEY, + sync_mode=True, + capture_compression="gzip", + capture_ai_compression="deflate", + timeout=9, + capture_ai_timeout=45, + ) with patch_capture_send("client") as mock_send: client.capture_ai("$ai_generation", distinct_id="d") client.capture("button_clicked", distinct_id="d") self.assertEqual( [ - (call.kwargs["path"], call.kwargs["compression"]) + ( + call.kwargs["path"], + call.kwargs["compression"], + call.kwargs["timeout"], + ) for call in mock_send.call_args_list ], [ - (_CAPTURE_AI_V1_PATH, CaptureCompression.NONE), - (_CAPTURE_V1_PATH, CaptureCompression.GZIP), + (_CAPTURE_AI_V1_PATH, CaptureCompression.DEFLATE, 45), + (_CAPTURE_V1_PATH, CaptureCompression.GZIP, 9), ], ) + def test_flush_drains_ai_lane_while_waiting_on_analytics(self): + release_analytics = threading.Event() + ai_sent = threading.Event() + + def send(api_key, host, batch, **kwargs): + if kwargs["path"] == _CAPTURE_AI_V1_PATH: + ai_sent.set() + else: + release_analytics.wait(5) + + client = Client(TEST_API_KEY, flush_interval=30) + with patch_capture_send("consumer", side_effect=send): + client.capture("button_clicked", distinct_id="d") + client.capture_ai("$ai_generation", distinct_id="d") + flusher = threading.Thread(target=client.flush, args=(10,)) + flusher.start() + try: + self.assertTrue(ai_sent.wait(2)) + finally: + release_analytics.set() + flusher.join(5) + client.join() + class TestAiLaneLazyStart(unittest.TestCase): def test_no_ai_consumers_until_first_capture_ai(self): diff --git a/posthog/test/test_consumer.py b/posthog/test/test_consumer.py index ca474882..f4cf3188 100644 --- a/posthog/test/test_consumer.py +++ b/posthog/test/test_consumer.py @@ -204,6 +204,27 @@ def test_max_msg_size_param_raises_per_event_ceiling(self) -> None: q.put(big_msg) self.assertEqual(consumer.next(), [big_msg]) + @parameterized.expand( + [ + # A small event serializes to 18 bytes, the "big" one to 81. + ("closes_before_overflow", 43, [0, 1, 2, 3], [[0, 1], [2, 3]]), + ("event_over_limit_goes_alone", 30, [0, "big", 1], [[0], ["big"], [1]]), + ] + ) + def test_batch_byte_limit_is_checked_before_appending( + self, _name, limit, ids, expected + ) -> None: + q = Queue() + consumer = Consumer(q, "", flush_at=10, flush_interval=0.01) + for i in ids: + q.put({"m": "x" * (60 if i == "big" else 1), "i": i}) + + with mock.patch("posthog.consumer.BATCH_SIZE_LIMIT", limit): + batches = [[e["i"] for e in consumer.next()] for _ in expected] + + self.assertEqual(batches, expected) + self.assertEqual(q.unfinished_tasks, len(ids)) + def test_upload(self) -> None: q = Queue() consumer = Consumer(q, TEST_API_KEY, flush_at=1) diff --git a/references/public_api_snapshot.txt b/references/public_api_snapshot.txt index 0cb0bc89..282645df 100644 --- a/references/public_api_snapshot.txt +++ b/references/public_api_snapshot.txt @@ -700,6 +700,7 @@ attribute posthog.capture_send.CaptureEventResult.details: Optional[str] = None attribute posthog.capture_send.CaptureEventResult.result: Optional[str] attribute posthog.capture_trace_context = False attribute posthog.client.Client.api_key = (project_api_key or '').strip() +attribute posthog.client.Client.capture_ai_compression = _resolve_capture_ai_compression(capture_ai_compression) attribute posthog.client.Client.capture_compression = _resolve_capture_compression(capture_compression) attribute posthog.client.Client.capture_exception_code_variables = capture_exception_code_variables attribute posthog.client.Client.capture_trace_context = capture_trace_context @@ -755,7 +756,9 @@ attribute posthog.code_variables_detect_secrets = DEFAULT_CODE_VARIABLES_DETECT_ attribute posthog.code_variables_ignore_patterns = DEFAULT_CODE_VARIABLES_IGNORE_PATTERNS attribute posthog.code_variables_mask_patterns = DEFAULT_CODE_VARIABLES_MASK_PATTERNS attribute posthog.code_variables_mask_url_credentials = DEFAULT_CODE_VARIABLES_MASK_URL_CREDENTIALS -attribute posthog.consumer.AI_MAX_MSG_SIZE = 8 * 1024 * 1024 +attribute posthog.consumer.AI_ENVELOPE_HEADROOM = 64 * 1024 +attribute posthog.consumer.AI_MAX_MSG_SIZE = AI_MAX_PROPERTIES_SIZE + AI_ENVELOPE_HEADROOM +attribute posthog.consumer.AI_MAX_PROPERTIES_SIZE = 8 * 1024 * 1024 attribute posthog.consumer.BATCH_SIZE_LIMIT = 5 * 1024 * 1024 attribute posthog.consumer.Consumer.api_key = api_key attribute posthog.consumer.Consumer.capture_compression = capture_compression @@ -1170,7 +1173,7 @@ class posthog.bucketed_rate_limiter.BucketedRateLimiter(bucket_size: Number, ref class posthog.capture_compression.CaptureCompression class posthog.capture_send.CaptureError(status: int | str, message: str, *, endpoint: str, retry_after: Optional[float] = None, request_id: Optional[str] = None, attempts: Optional[int] = None, retry_exhausted: Optional[list[str]] = None, drops: Optional[list[tuple[str, Optional[str]]]] = None, event_results: Optional[dict[str, CaptureEventResult]] = None) class posthog.capture_send.CaptureEventResult(result: Optional[str], details: Optional[str] = None) -class posthog.client.Client(project_api_key: str, host=None, *, debug=False, max_queue_size=10000, send=True, on_error=None, flush_at=100, flush_interval=5.0, max_retries=3, sync_mode=False, timeout=15, thread=1, poll_interval=30, personal_api_key=None, disabled=False, disable_geoip=True, is_server=True, historical_migration=False, feature_flags_request_timeout_seconds=3, feature_flags_request_max_retries=1, super_properties=None, enable_exception_autocapture=False, log_captured_exceptions=False, project_root=None, privacy_mode=False, before_send=None, flag_fallback_cache_url=None, enable_local_evaluation=True, flag_definition_cache_provider: Optional[FlagDefinitionCacheProvider] = None, capture_exception_code_variables=False, code_variables_mask_patterns=None, code_variables_ignore_patterns=None, code_variables_mask_url_credentials=None, code_variables_detect_secrets=None, in_app_modules: list[str] | None = None, enable_exception_autocapture_rate_limiting=False, exception_autocapture_bucket_size=ExceptionCapture.DEFAULT_BUCKET_SIZE, exception_autocapture_refill_rate=ExceptionCapture.DEFAULT_REFILL_RATE, exception_autocapture_refill_interval_seconds=ExceptionCapture.DEFAULT_REFILL_INTERVAL_SECONDS, capture_compression: Optional[Union[CaptureCompression, str]] = None, secret_key=None, metrics: Optional[dict] = None, enable_full_ai_capture=False, capture_trace_context=False, _use_ai_lane=False, _enable_multimodal_capture=False, traces: Optional[dict] = None) +class posthog.client.Client(project_api_key: str, host=None, *, debug=False, max_queue_size=10000, send=True, on_error=None, flush_at=100, flush_interval=5.0, max_retries=3, sync_mode=False, timeout=15, thread=1, poll_interval=30, personal_api_key=None, disabled=False, disable_geoip=True, is_server=True, historical_migration=False, feature_flags_request_timeout_seconds=3, feature_flags_request_max_retries=1, super_properties=None, enable_exception_autocapture=False, log_captured_exceptions=False, project_root=None, privacy_mode=False, before_send=None, flag_fallback_cache_url=None, enable_local_evaluation=True, flag_definition_cache_provider: Optional[FlagDefinitionCacheProvider] = None, capture_exception_code_variables=False, code_variables_mask_patterns=None, code_variables_ignore_patterns=None, code_variables_mask_url_credentials=None, code_variables_detect_secrets=None, in_app_modules: list[str] | None = None, enable_exception_autocapture_rate_limiting=False, exception_autocapture_bucket_size=ExceptionCapture.DEFAULT_BUCKET_SIZE, exception_autocapture_refill_rate=ExceptionCapture.DEFAULT_REFILL_RATE, exception_autocapture_refill_interval_seconds=ExceptionCapture.DEFAULT_REFILL_INTERVAL_SECONDS, capture_compression: Optional[Union[CaptureCompression, str]] = None, capture_ai_compression: Optional[Union[CaptureCompression, str]] = None, capture_ai_max_queue_size: int = 1000, capture_ai_timeout: float = 30, capture_ai_max_event_bytes: int = AI_MAX_MSG_SIZE, secret_key=None, metrics: Optional[dict] = None, enable_full_ai_capture=False, capture_trace_context=False, _use_ai_lane=False, _enable_multimodal_capture=False, traces: Optional[dict] = None) class posthog.consumer.Consumer(queue, api_key, flush_at=100, host=None, on_error=None, flush_interval=5.0, retries=3, timeout=15, historical_migration=False, endpoint=_CAPTURE_V1_PATH, max_msg_size=MAX_MSG_SIZE, capture_compression=CaptureCompression.NONE) class posthog.contexts.ContextScope(parent=None, fresh: bool = False, capture_exceptions: bool = True, client: Optional[Client] = None) class posthog.exception_capture.ExceptionCapture(client: Client, rate_limiting_enabled=False, bucket_size=DEFAULT_BUCKET_SIZE, refill_rate=DEFAULT_REFILL_RATE, refill_interval_seconds=DEFAULT_REFILL_INTERVAL_SECONDS) From 44d141dded88f6a37399852856cfc757122a540c Mon Sep 17 00:00:00 2001 From: Eli Reisman Date: Wed, 7 Oct 2026 10:41:10 -0700 Subject: [PATCH 2/3] chore: tighten AI lane comments --- posthog/client.py | 6 ++---- posthog/consumer.py | 7 ++----- 2 files changed, 4 insertions(+), 9 deletions(-) diff --git a/posthog/client.py b/posthog/client.py index 856692a7..4896000a 100644 --- a/posthog/client.py +++ b/posthog/client.py @@ -1128,10 +1128,8 @@ def __init__( capture_compression=self.capture_compression, eager_start=not sync_mode, ) - # The AI lane posts to its own endpoint so multi-MB AI events stay off - # the analytics endpoint's smaller caps, with its own queue, timeout, - # size guard and compression. Lazy start, so the many clients that never - # emit AI events pay for no extra threads. + # A separate endpoint keeps multi-MB AI events off the analytics caps. + # Lazy start, so clients that never send AI events pay for no threads. self._ai_lane = _Lane( name="ai", **lane_defaults, diff --git a/posthog/consumer.py b/posthog/consumer.py index a0779a3e..46c8c6ab 100644 --- a/posthog/consumer.py +++ b/posthog/consumer.py @@ -24,11 +24,8 @@ # The AI endpoint's per-event ceiling. The endpoint applies it to the # serialized properties alone. AI_MAX_PROPERTIES_SIZE = 8 * 1024 * 1024 -# The local guard measures the whole serialized event, so it adds headroom for -# the rest of the event. Without it, an event whose properties sit at the -# ceiling is refused here although the endpoint accepts it. The guard stays -# coarse: it skips a doomed multi-megabyte upload, it does not reproduce the -# endpoint's check. +# The local guard measures the whole event, not only its properties, so it +# allows this much more to keep events at the endpoint's ceiling. AI_ENVELOPE_HEADROOM = 64 * 1024 # The AI lane's per-event guard, and the upper bound for # `capture_ai_max_event_bytes`. From 682fd47fbcf29370c994ec1b2bdb8766f707aff3 Mon Sep 17 00:00:00 2001 From: Eli Reisman Date: Thu, 8 Oct 2026 16:38:54 -0700 Subject: [PATCH 3/3] fix: reject non-finite AI lane config values --- posthog/client.py | 10 ++++++++-- posthog/test/test_ai_capture_lane.py | 2 ++ 2 files changed, 10 insertions(+), 2 deletions(-) diff --git a/posthog/client.py b/posthog/client.py index 4896000a..a0af8531 100644 --- a/posthog/client.py +++ b/posthog/client.py @@ -4,6 +4,7 @@ import json from contextlib import contextmanager import logging +import math import os import sys import threading @@ -233,13 +234,18 @@ def _get_atexit_deadline() -> float: def _positive_config_value( name: str, value, *, integer: bool = False, maximum: Optional[int] = None ): - """Return ``value`` if it is positive and no larger than ``maximum``. + """Return ``value`` if it is finite, positive and no larger than ``maximum``. Bad lane config is a programming error, so it raises instead of falling back to a default. """ allowed = (int,) if integer else (int, float) - if isinstance(value, bool) or not isinstance(value, allowed) or value <= 0: + if ( + isinstance(value, bool) + or not isinstance(value, allowed) + or not math.isfinite(value) + or value <= 0 + ): kind = "integer" if integer else "number" raise ValueError(f"{name} must be a positive {kind}, got {value!r}") if maximum is not None and value > maximum: diff --git a/posthog/test/test_ai_capture_lane.py b/posthog/test/test_ai_capture_lane.py index 458d412a..2234de1b 100644 --- a/posthog/test/test_ai_capture_lane.py +++ b/posthog/test/test_ai_capture_lane.py @@ -287,6 +287,8 @@ def test_ai_lane_consumers_use_ai_config( ("queue_size_negative", "capture_ai_max_queue_size", -1), ("queue_size_string", "capture_ai_max_queue_size", "100"), ("timeout_zero", "capture_ai_timeout", 0), + ("timeout_nan", "capture_ai_timeout", float("nan")), + ("timeout_inf", "capture_ai_timeout", float("inf")), ("compression_unknown", "capture_ai_compression", "br"), ] )