Skip to content
Draft
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
14 changes: 2 additions & 12 deletions .github/workflows/sdk-compliance.yml
Original file line number Diff line number Diff line change
Expand Up @@ -28,21 +28,11 @@ jobs:
run: python -m pytest sdk_compliance_adapter/test_adapter.py --timeout=30

compliance:
name: PostHog SDK compliance tests (capture v0)
name: PostHog SDK compliance tests
uses: PostHog/posthog-sdk-test-harness/.github/workflows/test-sdk-action.yml@4593de8b423f61fa222115da592e5c18dc82ad3c # 1.11.0
with:
adapter-dockerfile: "sdk_compliance_adapter/Dockerfile"
adapter-context: "."
test-harness-version: "1.1.1"
continue-on-error: false
report-name: "sdk-compliance-report-v0"

compliance-v1:
name: PostHog SDK compliance tests (capture v1)
uses: PostHog/posthog-sdk-test-harness/.github/workflows/test-sdk-action.yml@4593de8b423f61fa222115da592e5c18dc82ad3c # 1.11.0
with:
adapter-dockerfile: "sdk_compliance_adapter/Dockerfile.v1"
adapter-context: "."
test-harness-version: "1.1.1"
continue-on-error: false
report-name: "sdk-compliance-report-v1"
report-name: "sdk-compliance-report"
5 changes: 5 additions & 0 deletions .sampo/changesets/capture-v1-major.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
pypi/posthog: major
---

Capture v1 is the only capture path. Events and AI events send to the capture v1 endpoints, and the legacy v0 capture path is removed. See the migration guide for breaking changes.
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`.
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
2 changes: 1 addition & 1 deletion CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ uv sync --extra dev --extra test

## CI-aligned checks

Run the smallest relevant tests first, for example `pytest posthog/test/test_capture_v1.py --timeout=30` for v1 transport changes. Then run these core CI-aligned checks from the repository root in the activated `.venv` populated by the setup commands above:
Run the smallest relevant tests first, for example `pytest posthog/test/test_capture_send.py --timeout=30` for v1 transport changes. Then run these core CI-aligned checks from the repository root in the activated `.venv` populated by the setup commands above:

```bash
ruff format --check .
Expand Down
13 changes: 6 additions & 7 deletions posthog/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@
OptionalSetArgs,
)
from posthog.capture_compression import CaptureCompression as CaptureCompression
from posthog.capture_mode import CaptureMode as CaptureMode
from posthog.capture_send import CaptureError as CaptureError
from posthog.client import Client
from posthog.tracing.span import Span as Span
from posthog.async_client import AsyncClient as AsyncClient
Expand Down Expand Up @@ -313,8 +313,11 @@ def get_tags() -> Dict[str, Any]:
disabled no-op global client.
host: PostHog ingestion host. Defaults to the US ingestion endpoint when not
set.
on_error: Optional callback invoked by background consumers when event upload
fails. Keep it short and non-blocking. Lifecycle methods can be called
on_error: Optional callback ``(error, batch)`` invoked when event upload
fails: by background consumers, or on the calling thread in ``sync_mode``.
Capture failures arrive as ``CaptureError``. Without it, each failed batch
logs one aggregate line. A capture that fails inside the callback logs that
line instead of calling it again. Keep it short and non-blocking. Lifecycle methods can be called
directly and will be deferred, but the callback must not wait for another
thread or task that calls ``flush()``, ``join()``, or ``shutdown()``.
debug: Enable verbose SDK logging and re-raise errors from public APIs.
Expand Down Expand Up @@ -424,9 +427,6 @@ def get_tags() -> Dict[str, Any]:
# We recommend setting this to False if you are only using the personalApiKey for evaluating remote config payloads via `get_remote_config_payload` and not using local evaluation.
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.
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.
enable_full_ai_capture = False # type: bool
Expand Down Expand Up @@ -1423,7 +1423,6 @@ def setup() -> Client:
exception_autocapture_bucket_size=exception_autocapture_bucket_size,
exception_autocapture_refill_rate=exception_autocapture_refill_rate,
exception_autocapture_refill_interval_seconds=exception_autocapture_refill_interval_seconds,
capture_mode=capture_mode,
)

# Always set in case user changes it. Preserve Client's auto-disabled state
Expand Down
108 changes: 41 additions & 67 deletions posthog/_async_consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,11 +9,11 @@
from dataclasses import dataclass
from typing import Any, Optional

