Skip to content

Commit d5eb0bc

Browse files
committed
feat: AsyncPosthog AI capture lane
1 parent fde9c61 commit d5eb0bc

5 files changed

Lines changed: 296 additions & 86 deletions

File tree

‎posthog/_async_consumer.py‎

Lines changed: 43 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -96,9 +96,11 @@ def __init__(
9696
flush_at: int,
9797
flush_interval: float,
9898
retries: int,
99-
timeout: int,
99+
timeout: float,
100100
historical_migration: bool,
101101
capture_compression: CaptureCompression,
102+
endpoint: str = _CAPTURE_V1_PATH,
103+
max_msg_size: int = MAX_MSG_SIZE,
102104
) -> None:
103105
self.queue = queue
104106
self.api_key = api_key
@@ -111,6 +113,8 @@ def __init__(
111113
self.timeout = timeout
112114
self.historical_migration = historical_migration
113115
self.capture_compression = capture_compression
116+
self.endpoint = endpoint
117+
self.max_msg_size = max_msg_size
114118
self._carryover: Optional[tuple[dict[str, Any], int]] = None
115119
self._flush_event = asyncio.Event()
116120

@@ -171,7 +175,7 @@ async def upload(self, batch: list[dict[str, Any]]) -> None:
171175
await self.request(batch)
172176
except Exception as error:
173177
await _report_capture_failure(
174-
self.on_error, self.log, error, batch, _CAPTURE_V1_PATH
178+
self.on_error, self.log, error, batch, self.endpoint
175179
)
176180
finally:
177181
for _ in batch:
@@ -229,12 +233,15 @@ async def next(self) -> tuple[list[dict[str, Any]], bool]:
229233
self.queue.task_done()
230234
continue
231235

232-
if item_size > MAX_MSG_SIZE:
236+
if item_size > self.max_msg_size:
237+
# Log only name and size: AI events may carry unredacted
238+
# multimodal payloads that must not leak into logs.
233239
self.log.error(
234-
"Event %s (%d bytes) exceeds the %dKiB limit, dropping.",
240+
"Event %s (%d bytes) exceeds the %dKiB limit for %s, dropping.",
235241
item.get("event"),
236242
item_size,
237-
MAX_MSG_SIZE // 1024,
243+
self.max_msg_size // 1024,
244+
self.endpoint,
238245
)
239246
self.queue.task_done()
240247
continue
@@ -258,4 +265,35 @@ async def request(self, batch: list[dict[str, Any]]) -> None:
258265
timeout=self.timeout,
259266
max_retries=self.retries,
260267
historical_migration=self.historical_migration,
268+
path=self.endpoint,
261269
)
270+
271+
272+
class _AsyncLane:
273+
"""One capture queue, the consumer tasks that drain it, and the endpoint they post to.
274+
275+
The client owns one lane per traffic class (analytics, AI), so each gets
276+
its own backpressure, timeout, size cap and compression.
277+
"""
278+
279+
def __init__(
280+
self,
281+
*,
282+
name: str,
283+
max_queue_size: int,
284+
endpoint: str,
285+
max_msg_size: int,
286+
timeout: float,
287+
capture_compression: CaptureCompression,
288+
) -> None:
289+
self.name = name
290+
self.queue: asyncio.Queue[Any] = asyncio.Queue(max_queue_size)
291+
self.endpoint = endpoint
292+
self.max_msg_size = max_msg_size
293+
self.timeout = timeout
294+
self.capture_compression = capture_compression
295+
self.consumers: list[_AsyncConsumer] = []
296+
self.worker_tasks: list[asyncio.Task[None]] = []
297+
298+
def pending_items(self) -> int:
299+
return int(getattr(self.queue, "_unfinished_tasks", self.queue.qsize()))

‎posthog/_async_request.py‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@
77
from urllib.parse import quote
88

99
from .capture_compression import CaptureCompression
10-
from .capture_send import _parse_retry_after, _send_v1_batch
10+
from .capture_send import _CAPTURE_V1_PATH, _parse_retry_after, _send_v1_batch
1111
from .request import (
1212
APIError,
1313
DatetimeSerializer,
@@ -150,9 +150,10 @@ async def async_send_v1_batch(
150150
batch: list[dict[str, Any]],
151151
*,
152152
compression: CaptureCompression,
153-
timeout: int,
153+
timeout: float,
154154
max_retries: int,
155155
historical_migration: bool,
156+
path: str = _CAPTURE_V1_PATH,
156157
) -> None:
157158
"""Run the existing capture-v1 submitter off-loop to preserve wire parity."""
158159
await asyncio.to_thread(
@@ -164,4 +165,5 @@ async def async_send_v1_batch(
164165
timeout=timeout,
165166
max_retries=max_retries,
166167
historical_migration=historical_migration,
168+
path=path,
167169
)

0 commit comments

Comments
 (0)