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.

Preserve v0 defaults/compatibility; strictly typed v1 options and `$set`/`$set_once` relocation; v1-only compression (zlib-wrapped deflate, optional zstd); 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`.
Preserve capture v1 as the default (v0 is opt-in for analytics only; `capture_ai` always posts v1 to `/i/v1/ai/events`); strictly typed v1 options and `$set`/`$set_once` relocation; v1-only compression (zlib-wrapped deflate, optional zstd); 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
2 changes: 1 addition & 1 deletion posthog/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -425,7 +425,7 @@ def get_tags() -> Dict[str, Any]:
enable_local_evaluation = True # type: bool
flag_definition_cache_provider = None # type: Optional[FlagDefinitionCacheProvider]
# Capture wire protocol for the global client. None defers to POSTHOG_CAPTURE_MODE
# then CaptureMode.V0. See posthog.capture_mode.CaptureMode.
# then CaptureMode.V1. See posthog.capture_mode.CaptureMode.
capture_mode = None # type: Optional[CaptureMode]
# Routes AI SDK wrapper events through the dedicated AI capture lane, skips
# truncation, and passes media unredacted. `privacy_mode` always wins.
Expand Down
17 changes: 8 additions & 9 deletions posthog/capture_mode.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,10 +13,9 @@
class CaptureMode(str, Enum):
"""Selects the capture wire protocol used for event ingestion.

``V0`` is the legacy ``POST /batch/`` endpoint and the default, so upgrading
is transparent to existing callers. ``V1`` opts into
``POST /i/v1/analytics/events`` (Bearer auth, per-event results, partial
retry). Inheriting from ``str`` keeps the members directly comparable to and
``V1`` is ``POST /i/v1/analytics/events`` (Bearer auth, per-event results,
partial retry) and the default. ``V0`` opts back into the legacy
``POST /batch/`` endpoint. Inheriting from ``str`` keeps the members directly comparable to and
serializable as their ``"v0"`` / ``"v1"`` values.
"""

Expand Down Expand Up @@ -61,24 +60,24 @@ def _resolve_capture_mode(
"""Resolve the effective capture mode.

Precedence: explicit ``capture_mode`` argument > ``POSTHOG_CAPTURE_MODE`` env
var > ``CaptureMode.V0``. An unrecognized env value logs a warning and falls
back to ``V0`` so a typo never silently flips the wire protocol.
var > ``CaptureMode.V1``. An unrecognized env value logs a warning and falls
back to ``V1`` so a typo never silently flips the wire protocol.
"""
if capture_mode is not None:
return _coerce_explicit(capture_mode)

raw = os.environ.get(CAPTURE_MODE_ENV_VAR)
if raw is None or raw.strip() == "":
return CaptureMode.V0
return CaptureMode.V1
Comment thread
eli-r-ph marked this conversation as resolved.

resolved = _ALIASES.get(raw.strip().lower())
if resolved is None:
log.warning(
"Unrecognized %s=%r; falling back to %s. Expected one of %s.",
CAPTURE_MODE_ENV_VAR,
raw,
CaptureMode.V0.value,
CaptureMode.V1.value,
sorted(_ALIASES),
)
return CaptureMode.V0
return CaptureMode.V1
return resolved
19 changes: 12 additions & 7 deletions posthog/capture_v1.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,10 @@
"""Serialization and transport for the Capture V1 wire protocol.

This module owns everything specific to ``POST /i/v1/analytics/events``: the
*transform* layer (legacy-shaped queued message -> v1 wire event + batch
envelope) and the *transport* layer (a single HTTP attempt, response parsing,
and the partial-retry send loop).
This module owns everything specific to the capture v1 endpoints
(``POST /i/v1/analytics/events`` and ``POST /i/v1/ai/events``, which share one
wire contract): the *transform* layer (legacy-shaped queued message -> v1 wire
event + batch envelope) and the *transport* layer (a single HTTP attempt,
response parsing, and the partial-retry send loop).

The v1 contract (see ``rust/capture/src/v1/analytics/types.rs``) differs from
the legacy ``/batch/`` shape in a few load-bearing ways that this module
Expand Down Expand Up @@ -66,6 +67,7 @@
__all__ = ["CaptureV1Error"]

_CAPTURE_V1_PATH = "/i/v1/analytics/events"
_CAPTURE_AI_V1_PATH = "/i/v1/ai/events"

# Required request/response headers for the v1 endpoint. Defined here as the
# single source of truth; the transport layer builds requests from them.
Expand Down Expand Up @@ -350,8 +352,9 @@ def _post_v1(
timeout: int = 15,
sdk_info: str = USER_AGENT,
session: Optional["requests.Session"] = None,
path: str = _CAPTURE_V1_PATH,
) -> "requests.Response":
"""Perform a single ``POST /i/v1/analytics/events`` attempt.
"""Perform a single capture v1 ``POST`` to ``path``.

Bearer-authed (no ``api_key`` in the body) with the required v1 headers.
``attempt`` (1-based) and the stable ``request_id`` are echoed via
Expand All @@ -361,7 +364,7 @@ def _post_v1(
the caller.
"""
trimmed_host = remove_trailing_slash(normalize_host(host))
url = trimmed_host + _CAPTURE_V1_PATH
url = trimmed_host + path
data = json.dumps(batch_body, cls=DatetimeSerializer)
headers = {
"Content-Type": "application/json",
Expand Down Expand Up @@ -473,8 +476,9 @@ def _send_v1_batch(
historical_migration: bool = False,
sdk_info: str = USER_AGENT,
session: Optional["requests.Session"] = None,
path: str = _CAPTURE_V1_PATH,
) -> None:
"""Deliver ``batch`` to the v1 endpoint with partial retry.
"""Deliver ``batch`` to the v1 endpoint at ``path`` with partial retry.

The v1 sibling of ``Consumer._send``: it loops up to ``max_retries + 1``
attempts, but unlike v0 it shrinks the batch to only the events the server
Expand Down Expand Up @@ -525,6 +529,7 @@ def _send_v1_batch(
timeout=timeout,
sdk_info=sdk_info,
session=session,
path=path,
)
except Exception as e:
# Transport-level failure (connection/timeout): retry like v0 does.
Expand Down
44 changes: 25 additions & 19 deletions posthog/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,11 @@
_resolve_capture_compression,
)
from posthog.capture_mode import CaptureMode, _resolve_capture_mode
from posthog.capture_v1 import _send_v1_batch
from posthog.capture_v1 import (
_CAPTURE_AI_V1_PATH,
_CAPTURE_V1_PATH,
_send_v1_batch,
)
from posthog.consumer import AI_MAX_MSG_SIZE, MAX_MSG_SIZE, Consumer, _DrainSignal
from posthog.contexts import (
_get_current_context,
Expand Down Expand Up @@ -85,7 +89,6 @@
from posthog.poller import Poller
from posthog.release_id import _resolve_release_id
from posthog.request import (
AI_EVENTS_ENDPOINT,
EVENTS_ENDPOINT,
USER_AGENT as _USER_AGENT,
APIError,
Expand Down Expand Up @@ -841,10 +844,10 @@ def __init__(
exception_autocapture_refill_interval_seconds: Seconds between
token refills for autocaptured exception rate limiting.
capture_mode: Capture wire protocol to use. Defaults to
``CaptureMode.V0`` (legacy ``/batch/``). Set ``CaptureMode.V1``
(or pass the string ``"v1"``) to opt into
``/i/v1/analytics/events``. When omitted, the
``POSTHOG_CAPTURE_MODE`` env var is consulted, then ``V0``.
``CaptureMode.V1`` (``/i/v1/analytics/events``). Set
``CaptureMode.V0`` (or pass the string ``"v0"``) to opt back
into the legacy ``/batch/`` endpoint. When omitted, the
``POSTHOG_CAPTURE_MODE`` env var is consulted, then ``V1``.
capture_compression: Request-body compression for capture-v1 uploads
(ignored in V0, which uses ``gzip``). ``CaptureCompression.GZIP``
or ``DEFLATE`` (or the strings ``"gzip"``/``"deflate"``). When
Expand Down Expand Up @@ -1085,24 +1088,27 @@ def __init__(
self._analytics_lane = _Lane(
name="analytics",
**lane_defaults,
endpoint=EVENTS_ENDPOINT,
endpoint=(
_CAPTURE_V1_PATH
if self.capture_mode == CaptureMode.V1
else EVENTS_ENDPOINT
),
max_msg_size=MAX_MSG_SIZE,
capture_mode=self.capture_mode,
capture_compression=self.capture_compression,
eager_start=not sync_mode,
)
# The AI lane is pinned to the v0 submitter: the AI endpoint has no v1
# form, and this keeps multi-MB AI events away from capture v1's
# smaller caps. The `capture_compression` pin is inert on v0 — its wire
# compression is the `gzip` flag, inherited from client config. Lazy
# start, so the many clients that never emit AI events pay for no
# extra threads.
# The AI lane always posts capture v1 to the AI endpoint, whatever the
# analytics `capture_mode`, 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.
self._ai_lane = _Lane(
name="ai",
**lane_defaults,
endpoint=AI_EVENTS_ENDPOINT,
endpoint=_CAPTURE_AI_V1_PATH,
max_msg_size=AI_MAX_MSG_SIZE,
capture_mode=CaptureMode.V0,
capture_mode=CaptureMode.V1,
capture_compression=CaptureCompression.NONE,
eager_start=False,
)
Expand Down Expand Up @@ -2438,19 +2444,19 @@ def _enqueue(self, msg, disable_geoip, lane=None, property_allowlist=None):
self.log.debug("enqueued with blocking %s.", msg["event"])

def send_sync() -> None:
# Sync mode bypasses the lane's queue but keeps its wire config:
# the AI lane is pinned to v0, so its events post to the AI
# endpoint regardless of `capture_mode`.
# Sync mode bypasses the lane's queue but keeps its wire config,
# so AI events post to the AI endpoint whatever `capture_mode`.
if lane.capture_mode == CaptureMode.V1:
_send_v1_batch(
self.api_key,
self.host,
[msg],
compression=self.capture_compression,
compression=lane.capture_compression,
timeout=self.timeout,
max_retries=self.max_retries,
historical_migration=self.historical_migration,
sdk_info=self._sdk_info,
path=lane.endpoint,
)
return

Expand Down
19 changes: 13 additions & 6 deletions posthog/consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
from posthog._logging import _configure_posthog_logging
from posthog.capture_compression import CaptureCompression
from posthog.capture_mode import CaptureMode
from posthog.capture_v1 import _backoff, _send_v1_batch
from posthog.capture_v1 import _CAPTURE_V1_PATH, _backoff, _send_v1_batch
from posthog.request import (
EVENTS_ENDPOINT,
USER_AGENT as _USER_AGENT,
Expand Down Expand Up @@ -113,9 +113,9 @@ def __init__(
retries=10,
timeout=15,
historical_migration=False,
endpoint=EVENTS_ENDPOINT,
endpoint=None,
max_msg_size=MAX_MSG_SIZE,
capture_mode=CaptureMode.V0,
capture_mode=CaptureMode.V1,
capture_compression=CaptureCompression.NONE,
):
"""Create a consumer thread."""
Expand All @@ -129,6 +129,12 @@ def __init__(
self.on_error = on_error
self.queue = queue
self.gzip = gzip
# Without an explicit endpoint, post to the analytics path of the
# selected protocol.
if endpoint is None:
endpoint = (
_CAPTURE_V1_PATH if capture_mode == CaptureMode.V1 else EVENTS_ENDPOINT
)
self.endpoint = endpoint
self.max_msg_size = max_msg_size
self.capture_mode = capture_mode
Expand Down Expand Up @@ -288,10 +294,10 @@ def next(self):
return items

def request(self, batch):
"""Upload the batch via the wire protocol selected by `capture_mode`.
"""Upload the batch to this consumer's `endpoint` via the wire protocol
selected by `capture_mode`.

V1 uses the partial-retry submitter (which posts to its own path); V0
posts the batch to this consumer's `endpoint`.
V1 uses the partial-retry submitter; V0 posts the whole batch.
"""
if self.capture_mode == CaptureMode.V1:
_send_v1_batch(
Expand All @@ -303,6 +309,7 @@ def request(self, batch):
max_retries=self.retries,
historical_migration=self.historical_migration,
sdk_info=self._sdk_info,
path=self.endpoint,
)
return
self._send(batch, self.endpoint)
Expand Down
58 changes: 58 additions & 0 deletions posthog/test/capture_helpers.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
"""Intercept capture uploads at the batch submitter for client-level tests.

Patching the submitter (not the HTTP layer) lets tests assert on the event
dicts the SDK built, before the wire encoding in ``capture_v1``. Wire shape is
covered by ``test_capture_v1``.
"""

import json
from unittest import mock

from requests import Response

_SUBMITTER = "_send_v1_batch"


def offline_v1_post(url: str, data=None, **kwargs) -> Response:
"""Stand-in for ``requests.Session.post`` that accepts every v1 event.

For subprocess tests with no server. Prints the uncompressed request body,
because the SDK never logs payloads, and answers ``ok`` for each event.
"""
print(f"capture request body: {data}", flush=True) # noqa: T201
events = json.loads(data)["batch"]
response = Response()
response.status_code = 200
response._content = json.dumps(
{"results": {event["uuid"]: {"result": "ok"} for event in events}}
).encode()
return response


def patch_capture_send(site: str = "client", **kwargs) -> "mock._patch":
"""Patch the submitter where ``posthog.<site>`` imported it.

``site="client"`` sees ``sync_mode`` uploads; ``site="consumer"`` sees
background consumer uploads.
"""
return mock.patch(f"posthog.{site}.{_SUBMITTER}", **kwargs)


def patch_async_capture_send(**kwargs) -> "mock._patch":
"""Patch the submitter the ``AsyncPosthog`` consumer awaits."""
return mock.patch("posthog._async_consumer.async_send_v1_batch", **kwargs)


def sent_batch(send_mock: mock.Mock, call_index: int = -1) -> list[dict]:
"""Return the event batch from one recorded upload (default: the last)."""
call = send_mock.call_args_list[call_index]
return call.args[2] if len(call.args) > 2 else call.kwargs["batch"]


def sent_events(send_mock: mock.Mock) -> list[dict]:
"""Return every event uploaded through ``send_mock``, in send order."""
return [
event
for index in range(len(send_mock.call_args_list))
for event in sent_batch(send_mock, index)
]
Loading
Loading