Skip to content

Commit 8da410d

Browse files
committed
feat: AI lane config and size-aware batching
1 parent 9e85baf commit 8da410d

7 files changed

Lines changed: 341 additions & 44 deletions

File tree

‎AGENTS.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@ Follow [Public API changes](./CONTRIBUTING.md#public-api-changes). As an agent,
2323

2424
Before changing capture configuration, serialization, routing, or retries, read the relevant implementation and tests.
2525

26-
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`.
26+
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`.
2727

2828
## Mirror and build safety
2929

‎posthog/capture_compression.py‎

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,7 @@ def _zstd_available() -> bool:
5454

5555
def _coerce_explicit(
5656
value: Union[CaptureCompression, str],
57+
name: str = "capture_compression",
5758
) -> CaptureCompression:
5859
"""Normalize an explicitly-supplied compression to a ``CaptureCompression``.
5960
@@ -68,7 +69,7 @@ def _coerce_explicit(
6869
if resolved is not None:
6970
return resolved
7071
raise ValueError(
71-
f"invalid capture_compression {value!r}; expected a CaptureCompression "
72+
f"invalid {name} {value!r}; expected a CaptureCompression "
7273
f"or one of {sorted(_ALIASES)}"
7374
)
7475

@@ -122,3 +123,22 @@ def _resolve_capture_compression(
122123
)
123124
return fallback
124125
return env_resolved
126+
127+
128+
def _resolve_capture_ai_compression(
129+
capture_ai_compression: Optional[Union[CaptureCompression, str]] = None,
130+
) -> CaptureCompression:
131+
"""Resolve the AI lane's request-body compression.
132+
133+
Explicit argument only, defaulting to ``NONE``. ``POSTHOG_CAPTURE_COMPRESSION``
134+
does not apply, so changing analytics compression never changes AI uploads.
135+
"""
136+
if capture_ai_compression is None:
137+
return CaptureCompression.NONE
138+
resolved = _coerce_explicit(capture_ai_compression, "capture_ai_compression")
139+
if resolved is CaptureCompression.ZSTD and not _zstd_available():
140+
raise ValueError(
141+
"capture_ai_compression 'zstd' requires the zstandard package; "
142+
"install posthog[zstd]"
143+
)
144+
return resolved

‎posthog/client.py‎

Lines changed: 120 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22
import hashlib as _hashlib
33
import inspect
44
import json
5+
from contextlib import contextmanager
56
import logging
67
import os
78
import sys
@@ -29,6 +30,7 @@
2930
from posthog.tracing._span import inert_span as _inert_span
3031
from posthog.capture_compression import (
3132
CaptureCompression,
33+
_resolve_capture_ai_compression,
3234
_resolve_capture_compression,
3335
)
3436
from posthog.capture_send import (
@@ -91,6 +93,7 @@
9193
from posthog.request import (
9294
USER_AGENT as _USER_AGENT,
9395
APIError,
96+
DatetimeSerializer as _DatetimeSerializer,
9497
QuotaLimitError,
9598
RequestsConnectionError,
9699
RequestsTimeout,
@@ -227,6 +230,23 @@ def _get_atexit_deadline() -> float:
227230
return _atexit_deadline
228231

229232

233+
def _positive_config_value(
234+
name: str, value, *, integer: bool = False, maximum: Optional[int] = None
235+
):
236+
"""Return ``value`` if it is positive and no larger than ``maximum``.
237+
238+
Bad lane config is a programming error, so it raises instead of falling
239+
back to a default.
240+
"""
241+
allowed = (int,) if integer else (int, float)
242+
if isinstance(value, bool) or not isinstance(value, allowed) or value <= 0:
243+
kind = "integer" if integer else "number"
244+
raise ValueError(f"{name} must be a positive {kind}, got {value!r}")
245+
if maximum is not None and value > maximum:
246+
raise ValueError(f"{name} must be at most {maximum}, got {value!r}")
247+
return value
248+
249+
230250
def get_identity_state(passed) -> tuple[str, bool]:
231251
"""Returns the distinct id to use, and whether this is a personless event or not"""
232252
stringified = stringify_id(passed)
@@ -708,6 +728,10 @@ def __init__(
708728
exception_autocapture_refill_rate=ExceptionCapture.DEFAULT_REFILL_RATE,
709729
exception_autocapture_refill_interval_seconds=ExceptionCapture.DEFAULT_REFILL_INTERVAL_SECONDS,
710730
capture_compression: Optional[Union[CaptureCompression, str]] = None,
731+
capture_ai_compression: Optional[Union[CaptureCompression, str]] = None,
732+
capture_ai_max_queue_size: int = 1000,
733+
capture_ai_timeout: float = 30,
734+
capture_ai_max_event_bytes: int = AI_MAX_MSG_SIZE,
711735
secret_key=None,
712736
metrics: Optional[dict] = None,
713737
enable_full_ai_capture=False,
@@ -726,7 +750,8 @@ def __init__(
726750
the corresponding ingestion host.
727751
debug: Enable verbose SDK logging and re-raise errors from public
728752
API methods.
729-
max_queue_size: Maximum number of events buffered before upload.
753+
max_queue_size: Maximum number of analytics events buffered before
754+
upload. AI events use ``capture_ai_max_queue_size``.
730755
send: If False, queueing succeeds but events are not sent.
731756
on_error: Optional callback ``(error, batch)`` invoked when an upload
732757
fails: by background consumers, or on the calling thread in
@@ -843,6 +868,21 @@ def __init__(
843868
strings ``"gzip"``/``"deflate"``). When omitted, the
844869
``POSTHOG_CAPTURE_COMPRESSION`` env var is consulted, then no
845870
compression.
871+
capture_ai_compression: Request-body compression for
872+
``capture_ai()`` uploads, set independently of
873+
``capture_compression``. Defaults to no compression, and the
874+
env var does not apply. ``CaptureCompression.ZSTD`` suits large
875+
AI payloads.
876+
capture_ai_max_queue_size: Maximum number of AI events buffered
877+
before upload. Defaults to 1000, lower than ``max_queue_size``
878+
because AI events are much larger.
879+
capture_ai_timeout: Seconds allowed for one AI upload request.
880+
Defaults to 30, longer than ``timeout`` because AI batches are
881+
much larger.
882+
capture_ai_max_event_bytes: Largest serialized AI event the SDK
883+
sends; a larger one is dropped with an error log. Defaults to
884+
the AI endpoint's ceiling plus envelope headroom, and may only
885+
be lowered.
846886
847887
Examples:
848888
```python
@@ -946,6 +986,21 @@ def __init__(
946986
self._library_version = VERSION
947987
self._sdk_info = f"{self._library_id}/{self._library_version}"
948988
self.capture_compression = _resolve_capture_compression(capture_compression)
989+
self.capture_ai_compression = _resolve_capture_ai_compression(
990+
capture_ai_compression
991+
)
992+
capture_ai_max_queue_size = _positive_config_value(
993+
"capture_ai_max_queue_size", capture_ai_max_queue_size, integer=True
994+
)
995+
capture_ai_timeout = _positive_config_value(
996+
"capture_ai_timeout", capture_ai_timeout
997+
)
998+
capture_ai_max_event_bytes = _positive_config_value(
999+
"capture_ai_max_event_bytes",
1000+
capture_ai_max_event_bytes,
1001+
integer=True,
1002+
maximum=AI_MAX_MSG_SIZE,
1003+
)
9491004
self.super_properties = super_properties
9501005
# Release id from POSTHOG_RELEASE_ID, attached to every event. Resolved
9511006
# here so the env var is read once per client.
@@ -1055,34 +1110,36 @@ def __init__(
10551110
api_key=self.api_key,
10561111
host=self.host,
10571112
on_error=on_error,
1058-
max_queue_size=max_queue_size,
10591113
thread_count=thread,
10601114
send=send,
10611115
flush_at=flush_at,
10621116
flush_interval=flush_interval,
10631117
max_retries=self.max_retries,
1064-
timeout=timeout,
10651118
historical_migration=historical_migration,
10661119
sdk_info=self._sdk_info,
10671120
)
10681121
self._analytics_lane = _Lane(
10691122
name="analytics",
10701123
**lane_defaults,
1124+
max_queue_size=max_queue_size,
1125+
timeout=timeout,
10711126
endpoint=_CAPTURE_V1_PATH,
10721127
max_msg_size=MAX_MSG_SIZE,
10731128
capture_compression=self.capture_compression,
10741129
eager_start=not sync_mode,
10751130
)
10761131
# The AI lane posts to its own endpoint so multi-MB AI events stay off
1077-
# the analytics endpoint's smaller caps. It sends uncompressed. Lazy
1078-
# start, so the many clients that never emit AI events pay for no extra
1079-
# threads.
1132+
# the analytics endpoint's smaller caps, with its own queue, timeout,
1133+
# size guard and compression. Lazy start, so the many clients that never
1134+
# emit AI events pay for no extra threads.
10801135
self._ai_lane = _Lane(
10811136
name="ai",
10821137
**lane_defaults,
1138+
max_queue_size=capture_ai_max_queue_size,
1139+
timeout=capture_ai_timeout,
10831140
endpoint=_CAPTURE_AI_V1_PATH,
1084-
max_msg_size=AI_MAX_MSG_SIZE,
1085-
capture_compression=CaptureCompression.NONE,
1141+
max_msg_size=capture_ai_max_event_bytes,
1142+
capture_compression=self.capture_ai_compression,
10861143
eager_start=False,
10871144
)
10881145
self._lanes = [self._analytics_lane, self._ai_lane]
@@ -2435,6 +2492,23 @@ def _enqueue(self, msg, disable_geoip, lane=None, property_allowlist=None):
24352492
if self.sync_mode:
24362493
self.log.debug("enqueued with blocking %s.", msg["event"])
24372494

2495+
try:
2496+
event_size = len(json.dumps(msg, cls=_DatetimeSerializer).encode())
2497+
except Exception:
2498+
self.log.error("Unable to serialize event for sizing, dropping.")
2499+
return None
2500+
if event_size > lane.max_msg_size:
2501+
# Log only name and size: AI events may carry unredacted
2502+
# multimodal payloads that must not leak into logs.
2503+
self.log.error(
2504+
"Event %s (%d bytes) exceeds the %dKiB limit for %s, dropping.",
2505+
msg["event"],
2506+
event_size,
2507+
lane.max_msg_size // 1024,
2508+
lane.endpoint,
2509+
)
2510+
return None
2511+
24382512
def send_sync() -> None:
24392513
# Sync mode bypasses the lane's queue but keeps its wire config,
24402514
# so AI events still post to the AI endpoint.
@@ -2443,7 +2517,7 @@ def send_sync() -> None:
24432517
self.host,
24442518
[msg],
24452519
compression=lane.capture_compression,
2446-
timeout=self.timeout,
2520+
timeout=lane.timeout,
24472521
max_retries=self.max_retries,
24482522
historical_migration=self.historical_migration,
24492523
sdk_info=self._sdk_info,
@@ -2679,13 +2753,14 @@ def flush(self, timeout_seconds: Optional[float] = 10) -> None:
26792753
# Spans drain with events: serverless handlers call flush(), not
26802754
# shutdown(), and leaving spans on their own timer would lose them.
26812755
span_flush = self._start_span_flush(timeout_seconds)
2682-
if timeout_seconds is None:
2683-
for lane in self._lanes:
2684-
lane.flush(None)
2685-
else:
2686-
deadline = time.monotonic() + timeout_seconds
2687-
for lane in self._lanes:
2688-
lane.flush(max(0.0, deadline - time.monotonic()))
2756+
with self._drain_lanes_together():
2757+
if timeout_seconds is None:
2758+
for lane in self._lanes:
2759+
lane.flush(None)
2760+
else:
2761+
deadline = time.monotonic() + timeout_seconds
2762+
for lane in self._lanes:
2763+
lane.flush(max(0.0, deadline - time.monotonic()))
26892764
if span_flush is not None:
26902765
# The last span request is bounded only by the request
26912766
# timeout, so the wait is not.
@@ -2822,7 +2897,36 @@ def _run_lifecycle_cleanup(
28222897
self.log.exception(log_message)
28232898
errors.append(error)
28242899

2900+
@contextmanager
2901+
def _drain_lanes_together(self):
2902+
"""Signal every lane to drain before waiting on any of them.
2903+
2904+
Lanes then drain in parallel under one budget. Otherwise a lane keeps
2905+
batching on its normal cadence while the client waits on the lane before it.
2906+
"""
2907+
signals: list[_DrainSignal] = []
2908+
try:
2909+
for lane in self._lanes:
2910+
signal = lane._drain_signal
2911+
signal.request()
2912+
signals.append(signal)
2913+
yield
2914+
finally:
2915+
for signal in signals:
2916+
signal.complete()
2917+
28252918
def _flush_or_discard_queues(self, errors: list[Exception]) -> None:
2919+
try:
2920+
with self._drain_lanes_together():
2921+
self._flush_or_discard_each_lane(errors)
2922+
return
2923+
except Exception as error:
2924+
self.log.exception("Failed to signal lane drains during lifecycle cleanup")
2925+
errors.append(error)
2926+
# Each lane's flush signals its own drain, so lanes still drain one by one.
2927+
self._flush_or_discard_each_lane(errors)
2928+
2929+
def _flush_or_discard_each_lane(self, errors: list[Exception]) -> None:
28262930
for lane in self._lanes:
28272931
try:
28282932
if any(consumer.is_alive() for consumer in lane.consumers):

‎posthog/consumer.py‎

Lines changed: 32 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -21,13 +21,21 @@
2121

2222
MAX_MSG_SIZE = 900 * 1024 # 900KiB per event
2323

24-
# AI events carry LLM inputs/outputs and post to a dedicated endpoint whose
25-
# pipeline accepts larger messages than analytics ingestion, so the AI lane
26-
# grants a higher per-event ceiling. `next()` appends an item before checking
27-
# BATCH_SIZE_LIMIT, so worst-case request body is BATCH_SIZE_LIMIT +
28-
# AI_MAX_MSG_SIZE (~13MiB) — keep that sum under the 20MiB server body cap.
29-
AI_MAX_MSG_SIZE = 8 * 1024 * 1024 # 8MiB per event
30-
24+
# The AI endpoint's per-event ceiling. The endpoint applies it to the
25+
# serialized properties alone.
26+
AI_MAX_PROPERTIES_SIZE = 8 * 1024 * 1024
27+
# The local guard measures the whole serialized event, so it adds headroom for
28+
# the rest of the event. Without it, an event whose properties sit at the
29+
# ceiling is refused here although the endpoint accepts it. The guard stays
30+
# coarse: it skips a doomed multi-megabyte upload, it does not reproduce the
31+
# endpoint's check.
32+
AI_ENVELOPE_HEADROOM = 64 * 1024
33+
# The AI lane's per-event guard, and the upper bound for
34+
# `capture_ai_max_event_bytes`.
35+
AI_MAX_MSG_SIZE = AI_MAX_PROPERTIES_SIZE + AI_ENVELOPE_HEADROOM
36+
37+
# A batch closes before it appends an event that would take it past this, so
38+
# a request carries at most this much event data, or one larger event alone.
3139
# The maximum request body size is currently 20MiB, let's be conservative
3240
# in case we want to lower it in the future.
3341
BATCH_SIZE_LIMIT = 5 * 1024 * 1024
@@ -265,6 +273,11 @@ def next(self):
265273
queue.task_done()
266274
pending_items -= 1
267275
continue
276+
if items and total_size + item_size > BATCH_SIZE_LIMIT:
277+
self._return_to_queue_head(item)
278+
pending_items -= 1
279+
self.log.debug("hit batch size limit (size: %d)", total_size)
280+
break
268281
items.append(item)
269282
total_size += item_size
270283
if total_size >= BATCH_SIZE_LIMIT:
@@ -284,6 +297,18 @@ def next(self):
284297

285298
return items
286299

300+
def _return_to_queue_head(self, item) -> None:
301+
"""Put a dequeued event back at the head of the queue for the next batch.
302+
303+
The event stays counted in ``unfinished_tasks``, because it was never
304+
marked done. Keeping it in the queue, not in the consumer, means a
305+
stop, a discard or a fork accounts for it like any other queued event.
306+
"""
307+
queue = self.queue
308+
with queue.not_empty:
309+
queue.queue.appendleft(item)
310+
queue.not_empty.notify()
311+
287312
def request(self, batch):
288313
"""Upload the batch to this consumer's `endpoint` with the capture v1
289314
partial-retry submitter."""

0 commit comments

Comments
 (0)