Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
22 changes: 21 additions & 1 deletion posthog/capture_compression.py
Original file line number Diff line number Diff line change
Expand Up @@ -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``.

Expand All @@ -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)}"
)

Expand Down Expand Up @@ -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
142 changes: 125 additions & 17 deletions posthog/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,9 @@
import hashlib as _hashlib
import inspect
import json
from contextlib import contextmanager
import logging
import math
import os
import sys
import threading
Expand All @@ -29,6 +31,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 (
Expand Down Expand Up @@ -91,6 +94,7 @@
from posthog.request import (
USER_AGENT as _USER_AGENT,
APIError,
DatetimeSerializer as _DatetimeSerializer,
QuotaLimitError,
RequestsConnectionError,
RequestsTimeout,
Expand Down Expand Up @@ -227,6 +231,28 @@ 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 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 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:
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)
Expand Down Expand Up @@ -708,6 +734,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,
Expand All @@ -726,7 +756,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
Expand Down Expand Up @@ -843,6 +874,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
Expand Down Expand Up @@ -946,6 +992,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.
Expand Down Expand Up @@ -1055,34 +1116,34 @@ 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.
# 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,
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]
Expand Down Expand Up @@ -2435,6 +2496,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.
Expand All @@ -2443,7 +2521,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,
Expand Down Expand Up @@ -2681,13 +2759,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.
Expand Down Expand Up @@ -2824,7 +2903,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):
Expand Down
36 changes: 29 additions & 7 deletions posthog/consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,18 @@

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 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`.
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
Expand Down Expand Up @@ -265,6 +270,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:
Expand All @@ -284,6 +294,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."""
Expand Down
Loading
Loading