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().
This commit is contained in:
liuhao1024
2026-09-18 08:14:56 +08:00
committed by Teknium
parent 9244ec1d77
commit 21d7283f3d
2 changed files with 81 additions and 0 deletions

View File

@@ -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

View File

@@ -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,
):