diff --git a/cron/jobs.py b/cron/jobs.py index 7d28d7b529..f2476f78c6 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -137,9 +137,12 @@ def self_removal_delivery_scope(job_id: str): def self_removal_delivery_allowed(job_id: str) -> bool: - """Whether the active run deleted exactly its own job record.""" + """Whether the active run deleted exactly its own job record and no record has since taken + its id (a replacement record belongs to another owner, so that stays fail-closed).""" marker = _self_removal_delivery.get() - return bool(marker is not None and marker.job_id == job_id and marker.removed) + if marker is None or marker.job_id != job_id or not marker.removed: + return False + return all(item.get("id") != job_id for item in load_jobs()) # Import-time snapshot so deliberate re-pointing of CRON_DIR/JOBS_FILE/OUTPUT_DIR (the documented # escape hatch for tests/embedders) is distinguishable from the constants merely being stale. @@ -382,11 +385,9 @@ def _under_fire_fence(job_id: str, fn: Callable[[], Any]) -> Any: @contextlib.contextmanager -def fire_claim_fence(job_id: str, *, expected_owner: str, allow_self_removed: bool = False): - """Hold a per-job fence while an owner performs an external side effect. - - A missing record is accepted only for the active run that removed this exact job. - """ +def fire_claim_fence(job_id: str, *, expected_owner: str): + """Hold a per-job fence while an owner performs an external side effect. A missing record + is accepted only for the active run that removed this exact job (#111039).""" with _fire_job_lock(job_id) as acquired: if not acquired: yield False @@ -395,7 +396,7 @@ def fire_claim_fence(job_id: str, *, expected_owner: str, allow_self_removed: bo job = next((item for item in load_jobs() if item.get("id") == job_id), None) claim = job.get("fire_claim") if isinstance(job, dict) else None owns_claim = isinstance(claim, dict) and claim.get("by") == expected_owner - if not owns_claim and job is None and allow_self_removed: + if job is None: owns_claim = self_removal_delivery_allowed(job_id) yield owns_claim diff --git a/cron/scheduler.py b/cron/scheduler.py index 33ff28ecaa..8c53704e32 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -2638,18 +2638,16 @@ class _FireOwnership: def side_effect_fence(self): if self.owner is None: return contextlib.nullcontext(True) - return fire_claim_fence( - self.job["id"], expected_owner=self.owner, - allow_self_removed=self_removal_delivery_allowed(self.job["id"]), - ) + return fire_claim_fence(self.job["id"], expected_owner=self.owner) 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"]): + # The run deleted its own record; there is no claim left to re-resolve. + return False try: if heartbeat_fire_claim(self.job["id"], expected_owner=self.owner): return False @@ -2796,8 +2794,10 @@ def _finish_completed_run(d: _RunDelivery, fire_owner: Optional[str], execution_ mark_kwargs["expected_fire_owner"] = fire_owner if d.blocked_config: mark_kwargs["status"] = "blocked_config" - marked = mark_job_run(job["id"], d.success, d.error, **mark_kwargs) - if fire_owner is not None and not marked and not self_removal_delivery_allowed(job["id"]): + # A run that removed its own record has nothing left to mark; the delivery above is its result. + marked = self_removal_delivery_allowed(job["id"]) or mark_job_run( + job["id"], d.success, d.error, **mark_kwargs) + if fire_owner is not None and not marked: finish_execution( execution_id, success=False, error="Fire claim ownership lost before terminal completion.") diff --git a/tests/cron/test_script_claim_heartbeat.py b/tests/cron/test_script_claim_heartbeat.py index 34219dd6c6..fe96c67571 100644 --- a/tests/cron/test_script_claim_heartbeat.py +++ b/tests/cron/test_script_claim_heartbeat.py @@ -418,49 +418,14 @@ def test_lost_fire_claim_stops_stale_delivery(monkeypatch): mark_run.assert_not_called() -def test_self_removed_job_still_delivers_its_completed_response(tmp_path, monkeypatch): - """Removing itself during a run releases its record, not this run's final response.""" +def _run_claimed_job_with_mid_run_action(tmp_path, monkeypatch, mid_run, *, execution_id): + """Fire a claimed job through run_one_job with a stubbed agent run that performs ``mid_run`` + on its own record, keeps working past one fire-claim heartbeat tick, then completes.""" import cron.jobs as jobs import cron.scheduler as scheduler def _run_job(job, **_kwargs): - assert jobs.remove_job(job["id"]) - return True, "saved output", "final response", None - - delivered = MagicMock(return_value=None) - finished = MagicMock() - 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-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-execution", success=True, error=None, delivery_outcome="delivered") - 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 + mid_run(jobs, job) time.sleep(0.3) return True, "saved output", "D1 is promoting", None @@ -479,18 +444,50 @@ def test_self_removed_job_still_delivers_after_post_removal_heartbeat(tmp_path, 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" + claimed["execution_id"] = execution_id 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 + return delivered, finished - 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_self_removed_job_still_delivers_after_post_removal_heartbeat(tmp_path, monkeypatch): + """A run that removes its own job (cronjob remove on its own id) and keeps working past a + heartbeat tick must still deliver its final response and complete its ledger row (#111039).""" + import cron.jobs as jobs + + delivered, finished = _run_claimed_job_with_mid_run_action( + tmp_path, monkeypatch, + lambda jobs_mod, job: jobs_mod.remove_job(job["id"]), + execution_id="self-removal-heartbeat-execution") + + delivered.assert_called_once() + finished.assert_called_once_with( + "self-removal-heartbeat-execution", + success=True, error=None, delivery_outcome="delivered") + with jobs.use_cron_store(tmp_path): + assert jobs.load_jobs() == [] + + +def test_self_removal_followed_by_replacement_record_stays_fail_closed(tmp_path, monkeypatch): + """Self-removal only excuses a MISSING record: once another owner's record reclaims the id, + the run is stale again and its result must be discarded, never delivered.""" + + def _remove_then_replace(jobs_mod, job): + assert jobs_mod.remove_job(job["id"]) + replacement = {k: v for k, v in job.items() if k != "execution_id"} + replacement["fire_claim"] = {"at": job["fire_claim"]["at"], "by": "other-machine:owner"} + jobs_mod.save_jobs(jobs_mod.load_jobs() + [replacement]) + + delivered, finished = _run_claimed_job_with_mid_run_action( + tmp_path, monkeypatch, _remove_then_replace, execution_id="replacement-execution") + + delivered.assert_not_called() + finished.assert_called_once_with( + "replacement-execution", success=False, + error="Fire claim ownership lost; stale result was discarded.") def test_initially_lost_fire_claim_finishes_execution_without_running(monkeypatch): diff --git a/website/docs/user-guide/features/cron.md b/website/docs/user-guide/features/cron.md index 04cf78df21..0cf0e84d14 100644 --- a/website/docs/user-guide/features/cron.md +++ b/website/docs/user-guide/features/cron.md @@ -314,6 +314,13 @@ cadence, or run a "cron librarian" job that reconciles the whole table job delivers nowhere). A job created by a scheduled agent can never point its output at a session that no longer exists. Explicit targets (`local`, `all`, `telegram:`) are honored verbatim. +- **A job may remove itself and still report.** The "watch for X, tell me + once, then stop" pattern — a recurring job whose run calls + `cronjob(action="remove", job_id=)` and then answers — delivers + that final response and records the run as `completed`; only the job record + is gone afterwards. Deleting the record from *outside* the run (another + process, or a replacement job reusing the id) still discards the stale + run's output, as before. Prefer prompts that update existing jobs (list first, then update by ID) over ones that create new jobs each run.