Skip to content
Merged
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
5 changes: 5 additions & 0 deletions .sampo/changesets/gevent-queue-compat.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
pypi/posthog: patch
---

fix: preserve event delivery when gevent monkey-patches `queue.Queue`, including in preloaded gunicorn workers
36 changes: 36 additions & 0 deletions examples/gevent_gunicorn.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
"""PostHog capture from preloaded gunicorn gevent workers.

Setup:
1. Set ``POSTHOG_PROJECT_API_KEY`` and optionally ``POSTHOG_HOST``.
2. Start gunicorn from the repository root::

uv run --with gunicorn --with 'gevent>=25.4.1' \
gunicorn --workers 2 --worker-class gevent --preload \
examples.gevent_gunicorn:app

3. Capture an event with ``curl http://localhost:8000``.

The SDK reinitializes its consumer after gunicorn forks, so no gunicorn
``post_fork`` hook or additional PostHog setup is required.
"""

import os

from posthog import Posthog


posthog = Posthog(
os.environ["POSTHOG_PROJECT_API_KEY"],
host=os.getenv("POSTHOG_HOST", "https://us.i.posthog.com"),
flush_at=1,
)


def app(environ, start_response):
"""Capture one event and return a minimal WSGI response."""
posthog.capture(
"gevent gunicorn example request",
distinct_id=environ.get("REMOTE_ADDR", "unknown"),
)
start_response("200 OK", [("Content-Type", "text/plain; charset=utf-8")])
return [b"Event queued\n"]
27 changes: 27 additions & 0 deletions posthog/_disabled_lane_queue.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
import threading
from queue import Empty, Full


class _DisabledLaneQueue:
"""Minimal empty queue used for lifecycle cleanup on an unavailable lane."""

def __init__(self, maxsize: int) -> None:
self.maxsize = maxsize
self.mutex = threading.Lock()
self.not_empty = threading.Condition(self.mutex)
self.unfinished_tasks = 0

def put(self, item, block: bool = True, timeout=None) -> None:
raise Full

def get_nowait(self):
raise Empty

def qsize(self) -> int:
return 0

def empty(self) -> bool:
return True

def task_done(self) -> None:
return None
91 changes: 83 additions & 8 deletions posthog/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,12 +10,13 @@
import weakref
from contextvars import ContextVar
from datetime import datetime, timedelta, timezone
from typing import Any, Callable, Dict, List, Mapping, Optional, Union
from typing import Any, Callable, Dict, List, Mapping, Optional, Union, cast
from uuid import UUID, uuid4

from typing_extensions import Unpack