from ._async_request import async_batch_post, async_send_v1_batch
from ._async_request import async_send_v1_batch
from .capture_compression import CaptureCompression
from .capture_mode import CaptureMode
from .capture_send import _CAPTURE_V1_PATH, _capture_loss_message
from .consumer import BATCH_SIZE_LIMIT, MAX_MSG_SIZE
from .request import APIError, DatetimeSerializer, EVENTS_ENDPOINT
from .request import DatetimeSerializer

_STOP = object()
_PROCESSING_EVENT = contextvars.ContextVar(
Expand Down Expand Up @@ -48,6 +48,33 @@ async def _invoke_callback(callback, *args):
return result


# True while an SDK-invoked `on_error` runs, so a capture that fails inside the
# callback logs instead of re-entering it.
_IN_ON_ERROR: contextvars.ContextVar[bool] = contextvars.ContextVar(
"posthog_in_on_error", default=False
)


async def _report_capture_failure(
on_error: Optional[Callable[..., Any]],
log: logging.Logger,
error: Exception,
batch: list[dict[str, Any]],
endpoint: str,
) -> None:
"""Hand a failed send to `on_error`, or log one aggregate line without one."""
if on_error is None or _IN_ON_ERROR.get():
log.error(_capture_loss_message(error, max(1, len(batch)), endpoint))
return
token = _IN_ON_ERROR.set(True)
try:
await _invoke_callback(on_error, error, batch)
except Exception as callback_error:
log.error("on_error handler failed (%s)", type(callback_error).__name__)
finally:
_IN_ON_ERROR.reset(token)


async def _serialized_event_size(event: dict[str, Any]) -> int:
serialized = await asyncio.to_thread(json.dumps, event, cls=DatetimeSerializer)
return len(serialized.encode())
Expand All @@ -68,13 +95,10 @@ def __init__(
process_event: Callable[[dict[str, Any]], Awaitable[Optional[dict[str, Any]]]],
flush_at: int,
flush_interval: float,
gzip: bool,
retries: int,
timeout: int,
historical_migration: bool,
capture_mode: CaptureMode,
capture_compression: CaptureCompression,
http_client: Optional[Any],
) -> None:
self.queue = queue
self.api_key = api_key
Expand All @@ -83,13 +107,10 @@ def __init__(
self.process_event = process_event
self.flush_at = flush_at
self.flush_interval = flush_interval
self.gzip = gzip
self.retries = max(0, retries)
self.timeout = timeout
self.historical_migration = historical_migration
self.capture_mode = capture_mode
self.capture_compression = capture_compression
self.http_client = http_client
self._carryover: Optional[tuple[dict[str, Any], int]] = None
self._flush_event = asyncio.Event()

Expand Down Expand Up @@ -149,18 +170,9 @@ async def upload(self, batch: list[dict[str, Any]]) -> None:
try:
await self.request(batch)
except Exception as error:
self.log.error(
"async capture upload failed (%s, status=%s)",
type(error).__name__,
getattr(error, "status", None),
await _report_capture_failure(
self.on_error, self.log, error, batch, _CAPTURE_V1_PATH
)
if self.on_error:
try:
await _invoke_callback(self.on_error, error, batch)
except Exception as callback_error:
self.log.error(
"on_error handler failed (%s)", type(callback_error).__name__
)
finally:
for _ in batch:
self.queue.task_done()
Expand Down Expand Up @@ -238,50 +250,12 @@ async def next(self) -> tuple[list[dict[str, Any]], bool]:
return items, stop

async def request(self, batch: list[dict[str, Any]]) -> None:
if self.capture_mode == CaptureMode.V1:
await async_send_v1_batch(
self.api_key,
self.host,
batch,
compression=self.capture_compression,
timeout=self.timeout,
max_retries=self.retries,
historical_migration=self.historical_migration,
)
return

last_error: Optional[Exception] = None
for attempt in range(self.retries + 1):
try:
await async_batch_post(
self.api_key,
self.host,
batch=batch,
path=EVENTS_ENDPOINT,
gzip=self.gzip,
timeout=self.timeout,
historical_migration=self.historical_migration,
client=self.http_client,
)
return
except Exception as error:
last_error = error
if not self._is_retryable(error) or attempt >= self.retries:
raise
retry_after = getattr(error, "retry_after", None)
delay = max(
min(2**attempt, 30),
min(retry_after, 30) if retry_after and retry_after > 0 else 0,
)
await asyncio.sleep(delay)

if last_error is not None: # pragma: no cover - loop always raises first
raise last_error

@staticmethod
def _is_retryable(error: Exception) -> bool:
if not isinstance(error, APIError):
return True
if not isinstance(error.status, int):
return False
return not (400 <= error.status < 500 and error.status not in (408, 429))
await async_send_v1_batch(
self.api_key,
self.host,
batch,
compression=self.capture_compression,
timeout=self.timeout,
max_retries=self.retries,
historical_migration=self.historical_migration,
)
114 changes: 2 additions & 112 deletions posthog/_async_request.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,16 +2,12 @@

import asyncio
import json
import logging
import zlib
from datetime import datetime, timezone
from gzip import GzipFile
from io import BytesIO
from typing import Any, Optional
from urllib.parse import quote, urljoin, urlsplit
from urllib.parse import quote

from .capture_compression import CaptureCompression
from .capture_v1 import _parse_retry_after, _send_v1_batch
from .capture_send import _parse_retry_after, _send_v1_batch
from .request import (
APIError,
DatetimeSerializer,
Expand Down Expand Up @@ -41,54 +37,6 @@ def _build_client(host: Optional[str] = None):
return httpx_module.AsyncClient(base_url=base_url, follow_redirects=False)


def _serialize_v0_body(
api_key: str, gzip_enabled: bool, body: dict[str, Any]
) -> tuple[str | bytes, dict[str, str]]:
payload = {
**body,
"sent_at": datetime.now(tz=timezone.utc).isoformat(),
"api_key": api_key,
}
serialized = json.dumps(payload, cls=DatetimeSerializer)
data: str | bytes = serialized
headers = {"Content-Type": "application/json", "User-Agent": USER_AGENT}

if gzip_enabled:
try:
buf = BytesIO()
with GzipFile(fileobj=buf, mode="w") as gz:
gz.write(serialized.encode("utf-8"))
data = buf.getvalue()
headers["Content-Encoding"] = "gzip"
except (OSError, zlib.error) as exc:
logging.getLogger("posthog").warning(
"failed to gzip async request body, sending uncompressed: %s", exc
)

return data, headers


def _origin(url: str) -> tuple[str, str, Optional[int]]:
parsed = urlsplit(url)
port = parsed.port
if port is None:
port = 443 if parsed.scheme.lower() == "https" else 80
return parsed.scheme.lower(), (parsed.hostname or "").lower(), port


def _same_origin_redirect_url(
base_url: str, current_url: str, location: str
) -> Optional[str]:
target = urlsplit(urljoin(current_url, location))
if _origin(target.geturl()) != _origin(base_url):
return None
return (
urlsplit(base_url)
._replace(path=target.path or "/", query=target.query, fragment="")
.geturl()
)


def _serialize_flags_body(
project_api_key: str, body: dict[str, Any]
) -> tuple[str, dict[str, str]]:
Expand Down Expand Up @@ -196,64 +144,6 @@ async def async_remote_config(
await http_client.aclose()


async def async_batch_post(
api_key: str,
host: Optional[str],
*,
batch: list[dict[str, Any]],
path: str,
gzip: bool = False,
timeout: int = 15,
historical_migration: bool = False,
client: Optional[Any] = None,
) -> None:
"""Post one legacy capture batch without blocking the event loop."""
if not path.startswith("/") or "://" in path:
raise ValueError("async capture paths must be relative")

data, headers = await asyncio.to_thread(
_serialize_v0_body,
api_key,
gzip,
{
"batch": batch,
"historical_migration": historical_migration,
},
)

owns_client = client is None
http_client = client or _build_client(host)
try:
logging.getLogger("posthog").debug("making async capture request")
base_url = remove_trailing_slash(normalize_host(host))
# Absolute URLs avoid reapplying an HTTPX base_url path on redirects.
request_url = f"{base_url}{path}"
for redirect_count in range(6):
response = await http_client.post(
request_url, content=data, headers=headers, timeout=timeout
)
if response.status_code not in (307, 308):
_process_response(response)
return

location = response.headers.get("Location") or response.headers.get(
"location"
)
redirect_url = (
_same_origin_redirect_url(base_url, request_url, location)
if location
else None
)
if redirect_url is None:
raise APIError(400, "Cross-origin or invalid redirect blocked")
if redirect_count >= 5:
raise APIError(400, "Too many capture redirects")
request_url = redirect_url
finally:
if owns_client:
await http_client.aclose()


async def async_send_v1_batch(
api_key: str,
host: Optional[str],
Expand Down
2 changes: 1 addition & 1 deletion posthog/ai/prompts.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@
from dataclasses import dataclass
from typing import Any, Dict, List, Literal, Optional, Union, overload

from posthog.capture_v1 import _parse_retry_after
from posthog.capture_send import _parse_retry_after
from posthog.request import USER_AGENT, _get_session
from posthog.utils import remove_trailing_slash

Expand Down
Loading