Download cron/executions.py from SaylorTwift/hermes-agent: direct link, hf CLI and curl.
- Browser
- Download file 14.5 kB
-
https://huggingface.co/SaylorTwift/hermes-agent/resolve/main/cron/executions.py
- Command line
-
hf download hf://SaylorTwift/hermes-agent/cron/executions.py
-
curl -L -o executions.py https://huggingface.co/SaylorTwift/hermes-agent/resolve/main/cron/executions.py
14.5 kB
| """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'" | |
| ) | |
| 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} | |