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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
48 changes: 43 additions & 5 deletions posthog/_async_consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -96,9 +96,11 @@ def __init__(
flush_at: int,
flush_interval: float,
retries: int,
timeout: int,
timeout: float,
historical_migration: bool,
capture_compression: CaptureCompression,
endpoint: str = _CAPTURE_V1_PATH,
max_msg_size: int = MAX_MSG_SIZE,
) -> None:
self.queue = queue
self.api_key = api_key
Expand All @@ -111,6 +113,8 @@ def __init__(
self.timeout = timeout
self.historical_migration = historical_migration
self.capture_compression = capture_compression
self.endpoint = endpoint
self.max_msg_size = max_msg_size
self._carryover: Optional[tuple[dict[str, Any], int]] = None
self._flush_event = asyncio.Event()

Expand Down Expand Up @@ -171,7 +175,7 @@ async def upload(self, batch: list[dict[str, Any]]) -> None:
await self.request(batch)
except Exception as error:
await _report_capture_failure(
self.on_error, self.log, error, batch, _CAPTURE_V1_PATH
self.on_error, self.log, error, batch, self.endpoint
)
finally:
for _ in batch:
Expand Down Expand Up @@ -229,12 +233,15 @@ async def next(self) -> tuple[list[dict[str, Any]], bool]:
self.queue.task_done()
continue

if item_size > MAX_MSG_SIZE:
if item_size > self.max_msg_size:
# Log only name and size: AI events may carry unredacted
# multimodal payloads that must not leak into logs.
self.log.error(
"Event %s (%d bytes) exceeds the %dKiB limit, dropping.",
"Event %s (%d bytes) exceeds the %dKiB limit for %s, dropping.",
item.get("event"),
item_size,
MAX_MSG_SIZE // 1024,
self.max_msg_size // 1024,
self.endpoint,
)
self.queue.task_done()
continue
Expand All @@ -258,4 +265,35 @@ async def request(self, batch: list[dict[str, Any]]) -> None:
timeout=self.timeout,
max_retries=self.retries,
historical_migration=self.historical_migration,
path=self.endpoint,
)


class _AsyncLane:
"""One capture queue, the consumer tasks that drain it, and the endpoint they post to.

The client owns one lane per traffic class (analytics, AI), so each gets
its own backpressure, timeout, size cap and compression.
"""

def __init__(
self,
*,
name: str,
max_queue_size: int,
endpoint: str,
max_msg_size: int,
timeout: float,
capture_compression: CaptureCompression,
) -> None:
self.name = name
self.queue: asyncio.Queue[Any] = asyncio.Queue(max_queue_size)
self.endpoint = endpoint
self.max_msg_size = max_msg_size
self.timeout = timeout
self.capture_compression = capture_compression
self.consumers: list[_AsyncConsumer] = []
self.worker_tasks: list[asyncio.Task[None]] = []

def pending_items(self) -> int:
return int(getattr(self.queue, "_unfinished_tasks", self.queue.qsize()))
6 changes: 4 additions & 2 deletions posthog/_async_request.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
from urllib.parse import quote

from .capture_compression import CaptureCompression
from .capture_send import _parse_retry_after, _send_v1_batch
from .capture_send import _CAPTURE_V1_PATH, _parse_retry_after, _send_v1_batch
from .request import (
APIError,
DatetimeSerializer,
Expand Down Expand Up @@ -150,9 +150,10 @@ async def async_send_v1_batch(
batch: list[dict[str, Any]],
*,
compression: CaptureCompression,
timeout: int,
timeout: float,
max_retries: int,
historical_migration: bool,
path: str = _CAPTURE_V1_PATH,
) -> None:
"""Run the existing capture-v1 submitter off-loop to preserve wire parity."""
await asyncio.to_thread(
Expand All @@ -164,4 +165,5 @@ async def async_send_v1_batch(
timeout=timeout,
max_retries=max_retries,
historical_migration=historical_migration,
path=path,
)
Loading
Loading