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:
@@ -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"}
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user