From 21d7283f3db98b2777154773276dc7576d6ca02b Mon Sep 17 00:00:00 2001 From: liuhao1024 Date: Fri, 18 Sep 2026 08:14:56 +0800 Subject: [PATCH] fix(cron): reap the external worker when the waiter returns on terminal ledger The waiter returns as soon as the execution ledger turns terminal, but the worker process may still be in final teardown at that point. The gateway stays the worker's parent, so nobody calling wait() afterwards pins a zombie (STAT=Z) under it until the gateway is restarted. The terminal early-return now hands the process to a short-lived daemon thread whose single job is the final wait(). --- cron/scheduler.py | 18 ++++++++ tests/cron/test_restart_safe_worker.py | 63 ++++++++++++++++++++++++++ 2 files changed, 81 insertions(+) diff --git a/cron/scheduler.py b/cron/scheduler.py index 9bffddc072..f0ed6c93ae 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -3096,6 +3096,23 @@ def _run_one_job_body( reset_terminal_scope(_terminal_scope_token) +def _reap_terminal_worker_in_background(process: subprocess.Popen) -> None: + """Keep the reap contract when the waiter returns before the worker exits. + + The ledger turning terminal lets the waiter return while the worker is + still in final teardown. The gateway remains the worker's parent, so if + nobody calls ``wait()`` afterwards the worker lingers as a zombie (STAT=Z) + under the gateway until it is restarted (#114509). A short-lived daemon + thread holds that single responsibility and ends with the process exit + it waits for. + """ + threading.Thread( + target=process.wait, + name=f"cron-worker-reap-{getattr(process, 'pid', '?')}", + daemon=True, + ).start() + + def _wait_for_external_cron_worker_body( process: subprocess.Popen, *, @@ -3122,6 +3139,7 @@ def _wait_for_external_cron_worker_body( returncode = process.wait(timeout=1.0) except subprocess.TimeoutExpired: if _is_terminal(): + _reap_terminal_worker_in_background(process) return True continue # The worker can commit its terminal row and exit between the first diff --git a/tests/cron/test_restart_safe_worker.py b/tests/cron/test_restart_safe_worker.py index 26998a3901..2af706a525 100644 --- a/tests/cron/test_restart_safe_worker.py +++ b/tests/cron/test_restart_safe_worker.py @@ -481,6 +481,69 @@ def test_external_worker_crash_recovers_uncertain_attempt(monkeypatch): assert get.call_count == 2 +def test_terminal_early_return_still_reaps_the_worker(monkeypatch): + """The ledger can turn terminal while the worker is still tearing down; the + waiter returns then, but the gateway stays the worker's parent, so the exit + must still be waited for somewhere — otherwise the worker lingers as a + zombie under the gateway until it is restarted (#114509).""" + import cron.scheduler as scheduler + + monkeypatch.setattr( + scheduler, + "get_execution", + lambda _execution_id: {"id": "exec-1", "status": "completed"}, + raising=False, + ) + + def wait(timeout=None): + if wait.calls == 0: + wait.calls += 1 + raise subprocess.TimeoutExpired(cmd="worker", timeout=timeout) + return 0 + + wait.calls = 0 + process = Mock() + process.pid = 4321 + process.wait.side_effect = wait + + assert scheduler._wait_for_external_cron_worker_body( + process, execution_id="exec-1" + ) is True + # the background reaper owns the second and final wait(); no third caller appears + deadline = time.monotonic() + 5.0 + while process.wait.call_count < 2 and time.monotonic() < deadline: + time.sleep(0.05) + assert process.wait.call_count == 2 + + +def test_terminal_early_return_reaps_a_real_worker_process(monkeypatch): + """End-to-end zombie guard: after the early return the real worker process + must be reaped without the test itself calling wait()/poll() — reading + ``Popen.returncode`` reaps nothing, so only the background thread can set + it (#114509).""" + import cron.scheduler as scheduler + + monkeypatch.setattr( + scheduler, + "get_execution", + lambda _execution_id: {"id": "exec-1", "status": "completed"}, + raising=False, + ) + process = subprocess.Popen( + [sys.executable, "-c", "import time; time.sleep(1.3)"], + stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, + ) + + assert scheduler._wait_for_external_cron_worker_body( + process, execution_id="exec-1" + ) is True + deadline = time.monotonic() + 8.0 + while process.returncode is None and time.monotonic() < deadline: + time.sleep(0.05) + assert process.returncode == 0 + + def test_launch_external_worker_stays_in_process_outside_managed_gateway( monkeypatch, ):