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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions posthog/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -313,8 +313,11 @@ def get_tags() -> Dict[str, Any]:
disabled no-op global client.
host: PostHog ingestion host. Defaults to the US ingestion endpoint when not
set.
on_error: Optional callback invoked by background consumers when event upload
fails. Keep it short and non-blocking. Lifecycle methods can be called
on_error: Optional callback ``(error, batch)`` invoked when event upload
fails: by background consumers, or on the calling thread in ``sync_mode``.
Capture failures arrive as ``CaptureError``. Without it, each failed batch
logs one aggregate line. A capture that fails inside the callback logs that
line instead of calling it again. Keep it short and non-blocking. Lifecycle methods can be called
directly and will be deferred, but the callback must not wait for another
thread or task that calls ``flush()``, ``join()``, or ``shutdown()``.
debug: Enable verbose SDK logging and re-raise errors from public APIs.
Expand Down
41 changes: 30 additions & 11 deletions posthog/_async_consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@

from ._async_request import async_send_v1_batch
from .capture_compression import CaptureCompression
from .capture_send import _CAPTURE_V1_PATH, _capture_loss_message
from .consumer import BATCH_SIZE_LIMIT, MAX_MSG_SIZE
from .request import DatetimeSerializer

Expand Down Expand Up @@ -47,6 +48,33 @@ async def _invoke_callback(callback, *args):
return result


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


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


async def _serialized_event_size(event: dict[str, Any]) -> int:
serialized = await asyncio.to_thread(json.dumps, event, cls=DatetimeSerializer)
return len(serialized.encode())
Expand Down Expand Up @@ -142,18 +170,9 @@ async def upload(self, batch: list[dict[str, Any]]) -> None:
try:
await self.request(batch)
except Exception as error:
self.log.error(
"async capture upload failed (%s, status=%s)",
type(error).__name__,
getattr(error, "status", None),
await _report_capture_failure(
self.on_error, self.log, error, batch, _CAPTURE_V1_PATH
)
if self.on_error:
try:
await _invoke_callback(self.on_error, error, batch)
except Exception as callback_error:
self.log.error(
"on_error handler failed (%s)", type(callback_error).__name__
)
finally:
for _ in batch:
self.queue.task_done()
Expand Down
29 changes: 14 additions & 15 deletions posthog/async_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,14 +11,15 @@
import weakref
from datetime import datetime, timezone
from typing import Any, Dict, Mapping, Optional, Union
from uuid import UUID, uuid4
from uuid import UUID

from typing_extensions import Unpack

from ._async_consumer import (
_STOP,
_AsyncConsumer,
_invoke_callback,
_report_capture_failure,
_is_processing_event,
_QueuedEvent,
_run_outside_processing_event,
Expand All @@ -34,6 +35,7 @@
CaptureCompression,
_resolve_capture_compression,
)
from .capture_send import _CAPTURE_V1_PATH
from .client import (
MAX_DICT_SIZE as _MAX_DICT_SIZE,
_MINIMAL_FLAG_CALLED_EVENT_PROPERTIES,
Expand Down Expand Up @@ -76,7 +78,13 @@
from .release_id import _resolve_release_id
from .request import QuotaLimitError, determine_server_host, normalize_host
from .types import FlagMetadata, FlagValue, normalize_flags_response
from .utils import SizeLimitedDict, _normalize_timestamp, clean, system_context
from .utils import (
SizeLimitedDict,
_normalize_timestamp,
_uuid7,
clean,
system_context,
)
from .version import VERSION

__all__ = ["AsyncClient", "AsyncPosthog"]
Expand Down Expand Up @@ -380,7 +388,7 @@ def _normalize_uuid(self, msg: dict[str, Any]) -> str:
msg["uuid"] = normalized
return normalized

normalized = str(uuid4())
normalized = str(_uuid7())
msg["uuid"] = normalized
return normalized

Expand Down Expand Up @@ -576,20 +584,11 @@ async def capture_immediate(
await consumer.request(error_batch)
return sent_uuid
except Exception as error:
if self.on_error:
try:
await _invoke_callback(self.on_error, error, error_batch)
except Exception as callback_error:
self.log.error(
"on_error handler failed (%s)", type(callback_error).__name__
)
await _report_capture_failure(
self.on_error, self.log, error, error_batch, _CAPTURE_V1_PATH
)
if self.debug:
raise
self.log.error(
"Immediate async capture failed (%s, status=%s)",
type(error).__name__,
getattr(error, "status", None),
)
return None
finally:
remaining_calls = self._immediate_callers[current] - 1
Expand Down
Loading
Loading