fix(cron): reconcile run history per execution and read it in the owner's scope

Review of the run-history fallback:

- Merge agent sessions with script-only output docs per execution instead of
  branching on 'any sessions exist': a job converted to no_agent keeps its id
  (update_job supports it), so a surviving historical agent session must not
  hide newer script-only fires. Same-execution rows dedupe by session span.
- Decode output-doc filenames and read the execution ledger inside the owner
  profile's home scope: hermes_time resolves the configured zone through the
  current HERMES_HOME, so a cross-profile request decoded with the
  dashboard's zone and read the dashboard's executions.db.
- Attach per-run status from the execution ledger (one row per fire) matched
  by claim window instead of last_run_at proximity, and surface terminal
  attempts with no surviving doc as their own rows — the latest failed run
  no longer disappears behind (or relabels) an older successful document.
- Read output docs with encoding='utf-8-sig' (Windows footgun: PowerShell and
  some editors BOM files; plain utf-8 breaks on BOM-prefixed docs).

Three regressions cover the converted-job merge, the per-profile timezone
decode, and the missing-newest-doc case.
This commit is contained in:
Hermes Agent
2026-09-25 18:27:06 -05:00
committed by brooklyn!
parent 4f11ac246c
commit 17c6a566f5
2 changed files with 390 additions and 56 deletions

View File

