Skip to content

Commit 4b1c7f6

Browse files
committed
fix: keep queue failures lane-local
1 parent 13dde49 commit 4b1c7f6

3 files changed

Lines changed: 49 additions & 17 deletions

File tree

‎posthog/_disabled_lane_queue.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33

44

55
class _DisabledLaneQueue:
6-
"""Minimal empty queue used for lifecycle cleanup on a disabled client."""
6+
"""Minimal empty queue used for lifecycle cleanup on an unavailable lane."""
77

88
def __init__(self, maxsize: int) -> None:
99
self.maxsize = maxsize

‎posthog/client.py‎

Lines changed: 19 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -142,12 +142,14 @@ def _supports_lane_synchronization(queue) -> bool:
142142

143143

144144
def _new_lane_queue(maxsize: int) -> Queue:
145-
"""Return a safe queue, disabling async capture instead of raising on failure."""
145+
"""Return a safe queue, disabling the lane instead of raising on failure."""
146146
log = logging.getLogger("posthog")
147147
try:
148148
queue: Queue = Queue(maxsize)
149149
except Exception:
150-
log.exception("Failed to initialize queue.Queue; disabling the PostHog client")
150+
log.exception(
151+
"Failed to initialize queue.Queue; disabling asynchronous capture for the lane"
152+
)
151153
return cast(Queue, _DisabledLaneQueue(maxsize))
152154

153155
if _supports_lane_synchronization(queue):
@@ -157,15 +159,16 @@ def _new_lane_queue(maxsize: int) -> Queue:
157159
if monkey is None:
158160
log.error(
159161
"queue.Queue lacks the synchronization interface required by PostHog "
160-
"and gevent.monkey is not loaded; disabling the PostHog client"
162+
"and gevent.monkey is not loaded; disabling asynchronous capture for the lane"
161163
)
162164
return cast(Queue, _DisabledLaneQueue(maxsize))
163165

164166
try:
165167
if not monkey.is_object_patched("queue", "Queue"):
166168
log.error(
167169
"queue.Queue lacks the synchronization interface required by PostHog "
168-
"but gevent does not report it as patched; disabling the PostHog client"
170+
"but gevent does not report it as patched; disabling asynchronous "
171+
"capture for the lane"
169172
)
170173
return cast(Queue, _DisabledLaneQueue(maxsize))
171174

@@ -174,7 +177,7 @@ def _new_lane_queue(maxsize: int) -> Queue:
174177
except Exception:
175178
log.exception(
176179
"Failed to restore the original queue.Queue after gevent monkey-patching; "
177-
"disabling the PostHog client"
180+
"disabling asynchronous capture for the lane"
178181
)
179182
return cast(Queue, _DisabledLaneQueue(maxsize))
180183

@@ -183,7 +186,7 @@ def _new_lane_queue(maxsize: int) -> Queue:
183186

184187
log.error(
185188
"The queue.Queue restored after gevent monkey-patching lacks the synchronization "
186-
"interface required by PostHog; disabling the PostHog client"
189+
"interface required by PostHog; disabling asynchronous capture for the lane"
187190
)
188191
return cast(Queue, _DisabledLaneQueue(maxsize))
189192

@@ -403,7 +406,7 @@ def __init__(
403406
self.start()
404407

405408
def _start_locked(self) -> None:
406-
if self._started or self._closed:
409+
if self._started or self._closed or not self.available:
407410
return
408411
for _ in range(self._thread_count):
409412
consumer = Consumer(
@@ -441,7 +444,7 @@ def start(self):
441444
def enqueue(self, msg) -> bool:
442445
"""Atomically admit and queue `msg`, starting the lane on its first event."""
443446
with self._start_lock:
444-
if self._closed:
447+
if self._closed or not self.available:
445448
return False
446449
self._start_locked()
447450
try:
@@ -1036,8 +1039,6 @@ def __init__(
10361039
eager_start=False,
10371040
)
10381041
self._lanes = [self._analytics_lane, self._ai_lane]
1039-
if not all(lane.available for lane in self._lanes):
1040-
self.disabled = True
10411042

10421043
if hasattr(os, "register_at_fork"):
10431044
weak_self = weakref.ref(self)
@@ -2137,8 +2138,6 @@ def _reinit_after_fork(self):
21372138
)
21382139
for lane in self._lanes:
21392140
lane.rebuild_after_fork(closed=terminal_requested)
2140-
if not all(lane.available for lane in self._lanes):
2141-
self.disabled = True
21422141

21432142
self._lifecycle_lock = threading.Lock()
21442143
self._lifecycle_condition = threading.Condition(self._lifecycle_lock)
@@ -2319,7 +2318,14 @@ def send_sync() -> None:
23192318
self.log.debug("enqueued %s.", msg["event"])
23202319
return sent_uuid
23212320

2322-
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:
23232329
self.log.warning(
23242330
"%s lane received event %s after shutdown, dropping it",
23252331
lane.name,

‎posthog/test/test_gevent_compat.py‎

Lines changed: 29 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -35,7 +35,7 @@ def test_disables_capture_if_no_compatible_queue_is_available(self):
3535
self.assertEqual(queue.unfinished_tasks, 0)
3636
self.assertIsNone(queue.task_done())
3737
self.assertIn("gevent.monkey is not loaded", logs.output[0])
38-
self.assertIn("disabling the PostHog client", logs.output[0])
38+
self.assertIn("disabling asynchronous capture for the lane", logs.output[0])
3939

4040
def test_logs_gevent_recovery_failure_before_disabling(self):
4141
incompatible_queue = mock.Mock(spec=[])
@@ -54,7 +54,7 @@ def test_logs_gevent_recovery_failure_before_disabling(self):
5454
self.assertIn("Failed to restore the original queue.Queue", logs.output[0])
5555
self.assertIn("broken gevent state", logs.output[0])
5656

57-
def test_disables_client_if_no_compatible_queue_is_available(self):
57+
def test_disables_only_async_capture_if_no_compatible_queue_is_available(self):
5858
incompatible_queue = mock.Mock(spec=[])
5959

6060
with self.assertLogs("posthog", level="ERROR"):
@@ -64,12 +64,38 @@ def test_disables_client_if_no_compatible_queue_is_available(self):
6464
):
6565
client = Client(FAKE_TEST_API_KEY)
6666

67-
self.assertTrue(client.disabled)
67+
self.assertFalse(client.disabled)
68+
self.assertFalse(client._analytics_lane.available)
69+
self.assertFalse(client._ai_lane.available)
6870
self.assertEqual(client.consumers, [])
6971
self.assertIsNone(client.capture("disabled-queue", distinct_id="distinct_id"))
72+
self.assertEqual(client.consumers, [])
7073
client.flush()
7174
client.shutdown()
7275

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+
7399

74100
@unittest.skipUnless(importlib.util.find_spec("gevent"), "gevent is not installed")
75101
class TestGeventCompatibility(unittest.TestCase):

0 commit comments

Comments
 (0)