fix(cron): surface external-worker dispatch failures through the incident path

A failed restart-safe handoff in run_one_job() recorded the failure on the
job and in the executions ledger, then returned before any incident or
delivery path ran: no cron_incidents row, no failure-lane notice. Route the
dispatch-failure branch through _deliver_crash_failure() so it opens the
same job+signature incident and delivers the same failure notice as any
other job failure, with the existing alerted-cooldown withholding repeats.
A notice-path exception no longer loses the bookkeeping: mark_job_run and
finish_execution still run with a "failed" delivery outcome (#123401).

(cherry picked from commit e5b5969df1b7ca212c5a6e27d30f4778fb835c52)
This commit is contained in:
liuhao1024
2026-09-26 11:51:50 +08:00
committed by kshitij
parent cb399b56d9
commit ecc5c263ec
2 changed files with 63 additions and 1 deletions

View File

@@ -2730,15 +2730,30 @@ def run_one_job(
logger.error("Job '%s': %s", job["id"], error)
claim = job.get("fire_claim")
owner = str(claim.get("by") or "") if isinstance(claim, dict) else ""
# The dispatch failure is a job failure like any other: it must open an
# incident and leave through the job's failure lane (#123401). Without
# this the outage is silent — no cron_incidents row, no ping — while
# executions.db keeps piling up failed rows.
delivery_error = None
delivery_outcome = "failed"
try:
delivery_error, delivery_outcome = _deliver_crash_failure(
job, error, adapters=adapters, loop=loop)
except Exception as notice_exc:
logger.error(
"Dispatch-failure notice failed for job %s: %s", job["id"], notice_exc)
try:
mark_job_run(
job["id"],
False,
error,
delivery_error=delivery_error,
**({"expected_fire_owner": owner} if owner else {}),
)
finally:
finish_execution(execution_id, success=False, error=error)
finish_execution(
execution_id, success=False, error=error,
delivery_outcome=delivery_outcome)
return True
if extra_prompt is None:
# Gateway-forwarded manual run stamps its prompt on the job via trigger_job; the fire that

View File

@@ -663,6 +663,53 @@ def test_shared_run_path_hands_gateway_fire_to_external_worker(monkeypatch):
run.assert_not_called()
def test_dispatch_failure_opens_incident_and_delivers_failure_notice(
execution_ledger, monkeypatch
):
"""A failed external-worker handoff must surface like any other job failure: one
incident row plus one failure-lane notice, with repeats withheld by the alerted
cooldown (#123401) — not a silent outage while executions.db piles up failed rows."""
import cron.incidents as incidents
import cron.scheduler as scheduler
def _handoff_boom(_job):
raise RuntimeError("worker exited before ownership acknowledgement")
monkeypatch.setattr(scheduler, "_launch_external_cron_worker", _handoff_boom)
monkeypatch.setattr(scheduler, "mark_job_run", lambda *_a, **_k: True)
delivered = []
monkeypatch.setattr(
scheduler, "_deliver_result",
lambda job, content, **_kw: delivered.append(content) or None)
record = execution_ledger.create_execution("job-dispatch", source="builtin")
job = {"id": "job-dispatch", "execution_id": record["id"],
"deliver": "telegram:123"}
assert scheduler.run_one_job(job, adapters=None) is True
rows = incidents.list_incidents()
assert len(rows) == 1
assert rows[0]["job_id"] == "job-dispatch"
assert rows[0]["state"] == "alerted"
assert "Restart-safe cron worker dispatch failed" in rows[0]["error"]
assert len(delivered) == 1 and delivered[0].strip()
finished = execution_ledger.get_execution(record["id"])
assert finished["status"] == "failed"
assert "Restart-safe cron worker dispatch failed" in finished["error"]
assert finished["delivery_outcome"] == "delivered"
# The same dispatch failure again: same signature -> incident dedup, and the
# alerted cooldown withholds the repeat ping.
repeat = execution_ledger.create_execution("job-dispatch", source="builtin")
job2 = {"id": "job-dispatch", "execution_id": repeat["id"],
"deliver": "telegram:123"}
assert scheduler.run_one_job(job2, adapters=None) is True
assert len(incidents.list_incidents()) == 1
assert len(delivered) == 1
assert execution_ledger.get_execution(repeat["id"])["delivery_outcome"] == \
"suppressed_acked"
def test_shutdown_does_not_interrupt_restart_safe_waiter():
import cron.scheduler as scheduler