from posthog._async_utils import _BackgroundEventLoopRunner
from posthog._disabled_lane_queue import _DisabledLaneQueue
from posthog.args import ID_TYPES, ExceptionArg, OptionalCaptureArgs, OptionalSetArgs
from posthog.metrics_capture import PostHogMetrics
from posthog.capture_compression import (
Expand Down Expand Up @@ -125,6 +126,71 @@
_atexit_deadline_lock = threading.Lock()


def _supports_lane_synchronization(queue) -> bool:
return all(
hasattr(queue, attribute)
for attribute in (
"mutex",
"not_empty",
"not_full",
"all_tasks_done",
"unfinished_tasks",
"_qsize",
"_get",
)
)


def _new_lane_queue(maxsize: int) -> Queue:
"""Return a safe queue, disabling the lane instead of raising on failure."""
log = logging.getLogger("posthog")
try:
queue: Queue = Queue(maxsize)
except Exception:
log.exception(
"Failed to initialize queue.Queue; disabling asynchronous capture for the lane"
)
return cast(Queue, _DisabledLaneQueue(maxsize))

if _supports_lane_synchronization(queue):
return queue

monkey = sys.modules.get("gevent.monkey")
if monkey is None:
log.error(
"queue.Queue lacks the synchronization interface required by PostHog "
"and gevent.monkey is not loaded; disabling asynchronous capture for the lane"
)
return cast(Queue, _DisabledLaneQueue(maxsize))

try:
if not monkey.is_object_patched("queue", "Queue"):
log.error(
"queue.Queue lacks the synchronization interface required by PostHog "
"but gevent does not report it as patched; disabling asynchronous "
"capture for the lane"
)
return cast(Queue, _DisabledLaneQueue(maxsize))

original_queue = monkey.get_original("queue", "Queue")
queue = cast(Queue, original_queue(maxsize))
except Exception:
log.exception(
"Failed to restore the original queue.Queue after gevent monkey-patching; "
"disabling asynchronous capture for the lane"
)
return cast(Queue, _DisabledLaneQueue(maxsize))

if _supports_lane_synchronization(queue):
return queue

log.error(
"The queue.Queue restored after gevent monkey-patching lacks the synchronization "
"interface required by PostHog; disabling asynchronous capture for the lane"
)
return cast(Queue, _DisabledLaneQueue(maxsize))


def _get_atexit_deadline() -> float:
global _atexit_deadline
with _atexit_deadline_lock:
Expand Down Expand Up @@ -327,19 +393,20 @@ def __init__(
self._max_queue_size = max_queue_size
self._thread_count = thread_count
self._eager_start = eager_start
self.queue: Queue = Queue(max_queue_size)
self.queue: Queue = _new_lane_queue(max_queue_size)
self.available = not isinstance(self.queue, _DisabledLaneQueue)
self.consumers: List[Consumer] = []
self._started = False
self._closed = False
self._active_sync_sends = 0
self._start_lock = threading.Lock()
self._sync_sends_done = threading.Condition(self._start_lock)
self._drain_signal = _DrainSignal(self.queue)
if eager_start:
if eager_start and self.available:
self.start()

def _start_locked(self) -> None:
if self._started or self._closed:
if self._started or self._closed or not self.available:
return
for _ in range(self._thread_count):
consumer = Consumer(
Expand Down Expand Up @@ -377,7 +444,7 @@ def start(self):
def enqueue(self, msg) -> bool:
"""Atomically admit and queue `msg`, starting the lane on its first event."""
with self._start_lock:
if self._closed:
if self._closed or not self.available:
return False
self._start_locked()
try:
Expand Down Expand Up @@ -549,13 +616,14 @@ def rebuild_after_fork(self, *, closed: bool) -> None:
the client's fork-visible lifecycle state. An eager open lane restarts
immediately; a lazy lane returns to not-started and restarts on next use.
"""
self.queue = Queue(self._max_queue_size)
self.queue = _new_lane_queue(self._max_queue_size)
self.available = not isinstance(self.queue, _DisabledLaneQueue)
self.reset_sync_send_state_after_fork()
self._drain_signal = _DrainSignal(self.queue)
self.consumers = []
self._started = False
self._closed = closed
if self._eager_start:
if self._eager_start and self.available:
self.start()


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

if lane._closed:
if not lane.available:
self.log.warning(
"%s lane is unavailable because a compatible queue could not be "
"initialized, dropping event %s",
lane.name,
msg["event"],
)
elif lane._closed:
self.log.warning(
"%s lane received event %s after shutdown, dropping it",
lane.name,
Expand Down
152 changes: 152 additions & 0 deletions posthog/test/test_gevent_compat.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,152 @@
"""Regression coverage for gevent monkey-patching compatibility."""

import importlib.util
import subprocess
import sys
import textwrap
import unittest
from queue import Full, Queue
from unittest import mock

from posthog.client import Client, _new_lane_queue
from posthog.test.test_utils import FAKE_TEST_API_KEY


class TestLaneQueueFallback(unittest.TestCase):
def test_uses_working_queue_without_loading_gevent(self):
with mock.patch.dict(sys.modules, {"gevent.monkey": None}):
queue = _new_lane_queue(10)

self.assertIsInstance(queue, Queue)

def test_disables_capture_if_no_compatible_queue_is_available(self):
incompatible_queue = mock.Mock(spec=[])

with self.assertLogs("posthog", level="ERROR") as logs:
with (
mock.patch("posthog.client.Queue", return_value=incompatible_queue),
mock.patch.dict(sys.modules, {"gevent.monkey": None}),
):
queue = _new_lane_queue(10)

with self.assertRaises(Full):
queue.put("event", block=False)
self.assertTrue(queue.empty())
self.assertEqual(queue.unfinished_tasks, 0)
self.assertIsNone(queue.task_done())
self.assertIn("gevent.monkey is not loaded", logs.output[0])
self.assertIn("disabling asynchronous capture for the lane", logs.output[0])

def test_logs_gevent_recovery_failure_before_disabling(self):
incompatible_queue = mock.Mock(spec=[])
monkey = mock.Mock()
monkey.is_object_patched.return_value = True
monkey.get_original.side_effect = RuntimeError("broken gevent state")

with self.assertLogs("posthog", level="ERROR") as logs:
with (
mock.patch("posthog.client.Queue", return_value=incompatible_queue),
mock.patch.dict(sys.modules, {"gevent.monkey": monkey}),
):
queue = _new_lane_queue(10)

self.assertTrue(queue.empty())
self.assertIn("Failed to restore the original queue.Queue", logs.output[0])
self.assertIn("broken gevent state", logs.output[0])

def test_disables_only_async_capture_if_no_compatible_queue_is_available(self):
incompatible_queue = mock.Mock(spec=[])

with self.assertLogs("posthog", level="ERROR"):
with (
mock.patch("posthog.client.Queue", return_value=incompatible_queue),
mock.patch.dict(sys.modules, {"gevent.monkey": None}),
):
client = Client(FAKE_TEST_API_KEY)

self.assertFalse(client.disabled)
self.assertFalse(client._analytics_lane.available)
self.assertFalse(client._ai_lane.available)
self.assertEqual(client.consumers, [])
self.assertIsNone(client.capture("disabled-queue", distinct_id="distinct_id"))
self.assertEqual(client.consumers, [])
client.flush()
client.shutdown()

def test_incompatible_queue_does_not_disable_queue_independent_capabilities(self):
incompatible_queue = mock.Mock(spec=[])

with self.assertLogs("posthog", level="ERROR"):
with (
mock.patch("posthog.client.Queue", return_value=incompatible_queue),
mock.patch.dict(sys.modules, {"gevent.monkey": None}),
mock.patch("posthog.client.batch_post") as mock_post,
mock.patch(
"posthog.client.flags",
return_value={"featureFlags": {"beta-feature": True}},
) as mock_flags,
):
client = Client(FAKE_TEST_API_KEY, sync_mode=True)
event_uuid = client.capture("sync-capture", distinct_id="distinct_id")
decision = client.get_flags_decision("distinct_id")

self.assertFalse(client.disabled)
self.assertIsNotNone(event_uuid)
mock_post.assert_called_once()
self.assertTrue(decision["flags"]["beta-feature"].enabled)
mock_flags.assert_called_once()


@unittest.skipUnless(importlib.util.find_spec("gevent"), "gevent is not installed")
class TestGeventCompatibility(unittest.TestCase):
def test_capture_and_flush_after_monkey_patching(self):
script = textwrap.dedent(
"""
import gevent.monkey

gevent.monkey.patch_all()

import queue

assert gevent.monkey.is_object_patched("queue", "Queue"), (
"gevent did not replace queue.Queue; the regression scenario "
"is not being exercised"
)

from unittest import mock

with mock.patch("posthog.consumer.batch_post") as mock_post:
from posthog.client import Client

client = Client("phc_test", flush_at=1, flush_interval=60)
original_queue = gevent.monkey.get_original("queue", "Queue")
assert isinstance(client.queue, original_queue)
assert not isinstance(client.queue, queue.Queue)
for attribute in (
"mutex",
"not_empty",
"not_full",
"all_tasks_done",
"unfinished_tasks",
"_qsize",
"_get",
):
assert hasattr(client.queue, attribute), attribute

client.capture("gevent-regression", distinct_id="distinct_id")
client.flush(timeout_seconds=10)
client.join()

assert mock_post.called, "batch_post was never called"
assert client.queue.empty(), "flush did not drain the queue"
"""
)

result = subprocess.run(
[sys.executable, "-c", script],
capture_output=True,
text=True,
timeout=60,
)

self.assertEqual(result.returncode, 0, result.stderr)
2 changes: 2 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,8 @@ test = [
"opentelemetry-exporter-otlp-proto-http>=1.20.0",
"pytest-bdd>=8.1.0",
"zstandard>=0.23.0",
# gevent 25.4.1+ replaces queue.Queue, exercising the compatibility path.
"gevent>=25.4.1; implementation_name == 'cpython'",
]

[tool.setuptools]
Expand Down
Loading
Loading