diff --git a/.github/workflows/sdk-compliance.yml b/.github/workflows/sdk-compliance.yml index 3891e105c..55cfc8ebe 100644 --- a/.github/workflows/sdk-compliance.yml +++ b/.github/workflows/sdk-compliance.yml @@ -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" diff --git a/AGENTS.md b/AGENTS.md index e05e0a255..6b7a8d74f 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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 diff --git a/posthog/__init__.py b/posthog/__init__.py index ccbc8db58..9e28ec5bf 100644 --- a/posthog/__init__.py +++ b/posthog/__init__.py @@ -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 @@ -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 @@ -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 diff --git a/posthog/_async_consumer.py b/posthog/_async_consumer.py index 04eeb082b..9352139ee 100644 --- a/posthog/_async_consumer.py +++ b/posthog/_async_consumer.py @@ -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( @@ -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 @@ -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() @@ -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( + self.api_key, + self.host, + batch, + compression=self.capture_compression, + timeout=self.timeout, + max_retries=self.retries, + historical_migration=self.historical_migration, + ) diff --git a/posthog/_async_request.py b/posthog/_async_request.py index ae750de9a..15b218f74 100644 --- a/posthog/_async_request.py +++ b/posthog/_async_request.py @@ -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 @@ -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]]: @@ -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], diff --git a/posthog/async_client.py b/posthog/async_client.py index 3f429e81a..e46f0e494 100644 --- a/posthog/async_client.py +++ b/posthog/async_client.py @@ -26,7 +26,6 @@ ) from ._async_request import ( _build_client, - _require_httpx, async_flags as _async_flags, async_remote_config as _async_remote_config, ) @@ -35,7 +34,6 @@ CaptureCompression, _resolve_capture_compression, ) -from .capture_mode import CaptureMode, _resolve_capture_mode from .client import ( MAX_DICT_SIZE as _MAX_DICT_SIZE, _MINIMAL_FLAG_CALLED_EVENT_PROPERTIES, @@ -98,13 +96,13 @@ def __init__( self, project_api_key: str, host: Optional[str] = None, + *, debug: bool = False, max_queue_size: int = 10000, send: bool = True, on_error=None, flush_at: int = 100, flush_interval: float = 5.0, - gzip: bool = False, max_retries: int = 3, timeout: int = 15, thread: int = 1, @@ -122,7 +120,6 @@ def __init__( code_variables_mask_url_credentials=None, code_variables_detect_secrets=None, in_app_modules: Optional[list[str]] = None, - capture_mode: Optional[Union[CaptureMode, str]] = None, capture_compression: Optional[Union[CaptureCompression, str]] = None, capture_trace_context: bool = False, secret_key: Optional[str] = None, @@ -141,7 +138,6 @@ def __init__( self.debug = debug self.send = send self.on_error = on_error - self.gzip = gzip self.max_retries = max(0, max_retries) self.timeout = timeout self.disabled = disabled or not self.api_key @@ -150,10 +146,7 @@ def __init__( self.historical_migration = historical_migration self.super_properties = super_properties self._release_id = _resolve_release_id() - self.capture_mode = _resolve_capture_mode(capture_mode) - self.capture_compression = _resolve_capture_compression( - capture_compression, gzip_fallback=gzip - ) + self.capture_compression = _resolve_capture_compression(capture_compression) self.capture_trace_context = capture_trace_context if personal_api_key is not None and secret_key is None: warnings.warn( @@ -310,9 +303,6 @@ def _get_http_client(self): return self._http_client def _new_consumer(self) -> _AsyncConsumer: - http_client = ( - self._get_http_client() if self.capture_mode == CaptureMode.V0 else None - ) return _AsyncConsumer( self._queue, self.api_key, @@ -321,19 +311,12 @@ def _new_consumer(self) -> _AsyncConsumer: process_event=self._process_event, flush_at=self._flush_at, flush_interval=self._flush_interval, - gzip=self.gzip, retries=self.max_retries, timeout=self.timeout, historical_migration=self.historical_migration, - capture_mode=self.capture_mode, capture_compression=self.capture_compression, - http_client=http_client, ) - def _validate_transport_available(self) -> None: - if self.capture_mode == CaptureMode.V0: - _require_httpx() - def _ensure_workers_started(self) -> None: if self.disabled or not self.send or self._closed or self._worker_tasks: return @@ -523,7 +506,6 @@ def capture( if not self.send: return sent_uuid - self._validate_transport_available() if not self._enqueue_prepared_event(prepared): return None self.log.debug("queued async event %s", event) @@ -745,7 +727,6 @@ def _enqueue_built_event( return None if not self.send: return sent_uuid - self._validate_transport_available() if not self._enqueue_prepared_event(prepared): return None return sent_uuid diff --git a/posthog/capture_compression.py b/posthog/capture_compression.py index 659c12fab..52a1d4c2d 100644 --- a/posthog/capture_compression.py +++ b/posthog/capture_compression.py @@ -19,10 +19,9 @@ class CaptureCompression(str, Enum): - """Selects the request-body compression for capture-v1 uploads. + """Selects the request-body compression for capture uploads. - Only honored when ``capture_mode`` is ``V1``; the legacy ``/batch/`` path - keeps using its own ``gzip`` flag. ``NONE`` sends the body uncompressed. + ``NONE`` sends the body uncompressed. ``GZIP`` and ``DEFLATE`` (zlib, RFC 1950) are both stdlib / zero-dependency; ``ZSTD`` is faster and compresses better but needs the optional zstandard package (``pip install posthog[zstd]``) until stdlib support lands in @@ -76,15 +75,12 @@ def _coerce_explicit( def _resolve_capture_compression( capture_compression: Optional[Union[CaptureCompression, str]] = None, - *, - gzip_fallback: bool = False, ) -> CaptureCompression: - """Resolve the effective v1 compression. + """Resolve the effective capture compression. Precedence: explicit ``capture_compression`` argument > - ``POSTHOG_CAPTURE_COMPRESSION`` env var > the legacy ``gzip`` flag - (``GZIP`` when set) > ``NONE``. An unrecognized env value logs a warning and - falls back to the ``gzip`` flag, so a typo never silently changes encoding. + ``POSTHOG_CAPTURE_COMPRESSION`` env var > ``NONE``. An unrecognized env + value logs a warning and falls back to ``NONE``. ``ZSTD`` requires the optional zstandard package: explicitly requesting it without the package raises ``ValueError`` (programming error, fail loud), @@ -100,7 +96,7 @@ def _resolve_capture_compression( ) return resolved - fallback = CaptureCompression.GZIP if gzip_fallback else CaptureCompression.NONE + fallback = CaptureCompression.NONE raw = os.environ.get(CAPTURE_COMPRESSION_ENV_VAR) if raw is None or raw.strip() == "": diff --git a/posthog/capture_mode.py b/posthog/capture_mode.py deleted file mode 100644 index 5f72e3702..000000000 --- a/posthog/capture_mode.py +++ /dev/null @@ -1,83 +0,0 @@ -import logging -import os -from enum import Enum -from typing import Optional, Union - -__all__ = ["CAPTURE_MODE_ENV_VAR", "CaptureMode"] - -log = logging.getLogger("posthog") - -CAPTURE_MODE_ENV_VAR = "POSTHOG_CAPTURE_MODE" - - -class CaptureMode(str, Enum): - """Selects the capture wire protocol used for event ingestion. - - ``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. - """ - - V0 = "v0" - V1 = "v1" - - -# Accepted spellings for both the explicit kwarg and the env var. Aliases mirror -# the posthog-go naming (``legacy`` / ``analytics_v1``) so the two SDKs are -# configured with the same vocabulary. -_ALIASES: dict[str, CaptureMode] = { - "v0": CaptureMode.V0, - "legacy": CaptureMode.V0, - "v1": CaptureMode.V1, - "analytics_v1": CaptureMode.V1, -} - - -def _coerce_explicit(value: Union[CaptureMode, str]) -> CaptureMode: - """Normalize an explicitly-supplied capture mode to a ``CaptureMode``. - - Accepts a ``CaptureMode`` or one of the string aliases. An explicit but - unrecognized value is a programming error, so it raises ``ValueError`` rather - than silently defaulting (unlike the env var, which is operator-supplied and - defaults defensively). - """ - if isinstance(value, CaptureMode): - return value - if isinstance(value, str): - resolved = _ALIASES.get(value.strip().lower()) - if resolved is not None: - return resolved - raise ValueError( - f"invalid capture_mode {value!r}; expected a CaptureMode or one of " - f"{sorted(_ALIASES)}" - ) - - -def _resolve_capture_mode( - capture_mode: Optional[Union[CaptureMode, str]] = None, -) -> CaptureMode: - """Resolve the effective capture mode. - - Precedence: explicit ``capture_mode`` argument > ``POSTHOG_CAPTURE_MODE`` env - 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.V1 - - 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.V1.value, - sorted(_ALIASES), - ) - return CaptureMode.V1 - return resolved diff --git a/posthog/client.py b/posthog/client.py index d2d705d52..083a6eb03 100644 --- a/posthog/client.py +++ b/posthog/client.py @@ -31,7 +31,6 @@ CaptureCompression, _resolve_capture_compression, ) -from posthog.capture_mode import CaptureMode, _resolve_capture_mode from posthog.capture_v1 import ( _CAPTURE_AI_V1_PATH, _CAPTURE_V1_PATH, @@ -89,7 +88,6 @@ from posthog.poller import Poller from posthog.release_id import _resolve_release_id from posthog.request import ( - EVENTS_ENDPOINT, USER_AGENT as _USER_AGENT, APIError, QuotaLimitError, @@ -97,7 +95,6 @@ RequestsTimeout, _get as _get_with_identity, _remote_config as _remote_config_with_identity, - batch_post, determine_server_host, flags, get, @@ -377,13 +374,11 @@ def __init__( send, flush_at, flush_interval, - gzip, max_retries, timeout, historical_migration, endpoint, max_msg_size, - capture_mode, capture_compression, sdk_info, eager_start, @@ -395,13 +390,11 @@ def __init__( self.send = send self.flush_at = flush_at self.flush_interval = flush_interval - self.gzip = gzip self.max_retries = max_retries self.timeout = timeout self.historical_migration = historical_migration self.endpoint = endpoint self.max_msg_size = max_msg_size - self.capture_mode = capture_mode self.capture_compression = capture_compression self.sdk_info = sdk_info self._max_queue_size = max_queue_size @@ -430,13 +423,11 @@ def _start_locked(self) -> None: on_error=self.on_error, flush_at=self.flush_at, flush_interval=self.flush_interval, - gzip=self.gzip, retries=self.max_retries, timeout=self.timeout, historical_migration=self.historical_migration, endpoint=self.endpoint, max_msg_size=self.max_msg_size, - capture_mode=self.capture_mode, capture_compression=self.capture_compression, ) consumer._sdk_info = self.sdk_info @@ -674,13 +665,13 @@ def __init__( self, project_api_key: str, host=None, + *, debug=False, max_queue_size=10000, send=True, on_error=None, flush_at=100, flush_interval=5.0, - gzip=False, max_retries=3, sync_mode=False, timeout=15, @@ -712,13 +703,10 @@ def __init__( exception_autocapture_bucket_size=ExceptionCapture.DEFAULT_BUCKET_SIZE, exception_autocapture_refill_rate=ExceptionCapture.DEFAULT_REFILL_RATE, exception_autocapture_refill_interval_seconds=ExceptionCapture.DEFAULT_REFILL_INTERVAL_SECONDS, - capture_mode: Optional[Union[CaptureMode, str]] = None, capture_compression: Optional[Union[CaptureCompression, str]] = None, secret_key=None, metrics: Optional[dict] = None, enable_full_ai_capture=False, - # Appended rather than grouped with the other `capture_*` options so - # existing positional arguments keep their slots. capture_trace_context=False, _use_ai_lane=False, _enable_multimodal_capture=False, @@ -744,7 +732,6 @@ def __init__( flush_at: Number of queued events that triggers a batch upload. flush_interval: Maximum seconds a background consumer waits before flushing a partial batch. - gzip: Whether to gzip event upload payloads. max_retries: Number of upload retries. Values below 0 are treated as 0. sync_mode: If True, send each event synchronously instead of using background worker threads. This blocks the calling thread; in @@ -843,16 +830,11 @@ def __init__( interval for each exception type's bucket. exception_autocapture_refill_interval_seconds: Seconds between token refills for autocaptured exception rate limiting. - capture_mode: Capture wire protocol to use. Defaults to - ``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 - omitted, the ``POSTHOG_CAPTURE_COMPRESSION`` env var is consulted, - then the legacy ``gzip`` flag, then no compression. + capture_compression: Request-body compression for ``capture()`` + uploads. ``CaptureCompression.GZIP`` or ``DEFLATE`` (or the + strings ``"gzip"``/``"deflate"``). When omitted, the + ``POSTHOG_CAPTURE_COMPRESSION`` env var is consulted, then no + compression. Examples: ```python @@ -893,7 +875,6 @@ def __init__( self.raw_host = normalize_host(host) self.host = determine_server_host(host) self._duplicate_client_registry_key: Optional[tuple[str, str]] = None - self.gzip = gzip self.timeout = timeout self.max_retries = max(0, max_retries) self._feature_flags: Optional[list[Any]] = ( @@ -953,18 +934,10 @@ def __init__( ) self.is_server = is_server self.historical_migration = historical_migration - # Selects the capture wire protocol (V0 legacy `/batch/` vs V1 - # `/i/v1/analytics/events`). Resolved here so the env-var fallback is - # applied once; V0 is the default and keeps upgrades transparent. - self.capture_mode = _resolve_capture_mode(capture_mode) self._library_id = "posthog-python" self._library_version = VERSION self._sdk_info = f"{self._library_id}/{self._library_version}" - # v1-only request compression; falls back to the legacy `gzip` flag when - # neither the kwarg nor POSTHOG_CAPTURE_COMPRESSION is set. - self.capture_compression = _resolve_capture_compression( - capture_compression, gzip_fallback=gzip - ) + self.capture_compression = _resolve_capture_compression(capture_compression) 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. @@ -1079,7 +1052,6 @@ def __init__( send=send, flush_at=flush_at, flush_interval=flush_interval, - gzip=gzip, max_retries=self.max_retries, timeout=timeout, historical_migration=historical_migration, @@ -1088,27 +1060,20 @@ def __init__( self._analytics_lane = _Lane( name="analytics", **lane_defaults, - endpoint=( - _CAPTURE_V1_PATH - if self.capture_mode == CaptureMode.V1 - else EVENTS_ENDPOINT - ), + endpoint=_CAPTURE_V1_PATH, max_msg_size=MAX_MSG_SIZE, - capture_mode=self.capture_mode, capture_compression=self.capture_compression, eager_start=not sync_mode, ) - # 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 + # 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. self._ai_lane = _Lane( name="ai", **lane_defaults, endpoint=_CAPTURE_AI_V1_PATH, max_msg_size=AI_MAX_MSG_SIZE, - capture_mode=CaptureMode.V1, capture_compression=CaptureCompression.NONE, eager_start=False, ) @@ -2445,30 +2410,17 @@ def _enqueue(self, msg, disable_geoip, lane=None, property_allowlist=None): def send_sync() -> None: # 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=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 - - batch_post( + # so AI events still post to the AI endpoint. + _send_v1_batch( self.api_key, self.host, - gzip=self.gzip, + [msg], + compression=lane.capture_compression, timeout=self.timeout, - batch=[msg], + max_retries=self.max_retries, historical_migration=self.historical_migration, + sdk_info=self._sdk_info, path=lane.endpoint, - **self._request_identity_kwargs(), ) if lane.run_sync_if_open(send_sync): diff --git a/posthog/consumer.py b/posthog/consumer.py index 8a69844cf..0b15431a8 100644 --- a/posthog/consumer.py +++ b/posthog/consumer.py @@ -6,14 +6,10 @@ from posthog._logging import _configure_posthog_logging from posthog.capture_compression import CaptureCompression -from posthog.capture_mode import CaptureMode -from posthog.capture_v1 import _CAPTURE_V1_PATH, _backoff, _send_v1_batch +from posthog.capture_v1 import _CAPTURE_V1_PATH, _send_v1_batch from posthog.request import ( - EVENTS_ENDPOINT, USER_AGENT as _USER_AGENT, - APIError, DatetimeSerializer, - batch_post, ) from queue import Empty @@ -109,13 +105,11 @@ def __init__( host=None, on_error=None, flush_interval=5.0, - gzip=False, retries=10, timeout=15, historical_migration=False, - endpoint=None, + endpoint=_CAPTURE_V1_PATH, max_msg_size=MAX_MSG_SIZE, - capture_mode=CaptureMode.V1, capture_compression=CaptureCompression.NONE, ): """Create a consumer thread.""" @@ -128,16 +122,8 @@ def __init__( self.host = host 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 self.capture_compression = capture_compression self._sdk_info = _USER_AGENT self._drain_signal: Optional[_DrainSignal] = None @@ -294,67 +280,16 @@ def next(self): return items def request(self, batch): - """Upload the batch to this consumer's `endpoint` via the wire protocol - selected by `capture_mode`. - - V1 uses the partial-retry submitter; V0 posts the whole batch. - """ - if self.capture_mode == CaptureMode.V1: - _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, - sdk_info=self._sdk_info, - path=self.endpoint, - ) - return - self._send(batch, self.endpoint) - - def _send(self, batch, path): - """Attempt to upload a single batch to `path`, retrying before raising an error""" - - def is_retryable(exc): - if isinstance(exc, APIError): - # retry on server errors and client errors - # with 408 (request timeout) or 429 (rate limited), - # don't retry on other client errors - if isinstance(exc.status, int): - return not ( - (400 <= exc.status < 500) and exc.status not in (408, 429) - ) - return False - else: - # retry on all other errors (eg. network) - return True - - last_exc = None - for attempt in range(self.retries + 1): - try: - batch_post( - self.api_key, - self.host, - gzip=self.gzip, - timeout=self.timeout, - batch=batch, - historical_migration=self.historical_migration, - path=path, - **( - {"_user_agent": self._sdk_info} - if self._sdk_info != _USER_AGENT - else {} - ), - ) - return - except Exception as e: - last_exc = e - if not is_retryable(e): - raise - if attempt < self.retries: - _backoff(attempt, getattr(e, "retry_after", None)) - - if last_exc: - raise last_exc + """Upload the batch to this consumer's `endpoint` with the capture v1 + partial-retry submitter.""" + _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, + sdk_info=self._sdk_info, + path=self.endpoint, + ) diff --git a/posthog/request.py b/posthog/request.py index 76df1fdc9..2f68fc38e 100644 --- a/posthog/request.py +++ b/posthog/request.py @@ -3,11 +3,8 @@ import re import socket import time -import zlib from dataclasses import dataclass from datetime import date, datetime, timezone -from gzip import GzipFile -from io import BytesIO from typing import Any, List, Optional, Tuple, Union, cast import requests @@ -216,7 +213,6 @@ def post( api_key: str, host: Optional[str] = None, path: Optional[str] = None, - gzip: bool = False, timeout: int = 15, session: Optional[requests.Session] = None, **kwargs, @@ -229,7 +225,7 @@ def post( trimmed_host = remove_trailing_slash(normalize_host(host)) url = trimmed_host + cast(str, path) body["api_key"] = api_key - data: str | bytes = json.dumps(body, cls=DatetimeSerializer) + data = json.dumps(body, cls=DatetimeSerializer) if log.isEnabledFor(logging.DEBUG): log.debug( "making request: %s to url: %s", @@ -237,17 +233,6 @@ def post( url, ) headers = {"Content-Type": "application/json", "User-Agent": user_agent} - if gzip: - try: - buf = BytesIO() - with GzipFile(fileobj=buf, mode="w") as gz: - # 'data' was produced by json.dumps(), - # whose default encoding is utf-8. - gz.write(cast(str, data).encode("utf-8")) - data = buf.getvalue() - headers["Content-Encoding"] = "gzip" - except (OSError, zlib.error) as exc: - log.warning("failed to gzip request body, sending uncompressed: %s", exc) res = (session or _get_session()).post( url, data=data, headers=headers, timeout=timeout @@ -314,7 +299,6 @@ def _feature_flags_retry_delay(failed_attempt: int) -> float: def flags( api_key: str, host: Optional[str] = None, - gzip: bool = False, timeout: int = 15, max_retries: int = 1, **kwargs, @@ -330,7 +314,6 @@ def flags( api_key, host, "/flags/?v=2", - gzip, timeout, session=_get_flags_session(), _user_agent=user_agent, @@ -382,25 +365,6 @@ def _remote_config( return response.data -EVENTS_ENDPOINT = "/batch/" -AI_EVENTS_ENDPOINT = "/i/v0/ai/batch/" - - -def batch_post( - api_key: str, - host: Optional[str] = None, - gzip: bool = False, - timeout: int = 15, - path: str = EVENTS_ENDPOINT, - **kwargs, -) -> requests.Response: - """Post the `kwargs` to the batch API endpoint for events""" - res = post(api_key, host, path, gzip, timeout, **kwargs) - return _process_response( - res, success_message="data uploaded successfully", return_json=False - ) - - def get( api_key: str, url: str, diff --git a/posthog/test/mcp/test_posthog_mcp.py b/posthog/test/mcp/test_posthog_mcp.py index 6ed662e81..8820a156a 100644 --- a/posthog/test/mcp/test_posthog_mcp.py +++ b/posthog/test/mcp/test_posthog_mcp.py @@ -9,7 +9,6 @@ import pytest from mcp.types import CallToolResult, ServerResult, TextContent, Tool -from posthog.capture_mode import CaptureMode from posthog.mcp import ( PostHogMCP, PreparedToolCall, @@ -134,20 +133,8 @@ def before_send(event): ) -def test_mcp_library_identity_reaches_capture_v0_header(): - response = mock.Mock(status_code=200) - client = PostHogMCP("phc_test", sync_mode=True) - - with mock.patch("posthog.request._session.post", return_value=response) as post: - client.capture("$mcp_custom") - - assert post.call_args.kwargs["headers"]["User-Agent"] == ( - f"posthog-python-mcp/{VERSION}" - ) - - def test_mcp_library_identity_reaches_capture_v1_header(): - client = PostHogMCP("phc_test", sync_mode=True, capture_mode=CaptureMode.V1) + client = PostHogMCP("phc_test", sync_mode=True) with mock.patch("posthog.client._send_v1_batch") as send: client.capture("$mcp_custom") diff --git a/posthog/test/snapshots/legacy_event_family.json b/posthog/test/snapshots/event_family.json similarity index 77% rename from posthog/test/snapshots/legacy_event_family.json rename to posthog/test/snapshots/event_family.json index 4e33c6bd8..0d72f4cf5 100644 --- a/posthog/test/snapshots/legacy_event_family.json +++ b/posthog/test/snapshots/event_family.json @@ -1,18 +1,16 @@ { "body": { - "api_key": "phc_snapshot_project", "batch": [ { "distinct_id": "user-123", "event": "order completed", + "options": {}, "properties": { "$geoip_disable": true, "$groups": { "company": "company-456" }, "$is_server": true, - "$lib": "posthog-python", - "$lib_version": "", "$os": "", "$os_distro": "", "$os_version": "", @@ -34,17 +32,16 @@ "uuid": "00000000-0000-4000-8000-000000000001" }, { - "$set": { - "email": "person@example.com", - "plan": "pro" - }, "distinct_id": "user-123", "event": "$set", + "options": {}, "properties": { "$geoip_disable": true, "$is_server": true, - "$lib": "posthog-python", - "$lib_version": "" + "$set": { + "email": "person@example.com", + "plan": "pro" + } }, "timestamp": "2026-01-02T03:04:05+00:00", "uuid": "00000000-0000-4000-8000-000000000002" @@ -52,11 +49,10 @@ { "distinct_id": "anonymous-789", "event": "$create_alias", + "options": {}, "properties": { "$geoip_disable": true, "$is_server": true, - "$lib": "posthog-python", - "$lib_version": "", "alias": "user-123", "distinct_id": "anonymous-789" }, @@ -66,6 +62,7 @@ { "distinct_id": "user-123", "event": "$groupidentify", + "options": {}, "properties": { "$geoip_disable": true, "$group_key": "company-456", @@ -74,21 +71,23 @@ "name": "Example Corp" }, "$group_type": "company", - "$is_server": true, - "$lib": "posthog-python", - "$lib_version": "" + "$is_server": true }, "timestamp": "2026-01-02T03:04:05+00:00", "uuid": "00000000-0000-4000-8000-000000000004" } ], - "historical_migration": false, - "sent_at": "2026-01-02T03:04:05+00:00" + "created_at": "2026-01-02T03:04:05+00:00" }, "headers": { + "Authorization": "Bearer phc_snapshot_project", "Content-Type": "application/json", + "PostHog-Attempt": "1", + "PostHog-Request-Id": "", + "PostHog-Request-Timestamp": "2026-01-02T03:04:05+00:00", + "PostHog-Sdk-Info": "posthog-python/", "User-Agent": "posthog-python/" }, "timeout": 15, - "url": "https://example.posthog.test/batch/" + "url": "https://example.posthog.test/i/v1/analytics/events" } diff --git a/posthog/test/snapshots/exception_event.json b/posthog/test/snapshots/exception_event.json index b971c7afe..59c72d814 100644 --- a/posthog/test/snapshots/exception_event.json +++ b/posthog/test/snapshots/exception_event.json @@ -1,10 +1,10 @@ { "body": { - "api_key": "phc_snapshot_project", "batch": [ { "distinct_id": "user-123", "event": "$exception", + "options": {}, "properties": { "$exception_list": [ { @@ -53,7 +53,7 @@ "", "def _exception_request():", " session = mock.MagicMock()", - " session.post.return_value = _successful_response()" + " session.post.side_effect = _capture_ok_response" ], "pre_context": [ "", @@ -113,8 +113,6 @@ "company": "company-456" }, "$is_server": true, - "$lib": "posthog-python", - "$lib_version": "", "$os": "", "$os_distro": "", "$os_version": "", @@ -127,13 +125,17 @@ "uuid": "00000000-0000-4000-8000-000000000005" } ], - "historical_migration": false, - "sent_at": "2026-01-02T03:04:05+00:00" + "created_at": "2026-01-02T03:04:05+00:00" }, "headers": { + "Authorization": "Bearer phc_snapshot_project", "Content-Type": "application/json", + "PostHog-Attempt": "1", + "PostHog-Request-Id": "", + "PostHog-Request-Timestamp": "2026-01-02T03:04:05+00:00", + "PostHog-Sdk-Info": "posthog-python/", "User-Agent": "posthog-python/" }, "timeout": 15, - "url": "https://example.posthog.test/batch/" + "url": "https://example.posthog.test/i/v1/analytics/events" } diff --git a/posthog/test/test_ai_capture_lane.py b/posthog/test/test_ai_capture_lane.py index 91b8c7256..6b8f45348 100644 --- a/posthog/test/test_ai_capture_lane.py +++ b/posthog/test/test_ai_capture_lane.py @@ -7,11 +7,9 @@ from posthog.ai.utils import _capture_ai_event, finalize_ai_content, with_privacy_mode from posthog.capture_compression import CaptureCompression -from posthog.capture_mode import CaptureMode from posthog.client import Client from posthog.consumer import AI_MAX_MSG_SIZE, MAX_MSG_SIZE from posthog.capture_v1 import _CAPTURE_AI_V1_PATH, _CAPTURE_V1_PATH -from posthog.request import EVENTS_ENDPOINT from posthog.version import VERSION from posthog.test.capture_helpers import patch_capture_send, sent_batch from posthog.test.test_utils import TEST_API_KEY @@ -116,7 +114,6 @@ def test_analytics_consumers_keep_todays_parameters(self): thread=2, flush_at=7, flush_interval=0.5, - gzip=True, max_retries=4, timeout=9, historical_migration=True, @@ -129,11 +126,9 @@ def test_analytics_consumers_keep_todays_parameters(self): self.assertEqual(consumer.max_msg_size, MAX_MSG_SIZE) self.assertEqual(consumer.flush_at, 7) self.assertEqual(consumer.flush_interval, 0.5) - self.assertTrue(consumer.gzip) self.assertEqual(consumer.retries, 4) self.assertEqual(consumer.timeout, 9) self.assertTrue(consumer.historical_migration) - self.assertEqual(consumer.capture_mode, client.capture_mode) self.assertEqual(consumer.capture_compression, client.capture_compression) client.join() @@ -195,67 +190,59 @@ def test_analytics_lane_rejects_events_over_900kib(self): self.assertTrue(client.queue.empty()) -class TestAiLaneAlwaysV1(unittest.TestCase): - """The AI lane posts capture v1 to the AI endpoint whatever `capture_mode`.""" +class TestAiLaneWireConfig(unittest.TestCase): + """The AI lane posts to the AI endpoint uncompressed, whatever the + analytics `capture_compression`.""" - def test_ai_lane_consumers_use_v1_and_ai_endpoint(self): - client = Client(TEST_API_KEY, send=False, capture_mode="v0", thread=2) + def test_ai_lane_consumers_use_ai_endpoint_without_compression(self): + client = Client(TEST_API_KEY, send=False, capture_compression="gzip", thread=2) client._ai_lane.start() self.assertEqual(len(client._ai_lane.consumers), 2) for consumer in client._ai_lane.consumers: self.assertIs(consumer.queue, client._ai_lane.queue) self.assertEqual(consumer.endpoint, _CAPTURE_AI_V1_PATH) self.assertEqual(consumer.max_msg_size, AI_MAX_MSG_SIZE) - self.assertEqual(consumer.capture_mode, CaptureMode.V1) self.assertEqual(consumer.capture_compression, CaptureCompression.NONE) - def test_async_ai_events_use_v1_with_capture_mode_v0(self): - client = Client( - TEST_API_KEY, - capture_mode="v0", - capture_compression="gzip", - flush_interval=0.05, - ) - with ( - mock.patch("posthog.consumer.batch_post") as mock_post, - patch_capture_send("consumer") as mock_v1, - ): + def test_async_lanes_keep_separate_path_and_compression(self): + client = Client(TEST_API_KEY, capture_compression="gzip", flush_interval=0.05) + with patch_capture_send("consumer") as mock_send: client.capture_ai("$ai_generation", distinct_id="d") client.capture("button_clicked", distinct_id="d") client.flush() + sends = { + call.kwargs["path"]: call.kwargs["compression"] + for call in mock_send.call_args_list + } self.assertEqual( - [call.kwargs["path"] for call in mock_post.call_args_list], - [EVENTS_ENDPOINT], + sends, + { + _CAPTURE_AI_V1_PATH: CaptureCompression.NONE, + _CAPTURE_V1_PATH: CaptureCompression.GZIP, + }, ) - mock_v1.assert_called_once() - self.assertEqual(mock_v1.call_args.kwargs["path"], _CAPTURE_AI_V1_PATH) self.assertEqual( - mock_v1.call_args.kwargs["compression"], CaptureCompression.NONE + _events_by_path(mock_send)[_CAPTURE_AI_V1_PATH][0]["event"], + "$ai_generation", ) - self.assertEqual([e["event"] for e in sent_batch(mock_v1)], ["$ai_generation"]) client.join() - def test_sync_ai_events_use_v1_with_capture_mode_v0(self): - client = Client( - TEST_API_KEY, - sync_mode=True, - capture_mode="v0", - capture_compression="gzip", - ) - with ( - mock.patch("posthog.client.batch_post") as mock_post, - patch_capture_send("client") as mock_v1, - ): + def test_sync_lanes_keep_separate_path_and_compression(self): + client = Client(TEST_API_KEY, sync_mode=True, capture_compression="gzip") + with patch_capture_send("client") as mock_send: client.capture_ai("$ai_generation", distinct_id="d") client.capture("button_clicked", distinct_id="d") - mock_post.assert_called_once() - self.assertEqual(mock_post.call_args.kwargs["path"], EVENTS_ENDPOINT) - mock_v1.assert_called_once() - self.assertEqual(mock_v1.call_args.kwargs["path"], _CAPTURE_AI_V1_PATH) self.assertEqual( - mock_v1.call_args.kwargs["compression"], CaptureCompression.NONE + [ + (call.kwargs["path"], call.kwargs["compression"]) + for call in mock_send.call_args_list + ], + [ + (_CAPTURE_AI_V1_PATH, CaptureCompression.NONE), + (_CAPTURE_V1_PATH, CaptureCompression.GZIP), + ], ) diff --git a/posthog/test/test_async_client.py b/posthog/test/test_async_client.py index 9c100bc76..34c098d7d 100644 --- a/posthog/test/test_async_client.py +++ b/posthog/test/test_async_client.py @@ -10,7 +10,7 @@ import pytest -from posthog import AsyncClient, AsyncPosthog, CaptureCompression, CaptureMode +from posthog import AsyncClient, AsyncPosthog, CaptureCompression from posthog.consumer import MAX_MSG_SIZE from posthog.contexts import ( new_context, @@ -306,7 +306,6 @@ async def test_capture_immediate_uses_capture_v1_without_building_httpx_client() ): client = AsyncPosthog( "test-key", - capture_mode=CaptureMode.V1, capture_compression=CaptureCompression.GZIP, ) event_uuid = await client.capture_immediate("event", distinct_id="user-1") @@ -319,23 +318,6 @@ async def test_capture_immediate_uses_capture_v1_without_building_httpx_client() assert send_v1.await_args.args[2][0]["uuid"] == event_uuid -@pytest.mark.asyncio -async def test_missing_async_extra_does_not_accept_undeliverable_events(): - with mock.patch( - "posthog.async_client._require_httpx", - side_effect=RuntimeError("install posthog[async]"), - ): - client = AsyncPosthog("test-key", capture_mode=CaptureMode.V0) - assert client.capture("event", distinct_id="user-1") is None - assert ( - client.set(distinct_id="user-1", properties={"email": "a@example.com"}) - is None - ) - assert client._pending_queue_items() == 0 - assert client._worker_tasks == [] - await client.shutdown() - - @pytest.mark.asyncio async def test_send_false_accepts_without_starting_workers_or_transport(): with mock.patch("posthog.async_client._build_client") as build_client: @@ -609,16 +591,17 @@ async def test_reuses_and_closes_instance_owned_http_client(): "posthog.async_client._build_client", return_value=http_client ) as build, mock.patch( - "posthog._async_consumer.async_batch_post", new=mock.AsyncMock() - ) as batch_post, + "posthog.async_client._async_flags", + new=mock.AsyncMock(return_value={"flags": {}}), + ) as flags, ): - client = AsyncPosthog("test-key", capture_mode=CaptureMode.V0) - await client.capture_immediate("first", distinct_id="user-1") - await client.capture_immediate("second", distinct_id="user-1") + client = AsyncPosthog("test-key") + await client._get_flags_decision("user-1") + await client._get_flags_decision("user-2") await client.shutdown() build.assert_called_once_with(client.host) - assert [call.kwargs["client"] for call in batch_post.await_args_list] == [ + assert [call.kwargs["client"] for call in flags.await_args_list] == [ http_client, http_client, ] diff --git a/posthog/test/test_async_consumer.py b/posthog/test/test_async_consumer.py index 9bb3f682f..4290a9419 100644 --- a/posthog/test/test_async_consumer.py +++ b/posthog/test/test_async_consumer.py @@ -1,17 +1,12 @@ from __future__ import annotations import asyncio -import json from unittest import mock -import httpx import pytest -from freezegun import freeze_time from posthog._async_consumer import _AsyncConsumer from posthog.capture_compression import CaptureCompression -from posthog.capture_mode import CaptureMode -from posthog.request import APIError def make_consumer(*, retries: int) -> _AsyncConsumer: @@ -23,137 +18,13 @@ def make_consumer(*, retries: int) -> _AsyncConsumer: process_event=mock.AsyncMock(side_effect=lambda event: event), flush_at=100, flush_interval=1, - gzip=False, retries=retries, timeout=3, historical_migration=False, - capture_mode=CaptureMode.V0, capture_compression=CaptureCompression.NONE, - http_client=mock.Mock(), ) -@pytest.mark.asyncio -@pytest.mark.parametrize( - ("failures", "retry_after", "expected_delays"), - [ - (1, None, [1]), - (2, None, [1, 2]), - (2, 5, [5, 5]), - ], -) -async def test_request_retries_transient_failures_until_success( - failures, retry_after, expected_delays -): - error = APIError(503, "temporary", retry_after=retry_after) - consumer = make_consumer(retries=failures) - - with ( - mock.patch( - "posthog._async_consumer.async_batch_post", - new=mock.AsyncMock(side_effect=[error] * failures + [None]), - ) as batch_post, - mock.patch( - "posthog._async_consumer.asyncio.sleep", new=mock.AsyncMock() - ) as sleep, - ): - await consumer.request([{"event": "test"}]) - - assert batch_post.await_count == failures + 1 - assert [call.args[0] for call in sleep.await_args_list] == expected_delays - - -@pytest.mark.asyncio -@pytest.mark.parametrize( - ("retry_after", "expected_delay"), - [ - ("Tue, 08 Sep 2026 00:00:10 GMT", 10), - ("Tue, 08 Sep 2026 00:01:00 GMT", 30), - ("Mon, 07 Sep 2026 23:59:59 GMT", 1), - ("5", 5), - ("0", 1), - ("-1", 1), - ("invalid", 1), - (None, 1), - ], -) -async def test_request_honors_retry_after_from_http_response( - retry_after, expected_delay -): - headers = {"Retry-After": retry_after} if retry_after is not None else {} - responses = [ - httpx.Response(503, headers=headers, json={"detail": "temporary"}), - httpx.Response(200, json={"ok": True}), - ] - requests = [] - - def handle_request(request): - requests.append(request) - return responses.pop(0) - - consumer = make_consumer(retries=1) - batch = [{"event": "test", "distinct_id": "test-user"}] - async with httpx.AsyncClient( - base_url="https://example.com", transport=httpx.MockTransport(handle_request) - ) as client: - consumer.http_client = client - with ( - freeze_time("2026-09-08 00:00:00", real_asyncio=True), - mock.patch( - "posthog._async_consumer.asyncio.sleep", new=mock.AsyncMock() - ) as sleep, - ): - await consumer.request(batch) - - sleep.assert_awaited_once_with(expected_delay) - assert len(requests) == 2 - assert [json.loads(request.content)["batch"] for request in requests] == [ - batch, - batch, - ] - - -@pytest.mark.asyncio -@pytest.mark.parametrize("status", [400, 401, 413]) -async def test_request_does_not_retry_terminal_client_errors(status): - consumer = make_consumer(retries=3) - - with ( - mock.patch( - "posthog._async_consumer.async_batch_post", - new=mock.AsyncMock(side_effect=APIError(status, "terminal")), - ) as batch_post, - mock.patch( - "posthog._async_consumer.asyncio.sleep", new=mock.AsyncMock() - ) as sleep, - pytest.raises(APIError), - ): - await consumer.request([{"event": "test"}]) - - batch_post.assert_awaited_once() - sleep.assert_not_awaited() - - -@pytest.mark.asyncio -async def test_request_stops_after_configured_retry_limit(): - consumer = make_consumer(retries=2) - - with ( - mock.patch( - "posthog._async_consumer.async_batch_post", - new=mock.AsyncMock(side_effect=APIError(503, "temporary")), - ) as batch_post, - mock.patch( - "posthog._async_consumer.asyncio.sleep", new=mock.AsyncMock() - ) as sleep, - pytest.raises(APIError), - ): - await consumer.request([{"event": "test"}]) - - assert batch_post.await_count == 3 - assert [call.args[0] for call in sleep.await_args_list] == [1, 2] - - @pytest.mark.asyncio @pytest.mark.parametrize("run_worker", [False, True], ids=["wait", "worker"]) async def test_get_or_flush_cancels_waiters_on_cancellation(run_worker): diff --git a/posthog/test/test_async_request.py b/posthog/test/test_async_request.py index d95da855a..f269db258 100644 --- a/posthog/test/test_async_request.py +++ b/posthog/test/test_async_request.py @@ -1,6 +1,5 @@ from __future__ import annotations -import asyncio import json import logging import subprocess @@ -13,7 +12,6 @@ from posthog._async_request import ( _build_client, _process_response, - async_batch_post, async_flags, async_remote_config, ) @@ -84,148 +82,6 @@ def test_build_client_scopes_requests_to_host_without_following_redirects(): ) -@pytest.mark.asyncio -async def test_async_batch_post_uses_configured_host_and_sanitized_logs(caplog): - caplog.set_level(logging.DEBUG, logger="posthog") - client = FakeAsyncClient() - - await async_batch_post( - "test-secret-key", - "https://example.com", - batch=[{"properties": {"password": "super-secret"}}], - path="/batch/", - client=client, - ) - - assert client.calls[0][1] == ("https://example.com/batch/",) - assert "super-secret" not in caplog.text - assert "test-secret-key" not in caplog.text - assert "https://example.com" not in caplog.text - - -@pytest.mark.asyncio -async def test_async_batch_post_follows_same_origin_temporary_redirect(): - client = FakeAsyncClient( - [ - FakeResponse(307, headers={"Location": "/redirected-batch/"}), - FakeResponse(200), - ] - ) - - await async_batch_post( - "test-key", - "https://example.com", - batch=[{"event": "event"}], - path="/batch/", - client=client, - ) - - assert [call[1] for call in client.calls] == [ - ("https://example.com/batch/",), - ("https://example.com/redirected-batch/",), - ] - - -@pytest.mark.asyncio -@pytest.mark.parametrize("status", [307, 308]) -@pytest.mark.parametrize( - "host", ["https://example.com/proxy", "https://example.com/proxy/"] -) -@pytest.mark.parametrize( - ("location", "redirected_path"), - [ - ("/proxy/redirected-batch/", "/proxy/redirected-batch/"), - ("https://example.com/proxy/redirected-batch/", "/proxy/redirected-batch/"), - ("/redirected-batch/", "/redirected-batch/"), - ("../redirected-batch/", "/proxy/redirected-batch/"), - ("?accepted=1", "/proxy/batch/?accepted=1"), - ], -) -async def test_async_batch_post_redirects_with_host_path_prefix( - status, host, location, redirected_path -): - requests = [] - responses = [ - httpx.Response(status, headers={"Location": location}), - httpx.Response(status, headers={"Location": "?attempt=2"}), - httpx.Response(200), - ] - - def handle_request(request): - requests.append(request) - return responses.pop(0) - - batch = [{"event": "test", "distinct_id": "test-user"}] - async with httpx.AsyncClient( - base_url=host, transport=httpx.MockTransport(handle_request) - ) as client: - await async_batch_post( - "test-key", host, batch=batch, path="/batch/", client=client - ) - - assert [str(request.url) for request in requests] == [ - "https://example.com/proxy/batch/", - f"https://example.com{redirected_path}", - f"https://example.com{redirected_path.split('?')[0]}?attempt=2", - ] - assert all(request.method == "POST" for request in requests) - assert all(request.content == requests[0].content for request in requests) - assert json.loads(requests[0].content)["batch"] == batch - - -@pytest.mark.asyncio -async def test_async_batch_post_rejects_cross_origin_temporary_redirect(): - client = FakeAsyncClient( - FakeResponse( - 307, - headers={"Location": "https://attacker.example/redirected-batch/"}, - ) - ) - - with pytest.raises(APIError): - await async_batch_post( - "test-key", - "https://example.com", - batch=[{"event": "event"}], - path="/batch/", - client=client, - ) - - assert len(client.calls) == 1 - - -@pytest.mark.asyncio -async def test_async_batch_post_serializes_off_event_loop(): - client = FakeAsyncClient() - real_to_thread = asyncio.to_thread - - with mock.patch( - "posthog._async_request.asyncio.to_thread", wraps=real_to_thread - ) as to_thread: - await async_batch_post( - "test-key", - "https://example.com", - batch=[{"event": "event"}], - path="/batch/", - gzip=True, - client=client, - ) - - to_thread.assert_awaited_once() - - -@pytest.mark.asyncio -async def test_async_batch_post_rejects_absolute_request_path(): - with pytest.raises(ValueError, match="relative"): - await async_batch_post( - "test-key", - "https://example.com", - batch=[], - path="https://attacker.example/batch/", - client=FakeAsyncClient(), - ) - - @pytest.mark.asyncio async def test_async_flags_sends_v2_request_payload(): client = FakeAsyncClient(FakeResponse(200, {"flags": {}})) diff --git a/posthog/test/test_capture_compression.py b/posthog/test/test_capture_compression.py index a4ce0154f..cb9376572 100644 --- a/posthog/test/test_capture_compression.py +++ b/posthog/test/test_capture_compression.py @@ -22,15 +22,9 @@ def setUp(self) -> None: self.addCleanup(patcher.stop) os.environ.pop(CAPTURE_COMPRESSION_ENV_VAR, None) - def test_defaults_to_none_with_no_kwarg_env_or_gzip(self) -> None: + def test_defaults_to_none_with_no_kwarg_or_env(self) -> None: self.assertIs(_resolve_capture_compression(None), CaptureCompression.NONE) - def test_gzip_fallback_used_when_nothing_else_set(self) -> None: - self.assertIs( - _resolve_capture_compression(None, gzip_fallback=True), - CaptureCompression.GZIP, - ) - @parameterized.expand( [ ("enum_gzip", CaptureCompression.GZIP, CaptureCompression.GZIP), @@ -48,12 +42,9 @@ def test_gzip_fallback_used_when_nothing_else_set(self) -> None: def test_explicit_kwarg_takes_precedence_and_coerces( self, _name, kwarg, expected ) -> None: - # Env names a different value and gzip_fallback is on, so each row proves - # the explicit kwarg wins over both lower-precedence sources. + # Env names a different value, so each row proves the explicit kwarg wins. with mock.patch.dict(os.environ, {CAPTURE_COMPRESSION_ENV_VAR: "deflate"}): - self.assertIs( - _resolve_capture_compression(kwarg, gzip_fallback=True), expected - ) + self.assertIs(_resolve_capture_compression(kwarg), expected) def test_invalid_kwarg_raises_even_with_valid_env(self) -> None: with mock.patch.dict(os.environ, {CAPTURE_COMPRESSION_ENV_VAR: "gzip"}): @@ -80,28 +71,16 @@ def test_env_var_resolution(self, _name, env_value, expected) -> None: with mock.patch.dict(os.environ, {CAPTURE_COMPRESSION_ENV_VAR: env_value}): self.assertIs(_resolve_capture_compression(None), expected) - def test_env_var_takes_precedence_over_gzip_fallback(self) -> None: - with mock.patch.dict(os.environ, {CAPTURE_COMPRESSION_ENV_VAR: "deflate"}): - self.assertIs( - _resolve_capture_compression(None, gzip_fallback=True), - CaptureCompression.DEFLATE, - ) - @parameterized.expand([("empty", ""), ("whitespace", " ")]) - def test_blank_env_var_falls_through_to_fallback(self, _name, env_value) -> None: + def test_blank_env_var_falls_through_to_none(self, _name, env_value) -> None: with mock.patch.dict(os.environ, {CAPTURE_COMPRESSION_ENV_VAR: env_value}): self.assertIs(_resolve_capture_compression(None), CaptureCompression.NONE) - self.assertIs( - _resolve_capture_compression(None, gzip_fallback=True), - CaptureCompression.GZIP, - ) - def test_unrecognized_env_var_warns_and_uses_fallback(self) -> None: + def test_unrecognized_env_var_warns_and_uses_none(self) -> None: with mock.patch.dict(os.environ, {CAPTURE_COMPRESSION_ENV_VAR: "bogus"}): with capture_message_only_logs() as stream: self.assertIs( - _resolve_capture_compression(None, gzip_fallback=True), - CaptureCompression.GZIP, + _resolve_capture_compression(None), CaptureCompression.NONE ) self.assertIn("bogus", stream.getvalue()) @@ -112,13 +91,12 @@ def test_explicit_zstd_without_package_raises(self, _name, kwarg) -> None: _resolve_capture_compression(kwarg) self.assertIn("posthog[zstd]", str(ctx.exception)) - def test_env_zstd_without_package_warns_and_uses_fallback(self) -> None: + def test_env_zstd_without_package_warns_and_uses_none(self) -> None: with mock.patch("posthog.capture_compression._zstandard", None): with mock.patch.dict(os.environ, {CAPTURE_COMPRESSION_ENV_VAR: "zstd"}): with capture_message_only_logs() as stream: self.assertIs( - _resolve_capture_compression(None, gzip_fallback=True), - CaptureCompression.GZIP, + _resolve_capture_compression(None), CaptureCompression.NONE ) self.assertIn("posthog[zstd]", stream.getvalue()) @@ -134,10 +112,6 @@ def test_client_defaults_to_none(self) -> None: client = Client(TEST_API_KEY, sync_mode=True) self.assertIs(client.capture_compression, CaptureCompression.NONE) - def test_client_gzip_flag_falls_back_to_gzip(self) -> None: - client = Client(TEST_API_KEY, sync_mode=True, gzip=True) - self.assertIs(client.capture_compression, CaptureCompression.GZIP) - @parameterized.expand( [ ("enum_deflate", CaptureCompression.DEFLATE, CaptureCompression.DEFLATE), @@ -145,11 +119,8 @@ def test_client_gzip_flag_falls_back_to_gzip(self) -> None: ("str_none", "none", CaptureCompression.NONE), ] ) - def test_client_kwarg_overrides_gzip_flag(self, _name, kwarg, expected) -> None: - # Even with the legacy gzip flag on, the explicit kwarg wins. - client = Client( - TEST_API_KEY, sync_mode=True, gzip=True, capture_compression=kwarg - ) + def test_client_kwarg_sets_compression(self, _name, kwarg, expected) -> None: + client = Client(TEST_API_KEY, sync_mode=True, capture_compression=kwarg) self.assertIs(client.capture_compression, expected) def test_client_propagates_to_consumers(self) -> None: diff --git a/posthog/test/test_capture_mode.py b/posthog/test/test_capture_mode.py deleted file mode 100644 index e63705d45..000000000 --- a/posthog/test/test_capture_mode.py +++ /dev/null @@ -1,109 +0,0 @@ -import os -import unittest -from unittest import mock - -from parameterized import parameterized - -from posthog.capture_mode import ( - CAPTURE_MODE_ENV_VAR, - CaptureMode, - _resolve_capture_mode, -) -from posthog.client import Client -from posthog.consumer import Consumer -from posthog.test.logging_helpers import capture_message_only_logs -from posthog.test.test_utils import TEST_API_KEY - - -class TestResolveCaptureMode(unittest.TestCase): - def test_defaults_to_v1_with_no_kwarg_and_no_env(self) -> None: - with mock.patch.dict(os.environ, {}, clear=False): - os.environ.pop(CAPTURE_MODE_ENV_VAR, None) - self.assertIs(_resolve_capture_mode(None), CaptureMode.V1) - - @parameterized.expand( - [ - # (name, kwarg, expected, opposite_env): the env always names the - # mode the kwarg must override, so every row proves the kwarg wins. - ("enum_v0", CaptureMode.V0, CaptureMode.V0, "v1"), - ("enum_v1", CaptureMode.V1, CaptureMode.V1, "v0"), - ("str_v0", "v0", CaptureMode.V0, "v1"), - ("str_v1", "v1", CaptureMode.V1, "v0"), - ("str_legacy_alias", "legacy", CaptureMode.V0, "v1"), - ("str_analytics_v1_alias", "analytics_v1", CaptureMode.V1, "v0"), - ("str_upper_and_padded", " V1 ", CaptureMode.V1, "v0"), - ] - ) - def test_explicit_kwarg_takes_precedence_and_coerces( - self, _name, kwarg, expected, opposite_env - ) -> None: - with mock.patch.dict(os.environ, {CAPTURE_MODE_ENV_VAR: opposite_env}): - self.assertIs(_resolve_capture_mode(kwarg), expected) - - def test_invalid_kwarg_raises_even_with_valid_env(self) -> None: - # The kwarg path is consulted before the env, so an invalid kwarg raises - # rather than silently falling back to a valid env value. - with mock.patch.dict(os.environ, {CAPTURE_MODE_ENV_VAR: "v1"}): - with self.assertRaises(ValueError): - _resolve_capture_mode("bogus") - - @parameterized.expand( - [ - ("v0", "v0", CaptureMode.V0), - ("legacy", "legacy", CaptureMode.V0), - ("v1", "v1", CaptureMode.V1), - ("analytics_v1", "analytics_v1", CaptureMode.V1), - ("uppercase", "V1", CaptureMode.V1), - ("padded", " v1 ", CaptureMode.V1), - ] - ) - def test_env_var_resolution(self, _name, env_value, expected) -> None: - with mock.patch.dict(os.environ, {CAPTURE_MODE_ENV_VAR: env_value}): - self.assertIs(_resolve_capture_mode(None), expected) - - @parameterized.expand([("empty", ""), ("whitespace", " ")]) - def test_blank_env_var_defaults_to_v1(self, _name, env_value) -> None: - with mock.patch.dict(os.environ, {CAPTURE_MODE_ENV_VAR: env_value}): - self.assertIs(_resolve_capture_mode(None), CaptureMode.V1) - - def test_unrecognized_env_var_warns_and_defaults_to_v1(self) -> None: - with mock.patch.dict(os.environ, {CAPTURE_MODE_ENV_VAR: "bogus"}): - with capture_message_only_logs() as stream: - self.assertIs(_resolve_capture_mode(None), CaptureMode.V1) - self.assertIn("bogus", stream.getvalue()) - - @parameterized.expand([("bad_str", "bogus"), ("wrong_type", 1)]) - def test_invalid_explicit_kwarg_raises(self, _name, value) -> None: - with self.assertRaises(ValueError): - _resolve_capture_mode(value) - - -class TestCaptureModePlumbing(unittest.TestCase): - def test_client_resolves_and_stores_default_v1(self) -> None: - with mock.patch.dict(os.environ, {}, clear=False): - os.environ.pop(CAPTURE_MODE_ENV_VAR, None) - client = Client(TEST_API_KEY, sync_mode=True) - self.assertIs(client.capture_mode, CaptureMode.V1) - - @parameterized.expand( - [ - ("enum_v1", CaptureMode.V1, CaptureMode.V1), - ("str_v1", "v1", CaptureMode.V1), - ("enum_v0", CaptureMode.V0, CaptureMode.V0), - ] - ) - def test_client_kwarg_sets_mode(self, _name, kwarg, expected) -> None: - client = Client(TEST_API_KEY, sync_mode=True, capture_mode=kwarg) - self.assertIs(client.capture_mode, expected) - - def test_client_propagates_mode_to_consumers(self) -> None: - # Async (non-sync) client builds Consumer threads; assert each carries - # the resolved mode. - client = Client(TEST_API_KEY, capture_mode=CaptureMode.V0, send=False, thread=2) - self.assertEqual(len(client.consumers), 2) - for consumer in client.consumers: - self.assertIs(consumer.capture_mode, CaptureMode.V0) - - def test_consumer_defaults_to_v1(self) -> None: - consumer = Consumer(None, TEST_API_KEY) - self.assertIs(consumer.capture_mode, CaptureMode.V1) diff --git a/posthog/test/test_client.py b/posthog/test/test_client.py index faef807c4..10474a25c 100644 --- a/posthog/test/test_client.py +++ b/posthog/test/test_client.py @@ -3503,8 +3503,12 @@ def test_numeric_distinct_id(self): def test_debug(self): Client("bad_key", debug=True) - def test_gzip(self): - client = Client(FAKE_TEST_API_KEY, on_error=self.fail, gzip=True) + def test_gzip_compression(self): + client = Client( + FAKE_TEST_API_KEY, + on_error=self.fail, + capture_compression=CaptureCompression.GZIP, + ) for _ in range(10): client.capture( "event", distinct_id="distinct_id", properties={"trait": "value"} @@ -4660,30 +4664,18 @@ def test_debug_flag_re_raises_exceptions(self, mock_enqueue): class TestClientCaptureRetrySemantics(unittest.TestCase): - @parameterized.expand( - [ - ("v0_sync", "v0", True), - ("v1_sync", "v1", True), - ("v0_async", "v0", False), - ("v1_async", "v1", False), - ] - ) - def test_negative_max_retries_still_attempts_delivery_once( - self, _name, capture_mode, sync_mode - ): + @parameterized.expand([("sync", True), ("async", False)]) + def test_negative_max_retries_still_attempts_delivery_once(self, _name, sync_mode): response = mock.Mock(status_code=200, headers={}, text="") response.json.return_value = {"results": {}} client = None - with ( - mock.patch("posthog.client.batch_post") as sync_v0_post, - mock.patch("posthog.consumer.batch_post") as async_v0_post, - mock.patch("posthog.capture_v1._post_v1", return_value=response) as v1_post, - ): + with mock.patch( + "posthog.capture_v1._post_v1", return_value=response + ) as v1_post: try: client = Client( FAKE_TEST_API_KEY, - capture_mode=capture_mode, sync_mode=sync_mode, max_retries=-1, flush_at=1, @@ -4694,61 +4686,30 @@ def test_negative_max_retries_still_attempts_delivery_once( client.flush() self.assertEqual(client.max_retries, 0) - if capture_mode == "v1": - v1_post.assert_called_once() - sync_v0_post.assert_not_called() - async_v0_post.assert_not_called() - elif sync_mode: - sync_v0_post.assert_called_once() - async_v0_post.assert_not_called() - v1_post.assert_not_called() - else: - async_v0_post.assert_called_once() - sync_v0_post.assert_not_called() - v1_post.assert_not_called() + v1_post.assert_called_once() finally: if client is not None and not sync_mode: client.shutdown() -class TestClientSyncCaptureMode(unittest.TestCase): - """Sync-mode `_enqueue` selects the analytics submitter by `capture_mode`.""" +class TestClientSyncCapture(unittest.TestCase): + """Sync-mode `_enqueue` sends analytics events through the v1 submitter.""" def _client(self, **kwargs): return Client(FAKE_TEST_API_KEY, sync_mode=True, **kwargs) - @parameterized.expand( - [ - ("default", None, True), - ("v1", "v1", True), - ("v0", "v0", False), - ] - ) - def test_capture_mode_selects_sync_submitter(self, _name, capture_mode, expects_v1): - kwargs = {"capture_mode": capture_mode} if capture_mode else {} - with ( - mock.patch("posthog.client.batch_post") as mock_post, - mock.patch("posthog.client._send_v1_batch") as mock_v1, - ): - self._client(**kwargs).capture("evt", distinct_id="d") - if expects_v1: - mock_post.assert_not_called() - mock_v1.assert_called_once() - batch = mock_v1.call_args.args[2] - self.assertEqual(len(batch), 1) - self.assertEqual(batch[0]["event"], "evt") - self.assertEqual(mock_v1.call_args.kwargs["path"], _CAPTURE_V1_PATH) - else: - mock_v1.assert_not_called() - mock_post.assert_called_once() - - def test_v1_sync_forwards_config_to_submitter(self): - with ( - mock.patch("posthog.client.batch_post"), - mock.patch("posthog.client._send_v1_batch") as mock_v1, - ): + def test_sync_capture_posts_to_analytics_endpoint(self): + with mock.patch("posthog.client._send_v1_batch") as mock_v1: + self._client().capture("evt", distinct_id="d") + mock_v1.assert_called_once() + batch = mock_v1.call_args.args[2] + self.assertEqual(len(batch), 1) + self.assertEqual(batch[0]["event"], "evt") + self.assertEqual(mock_v1.call_args.kwargs["path"], _CAPTURE_V1_PATH) + + def test_sync_forwards_config_to_submitter(self): + with mock.patch("posthog.client._send_v1_batch") as mock_v1: self._client( - capture_mode="v1", capture_compression=CaptureCompression.GZIP, max_retries=4, historical_migration=True, @@ -4758,29 +4719,15 @@ def test_v1_sync_forwards_config_to_submitter(self): self.assertEqual(kwargs["max_retries"], 4) self.assertEqual(kwargs["historical_migration"], True) - def test_v1_sync_gzip_flag_falls_back_to_gzip_compression(self): - # Legacy `gzip=True` with no explicit capture_compression -> GZIP on v1. - with ( - mock.patch("posthog.client.batch_post"), - mock.patch("posthog.client._send_v1_batch") as mock_v1, - ): - self._client(capture_mode="v1", gzip=True).capture("evt", distinct_id="d") - self.assertEqual( - mock_v1.call_args.kwargs["compression"], CaptureCompression.GZIP - ) - - def test_v1_sync_ai_named_event_through_capture_uses_v1(self): + def test_sync_ai_named_event_through_capture_uses_analytics_endpoint(self): # `capture()` never special-cases AI events: an `$ai_*`-named event - # follows `capture_mode` and rides the v1 submitter like any analytics - # event. Only `capture_ai()` reaches the AI lane. - with ( - mock.patch("posthog.client.batch_post") as mock_post, - mock.patch("posthog.client._send_v1_batch") as mock_v1, - ): - client = self._client(capture_mode="v1") + # rides the analytics endpoint like any other event. Only + # `capture_ai()` reaches the AI lane. + with mock.patch("posthog.client._send_v1_batch") as mock_v1: + client = self._client() client.capture("$ai_generation", distinct_id="d") - mock_post.assert_not_called() mock_v1.assert_called_once() batch = mock_v1.call_args.args[2] self.assertEqual(len(batch), 1) self.assertEqual(batch[0]["event"], "$ai_generation") + self.assertEqual(mock_v1.call_args.kwargs["path"], _CAPTURE_V1_PATH) diff --git a/posthog/test/test_consumer.py b/posthog/test/test_consumer.py index 1199d2da1..92cd9d888 100644 --- a/posthog/test/test_consumer.py +++ b/posthog/test/test_consumer.py @@ -2,9 +2,6 @@ import threading import time import unittest -from datetime import datetime, timedelta, timezone -from email.utils import format_datetime -from typing import Any from unittest import mock from parameterized import parameterized @@ -15,10 +12,8 @@ from Queue import Queue from posthog.capture_compression import CaptureCompression -from posthog.capture_mode import CaptureMode from posthog.capture_v1 import _CAPTURE_AI_V1_PATH, _CAPTURE_V1_PATH from posthog.consumer import MAX_MSG_SIZE, Consumer, _DrainSignal -from posthog.request import AI_EVENTS_ENDPOINT, EVENTS_ENDPOINT, APIError from posthog.test.capture_helpers import patch_capture_send, sent_batch from posthog.test.logging_helpers import capture_message_only_logs from posthog.test.test_utils import TEST_API_KEY @@ -329,88 +324,15 @@ def record_batch(api_key, host, batch, **kwargs): consumer.join(15) self.assertFalse(consumer.is_alive()) - def test_request(self) -> None: - consumer = Consumer(None, TEST_API_KEY, capture_mode=CaptureMode.V0) - batch = [_track_event()] - with mock.patch("posthog.consumer.batch_post") as post: - consumer.request(batch) - post.assert_called_once_with( - TEST_API_KEY, - None, - gzip=False, - timeout=15, - batch=batch, - historical_migration=False, - path="/batch/", - ) - - def _run_retry_test( - self, - exception: Exception, - exception_count: int, - retries: int = 10, - expected_attempts: int = 3, - raises: bool = False, - ) -> None: - call_count = [0] - - def mock_post(*args: Any, **kwargs: Any) -> None: - call_count[0] += 1 - if call_count[0] <= exception_count: - raise exception - - consumer = Consumer( - None, TEST_API_KEY, retries=retries, capture_mode=CaptureMode.V0 - ) - batch = [_track_event()] - with ( - mock.patch("posthog.consumer.batch_post", side_effect=mock_post) as post, - mock.patch("posthog.consumer.time.sleep"), - ): - if raises: - with self.assertRaises(type(exception)) as raised: - consumer.request(batch) - self.assertIs(raised.exception, exception) - else: - consumer.request(batch) - self.assertEqual(post.call_count, expected_attempts) - for call in post.call_args_list: - self.assertEqual(call.kwargs["batch"], batch) - - @parameterized.expand( - [ - ("general_errors", Exception("generic exception"), 2), - ("server_errors", APIError(500, "Internal Server Error"), 2), - ("rate_limit_errors", APIError(429, "Too Many Requests"), 2), - ] - ) - def test_request_retries_on_retriable_errors( - self, _name: str, exception: Exception, exception_count: int - ) -> None: - self._run_retry_test(exception, exception_count) - - def test_request_does_not_retry_client_errors(self) -> None: - self._run_retry_test( - APIError(400, "Client Errors"), 1, expected_attempts=1, raises=True - ) - - def test_request_fails_when_exceptions_exceed_retries(self) -> None: - self._run_retry_test( - APIError(500, "Internal Server Error"), - 4, - retries=3, - expected_attempts=4, - raises=True, - ) - def test_negative_retries_still_attempts_delivery_once(self) -> None: - consumer = Consumer(None, TEST_API_KEY, retries=-1, capture_mode=CaptureMode.V0) + consumer = Consumer(None, TEST_API_KEY, retries=-1) - with mock.patch("posthog.consumer.batch_post") as mock_post: + with patch_capture_send("consumer") as mock_send: consumer.request([_track_event()]) self.assertEqual(consumer.retries, 0) - mock_post.assert_called_once() + mock_send.assert_called_once() + self.assertEqual(mock_send.call_args.kwargs["max_retries"], 0) def test_pause(self) -> None: consumer = Consumer(None, TEST_API_KEY) @@ -624,101 +546,6 @@ def test_max_batch_size(self) -> None: consumer.join(5) self.assertFalse(consumer.is_alive()) - def test_request_sleeps_with_retry_after(self) -> None: - error = APIError(429, "Too Many Requests", retry_after=5.0) - call_count = [0] - - def mock_post(*args: Any, **kwargs: Any) -> None: - call_count[0] += 1 - if call_count[0] <= 1: - raise error - - consumer = Consumer(None, TEST_API_KEY, retries=3, capture_mode=CaptureMode.V0) - with ( - mock.patch("posthog.consumer.batch_post", side_effect=mock_post), - mock.patch("posthog.consumer.time.sleep") as mock_sleep, - ): - consumer.request([_track_event()]) - mock_sleep.assert_called_once_with(5.0) - - def test_request_uses_exponential_backoff_without_retry_after(self) -> None: - error = APIError(503, "Service Unavailable") - call_count = [0] - - def mock_post(*args: Any, **kwargs: Any) -> None: - call_count[0] += 1 - if call_count[0] <= 3: - raise error - - consumer = Consumer(None, TEST_API_KEY, retries=3, capture_mode=CaptureMode.V0) - with ( - mock.patch("posthog.consumer.batch_post", side_effect=mock_post), - mock.patch("posthog.consumer.time.sleep") as mock_sleep, - ): - consumer.request([_track_event()]) - self.assertEqual( - mock_sleep.call_args_list, - [ - mock.call(1), # 2^0 - mock.call(2), # 2^1 - mock.call(4), # 2^2 - ], - ) - - @parameterized.expand( - [ - ("huge_numeric", "1000000000", [30, 30]), - ("small_numeric", "0.25", [1, 2]), - ("huge_date", "Fri, 01 Jan 2100 00:00:00 GMT", [30, 30]), - ("small_date", None, [1, 2]), - ] - ) - def test_request_bounds_retry_after_without_reducing_attempts( - self, _name: str, retry_after_header: str | None, expected_sleeps: list[int] - ) -> None: - if retry_after_header is None: - retry_after_header = format_datetime( - datetime.now(timezone.utc) + timedelta(seconds=1), usegmt=True - ) - - retry_response = mock.Mock( - status_code=503, - headers={"Retry-After": retry_after_header}, - text="Service Unavailable", - ) - retry_response.json.return_value = {"detail": "Service Unavailable"} - success_response = mock.Mock(status_code=200) - session = mock.Mock() - session.post.side_effect = [retry_response, retry_response, success_response] - - consumer = Consumer(None, TEST_API_KEY, retries=2, capture_mode=CaptureMode.V0) - with ( - mock.patch("posthog.request._get_session", return_value=session), - mock.patch("posthog.consumer.time.sleep") as mock_sleep, - ): - consumer.request([_track_event()]) - - self.assertEqual(session.post.call_count, 3) - self.assertEqual( - [call.args[0] for call in mock_sleep.call_args_list], expected_sleeps - ) - - def test_request_retries_on_408(self) -> None: - call_count = [0] - - def mock_post(*args: Any, **kwargs: Any) -> None: - call_count[0] += 1 - if call_count[0] <= 1: - raise APIError(408, "Request Timeout") - - consumer = Consumer(None, TEST_API_KEY, retries=3, capture_mode=CaptureMode.V0) - with ( - mock.patch("posthog.consumer.batch_post", side_effect=mock_post), - mock.patch("posthog.consumer.time.sleep"), - ): - consumer.request([_track_event()]) - self.assertEqual(call_count[0], 2) - @parameterized.expand( [ ("on_error_succeeds", False), @@ -755,51 +582,28 @@ def _ai_event(event_name: str = "$ai_generation") -> dict[str, str]: return {"type": "track", "event": event_name, "distinct_id": "distinct_id"} -class TestConsumerCaptureModeRouting(unittest.TestCase): - """`capture_mode` selects the submitter; both post to the consumer's `endpoint`.""" +class TestConsumerSubmitterRouting(unittest.TestCase): + """Every consumer sends through the capture v1 submitter to its `endpoint`.""" - @parameterized.expand( - [ - ("default", None, True), - ("v0", CaptureMode.V0, False), - ("v1", CaptureMode.V1, True), - ] - ) - def test_capture_mode_selects_analytics_submitter( - self, _name, mode, expects_v1 - ) -> None: - kwargs = {"capture_mode": mode} if mode else {} - consumer = Consumer(None, TEST_API_KEY, **kwargs) + def test_default_posts_to_analytics_endpoint(self) -> None: + consumer = Consumer(None, TEST_API_KEY) batch = [_track_event()] - with ( - mock.patch("posthog.consumer.batch_post") as mock_post, - mock.patch("posthog.consumer._send_v1_batch") as mock_v1, - ): + with patch_capture_send("consumer") as mock_v1: consumer.request(batch) - if expects_v1: - mock_post.assert_not_called() - mock_v1.assert_called_once() - self.assertEqual(mock_v1.call_args.args[2], batch) - self.assertEqual(mock_v1.call_args.kwargs["path"], _CAPTURE_V1_PATH) - else: - mock_v1.assert_not_called() - mock_post.assert_called_once() - self.assertEqual(mock_post.call_args.kwargs["path"], EVENTS_ENDPOINT) - - def test_v1_forwards_consumer_config_to_submitter(self) -> None: + mock_v1.assert_called_once() + self.assertEqual(sent_batch(mock_v1), batch) + self.assertEqual(mock_v1.call_args.kwargs["path"], _CAPTURE_V1_PATH) + + def test_forwards_consumer_config_to_submitter(self) -> None: consumer = Consumer( None, TEST_API_KEY, - capture_mode=CaptureMode.V1, capture_compression=CaptureCompression.DEFLATE, timeout=7, retries=4, historical_migration=True, ) - with ( - mock.patch("posthog.consumer.batch_post"), - mock.patch("posthog.consumer._send_v1_batch") as mock_v1, - ): + with patch_capture_send("consumer") as mock_v1: consumer.request([_track_event()]) kwargs = mock_v1.call_args.kwargs self.assertEqual(kwargs["compression"], CaptureCompression.DEFLATE) @@ -807,7 +611,7 @@ def test_v1_forwards_consumer_config_to_submitter(self) -> None: self.assertEqual(kwargs["max_retries"], 4) self.assertEqual(kwargs["historical_migration"], True) - def test_v1_posts_to_configured_endpoint(self) -> None: + def test_posts_to_configured_endpoint(self) -> None: consumer = Consumer(None, TEST_API_KEY, endpoint=_CAPTURE_AI_V1_PATH) batch = [_ai_event()] with patch_capture_send("consumer") as mock_v1: @@ -815,31 +619,3 @@ def test_v1_posts_to_configured_endpoint(self) -> None: mock_v1.assert_called_once() self.assertEqual(mock_v1.call_args.kwargs["path"], _CAPTURE_AI_V1_PATH) self.assertEqual(sent_batch(mock_v1), batch) - - def test_v0_posts_to_configured_endpoint(self) -> None: - consumer = Consumer( - None, - TEST_API_KEY, - endpoint=AI_EVENTS_ENDPOINT, - capture_mode=CaptureMode.V0, - ) - batch = [_ai_event()] - with mock.patch("posthog.consumer.batch_post") as mock_post: - consumer.request(batch) - mock_post.assert_called_once() - self.assertEqual(mock_post.call_args.kwargs["path"], AI_EVENTS_ENDPOINT) - self.assertEqual(mock_post.call_args.kwargs["batch"], batch) - - def test_v1_routes_whole_batch_through_v1_submitter(self) -> None: - # A consumer doesn't know AI events exist: with `capture_mode` v1, the - # whole batch (including `$ai_*`-named events) rides the v1 submitter. - consumer = Consumer(None, TEST_API_KEY, capture_mode=CaptureMode.V1) - batch = [_ai_event(), _track_event()] - with ( - mock.patch("posthog.consumer.batch_post") as mock_post, - mock.patch("posthog.consumer._send_v1_batch") as mock_v1, - ): - consumer.request(batch) - mock_v1.assert_called_once() - self.assertEqual(mock_v1.call_args.args[2], batch) - mock_post.assert_not_called() diff --git a/posthog/test/test_flag_definition_cache.py b/posthog/test/test_flag_definition_cache.py index 2d05997d4..d221245d9 100644 --- a/posthog/test/test_flag_definition_cache.py +++ b/posthog/test/test_flag_definition_cache.py @@ -17,6 +17,7 @@ ) from posthog.request import GetResponse from posthog.test.test_utils import FAKE_TEST_API_KEY +from posthog.test.capture_helpers import patch_capture_send class MockCacheProvider: @@ -94,8 +95,8 @@ class TestFlagDefinitionCacheProvider(unittest.TestCase): @classmethod def setUpClass(cls): # Prevent real HTTP requests - cls.client_post_patcher = mock.patch("posthog.client.batch_post") - cls.consumer_post_patcher = mock.patch("posthog.consumer.batch_post") + cls.client_post_patcher = patch_capture_send("client") + cls.consumer_post_patcher = patch_capture_send("consumer") cls.client_post_patcher.start() cls.consumer_post_patcher.start() diff --git a/posthog/test/test_request.py b/posthog/test/test_request.py index d4ad7c9da..c16e08f47 100644 --- a/posthog/test/test_request.py +++ b/posthog/test/test_request.py @@ -1,7 +1,5 @@ -import gzip import json import unittest -import zlib from datetime import date, datetime, timedelta from unittest import mock @@ -18,7 +16,6 @@ KEEP_ALIVE_SOCKET_OPTIONS, QuotaLimitError, _mask_tokens_in_url, - batch_post, determine_server_host, disable_connection_reuse, enable_keep_alive, @@ -132,47 +129,6 @@ def test_message_only_debug_logs_include_posthog_prefix(): class TestRequests(unittest.TestCase): - def test_valid_request(self): - response = requests.Response() - response.status_code = 200 - session = mock.Mock() - session.post.return_value = response - batch = [ - {"distinct_id": "distinct_id", "event": "python event", "type": "track"} - ] - - res = batch_post(TEST_API_KEY, batch=batch, session=session) - - self.assertIs(res, response) - session.post.assert_called_once() - self.assertTrue(session.post.call_args.args[0].endswith("/batch/")) - body = json.loads(session.post.call_args.kwargs["data"]) - self.assertEqual(body["batch"], batch) - self.assertEqual(body["api_key"], TEST_API_KEY) - - def test_invalid_request_error(self): - response = requests.Response() - response.status_code = 400 - response._content = b'{"detail": "Invalid batch"}' - session = mock.Mock() - session.post.return_value = response - - with self.assertRaises(APIError) as raised: - batch_post("testsecret", batch=[], session=session) - - self.assertEqual(raised.exception.status, 400) - self.assertEqual(raised.exception.message, "Invalid batch") - session.post.assert_called_once() - - def test_invalid_host(self): - self.assertRaises( - requests.exceptions.MissingSchema, - batch_post, - "testsecret", - "t.posthog.com/", - batch=[], - ) - def test_post_without_path_preserves_type_error(self): mock_session = mock.MagicMock() @@ -185,7 +141,7 @@ def test_post_without_path_preserves_type_error(self): mock_session.post.assert_not_called() - def test_post_sends_string_payload_without_gzip(self): + def test_post_sends_string_payload(self): mock_response = requests.Response() mock_response.status_code = 200 mock_session = mock.MagicMock() @@ -205,54 +161,6 @@ def test_post_sends_string_payload_without_gzip(self): self.assertEqual(url, "https://test.posthog.com/batch/") self.assertIsInstance(data, str) - def test_post_sends_bytes_payload_with_gzip(self): - mock_response = requests.Response() - mock_response.status_code = 200 - mock_session = mock.MagicMock() - mock_session.post.return_value = mock_response - - request_module.post( - TEST_API_KEY, - host="https://test.posthog.com", - path="/batch/", - gzip=True, - session=mock_session, - batch=[], - ) - - data = mock_session.post.call_args.kwargs["data"] - headers = mock_session.post.call_args.kwargs["headers"] - self.assertIsInstance(data, bytes) - self.assertEqual(headers["Content-Encoding"], "gzip") - body = json.loads(gzip.decompress(data)) - self.assertEqual(body["batch"], []) - self.assertEqual(body["api_key"], TEST_API_KEY) - - def test_post_falls_back_to_uncompressed_payload_when_gzip_fails(self): - for compression_error in [OSError("boom"), zlib.error("boom")]: - with self.subTest(compression_error=type(compression_error)): - mock_response = requests.Response() - mock_response.status_code = 200 - mock_session = mock.MagicMock() - mock_session.post.return_value = mock_response - - with mock.patch.object( - request_module, "GzipFile", side_effect=compression_error - ): - request_module.post( - TEST_API_KEY, - host="https://test.posthog.com", - path="/batch/", - gzip=True, - session=mock_session, - batch=[], - ) - - data = mock_session.post.call_args.kwargs["data"] - headers = mock_session.post.call_args.kwargs["headers"] - self.assertIsInstance(data, str) - self.assertNotIn("Content-Encoding", headers) - def test_datetime_serialization(self): data = {"created": datetime(2012, 3, 4, 5, 6, 7, 891011)} result = json.dumps(data, cls=DatetimeSerializer) @@ -265,30 +173,6 @@ def test_date_serialization(self): expected = '{"created": "%s"}' % today.isoformat() self.assertEqual(result, expected) - def test_should_not_timeout(self): - response = requests.Response() - response.status_code = 200 - session = mock.Mock() - session.post.return_value = response - - self.assertIs( - batch_post(TEST_API_KEY, batch=[], timeout=7, session=session), response - ) - session.post.assert_called_once() - self.assertEqual(session.post.call_args.kwargs["timeout"], 7) - - def test_should_timeout(self): - error = requests.ReadTimeout("response deadline exceeded") - session = mock.Mock() - session.post.side_effect = error - - with self.assertRaises(requests.ReadTimeout) as raised: - batch_post("key", batch=[], timeout=1, session=session) - - self.assertIs(raised.exception, error) - session.post.assert_called_once() - self.assertEqual(session.post.call_args.kwargs["timeout"], 1) - def test_quota_limited_flags_response(self): mock_response = requests.Response() mock_response.status_code = 200 diff --git a/posthog/test/test_server_payload_snapshots.py b/posthog/test/test_server_payload_snapshots.py index 2874e0761..80e03415d 100644 --- a/posthog/test/test_server_payload_snapshots.py +++ b/posthog/test/test_server_payload_snapshots.py @@ -29,6 +29,16 @@ def _successful_response() -> requests.Response: return response +def _capture_ok_response(url, data=None, **kwargs) -> requests.Response: + response = requests.Response() + response.status_code = 200 + uuids = [event["uuid"] for event in json.loads(data)["batch"]] + response._content = json.dumps( + {"results": {uuid: {"result": "ok"} for uuid in uuids}} + ).encode() + return response + + def _has_test_file_suffix(value: str) -> bool: return value.replace("\\", "/").endswith(_TEST_FILE_SUFFIX) @@ -43,8 +53,10 @@ def _normalize_snapshot_value(value): for key, item in value.items(): if key == "$lib_version": normalized[key] = "" + elif key == "PostHog-Request-Id": + normalized[key] = "" elif ( - key == "User-Agent" + key in ("User-Agent", "PostHog-Sdk-Info") and isinstance(item, str) and _USER_AGENT_PATTERN.fullmatch(item) ): @@ -99,19 +111,18 @@ def _assert_json_snapshot(name: str, value) -> None: assert actual == expected -def _legacy_event_family_request(): +def _event_family_request(): session = mock.MagicMock() - session.post.return_value = _successful_response() + session.post.side_effect = _capture_ok_response with ( freeze_time(_FIXED_TIME), - mock.patch("posthog.request._get_session", return_value=session), + mock.patch("posthog.capture_v1._get_session", return_value=session), mock.patch("posthog.client.system_context", return_value=_RUNTIME_CONTEXT), ): client = Client( "phc_snapshot_project", host="https://example.posthog.test", - capture_mode="v0", flush_at=100, flush_interval=100, ) @@ -167,18 +178,17 @@ def _raise_snapshot_exception() -> None: def _exception_request(): session = mock.MagicMock() - session.post.return_value = _successful_response() + session.post.side_effect = _capture_ok_response with ( freeze_time(_FIXED_TIME), - mock.patch("posthog.request._get_session", return_value=session), + mock.patch("posthog.capture_v1._get_session", return_value=session), mock.patch("posthog.client.system_context", return_value=_RUNTIME_CONTEXT), mock.patch("posthog.client._get_current_otel_span_properties", return_value={}), ): client = Client( "phc_snapshot_project", host="https://example.posthog.test", - capture_mode="v0", sync_mode=True, project_root=str(Path(__file__).parents[2]), capture_exception_code_variables=False, @@ -257,8 +267,8 @@ def test_does_not_normalize_unexpected_user_agent(user_agent): } -def test_legacy_capture_identify_alias_and_group_identify_request_snapshot(): - _assert_json_snapshot("legacy_event_family", _legacy_event_family_request()) +def test_capture_identify_alias_and_group_identify_request_snapshot(): + _assert_json_snapshot("event_family", _event_family_request()) def test_complete_exception_request_snapshot(): diff --git a/pyproject.toml b/pyproject.toml index 17052d31f..7ad70ea0a 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -24,7 +24,6 @@ classifiers = [ ] dependencies = [ "requests>=2.7,<3.0", - "backoff>=1.10.0", "distro>=1.5.0", "typing-extensions>=4.2.0", ] diff --git a/references/public_api_snapshot.txt b/references/public_api_snapshot.txt index 319451d6a..f412a2113 100644 --- a/references/public_api_snapshot.txt +++ b/references/public_api_snapshot.txt @@ -7,7 +7,6 @@ alias posthog.AsyncClient -> posthog.async_client.AsyncClient alias posthog.AsyncPosthog -> posthog.async_client.AsyncPosthog alias posthog.BeforeSendCallback -> posthog.types.BeforeSendCallback alias posthog.CaptureCompression -> posthog.capture_compression.CaptureCompression -alias posthog.CaptureMode -> posthog.capture_mode.CaptureMode alias posthog.Client -> posthog.client.Client alias posthog.DEFAULT_CODE_VARIABLES_DETECT_SECRETS -> posthog.exception_utils.DEFAULT_CODE_VARIABLES_DETECT_SECRETS alias posthog.DEFAULT_CODE_VARIABLES_IGNORE_PATTERNS -> posthog.exception_utils.DEFAULT_CODE_VARIABLES_IGNORE_PATTERNS @@ -235,13 +234,11 @@ alias posthog.args.SendFeatureFlagsOptions -> posthog.types.SendFeatureFlagsOpti alias posthog.client.AI_MAX_MSG_SIZE -> posthog.consumer.AI_MAX_MSG_SIZE alias posthog.client.APIError -> posthog.request.APIError alias posthog.client.CaptureCompression -> posthog.capture_compression.CaptureCompression -alias posthog.client.CaptureMode -> posthog.capture_mode.CaptureMode alias posthog.client.Consumer -> posthog.consumer.Consumer alias posthog.client.DEFAULT_CODE_VARIABLES_DETECT_SECRETS -> posthog.exception_utils.DEFAULT_CODE_VARIABLES_DETECT_SECRETS alias posthog.client.DEFAULT_CODE_VARIABLES_IGNORE_PATTERNS -> posthog.exception_utils.DEFAULT_CODE_VARIABLES_IGNORE_PATTERNS alias posthog.client.DEFAULT_CODE_VARIABLES_MASK_PATTERNS -> posthog.exception_utils.DEFAULT_CODE_VARIABLES_MASK_PATTERNS alias posthog.client.DEFAULT_CODE_VARIABLES_MASK_URL_CREDENTIALS -> posthog.exception_utils.DEFAULT_CODE_VARIABLES_MASK_URL_CREDENTIALS -alias posthog.client.EVENTS_ENDPOINT -> posthog.request.EVENTS_ENDPOINT alias posthog.client.ExceptionArg -> posthog.args.ExceptionArg alias posthog.client.ExceptionCapture -> posthog.exception_capture.ExceptionCapture alias posthog.client.FeatureFlag -> posthog.types.FeatureFlag @@ -272,7 +269,6 @@ alias posthog.client.SendFeatureFlagsOptions -> posthog.types.SendFeatureFlagsOp alias posthog.client.SizeLimitedDict -> posthog.utils.SizeLimitedDict alias posthog.client.Span -> posthog.tracing.span.Span alias posthog.client.VERSION -> posthog.version.VERSION -alias posthog.client.batch_post -> posthog.request.batch_post alias posthog.client.clean -> posthog.utils.clean alias posthog.client.determine_server_host -> posthog.request.determine_server_host alias posthog.client.exc_info_from_error -> posthog.exception_utils.exc_info_from_error @@ -303,12 +299,8 @@ alias posthog.client.to_flags_and_payloads -> posthog.types.to_flags_and_payload alias posthog.client.to_payloads -> posthog.types.to_payloads alias posthog.client.to_values -> posthog.types.to_values alias posthog.client.try_attach_code_variables_to_frames -> posthog.exception_utils.try_attach_code_variables_to_frames -alias posthog.consumer.APIError -> posthog.request.APIError alias posthog.consumer.CaptureCompression -> posthog.capture_compression.CaptureCompression -alias posthog.consumer.CaptureMode -> posthog.capture_mode.CaptureMode alias posthog.consumer.DatetimeSerializer -> posthog.request.DatetimeSerializer -alias posthog.consumer.EVENTS_ENDPOINT -> posthog.request.EVENTS_ENDPOINT -alias posthog.consumer.batch_post -> posthog.request.batch_post alias posthog.contexts.Client -> posthog.client.Client alias posthog.disable_connection_reuse -> posthog.request.disable_connection_reuse alias posthog.enable_keep_alive -> posthog.request.enable_keep_alive @@ -659,9 +651,8 @@ attribute posthog.args.OptionalSetArgs.timestamp: NotRequired[Optional[Union[dat attribute posthog.args.OptionalSetArgs.uuid: NotRequired[Optional[Union[str, UUID]]] attribute posthog.async_client.AsyncClient.api_key = (project_api_key or '').strip() attribute posthog.async_client.AsyncClient.before_send = before_send -attribute posthog.async_client.AsyncClient.capture_compression = _resolve_capture_compression(capture_compression, gzip_fallback=gzip) +attribute posthog.async_client.AsyncClient.capture_compression = _resolve_capture_compression(capture_compression) attribute posthog.async_client.AsyncClient.capture_exception_code_variables = capture_exception_code_variables -attribute posthog.async_client.AsyncClient.capture_mode = _resolve_capture_mode(capture_mode) attribute posthog.async_client.AsyncClient.capture_trace_context = capture_trace_context attribute posthog.async_client.AsyncClient.code_variables_detect_secrets = code_variables_detect_secrets if code_variables_detect_secrets is not None else DEFAULT_CODE_VARIABLES_DETECT_SECRETS attribute posthog.async_client.AsyncClient.code_variables_ignore_patterns = code_variables_ignore_patterns if code_variables_ignore_patterns is not None else DEFAULT_CODE_VARIABLES_IGNORE_PATTERNS @@ -673,7 +664,6 @@ attribute posthog.async_client.AsyncClient.disabled = disabled or not self.api_k attribute posthog.async_client.AsyncClient.distinct_ids_feature_flags_reported = SizeLimitedDict(_MAX_DICT_SIZE, set) attribute posthog.async_client.AsyncClient.feature_flags_request_max_retries = max(0, feature_flags_request_max_retries) attribute posthog.async_client.AsyncClient.feature_flags_request_timeout_seconds = feature_flags_request_timeout_seconds -attribute posthog.async_client.AsyncClient.gzip = gzip attribute posthog.async_client.AsyncClient.historical_migration = historical_migration attribute posthog.async_client.AsyncClient.host = determine_server_host(host) attribute posthog.async_client.AsyncClient.in_app_modules = in_app_modules @@ -699,18 +689,14 @@ attribute posthog.capture_compression.CaptureCompression.GZIP = 'gzip' attribute posthog.capture_compression.CaptureCompression.NONE = 'none' attribute posthog.capture_compression.CaptureCompression.ZSTD = 'zstd' attribute posthog.capture_exception_code_variables = False -attribute posthog.capture_mode.CAPTURE_MODE_ENV_VAR = 'POSTHOG_CAPTURE_MODE' -attribute posthog.capture_mode.CaptureMode.V0 = 'v0' -attribute posthog.capture_mode.CaptureMode.V1 = 'v1' attribute posthog.capture_trace_context = False attribute posthog.capture_v1.CaptureV1Error.attempts = attempts attribute posthog.capture_v1.CaptureV1Error.drops = drops or [] attribute posthog.capture_v1.CaptureV1Error.request_id = request_id attribute posthog.capture_v1.CaptureV1Error.retry_exhausted = retry_exhausted or [] attribute posthog.client.Client.api_key = (project_api_key or '').strip() -attribute posthog.client.Client.capture_compression = _resolve_capture_compression(capture_compression, gzip_fallback=gzip) +attribute posthog.client.Client.capture_compression = _resolve_capture_compression(capture_compression) attribute posthog.client.Client.capture_exception_code_variables = capture_exception_code_variables -attribute posthog.client.Client.capture_mode = _resolve_capture_mode(capture_mode) attribute posthog.client.Client.capture_trace_context = capture_trace_context attribute posthog.client.Client.code_variables_detect_secrets = code_variables_detect_secrets if code_variables_detect_secrets is not None else DEFAULT_CODE_VARIABLES_DETECT_SECRETS attribute posthog.client.Client.code_variables_ignore_patterns = code_variables_ignore_patterns if code_variables_ignore_patterns is not None else DEFAULT_CODE_VARIABLES_IGNORE_PATTERNS @@ -738,7 +724,6 @@ attribute posthog.client.Client.flag_cache = self._initialize_flag_cache(flag_fa attribute posthog.client.Client.flag_definition_version = 0 attribute posthog.client.Client.flag_fallback_cache_url = flag_fallback_cache_url attribute posthog.client.Client.group_type_mapping: Optional[dict[str, str]] = None -attribute posthog.client.Client.gzip = gzip attribute posthog.client.Client.historical_migration = historical_migration attribute posthog.client.Client.host = determine_server_host(host) attribute posthog.client.Client.in_app_modules = in_app_modules @@ -769,12 +754,10 @@ attribute posthog.consumer.AI_MAX_MSG_SIZE = 8 * 1024 * 1024 attribute posthog.consumer.BATCH_SIZE_LIMIT = 5 * 1024 * 1024 attribute posthog.consumer.Consumer.api_key = api_key attribute posthog.consumer.Consumer.capture_compression = capture_compression -attribute posthog.consumer.Consumer.capture_mode = capture_mode attribute posthog.consumer.Consumer.daemon = True attribute posthog.consumer.Consumer.endpoint = endpoint attribute posthog.consumer.Consumer.flush_at = flush_at attribute posthog.consumer.Consumer.flush_interval = flush_interval -attribute posthog.consumer.Consumer.gzip = gzip attribute posthog.consumer.Consumer.historical_migration = historical_migration attribute posthog.consumer.Consumer.host = host attribute posthog.consumer.Consumer.log = logging.getLogger('posthog') @@ -1004,13 +987,11 @@ attribute posthog.privacy_mode = False attribute posthog.project_api_key = None attribute posthog.project_root = None attribute posthog.release_id.RELEASE_ID_ENV_VAR = 'POSTHOG_RELEASE_ID' -attribute posthog.request.AI_EVENTS_ENDPOINT = '/i/v0/ai/batch/' attribute posthog.request.APIError.message = message attribute posthog.request.APIError.retry_after = retry_after attribute posthog.request.APIError.status = status attribute posthog.request.DEFAULT_HOST = US_INGESTION_ENDPOINT attribute posthog.request.EU_INGESTION_ENDPOINT = 'https://eu.i.posthog.com' -attribute posthog.request.EVENTS_ENDPOINT = '/batch/' attribute posthog.request.GetResponse.data: Any attribute posthog.request.GetResponse.etag: Optional[str] = None attribute posthog.request.GetResponse.not_modified: bool = False @@ -1178,14 +1159,13 @@ class posthog.ai.types.TokenUsage class posthog.ai.types.ToolInProgress class posthog.args.OptionalCaptureArgs class posthog.args.OptionalSetArgs -class posthog.async_client.AsyncClient(project_api_key: str, host: Optional[str] = None, debug: bool = False, max_queue_size: int = 10000, send: bool = True, on_error=None, flush_at: int = 100, flush_interval: float = 5.0, gzip: bool = False, max_retries: int = 3, timeout: int = 15, thread: int = 1, disabled: bool = False, disable_geoip: bool = True, is_server: bool = True, historical_migration: bool = False, super_properties: Optional[dict[str, Any]] = None, before_send=None, log_captured_exceptions: bool = False, project_root: Optional[str] = None, capture_exception_code_variables: bool = False, code_variables_mask_patterns=None, code_variables_ignore_patterns=None, code_variables_mask_url_credentials=None, code_variables_detect_secrets=None, in_app_modules: Optional[list[str]] = None, capture_mode: Optional[Union[CaptureMode, str]] = None, capture_compression: Optional[Union[CaptureCompression, str]] = None, capture_trace_context: bool = False, secret_key: Optional[str] = None, personal_api_key: Optional[str] = None, feature_flags_request_timeout_seconds: int = 3, feature_flags_request_max_retries: int = 1) +class posthog.async_client.AsyncClient(project_api_key: str, host: Optional[str] = None, *, debug: bool = False, max_queue_size: int = 10000, send: bool = True, on_error=None, flush_at: int = 100, flush_interval: float = 5.0, max_retries: int = 3, timeout: int = 15, thread: int = 1, disabled: bool = False, disable_geoip: bool = True, is_server: bool = True, historical_migration: bool = False, super_properties: Optional[dict[str, Any]] = None, before_send=None, log_captured_exceptions: bool = False, project_root: Optional[str] = None, capture_exception_code_variables: bool = False, code_variables_mask_patterns=None, code_variables_ignore_patterns=None, code_variables_mask_url_credentials=None, code_variables_detect_secrets=None, in_app_modules: Optional[list[str]] = None, capture_compression: Optional[Union[CaptureCompression, str]] = None, capture_trace_context: bool = False, secret_key: Optional[str] = None, personal_api_key: Optional[str] = None, feature_flags_request_timeout_seconds: int = 3, feature_flags_request_max_retries: int = 1) class posthog.async_client.AsyncPosthog class posthog.bucketed_rate_limiter.BucketedRateLimiter(bucket_size: Number, refill_rate: Number, refill_interval_seconds: Number, on_bucket_rate_limited: Optional[Callable[[Hashable], None]] = None, clock: Callable[[], float] = time.monotonic) class posthog.capture_compression.CaptureCompression -class posthog.capture_mode.CaptureMode class posthog.capture_v1.CaptureV1Error(status: int | str, message: str, *, retry_after: Optional[float] = None, request_id: Optional[str] = None, attempts: Optional[int] = None, retry_exhausted: Optional[list[str]] = None, drops: Optional[list[tuple[str, Optional[str]]]] = None) -class posthog.client.Client(project_api_key: str, host=None, debug=False, max_queue_size=10000, send=True, on_error=None, flush_at=100, flush_interval=5.0, gzip=False, max_retries=3, sync_mode=False, timeout=15, thread=1, poll_interval=30, personal_api_key=None, disabled=False, disable_geoip=True, is_server=True, historical_migration=False, feature_flags_request_timeout_seconds=3, feature_flags_request_max_retries=1, super_properties=None, enable_exception_autocapture=False, log_captured_exceptions=False, project_root=None, privacy_mode=False, before_send=None, flag_fallback_cache_url=None, enable_local_evaluation=True, flag_definition_cache_provider: Optional[FlagDefinitionCacheProvider] = None, capture_exception_code_variables=False, code_variables_mask_patterns=None, code_variables_ignore_patterns=None, code_variables_mask_url_credentials=None, code_variables_detect_secrets=None, in_app_modules: list[str] | None = None, enable_exception_autocapture_rate_limiting=False, exception_autocapture_bucket_size=ExceptionCapture.DEFAULT_BUCKET_SIZE, exception_autocapture_refill_rate=ExceptionCapture.DEFAULT_REFILL_RATE, exception_autocapture_refill_interval_seconds=ExceptionCapture.DEFAULT_REFILL_INTERVAL_SECONDS, capture_mode: Optional[Union[CaptureMode, str]] = None, capture_compression: Optional[Union[CaptureCompression, str]] = None, secret_key=None, metrics: Optional[dict] = None, enable_full_ai_capture=False, capture_trace_context=False, _use_ai_lane=False, _enable_multimodal_capture=False, traces: Optional[dict] = None) -class posthog.consumer.Consumer(queue, api_key, flush_at=100, host=None, on_error=None, flush_interval=5.0, gzip=False, retries=10, timeout=15, historical_migration=False, endpoint=None, max_msg_size=MAX_MSG_SIZE, capture_mode=CaptureMode.V1, capture_compression=CaptureCompression.NONE) +class posthog.client.Client(project_api_key: str, host=None, *, debug=False, max_queue_size=10000, send=True, on_error=None, flush_at=100, flush_interval=5.0, max_retries=3, sync_mode=False, timeout=15, thread=1, poll_interval=30, personal_api_key=None, disabled=False, disable_geoip=True, is_server=True, historical_migration=False, feature_flags_request_timeout_seconds=3, feature_flags_request_max_retries=1, super_properties=None, enable_exception_autocapture=False, log_captured_exceptions=False, project_root=None, privacy_mode=False, before_send=None, flag_fallback_cache_url=None, enable_local_evaluation=True, flag_definition_cache_provider: Optional[FlagDefinitionCacheProvider] = None, capture_exception_code_variables=False, code_variables_mask_patterns=None, code_variables_ignore_patterns=None, code_variables_mask_url_credentials=None, code_variables_detect_secrets=None, in_app_modules: list[str] | None = None, enable_exception_autocapture_rate_limiting=False, exception_autocapture_bucket_size=ExceptionCapture.DEFAULT_BUCKET_SIZE, 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, secret_key=None, metrics: Optional[dict] = None, enable_full_ai_capture=False, capture_trace_context=False, _use_ai_lane=False, _enable_multimodal_capture=False, traces: Optional[dict] = None) +class posthog.consumer.Consumer(queue, api_key, flush_at=100, host=None, on_error=None, flush_interval=5.0, retries=10, timeout=15, historical_migration=False, endpoint=_CAPTURE_V1_PATH, max_msg_size=MAX_MSG_SIZE, capture_compression=CaptureCompression.NONE) class posthog.contexts.ContextScope(parent=None, fresh: bool = False, capture_exceptions: bool = True, client: Optional[Client] = None) class posthog.exception_capture.ExceptionCapture(client: Client, rate_limiting_enabled=False, bucket_size=DEFAULT_BUCKET_SIZE, refill_rate=DEFAULT_REFILL_RATE, refill_interval_seconds=DEFAULT_REFILL_INTERVAL_SECONDS) class posthog.exception_utils.AnnotatedValue(value, metadata) @@ -1418,14 +1398,13 @@ function posthog.mcp.session_token.encode_session_id(payload: SessionTokenPayloa function posthog.mcp.session_token.read_mcp_session_header(headers: Any) -> Optional[str] function posthog.mcp.tools.get_more_tools_result() -> Dict[str, Any] function posthog.new_context(fresh: bool = False, capture_exceptions: Optional[bool] = None, client: Optional[Client] = None) -function posthog.request.batch_post(api_key: str, host: Optional[str] = None, gzip: bool = False, timeout: int = 15, path: str = EVENTS_ENDPOINT, **kwargs) -> requests.Response function posthog.request.determine_server_host(host: Optional[str]) -> str function posthog.request.disable_connection_reuse() -> None function posthog.request.enable_keep_alive() -> None -function posthog.request.flags(api_key: str, host: Optional[str] = None, gzip: bool = False, timeout: int = 15, max_retries: int = 1, **kwargs) -> Any +function posthog.request.flags(api_key: str, host: Optional[str] = None, timeout: int = 15, max_retries: int = 1, **kwargs) -> Any function posthog.request.get(api_key: str, url: str, host: Optional[str] = None, timeout: Optional[int] = None, etag: Optional[str] = None) -> GetResponse function posthog.request.normalize_host(host: Optional[str]) -> str -function posthog.request.post(api_key: str, host: Optional[str] = None, path: Optional[str] = None, gzip: bool = False, timeout: int = 15, session: Optional[requests.Session] = None, **kwargs) -> requests.Response +function posthog.request.post(api_key: str, host: Optional[str] = None, path: Optional[str] = None, timeout: int = 15, session: Optional[requests.Session] = None, **kwargs) -> requests.Response function posthog.request.remote_config(personal_api_key: str, project_api_key: str, host: Optional[str] = None, key: str = '', timeout: int = 15) -> Any function posthog.request.reset_sessions() -> None function posthog.request.set_socket_options(socket_options: Optional[SocketOptions]) -> None @@ -1762,7 +1741,6 @@ module posthog.args module posthog.async_client module posthog.bucketed_rate_limiter module posthog.capture_compression -module posthog.capture_mode module posthog.capture_v1 module posthog.client module posthog.consumer diff --git a/sdk_compliance_adapter/Dockerfile.v1 b/sdk_compliance_adapter/Dockerfile.v1 deleted file mode 100644 index 6891837a5..000000000 --- a/sdk_compliance_adapter/Dockerfile.v1 +++ /dev/null @@ -1,25 +0,0 @@ -FROM python:3.12-slim - -WORKDIR /app - -# Copy the SDK source code -COPY posthog/ /app/sdk/posthog/ -COPY setup.py pyproject.toml README.md LICENSE /app/sdk/ - -# Install the SDK from source -RUN cd /app/sdk && pip install --no-cache-dir -e . - -# Install adapter dependencies -RUN pip install --no-cache-dir flask python-dateutil - -# Copy adapter code -COPY sdk_compliance_adapter/adapter.py /app/adapter.py - -# Select the capture-v1 protocol; same adapter code, different runtime mode. -ENV CAPTURE_MODE=v1 - -# Expose port 8080 -EXPOSE 8080 - -# Run the adapter -CMD ["python", "/app/adapter.py"] diff --git a/sdk_compliance_adapter/README.md b/sdk_compliance_adapter/README.md index eae25916b..6b17bf5dd 100644 --- a/sdk_compliance_adapter/README.md +++ b/sdk_compliance_adapter/README.md @@ -30,7 +30,7 @@ The adapter implements the standard SDK adapter interface defined in the [test h ### Key Implementation Details -**Request Tracking**: The adapter monkey-patches `batch_post` to track all HTTP requests made by the SDK, including retries. +**Request Tracking**: The adapter monkey-patches the capture v1 `_post_v1` to track all HTTP requests made by the SDK, including retries. **State Management**: Thread-safe state tracking for events captured vs sent, retry attempts, and errors. diff --git a/sdk_compliance_adapter/adapter.py b/sdk_compliance_adapter/adapter.py index 2a68dbf56..b663193cb 100644 --- a/sdk_compliance_adapter/adapter.py +++ b/sdk_compliance_adapter/adapter.py @@ -17,8 +17,7 @@ from posthog.capture_compression import CaptureCompression from posthog.capture_v1 import _CAPTURE_V1_PATH from posthog.capture_v1 import _post_v1 as original_post_v1 -from posthog.request import EVENTS_ENDPOINT, USER_AGENT -from posthog.request import batch_post as original_batch_post +from posthog.request import USER_AGENT from posthog.version import VERSION # Configure logging @@ -29,16 +28,6 @@ app = Flask(__name__) -# Selects which capture protocol this adapter process speaks. Baked at build -# time via the CAPTURE_MODE env var ("v1" => capture-v1, anything else => legacy -# v0), mirroring the v0/v1 Dockerfile split. One process speaks one mode and -# advertises it via /health capabilities. -CAPTURE_MODE = os.environ.get("CAPTURE_MODE", "") - - -def is_v1() -> bool: - return CAPTURE_MODE == "v1" - class RequestInfo: """Information about an HTTP request made by the SDK""" @@ -81,7 +70,6 @@ def __init__(self): self.client: Optional[Client] = None self.remote_client: Client | None = None self.reload_thread: threading.Thread | None = None - self.retry_attempts: Dict[str, int] = {} # Track retry attempts by batch ID def reset(self): """Reset all state""" @@ -107,7 +95,6 @@ def reset(self): self.total_retries = 0 self.last_error = None self.requests_made = [] - self.retry_attempts = {} def increment_captured(self): """Increment total events captured""" @@ -115,37 +102,6 @@ def increment_captured(self): self.total_events_captured += 1 self.pending_events += 1 - def record_request(self, status_code: int, batch: List[Dict], batch_id: str): - """Record an HTTP request made by the SDK""" - with self.lock: - # Determine retry attempt for this batch - retry_attempt = self.retry_attempts.get(batch_id, 0) - - # Extract UUIDs from batch - uuid_list = [event.get("uuid", "") for event in batch] - - request_info = RequestInfo( - timestamp_ms=int(time.time() * 1000), - status_code=status_code, - retry_attempt=retry_attempt, - event_count=len(batch), - uuid_list=uuid_list, - ) - self.requests_made.append(request_info) - - # Update counters - if status_code == 200: - # Success - clear pending events - self.total_events_sent += len(batch) - self.pending_events = max(0, self.pending_events - len(batch)) - # Remove batch from retry tracking - self.retry_attempts.pop(batch_id, None) - else: - # Failure - increment retry count - self.retry_attempts[batch_id] = retry_attempt + 1 - if retry_attempt > 0: - self.total_retries += 1 - def record_request_v1( self, status_code: int, batch: List[Dict], attempt: int, terminal_count: int ): @@ -195,40 +151,6 @@ def get_state(self) -> Dict[str, Any]: state = SDKState() -def create_batch_id(batch: List[Dict]) -> str: - """Create a unique ID for a batch based on UUIDs""" - uuids = sorted([event.get("uuid", "") for event in batch]) - return "-".join(uuids[:3]) # Use first 3 UUIDs as batch ID - - -def patched_batch_post( - api_key: str, - host: Optional[str] = None, - gzip: bool = False, - timeout: int = 15, - path: str = EVENTS_ENDPOINT, - **kwargs, -): - """Patched version of batch_post that tracks requests""" - batch = kwargs.get("batch", []) - batch_id = create_batch_id(batch) - - try: - # Call original batch_post - response = original_batch_post(api_key, host, gzip, timeout, path, **kwargs) - # Record successful request - state.record_request(200, batch, batch_id) - return response - except Exception as e: - # Record failed request - status_code = ( - getattr(e, "status_code", 500) if hasattr(e, "status_code") else 500 - ) - state.record_request(status_code, batch, batch_id) - state.record_error(str(e)) - raise - - def patched_post_v1( api_key: str, host: Optional[str], @@ -244,8 +166,8 @@ def patched_post_v1( ): """Patched version of _post_v1 that records requests for /state assertions. - Mirrors the legacy `patched_batch_post`, but reads the retry attempt from the - call (1-based) and counts only terminal per-event results as sent. + Reads the retry attempt from the call (1-based) and counts only terminal + per-event results as sent. """ batch = batch_body.get("batch", []) try: @@ -287,16 +209,6 @@ def patched_post_v1( return response -# Monkey-patch the batch_post function -import posthog.request # noqa: E402 - -posthog.request.batch_post = patched_batch_post - -# Also patch in consumer module -import posthog.consumer # noqa: E402 - -posthog.consumer.batch_post = patched_batch_post - # Patch the capture-v1 submitter. `_send_v1_batch` resolves `_post_v1` as a module # global at call time, so patching it here covers both the async consumer and the # sync client paths. @@ -310,10 +222,11 @@ def health(): """Health check endpoint""" # No AI capture capability: `capture_ai` posts capture v1 to # /i/v1/ai/events, which this harness version has no suite for. - capabilities = ( - ["capture_v1", "encoding_gzip"] if is_v1() else ["capture_v0", "encoding_gzip"] - ) - capabilities.append("feature_flags_local_evaluation_v1") + capabilities = [ + "capture_v1", + "encoding_gzip", + "feature_flags_local_evaluation_v1", + ] return jsonify( { "sdk_name": "posthog-python", @@ -354,9 +267,6 @@ def init(): # Convert flush_interval from ms to seconds flush_interval = flush_interval_ms / 1000.0 - # One adapter process speaks one capture protocol, selected by CAPTURE_MODE. - capture_mode = "v1" if is_v1() else "v0" - # Explicit reloads exercise the real loader without background polling # racing the harness's per-test definition snapshots. client_options = { @@ -364,12 +274,15 @@ def init(): "host": host, "flush_at": flush_at, "flush_interval": flush_interval, - "gzip": enable_compression, + "capture_compression": ( + CaptureCompression.GZIP + if enable_compression + else CaptureCompression.NONE + ), "max_retries": max_retries, "debug": False, "disable_geoip": disable_geoip, "historical_migration": historical_migration, - "capture_mode": capture_mode, "enable_local_evaluation": False, } personal_api_key = data.get("personal_api_key") @@ -384,8 +297,8 @@ def init(): logger.info( f"Initialized SDK with api_key={api_key[:10]}..., host={host}, " f"flush_at={flush_at}, flush_interval={flush_interval}, " - f"max_retries={max_retries}, gzip={enable_compression}, " - f"capture_mode={capture_mode}, disable_geoip={disable_geoip}, " + f"max_retries={max_retries}, compression={enable_compression}, " + f"disable_geoip={disable_geoip}, " f"historical_migration={historical_migration}" ) @@ -417,9 +330,8 @@ def capture(): # Fold capture-v1 options back into the magic `$`-prefixed properties the # SDK lifts onto the wire `options` object. Renamed keys mirror the SDK's - # sentinel table; unknown keys get a bare `$` prefix. v0 has no wire - # options object, so this only applies in v1 mode. - if options and is_v1(): + # sentinel table; unknown keys get a bare `$` prefix. + if options: properties = dict(properties or {}) option_to_property = { "cookieless_mode": "$cookieless_mode", @@ -473,7 +385,7 @@ def capture_ai(): if not event: return jsonify({"error": "event is required"}), 400 - if options and is_v1(): + if options: properties = dict(properties or {}) option_to_property = { "cookieless_mode": "$cookieless_mode", diff --git a/sdk_compliance_adapter/docker-compose.yml b/sdk_compliance_adapter/docker-compose.yml index c679f5c4f..56ef72680 100644 --- a/sdk_compliance_adapter/docker-compose.yml +++ b/sdk_compliance_adapter/docker-compose.yml @@ -1,7 +1,7 @@ version: "3.8" services: - # PostHog Python SDK adapter (capture v0) + # PostHog Python SDK adapter sdk-adapter: build: context: .. @@ -11,16 +11,6 @@ services: networks: - test-network - # PostHog Python SDK adapter (capture v1) - sdk-adapter-v1: - build: - context: .. - dockerfile: sdk_compliance_adapter/Dockerfile.v1 - ports: - - "8082:8080" - networks: - - test-network - # Test harness test-harness: image: ghcr.io/posthog/sdk-test-harness:1.1.1 diff --git a/sdk_compliance_adapter/test_adapter.py b/sdk_compliance_adapter/test_adapter.py index ac1ddc294..2cd484a5d 100644 --- a/sdk_compliance_adapter/test_adapter.py +++ b/sdk_compliance_adapter/test_adapter.py @@ -20,8 +20,6 @@ def adapter(monkeypatch): # Importing the adapter installs transport instrumentation. Restore it after # every test so collecting these tests alongside SDK tests is safe. for module, name in [ - (posthog.request, "batch_post"), - (posthog.consumer, "batch_post"), (posthog.capture_v1, "_post_v1"), ]: monkeypatch.setattr(module, name, getattr(module, name)) @@ -79,14 +77,11 @@ def initialize(adapter, **overrides): return adapter.app.test_client() -@pytest.mark.parametrize("mode,capability", [("", "capture_v0"), ("v1", "capture_v1")]) -def test_health_opts_into_local_evaluation_without_losing_capture( - adapter, monkeypatch, mode, capability -): - monkeypatch.setattr(adapter, "CAPTURE_MODE", mode) +def test_health_opts_into_local_evaluation_without_losing_capture(adapter): capabilities = adapter.app.test_client().get("/health").json["capabilities"] assert "feature_flags_local_evaluation_v1" in capabilities - assert capability in capabilities + assert "capture_v1" in capabilities + assert "capture_v0" not in capabilities assert "capture_ai_v0" not in capabilities diff --git a/uv.lock b/uv.lock index 60046d858..0a2391359 100644 --- a/uv.lock +++ b/uv.lock @@ -319,15 +319,6 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/fb/95/adcb68e20c34162e9135f370d6e31737719c2b6f94bc953fe7ed1f10fe21/authlib-1.7.2-py2.py3-none-any.whl", hash = "sha256:3e1faedc9d87e7d56a164eca3ccb6ace0d61b94abe83e92242f8dc8bba9b4a9f", size = 259548, upload-time = "2026-05-06T08:10:21.436Z" }, ] -[[package]] -name = "backoff" -version = "2.2.1" -source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/47/d7/5bbeb12c44d7c4f2fb5b56abce497eb5ed9f34d85701de869acedd602619/backoff-2.2.1.tar.gz", hash = "sha256:03f829f5bb1923180821643f8753b0502c3b682293992485b0eef2807afa5cba", size = 17001, upload-time = "2022-10-05T19:19:32.061Z" } -wheels = [ - { url = "https://files.pythonhosted.org/packages/df/73/b6e24bd22e6720ca8ee9a85a0c4a2971af8497d8f3193fa05390cbd46e09/backoff-2.2.1-py3-none-any.whl", hash = "sha256:63579f9a0628e06278f7e47b7d7d5b6ce20dc65c5e96a6f3ca99a6adca0396e8", size = 15148, upload-time = "2022-10-05T19:19:30.546Z" }, -] - [[package]] name = "backports-asyncio-runner" version = "1.2.0" @@ -2773,7 +2764,6 @@ name = "posthog" version = "7.64.1" source = { editable = "." } dependencies = [ - { name = "backoff" }, { name = "distro" }, { name = "requests" }, { name = "typing-extensions" }, @@ -2852,7 +2842,6 @@ dev = [ [package.metadata] requires-dist = [ { name = "anthropic", marker = "extra == 'test'", specifier = ">=0.72" }, - { name = "backoff", specifier = ">=1.10.0" }, { name = "claude-agent-sdk", marker = "extra == 'test'" }, { name = "coverage", marker = "extra == 'test'" }, { name = "distro", specifier = ">=1.5.0" },