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
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
Comment thread
eli-r-ph marked this conversation as resolved.
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"
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 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`.
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`.

## Mirror and build safety

Expand Down
5 changes: 0 additions & 5 deletions posthog/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,6 @@
OptionalSetArgs,
)
from posthog.capture_compression import CaptureCompression as CaptureCompression
from posthog.capture_mode import CaptureMode as CaptureMode
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 @@ -424,9 +423,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.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.
enable_full_ai_capture = False # type: bool
Expand Down Expand Up @@ -1423,7 +1419,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
67 changes: 11 additions & 56 deletions posthog/_async_consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,11 +9,10 @@
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 .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 @@ -68,13 +67,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 +79,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 @@ -238,50 +231,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(
Comment thread
eli-r-ph marked this conversation as resolved.
self.api_key,
self.host,
batch,
compression=self.capture_compression,
timeout=self.timeout,
max_retries=self.retries,
historical_migration=self.historical_migration,
)
112 changes: 1 addition & 111 deletions posthog/_async_request.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,13 +2,9 @@

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
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
Loading
Loading