fix(cron): deliver alerts for escaped run failures

This commit is contained in:
fangliquanflq
2026-08-15 07:03:51 +08:00
committed by Teknium
parent ec1aa8977c
commit 4668750fad
2 changed files with 130 additions and 2 deletions
+47 -2
View File
@@ -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",
+83
View File
@@ -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