@@ -9,6 +9,7 @@ import asyncio
import functools
import re
import time
from contextlib import contextmanager
from datetime import datetime
from pathlib import Path
from typing import Any, Dict, List, Optional
@@ -24,7 +25,11 @@ from hermes_cli.web_server_cron import (
from hermes_cli.web_models import AutomationBlueprintInstantiate, CronJobCreate, CronJobUpdate
from hermes_cli.web_routers._common import log as _log
from hermes_time import get_timezone as _get_timezone
from hermes_constants import get_hermes_home as _get_hermes_home
from hermes_constants import (
get_hermes_home as _get_hermes_home,
reset_hermes_home_override as _reset_hermes_home_override,
set_hermes_home_override as _set_hermes_home_override,
)
router = APIRouter()
@@ -168,6 +173,34 @@ def _cron_output_runs_dir(profile: Optional[str], job_id: str) -> Path:
return Path(profile_home) / "cron" / "output" / job_id
@contextmanager
def _owner_home_scope(profile: Optional[str]):
"""Keep reads in the owner's profile context for the whole run-history build.
Filename stems, the execution ledger and ``hermes_time``'s configured zone all
resolve through the *current* ``get_hermes_home()``; without this scope a
cross-profile listing decodes another profile's docs with the dashboard's
zone (an hours-off timestamp) and reads the dashboard's executions.db. The
dashboard's own profile needs no override (its home is already current).
"""
if not profile:
yield None
return
try:
_, profile_home = _cron_profile_home(profile)
except Exception:
yield None
return
if Path(profile_home).resolve() == _get_hermes_home().resolve():
yield None
return
token = _set_hermes_home_override(str(profile_home))
try:
yield profile_home
finally:
_reset_hermes_home_override(token)
def _cron_output_run_timestamp(path: Path) -> Optional[float]:
"""Epoch seconds for an output filename's wall time.
@@ -191,7 +224,7 @@ def _cron_output_run_timestamp(path: Path) -> Optional[float]:
def _cron_output_run_preview(path: Path, max_chars: int = 180) -> str:
try:
raw = path.read_text(encoding="utf-8", errors="replace")
raw = path.read_text(encoding="utf-8-sig", errors="replace")
except OSError:
return ""
preview = re.sub(r"\s+", " ", raw).strip()
@@ -226,6 +259,90 @@ def _cron_output_status_label(job: Optional[Dict[str, Any]]) -> str:
return status.replace("_", " ").upper()
def _cron_output_run_row(started_at: float, title: str, preview: Optional[str]) -> Dict[str, Any]:
return {
"title": title,
"preview": preview or None,
"source": "cron_output",
"started_at": started_at,
"last_active": started_at,
"ended_at": started_at,
"input_tokens": 0,
"output_tokens": 0,
"message_count": 0,
"tool_call_count": 0,
"model": None,
"cwd": None,
"archived": False,
"is_active": False,
}
def _iso_to_epoch(text: Any) -> Optional[float]:
if not isinstance(text, str) or not text.strip():
return None
try:
return datetime.fromisoformat(text.strip().replace("Z", "+00:00")).timestamp()
except ValueError:
return None
def _owner_profile_executions(canonical_job_id: str) -> List[Dict[str, Any]]:
"""Terminal execution-ledger rows for the job, newest first.
Each script-only fire creates exactly one ledger row (claimed → completed /
failed), so the ledger — not filename proximity to ``last_run_at`` — is the
per-attempt record that pairs an output doc with its real status. Read inside
the owner-home scope so it opens the OWNER's cron/executions.db.
"""
try:
from cron.executions import list_executions
rows = list_executions(job_id=canonical_job_id, limit=100)
except Exception:
return []
terminal: List[Dict[str, Any]] = []
for row in rows:
if str(row.get("status") or "") not in ("completed", "failed", "unknown"):
continue
terminal.append({
"claimed_at": _iso_to_epoch(row.get("claimed_at")) or 0.0,
"finished_at": _iso_to_epoch(row.get("finished_at")),
"status": str(row.get("status") or ""),
"error": str(row.get("error") or "").strip(),
})
terminal.sort(key=lambda r: r["claimed_at"], reverse=True)
return terminal
def _execution_contains(
attempt: Dict[str, Any], started_at: float, grace_seconds: float = 300.0,
) -> bool:
"""Whether an output doc's timestamp falls inside a ledger attempt's window.
``save_job_output`` runs after the script finishes but before
``finish_execution`` closes the attempt, and both clocks are
``hermes_time.now()`` in the owner's zone, so a doc written by an attempt
lands within ``[claimed_at, finished_at + grace]`` of THAT attempt (a still-
open attempt has no ``finished_at``). Older attempts closed before this doc
was written, which is what keeps a fast-firing job's status from bleeding
onto an earlier run's row.
"""
if started_at + 0.001 < attempt["claimed_at"]:
return False
finish = attempt["finished_at"]
return finish is None or started_at <= finish + grace_seconds
def _execution_status_title(status: str, error: str, fallback: str) -> str:
label = (status or "").replace("_", " ").upper()
if label == "UNKNOWN":
return fallback
if not label:
return fallback
return f"{label} · {error}" if (label == "FAILED" and error) else f"{label} · {fallback}"
def _list_cron_output_runs(
job: Optional[Dict[str, Any]],
canonical_job_id: str,
@@ -240,6 +357,11 @@ def _list_cron_output_runs(
per-run record. Rows mirror /api/sessions shape with source='cron_output'
so the frontend reuses SessionInfo; ids use a cron_output: prefix that can
never collide with SessionDB cron_{job_id}_* session ids.
Status comes from the execution ledger (one row per fire), matched to a doc
by claim window — never from ``last_run_at`` proximity, which mislabels an
older doc when the newest run's doc is missing. Terminal attempts with no
surviving doc get their own metadata rows.
"""
output_dir = _cron_output_runs_dir(profile, canonical_job_id)
try:
@@ -251,11 +373,11 @@ def _list_cron_output_runs(
except OSError:
files = []
latest_ts = _cron_job_last_run_timestamp(job)
latest_status = _cron_output_status_label(job)
executions = _owner_profile_executions(canonical_job_id)
represented: set = set()
runs: List[Dict[str, Any]] = []
for index, path in enumerate(files[:limit]):
for path in files[:limit]:
started_at = _cron_output_run_timestamp(path)
if started_at is None:
try:
@@ -264,36 +386,49 @@ def _list_cron_output_runs(
started_at = 0.0
preview = _cron_output_run_preview(path)
title = preview or "Script-only run"
if (
index == 0
and latest_status
and (latest_ts is None or abs(latest_ts - started_at) <= 120)
):
title = f"{latest_status} · {title}"
for index, attempt in enumerate(executions):
if attempt["status"] not in ("completed", "failed"):
continue
if _execution_contains(attempt, started_at):
represented.add(index)
title = _execution_status_title(attempt["status"], attempt["error"], title)
break
runs.append({
"id": f"cron_output:{canonical_job_id}:{path.stem}",
"title": title,
"preview": preview or None,
"source": "cron_output",
"started_at": started_at,
"last_active": started_at,
"ended_at": started_at,
"input_tokens": 0,
"output_tokens": 0,
"message_count": 0,
"tool_call_count": 0,
"model": None,
"cwd": None,
"archived": False,
"is_active": False,
**_cron_output_run_row(started_at, title, preview),
})
# Terminal attempts whose output doc is gone (pruned, or the run never
# wrote one — e.g. a failed fire) are still real executions: give each its
# own metadata row from the ledger instead of letting the latest attempt's
# status bleed onto an older surviving document.
for index, attempt in enumerate(executions):
if index in represented or attempt["status"] not in ("completed", "failed"):
continue
if len(runs) >= limit:
break
when = attempt["finished_at"] or attempt["claimed_at"]
error = attempt["error"]
runs.append({
"id": f"cron_output:{canonical_job_id}:exec:{index}",
**_cron_output_run_row(
when,
_execution_status_title(attempt["status"], error, "Script-only run"),
error or None,
),
# Internal: the attempt's claim window, so the caller can drop this
# row when a session already represents the same execution.
"_claim_window": (attempt["claimed_at"], attempt["finished_at"]),
})
runs.sort(key=lambda r: float(r.get("started_at") or 0), reverse=True)
if runs:
return runs
# No output docs survived (pruned or never written) but the job HAS run:
# surface one metadata-only row from last_run_at/last_status instead of the
# bare "No runs" the issue reports.
latest_ts = _cron_job_last_run_timestamp(job)
if latest_ts is None:
return []
@@ -305,33 +440,28 @@ def _list_cron_output_runs(
title = f"{title} · {preview}"
return [{
"id": f"cron_output:{canonical_job_id}:latest",
"title": title,
"preview": preview or None,
"source": "cron_output",
"started_at": latest_ts,
"last_active": latest_ts,
"ended_at": latest_ts,
"input_tokens": 0,
"output_tokens": 0,
"message_count": 0,
"tool_call_count": 0,
"model": None,
"cwd": None,
"archived": False,
"is_active": False,
**_cron_output_run_row(latest_ts, title, preview),
}]
def _list_cron_job_runs_sync(job_id: str, profile: Optional[str] = None, limit: int = 20):
"""Run sessions produced by a cron job, newest first.
"""Run history for a cron job, newest first: agent sessions PLUS script-only fires.
Runs are ordinary sessions with id ``cron_{job_id}_{timestamp}`` (see
Agent runs are ordinary sessions with id ``cron_{job_id}_{timestamp}`` (see
cron/scheduler.run_job); ``source='cron'`` plus the id prefix binds them to
this job. Same row shape as ``/api/sessions`` so the frontend reuses
SessionInfo. Backed by ``SessionDB.list_cron_job_runs`` — a bounded id-range
this job. Backed by ``SessionDB.list_cron_job_runs`` — a bounded id-range
scan, so cost scales with the requested window, not total cron history.
Script-only (no_agent) jobs never write sessions: their history falls back
to per-fire output docs (``_list_cron_output_runs``).
Script-only (no_agent) jobs never write sessions; their per-fire output docs
(``_list_cron_output_runs``) fill the gaps. The two representations are
reconciled PER EXECUTION, not by branch: a job can switch modes without
changing its id (update_job supports no_agent on an existing record), so a
surviving historical agent session must not hide newer script-only fires
and a recent session must not hide older script output. Rows carrying the
same wall-clock execution (an agent fire writes BOTH a session and an
output doc) are de-duplicated by timestamp window; script-only rows that
match no session are added in. All owner-scoped reads (SessionDB, output
docs, execution ledger, filename timezone) run inside the owner's profile
scope so the dashboard's own zone/store is never borrowed cross-profile.
"""
selected = _job_owner_profile(job_id, profile)
# job_id may be a human name; resolve to the canonical id used in run-session ids.
@@ -347,20 +477,64 @@ def _list_cron_job_runs_sync(job_id: str, profile: Optional[str] = None, limit:
except (TypeError, ValueError):
limit_n = 20
db = _open_session_db_for_profile(selected, read_only=True)
try:
runs = db.list_cron_job_runs(canonical, limit=limit_n, offset=0)
if not runs:
return {"runs": _list_cron_output_runs(job, canonical, selected, limit_n), "limit": limit_n}
with _owner_home_scope(selected):
db = _open_session_db_for_profile(selected, read_only=True)
try:
session_runs = db.list_cron_job_runs(canonical, limit=limit_n, offset=0)
finally:
db.close()
now = time.time()
for s in runs:
for s in session_runs:
s["is_active"] = s.get("ended_at") is None and (now - s.get("last_active", s.get("started_at", 0))) < 300
s["archived"] = bool(s.get("archived"))
if selected:
s["profile"] = selected
return {"runs": runs, "limit": limit_n}
finally:
db.close()
doc_runs = _list_cron_output_runs(job, canonical, selected, limit_n)
return {"runs": _reconcile_cron_runs(session_runs, doc_runs, limit_n), "limit": limit_n}
def _reconcile_cron_runs(
session_runs: List[Dict[str, Any]],
doc_runs: List[Dict[str, Any]],
limit: int,
) -> List[Dict[str, Any]]:
"""Merge session rows and output-doc rows per execution, newest first.
An agent fire writes BOTH a session and an output doc for the same
execution; keeping both would double-render the run. Sessions carry the
richer record (tokens, message counts, chat navigation), so a doc row
within ``_SESSION_DOC_MATCH_SECONDS`` of a session row is the same
execution and the session wins. Doc rows matching no session are distinct
executions — script-only fires — and are kept. A session row must never be
dropped in favor of a doc: the doc is the fallback representation.
"""
_SESSION_DOC_MATCH_SECONDS = 300.0
merged = list(session_runs)
for doc in doc_runs:
doc = {k: v for k, v in doc.items() if not k.startswith("_")}
doc_ts = float(doc.get("started_at") or 0)
if any(_doc_matches_session(doc_ts, s, _SESSION_DOC_MATCH_SECONDS) for s in session_runs):
continue
merged.append(doc)
merged.sort(key=lambda r: float(r.get("started_at") or 0), reverse=True)
return merged[:limit]
def _doc_matches_session(doc_ts: float, session: Dict[str, Any], grace_seconds: float) -> bool:
"""Whether an output doc belongs to a session's run.
The doc is written when the run FINISHES, so its filename timestamp sits
at the run's end — compare against the session's full span, not its start:
``[started_at - grace, last_active + grace]`` (``last_active`` covers the
still-open case where ``ended_at`` is None).
"""
start = float(session.get("started_at") or 0)
end = float(session.get("ended_at") or session.get("last_active") or start)
return start - grace_seconds <= doc_ts <= end + grace_seconds
_EXECUTION_FIELDS = {"prompt", "skill", "skills", "script", "no_agent"}

View File

@@ -1285,7 +1285,10 @@ class TestCronRunHistoryFallback:
f"cron_output:{job_id}:2026-07-08_09-00-00",
]
assert runs[0]["source"] == "cron_output"
assert runs[0]["title"].startswith("OK · latest output")
# Status is attached per execution from the ledger, never from
# last_run_at/last_status proximity (#42433 review): with no ledger
# rows the doc's own output is the honest title.
assert runs[0]["title"] == "latest output"
assert runs[1]["title"] == "older output"
assert all(run["is_active"] is False for run in runs)
@@ -1331,15 +1334,25 @@ class TestCronRunHistoryFallback:
def test_keeps_session_runs_when_db_history_exists(
self, isolated_profiles, monkeypatch
):
"""An agent fire writes BOTH a session and an output doc; the session row
is the richer representation, so the doc for the same execution is
dropped and the session survives the merge."""
from zoneinfo import ZoneInfo
job_id = "job-with-sessions"
home = isolated_profiles["default"]
# Doc written at the END of the session's run (2026 wall time): the
# session span covers it, so it is the same execution.
doc_epoch = datetime.strptime(
"2026-07-08_09-05-00", "%Y-%m-%d_%H-%M-%S").replace(tzinfo=ZoneInfo("UTC")).timestamp()
self._write_output_doc(monkeypatch, home, job_id, "2026-07-08_09-05-00.md", "fallback output\n")
class _FakeDB:
def list_cron_job_runs(self, canonical, limit=20, offset=0):
return [{
"id": f"cron_{canonical}_1", "source": "cron",
"started_at": 123.0, "last_active": 125.0, "ended_at": 126.0, "archived": False,
"started_at": doc_epoch - 120.0, "last_active": doc_epoch - 10.0,
"ended_at": doc_epoch - 5.0, "archived": False,
}]
def close(self):
@@ -1401,3 +1414,150 @@ class TestCronRunHistoryFallback:
result = _rt_cron._list_cron_job_runs_sync(job_id, limit=5)
assert result["runs"] == []
def test_converted_job_merges_agent_sessions_with_script_fires(
self, isolated_profiles, monkeypatch, no_session_rows
):
"""Review of #42433: updating an existing job to no_agent=True keeps the
job id, so a surviving historical agent session must not hide later
script-only executions. Reconcile per execution: the newer script run
appears first, the older agent session stays exactly once."""
from zoneinfo import ZoneInfo
job_id = "j1"
home = isolated_profiles["default"]
# Old agent session (survives in SessionDB) and a NEWER script-only fire.
agent_started = datetime.strptime(
"2026-09-24_08-00-00", "%Y-%m-%d_%H-%M-%S").replace(tzinfo=ZoneInfo("UTC")).timestamp()
self._write_output_doc(monkeypatch, home, job_id, "2026-09-25_00-14-10.md", "script run output\n")
class _FakeDB:
def list_cron_job_runs(self, canonical, limit=20, offset=0):
return [{
"id": f"cron_{canonical}_old", "source": "cron",
"started_at": agent_started, "last_active": agent_started + 60.0,
"ended_at": agent_started + 60.0, "archived": False,
}]
def close(self):
pass
monkeypatch.setattr(
_rt_cron, "_open_session_db_for_profile", lambda profile, *, read_only: _FakeDB()
)
monkeypatch.setattr(_rt_cron, "_job_owner_profile", lambda _job_id, _profile: "default")
monkeypatch.setattr(
_rt_cron, "_call_cron_for_profile",
lambda _profile, cmd, *_args, **_kwargs: {
**self._script_only_job(job_id),
} if cmd == "get_job" else None,
)
result = _rt_cron._list_cron_job_runs_sync(job_id, limit=10)
runs = result["runs"]
assert [run["id"] for run in runs] == [
f"cron_output:{job_id}:2026-09-25_00-14-10",
f"cron_{job_id}_old",
]
# The older agent session appears exactly once — not shadowed, not doubled.
assert sum(1 for run in runs if run["id"] == f"cron_{job_id}_old") == 1
def test_output_doc_timestamps_decode_in_the_owners_profile_timezone(
self, isolated_profiles, monkeypatch, no_session_rows
):
"""Review of #42433: the filename stem is written with the OWNER
profile's configured zone. A cross-profile request must decode it in
that zone — not the dashboard's — with no global HERMES_TIMEZONE set
(both profiles configure config.yaml timezones, so a shared env var
cannot mask the mix-up)."""
from zoneinfo import ZoneInfo
from hermes_time import _tz_cache
job_id = "job-owner-tz"
worker_home = isolated_profiles["worker_alpha"]
# worker_alpha runs in Asia/Taipei; the dashboard profile stays UTC.
(worker_home / "config.yaml").write_text(
"model: test-model\ntimezone: Asia/Taipei\n", encoding="utf-8")
# Wall time 2026-09-25_00-14-10 in Asia/Taipei (+08:00) = 2026-09-24 16:14:10 UTC.
self._write_output_doc(monkeypatch, worker_home, job_id, "2026-09-25_00-14-10.md", "tz output\n")
_tz_cache.clear()
monkeypatch.setattr(_rt_cron, "_job_owner_profile", lambda _job_id, _profile: "worker_alpha")
monkeypatch.setattr(
_rt_cron, "_call_cron_for_profile",
lambda _profile, cmd, *_args, **_kwargs: {
**self._script_only_job(job_id), "last_status": "ok",
} if cmd == "get_job" else None,
)
try:
result = _rt_cron._list_cron_job_runs_sync(job_id, profile="default", limit=5)
finally:
_tz_cache.clear()
naive = datetime.strptime("2026-09-25_00-14-10", "%Y-%m-%d_%H-%M-%S")
expected = naive.replace(tzinfo=ZoneInfo("Asia/Taipei")).timestamp()
started = result["runs"][0]["started_at"]
assert abs(started - expected) < 1, (
f"Decoded with the dashboard's zone instead of the owner's: got epoch {started}, expected {expected}"
)
def test_latest_attempt_without_a_doc_gets_its_own_row(
self, isolated_profiles, monkeypatch, no_session_rows
):
"""Review of #42433: two runs, only the newest run's output doc missing.
The surviving older doc must keep its own status — not inherit the
latest attempt's — and the failed latest attempt must appear as its own
row carrying its timestamp and error (a 120s filename window cannot
establish that the newest file belongs to last_run_at)."""
from cron import executions as cron_executions
from zoneinfo import ZoneInfo
job_id = "job-missing-newest-doc"
home = isolated_profiles["default"]
ledger_path = home / "cron" / "executions.db"
monkeypatch.setattr(cron_executions, "EXECUTIONS_FILE", ledger_path)
def _at(name):
return datetime.strptime(name, "%Y-%m-%d_%H-%M-%S").replace(tzinfo=ZoneInfo("UTC"))
# Attempt 1 (10:00): succeeded; its doc survives.
monkeypatch.setattr(cron_executions, "_hermes_now", lambda: _at("2026-09-25_10-00-00"))
first = cron_executions.create_execution(job_id, source="direct")
monkeypatch.setattr(cron_executions, "_hermes_now", lambda: _at("2026-09-25_10-00-30"))
cron_executions.finish_execution(first["id"], success=True)
# Attempt 2 (10:01): failed; its doc is gone.
monkeypatch.setattr(cron_executions, "_hermes_now", lambda: _at("2026-09-25_10-01-00"))
second = cron_executions.create_execution(job_id, source="direct")
monkeypatch.setattr(cron_executions, "_hermes_now", lambda: _at("2026-09-25_10-01-20"))
cron_executions.finish_execution(second["id"], success=False, error="new run failed")
self._write_output_doc(monkeypatch, home, job_id, "2026-09-25_10-00-00.md", "successful older run\n")
monkeypatch.setattr(_rt_cron, "_job_owner_profile", lambda _job_id, _profile: "default")
monkeypatch.setattr(
_rt_cron, "_call_cron_for_profile",
lambda _profile, cmd, *_args, **_kwargs: {
**self._script_only_job(job_id),
"last_run_at": "2026-09-25T10:01:00+00:00",
"last_status": "error",
"last_error": "new run failed",
} if cmd == "get_job" else None,
)
result = _rt_cron._list_cron_job_runs_sync(job_id, limit=10)
runs = result["runs"]
assert len(runs) == 2
newest, older = runs[0], runs[1]
# The failed latest attempt is present as its own row, with its own
# timestamp and error — not silently dropped behind the older doc.
assert newest["id"].startswith(f"cron_output:{job_id}:exec:")
assert newest["started_at"] == _at("2026-09-25_10-01-20").timestamp()
assert newest["preview"] == "new run failed"
assert newest["title"].startswith("FAILED")
# The surviving older doc keeps its own output and its own completed
# status — not relabeled with the newer attempt's error.
assert older["id"] == f"cron_output:{job_id}:2026-09-25_10-00-00"
assert older["title"].startswith("COMPLETED · successful older run")
assert older["started_at"] == _at("2026-09-25_10-00-00").timestamp()