From b6ae3e4b62c39cb3106c5ca93cf2a2ef25f70430 Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 23 Jul 2026 14:02:54 +0000 Subject: [PATCH 1/5] fix(python): Handle worker signals without interrupts Co-Authored-By: Markus Unterwaditzer --- .../src/taskbroker_client/worker/worker.py | 25 +++---- clients/python/tests/worker/test_worker.py | 69 +++++++++++++++++++ 2 files changed, 80 insertions(+), 14 deletions(-) diff --git a/clients/python/src/taskbroker_client/worker/worker.py b/clients/python/src/taskbroker_client/worker/worker.py index 0cce2164..e8dfd18b 100644 --- a/clients/python/src/taskbroker_client/worker/worker.py +++ b/clients/python/src/taskbroker_client/worker/worker.py @@ -415,16 +415,16 @@ def start(self) -> int: self.worker_pool.start_result_thread() self.worker_pool.start_spawn_children_thread() - # Convert signals into KeyboardInterrupt. - # Running shutdown() within the signal handler can lead to deadlocks + # Signal shutdown without raising an exception while multiprocessing + # synchronization primitives may be held. server: grpc.Server | None = None health_servicer: health.HealthServicer | None = None def signal_handler(*args: Any) -> None: + self._grpc_sync_event.set() if server: server.stop(grace=5) - raise KeyboardInterrupt() signal.signal(signal.SIGINT, signal_handler) signal.signal(signal.SIGTERM, signal_handler) @@ -477,11 +477,7 @@ def signal_handler(*args: Any) -> None: logger.info("taskworker.grpc_server.started", extra={"port": self._grpc_port}) self._start_health_check_thread() - try: - server.wait_for_termination() - except KeyboardInterrupt: - # Signals are converted to KeyboardInterrupt, swallow for exit code 0 - pass + server.wait_for_termination() finally: if health_servicer is not None: @@ -595,20 +591,21 @@ def start(self) -> int: self.worker_pool.start_result_thread() self.worker_pool.start_spawn_children_thread() - # Convert signals into KeyboardInterrupt. - # Running shutdown() within the signal handler can lead to deadlocks + # Signal shutdown without raising an exception while multiprocessing + # synchronization primitives may be held. def signal_handler(*args: Any) -> None: - raise KeyboardInterrupt() + self._grpc_sync_event.set() signal.signal(signal.SIGINT, signal_handler) signal.signal(signal.SIGTERM, signal_handler) try: - while True: + while not self._grpc_sync_event.is_set(): self.run_once() - except KeyboardInterrupt: + finally: self.shutdown() - raise + + return 0 def run_once(self) -> None: """Access point for tests to run a single worker loop""" diff --git a/clients/python/tests/worker/test_worker.py b/clients/python/tests/worker/test_worker.py index 7253a74e..2ffc038a 100644 --- a/clients/python/tests/worker/test_worker.py +++ b/clients/python/tests/worker/test_worker.py @@ -710,6 +710,35 @@ def update_task_response(*args: Any, **kwargs: Any) -> None: assert redis.get("no-retries-remaining"), "key should exist if except block was hit" redis.delete("no-retries-remaining") + def test_start_handles_sigterm_without_raising(self) -> None: + taskworker = TaskWorker( + app_module="examples.app:app", + broker_hosts=["127.0.0.1:50051"], + max_child_task_count=100, + process_type="fork", + ) + handlers: dict[signal.Signals, Callable[..., None]] = {} + + def install_handler(signum: signal.Signals, handler: Callable[..., None]) -> None: + handlers[signum] = handler + + def request_shutdown() -> None: + handlers[signal.SIGTERM]() + + with ( + mock.patch("taskbroker_client.worker.worker.signal.signal", side_effect=install_handler), + mock.patch.object(taskworker.worker_pool, "start_metrics_thread"), + mock.patch.object(taskworker.worker_pool, "start_result_thread"), + mock.patch.object(taskworker.worker_pool, "start_spawn_children_thread"), + mock.patch.object(taskworker.worker_pool, "shutdown") as shutdown, + mock.patch.object(taskworker, "run_once", side_effect=request_shutdown), + ): + exitcode = taskworker.start() + + assert exitcode == 0 + assert taskworker._grpc_sync_event.is_set() + shutdown.assert_called_once_with() + def test_constructor_push_mode(self) -> None: taskworker = PushTaskWorker( app_module="examples.app:app", @@ -831,6 +860,46 @@ def warm_up() -> None: assert timeout_calls == [] +def test_push_worker_handles_sigterm_without_raising() -> None: + taskworker = _make_push_worker(concurrency=2, warmup_timeout=5) + handlers: dict[signal.Signals, Callable[..., None]] = {} + + def install_handler(signum: signal.Signals, handler: Callable[..., None]) -> None: + handlers[signum] = handler + + fake_health = mock.MagicMock() + fake_server = mock.MagicMock() + fake_server.wait_for_termination.side_effect = lambda: handlers[signal.SIGTERM]() + + with ( + mock.patch("taskbroker_client.worker.worker.signal.signal", side_effect=install_handler), + mock.patch.object(taskworker.worker_pool, "start_metrics_thread"), + mock.patch.object(taskworker.worker_pool, "start_result_thread"), + mock.patch.object(taskworker.worker_pool, "start_spawn_children_thread"), + mock.patch.object(taskworker.worker_pool, "shutdown") as shutdown, + mock.patch.object(taskworker, "_start_health_check_thread"), + mock.patch.object(taskworker, "_stop_health_check_thread"), + mock.patch.object( + TaskWorkerProcessingPool, "ready_count", new_callable=mock.PropertyMock, return_value=2 + ), + mock.patch("taskbroker_client.worker.worker.grpc.server", return_value=fake_server), + mock.patch( + "taskbroker_client.worker.worker.health.HealthServicer", return_value=fake_health + ), + mock.patch("taskbroker_client.worker.worker.health_pb2_grpc.add_HealthServicer_to_server"), + mock.patch( + "taskbroker_client.worker.worker.taskbroker_pb2_grpc" + ".add_WorkerServiceServicer_to_server" + ), + ): + exitcode = taskworker.start() + + assert exitcode == 0 + assert taskworker._grpc_sync_event.is_set() + fake_server.stop.assert_called_with(grace=5) + shutdown.assert_called_once_with() + + def test_start_does_not_serve_when_shutdown_during_warmup() -> None: from grpc_health.v1 import health_pb2 From b071382a2c89f13be2440ff88f258b3fb52bbeee Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 23 Jul 2026 14:14:01 +0000 Subject: [PATCH 2/5] fix(python): Stop fetching during shutdown Co-Authored-By: Markus Unterwaditzer --- .../src/taskbroker_client/worker/worker.py | 3 +++ clients/python/tests/worker/test_worker.py | 23 +++++++++++++++++-- 2 files changed, 24 insertions(+), 2 deletions(-) diff --git a/clients/python/src/taskbroker_client/worker/worker.py b/clients/python/src/taskbroker_client/worker/worker.py index e8dfd18b..4e772bbb 100644 --- a/clients/python/src/taskbroker_client/worker/worker.py +++ b/clients/python/src/taskbroker_client/worker/worker.py @@ -688,6 +688,9 @@ def _send_update_task( def fetch_task(self) -> InflightTaskActivation | None: self._grpc_sync_event.wait(self._gettask_backoff_seconds) + if self._grpc_sync_event.is_set(): + return None + try: activation = self.client.get_task(self._namespace) except grpc.RpcError as e: diff --git a/clients/python/tests/worker/test_worker.py b/clients/python/tests/worker/test_worker.py index 2ffc038a..39e29865 100644 --- a/clients/python/tests/worker/test_worker.py +++ b/clients/python/tests/worker/test_worker.py @@ -453,6 +453,21 @@ def test_fetch_task(self) -> None: assert task assert task.activation.id == SIMPLE_TASK.activation.id + def test_fetch_task_skips_request_during_shutdown(self) -> None: + taskworker = TaskWorker( + app_module="examples.app:app", + broker_hosts=["127.0.0.1:50051"], + max_child_task_count=100, + process_type="fork", + ) + taskworker._grpc_sync_event.set() + + with mock.patch.object(taskworker.client, "get_task") as mock_get: + task = taskworker.fetch_task() + + assert task is None + mock_get.assert_not_called() + def test_fetch_no_task(self) -> None: taskworker = TaskWorker( app_module="examples.app:app", @@ -726,7 +741,9 @@ def request_shutdown() -> None: handlers[signal.SIGTERM]() with ( - mock.patch("taskbroker_client.worker.worker.signal.signal", side_effect=install_handler), + mock.patch( + "taskbroker_client.worker.worker.signal.signal", side_effect=install_handler + ), mock.patch.object(taskworker.worker_pool, "start_metrics_thread"), mock.patch.object(taskworker.worker_pool, "start_result_thread"), mock.patch.object(taskworker.worker_pool, "start_spawn_children_thread"), @@ -872,7 +889,9 @@ def install_handler(signum: signal.Signals, handler: Callable[..., None]) -> Non fake_server.wait_for_termination.side_effect = lambda: handlers[signal.SIGTERM]() with ( - mock.patch("taskbroker_client.worker.worker.signal.signal", side_effect=install_handler), + mock.patch( + "taskbroker_client.worker.worker.signal.signal", side_effect=install_handler + ), mock.patch.object(taskworker.worker_pool, "start_metrics_thread"), mock.patch.object(taskworker.worker_pool, "start_result_thread"), mock.patch.object(taskworker.worker_pool, "start_spawn_children_thread"), From 8aacfaa187167277be3a94d51d76f1d6b66fad7f Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 23 Jul 2026 14:17:09 +0000 Subject: [PATCH 3/5] fix(python): Stop push worker startup on shutdown --- .../src/taskbroker_client/worker/worker.py | 11 ++++- clients/python/tests/worker/test_worker.py | 41 ++++++++++++++++++- 2 files changed, 49 insertions(+), 3 deletions(-) diff --git a/clients/python/src/taskbroker_client/worker/worker.py b/clients/python/src/taskbroker_client/worker/worker.py index 4e772bbb..72b28490 100644 --- a/clients/python/src/taskbroker_client/worker/worker.py +++ b/clients/python/src/taskbroker_client/worker/worker.py @@ -419,11 +419,12 @@ def start(self) -> int: # synchronization primitives may be held. server: grpc.Server | None = None + server_started = False health_servicer: health.HealthServicer | None = None def signal_handler(*args: Any) -> None: self._grpc_sync_event.set() - if server: + if server is not None and server_started: server.stop(grace=5) signal.signal(signal.SIGINT, signal_handler) @@ -459,7 +460,13 @@ def signal_handler(*args: Any) -> None: health_servicer.set(WORKER_SERVICE_NAME, health_pb2.HealthCheckResponse.NOT_SERVING) server.add_insecure_port(f"[::]:{self._grpc_port}") + if self._grpc_sync_event.is_set(): + return 0 + server.start() + server_started = True + if self._grpc_sync_event.is_set(): + return 0 # Hold NOT_SERVING until children are warm so the pod stays out of # the NEG/readiness set while its child processes are still loading. @@ -484,7 +491,7 @@ def signal_handler(*args: Any) -> None: health_servicer.set("", health_pb2.HealthCheckResponse.NOT_SERVING) health_servicer.set(WORKER_SERVICE_NAME, health_pb2.HealthCheckResponse.NOT_SERVING) - if server is not None: + if server is not None and server_started: server.stop(grace=5) self._stop_health_check_thread() diff --git a/clients/python/tests/worker/test_worker.py b/clients/python/tests/worker/test_worker.py index 39e29865..38ab5002 100644 --- a/clients/python/tests/worker/test_worker.py +++ b/clients/python/tests/worker/test_worker.py @@ -919,6 +919,44 @@ def install_handler(signum: signal.Signals, handler: Callable[..., None]) -> Non shutdown.assert_called_once_with() +def test_push_worker_does_not_start_server_when_signal_arrives_during_setup() -> None: + taskworker = _make_push_worker(concurrency=2, warmup_timeout=5) + handlers: dict[signal.Signals, Callable[..., None]] = {} + + def install_handler(signum: signal.Signals, handler: Callable[..., None]) -> None: + handlers[signum] = handler + + fake_health = mock.MagicMock() + fake_server = mock.MagicMock() + fake_server.add_insecure_port.side_effect = lambda *args: handlers[signal.SIGTERM]() + + with ( + mock.patch( + "taskbroker_client.worker.worker.signal.signal", side_effect=install_handler + ), + mock.patch.object(taskworker.worker_pool, "start_metrics_thread"), + mock.patch.object(taskworker.worker_pool, "start_result_thread"), + mock.patch.object(taskworker.worker_pool, "start_spawn_children_thread"), + mock.patch.object(taskworker.worker_pool, "shutdown") as shutdown, + mock.patch.object(taskworker, "_stop_health_check_thread"), + mock.patch("taskbroker_client.worker.worker.grpc.server", return_value=fake_server), + mock.patch( + "taskbroker_client.worker.worker.health.HealthServicer", return_value=fake_health + ), + mock.patch("taskbroker_client.worker.worker.health_pb2_grpc.add_HealthServicer_to_server"), + mock.patch( + "taskbroker_client.worker.worker.taskbroker_pb2_grpc" + ".add_WorkerServiceServicer_to_server" + ), + ): + exitcode = taskworker.start() + + assert exitcode == 0 + fake_server.start.assert_not_called() + fake_server.stop.assert_not_called() + shutdown.assert_called_once_with() + + def test_start_does_not_serve_when_shutdown_during_warmup() -> None: from grpc_health.v1 import health_pb2 @@ -954,7 +992,8 @@ def test_start_does_not_serve_when_shutdown_during_warmup() -> None: if c.args[1] == health_pb2.HealthCheckResponse.SERVING ] assert serving_calls == [] - # We never reached server.wait_for_termination() (returned before it). + # We never started the server or reached wait_for_termination(). + fake_server.start.assert_not_called() fake_server.wait_for_termination.assert_not_called() From 55e29eecbd0a16bf8dc022de7f8df44496e079c9 Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 23 Jul 2026 14:23:36 +0000 Subject: [PATCH 4/5] fix(python): Avoid fetching tasks during result drain --- .../src/taskbroker_client/worker/worker.py | 2 ++ clients/python/tests/worker/test_worker.py | 18 ++++++++++++++++++ 2 files changed, 20 insertions(+) diff --git a/clients/python/src/taskbroker_client/worker/worker.py b/clients/python/src/taskbroker_client/worker/worker.py index 72b28490..a46a1310 100644 --- a/clients/python/src/taskbroker_client/worker/worker.py +++ b/clients/python/src/taskbroker_client/worker/worker.py @@ -671,6 +671,8 @@ def _send_update_task( ) self._grpc_sync_event.wait(self._setstatus_backoff_seconds) + if self._grpc_sync_event.is_set(): + fetch_next = None try: next_task = self.client.update_task(result, fetch_next) diff --git a/clients/python/tests/worker/test_worker.py b/clients/python/tests/worker/test_worker.py index 38ab5002..e4d1f725 100644 --- a/clients/python/tests/worker/test_worker.py +++ b/clients/python/tests/worker/test_worker.py @@ -30,6 +30,7 @@ TASK_ACTIVATION_STATUS_COMPLETE, TASK_ACTIVATION_STATUS_FAILURE, TASK_ACTIVATION_STATUS_RETRY, + FetchNextTask, PushTaskRequest, PushTaskResponse, RetryState, @@ -468,6 +469,23 @@ def test_fetch_task_skips_request_during_shutdown(self) -> None: assert task is None mock_get.assert_not_called() + def test_send_update_does_not_fetch_next_during_shutdown(self) -> None: + taskworker = TaskWorker( + app_module="examples.app:app", + broker_hosts=["127.0.0.1:50051"], + max_child_task_count=100, + process_type="fork", + ) + taskworker._grpc_sync_event.set() + result = _make_processing_result("completed") + + with mock.patch.object( + taskworker.client, "update_task", return_value=None + ) as update: + taskworker._send_update_task(result, FetchNextTask(namespace="examples")) + + update.assert_called_once_with(result, None) + def test_fetch_no_task(self) -> None: taskworker = TaskWorker( app_module="examples.app:app", From e6299d618b1c282f411d33b3c0f9465bff6474dc Mon Sep 17 00:00:00 2001 From: "sentry-junior[bot]" <264270552+sentry-junior[bot]@users.noreply.github.com> Date: Thu, 23 Jul 2026 14:30:47 +0000 Subject: [PATCH 5/5] style(python): Format worker tests Co-Authored-By: Markus Unterwaditzer --- clients/python/tests/worker/test_worker.py | 12 +++--------- 1 file changed, 3 insertions(+), 9 deletions(-) diff --git a/clients/python/tests/worker/test_worker.py b/clients/python/tests/worker/test_worker.py index e4d1f725..640dcc1f 100644 --- a/clients/python/tests/worker/test_worker.py +++ b/clients/python/tests/worker/test_worker.py @@ -479,9 +479,7 @@ def test_send_update_does_not_fetch_next_during_shutdown(self) -> None: taskworker._grpc_sync_event.set() result = _make_processing_result("completed") - with mock.patch.object( - taskworker.client, "update_task", return_value=None - ) as update: + with mock.patch.object(taskworker.client, "update_task", return_value=None) as update: taskworker._send_update_task(result, FetchNextTask(namespace="examples")) update.assert_called_once_with(result, None) @@ -907,9 +905,7 @@ def install_handler(signum: signal.Signals, handler: Callable[..., None]) -> Non fake_server.wait_for_termination.side_effect = lambda: handlers[signal.SIGTERM]() with ( - mock.patch( - "taskbroker_client.worker.worker.signal.signal", side_effect=install_handler - ), + mock.patch("taskbroker_client.worker.worker.signal.signal", side_effect=install_handler), mock.patch.object(taskworker.worker_pool, "start_metrics_thread"), mock.patch.object(taskworker.worker_pool, "start_result_thread"), mock.patch.object(taskworker.worker_pool, "start_spawn_children_thread"), @@ -949,9 +945,7 @@ def install_handler(signum: signal.Signals, handler: Callable[..., None]) -> Non fake_server.add_insecure_port.side_effect = lambda *args: handlers[signal.SIGTERM]() with ( - mock.patch( - "taskbroker_client.worker.worker.signal.signal", side_effect=install_handler - ), + mock.patch("taskbroker_client.worker.worker.signal.signal", side_effect=install_handler), mock.patch.object(taskworker.worker_pool, "start_metrics_thread"), mock.patch.object(taskworker.worker_pool, "start_result_thread"), mock.patch.object(taskworker.worker_pool, "start_spawn_children_thread"),