From 1e907343c8d63b8387787f43ec6069d45a349770 Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Tue, 21 Jul 2026 13:12:48 +0200 Subject: [PATCH 01/11] test: Initialize SDK in relevant tests --- clients/python/tests/conftest.py | 37 ++++++ clients/python/tests/test_task.py | 20 ++- clients/python/tests/worker/test_worker.py | 135 +++++++++++++++++---- 3 files changed, 165 insertions(+), 27 deletions(-) diff --git a/clients/python/tests/conftest.py b/clients/python/tests/conftest.py index 2d777a25..a43f7ae0 100644 --- a/clients/python/tests/conftest.py +++ b/clients/python/tests/conftest.py @@ -1,7 +1,14 @@ +from collections.abc import Callable, Generator from datetime import UTC, datetime +from typing import Any +import pytest +import sentry_sdk import time_machine from arroyo.backends.kafka import KafkaProducer +from pytest import FixtureRequest +from sentry_sdk.envelope import Envelope +from sentry_sdk.transport import Transport from taskbroker_client.types import AtMostOnceStore @@ -21,6 +28,36 @@ def freeze_time(t: str | datetime | None = None) -> time_machine.travel: return time_machine.travel(t, tick=False) +@pytest.fixture +def sentry_init(request: FixtureRequest) -> Generator[Callable[..., None], None, None]: + def inner(*a: Any, **kw: Any) -> None: + kw.setdefault("transport", TestTransport()) + client = sentry_sdk.Client(*a, **kw) + sentry_sdk.get_global_scope().set_client(client) + + if request.node.get_closest_marker("forked"): + # Do not run isolation if the test is already running in + # ultimate isolation (seems to be required for celery tests that + # fork) + yield inner + else: + old_client = sentry_sdk.get_global_scope().client + try: + sentry_sdk.get_current_scope().set_client(None) + yield inner + finally: + sentry_sdk.get_global_scope().set_client(old_client) + + +class TestTransport(Transport): + def __init__(self) -> None: + Transport.__init__(self) + + def capture_envelope(self, _: Envelope) -> None: + """No-op capture_envelope for tests""" + pass + + class StubAtMostOnce(AtMostOnceStore): def __init__(self) -> None: self._keys: dict[str, str] = {} diff --git a/clients/python/tests/test_task.py b/clients/python/tests/test_task.py index 2556b9e7..c9e8501e 100644 --- a/clients/python/tests/test_task.py +++ b/clients/python/tests/test_task.py @@ -2,7 +2,7 @@ import datetime from collections.abc import MutableMapping from concurrent.futures import Future -from typing import Any +from typing import Any, Callable from unittest.mock import patch import msgpack @@ -297,7 +297,11 @@ def with_parameters(one: str, two: int, org_id: int) -> None: assert activation.parameters == "" -def test_create_activation_tracing(task_namespace: TaskNamespace) -> None: +def test_create_activation_tracing( + sentry_init: Callable[..., None], task_namespace: TaskNamespace +) -> None: + sentry_init(traces_sample_rate=1.0) + @task_namespace.register(name="test.parameters") def with_parameters(one: str, two: int, org_id: int) -> None: raise NotImplementedError @@ -310,7 +314,11 @@ def with_parameters(one: str, two: int, org_id: int) -> None: assert "baggage" in headers -def test_create_activation_tracing_headers(task_namespace: TaskNamespace) -> None: +def test_create_activation_tracing_headers( + sentry_init: Callable[..., None], task_namespace: TaskNamespace +) -> None: + sentry_init(traces_sample_rate=1.0) + @task_namespace.register(name="test.parameters") def with_parameters(one: str, two: int, org_id: int) -> None: raise NotImplementedError @@ -326,7 +334,11 @@ def with_parameters(one: str, two: int, org_id: int) -> None: assert headers["key"] == "value" -def test_create_activation_tracing_disable(task_namespace: TaskNamespace) -> None: +def test_create_activation_tracing_disable( + sentry_init: Callable[..., None], task_namespace: TaskNamespace +) -> None: + sentry_init(traces_sample_rate=1.0) + @task_namespace.register(name="test.parameters") def with_parameters(one: str, two: int, org_id: int) -> None: raise NotImplementedError diff --git a/clients/python/tests/worker/test_worker.py b/clients/python/tests/worker/test_worker.py index 7253a74e..5769228d 100644 --- a/clients/python/tests/worker/test_worker.py +++ b/clients/python/tests/worker/test_worker.py @@ -1022,7 +1022,11 @@ def test_push_task_worker_busy(self) -> None: @mock.patch("taskbroker_client.worker.workerchild.capture_checkin") -def test_child_process_complete(mock_capture_checkin: mock.MagicMock) -> None: +def test_child_process_complete( + sentry_init: Callable[..., None], mock_capture_checkin: mock.MagicMock +) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1047,7 +1051,11 @@ def test_child_process_complete(mock_capture_checkin: mock.MagicMock) -> None: assert mock_capture_checkin.call_count == 0 -def test_child_process_canary_task(capsys: pytest.CaptureFixture[str]) -> None: +def test_child_process_canary_task( + sentry_init: Callable[..., None], capsys: pytest.CaptureFixture[str] +) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1073,7 +1081,9 @@ def test_child_process_canary_task(capsys: pytest.CaptureFixture[str]) -> None: assert capsys.readouterr().out == "Done running canary task!\n" -def test_child_process_emits_running_message() -> None: +def test_child_process_emits_running_message(sentry_init: Callable[..., None]) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1106,8 +1116,11 @@ def test_child_process_emits_running_message() -> None: @mock.patch("taskbroker_client.worker.workerchild.capture_checkin") def test_child_process_emits_exiting_once_and_continues_until_release( + sentry_init: Callable[..., None], mock_capture_checkin: mock.MagicMock, ) -> None: + sentry_init(traces_sample_rate=1.0) + shutdown = Event() ctx = get_context("fork") child_id = uuid4() @@ -1167,7 +1180,9 @@ def test_child_process_emits_exiting_once_and_continues_until_release( assert mock_capture_checkin.call_count == 0 -def test_child_process_emits_busy_and_idle_messages() -> None: +def test_child_process_emits_busy_and_idle_messages(sentry_init: Callable[..., None]) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1199,7 +1214,9 @@ def test_child_process_emits_busy_and_idle_messages() -> None: assert processed.get(timeout=1).task_id == SIMPLE_TASK.activation.id -def test_child_process_remove_start_time_kwargs() -> None: +def test_child_process_remove_start_time_kwargs(sentry_init: Callable[..., None]) -> None: + sentry_init(traces_sample_rate=1.0) + activation = InflightTaskActivation( host="localhost:50051", receive_timestamp=0, @@ -1236,7 +1253,9 @@ def test_child_process_remove_start_time_kwargs() -> None: assert result.status == TASK_ACTIVATION_STATUS_COMPLETE -def test_child_process_retry_task() -> None: +def test_child_process_retry_task(sentry_init: Callable[..., None]) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1263,8 +1282,10 @@ def test_child_process_retry_task() -> None: @mock.patch("taskbroker_client.worker.workerchild.logger") @mock.patch("taskbroker_client.worker.workerchild.sentry_sdk.capture_exception") def test_child_process_retry_task_max_attempts( - mock_capture: mock.Mock, mock_logger: mock.Mock + sentry_init: Callable[..., None], mock_capture: mock.Mock, mock_logger: mock.Mock ) -> None: + sentry_init(traces_sample_rate=1.0) + # Create an activation that is on its final attempt and # will raise an error again. activation = InflightTaskActivation( @@ -1324,7 +1345,9 @@ def test_child_process_retry_task_max_attempts( assert extra["retry_max_attempts"] == 3 -def test_child_process_failure_task() -> None: +def test_child_process_failure_task(sentry_init: Callable[..., None]) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1348,7 +1371,9 @@ def test_child_process_failure_task() -> None: assert result.status == TASK_ACTIVATION_STATUS_FAILURE -def test_child_process_shutdown() -> None: +def test_child_process_shutdown(sentry_init: Callable[..., None]) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1372,7 +1397,9 @@ def test_child_process_shutdown() -> None: assert processed.qsize() == 0 -def test_child_process_unknown_task() -> None: +def test_child_process_unknown_task(sentry_init: Callable[..., None]) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1400,7 +1427,9 @@ def test_child_process_unknown_task() -> None: assert result.status == TASK_ACTIVATION_STATUS_COMPLETE -def test_child_process_at_most_once() -> None: +def test_child_process_at_most_once(sentry_init: Callable[..., None]) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1431,7 +1460,11 @@ def test_child_process_at_most_once() -> None: @mock.patch("taskbroker_client.worker.workerchild.capture_checkin") -def test_child_process_record_checkin(mock_capture_checkin: mock.Mock) -> None: +def test_child_process_record_checkin( + sentry_init: Callable[..., None], mock_capture_checkin: mock.Mock +) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1463,8 +1496,10 @@ def test_child_process_record_checkin(mock_capture_checkin: mock.Mock) -> None: ) -def test_child_process_pass_headers() -> None: +def test_child_process_pass_headers(sentry_init: Callable[..., None]) -> None: """Task with pass_headers=True receives headers from the activation.""" + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1494,7 +1529,11 @@ def test_child_process_pass_headers() -> None: @mock.patch("taskbroker_client.worker.workerchild.logger") -def test_child_process_terminate_task(mock_logger: mock.Mock) -> None: +def test_child_process_terminate_task( + sentry_init: Callable[..., None], mock_logger: mock.Mock +) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1541,7 +1580,10 @@ def test_child_process_terminate_task(mock_logger: mock.Mock) -> None: @mock.patch("taskbroker_client.worker.workerchild.capture_checkin") -def test_child_process_decompression(mock_capture_checkin: mock.MagicMock) -> None: +def test_child_process_decompression( + sentry_init: Callable[..., None], mock_capture_checkin: mock.MagicMock +) -> None: + sentry_init(traces_sample_rate=1.0) todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() @@ -1567,8 +1609,10 @@ def test_child_process_decompression(mock_capture_checkin: mock.MagicMock) -> No assert mock_capture_checkin.call_count == 0 -def test_child_process_context_hooks() -> None: +def test_child_process_context_hooks(sentry_init: Callable[..., None]) -> None: """Context hooks' on_execute is called with activation headers during task execution.""" + sentry_init(traces_sample_rate=1.0) + executed_headers: list[dict[str, str]] = [] class RecordingHook: @@ -1625,7 +1669,11 @@ def on_execute(self, headers: dict[str, str]) -> contextlib.AbstractContextManag @mock.patch("taskbroker_client.worker.workerchild.logger") -def test_child_process_silenced_timeout(mock_logger: mock.Mock) -> None: +def test_child_process_silenced_timeout( + sentry_init: Callable[..., None], mock_logger: mock.Mock +) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1656,7 +1704,11 @@ def test_child_process_silenced_timeout(mock_logger: mock.Mock) -> None: @mock.patch("taskbroker_client.worker.workerchild.sentry_sdk.capture_exception") -def test_child_process_silenced_exception_with_retries(mock_capture: mock.Mock) -> None: +def test_child_process_silenced_exception_with_retries( + sentry_init: Callable[..., None], mock_capture: mock.Mock +) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1684,7 +1736,11 @@ def test_child_process_silenced_exception_with_retries(mock_capture: mock.Mock) @mock.patch("taskbroker_client.worker.workerchild.sentry_sdk.capture_exception") -def test_child_process_expected_ignored_exception_max_attempts(mock_capture: mock.Mock) -> None: +def test_child_process_expected_ignored_exception_max_attempts( + sentry_init: Callable[..., None], mock_capture: mock.Mock +) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1714,9 +1770,11 @@ def test_child_process_expected_ignored_exception_max_attempts(mock_capture: moc @mock.patch("taskbroker_client.worker.workerchild.logger") @mock.patch("taskbroker_client.worker.workerchild.sentry_sdk.capture_exception") def test_child_process_silenced_exception_max_attempts( - mock_capture: mock.Mock, mock_logger: mock.Mock + sentry_init: Callable[..., None], mock_capture: mock.Mock, mock_logger: mock.Mock ) -> None: """Silenced exceptions do not raise on retry exhaustion.""" + sentry_init(traces_sample_rate=1.0) + activation = InflightTaskActivation( host="localhost:50051", receive_timestamp=0, @@ -1768,7 +1826,11 @@ def test_child_process_silenced_exception_max_attempts( @mock.patch("taskbroker_client.worker.workerchild.logger") -def test_child_process_retry_on_deadline_exceeded(mock_logger: mock.Mock) -> None: +def test_child_process_retry_on_deadline_exceeded( + sentry_init: Callable[..., None], mock_logger: mock.Mock +) -> None: + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1801,8 +1863,12 @@ def test_child_process_retry_on_deadline_exceeded(mock_logger: mock.Mock) -> Non @mock.patch("taskbroker_client.worker.workerchild.logger") -def test_child_process_general_exception_logs_task_failed(mock_logger: mock.Mock) -> None: +def test_child_process_general_exception_logs_task_failed( + sentry_init: Callable[..., None], mock_logger: mock.Mock +) -> None: """A non-retriable Exception emits taskworker.task.failed with all fields.""" + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1838,10 +1904,13 @@ def test_child_process_general_exception_logs_task_failed(mock_logger: mock.Mock @mock.patch("taskbroker_client.worker.workerchild.logger") def test_child_process_silenced_exception_does_not_log_task_failed( + sentry_init: Callable[..., None], mock_logger: mock.Mock, ) -> None: """When err is in silenced_exceptions, taskworker.task.failed is NOT logged. Preserves the silencing semantics added in #608.""" + sentry_init(traces_sample_rate=1.0) + todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1936,10 +2005,13 @@ def _producing_task(task_id: str = "task-with-futures") -> InflightTaskActivatio @pytest.mark.parametrize("producer_cls", _PRODUCER_CLASSES) def test_child_process_tracks_producer_futures( + sentry_init: Callable[..., None], producer_cls: type, clear_pending_futures: None, restore_signal_handlers: None, ) -> None: + sentry_init(traces_sample_rate=1.0) + task = _producing_task() todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() @@ -1974,10 +2046,13 @@ def test_child_process_tracks_producer_futures( @pytest.mark.parametrize("producer_cls", _PRODUCER_CLASSES) def test_child_process_holds_result_until_futures_done( + sentry_init: Callable[..., None], producer_cls: type, clear_pending_futures: None, restore_signal_handlers: None, ) -> None: + sentry_init(traces_sample_rate=1.0) + task = _producing_task() todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() @@ -2029,10 +2104,13 @@ def observe_and_resolve() -> None: @pytest.mark.parametrize("producer_cls", _PRODUCER_CLASSES) def test_child_process_skip_awaiting_futures_places_result_immediately( + sentry_init: Callable[..., None], producer_cls: type, clear_pending_futures: None, restore_signal_handlers: None, ) -> None: + sentry_init(traces_sample_rate=1.0) + task = _producing_task() todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() @@ -2091,10 +2169,13 @@ def observe_and_resolve() -> None: @pytest.mark.parametrize("producer_cls", _PRODUCER_CLASSES) def test_child_process_drains_pending_futures_on_sigterm( + sentry_init: Callable[..., None], producer_cls: type, clear_pending_futures: None, restore_signal_handlers: None, ) -> None: + sentry_init(traces_sample_rate=1.0) + task = _producing_task() todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() @@ -2139,10 +2220,13 @@ def deliver_sigterm() -> None: @pytest.mark.parametrize("producer_cls", _PRODUCER_CLASSES) def test_child_process_retries_on_failed_future( + sentry_init: Callable[..., None], producer_cls: type, clear_pending_futures: None, restore_signal_handlers: None, ) -> None: + sentry_init(traces_sample_rate=1.0) + retriable_task = InflightTaskActivation( host="localhost:50051", receive_timestamp=0, @@ -2189,10 +2273,13 @@ def test_child_process_retries_on_failed_future( @pytest.mark.parametrize("pending_registry", _PENDING_REGISTRIES) def test_child_process_clears_pending_futures_when_task_fails( + sentry_init: Callable[..., None], pending_registry: Any, clear_pending_futures: None, restore_signal_handlers: None, ) -> None: + sentry_init(traces_sample_rate=1.0) + leftover_future: Future[BrokerValue[KafkaPayload]] = Future() leftover_future.set_result(_make_broker_value()) pending_registry["test.producer"].append(leftover_future) @@ -2226,9 +2313,11 @@ def test_child_process_clears_pending_futures_when_task_fails( def test_child_process_uses_configured_future_checking_frequency( - clear_pending_futures: None, restore_signal_handlers: None + sentry_init: Callable[..., None], clear_pending_futures: None, restore_signal_handlers: None ) -> None: """The idle future-checking loop polls on the configured interval.""" + sentry_init(traces_sample_rate=1.0) + # A task that runs long enough for the idle future-checking loop to poll a # few times before max_task_count triggers shutdown. slow_task = InflightTaskActivation( From 5de3e0dd777d5b8b5a843277ac8951ffc57ed95a Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Tue, 21 Jul 2026 13:13:33 +0200 Subject: [PATCH 02/11] remove dead code --- clients/python/tests/conftest.py | 16 +++++----------- 1 file changed, 5 insertions(+), 11 deletions(-) diff --git a/clients/python/tests/conftest.py b/clients/python/tests/conftest.py index a43f7ae0..a4ad61a8 100644 --- a/clients/python/tests/conftest.py +++ b/clients/python/tests/conftest.py @@ -35,18 +35,12 @@ def inner(*a: Any, **kw: Any) -> None: client = sentry_sdk.Client(*a, **kw) sentry_sdk.get_global_scope().set_client(client) - if request.node.get_closest_marker("forked"): - # Do not run isolation if the test is already running in - # ultimate isolation (seems to be required for celery tests that - # fork) + old_client = sentry_sdk.get_global_scope().client + try: + sentry_sdk.get_current_scope().set_client(None) yield inner - else: - old_client = sentry_sdk.get_global_scope().client - try: - sentry_sdk.get_current_scope().set_client(None) - yield inner - finally: - sentry_sdk.get_global_scope().set_client(old_client) + finally: + sentry_sdk.get_global_scope().set_client(old_client) class TestTransport(Transport): From d570be1cdb1411dde2c19064e4cf75f8decc54fa Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Tue, 21 Jul 2026 13:26:26 +0200 Subject: [PATCH 03/11] ref(o11y): Remove redundant span status assignment --- .../taskbroker_client/worker/workerchild.py | 45 +++++++++---------- 1 file changed, 20 insertions(+), 25 deletions(-) diff --git a/clients/python/src/taskbroker_client/worker/workerchild.py b/clients/python/src/taskbroker_client/worker/workerchild.py index d8f206e1..726b4722 100644 --- a/clients/python/src/taskbroker_client/worker/workerchild.py +++ b/clients/python/src/taskbroker_client/worker/workerchild.py @@ -30,7 +30,7 @@ TaskActivation, TaskActivationStatus, ) -from sentry_sdk.consts import OP, SPANDATA, SPANSTATUS +from sentry_sdk.consts import OP, SPANDATA from sentry_sdk.crons import MonitorStatus, capture_checkin from taskbroker_client.app import import_app @@ -658,30 +658,25 @@ def _execute_activation( if "__start_time" in kwargs: kwargs.pop("__start_time") - try: - with contextlib.ExitStack() as stack: - with metrics.timer( - "taskworker.worker.context_rebuild.duration", - tags={ - "namespace": activation.namespace, - "taskname": activation.taskname, - }, - ): - for hook in context_hooks: - stack.enter_context(hook.on_execute(headers)) - if task_func.pass_headers: - if "headers" in kwargs: - raise TypeError( - f"Task '{task_func.name}' has pass_headers=True, but 'headers' was passed in kwargs. " - "The 'headers' parameter is injected by the worker and cannot be passed by the caller." - ) - task_func(*args, headers=headers, **kwargs) - else: - task_func(*args, **kwargs) - transaction.set_status(SPANSTATUS.OK) - except Exception: - transaction.set_status(SPANSTATUS.INTERNAL_ERROR) - raise + with contextlib.ExitStack() as stack: + with metrics.timer( + "taskworker.worker.context_rebuild.duration", + tags={ + "namespace": activation.namespace, + "taskname": activation.taskname, + }, + ): + for hook in context_hooks: + stack.enter_context(hook.on_execute(headers)) + if task_func.pass_headers: + if "headers" in kwargs: + raise TypeError( + f"Task '{task_func.name}' has pass_headers=True, but 'headers' was passed in kwargs. " + "The 'headers' parameter is injected by the worker and cannot be passed by the caller." + ) + task_func(*args, headers=headers, **kwargs) + else: + task_func(*args, **kwargs) def record_task_execution( activation: TaskActivation, From 49883c800345552eea6ab397abb624e24e07bfe1 Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Tue, 21 Jul 2026 14:03:50 +0200 Subject: [PATCH 04/11] move fixtures after mocks --- clients/python/tests/worker/test_worker.py | 39 ++++++++++++++-------- 1 file changed, 26 insertions(+), 13 deletions(-) diff --git a/clients/python/tests/worker/test_worker.py b/clients/python/tests/worker/test_worker.py index 5769228d..ad8677d3 100644 --- a/clients/python/tests/worker/test_worker.py +++ b/clients/python/tests/worker/test_worker.py @@ -1023,7 +1023,8 @@ def test_push_task_worker_busy(self) -> None: @mock.patch("taskbroker_client.worker.workerchild.capture_checkin") def test_child_process_complete( - sentry_init: Callable[..., None], mock_capture_checkin: mock.MagicMock + mock_capture_checkin: mock.MagicMock, + sentry_init: Callable[..., None], ) -> None: sentry_init(traces_sample_rate=1.0) @@ -1116,8 +1117,8 @@ def test_child_process_emits_running_message(sentry_init: Callable[..., None]) - @mock.patch("taskbroker_client.worker.workerchild.capture_checkin") def test_child_process_emits_exiting_once_and_continues_until_release( - sentry_init: Callable[..., None], mock_capture_checkin: mock.MagicMock, + sentry_init: Callable[..., None], ) -> None: sentry_init(traces_sample_rate=1.0) @@ -1282,7 +1283,9 @@ def test_child_process_retry_task(sentry_init: Callable[..., None]) -> None: @mock.patch("taskbroker_client.worker.workerchild.logger") @mock.patch("taskbroker_client.worker.workerchild.sentry_sdk.capture_exception") def test_child_process_retry_task_max_attempts( - sentry_init: Callable[..., None], mock_capture: mock.Mock, mock_logger: mock.Mock + mock_capture: mock.Mock, + mock_logger: mock.Mock, + sentry_init: Callable[..., None], ) -> None: sentry_init(traces_sample_rate=1.0) @@ -1461,7 +1464,8 @@ def test_child_process_at_most_once(sentry_init: Callable[..., None]) -> None: @mock.patch("taskbroker_client.worker.workerchild.capture_checkin") def test_child_process_record_checkin( - sentry_init: Callable[..., None], mock_capture_checkin: mock.Mock + mock_capture_checkin: mock.Mock, + sentry_init: Callable[..., None], ) -> None: sentry_init(traces_sample_rate=1.0) @@ -1530,7 +1534,8 @@ def test_child_process_pass_headers(sentry_init: Callable[..., None]) -> None: @mock.patch("taskbroker_client.worker.workerchild.logger") def test_child_process_terminate_task( - sentry_init: Callable[..., None], mock_logger: mock.Mock + mock_logger: mock.Mock, + sentry_init: Callable[..., None], ) -> None: sentry_init(traces_sample_rate=1.0) @@ -1581,7 +1586,8 @@ def test_child_process_terminate_task( @mock.patch("taskbroker_client.worker.workerchild.capture_checkin") def test_child_process_decompression( - sentry_init: Callable[..., None], mock_capture_checkin: mock.MagicMock + mock_capture_checkin: mock.MagicMock, + sentry_init: Callable[..., None], ) -> None: sentry_init(traces_sample_rate=1.0) @@ -1670,7 +1676,8 @@ def on_execute(self, headers: dict[str, str]) -> contextlib.AbstractContextManag @mock.patch("taskbroker_client.worker.workerchild.logger") def test_child_process_silenced_timeout( - sentry_init: Callable[..., None], mock_logger: mock.Mock + mock_logger: mock.Mock, + sentry_init: Callable[..., None], ) -> None: sentry_init(traces_sample_rate=1.0) @@ -1705,7 +1712,8 @@ def test_child_process_silenced_timeout( @mock.patch("taskbroker_client.worker.workerchild.sentry_sdk.capture_exception") def test_child_process_silenced_exception_with_retries( - sentry_init: Callable[..., None], mock_capture: mock.Mock + mock_capture: mock.Mock, + sentry_init: Callable[..., None], ) -> None: sentry_init(traces_sample_rate=1.0) @@ -1737,7 +1745,8 @@ def test_child_process_silenced_exception_with_retries( @mock.patch("taskbroker_client.worker.workerchild.sentry_sdk.capture_exception") def test_child_process_expected_ignored_exception_max_attempts( - sentry_init: Callable[..., None], mock_capture: mock.Mock + mock_capture: mock.Mock, + sentry_init: Callable[..., None], ) -> None: sentry_init(traces_sample_rate=1.0) @@ -1770,7 +1779,9 @@ def test_child_process_expected_ignored_exception_max_attempts( @mock.patch("taskbroker_client.worker.workerchild.logger") @mock.patch("taskbroker_client.worker.workerchild.sentry_sdk.capture_exception") def test_child_process_silenced_exception_max_attempts( - sentry_init: Callable[..., None], mock_capture: mock.Mock, mock_logger: mock.Mock + mock_capture: mock.Mock, + mock_logger: mock.Mock, + sentry_init: Callable[..., None], ) -> None: """Silenced exceptions do not raise on retry exhaustion.""" sentry_init(traces_sample_rate=1.0) @@ -1827,7 +1838,8 @@ def test_child_process_silenced_exception_max_attempts( @mock.patch("taskbroker_client.worker.workerchild.logger") def test_child_process_retry_on_deadline_exceeded( - sentry_init: Callable[..., None], mock_logger: mock.Mock + mock_logger: mock.Mock, + sentry_init: Callable[..., None], ) -> None: sentry_init(traces_sample_rate=1.0) @@ -1864,7 +1876,8 @@ def test_child_process_retry_on_deadline_exceeded( @mock.patch("taskbroker_client.worker.workerchild.logger") def test_child_process_general_exception_logs_task_failed( - sentry_init: Callable[..., None], mock_logger: mock.Mock + mock_logger: mock.Mock, + sentry_init: Callable[..., None], ) -> None: """A non-retriable Exception emits taskworker.task.failed with all fields.""" sentry_init(traces_sample_rate=1.0) @@ -1904,8 +1917,8 @@ def test_child_process_general_exception_logs_task_failed( @mock.patch("taskbroker_client.worker.workerchild.logger") def test_child_process_silenced_exception_does_not_log_task_failed( - sentry_init: Callable[..., None], mock_logger: mock.Mock, + sentry_init: Callable[..., None], ) -> None: """When err is in silenced_exceptions, taskworker.task.failed is NOT logged. Preserves the silencing semantics added in #608.""" From d6ea80956c3a6e45313cc4dbd2aacc3aef6e70a6 Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Tue, 21 Jul 2026 14:10:52 +0200 Subject: [PATCH 05/11] remove sentry_init from frequency test --- clients/python/tests/worker/test_worker.py | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/clients/python/tests/worker/test_worker.py b/clients/python/tests/worker/test_worker.py index ad8677d3..5e8ea8ea 100644 --- a/clients/python/tests/worker/test_worker.py +++ b/clients/python/tests/worker/test_worker.py @@ -2326,11 +2326,9 @@ def test_child_process_clears_pending_futures_when_task_fails( def test_child_process_uses_configured_future_checking_frequency( - sentry_init: Callable[..., None], clear_pending_futures: None, restore_signal_handlers: None + clear_pending_futures: None, restore_signal_handlers: None ) -> None: """The idle future-checking loop polls on the configured interval.""" - sentry_init(traces_sample_rate=1.0) - # A task that runs long enough for the idle future-checking loop to poll a # few times before max_task_count triggers shutdown. slow_task = InflightTaskActivation( From 9c64fa8567d907d6b05a279ed147b4373c7422d9 Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Tue, 21 Jul 2026 14:32:32 +0200 Subject: [PATCH 06/11] configure sdk to avoid sleeps --- clients/python/tests/worker/test_worker.py | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/clients/python/tests/worker/test_worker.py b/clients/python/tests/worker/test_worker.py index 5e8ea8ea..4fa0b5bc 100644 --- a/clients/python/tests/worker/test_worker.py +++ b/clients/python/tests/worker/test_worker.py @@ -2326,9 +2326,14 @@ def test_child_process_clears_pending_futures_when_task_fails( def test_child_process_uses_configured_future_checking_frequency( - clear_pending_futures: None, restore_signal_handlers: None + sentry_init: Callable[..., None], clear_pending_futures: None, restore_signal_handlers: None ) -> None: """The idle future-checking loop polls on the configured interval.""" + sentry_init( + traces_sample_rate=1.0, + enable_backpressure_handling=False, # To avoid time.sleep which the test patches. + ) + # A task that runs long enough for the idle future-checking loop to poll a # few times before max_task_count triggers shutdown. slow_task = InflightTaskActivation( From ea6dba32c9d737b7b96275486f8cf8fe460f8adc Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Tue, 21 Jul 2026 14:54:35 +0200 Subject: [PATCH 07/11] fix import ordering problem --- clients/python/tests/worker/test_worker.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/clients/python/tests/worker/test_worker.py b/clients/python/tests/worker/test_worker.py index 4fa0b5bc..40b4f13d 100644 --- a/clients/python/tests/worker/test_worker.py +++ b/clients/python/tests/worker/test_worker.py @@ -2360,6 +2360,8 @@ def recording_sleep(seconds: float) -> None: idle_sleeps.append(seconds) real_sleep(seconds) + import examples.tasks # noqa: F401 — ensure sleep ref is bound before patching + # time.sleep is only used by the idle branch of check_task_future_completion # inside workerchild, so every recorded call comes from that loop. The task's # own sleep uses a separate `from time import sleep` import in examples.tasks. From 32fd19dcc7d8e60bd8c6bc6912f847b5394c4533 Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Tue, 21 Jul 2026 15:00:15 +0200 Subject: [PATCH 08/11] cleanup in fixture --- clients/python/tests/conftest.py | 3 +++ clients/python/tests/worker/test_worker.py | 2 +- 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/clients/python/tests/conftest.py b/clients/python/tests/conftest.py index a4ad61a8..8c945d73 100644 --- a/clients/python/tests/conftest.py +++ b/clients/python/tests/conftest.py @@ -40,6 +40,9 @@ def inner(*a: Any, **kw: Any) -> None: sentry_sdk.get_current_scope().set_client(None) yield inner finally: + current = sentry_sdk.get_global_scope().client + if current is not None: + current.close() sentry_sdk.get_global_scope().set_client(old_client) diff --git a/clients/python/tests/worker/test_worker.py b/clients/python/tests/worker/test_worker.py index 40b4f13d..858ad939 100644 --- a/clients/python/tests/worker/test_worker.py +++ b/clients/python/tests/worker/test_worker.py @@ -2360,7 +2360,7 @@ def recording_sleep(seconds: float) -> None: idle_sleeps.append(seconds) real_sleep(seconds) - import examples.tasks # noqa: F401 — ensure sleep ref is bound before patching + import examples.tasks # noqa: F401; Ensure time.sleep reference is set before patching. # time.sleep is only used by the idle branch of check_task_future_completion # inside workerchild, so every recorded call comes from that loop. The task's From 265e1266df6aace2d591aa8bc207eebcfd8423fd Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Tue, 21 Jul 2026 15:16:48 +0200 Subject: [PATCH 09/11] remove unused parameter --- clients/python/tests/conftest.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/clients/python/tests/conftest.py b/clients/python/tests/conftest.py index 8c945d73..db436d8e 100644 --- a/clients/python/tests/conftest.py +++ b/clients/python/tests/conftest.py @@ -6,7 +6,6 @@ import sentry_sdk import time_machine from arroyo.backends.kafka import KafkaProducer -from pytest import FixtureRequest from sentry_sdk.envelope import Envelope from sentry_sdk.transport import Transport @@ -29,7 +28,7 @@ def freeze_time(t: str | datetime | None = None) -> time_machine.travel: @pytest.fixture -def sentry_init(request: FixtureRequest) -> Generator[Callable[..., None], None, None]: +def sentry_init() -> Generator[Callable[..., None], None, None]: def inner(*a: Any, **kw: Any) -> None: kw.setdefault("transport", TestTransport()) client = sentry_sdk.Client(*a, **kw) From 2b52bbcfb0f38988b0d34f9a1357dc9751899bed Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Tue, 21 Jul 2026 15:22:40 +0200 Subject: [PATCH 10/11] simplify transport --- clients/python/tests/conftest.py | 3 --- 1 file changed, 3 deletions(-) diff --git a/clients/python/tests/conftest.py b/clients/python/tests/conftest.py index db436d8e..7877d095 100644 --- a/clients/python/tests/conftest.py +++ b/clients/python/tests/conftest.py @@ -46,9 +46,6 @@ def inner(*a: Any, **kw: Any) -> None: class TestTransport(Transport): - def __init__(self) -> None: - Transport.__init__(self) - def capture_envelope(self, _: Envelope) -> None: """No-op capture_envelope for tests""" pass From 35831f8eee61b7e909b57d95f37bd35cc74f0c11 Mon Sep 17 00:00:00 2001 From: Alexander Alderman Webb Date: Wed, 22 Jul 2026 16:33:09 +0200 Subject: [PATCH 11/11] revert test changes --- clients/python/tests/conftest.py | 30 ---- clients/python/tests/test_task.py | 20 +-- clients/python/tests/worker/test_worker.py | 153 ++++----------------- 3 files changed, 27 insertions(+), 176 deletions(-) diff --git a/clients/python/tests/conftest.py b/clients/python/tests/conftest.py index 7877d095..2d777a25 100644 --- a/clients/python/tests/conftest.py +++ b/clients/python/tests/conftest.py @@ -1,13 +1,7 @@ -from collections.abc import Callable, Generator from datetime import UTC, datetime -from typing import Any -import pytest -import sentry_sdk import time_machine from arroyo.backends.kafka import KafkaProducer -from sentry_sdk.envelope import Envelope -from sentry_sdk.transport import Transport from taskbroker_client.types import AtMostOnceStore @@ -27,30 +21,6 @@ def freeze_time(t: str | datetime | None = None) -> time_machine.travel: return time_machine.travel(t, tick=False) -@pytest.fixture -def sentry_init() -> Generator[Callable[..., None], None, None]: - def inner(*a: Any, **kw: Any) -> None: - kw.setdefault("transport", TestTransport()) - client = sentry_sdk.Client(*a, **kw) - sentry_sdk.get_global_scope().set_client(client) - - old_client = sentry_sdk.get_global_scope().client - try: - sentry_sdk.get_current_scope().set_client(None) - yield inner - finally: - current = sentry_sdk.get_global_scope().client - if current is not None: - current.close() - sentry_sdk.get_global_scope().set_client(old_client) - - -class TestTransport(Transport): - def capture_envelope(self, _: Envelope) -> None: - """No-op capture_envelope for tests""" - pass - - class StubAtMostOnce(AtMostOnceStore): def __init__(self) -> None: self._keys: dict[str, str] = {} diff --git a/clients/python/tests/test_task.py b/clients/python/tests/test_task.py index c9e8501e..2556b9e7 100644 --- a/clients/python/tests/test_task.py +++ b/clients/python/tests/test_task.py @@ -2,7 +2,7 @@ import datetime from collections.abc import MutableMapping from concurrent.futures import Future -from typing import Any, Callable +from typing import Any from unittest.mock import patch import msgpack @@ -297,11 +297,7 @@ def with_parameters(one: str, two: int, org_id: int) -> None: assert activation.parameters == "" -def test_create_activation_tracing( - sentry_init: Callable[..., None], task_namespace: TaskNamespace -) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_create_activation_tracing(task_namespace: TaskNamespace) -> None: @task_namespace.register(name="test.parameters") def with_parameters(one: str, two: int, org_id: int) -> None: raise NotImplementedError @@ -314,11 +310,7 @@ def with_parameters(one: str, two: int, org_id: int) -> None: assert "baggage" in headers -def test_create_activation_tracing_headers( - sentry_init: Callable[..., None], task_namespace: TaskNamespace -) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_create_activation_tracing_headers(task_namespace: TaskNamespace) -> None: @task_namespace.register(name="test.parameters") def with_parameters(one: str, two: int, org_id: int) -> None: raise NotImplementedError @@ -334,11 +326,7 @@ def with_parameters(one: str, two: int, org_id: int) -> None: assert headers["key"] == "value" -def test_create_activation_tracing_disable( - sentry_init: Callable[..., None], task_namespace: TaskNamespace -) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_create_activation_tracing_disable(task_namespace: TaskNamespace) -> None: @task_namespace.register(name="test.parameters") def with_parameters(one: str, two: int, org_id: int) -> None: raise NotImplementedError diff --git a/clients/python/tests/worker/test_worker.py b/clients/python/tests/worker/test_worker.py index 858ad939..7253a74e 100644 --- a/clients/python/tests/worker/test_worker.py +++ b/clients/python/tests/worker/test_worker.py @@ -1022,12 +1022,7 @@ def test_push_task_worker_busy(self) -> None: @mock.patch("taskbroker_client.worker.workerchild.capture_checkin") -def test_child_process_complete( - mock_capture_checkin: mock.MagicMock, - sentry_init: Callable[..., None], -) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_complete(mock_capture_checkin: mock.MagicMock) -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1052,11 +1047,7 @@ def test_child_process_complete( assert mock_capture_checkin.call_count == 0 -def test_child_process_canary_task( - sentry_init: Callable[..., None], capsys: pytest.CaptureFixture[str] -) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_canary_task(capsys: pytest.CaptureFixture[str]) -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1082,9 +1073,7 @@ def test_child_process_canary_task( assert capsys.readouterr().out == "Done running canary task!\n" -def test_child_process_emits_running_message(sentry_init: Callable[..., None]) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_emits_running_message() -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1118,10 +1107,7 @@ def test_child_process_emits_running_message(sentry_init: Callable[..., None]) - @mock.patch("taskbroker_client.worker.workerchild.capture_checkin") def test_child_process_emits_exiting_once_and_continues_until_release( mock_capture_checkin: mock.MagicMock, - sentry_init: Callable[..., None], ) -> None: - sentry_init(traces_sample_rate=1.0) - shutdown = Event() ctx = get_context("fork") child_id = uuid4() @@ -1181,9 +1167,7 @@ def test_child_process_emits_exiting_once_and_continues_until_release( assert mock_capture_checkin.call_count == 0 -def test_child_process_emits_busy_and_idle_messages(sentry_init: Callable[..., None]) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_emits_busy_and_idle_messages() -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1215,9 +1199,7 @@ def test_child_process_emits_busy_and_idle_messages(sentry_init: Callable[..., N assert processed.get(timeout=1).task_id == SIMPLE_TASK.activation.id -def test_child_process_remove_start_time_kwargs(sentry_init: Callable[..., None]) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_remove_start_time_kwargs() -> None: activation = InflightTaskActivation( host="localhost:50051", receive_timestamp=0, @@ -1254,9 +1236,7 @@ def test_child_process_remove_start_time_kwargs(sentry_init: Callable[..., None] assert result.status == TASK_ACTIVATION_STATUS_COMPLETE -def test_child_process_retry_task(sentry_init: Callable[..., None]) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_retry_task() -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1283,12 +1263,8 @@ def test_child_process_retry_task(sentry_init: Callable[..., None]) -> None: @mock.patch("taskbroker_client.worker.workerchild.logger") @mock.patch("taskbroker_client.worker.workerchild.sentry_sdk.capture_exception") def test_child_process_retry_task_max_attempts( - mock_capture: mock.Mock, - mock_logger: mock.Mock, - sentry_init: Callable[..., None], + mock_capture: mock.Mock, mock_logger: mock.Mock ) -> None: - sentry_init(traces_sample_rate=1.0) - # Create an activation that is on its final attempt and # will raise an error again. activation = InflightTaskActivation( @@ -1348,9 +1324,7 @@ def test_child_process_retry_task_max_attempts( assert extra["retry_max_attempts"] == 3 -def test_child_process_failure_task(sentry_init: Callable[..., None]) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_failure_task() -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1374,9 +1348,7 @@ def test_child_process_failure_task(sentry_init: Callable[..., None]) -> None: assert result.status == TASK_ACTIVATION_STATUS_FAILURE -def test_child_process_shutdown(sentry_init: Callable[..., None]) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_shutdown() -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1400,9 +1372,7 @@ def test_child_process_shutdown(sentry_init: Callable[..., None]) -> None: assert processed.qsize() == 0 -def test_child_process_unknown_task(sentry_init: Callable[..., None]) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_unknown_task() -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1430,9 +1400,7 @@ def test_child_process_unknown_task(sentry_init: Callable[..., None]) -> None: assert result.status == TASK_ACTIVATION_STATUS_COMPLETE -def test_child_process_at_most_once(sentry_init: Callable[..., None]) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_at_most_once() -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1463,12 +1431,7 @@ def test_child_process_at_most_once(sentry_init: Callable[..., None]) -> None: @mock.patch("taskbroker_client.worker.workerchild.capture_checkin") -def test_child_process_record_checkin( - mock_capture_checkin: mock.Mock, - sentry_init: Callable[..., None], -) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_record_checkin(mock_capture_checkin: mock.Mock) -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1500,10 +1463,8 @@ def test_child_process_record_checkin( ) -def test_child_process_pass_headers(sentry_init: Callable[..., None]) -> None: +def test_child_process_pass_headers() -> None: """Task with pass_headers=True receives headers from the activation.""" - sentry_init(traces_sample_rate=1.0) - todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1533,12 +1494,7 @@ def test_child_process_pass_headers(sentry_init: Callable[..., None]) -> None: @mock.patch("taskbroker_client.worker.workerchild.logger") -def test_child_process_terminate_task( - mock_logger: mock.Mock, - sentry_init: Callable[..., None], -) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_terminate_task(mock_logger: mock.Mock) -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1585,11 +1541,7 @@ def test_child_process_terminate_task( @mock.patch("taskbroker_client.worker.workerchild.capture_checkin") -def test_child_process_decompression( - mock_capture_checkin: mock.MagicMock, - sentry_init: Callable[..., None], -) -> None: - sentry_init(traces_sample_rate=1.0) +def test_child_process_decompression(mock_capture_checkin: mock.MagicMock) -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() @@ -1615,10 +1567,8 @@ def test_child_process_decompression( assert mock_capture_checkin.call_count == 0 -def test_child_process_context_hooks(sentry_init: Callable[..., None]) -> None: +def test_child_process_context_hooks() -> None: """Context hooks' on_execute is called with activation headers during task execution.""" - sentry_init(traces_sample_rate=1.0) - executed_headers: list[dict[str, str]] = [] class RecordingHook: @@ -1675,12 +1625,7 @@ def on_execute(self, headers: dict[str, str]) -> contextlib.AbstractContextManag @mock.patch("taskbroker_client.worker.workerchild.logger") -def test_child_process_silenced_timeout( - mock_logger: mock.Mock, - sentry_init: Callable[..., None], -) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_silenced_timeout(mock_logger: mock.Mock) -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1711,12 +1656,7 @@ def test_child_process_silenced_timeout( @mock.patch("taskbroker_client.worker.workerchild.sentry_sdk.capture_exception") -def test_child_process_silenced_exception_with_retries( - mock_capture: mock.Mock, - sentry_init: Callable[..., None], -) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_silenced_exception_with_retries(mock_capture: mock.Mock) -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1744,12 +1684,7 @@ def test_child_process_silenced_exception_with_retries( @mock.patch("taskbroker_client.worker.workerchild.sentry_sdk.capture_exception") -def test_child_process_expected_ignored_exception_max_attempts( - mock_capture: mock.Mock, - sentry_init: Callable[..., None], -) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_expected_ignored_exception_max_attempts(mock_capture: mock.Mock) -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1779,13 +1714,9 @@ def test_child_process_expected_ignored_exception_max_attempts( @mock.patch("taskbroker_client.worker.workerchild.logger") @mock.patch("taskbroker_client.worker.workerchild.sentry_sdk.capture_exception") def test_child_process_silenced_exception_max_attempts( - mock_capture: mock.Mock, - mock_logger: mock.Mock, - sentry_init: Callable[..., None], + mock_capture: mock.Mock, mock_logger: mock.Mock ) -> None: """Silenced exceptions do not raise on retry exhaustion.""" - sentry_init(traces_sample_rate=1.0) - activation = InflightTaskActivation( host="localhost:50051", receive_timestamp=0, @@ -1837,12 +1768,7 @@ def test_child_process_silenced_exception_max_attempts( @mock.patch("taskbroker_client.worker.workerchild.logger") -def test_child_process_retry_on_deadline_exceeded( - mock_logger: mock.Mock, - sentry_init: Callable[..., None], -) -> None: - sentry_init(traces_sample_rate=1.0) - +def test_child_process_retry_on_deadline_exceeded(mock_logger: mock.Mock) -> None: todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1875,13 +1801,8 @@ def test_child_process_retry_on_deadline_exceeded( @mock.patch("taskbroker_client.worker.workerchild.logger") -def test_child_process_general_exception_logs_task_failed( - mock_logger: mock.Mock, - sentry_init: Callable[..., None], -) -> None: +def test_child_process_general_exception_logs_task_failed(mock_logger: mock.Mock) -> None: """A non-retriable Exception emits taskworker.task.failed with all fields.""" - sentry_init(traces_sample_rate=1.0) - todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -1918,12 +1839,9 @@ def test_child_process_general_exception_logs_task_failed( @mock.patch("taskbroker_client.worker.workerchild.logger") def test_child_process_silenced_exception_does_not_log_task_failed( mock_logger: mock.Mock, - sentry_init: Callable[..., None], ) -> None: """When err is in silenced_exceptions, taskworker.task.failed is NOT logged. Preserves the silencing semantics added in #608.""" - sentry_init(traces_sample_rate=1.0) - todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() shutdown = Event() @@ -2018,13 +1936,10 @@ def _producing_task(task_id: str = "task-with-futures") -> InflightTaskActivatio @pytest.mark.parametrize("producer_cls", _PRODUCER_CLASSES) def test_child_process_tracks_producer_futures( - sentry_init: Callable[..., None], producer_cls: type, clear_pending_futures: None, restore_signal_handlers: None, ) -> None: - sentry_init(traces_sample_rate=1.0) - task = _producing_task() todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() @@ -2059,13 +1974,10 @@ def test_child_process_tracks_producer_futures( @pytest.mark.parametrize("producer_cls", _PRODUCER_CLASSES) def test_child_process_holds_result_until_futures_done( - sentry_init: Callable[..., None], producer_cls: type, clear_pending_futures: None, restore_signal_handlers: None, ) -> None: - sentry_init(traces_sample_rate=1.0) - task = _producing_task() todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() @@ -2117,13 +2029,10 @@ def observe_and_resolve() -> None: @pytest.mark.parametrize("producer_cls", _PRODUCER_CLASSES) def test_child_process_skip_awaiting_futures_places_result_immediately( - sentry_init: Callable[..., None], producer_cls: type, clear_pending_futures: None, restore_signal_handlers: None, ) -> None: - sentry_init(traces_sample_rate=1.0) - task = _producing_task() todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() @@ -2182,13 +2091,10 @@ def observe_and_resolve() -> None: @pytest.mark.parametrize("producer_cls", _PRODUCER_CLASSES) def test_child_process_drains_pending_futures_on_sigterm( - sentry_init: Callable[..., None], producer_cls: type, clear_pending_futures: None, restore_signal_handlers: None, ) -> None: - sentry_init(traces_sample_rate=1.0) - task = _producing_task() todo: queue.Queue[InflightTaskActivation] = queue.Queue() processed: queue.Queue[ProcessingResult] = queue.Queue() @@ -2233,13 +2139,10 @@ def deliver_sigterm() -> None: @pytest.mark.parametrize("producer_cls", _PRODUCER_CLASSES) def test_child_process_retries_on_failed_future( - sentry_init: Callable[..., None], producer_cls: type, clear_pending_futures: None, restore_signal_handlers: None, ) -> None: - sentry_init(traces_sample_rate=1.0) - retriable_task = InflightTaskActivation( host="localhost:50051", receive_timestamp=0, @@ -2286,13 +2189,10 @@ def test_child_process_retries_on_failed_future( @pytest.mark.parametrize("pending_registry", _PENDING_REGISTRIES) def test_child_process_clears_pending_futures_when_task_fails( - sentry_init: Callable[..., None], pending_registry: Any, clear_pending_futures: None, restore_signal_handlers: None, ) -> None: - sentry_init(traces_sample_rate=1.0) - leftover_future: Future[BrokerValue[KafkaPayload]] = Future() leftover_future.set_result(_make_broker_value()) pending_registry["test.producer"].append(leftover_future) @@ -2326,14 +2226,9 @@ def test_child_process_clears_pending_futures_when_task_fails( def test_child_process_uses_configured_future_checking_frequency( - sentry_init: Callable[..., None], clear_pending_futures: None, restore_signal_handlers: None + clear_pending_futures: None, restore_signal_handlers: None ) -> None: """The idle future-checking loop polls on the configured interval.""" - sentry_init( - traces_sample_rate=1.0, - enable_backpressure_handling=False, # To avoid time.sleep which the test patches. - ) - # A task that runs long enough for the idle future-checking loop to poll a # few times before max_task_count triggers shutdown. slow_task = InflightTaskActivation( @@ -2360,8 +2255,6 @@ def recording_sleep(seconds: float) -> None: idle_sleeps.append(seconds) real_sleep(seconds) - import examples.tasks # noqa: F401; Ensure time.sleep reference is set before patching. - # time.sleep is only used by the idle branch of check_task_future_completion # inside workerchild, so every recorded call comes from that loop. The task's # own sleep uses a separate `from time import sleep` import in examples.tasks.