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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
63 changes: 61 additions & 2 deletions posthog/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,12 @@
from posthog.contexts import (
get_tags as inner_get_tags,
)
from posthog.contexts import (
set_context_option as inner_set_context_option,
)
from posthog.contexts import (
get_context_options as inner_get_context_options,
)
from posthog.exception_utils import (
DEFAULT_CODE_VARIABLES_DETECT_SECRETS,
DEFAULT_CODE_VARIABLES_IGNORE_PATTERNS,
Expand Down Expand Up @@ -299,6 +305,44 @@ def get_tags() -> Dict[str, Any]:
return inner_get_tags()


def set_context_option(key: str, value: Any) -> None:
"""
Set a capture option for every event captured within the current context.

Context options fill options an event leaves unset, before ``before_send``
runs, so the hook sees them. They override ``super_options``. An event's
own ``options`` override them, and ``before_send`` can change them.

Args:
key: The option name, such as ``"process_person_profile"``
value: The option value, sent as given

Examples:
```python
from posthog import new_context, set_context_option
with new_context():
set_context_option("process_person_profile", False)
```

Category:
Contexts
"""
return inner_set_context_option(key, value)


def get_context_options() -> Dict[str, Any]:
"""
Get all capture options from the current context.

Returns:
Dict of all capture options in the current context

Category:
Contexts
"""
return inner_get_context_options()


"""Settings.

These module-level settings configure the legacy global PostHog client used by
Expand Down Expand Up @@ -340,7 +384,18 @@ def get_tags() -> Dict[str, Any]:
feature_flags_request_max_retries: Number of retries for feature flag
requests after network, transport, or timeout failures. Defaults to 1.
Set to 0 to disable retries.
super_properties: Properties merged into every captured event.
super_properties: Properties for every captured event. They fill only
keys the event leaves unset, before ``before_send`` runs. ``$set``,
``$set_once``, ``$groups`` and ``$group_set`` fill one level deep. An
event's own properties and context tags override them, and
``before_send`` can change or remove them.
super_options: Capture options for every captured event, such as
``{"cookieless_mode": True}``. They fill only options the event leaves
unset, before ``before_send`` runs. An event's own ``options`` and
context options override them, and ``before_send`` can change them. They also
win over an event's legacy property for the same key, such as
``$cookieless_mode``, so pass per-event overrides of that key as
``options``.
metrics: Config dict for the ``client.metrics`` API (``service_name``,
``service_version``, ``environment``, ``flush_interval``, ...). Applied
when ``setup()`` builds the global client, or on a later ``setup()``
Expand All @@ -364,7 +419,9 @@ def get_tags() -> Dict[str, Any]:
project_root: Root path used to determine in-app exception stack frames.
privacy_mode: Capture AI usage metadata without prompt inputs or outputs.
before_send: Optional callback that can modify or drop events before upload.
Return ``None`` to drop an event.
Return ``None`` to drop an event. Context tags, context options,
``super_properties`` and ``super_options`` fill in before it runs, so
it can change or remove them.
enable_local_evaluation: Whether to poll feature flag definitions for local
evaluation when a personal API key is configured.
flag_definition_cache_provider: Optional external cache provider for sharing
Expand Down Expand Up @@ -414,6 +471,7 @@ def get_tags() -> Dict[str, Any]:
feature_flags_request_timeout_seconds = 3 # type: int
feature_flags_request_max_retries = 1 # type: int
super_properties = None # type: Optional[Dict]
super_options = None # type: Optional[Dict]
metrics = None # type: Optional[Dict]
traces = None # type: Optional[Dict]
enable_exception_autocapture = False # type: bool
Expand Down Expand Up @@ -1406,6 +1464,7 @@ def setup() -> Client:
feature_flags_request_timeout_seconds=feature_flags_request_timeout_seconds,
feature_flags_request_max_retries=feature_flags_request_max_retries,
super_properties=super_properties,
super_options=super_options,
metrics=metrics,
traces=traces,
# TODO: Currently this monitoring begins only when the Client is initialised (which happens when you do something with the SDK)
Expand Down
13 changes: 9 additions & 4 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_event import _EventDefaults
from .capture_send import _CAPTURE_V1_PATH, _capture_loss_message
from .consumer import BATCH_SIZE_LIMIT, MAX_MSG_SIZE
from .request import DatetimeSerializer
Expand All @@ -37,6 +38,7 @@ async def _run_outside_processing_event(awaitable):
class _QueuedEvent:
event: dict[str, Any]
context: contextvars.Context
defaults: Optional[_EventDefaults] = None


async def _invoke_callback(callback, *args):
Expand Down Expand Up @@ -92,7 +94,10 @@ def __init__(
*,
host: Optional[str],
on_error: Optional[Callable[[Exception, list[dict[str, Any]]], Any]],
process_event: Callable[[dict[str, Any]], Awaitable[Optional[dict[str, Any]]]],
process_event: Callable[
[dict[str, Any], Optional[_EventDefaults]],
Awaitable[Optional[dict[str, Any]]],
],
flush_at: int,
flush_interval: float,
retries: int,
Expand Down Expand Up @@ -158,11 +163,11 @@ async def _get_or_flush(self, timeout: float) -> tuple[Any, bool]:
return None, False

async def _process_queued_event(
self, event: dict[str, Any]
self, queued: _QueuedEvent
) -> Optional[dict[str, Any]]:
token = _PROCESSING_EVENT.set(True)
try:
return await self.process_event(event)
return await self.process_event(queued.event, queued.defaults)
finally:
_PROCESSING_EVENT.reset(token)

Expand Down Expand Up @@ -207,7 +212,7 @@ async def next(self) -> tuple[list[dict[str, Any]], bool]:

try:
process_task = queued.context.run(
asyncio.create_task, self._process_queued_event(queued.event)
asyncio.create_task, self._process_queued_event(queued)
)
item = await process_task
except Exception as error:
Expand Down
Loading
Loading