fix(cron): deliver alerts for escaped run failures
This commit is contained in:
+47
-2
@@ -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",
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user