Files
hermes-agent/tests/cron/test_execution_ledger.py
alt-glitch 318f6ac63f fix(cron): order execution, delivery and incident history by instant, not by text
The cron SQLite stores write hermes_time.now().isoformat(), an ISO string
that carries the local UTC offset. That offset changes at a DST transition
and after a timezone config change, so the text order of these columns is
not their time order. Across a fall-back, 01:10-05:00 is 20 minutes after
01:50-04:00 but sorts before it. latest_execution(), latest_executions()
(which fills job["latest_execution"] for `hermes cron list` and the
dashboard), list_executions() and its before_claimed_at cursor, the
delivery queue's pending-claim order, delivery and execution pruning, and
list_incidents() all returned the older row first.

Every ORDER BY and the cursor comparison on these columns now use
julianday(col), which SQLite resolves to the same instant for any offset.
The raw text stays as the next sort key, so rows within the same
millisecond (julianday's resolution) keep their microsecond order, and the
cursor compares the same (instant, text) key on both sides. Stored values
and the schema are unchanged, so existing rows with mixed offsets order
correctly without a migration.

latest_executions() moves from a correlated per-row subquery to one
ROW_NUMBER() window. With julianday() in the ORDER BY the index on
claimed_at no longer serves the per-row LIMIT 1, and the correlated form
took about 90 ms at 1000 rows; the window form takes about 1 ms.

Refs #86520
Thanks to #86889 (zhao0112) for the executions-ledger diagnosis.
2026-09-28 19:07:43 +05:30

517 lines
20 KiB
Python

"""Durable cron execution-ledger behavior."""
from __future__ import annotations
import json
import os
import sqlite3
import subprocess
import sys
from pathlib import Path
def _point_ledger(monkeypatch, tmp_path):
import cron.executions as executions
monkeypatch.setattr(executions, "EXECUTIONS_FILE", tmp_path / "cron" / "executions.db")
return executions
def test_execution_transitions_are_durable(monkeypatch, tmp_path):
executions = _point_ledger(monkeypatch, tmp_path)
claimed = executions.create_execution("job-1", source="builtin")
assert claimed["status"] == "claimed"
assert claimed["claimed_at"]
assert claimed["started_at"] is None
assert claimed["finished_at"] is None
running = executions.mark_execution_running(claimed["id"])
assert running["status"] == "running"
assert running["started_at"]
completed = executions.finish_execution(claimed["id"], success=True)
assert completed["status"] == "completed"
assert completed["finished_at"]
assert completed["error"] is None
persisted = executions.list_executions(job_id="job-1")
assert persisted == [completed]
def test_execution_can_be_loaded_by_exact_attempt_id(monkeypatch, tmp_path):
executions = _point_ledger(monkeypatch, tmp_path)
first = executions.create_execution("same-job", source="builtin")
second = executions.create_execution("same-job", source="builtin")
assert executions.get_execution(first["id"]) == first
assert executions.get_execution(second["id"]) == second
assert executions.get_execution("missing") is None
def test_fresh_external_handoff_is_not_recovered_before_worker_adopts(
monkeypatch, tmp_path
):
executions = _point_ledger(monkeypatch, tmp_path)
record = executions.create_execution("handoff-job", source="builtin")
assert executions.mark_execution_handoff_pending(record["id"]) is not None
monkeypatch.setattr(executions, "_PROCESS_ID", "replacement-gateway")
monkeypatch.setattr(executions, "_owner_is_live", lambda _pid, _started: False)
assert executions.recover_interrupted_executions() == 0
assert executions.get_execution(record["id"])["status"] == "claimed"
adopted = executions.adopt_claimed_execution(record["id"])
assert adopted["status"] == "running"
assert adopted["handoff_pending"] == 0
def test_stale_external_handoff_is_recovered_unknown(monkeypatch, tmp_path):
executions = _point_ledger(monkeypatch, tmp_path)
record = executions.create_execution("handoff-job", source="builtin")
pending = executions.mark_execution_handoff_pending(record["id"])
monkeypatch.setattr(executions, "_PROCESS_ID", "replacement-gateway")
monkeypatch.setattr(executions, "_owner_is_live", lambda _pid, _started: False)
monkeypatch.setattr(
executions.time,
"time",
lambda: pending["handoff_started_at"]
+ executions.HANDOFF_ADOPTION_GRACE_SECONDS
+ 1,
)
assert executions.recover_interrupted_executions() == 1
recovered = executions.get_execution(record["id"])
assert recovered["status"] == "unknown"
assert recovered["handoff_pending"] == 0
def test_recovery_does_not_overwrite_concurrent_worker_adoption(monkeypatch, tmp_path):
executions = _point_ledger(monkeypatch, tmp_path)
record = executions.create_execution("adoption-race", source="builtin")
pending = executions.mark_execution_handoff_pending(record["id"])
assert pending is not None
monkeypatch.setattr(executions, "_PROCESS_ID", "replacement-scheduler")
monkeypatch.setattr(
executions.time,
"time",
lambda: pending["handoff_started_at"]
+ executions.HANDOFF_ADOPTION_GRACE_SECONDS
+ 1,
)
def adopt_while_liveness_is_checked(_pid, _started_at):
monkeypatch.setattr(executions, "_PROCESS_ID", "external-worker")
monkeypatch.setattr(executions.os, "getpid", lambda: 4242)
monkeypatch.setattr(executions, "_process_start_time", lambda _pid: 9876)
assert executions.adopt_claimed_execution(record["id"]) is not None
return False
monkeypatch.setattr(executions, "_owner_is_live", adopt_while_liveness_is_checked)
assert executions.recover_interrupted_executions() == 0
current = executions.get_execution(record["id"])
assert current is not None
assert current["status"] == "running"
assert current["process_id"] == "external-worker"
assert current["pid"] == 4242
def test_foreign_process_cannot_start_or_finish_execution(monkeypatch, tmp_path):
executions = _point_ledger(monkeypatch, tmp_path)
record = executions.create_execution("owner-fence", source="builtin")
original_process_id = executions._PROCESS_ID
original_pid = record["pid"]
monkeypatch.setattr(executions, "_PROCESS_ID", "foreign-process")
monkeypatch.setattr(executions.os, "getpid", lambda: original_pid + 1)
assert executions.mark_execution_running(record["id"]) is None
assert executions.finish_execution(record["id"], success=True) is None
monkeypatch.setattr(executions, "_PROCESS_ID", original_process_id)
monkeypatch.setattr(executions.os, "getpid", lambda: original_pid)
assert executions.mark_execution_running(record["id"]) is not None
assert executions.finish_execution(record["id"], success=True) is not None
def test_execution_ledger_follows_the_current_profile_home(monkeypatch, tmp_path):
import cron.executions as executions
current_home = {"path": tmp_path / "default"}
monkeypatch.setattr(executions, "EXECUTIONS_FILE", None)
monkeypatch.setattr(executions, "get_hermes_home", lambda: current_home["path"])
default_row = executions.create_execution("default-job", source="builtin")
current_home["path"] = tmp_path / "worker"
worker_row = executions.create_execution("worker-job", source="builtin")
assert executions.list_executions() == [worker_row]
current_home["path"] = tmp_path / "default"
assert executions.list_executions() == [default_row]
assert (tmp_path / "default" / "cron" / "executions.db").is_file()
assert (tmp_path / "worker" / "cron" / "executions.db").is_file()
def test_terminal_execution_cannot_be_rewritten(monkeypatch, tmp_path):
executions = _point_ledger(monkeypatch, tmp_path)
record = executions.create_execution("immutable", source="builtin")
executions.mark_execution_running(record["id"])
executions.finish_execution(record["id"], success=True)
assert executions.finish_execution(
record["id"], success=False, error="late writer"
) is None
assert executions.latest_execution("immutable")["status"] == "completed"
def test_retention_bounds_terminal_history_but_preserves_inflight(monkeypatch, tmp_path):
executions = _point_ledger(monkeypatch, tmp_path)
monkeypatch.setattr(executions, "MAX_TERMINAL_EXECUTIONS", 3)
inflight = executions.create_execution("live", source="builtin")
executions.mark_execution_running(inflight["id"])
for index in range(8):
row = executions.create_execution(f"done-{index}", source="builtin")
executions.finish_execution(row["id"], success=True)
records = executions.list_executions(limit=100)
assert len([row for row in records if row["status"] == "completed"]) == 3
assert executions.latest_execution("live")["status"] == "running"
def test_recently_finished_long_running_execution_survives_retention(
monkeypatch, tmp_path
):
executions = _point_ledger(monkeypatch, tmp_path)
monkeypatch.setattr(executions, "MAX_TERMINAL_EXECUTIONS", 1)
long_running = executions.create_execution("long-running", source="builtin")
assert executions.mark_execution_running(long_running["id"]) is not None
newer = executions.create_execution("newer", source="builtin")
assert executions.finish_execution(newer["id"], success=True) is not None
finished = executions.finish_execution(long_running["id"], success=True)
assert finished is not None
assert finished["status"] == "completed"
assert executions.get_execution(long_running["id"])["status"] == "completed"
assert executions.get_execution(newer["id"]) is None
def test_corrupt_store_fails_closed_without_overwrite(monkeypatch, tmp_path):
executions = _point_ledger(monkeypatch, tmp_path)
executions.EXECUTIONS_FILE.parent.mkdir(parents=True)
executions.EXECUTIONS_FILE.write_bytes(b"not a sqlite database")
with __import__("pytest").raises(sqlite3.DatabaseError):
executions.create_execution("new", source="builtin")
assert executions.EXECUTIONS_FILE.read_bytes() == b"not a sqlite database"
def test_cron_runs_cli_prints_execution_history(monkeypatch, tmp_path, capsys):
executions = _point_ledger(monkeypatch, tmp_path)
row = executions.create_execution("cli-job", source="builtin")
executions.finish_execution(row["id"], success=False, error="boom")
from hermes_cli.cron import cron_runs
cron_runs("cli-job", limit=10)
output = capsys.readouterr().out
assert row["id"] in output
assert "failed" in output
assert "boom" in output
def test_quick_backup_includes_execution_ledger():
from hermes_cli.backup import _QUICK_STATE_FILES
assert "cron/executions.db" in _QUICK_STATE_FILES
def test_failed_execution_keeps_error(monkeypatch, tmp_path):
executions = _point_ledger(monkeypatch, tmp_path)
record = executions.create_execution("job-2", source="external")
failed = executions.finish_execution(record["id"], success=False, error="provider exploded")
assert failed["status"] == "failed"
assert failed["error"] == "provider exploded"
def test_recovery_does_not_mark_live_process_execution_unknown(monkeypatch, tmp_path):
executions = _point_ledger(monkeypatch, tmp_path)
record = executions.create_execution("still-live", source="builtin")
executions.mark_execution_running(record["id"])
assert executions.recover_interrupted_executions() == 0
assert executions.latest_execution("still-live")["status"] == "running"
def test_restart_marks_interrupted_execution_unknown_without_requeue(tmp_path):
"""Real temp-HERMES_HOME subprocess restart: in-flight is audit-only unknown."""
home = tmp_path / "home"
repo = Path(__file__).resolve().parents[2]
env = os.environ.copy()
env["HERMES_HOME"] = str(home)
env["PYTHONPATH"] = str(repo)
create = subprocess.run(
[
sys.executable,
"-c",
"from cron.executions import create_execution, mark_execution_running; "
"r=create_execution('restart-job', source='builtin'); "
"mark_execution_running(r['id']); print(r['id'])",
],
cwd=repo,
env=env,
text=True,
capture_output=True,
check=True,
)
execution_id = create.stdout.strip()
recover = subprocess.run(
[
sys.executable,
"-c",
"import json; from cron.executions import recover_interrupted_executions, list_executions; "
"print(recover_interrupted_executions()); "
"print(json.dumps(list_executions(job_id='restart-job'))) ",
],
cwd=repo,
env=env,
text=True,
capture_output=True,
check=True,
)
lines = recover.stdout.strip().splitlines()
assert lines[0] == "1"
records = json.loads(lines[1])
assert len(records) == 1
assert records[0]["id"] == execution_id
assert records[0]["status"] == "unknown"
assert records[0]["finished_at"]
assert "restart" in records[0]["error"].lower()
# Recovery only classifies the old attempt. It must not manufacture a new
# claimed record (which would imply an automatic retry).
assert [r["status"] for r in records] == ["unknown"]
def test_generic_submit_failure_finishes_attempt_and_releases_guard(monkeypatch):
import cron.scheduler as scheduler
class BrokenPool:
def submit(self, _callable):
raise ValueError("executor rejected")
finished = []
monkeypatch.setattr(
scheduler, "create_execution",
lambda *_args, **_kwargs: {"id": "exec-submit-fail"},
)
monkeypatch.setattr(
scheduler, "finish_execution",
lambda execution_id, **kwargs: finished.append((execution_id, kwargs)),
)
monkeypatch.setattr(scheduler, "get_due_jobs", lambda: [{"id": "submit-fail"}])
monkeypatch.setattr(scheduler, "claim_job_for_fire", lambda _job_id: True)
monkeypatch.setattr(scheduler, "_get_parallel_pool", lambda _workers: BrokenPool())
assert scheduler.tick(verbose=False, sync=False) == 0
assert finished == [
("exec-submit-fail", {
"success": False,
"error": "Executor dispatch failed: executor rejected",
})
]
assert "submit-fail" not in scheduler.get_running_job_ids()
def test_run_one_job_records_running_then_terminal(monkeypatch):
import cron.scheduler as scheduler
events = []
run_execution_ids = []
monkeypatch.setattr(
scheduler,
"mark_execution_running",
lambda execution_id: events.append(("running", execution_id)) or {},
raising=False,
)
monkeypatch.setattr(
scheduler,
"finish_execution",
lambda execution_id, **kwargs: events.append(("finish", execution_id, kwargs)),
raising=False,
)
monkeypatch.setattr(scheduler, "claim_dispatch", lambda _job_id: True)
def fake_run_job(job, *, defer_agent_teardown=None, execution_id=None, **_kw):
run_execution_ids.append(execution_id)
return True, "output", "response", None
monkeypatch.setattr(scheduler, "run_job", fake_run_job)
monkeypatch.setattr(scheduler, "save_job_output", lambda *_args: None)
monkeypatch.setattr(scheduler, "_deliver_result", lambda *_args, **_kwargs: None)
monkeypatch.setattr(scheduler, "mark_job_run", lambda *_args, **_kwargs: None)
assert scheduler.run_one_job({"id": "job-3", "execution_id": "exec-3"}) is True
assert run_execution_ids == ["exec-3"]
assert events[0] == ("running", "exec-3")
assert events[-1][0:2] == ("finish", "exec-3")
assert events[-1][2]["success"] is True
def test_provider_start_recovers_interrupted_records_before_tick(monkeypatch):
import cron.scheduler_provider as provider
events = []
stop = __import__("threading").Event()
stop.set()
monkeypatch.setattr(
"cron.executions.recover_interrupted_executions",
lambda: events.append("recover") or 0,
raising=False,
)
monkeypatch.setattr("cron.jobs.record_ticker_heartbeat", lambda **_kwargs: events.append("heartbeat"))
provider.InProcessCronScheduler().start(stop, interval=1)
assert events[:2] == ["recover", "heartbeat"]
def test_external_provider_start_recovers_interrupted_records(monkeypatch):
from plugins.cron_providers.chronos import ChronosCronScheduler
provider = ChronosCronScheduler()
provider._client = type("Client", (), {"arm": lambda self, **kwargs: None})()
events = []
monkeypatch.setattr(
"cron.executions.recover_interrupted_executions",
lambda: events.append("recover") or 0,
)
monkeypatch.setattr(provider, "reconcile", lambda: events.append("reconcile"))
provider.start(__import__("threading").Event())
assert events == ["recover", "reconcile"]
class _TrackingConnection:
"""Delegates to a real sqlite3.Connection while recording close() calls.
sqlite3.Connection is a static C type: it has no per-instance __dict__
and its class methods can't be monkeypatched, so open/close tracking is
done via a delegating wrapper returned in place of the real connection.
"""
def __init__(self, real, closed_ids):
object.__setattr__(self, "_real", real)
object.__setattr__(self, "_closed_ids", closed_ids)
def close(self):
self._closed_ids.append(id(self._real))
self._real.close()
def __enter__(self):
self._real.__enter__()
return self
def __exit__(self, exc_type, exc, tb):
return self._real.__exit__(exc_type, exc, tb)
def __getattr__(self, name):
return getattr(self._real, name)
def __setattr__(self, name, value):
setattr(self._real, name, value)
def _count_open_connections(executions, monkeypatch):
"""Wrap sqlite3.connect to track open/close balance for the ledger module."""
opened_ids = []
closed_ids = []
real_connect = sqlite3.connect
def tracking_connect(*args, **kwargs):
conn = real_connect(*args, **kwargs)
opened_ids.append(id(conn))
return _TrackingConnection(conn, closed_ids)
monkeypatch.setattr(executions.sqlite3, "connect", tracking_connect)
return opened_ids, closed_ids
def test_ledger_operations_close_every_connection(monkeypatch, tmp_path):
"""Regression for #69567: every ledger call must close its connection
deterministically instead of relying on garbage collection."""
executions = _point_ledger(monkeypatch, tmp_path)
opened, closed = _count_open_connections(executions, monkeypatch)
record = executions.create_execution("leak-check", source="builtin")
executions.mark_execution_running(record["id"])
executions.finish_execution(record["id"], success=True)
executions.list_executions(job_id="leak-check")
executions.latest_executions(["leak-check"])
executions.recover_interrupted_executions()
assert opened
assert sorted(opened) == sorted(closed)
def test_early_return_still_closes_connection(monkeypatch, tmp_path):
"""mark_execution_running returns None mid-block on a bad transition;
the connection must still be closed rather than leaked."""
executions = _point_ledger(monkeypatch, tmp_path)
opened, closed = _count_open_connections(executions, monkeypatch)
assert executions.mark_execution_running("does-not-exist") is None
assert len(opened) == 1
assert len(closed) == 1
def test_job_listing_exposes_latest_execution(monkeypatch, tmp_path):
import cron.jobs as jobs
monkeypatch.setattr(jobs, "CRON_DIR", tmp_path / "cron")
monkeypatch.setattr(jobs, "JOBS_FILE", tmp_path / "cron" / "jobs.json")
monkeypatch.setattr(jobs, "OUTPUT_DIR", tmp_path / "cron" / "output")
executions = _point_ledger(monkeypatch, tmp_path)
job = jobs.create_job(prompt="audit me", schedule="every 1h", name="audit")
record = executions.create_execution(job["id"], source="builtin")
executions.mark_execution_running(record["id"])
listed = jobs.list_jobs(include_disabled=True)
assert listed[0]["latest_execution"]["id"] == record["id"]
assert listed[0]["latest_execution"]["status"] == "running"
def test_history_orders_by_instant_across_dst_fall_back(monkeypatch, tmp_path):
"""01:10-05:00 is 20 minutes after 01:50-04:00 but sorts first as text."""
from datetime import datetime
from zoneinfo import ZoneInfo
executions = _point_ledger(monkeypatch, tmp_path)
new_york = ZoneInfo("America/New_York")
first_pass = datetime(2026, 11, 1, 1, 50, tzinfo=new_york, fold=0)
second_pass = datetime(2026, 11, 1, 1, 10, tzinfo=new_york, fold=1)
monkeypatch.setattr(executions, "_hermes_now", lambda: first_pass)
earlier = executions.create_execution("dst-job", source="builtin")
monkeypatch.setattr(executions, "_hermes_now", lambda: second_pass)
later = executions.create_execution("dst-job", source="builtin")
assert later["claimed_at"] < earlier["claimed_at"]
assert executions.latest_execution("dst-job")["id"] == later["id"]
assert executions.latest_executions(["dst-job"])["dst-job"]["id"] == later["id"]
assert [r["id"] for r in executions.list_executions(job_id="dst-job")] == [
later["id"], earlier["id"],
]
page = executions.list_executions(job_id="dst-job", before_claimed_at=later["claimed_at"])
assert [r["id"] for r in page] == [earlier["id"]]