diff --git a/cron/scheduler.py b/cron/scheduler.py index 977e77eb8c..3e6d22bd14 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -5581,6 +5581,8 @@ def _run_one_job_body( execution_id = job.get("execution_id") if not execution_id: execution_id = create_execution(job["id"], source="direct")["id"] + delivery_attempted = False + delivery_error = None try: # Pre-run dispatch claim (issue #38758): atomically commit a finite # one-shot's dispatch BEFORE its side effect runs, so a tick that dies @@ -5692,7 +5694,6 @@ def _run_one_job_body( # / empty-response computation, or _deliver_result itself — raises, the # deferred agent is still torn down. Otherwise the outer `except` would # swallow the error and leak the agent's subprocesses/clients (#10200). - delivery_error = None blocked_config = False side_effect_ownership_lost = False try: @@ -5803,6 +5804,7 @@ def _run_one_job_body( with _side_effect_fence() as owns_delivery: if not owns_delivery: raise _FireClaimLostDuringSideEffect + delivery_attempted = True delivery_error = _deliver_result( job, deliver_content, @@ -5925,11 +5927,49 @@ def _run_one_job_body( # a stale worker must not record over a replacement claim owner. _err_text = str(e) or type(e).__name__ logger.error("Error processing job %s: %s", job['id'], _err_text) + delivery_outcome = "suppressed" + # Owner fencing: a stale worker whose fire claim was taken over (or a + # transport-cancelled worker) must not send a failure alert on top of + # the replacement run's own delivery — fall through silently and let + # the fenced bookkeeping below decide what (if anything) to record. + if ( + isinstance(e, Exception) + and not delivery_attempted + and not isinstance(e, _FireClaimLostDuringSideEffect) + and not _fire_claim_ownership_lost() + ): + normalized_deliver = _normalize_deliver_value( + job.get("deliver", "local") + ) + unresolved_origin = False + try: + delivery_attempted = True + delivery_error = _deliver_result( + job, + _summarize_cron_failure_for_delivery(job, _err_text), + adapters=adapters, + loop=loop, + ) + except Exception as delivery_exc: + delivery_error = str(delivery_exc) + logger.error( + "Delivery failed for job %s: %s", job["id"], delivery_exc + ) + if not delivery_error and normalized_deliver == "origin": + unresolved_origin = not _resolve_delivery_targets(job) + if delivery_error: + delivery_outcome = "failed" + elif unresolved_origin: + delivery_outcome = "not_configured" + elif normalized_deliver != "local": + delivery_outcome = "delivered" try: if not _consume_interrupted_flag(job["id"], execution_token): mark_kwargs = {} if fire_owner is not None: mark_kwargs["expected_fire_owner"] = fire_owner + if isinstance(e, Exception): + mark_kwargs["delivery_error"] = delivery_error mark_job_run(job["id"], False, _err_text, **mark_kwargs) except Exception as record_err: # Never let bookkeeping mask the original interruption. @@ -5938,7 +5978,12 @@ def _run_one_job_body( job["id"], record_err, ) try: - finish_execution(execution_id, success=False, error=_err_text) + finish_execution( + execution_id, + success=False, + error=_err_text, + delivery_outcome=delivery_outcome, + ) except Exception as record_err: logger.error( "Failed to finish execution record for job %s: %s", diff --git a/tests/cron/test_run_one_job.py b/tests/cron/test_run_one_job.py index 462416fe1b..d63fe0c935 100644 --- a/tests/cron/test_run_one_job.py +++ b/tests/cron/test_run_one_job.py @@ -76,6 +76,89 @@ def test_run_one_job_success_sequence(monkeypatch): assert calls[-1] == ("mark", "j2", True) +def test_run_one_job_exception_delivers_failure_alert(monkeypatch): + """An exception escaping the run body must not become a silent error row.""" + delivered = [] + marked = [] + finished = [] + + monkeypatch.setattr( + s, "create_execution", lambda *_a, **_kw: {"id": "exec-j3"} + ) + monkeypatch.setattr(s, "claim_dispatch", lambda _job_id: True) + monkeypatch.setattr(s, "mark_execution_running", lambda _execution_id: None) + monkeypatch.setattr( + s, + "run_job", + lambda *_a, **_kw: (_ for _ in ()).throw( + RuntimeError("Gemini HTTP 503 (UNAVAILABLE)") + ), + ) + monkeypatch.setattr( + s, + "_deliver_result", + lambda job, content, **_kw: delivered.append((job["id"], content)) or None, + ) + monkeypatch.setattr( + s, + "mark_job_run", + lambda *args, **kwargs: marked.append((args, kwargs)), + ) + monkeypatch.setattr( + s, + "finish_execution", + lambda *args, **kwargs: finished.append((args, kwargs)), + ) + + ok = s.run_one_job({"id": "j3", "name": "morning", "deliver": "telegram"}) + + assert ok is False + assert delivered == [ + ("j3", "⚠️ Cron 'morning' failed: Gemini HTTP 503 (UNAVAILABLE)") + ] + assert marked == [ + (("j3", False, "Gemini HTTP 503 (UNAVAILABLE)"), {"delivery_error": None}) + ] + assert finished == [ + ( + ("exec-j3",), + { + "success": False, + "error": "Gemini HTTP 503 (UNAVAILABLE)", + "delivery_outcome": "delivered", + }, + ) + ] + + +def test_run_one_job_exception_records_failure_alert_delivery_error(monkeypatch): + """A failed fallback alert must populate last_delivery_error.""" + marked = [] + + monkeypatch.setattr( + s, "create_execution", lambda *_a, **_kw: {"id": "exec-j4"} + ) + monkeypatch.setattr(s, "claim_dispatch", lambda _job_id: True) + monkeypatch.setattr(s, "mark_execution_running", lambda _execution_id: None) + monkeypatch.setattr( + s, + "run_job", + lambda *_a, **_kw: (_ for _ in ()).throw(RuntimeError("provider failed")), + ) + monkeypatch.setattr(s, "_deliver_result", lambda *_a, **_kw: "send failed: 502") + monkeypatch.setattr( + s, + "mark_job_run", + lambda *args, **kwargs: marked.append((args, kwargs)), + ) + monkeypatch.setattr(s, "finish_execution", lambda *_a, **_kw: None) + + assert s.run_one_job({"id": "j4", "deliver": "telegram"}) is False + assert marked == [ + (("j4", False, "provider failed"), {"delivery_error": "send failed: 502"}) + ] + + def test_run_one_job_installs_secret_scope_under_multiplex(monkeypatch, tmp_path): """Regression: under profile isolation (multiplex active), run_one_job must execute run_job inside a profile secret scope so credential reads