Skip to content

Commit 9beed86

Browse files
authored
fix: Preserve event delivery with gevent queues (#867)
* fix: preserve event delivery with gevent queues * fix: disable capture when queues are unavailable * refactor: isolate the disabled lane queue * refactor: reuse the loaded gevent monkey module * fix: explain queue-related client disabling * fix: keep queue failures lane-local
1 parent 9f64299 commit 9beed86

7 files changed

Lines changed: 428 additions & 8 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
pypi/posthog: patch
3+
---
4+
5+
fix: preserve event delivery when gevent monkey-patches `queue.Queue`, including in preloaded gunicorn workers

‎examples/gevent_gunicorn.py‎

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,36 @@
1+
"""PostHog capture from preloaded gunicorn gevent workers.
2+
3+
Setup:
4+
1. Set ``POSTHOG_PROJECT_API_KEY`` and optionally ``POSTHOG_HOST``.
5+
2. Start gunicorn from the repository root::
6+
7+
uv run --with gunicorn --with 'gevent>=25.4.1' \
8+
gunicorn --workers 2 --worker-class gevent --preload \
9+
examples.gevent_gunicorn:app
10+
11+
3. Capture an event with ``curl http://localhost:8000``.
12+
13+
The SDK reinitializes its consumer after gunicorn forks, so no gunicorn
14+
``post_fork`` hook or additional PostHog setup is required.
15+
"""
16+
17+
import os
18+
19+
from posthog import Posthog
20+
21+
22+
posthog = Posthog(
23+
os.environ["POSTHOG_PROJECT_API_KEY"],
24+
host=os.getenv("POSTHOG_HOST", "https://us.i.posthog.com"),
25+
flush_at=1,
26+
)
27+
28+
29+
def app(environ, start_response):
30+
"""Capture one event and return a minimal WSGI response."""
31+
posthog.capture(
32+
"gevent gunicorn example request",
33+
distinct_id=environ.get("REMOTE_ADDR", "unknown"),
34+
)
35+
start_response("200 OK", [("Content-Type", "text/plain; charset=utf-8")])
36+
return [b"Event queued\n"]

‎posthog/_disabled_lane_queue.py‎

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,27 @@
1+
import threading
2+
from queue import Empty, Full
3+
4+
5+
class _DisabledLaneQueue:
6+
"""Minimal empty queue used for lifecycle cleanup on an unavailable lane."""
7+
8+
def __init__(self, maxsize: int) -> None:
9+
self.maxsize = maxsize
10+
self.mutex = threading.Lock()
11+
self.not_empty = threading.Condition(self.mutex)
12+
self.unfinished_tasks = 0
13+
14+
def put(self, item, block: bool = True, timeout=None) -> None:
15+
raise Full
16+
17+
def get_nowait(self):
18+
raise Empty
19+
20+
def qsize(self) -> int:
21+
return 0
22+
23+
def empty(self) -> bool:
24+
return True
25+
26+
def task_done(self) -> None:
27+
return None

‎posthog/client.py‎

Lines changed: 83 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -10,12 +10,13 @@
1010
import weakref
1111
from contextvars import ContextVar
1212
from datetime import datetime, timedelta, timezone
13-
from typing import Any, Callable, Dict, List, Mapping, Optional, Union
13+
from typing import Any, Callable, Dict, List, Mapping, Optional, Union, cast
1414
from uuid import UUID, uuid4
1515

1616
from typing_extensions import Unpack
1717

1818
from posthog._async_utils import _BackgroundEventLoopRunner
19+
from posthog._disabled_lane_queue import _DisabledLaneQueue
1920
from posthog.args import ID_TYPES, ExceptionArg, OptionalCaptureArgs, OptionalSetArgs
2021
from posthog.metrics_capture import PostHogMetrics
2122
from posthog.capture_compression import (
@@ -125,6 +126,71 @@
125126
_atexit_deadline_lock = threading.Lock()
126127

127128

129+
def _supports_lane_synchronization(queue) -> bool:
130+
return all(
131+
hasattr(queue, attribute)
132+
for attribute in (
133+
"mutex",
134+
"not_empty",
135+
"not_full",
136+
"all_tasks_done",
137+
"unfinished_tasks",
138+
"_qsize",
139+
"_get",
140+
)
141+
)
142+
143+
144+
def _new_lane_queue(maxsize: int) -> Queue:
145+
"""Return a safe queue, disabling the lane instead of raising on failure."""
146+
log = logging.getLogger("posthog")
147+
try:
148+
queue: Queue = Queue(maxsize)
149+
except Exception:
150+
log.exception(
151+
"Failed to initialize queue.Queue; disabling asynchronous capture for the lane"
152+
)
153+
return cast(Queue, _DisabledLaneQueue(maxsize))
154+
155+
if _supports_lane_synchronization(queue):
156+
return queue
157+
158+
monkey = sys.modules.get("gevent.monkey")
159+
if monkey is None:
160+
log.error(
161+
"queue.Queue lacks the synchronization interface required by PostHog "
162+
"and gevent.monkey is not loaded; disabling asynchronous capture for the lane"
163+
)
164+
return cast(Queue, _DisabledLaneQueue(maxsize))
165+
166+
try:
167+
if not monkey.is_object_patched("queue", "Queue"):
168+
log.error(
169+
"queue.Queue lacks the synchronization interface required by PostHog "
170+
"but gevent does not report it as patched; disabling asynchronous "
171+
"capture for the lane"
172+
)
173+
return cast(Queue, _DisabledLaneQueue(maxsize))
174+
175+
original_queue = monkey.get_original("queue", "Queue")
176+
queue = cast(Queue, original_queue(maxsize))
177+
except Exception:
178+
log.exception(
179+
"Failed to restore the original queue.Queue after gevent monkey-patching; "
180+
"disabling asynchronous capture for the lane"
181+
)
182+
return cast(Queue, _DisabledLaneQueue(maxsize))
183+
184+
if _supports_lane_synchronization(queue):
185+
return queue
186+
187+
log.error(
188+
"The queue.Queue restored after gevent monkey-patching lacks the synchronization "
189+
"interface required by PostHog; disabling asynchronous capture for the lane"
190+
)
191+
return cast(Queue, _DisabledLaneQueue(maxsize))
192+
193+
128194
def _get_atexit_deadline() -> float:
129195
global _atexit_deadline
130196
with _atexit_deadline_lock:
@@ -327,19 +393,20 @@ def __init__(
327393
self._max_queue_size = max_queue_size
328394
self._thread_count = thread_count
329395
self._eager_start = eager_start
330-
self.queue: Queue = Queue(max_queue_size)
396+
self.queue: Queue = _new_lane_queue(max_queue_size)
397+
self.available = not isinstance(self.queue, _DisabledLaneQueue)
331398
self.consumers: List[Consumer] = []
332399
self._started = False
333400
self._closed = False
334401
self._active_sync_sends = 0
335402
self._start_lock = threading.Lock()
336403
self._sync_sends_done = threading.Condition(self._start_lock)
337404
self._drain_signal = _DrainSignal(self.queue)
338-
if eager_start:
405+
if eager_start and self.available:
339406
self.start()
340407

341408
def _start_locked(self) -> None:
342-
if self._started or self._closed:
409+
if self._started or self._closed or not self.available:
343410
return
344411
for _ in range(self._thread_count):
345412
consumer = Consumer(
@@ -377,7 +444,7 @@ def start(self):
377444
def enqueue(self, msg) -> bool:
378445
"""Atomically admit and queue `msg`, starting the lane on its first event."""
379446
with self._start_lock:
380-
if self._closed:
447+
if self._closed or not self.available:
381448
return False
382449
self._start_locked()
383450
try:
@@ -549,13 +616,14 @@ def rebuild_after_fork(self, *, closed: bool) -> None:
549616
the client's fork-visible lifecycle state. An eager open lane restarts
550617
immediately; a lazy lane returns to not-started and restarts on next use.
551618
"""
552-
self.queue = Queue(self._max_queue_size)
619+
self.queue = _new_lane_queue(self._max_queue_size)
620+
self.available = not isinstance(self.queue, _DisabledLaneQueue)
553621
self.reset_sync_send_state_after_fork()
554622
self._drain_signal = _DrainSignal(self.queue)
555623
self.consumers = []
556624
self._started = False
557625
self._closed = closed
558-
if self._eager_start:
626+
if self._eager_start and self.available:
559627
self.start()
560628

561629

@@ -2250,7 +2318,14 @@ def send_sync() -> None:
22502318
self.log.debug("enqueued %s.", msg["event"])
22512319
return sent_uuid
22522320

2253-
if lane._closed:
2321+
if not lane.available:
2322+
self.log.warning(
2323+
"%s lane is unavailable because a compatible queue could not be "
2324+
"initialized, dropping event %s",
2325+
lane.name,
2326+
msg["event"],
2327+
)
2328+
elif lane._closed:
22542329
self.log.warning(
22552330
"%s lane received event %s after shutdown, dropping it",
22562331
lane.name,

‎posthog/test/test_gevent_compat.py‎

Lines changed: 152 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,152 @@
1+
"""Regression coverage for gevent monkey-patching compatibility."""
2+
3+
import importlib.util
4+
import subprocess
5+
import sys
6+
import textwrap
7+
import unittest
8+
from queue import Full, Queue
9+
from unittest import mock
10+
11+
from posthog.client import Client, _new_lane_queue
12+
from posthog.test.test_utils import FAKE_TEST_API_KEY
13+
14+
15+
class TestLaneQueueFallback(unittest.TestCase):
16+
def test_uses_working_queue_without_loading_gevent(self):
17+
with mock.patch.dict(sys.modules, {"gevent.monkey": None}):
18+
queue = _new_lane_queue(10)
19+
20+
self.assertIsInstance(queue, Queue)
21+
22+
def test_disables_capture_if_no_compatible_queue_is_available(self):
23+
incompatible_queue = mock.Mock(spec=[])
24+
25+
with self.assertLogs("posthog", level="ERROR") as logs:
26+
with (
27+
mock.patch("posthog.client.Queue", return_value=incompatible_queue),
28+
mock.patch.dict(sys.modules, {"gevent.monkey": None}),
29+
):
30+
queue = _new_lane_queue(10)
31+
32+
with self.assertRaises(Full):
33+
queue.put("event", block=False)
34+
self.assertTrue(queue.empty())
35+
self.assertEqual(queue.unfinished_tasks, 0)
36+
self.assertIsNone(queue.task_done())
37+
self.assertIn("gevent.monkey is not loaded", logs.output[0])
38+
self.assertIn("disabling asynchronous capture for the lane", logs.output[0])
39+
40+
def test_logs_gevent_recovery_failure_before_disabling(self):
41+
incompatible_queue = mock.Mock(spec=[])
42+
monkey = mock.Mock()
43+
monkey.is_object_patched.return_value = True
44+
monkey.get_original.side_effect = RuntimeError("broken gevent state")
45+
46+
with self.assertLogs("posthog", level="ERROR") as logs:
47+
with (
48+
mock.patch("posthog.client.Queue", return_value=incompatible_queue),
49+
mock.patch.dict(sys.modules, {"gevent.monkey": monkey}),
50+
):
51+
queue = _new_lane_queue(10)
52+
53+
self.assertTrue(queue.empty())
54+
self.assertIn("Failed to restore the original queue.Queue", logs.output[0])
55+
self.assertIn("broken gevent state", logs.output[0])
56+
57+
def test_disables_only_async_capture_if_no_compatible_queue_is_available(self):
58+
incompatible_queue = mock.Mock(spec=[])
59+
60+
with self.assertLogs("posthog", level="ERROR"):
61+
with (
62+
mock.patch("posthog.client.Queue", return_value=incompatible_queue),
63+
mock.patch.dict(sys.modules, {"gevent.monkey": None}),
64+
):
65+
client = Client(FAKE_TEST_API_KEY)
66+
67+
self.assertFalse(client.disabled)
68+
self.assertFalse(client._analytics_lane.available)
69+
self.assertFalse(client._ai_lane.available)
70+
self.assertEqual(client.consumers, [])
71+
self.assertIsNone(client.capture("disabled-queue", distinct_id="distinct_id"))
72+
self.assertEqual(client.consumers, [])
73+
client.flush()
74+
client.shutdown()
75+
76+
def test_incompatible_queue_does_not_disable_queue_independent_capabilities(self):
77+
incompatible_queue = mock.Mock(spec=[])
78+
79+
with self.assertLogs("posthog", level="ERROR"):
80+
with (
81+
mock.patch("posthog.client.Queue", return_value=incompatible_queue),
82+
mock.patch.dict(sys.modules, {"gevent.monkey": None}),
83+
mock.patch("posthog.client.batch_post") as mock_post,
84+
mock.patch(
85+
"posthog.client.flags",
86+
return_value={"featureFlags": {"beta-feature": True}},
87+
) as mock_flags,
88+
):
89+
client = Client(FAKE_TEST_API_KEY, sync_mode=True)
90+
event_uuid = client.capture("sync-capture", distinct_id="distinct_id")
91+
decision = client.get_flags_decision("distinct_id")
92+
93+
self.assertFalse(client.disabled)
94+
self.assertIsNotNone(event_uuid)
95+
mock_post.assert_called_once()
96+
self.assertTrue(decision["flags"]["beta-feature"].enabled)
97+
mock_flags.assert_called_once()
98+
99+
100+
@unittest.skipUnless(importlib.util.find_spec("gevent"), "gevent is not installed")
101+
class TestGeventCompatibility(unittest.TestCase):
102+
def test_capture_and_flush_after_monkey_patching(self):
103+
script = textwrap.dedent(
104+
"""
105+
import gevent.monkey
106+
107+
gevent.monkey.patch_all()
108+
109+
import queue
110+
111+
assert gevent.monkey.is_object_patched("queue", "Queue"), (
112+
"gevent did not replace queue.Queue; the regression scenario "
113+
"is not being exercised"
114+
)
115+
116+
from unittest import mock
117+
118+
with mock.patch("posthog.consumer.batch_post") as mock_post:
119+
from posthog.client import Client
120+
121+
client = Client("phc_test", flush_at=1, flush_interval=60)
122+
original_queue = gevent.monkey.get_original("queue", "Queue")
123+
assert isinstance(client.queue, original_queue)
124+
assert not isinstance(client.queue, queue.Queue)
125+
for attribute in (
126+
"mutex",
127+
"not_empty",
128+
"not_full",
129+
"all_tasks_done",
130+
"unfinished_tasks",
131+
"_qsize",
132+
"_get",
133+
):
134+
assert hasattr(client.queue, attribute), attribute
135+
136+
client.capture("gevent-regression", distinct_id="distinct_id")
137+
client.flush(timeout_seconds=10)
138+
client.join()
139+
140+
assert mock_post.called, "batch_post was never called"
141+
assert client.queue.empty(), "flush did not drain the queue"
142+
"""
143+
)
144+
145+
result = subprocess.run(
146+
[sys.executable, "-c", script],
147+
capture_output=True,
148+
text=True,
149+
timeout=60,
150+
)
151+
152+
self.assertEqual(result.returncode, 0, result.stderr)

‎pyproject.toml‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -94,6 +94,8 @@ test = [
9494
"opentelemetry-exporter-otlp-proto-http>=1.20.0",
9595
"pytest-bdd>=8.1.0",
9696
"zstandard>=0.23.0",
97+
# gevent 25.4.1+ replaces queue.Queue, exercising the compatibility path.
98+
"gevent>=25.4.1; implementation_name == 'cpython'",
9799
]
98100

99101
[tool.setuptools]

0 commit comments

Comments
 (0)