From ec1aa8977c7e5afda74a3ebdda855aa2e72d2f20 Mon Sep 17 00:00:00 2001 From: zhao <1615063567@qq.com> Date: Sat, 15 Aug 2026 09:50:42 +0800 Subject: [PATCH] [verified] fix: clear running-job lock when execution creation fails --- cron/scheduler.py | 11 +++++-- tests/cron/test_parallel_pool.py | 55 ++++++++++++++++++++++++++++++++ 2 files changed, 63 insertions(+), 3 deletions(-) diff --git a/cron/scheduler.py b/cron/scheduler.py index fb801fee28..977e77eb8c 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -6220,14 +6220,19 @@ def tick( execution = create_execution(job_id, source="builtin") dispatched_job = dict(job, execution_id=execution["id"]) _ctx = contextvars.copy_context() - except BaseException: + except Exception as execution_err: # Init/creation failure between the claim and the submit — # release the in-flight claim immediately so the next tick can # retry instead of wedging on 'already running' forever (the # audit requirement: every add is paired with guaranteed - # cleanup). Re-raise so the caller sees the failure. + # cleanup). release_running_job(job_id) - raise + logger.exception( + "Job '%s' not dispatched: execution creation failed: %s", + job.get("name", job_id), + execution_err, + ) + return None def _run_and_release(j=dispatched_job, ctx=_ctx): try: diff --git a/tests/cron/test_parallel_pool.py b/tests/cron/test_parallel_pool.py index 41d3cff547..70c3b3b16d 100644 --- a/tests/cron/test_parallel_pool.py +++ b/tests/cron/test_parallel_pool.py @@ -135,6 +135,61 @@ class TestRunningJobGuard: assert "queued-job" not in sched._running_job_ids + def test_create_execution_failure_does_not_wedge_running_set(self, tmp_path, monkeypatch): + """create_execution failures clear the running lock and still allow next jobs.""" + import cron.scheduler as sched + + sched._parallel_pool = None + sched._parallel_pool_max_workers = None + sched._running_job_ids.clear() + + failing_job = { + "id": "failing-job", + "name": "failing-job", + "prompt": "test", + "schedule": "every 5m", + "enabled": True, + "next_run_at": "2020-01-01T00:00:00", + "deliver": "local", + } + healthy_job = { + "id": "healthy-job", + "name": "healthy-job", + "prompt": "test", + "schedule": "every 5m", + "enabled": True, + "next_run_at": "2020-01-01T00:00:00", + "deliver": "local", + } + + called = [] + + def create_execution_side_effect(job_id, source): + if job_id == "failing-job": + raise RuntimeError("execution ledger unavailable") + return {"id": f"{job_id}-execution"} + + monkeypatch.setattr(sched, "get_due_jobs", lambda: [failing_job, healthy_job]) + monkeypatch.setattr(sched, "advance_next_runs", lambda *_a, **_kw: 0) + monkeypatch.setattr(sched, "create_execution", create_execution_side_effect) + monkeypatch.setattr(sched, "run_job", lambda j, **_kw: called.append(j["id"]) or (True, "out", "resp", None)) + monkeypatch.setattr(sched, "save_job_output", lambda *_a, **_kw: None) + monkeypatch.setattr(sched, "mark_job_run", lambda *_a, **_kw: None) + monkeypatch.setattr(sched, "_deliver_result", lambda *_a, **_kw: None) + monkeypatch.setattr(sched, "finish_execution", lambda *_a, **_kw: None) + monkeypatch.setattr(sched, "claim_dispatch", lambda *_a, **_kw: True) + monkeypatch.setattr(sched, "mark_execution_running", lambda *_a, **_kw: None) + + n = sched.tick(verbose=False) + + assert n == 1 + assert called == ["healthy-job"] + assert "failing-job" not in sched._running_job_ids + assert "healthy-job" not in sched._running_job_ids + + sched._shutdown_parallel_pool() + + class TestSyncMode: """tick() blocks by default (sync=True); tick(sync=False) returns immediately."""