File size: 14,495 Bytes
f76c374 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 | """Profile-local durable audit ledger for cron execution attempts.
The ledger records what is known about each attempt; it is not a retry queue. Interrupted attempts
become ``unknown`` only after their exact owner process is proved gone. Terminal states are
immutable.
"""
from __future__ import annotations
import os
import sqlite3
import threading
import time
import uuid
from contextlib import contextmanager
from pathlib import Path
from typing import Any, Dict, Iterator, List, Optional
from hermes_constants import get_hermes_home
from hermes_time import now as _hermes_now
# Optional test override. Production resolves the path at transaction time so dashboard operations
# that temporarily enter another profile cannot leak that profile's records into the import-time
# home.
EXECUTIONS_FILE: Optional[Path] = None
MAX_TERMINAL_EXECUTIONS = 1000
HANDOFF_ADOPTION_GRACE_SECONDS = 30.0
_TERMINAL_STATES = ("completed", "failed", "unknown")
_lock = threading.RLock()
_PROCESS_ID = uuid.uuid4().hex
# --- executions ledger --------------------------------------------------------------------------
def _connect() -> sqlite3.Connection:
# Late imports: a scheduler daemon that outlives an on-disk upgrade already has the OLD
# ``hermes_cli.sqlite_util`` / ``cron.jobs`` cached, so new names must be resolved at call time,
# not at import time (the guarantee cron/ledger.py used to carry, see e24c8499).
from cron.jobs import _ensure_cron_dir
from hermes_cli.sqlite_util import open_db
path = EXECUTIONS_FILE or (get_hermes_home().resolve() / "cron" / "executions.db")
_ensure_cron_dir(path.parent)
return open_db(path, db_label="cron/executions.db", synchronous_full=True, initialize=_initialize_schema)
def _initialize_schema(conn: sqlite3.Connection) -> None:
from hermes_cli.sqlite_util import add_column_if_missing
conn.execute(
"""CREATE TABLE IF NOT EXISTS executions (
id TEXT PRIMARY KEY,
job_id TEXT NOT NULL,
source TEXT NOT NULL,
process_id TEXT NOT NULL,
pid INTEGER NOT NULL,
process_started_at INTEGER,
status TEXT NOT NULL CHECK(status IN
('claimed','running','completed','failed','unknown')),
handoff_pending INTEGER NOT NULL DEFAULT 0,
handoff_started_at REAL,
claimed_at TEXT NOT NULL,
started_at TEXT,
finished_at TEXT,
error TEXT
)"""
)
add_column_if_missing(
conn, "executions", "handoff_pending",
"handoff_pending INTEGER NOT NULL DEFAULT 0",
)
add_column_if_missing(
conn, "executions", "handoff_started_at", "handoff_started_at REAL"
)
conn.execute(
"CREATE INDEX IF NOT EXISTS idx_executions_job_claimed "
"ON executions(job_id, claimed_at DESC, id DESC)"
)
conn.execute(
"CREATE INDEX IF NOT EXISTS idx_executions_status_claimed "
"ON executions(status, claimed_at DESC, id DESC)"
)
add_column_if_missing(conn, "executions", "delivery_outcome", "delivery_outcome TEXT")
add_column_if_missing(conn, "executions", "scheduled_instant", "scheduled_instant TEXT")
conn.execute(
"CREATE INDEX IF NOT EXISTS idx_executions_occurrence "
"ON executions(job_id, scheduled_instant) WHERE status='completed'"
)
@contextmanager
def _transaction() -> Iterator[sqlite3.Connection]:
from hermes_cli.sqlite_util import transaction
with _lock, transaction(_connect()) as conn:
yield conn
def _fetch(conn: sqlite3.Connection, execution_id: str) -> Optional[Dict[str, Any]]:
row = conn.execute("SELECT * FROM executions WHERE id=?", (execution_id,)).fetchone()
return dict(row) if row is not None else None
def _emit_execution_state(
record: Optional[Dict[str, Any]], *, delivery_outcome: Optional[str] = None
) -> None:
"""Project durable state to monitoring without affecting ledger behavior."""
try:
from agent.monitoring.cron_health import emit_execution_state
emit_execution_state(record, delivery_outcome=delivery_outcome)
except Exception:
pass
def _process_start_time(pid: int) -> Optional[int]:
try:
from gateway.status import get_process_start_time
return get_process_start_time(pid)
except Exception:
return None
def _owner_is_live(pid: int, started_at: Optional[int]) -> bool:
try:
from gateway.status import _pid_exists
if not _pid_exists(pid):
return False
except Exception:
return True # fail safe: inability to prove death must not rewrite state
if started_at is None:
return pid == os.getpid()
current = _process_start_time(pid)
return current is not None and current == started_at
def _prune_unlocked(conn: sqlite3.Connection) -> None:
conn.execute(
"""DELETE FROM executions WHERE id IN (
SELECT id FROM executions
WHERE status IN ('completed','failed','unknown')
ORDER BY finished_at DESC, claimed_at DESC, id DESC LIMIT -1 OFFSET ?
)""",
(max(0, int(MAX_TERMINAL_EXECUTIONS)),),
)
def create_execution(
job_id: str, *, source: str, scheduled_instant: Optional[str] = None,
) -> Dict[str, Any]:
"""Persist a claimed attempt before executor/provider dispatch."""
from cron.occurrences import scheduled_instant as canonical_instant
now = _hermes_now().isoformat()
execution_id = uuid.uuid4().hex
pid = os.getpid()
with _transaction() as conn:
conn.execute(
"""INSERT INTO executions
(id, job_id, source, process_id, pid, process_started_at,
status, claimed_at, scheduled_instant)
VALUES (?, ?, ?, ?, ?, ?, 'claimed', ?, ?)""",
(execution_id, str(job_id), str(source), _PROCESS_ID, pid,
_process_start_time(pid), now, canonical_instant(scheduled_instant)),
)
record = _fetch(conn, execution_id)
_emit_execution_state(record)
return record # type: ignore[return-value]
def set_execution_occurrence(execution_id: str, instant: Optional[str]) -> None:
"""Bind the store-claimed snapshot before a provider hands it to a worker."""
from cron.occurrences import scheduled_instant
with _transaction() as conn:
cur = conn.execute(
"UPDATE executions SET scheduled_instant=? WHERE id=? AND status='claimed' "
"AND handoff_pending=0 AND process_id=? AND pid=?",
(scheduled_instant(instant), execution_id, _PROCESS_ID, os.getpid()),
)
if cur.rowcount != 1:
raise RuntimeError("Cron occurrence could not be bound before dispatch")
def mark_execution_handoff_pending(execution_id: str) -> Optional[Dict[str, Any]]:
"""Fence restart recovery while an external worker is adopting a claim."""
with _transaction() as conn:
cur = conn.execute(
"""UPDATE executions
SET handoff_pending=1, handoff_started_at=?
WHERE id=? AND status='claimed'
AND process_id=? AND pid=?""",
(time.time(), execution_id, _PROCESS_ID, os.getpid()),
)
if cur.rowcount != 1:
return None
record = _fetch(conn, execution_id)
_emit_execution_state(record)
return record
def adopt_claimed_execution(execution_id: str) -> Optional[Dict[str, Any]]:
"""Atomically transfer and start an attempt in its worker process.
The dispatching gateway creates the row before spawning a restart-safe
worker. Adoption is the single ``claimed`` → ``running`` gate: only the
winner may acknowledge ownership or run side effects.
"""
pid = os.getpid()
process_started_at = _process_start_time(pid)
now = _hermes_now().isoformat()
with _transaction() as conn:
cur = conn.execute(
"""UPDATE executions
SET process_id=?, pid=?, process_started_at=?,
status='running', started_at=?, handoff_pending=0,
handoff_started_at=NULL
WHERE id=? AND status='claimed' AND handoff_pending=1""",
(_PROCESS_ID, pid, process_started_at, now, execution_id),
)
if cur.rowcount != 1:
return None
record = _fetch(conn, execution_id)
_emit_execution_state(record)
return record
def mark_execution_running(execution_id: str) -> Optional[Dict[str, Any]]:
"""Transition one claimed attempt to running exactly once."""
now = _hermes_now().isoformat()
with _transaction() as conn:
cur = conn.execute(
"""UPDATE executions
SET status='running', started_at=?, handoff_pending=0,
handoff_started_at=NULL
WHERE id=? AND status='claimed' AND handoff_pending=0
AND process_id=? AND pid=?""",
(now, execution_id, _PROCESS_ID, os.getpid()),
)
if cur.rowcount != 1:
return None
record = _fetch(conn, execution_id)
_emit_execution_state(record)
return record
def finish_execution(
execution_id: str, *, success: bool, error: Optional[str] = None,
delivery_outcome: Optional[str] = None,
) -> Optional[Dict[str, Any]]:
"""Write a terminal result once; terminal attempts cannot be rewritten."""
now = _hermes_now().isoformat()
status = "completed" if success else "failed"
detail = None if success else (str(error) if error else "unknown failure")
with _transaction() as conn:
cur = conn.execute(
"""UPDATE executions
SET status=?, finished_at=?, error=?, handoff_pending=0,
handoff_started_at=NULL, delivery_outcome=?
WHERE id=? AND status IN ('claimed','running')
AND process_id=? AND pid=?""",
(status, now, detail, delivery_outcome, execution_id, _PROCESS_ID, os.getpid()),
)
if cur.rowcount != 1:
return None
_prune_unlocked(conn)
record = _fetch(conn, execution_id)
_emit_execution_state(record, delivery_outcome=delivery_outcome)
return record
def recover_interrupted_executions() -> int:
"""Mark provably abandoned attempts unknown without scheduling retries."""
now = _hermes_now().isoformat()
changed = 0
recovered: List[Dict[str, Any]] = []
with _transaction() as conn:
rows = conn.execute(
"""SELECT id, status, process_id, pid, process_started_at,
handoff_pending, handoff_started_at
FROM executions
WHERE status IN ('claimed','running')"""
).fetchall()
for row in rows:
if row["process_id"] == _PROCESS_ID:
continue
if _owner_is_live(int(row["pid"]), row["process_started_at"]):
continue
handoff_started_at = row["handoff_started_at"]
if (
row["handoff_pending"]
and handoff_started_at is not None
and time.time() - float(handoff_started_at)
< HANDOFF_ADOPTION_GRACE_SECONDS
):
continue
cur = conn.execute(
"""UPDATE executions
SET status='unknown', finished_at=?, error=?,
handoff_pending=0, handoff_started_at=NULL
WHERE id=? AND status=? AND process_id=? AND pid=?
AND handoff_pending=?
AND handoff_started_at IS ?""",
(now,
"Scheduler restarted after this execution's owner exited before a durable "
"terminal state; whether side effects ran is unknown.",
row["id"], row["status"], row["process_id"], row["pid"],
row["handoff_pending"], row["handoff_started_at"]),
)
changed += cur.rowcount
if cur.rowcount:
record = _fetch(conn, row["id"])
if record is not None:
recovered.append(record)
if changed:
_prune_unlocked(conn)
for record in recovered:
_emit_execution_state(record)
return changed
def list_executions(
*, job_id: Optional[str] = None, limit: int = 50, before_claimed_at: Optional[str] = None,
) -> List[Dict[str, Any]]:
"""Return indexed, newest-first execution history with cursor pagination."""
clauses: List[str] = []
params: List[Any] = []
if job_id is not None:
clauses.append("job_id=?")
params.append(str(job_id))
if before_claimed_at is not None:
clauses.append("claimed_at < ?")
params.append(str(before_claimed_at))
where = " WHERE " + " AND ".join(clauses) if clauses else ""
params.append(max(1, min(int(limit), 500)))
with _transaction() as conn:
rows = conn.execute(
"SELECT * FROM executions" + where
+ " ORDER BY claimed_at DESC, id DESC LIMIT ?",
params,
).fetchall()
return [dict(row) for row in rows]
def get_execution(execution_id: str) -> Optional[Dict[str, Any]]:
"""Return one exact execution attempt, or ``None`` when it is absent."""
with _transaction() as conn:
row = conn.execute(
"SELECT * FROM executions WHERE id=?",
(str(execution_id),),
).fetchone()
return dict(row) if row is not None else None
def latest_execution(job_id: str) -> Optional[Dict[str, Any]]:
rows = list_executions(job_id=job_id, limit=1)
return rows[0] if rows else None
def latest_executions(job_ids: List[str]) -> Dict[str, Dict[str, Any]]:
"""Load latest execution for many jobs in one indexed query."""
clean = [str(job_id) for job_id in dict.fromkeys(job_ids) if job_id]
if not clean:
return {}
placeholders = ",".join("?" for _ in clean)
with _transaction() as conn:
rows = conn.execute(
f"""SELECT e.* FROM executions e
WHERE e.job_id IN ({placeholders})
AND e.id=(SELECT e2.id FROM executions e2
WHERE e2.job_id=e.job_id
ORDER BY e2.claimed_at DESC, e2.id DESC LIMIT 1)""",
clean,
).fetchall()
return {row["job_id"]: dict(row) for row in rows}
|