From 8b3059fc6e01e8511a153690f9302f7f705db097 Mon Sep 17 00:00:00 2001 From: KoNit-K <124019182+KoNit-K@users.noreply.github.com> Date: Tue, 15 Sep 2026 01:37:56 +0800 Subject: [PATCH] fix(cron): keep self-removed runs alive across post-removal heartbeats A run that deletes its own job after the first heartbeat interval was still marked stale because the fire-claim loop treated a missing record as lost ownership before the self-removal marker could win. Co-authored-by: Cursor --- cron/scheduler.py | 7 ++-- tests/cron/test_script_claim_heartbeat.py | 39 +++++++++++++++++++++++ 2 files changed, 44 insertions(+), 2 deletions(-) diff --git a/cron/scheduler.py b/cron/scheduler.py index 21f1496ebe..33ff28ecaa 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -2434,6 +2434,9 @@ def _run_with_fire_claim_heartbeat(job: dict, run) -> bool: while not stop.wait(_RUN_CLAIM_HEARTBEAT_SECONDS): try: if not heartbeat_fire_claim(job_id, expected_owner=owner): + if self_removal_delivery_allowed(job_id): + # Record dropped by this run; nothing left to keep fresh. + continue lost_ownership.set() logger.warning( "Job '%s': fire claim ownership lost; interrupting stale run", @@ -2641,12 +2644,12 @@ class _FireOwnership: ) def lost(self) -> bool: + if self_removal_delivery_allowed(self.job["id"]): + return False if self.fire_claim_lost is not None and self.fire_claim_lost.is_set(): return True if self.owner is None: return False - if self_removal_delivery_allowed(self.job["id"]): - return False try: if heartbeat_fire_claim(self.job["id"], expected_owner=self.owner): return False diff --git a/tests/cron/test_script_claim_heartbeat.py b/tests/cron/test_script_claim_heartbeat.py index 79a4bfb779..34219dd6c6 100644 --- a/tests/cron/test_script_claim_heartbeat.py +++ b/tests/cron/test_script_claim_heartbeat.py @@ -454,6 +454,45 @@ def test_self_removed_job_still_delivers_its_completed_response(tmp_path, monkey assert jobs.get_job(job["id"]) is None +def test_self_removed_job_still_delivers_after_post_removal_heartbeat(tmp_path, monkeypatch): + """A run that keeps working past one heartbeat after self-removal must still deliver.""" + import cron.jobs as jobs + import cron.scheduler as scheduler + + def _run_job(job, **_kwargs): + assert jobs.remove_job(job["id"]) is True + time.sleep(0.3) + return True, "saved output", "D1 is promoting", None + + delivered = MagicMock(return_value=None) + finished = MagicMock() + monkeypatch.setattr(scheduler, "_RUN_CLAIM_HEARTBEAT_SECONDS", 0.05) + monkeypatch.setattr(scheduler, "run_job", _run_job) + monkeypatch.setattr(scheduler, "claim_dispatch", lambda *_args: True) + monkeypatch.setattr(scheduler, "mark_execution_running", lambda *_args: {}) + monkeypatch.setattr(scheduler, "finish_execution", finished) + monkeypatch.setattr(scheduler, "save_job_output", lambda *_args: "output.md") + monkeypatch.setattr(scheduler, "_deliver_result", delivered) + + with jobs.use_cron_store(tmp_path): + job = jobs.create_job( + prompt="work", schedule="every 5m", name="remove self", deliver="telegram") + assert jobs.claim_job_for_fire(job["id"]) + claimed = jobs.get_job(job["id"]) + claimed["execution_id"] = "self-removal-heartbeat-execution" + + with patch("agent.secret_scope.set_secret_scope", return_value=None), \ + patch("agent.secret_scope.build_profile_secret_scope", return_value=None), \ + patch("agent.secret_scope.reset_secret_scope"): + assert scheduler.run_one_job(claimed) is True + + delivered.assert_called_once() + finished.assert_called_once_with( + "self-removal-heartbeat-execution", + success=True, error=None, delivery_outcome="delivered") + assert jobs.get_job(job["id"]) is None + + def test_initially_lost_fire_claim_finishes_execution_without_running(monkeypatch): """A stale claimed snapshot rejected before body entry must close its ledger row.""" import cron.scheduler as scheduler