diff --git a/cron/jobs.py b/cron/jobs.py index 91e1ab67a0..d867658a9e 100644 --- a/cron/jobs.py +++ b/cron/jobs.py @@ -52,9 +52,8 @@ def _ensure_croniter() -> bool: HAS_CRONITER = False return bool(HAS_CRONITER) -# ============================================================================= -# Configuration -# ============================================================================= + +# --- Configuration --- # Cron is per-profile by design: anchor at get_hermes_home() (active profile home), NOT # get_default_hermes_root() — the shared root would funnel every profile's jobs into one jobs.json @@ -64,10 +63,9 @@ HERMES_DIR = get_hermes_home().resolve() # scope paths with use_cron_store() instead of mutating these process-wide. CRON_DIR = HERMES_DIR / "cron" JOBS_FILE = CRON_DIR / "jobs.json" -# Touched by the ticker every loop so `hermes cron status` can tell the ticker THREAD is alive, -# not just the gateway PROCESS (a silently dead ticker would otherwise report healthy). +# Heartbeat: touched every ticker loop so `hermes cron status` can tell the ticker THREAD is alive, +# not just the gateway PROCESS. Success: last tick that completed WITHOUT raising. TICKER_HEARTBEAT_FILE = CRON_DIR / "ticker_heartbeat" -# Last tick that completed WITHOUT raising — lets status detect a ticker alive but failing every tick. TICKER_SUCCESS_FILE = CRON_DIR / "ticker_last_success" # Single source of truth for the ticker interval (scheduler_provider.py) and the staleness # threshold in `hermes cron status` (hermes_cli/cron.py), so they never drift apart. @@ -101,10 +99,7 @@ class _CronStorePaths: _cron_store_override: ContextVar[Optional[_CronStorePaths]] = ContextVar( - "cron_store_override", - default=None, -) - + "cron_store_override", default=None) # Import-time snapshot so deliberate re-pointing of CRON_DIR/JOBS_FILE/OUTPUT_DIR (the documented # escape hatch for tests/embedders) is distinguishable from the constants merely being stale. @@ -112,12 +107,11 @@ _IMPORT_STORE = _CronStorePaths(CRON_DIR, JOBS_FILE, OUTPUT_DIR) def _current_cron_store() -> _CronStorePaths: - """Return paths pinned to this execution context's profile. + """Paths pinned to this execution context's profile. - Precedence: (1) active use_cron_store() override; (2) deliberately re-pointed module constants - (differ from import-time values → honor the compatibility surface); (3) the ACTIVE profile home - resolved fresh via get_hermes_home(), so re-pointing HERMES_HOME after import reads/writes ITS - OWN store rather than the user's real jobs.json frozen at import; (4) the import-time constants. + Precedence: (1) active use_cron_store() override; (2) deliberately re-pointed module constants; + (3) the ACTIVE profile home via get_hermes_home(), so re-pointing HERMES_HOME after import uses + ITS OWN store rather than the user's real jobs.json frozen at import; (4) import-time constants. """ override = _cron_store_override.get() if override is not None: @@ -134,7 +128,8 @@ def _current_cron_store() -> _CronStorePaths: @contextlib.contextmanager def use_cron_store(home: Union[str, Path]): """Route cron storage to ``home`` without mutating process globals.""" - token = _cron_store_override.set(_CronStorePaths.for_dir(Path(home).expanduser().resolve() / "cron")) + token = _cron_store_override.set( + _CronStorePaths.for_dir(Path(home).expanduser().resolve() / "cron")) try: yield finally: @@ -159,12 +154,8 @@ _DEFAULT_CRON_INACTIVITY_TIMEOUT = 600.0 def _oneshot_run_claim_ttl_seconds() -> float: - """Resolve the one-shot running-claim TTL from ``HERMES_CRON_TIMEOUT``. - - unset/invalid → 600s → 1800s; ``0`` (unlimited) → ``ONESHOT_RUN_CLAIM_TTL_SECONDS``; - positive N → ``max(N * headroom, ONESHOT_RUN_CLAIM_TTL_SECONDS)`` so a tiny timeout never - expires a claim mid-run. - """ + """One-shot running-claim TTL from ``HERMES_CRON_TIMEOUT``: unset/invalid → 600s → 1800s; + ``0`` (unlimited) → the fixed floor; positive N → ``max(N * headroom, floor)``.""" raw = os.getenv("HERMES_CRON_TIMEOUT", "").strip() timeout = _DEFAULT_CRON_INACTIVITY_TIMEOUT if raw: @@ -173,20 +164,15 @@ def _oneshot_run_claim_ttl_seconds() -> float: except (ValueError, TypeError): timeout = _DEFAULT_CRON_INACTIVITY_TIMEOUT if timeout <= 0: - # Unlimited runs — cannot bound; use the fixed fallback floor. return float(ONESHOT_RUN_CLAIM_TTL_SECONDS) - return max( - timeout * _ONESHOT_RUN_CLAIM_TTL_HEADROOM, - float(ONESHOT_RUN_CLAIM_TTL_SECONDS), - ) + return max(timeout * _ONESHOT_RUN_CLAIM_TTL_HEADROOM, float(ONESHOT_RUN_CLAIM_TTL_SECONDS)) def _job_running_in_this_process(job_id: str) -> bool: - """Return True when the scheduler in THIS process is still running ``job_id``. + """True when the scheduler in THIS process is still running ``job_id``. - The run_claim TTL alone cannot distinguish "claiming tick died" from "alive but slow" (a run - stalled on I/O or a slept laptop legitimately outlives it); the in-process running set settles - the single-gateway case. Lazy import: the scheduler imports this module, so top-level would be circular. + The run_claim TTL alone cannot distinguish "claiming tick died" from "alive but slow"; the + in-process running set settles the single-gateway case. Lazy import: scheduler imports us. """ try: from cron.scheduler import get_running_job_ids @@ -195,9 +181,7 @@ def _job_running_in_this_process(job_id: str) -> bool: logger.warning( "Cron running-set liveness check failed for job %r; keeping the " "entry to avoid deleting a possibly live one-shot run", - job_id, - exc_info=True, - ) + job_id, exc_info=True) return True @@ -279,13 +263,9 @@ def _jobs_lock(): "jobs lock (%s) — another process is holding " "it. Proceeding with in-process locking only " "so the scheduler stays alive (#60703).", - _JOBS_LOCK_TIMEOUT_SECONDS, - _jobs_lock_file(), - ) - try: + _JOBS_LOCK_TIMEOUT_SECONDS, _jobs_lock_file()) + with contextlib.suppress(OSError): lock_fd.close() - except OSError: - pass lock_fd = None except (OSError, IOError) as e: # A locking failure must never take down cron writes — in-process lock still held. @@ -331,8 +311,7 @@ def _fire_job_lock(job_id: str): try: ensure_dirs() - lock_name = uuid.uuid5(uuid.NAMESPACE_URL, lock_key).hex - lock_path = cron_dir / f".fire-{lock_name}.lock" + lock_path = cron_dir / f".fire-{uuid.uuid5(uuid.NAMESPACE_URL, lock_key).hex}.lock" lock_fd = None acquired = False try: @@ -361,6 +340,14 @@ def _fire_job_lock(job_id: str): local_lock.release() +def _under_fire_fence(job_id: str, fn: Callable[[], Any]) -> Any: + """Run ``fn()`` holding the job's fire fence; False (fail closed) when it can't be acquired.""" + with _fire_job_lock(job_id) as acquired: + if not acquired: + return False + return fn() + + @contextlib.contextmanager def fire_claim_fence(job_id: str, *, expected_owner: str): """Hold a per-job fence while an owner performs an external side effect.""" @@ -371,22 +358,18 @@ def fire_claim_fence(job_id: str, *, expected_owner: str): with _jobs_lock(): job = next((item for item in load_jobs() if item.get("id") == job_id), None) claim = job.get("fire_claim") if isinstance(job, dict) else None - owns_claim = ( - isinstance(claim, dict) and claim.get("by") == expected_owner - ) + owns_claim = isinstance(claim, dict) and claim.get("by") == expected_owner yield owns_claim + # Fields that must never change after creation: ``id`` is a path component under OUTPUT_DIR, so an # update could leak ``../escape``/absolute/nested values into output writes/deletes. _IMMUTABLE_JOB_FIELDS = frozenset({"id"}) def _job_output_dir(job_id: str) -> Path: - """Resolve a job's output directory, rejecting any path-escape attempt. - - IDs containing ``..``, absolute paths, or separators would let output writes/deletes escape the - sandbox; only a single safe path component is accepted. - """ + """Resolve a job's output directory, rejecting any path-escape attempt (``..``, absolute + paths, separators): only a single safe path component is accepted.""" text = str(job_id or "").strip() if not text or text in {".", ".."} or "/" in text or "\\" in text: raise ValueError(f"Invalid cron job id for output path: {job_id!r}") @@ -403,7 +386,6 @@ def _normalize_skill_list(skill: Optional[str] = None, skills: Optional[Any] = N raw_items = [skills] else: raw_items = list(skills) - normalized: List[str] = [] for item in raw_items: text = str(item or "").strip() @@ -423,9 +405,7 @@ def _apply_skill_fields(job: Dict[str, Any]) -> Dict[str, Any]: def _coerce_job_text(value: Any, fallback: str = "") -> str: """Coerce legacy/hand-edited nullable cron fields to strings for readers.""" - if value is None: - return fallback - return str(value) + return fallback if value is None else str(value) # Fields whose presence in an update can turn a runnable job into an empty one. @@ -443,17 +423,12 @@ NO_AGENT_WITHOUT_SCRIPT_ERROR = ( def job_payload_is_empty(job: Dict[str, Any]) -> bool: - """True when a job record has nothing runnable (blank prompt, no script, no skills). - - ``no_agent`` needs no special case — it already requires a script. - """ - if _coerce_job_text(job.get("prompt")).strip(): - return False - if _coerce_job_text(job.get("script")).strip(): + """True when a job record has nothing runnable (blank prompt, no script, no skills) AND at + least one payload field is explicitly present. ``no_agent`` already requires a script.""" + if _coerce_job_text(job.get("prompt")).strip() or _coerce_job_text(job.get("script")).strip(): return False if _normalize_skill_list(job.get("skill"), job.get("skills")): return False - # Only flag if at least one payload field is explicitly present in the record return any(k in job for k in ("prompt", "script", "skill", "skills")) @@ -461,7 +436,6 @@ def _schedule_display_for_job(job: Dict[str, Any]) -> str: display = _coerce_job_text(job.get("schedule_display")).strip() if display: return display - schedule = job.get("schedule") if isinstance(schedule, dict): for key in ("display", "value", "expr", "run_at"): @@ -470,53 +444,44 @@ def _schedule_display_for_job(job: Dict[str, Any]) -> str: return text elif schedule is not None: return str(schedule) - return "?" def _normalize_job_record(job: Dict[str, Any]) -> Dict[str, Any]: - """Return a read-safe job shape: legacy/hand-edited records may have nullable ``prompt``, - ``name``, ``schedule_display``. Storage is untouched; consumers never crash on formatting.""" + """Read-safe job shape: legacy/hand-edited records may have nullable ``prompt``, ``name``, + ``schedule_display``. Storage is untouched; consumers never crash on formatting.""" normalized = _apply_skill_fields(job) job_id = _coerce_job_text(normalized.get("id"), "unknown") prompt = _coerce_job_text(normalized.get("prompt")) normalized["id"] = job_id normalized["prompt"] = prompt - name = _coerce_job_text(normalized.get("name")).strip() if not name: - script = _coerce_job_text(normalized.get("script")).strip() label_source = ( prompt or (normalized["skills"][0] if normalized.get("skills") else "") - or script + or _coerce_job_text(normalized.get("script")).strip() or job_id or "cron job" ) name = label_source[:50].strip() or "cron job" normalized["name"] = name normalized["schedule_display"] = _schedule_display_for_job(normalized) - # Derived from the scheduler-honoured ``enabled`` flag so a half-paused record cannot render # "paused" while still firing. See effective_job_state(). normalized["state"] = effective_job_state(normalized) - return normalized def _has_pause_marker(job: Dict[str, Any]) -> bool: """True when the record carries any operator-facing pause signal.""" - if _coerce_job_text(job.get("state")).strip() == "paused": - return True - return bool(job.get("paused_at")) + return _coerce_job_text(job.get("state")).strip() == "paused" or bool(job.get("paused_at")) def is_job_runnable(job: Dict[str, Any]) -> bool: """True iff the scheduler may fire this job: ``enabled`` plus pause markers as a second gate so a contradictory half-paused record never fires even before self-heal runs.""" - if not job.get("enabled", True): - return False - return not _has_pause_marker(job) + return bool(job.get("enabled", True)) and not _has_pause_marker(job) def effective_job_state(job: Dict[str, Any]) -> str: @@ -544,11 +509,9 @@ def _is_recoverable_error_job(job: Dict[str, Any]) -> bool: """True for a recurring job stuck in ``state=error``. ``state=error`` is set ONLY when ``compute_next_run()`` fails for a cron/interval job (croniter - missing, malformed schedule); recurring jobs must NEVER be silently disabled, and unlike - ``completed`` such a job still has future occurrences once the issue resolves. Treating it as - terminal would block the due-scan self-heal, at-most-once pre-advance, dispatch claim, and - ``resume_job`` — wedging it forever. Use ``is_terminal_job()`` for "truly done"; exclude this - case for "can still reach a future occurrence". + missing, malformed schedule); such a job still has future occurrences once the issue resolves, + so treating it as terminal would block due-scan self-heal, pre-advance, dispatch claim and + ``resume_job`` — wedging it forever. ``is_terminal_job()`` alone means "truly done". """ return ( job.get("state") == "error" @@ -557,20 +520,16 @@ def _is_recoverable_error_job(job: Dict[str, Any]) -> bool: def _secure_dir(path: Path): - """Set directory to owner-only access (0700). No-op on Windows.""" - try: + """Set directory to owner-only access (0700). No-op where chmod is unsupported (Windows).""" + with contextlib.suppress(OSError, NotImplementedError): os.chmod(path, 0o700) - except (OSError, NotImplementedError): - pass # Windows or other platforms where chmod is not supported def _secure_file(path: Path): - """Set file to owner-only read/write (0600). No-op on Windows.""" - try: + """Set file to owner-only read/write (0600). No-op where chmod is unsupported (Windows).""" + with contextlib.suppress(OSError, NotImplementedError): if path.exists(): os.chmod(path, 0o600) - except (OSError, NotImplementedError): - pass def _preserve_file_ownership(path: Path, before: Optional[os.stat_result]) -> None: @@ -582,47 +541,32 @@ def _preserve_file_ownership(path: Path, before: Optional[os.stat_result]) -> No """ if before is None or os.name != "posix": return - geteuid = getattr(os, "geteuid", None) - getegid = getattr(os, "getegid", None) - if geteuid is None or getegid is None: - return try: - euid = geteuid() - if euid != 0: - return # unprivileged writer — nothing to (or we could) restore - if (before.st_uid, before.st_gid) == (euid, getegid()): - return # already ours before the rewrite — nothing changed + euid = os.geteuid() + if euid != 0 or (before.st_uid, before.st_gid) == (euid, os.getegid()): + return # unprivileged writer, or already ours before the rewrite os.chown(path, before.st_uid, before.st_gid) except OSError as e: logger.warning( "Could not restore ownership of %s to uid=%s gid=%s after rewrite: %s " "— if the gateway runs as a different user, its cron ticker may now " "be locked out (see issue #68483).", - path, before.st_uid, before.st_gid, e, - ) + path, before.st_uid, before.st_gid, e) def _is_named_profile_path(path: Path) -> bool: """True if *path* is under ``/profiles//`` (default/custom homes are not). - - Checks the resolved path (symlinked parents) and the raw path (symlinked profile homes whose - target no longer contains ``profiles``). - """ - try: + Checks the resolved path (symlinked parents) and the raw path (symlinked profile homes).""" + with contextlib.suppress(OSError, RuntimeError): if "profiles" in path.resolve().parts: return True - except (OSError, RuntimeError): - pass return "profiles" in path.parts def _ensure_cron_dir(cron_dir: Path) -> None: - """Create a cron directory without resurrecting a deleted profile home. - - A stale multiplex scheduler may still hold a deleted profile's path; ``parents=False`` makes - that race fail closed instead of restoring the tree. Default/custom homes keep ``parents=True`` - so first-run creation works. - """ + """Create a cron directory without resurrecting a deleted profile home: a stale multiplex + scheduler may still hold a deleted profile's path, so named profiles use ``parents=False`` and + fail closed. Default/custom homes keep ``parents=True`` so first-run creation works.""" if _is_named_profile_path(cron_dir): cron_dir.mkdir(exist_ok=True) return @@ -638,16 +582,11 @@ def ensure_dirs(): _secure_dir(store.output_dir) -# ============================================================================= -# Schedule Parsing -# ============================================================================= +# --- Schedule Parsing --- def normalize_repeat_value(repeat: Any) -> Optional[int]: - """Coerce a repeat value (int or user-facing string) into ``Optional[int]``. - - ``'forever'``-family -> None (infinite), ``'once'``-family -> 1, numeric -> int, 0/negative -> - None, anything else -> ValueError (never store garbage that breaks ``mark_job_run`` later). - """ + """Coerce a repeat value (int or user-facing string) into ``Optional[int]``: ``'forever'``-family + -> None, ``'once'``-family -> 1, numeric -> int, 0/negative -> None, else ValueError.""" if repeat is None: return None if isinstance(repeat, str): @@ -666,6 +605,9 @@ def normalize_repeat_value(repeat: Any) -> Optional[int]: return None if repeat <= 0 else int(repeat) +_DURATION_MULTIPLIERS = {'m': 1, 'h': 60, 'd': 1440} + + def parse_duration(s: str) -> int: """Parse a duration string into minutes: "30m" → 30, "2h" → 120, "1d" → 1440, bare "hour" → 60.""" s = s.strip().lower() @@ -673,14 +615,9 @@ def parse_duration(s: str) -> int: if not match: raise ValueError( f"Invalid duration: '{s}'. Use format like '30m', '2h', '1d', " - "or a bare unit like 'hour' (defaults to 1)." - ) - + "or a bare unit like 'hour' (defaults to 1).") value = int(match.group(1)) if match.group(1) else 1 - unit = match.group(2)[0] # First char: m, h, or d - - multipliers = {'m': 1, 'h': 60, 'd': 1440} - return value * multipliers[unit] + return value * _DURATION_MULTIPLIERS[match.group(2)[0]] # Day-spec phrases for "every monday 9am" / "every day at 9am". Cron weekday numbering is @@ -724,7 +661,7 @@ def _parse_clock_time(text: str) -> Optional[tuple]: return None if meridiem == "am": hour = 0 if hour == 12 else hour - else: # pm + else: hour = 12 if hour == 12 else hour + 12 if hour > 23 or minute > 59: return None @@ -737,10 +674,8 @@ def _natural_every_to_cron(rest: str) -> Optional[str]: tokens = rest.lower().replace(",", " ").split() if not tokens: return None - # Leading day tokens: a keyword spec ("weekdays") or a comma/"and"-separated weekday list. - day_token = tokens[0] - dow = _DAYSPEC_TO_CRON_DOW.get(day_token) + dow = _DAYSPEC_TO_CRON_DOW.get(tokens[0]) idx = 1 if dow is None: days = [] @@ -759,14 +694,11 @@ def _natural_every_to_cron(rest: str) -> Optional[str]: return None dow = ",".join(days) idx -= 1 - time_tokens = tokens[idx:] - # Optional "at" separator: "every day at 9am". - if time_tokens and time_tokens[0] == "at": + if time_tokens and time_tokens[0] == "at": # optional separator: "every day at 9am" time_tokens = time_tokens[1:] if not time_tokens: return None - parsed = _parse_clock_time(" ".join(time_tokens)) if parsed is None: return None @@ -785,6 +717,10 @@ def _cron_schedule(expr: str, display: str, missing_croniter: str, invalid_label return {"kind": "cron", "expr": expr, "display": display} +def _interval_schedule(minutes: int) -> Dict[str, Any]: + return {"kind": "interval", "minutes": minutes, "display": f"every {minutes}m"} + + def parse_schedule(schedule: str) -> Dict[str, Any]: """Parse a schedule string into ``{"kind": "once"|"interval"|"cron", ...}`` with ``run_at`` / ``minutes`` / ``expr``. "30m" and "every 30m" are recurring intervals; "every monday 9am" and @@ -801,10 +737,8 @@ def parse_schedule(schedule: str) -> Dict[str, Any]: return _cron_schedule( cron_expr, original, "Weekday/time schedules like 'every monday 9am' require the 'croniter' package.", - "schedule", - ) - minutes = parse_duration(rest) - return {"kind": "interval", "minutes": minutes, "display": f"every {minutes}m"} + "schedule") + return _interval_schedule(parse_duration(rest)) # No-"every" phrases ("weekdays at 9am", "daily at 7am"): same shape sans prefix. cron_expr = _natural_every_to_cron(schedule_lower) @@ -812,24 +746,21 @@ def parse_schedule(schedule: str) -> Dict[str, Any]: return _cron_schedule( cron_expr, original, "Weekday/time schedules like 'weekdays at 9am' require the 'croniter' package.", - "schedule", - ) + "schedule") # Cron expression (5-6 fields). Letters are allowed so named months/weekdays (JAN-DEC, MON-FRI) - # reach croniter, which supports them; a digit-only pattern silently rejected them. + # reach croniter, which supports them. parts = schedule.split() if len(parts) >= 5 and all(re.match(r'^[A-Za-z\d\*\-,/]+$', p) for p in parts[:5]): return _cron_schedule( - schedule, schedule, "Cron expressions require 'croniter' package.", "cron expression" - ) + schedule, schedule, "Cron expressions require 'croniter' package.", "cron expression") # ISO timestamp (contains T or looks like date) if 'T' in schedule or re.match(r'^\d{4}-\d{2}-\d{2}', schedule): try: dt = datetime.fromisoformat(schedule.replace('Z', '+00:00')) - # Make naive timestamps aware at parse time, anchored to the CONFIGURED Hermes timezone - # (not server-local): the due-check compares against hermes_time.now(), so a - # server-local reading would land hours off the user's wall-clock intent. + # Naive timestamps become aware in the CONFIGURED Hermes timezone (not server-local): + # the due-check compares against hermes_time.now(). if dt.tzinfo is None: dt = dt.replace(tzinfo=_hermes_now().tzinfo) return { @@ -840,23 +771,19 @@ def parse_schedule(schedule: str) -> Dict[str, Any]: except ValueError as e: raise ValueError(f"Invalid timestamp '{schedule}': {e}") - # Bare duration ("30m") → RECURRING interval per the documented tool contract; the explicit - # one-shot-by-duration form is "in 30m"/"in 2h". + # "in 30m"/"in 2h" is the explicit one-shot-by-duration form; a bare duration ("30m") is a + # RECURRING interval per the documented tool contract. if schedule_lower.startswith("in "): duration_str = schedule[3:].strip() try: minutes = parse_duration(duration_str) except ValueError: raise ValueError( - f"Invalid duration '{duration_str}' after 'in '. Use e.g. 'in 30m', 'in 2h'." - ) + f"Invalid duration '{duration_str}' after 'in '. Use e.g. 'in 30m', 'in 2h'.") run_at = _hermes_now() + timedelta(minutes=minutes) return {"kind": "once", "run_at": run_at.isoformat(), "display": f"once in {duration_str}"} - try: - minutes = parse_duration(schedule) - return {"kind": "interval", "minutes": minutes, "display": f"every {minutes}m"} - except ValueError: - pass + with contextlib.suppress(ValueError): + return _interval_schedule(parse_duration(schedule)) raise ValueError( f"Invalid schedule '{original}'. Use:\n" @@ -866,14 +793,15 @@ def parse_schedule(schedule: str) -> Dict[str, Any]: f" - Cron: '0 9 * * *' (cron expression)\n" f" - Timestamp: '2026-02-03T14:00:00' (one-shot at time)" ) + + def _ensure_aware(dt: datetime) -> datetime: - """Return an aware datetime in the configured Hermes timezone. Legacy naive values are read as + """Aware datetime in the configured Hermes timezone. Legacy naive values are read as *system-local* wall time (what created them) then converted, preserving ordering across timezone changes and avoiding false not-due results.""" target_tz = _hermes_now().tzinfo if dt.tzinfo is None: - local_tz = datetime.now().astimezone().tzinfo - return dt.replace(tzinfo=local_tz).astimezone(target_tz) + return dt.replace(tzinfo=datetime.now().astimezone().tzinfo).astimezone(target_tz) return dt.astimezone(target_tz) @@ -894,27 +822,19 @@ def _timezone_offset_mismatch(stored: datetime, current: datetime) -> bool: def _stored_wall_clock_is_future(stored: datetime, current: datetime) -> bool: - """True when the stored local wall-clock time has not arrived yet. - - Cron expresses wall-clock intent; after a timezone change an old offset can make a future run - look due (21:00+10 → 13:00+02). Comparing naive wall clocks separates that from a genuine miss. - """ + """True when the stored local wall-clock time has not arrived yet. Cron expresses wall-clock + intent; after a timezone change an old offset can make a future run look due (21:00+10 → + 13:00+02). Comparing naive wall clocks separates that from a genuine miss.""" return stored.replace(tzinfo=None) > current.replace(tzinfo=None) def _recoverable_oneshot_run_at( - schedule: Dict[str, Any], - now: datetime, - *, - last_run_at: Optional[str] = None, + schedule: Dict[str, Any], now: datetime, *, last_run_at: Optional[str] = None, ) -> Optional[str]: - """Return a one-shot run time if still eligible: a small grace window covers jobs created just - after their minute; once run, a one-shot is never eligible again.""" - if not isinstance(schedule, dict) or schedule.get("kind") != "once": + """One-shot run time if still eligible: a small grace window covers jobs created just after + their minute; once run, a one-shot is never eligible again.""" + if not isinstance(schedule, dict) or schedule.get("kind") != "once" or last_run_at: return None - if last_run_at: - return None - run_at = schedule.get("run_at") run_at_dt = _parse_aware(run_at) if run_at else None if run_at_dt is not None and run_at_dt >= now - timedelta(seconds=ONESHOT_GRACE_SECONDS): @@ -922,17 +842,17 @@ def _recoverable_oneshot_run_at( return None +_MIN_GRACE_SECONDS = 120 +_MAX_GRACE_SECONDS = 7200 + + def _compute_grace_seconds(schedule: dict) -> int: """How late a job can be and still catch up rather than fast-forward: half the period, clamped to [120s, 2h], so daily jobs catch up but frequent jobs fast-forward quickly.""" - MIN_GRACE = 120 - MAX_GRACE = 7200 # 2 hours - period_seconds = _schedule_cadence_seconds(schedule) if not period_seconds: - return MIN_GRACE - grace = int(period_seconds) // 2 - return max(MIN_GRACE, min(grace, MAX_GRACE)) + return _MIN_GRACE_SECONDS + return max(_MIN_GRACE_SECONDS, min(int(period_seconds) // 2, _MAX_GRACE_SECONDS)) # A recurring dispatch within this many seconds of schedule renders "on time": a busy once-a-minute @@ -958,17 +878,12 @@ _persisted_error_recoveries_recent: list = [] def _job_is_stale_error_recurring( - job: Dict[str, Any], - schedule: Dict[str, Any], - now: datetime, + job: Dict[str, Any], schedule: Dict[str, Any], now: datetime, ) -> bool: - """True when a recurring job (caller-checked) is wedged in a stale persisted error state. - - Requires ``last_status == "error"``, not running in this process (a live run must never be - re-armed underneath itself), and ``last_run_at`` older than ``cadence + grace``. The age is the - key discriminator: mark_job_run stamps ``last_run_at`` on every fire, so a job merely - erroring-and-retrying on schedule stays fresh and is not flagged. - """ + """True when a recurring job (caller-checked) is wedged in a stale persisted error state: + ``last_status == "error"``, not running in this process (a live run must never be re-armed + underneath itself), and ``last_run_at`` older than ``cadence + grace``. Age is the key + discriminator: a job merely erroring-and-retrying on schedule stays fresh and is not flagged.""" if job.get("last_status") != "error": return False if _job_running_in_this_process(str(job.get("id") or "")): @@ -980,19 +895,22 @@ def _job_is_stale_error_recurring( age_seconds = (now - last_run_dt).total_seconds() if age_seconds < 0: return False + grace = _compute_grace_seconds(schedule) cadence_seconds = _schedule_cadence_seconds(schedule) if cadence_seconds is None: # Unknown cadence: fall back to the grace window, never re-arming anything younger than it. - cadence_seconds = _compute_grace_seconds(schedule) - return age_seconds > (cadence_seconds + _compute_grace_seconds(schedule)) + cadence_seconds = grace + return age_seconds > (cadence_seconds + grace) + + +# Per-expr cache for _schedule_cadence_seconds' croniter measurements. +_cron_cadence_cache: Dict[str, Optional[float]] = {} def _schedule_cadence_seconds(schedule: Dict[str, Any]) -> Optional[float]: - """Approximate schedule period in seconds, or None (croniter missing / malformed expr). - - Cron results are cached per expr because this runs under ``_jobs_lock`` every tick; the gap can - vary with base time for irregular exprs, acceptable for a staleness *threshold*. - """ + """Approximate schedule period in seconds, or None (croniter missing / malformed expr). Cron + results are cached per expr because this runs under ``_jobs_lock`` every tick; the gap can vary + with base time for irregular exprs, acceptable for a staleness *threshold*.""" if not isinstance(schedule, dict): return None kind = schedule.get("kind") @@ -1002,33 +920,25 @@ def _schedule_cadence_seconds(schedule: Dict[str, Any]) -> Optional[float]: return float(minutes) * 60.0 if minutes else None except (TypeError, ValueError): return None - if kind == "cron": - if not _ensure_croniter(): - return None - expr = schedule.get("expr") - if not expr: - return None - if expr in _cron_cadence_cache: - return _cron_cadence_cache[expr] - try: - base = _hermes_now() - it = croniter(expr, base) - first = it.get_next(datetime) - second = it.get_next(datetime) - gap = (second - first).total_seconds() - result = gap if gap > 0 else None - except Exception: - result = None - # Hard bound so deleted/edited exprs can't grow the cache unboundedly in a long-lived gateway. - if len(_cron_cadence_cache) >= 256: - _cron_cadence_cache.clear() - _cron_cadence_cache[expr] = result - return result - return None - - -# Per-expr cache for _schedule_cadence_seconds' croniter measurements. -_cron_cadence_cache: Dict[str, Optional[float]] = {} + if kind != "cron" or not _ensure_croniter(): + return None + expr = schedule.get("expr") + if not expr: + return None + if expr in _cron_cadence_cache: + return _cron_cadence_cache[expr] + try: + it = croniter(expr, _hermes_now()) + first = it.get_next(datetime) + gap = (it.get_next(datetime) - first).total_seconds() + result = gap if gap > 0 else None + except Exception: + result = None + # Hard bound so deleted/edited exprs can't grow the cache unboundedly in a long-lived gateway. + if len(_cron_cadence_cache) >= 256: + _cron_cadence_cache.clear() + _cron_cadence_cache[expr] = result + return result def _append_telemetry_record(filename: str, entry: Dict[str, Any], recent: list) -> None: @@ -1057,8 +967,7 @@ def _record_persisted_error_recovery(job: Dict[str, Any], previous_next_run: str } _persisted_error_recoveries += 1 _append_telemetry_record( - "persisted_error_recoveries.jsonl", entry, _persisted_error_recoveries_recent - ) + "persisted_error_recoveries.jsonl", entry, _persisted_error_recoveries_recent) def get_persisted_error_recovery_stats() -> Dict[str, Any]: @@ -1069,16 +978,10 @@ def get_persisted_error_recovery_stats() -> Dict[str, Any]: } -def _cron_next_run_matches_expr( - schedule: Dict[str, Any], - next_run_dt: datetime, -) -> bool: - """Whether ``next_run_dt`` is an occurrence of the schedule's current expr. - - Detects a hand-edited ``schedule.expr`` whose stored ``next_run_at`` came from the old one. - Best-effort: anything uncheckable (non-cron, no expr, no croniter, malformed) reports a match - so the fire path keeps its existing semantics. - """ +def _cron_next_run_matches_expr(schedule: Dict[str, Any], next_run_dt: datetime) -> bool: + """Whether ``next_run_dt`` is an occurrence of the schedule's current expr (detects a hand-edited + ``schedule.expr`` whose stored ``next_run_at`` came from the old one). Best-effort: anything + uncheckable (non-cron, no expr, no croniter, malformed) reports a match.""" if schedule.get("kind") != "cron": return True expr = schedule.get("expr") @@ -1087,8 +990,7 @@ def _cron_next_run_matches_expr( try: # Last occurrence at-or-before the instant: base one second past it so an exact hit is # included, then compare at second granularity (croniter is second-precision). - base = next_run_dt + timedelta(seconds=1) - prev = croniter(str(expr), base).get_prev(datetime) + prev = croniter(str(expr), next_run_dt + timedelta(seconds=1)).get_prev(datetime) return abs((prev - next_run_dt).total_seconds()) < 1.0 except Exception: return True @@ -1101,25 +1003,20 @@ STALE_CRON_EXPR_EDIT = "expr_edit" def _classify_stale_cron_next_run( - schedule: Dict[str, Any], - raw_next_run_dt: datetime, - next_run_dt: datetime, + schedule: Dict[str, Any], raw_next_run_dt: datetime, next_run_dt: datetime, ) -> str: """Explain WHY a stored ``next_run_at`` misses the current cron lattice; the causes need opposite actions. - ``expr_edit``: a hand edit changed ``schedule.expr`` — the instant is excluded, so re-anchor - WITHOUT firing. ``timezone_migration``: only the offset representation changed (legacy UTC rows - normalized into the profile tz), and treating it as an edit would skip a due, never-fired - occurrence. Discriminator: whether normalization itself moved the wall clock — a stored instant - whose OWN wall clock is a legal occurrence is a migration. When offsets agree the wall clock is - unchanged, so a genuine expr edit can never be misread as a migration. + ``expr_edit``: a hand edit changed ``schedule.expr`` — re-anchor WITHOUT firing. + ``timezone_migration``: only the offset representation changed (legacy UTC rows normalized into + the profile tz); treating it as an edit would skip a due, never-fired occurrence. Discriminator: + a stored instant whose OWN wall clock is a legal occurrence, when normalization moved the wall + clock, is a migration. When offsets agree a genuine expr edit can never be misread as one. """ if _cron_next_run_matches_expr(schedule, next_run_dt): return STALE_CRON_MATCH - wall_clock_shifted = ( - raw_next_run_dt.replace(tzinfo=None) != next_run_dt.replace(tzinfo=None) - ) + wall_clock_shifted = raw_next_run_dt.replace(tzinfo=None) != next_run_dt.replace(tzinfo=None) if wall_clock_shifted and _cron_next_run_matches_expr(schedule, raw_next_run_dt): return STALE_CRON_TIMEZONE_MIGRATION return STALE_CRON_EXPR_EDIT @@ -1131,9 +1028,7 @@ _timezone_migration_catchups_recent: list = [] def _record_timezone_migration_catchup( - job: Dict[str, Any], - raw_next_run_dt: datetime, - next_run_dt: datetime, + job: Dict[str, Any], raw_next_run_dt: datetime, next_run_dt: datetime, ) -> None: """Persist a countable signal for one offset-migration catch-up fire.""" global _timezone_migration_catchups @@ -1147,8 +1042,7 @@ def _record_timezone_migration_catchup( } _timezone_migration_catchups += 1 _append_telemetry_record( - "timezone_migration_catchups.jsonl", entry, _timezone_migration_catchups_recent - ) + "timezone_migration_catchups.jsonl", entry, _timezone_migration_catchups_recent) def get_timezone_migration_catchup_stats() -> Dict[str, Any]: @@ -1162,24 +1056,19 @@ def get_timezone_migration_catchup_stats() -> Dict[str, Any]: def compute_next_run(schedule: Dict[str, Any], last_run_at: Optional[str] = None) -> Optional[str]: """Compute the next run time for a schedule as an ISO string, or None if no more runs.""" now = _hermes_now() - if not isinstance(schedule, dict): return None kind = schedule.get("kind") - if kind is None: - return None - if kind == "once": return _recoverable_oneshot_run_at(schedule, now, last_run_at=last_run_at) - - elif kind == "interval": + # Recurring kinds anchor on last_run_at so a restart doesn't re-anchor the schedule. + base_time = (_parse_aware(last_run_at) if last_run_at else None) or now + if kind == "interval": minutes = schedule.get("minutes") if minutes is None: return None - last = _parse_aware(last_run_at) if last_run_at else None - return ((last or now) + timedelta(minutes=minutes)).isoformat() - - elif kind == "cron": + return (base_time + timedelta(minutes=minutes)).isoformat() + if kind == "cron": expr = schedule.get("expr") if not expr: return None @@ -1189,15 +1078,9 @@ def compute_next_run(schedule: Dict[str, Any], last_run_at: Optional[str] = None "not installed. croniter is a core dependency as of v0.9.x; " "reinstall hermes-agent or run 'pip install croniter' in your " "runtime env.", - expr, - ) + expr) return None - # Anchor on last_run_at (like interval jobs) so a restart doesn't re-anchor the schedule. - base_time = (_parse_aware(last_run_at) if last_run_at else None) or now - cron = croniter(expr, base_time) - next_run = cron.get_next(datetime) - return next_run.isoformat() - + return croniter(expr, base_time).get_next(datetime).isoformat() return None @@ -1220,9 +1103,10 @@ def record_ticker_heartbeat(success: bool = False) -> None: _write_marker("ticker_last_success", str(time.time()), ".hb_") -def _epoch_file_age(path: Path) -> Optional[float]: +def _epoch_file_age(name: str) -> Optional[float]: + """Seconds since the epoch stamp stored in ``/``; None = missing/unreadable.""" try: - raw = path.read_text(encoding="utf-8").strip() + raw = (_current_cron_store().cron_dir / name).read_text(encoding="utf-8").strip() return max(0.0, time.time() - float(raw)) except Exception: return None @@ -1230,12 +1114,12 @@ def _epoch_file_age(path: Path) -> Optional[float]: def get_ticker_heartbeat_age() -> Optional[float]: """Seconds since the ticker loop last iterated; None = missing/unreadable ("cannot determine", not "dead").""" - return _epoch_file_age(_current_cron_store().cron_dir / "ticker_heartbeat") + return _epoch_file_age("ticker_heartbeat") def get_ticker_success_age() -> Optional[float]: """Seconds since the ticker last completed a tick WITHOUT raising, or None.""" - return _epoch_file_age(_current_cron_store().cron_dir / "ticker_last_success") + return _epoch_file_age("ticker_last_success") def get_catch_up_occurrence_count() -> int: @@ -1259,25 +1143,20 @@ def record_ticker_error(message: str) -> None: def clear_ticker_error() -> None: """Remove the last-tick-error marker after a successful tick. Best-effort.""" - store = _current_cron_store() - try: - (store.cron_dir / "ticker_last_error").unlink() - except OSError: - pass + with contextlib.suppress(OSError): + (_current_cron_store().cron_dir / "ticker_last_error").unlink() def get_ticker_last_error() -> Optional[str]: """Return the most recent recorded tick error message, or None.""" - store = _current_cron_store() try: - raw = (store.cron_dir / "ticker_last_error").read_text(encoding="utf-8") + raw = (_current_cron_store().cron_dir / "ticker_last_error").read_text(encoding="utf-8") except Exception: return None lines = raw.splitlines() if len(lines) < 2: return None - message = "\n".join(lines[1:]).strip() - return message or None + return "\n".join(lines[1:]).strip() or None # --- Job CRUD Operations --- @@ -1323,10 +1202,8 @@ def load_jobs() -> List[Dict[str, Any]]: if skipped: logger.warning( "Skipping %d non-dict entr%s in id-keyed jobs map: %s", - len(skipped), - "y" if len(skipped) == 1 else "ies", - ", ".join(map(repr, skipped)), - ) + len(skipped), "y" if len(skipped) == 1 else "ies", + ", ".join(map(repr, skipped))) jobs = [{**v, "id": v.get("id") or k} for k, v in jobs.items() if isinstance(v, dict)] repair = "id-keyed jobs map flattened to list" elif isinstance(data, list): @@ -1334,8 +1211,7 @@ def load_jobs() -> List[Dict[str, Any]]: repair = "bare list wrapped as dict" else: raise RuntimeError( - f"Cron database corrupted: expected {{'jobs': [...]}}, got {type(data).__name__}" - ) + f"Cron database corrupted: expected {{'jobs': [...]}}, got {type(data).__name__}") if jobs and repair: save_jobs(jobs) logger.warning("Auto-repaired jobs.json (%s)", repair) @@ -1356,9 +1232,7 @@ def _peek_jobs_unlocked() -> Optional[List[Dict[str, Any]]]: if isinstance(data, dict): jobs = data.get("jobs", []) return jobs if isinstance(jobs, list) else None - if isinstance(data, list): - return data - return None + return data if isinstance(data, list) else None def _jobs_file_stamp(jobs_file: Path) -> Optional[Tuple[int, int, int]]: @@ -1375,9 +1249,8 @@ def _record_load_stamp(stamp: Optional[Tuple[int, int, int]]) -> None: """Remember jobs.json's stamp for the enclosing _jobs_lock() section (no-op outside one) so the save path can skip the shrink-merge when disk provably hasn't changed. Capture it BEFORE reading: a mid-read sibling then mismatches (fail-safe); stamping after would certify an unseen write.""" - if not getattr(_jobs_lock_state, "depth", 0): - return - _jobs_lock_state.load_stamp = stamp + if getattr(_jobs_lock_state, "depth", 0): + _jobs_lock_state.load_stamp = stamp def _unmerged_disk_jobs( @@ -1405,9 +1278,7 @@ def _unmerged_disk_jobs( def _merge_unexpected_disk_jobs( - jobs: List[Dict[str, Any]], - *, - removed_ids: Optional[Collection[str]] = None, + jobs: List[Dict[str, Any]], *, removed_ids: Optional[Collection[str]] = None, ) -> List[Dict[str, Any]]: """*jobs* plus on-disk jobs absent from the payload (under the degraded flock-timeout path a stale writer would otherwise clobber concurrent creates). Deletes pass ``removed_ids``; never mutates *jobs*.""" @@ -1418,19 +1289,14 @@ def _merge_unexpected_disk_jobs( "Preserved %d cron job(s) present on disk but missing from the " "in-memory save payload (concurrent create under degraded lock " "or stale writer) (#80624): %s", - len(recovered), - [j.get("id") for j in recovered], - ) + len(recovered), [j.get("id") for j in recovered]) return jobs + recovered def _unlink_quiet(path: Optional[str]) -> None: - if path is None: - return - try: - os.unlink(path) - except OSError: - pass + if path is not None: + with contextlib.suppress(OSError): + os.unlink(path) def _stage_jobs_payload(jobs_file: Path, jobs: List[Dict[str, Any]]) -> str: @@ -1440,10 +1306,7 @@ def _stage_jobs_payload(jobs_file: Path, jobs: List[Dict[str, Any]]) -> str: with os.fdopen(fd, "w", encoding="utf-8") as f: json.dump( {"jobs": jobs, "updated_at": _hermes_now().isoformat()}, - f, - indent=2, - ensure_ascii=False, - ) + f, indent=2, ensure_ascii=False) f.flush() os.fsync(f.fileno()) except BaseException: @@ -1501,6 +1364,8 @@ def _save_jobs_unlocked( except BaseException: _unlink_quiet(tmp_path) raise + + def save_jobs( jobs: List[Dict[str, Any]], *, @@ -1553,8 +1418,7 @@ def _normalize_workdir(workdir: Optional[str]) -> Optional[str]: if not expanded.is_absolute(): raise ValueError( f"Cron workdir must be an absolute path (got {raw!r}). " - f"Cron jobs run detached from any shell cwd, so relative paths are ambiguous." - ) + f"Cron jobs run detached from any shell cwd, so relative paths are ambiguous.") resolved = expanded.resolve() if not resolved.exists(): raise ValueError(f"Cron workdir does not exist: {resolved}") @@ -1573,11 +1437,9 @@ def _resolve_default_model_snapshot() -> Optional[str]: if not cfg_path.exists(): return None cfg = read_user_config_raw(cfg_path) - try: + with contextlib.suppress(Exception): from hermes_cli import managed_scope cfg = managed_scope.apply_managed_overlay(cfg) - except Exception: - pass cfg = _expand_env_vars(cfg) cron_cfg = cfg.get("cron") or {} if isinstance(cron_cfg, dict): @@ -1624,12 +1486,33 @@ def _normalize_context_from(value: Any) -> Optional[List[str]]: def _normalize_failure_deliver(value: Any) -> Optional[str]: """failure_deliver shares deliver's value grammar; flatten str/list like the tool layer's _normalize_deliver_param for direct create_job callers. Semantic validation happens at - resolution time via the shared deliver path (NS-788).""" + resolution time via the shared deliver path.""" if isinstance(value, (list, tuple)): return ",".join(str(p).strip() for p in value if str(p).strip()) or None return _normalize_job_optional_text(value) +def _normalize_reasoning_effort(value: Any) -> Optional[str]: + """Spelling-only validation via the shared parser (cron knob never stricter/looser than + config.yaml); model capability is deliberately NOT checked (model unknowable at create time, + transports clamp at send time). None for unset, lowercase level, or ValueError.""" + if value is None: + return None + text = str(value).strip().lower() + if not text: + return None + from hermes_constants import parse_reasoning_effort + + if parse_reasoning_effort(text) is None: + raise ValueError( + f"Invalid reasoning_effort {value!r}. Valid levels: " + "none, minimal, low, medium, high, xhigh, max, ultra " + "(empty string clears the override).") + if text in {"false", "disabled"}: + return "none" + return text + + # Normalizers for create_job (all fields) / update_job (present fields). Invalid values raise BEFORE storing. _CREATE_FIELD_NORMALIZERS: Dict[str, Callable[[Any], Any]] = { "model": _normalize_job_optional_text, @@ -1648,37 +1531,12 @@ _UPDATE_FIELD_NORMALIZERS: Dict[str, Callable[[Any], Any]] = { "workdir": lambda v: None if v in {None, "", False} else _normalize_workdir(v), "monitor_script": _normalize_job_optional_text, "monitor_url": _normalize_job_optional_text, + "reasoning_effort": _normalize_reasoning_effort, } -def _normalize_reasoning_effort(value: Any) -> Optional[str]: - """Spelling-only validation via the shared parser (cron knob never stricter/looser than - config.yaml); model capability is deliberately NOT checked (model unknowable at create time, - transports clamp at send time). None for unset, lowercase level, or ValueError.""" - if value is None: - return None - text = str(value).strip().lower() - if not text: - return None - from hermes_constants import parse_reasoning_effort - - if parse_reasoning_effort(text) is None: - raise ValueError( - f"Invalid reasoning_effort {value!r}. Valid levels: " - "none, minimal, low, medium, high, xhigh, max, ultra " - "(empty string clears the override)." - ) - if text in {"false", "disabled"}: - return "none" - return text - - def _compute_provider_model_snapshots( - *, - provider: Any, - model: Any, - base_url: Any, - no_agent: Any, + *, provider: Any, model: Any, base_url: Any, no_agent: Any, ) -> Tuple[Optional[str], Optional[str]]: """Snapshot unpinned provider/model resolution so a later global switch fails closed at fire time instead of silently changing spend. Pinned axes and no-agent jobs carry no snapshot.""" @@ -1698,8 +1556,7 @@ def _compute_provider_model_snapshots( if normalized_base_url: runtime_kwargs["explicit_base_url"] = normalized_base_url snap = resolve_runtime_provider(**runtime_kwargs) - snap_provider = str(snap.get("provider") or "").strip().lower() - provider_snapshot = snap_provider or None + provider_snapshot = str(snap.get("provider") or "").strip().lower() or None except Exception: provider_snapshot = None if normalized_model is None: @@ -1721,23 +1578,18 @@ def _normalized_inference_axes(job: Dict[str, Any]) -> Tuple[Optional[str], Opti def _validate_job_mode_invariants( - monitor_script: Optional[str], - monitor_url: Optional[str], - no_agent: bool, - script: Optional[str], + monitor_script: Optional[str], monitor_url: Optional[str], no_agent: bool, script: Optional[str], ) -> None: """Execution-mode invariants shared by create_job and update_job (no bypass via the update door).""" if monitor_script and monitor_url: raise ValueError( "monitor_script and monitor_url are mutually exclusive — a job " - "can only have one monitor source." - ) + "can only have one monitor source.") if (monitor_script or monitor_url) and no_agent: raise ValueError( "monitor_script/monitor_url cannot be combined with no_agent=True — " "the whole point of a monitor job is to suppress or wake the AGENT " - "based on source changes. Use a plain no_agent script job instead." - ) + "based on source changes. Use a plain no_agent script job instead.") if no_agent and not script: raise ValueError(NO_AGENT_WITHOUT_SCRIPT_ERROR) @@ -1745,8 +1597,22 @@ def _validate_job_mode_invariants( def _oneshot_past_grace_error(run_at: Any) -> ValueError: return ValueError( f"Requested one-shot time {run_at} is more than " - f"{ONESHOT_GRACE_SECONDS}s in the past and cannot be scheduled." - ) + f"{ONESHOT_GRACE_SECONDS}s in the past and cannot be scheduled.") + + +def _next_run_or_reject_past_oneshot( + parsed_schedule: Dict[str, Any], label: str, fallback_run_at: Any, what: str, +) -> Optional[str]: + """``compute_next_run`` that raises (after a warning log) for a one-shot outside the grace window, + so a ghost job with ``next_run_at=None`` can never be stored.""" + next_run_at = compute_next_run(parsed_schedule) + if parsed_schedule.get("kind") == "once" and next_run_at is None: + run_at = parsed_schedule.get("run_at") or fallback_run_at + logger.warning( + "Rejecting one-shot cron job %s'%s': run_at %s is outside the %ss grace window", + what, label, run_at, ONESHOT_GRACE_SECONDS) + raise _oneshot_past_grace_error(run_at) + return next_run_at def create_job( @@ -1774,28 +1640,19 @@ def create_job( ) -> Dict[str, Any]: """Create a new cron job and return the stored record. - prompt: self-contained prompt (or task instruction when skills are set); with ``no_agent`` only - a name hint. deliver: defaults to "origin" when ``origin`` is given, else "local". repeat: None = - forever. skill/skills: legacy single / ordered list loaded before the prompt. script: stdout is - injected as prompt context, or with ``no_agent=True`` IS the job (stdout delivered verbatim; - empty = silent; requires ``script``); relative paths resolve under ~/.hermes/scripts/, - .sh/.bash run via bash else Python. context_from: job id(s) whose latest output is injected - (chaining). enabled_toolsets: restrict the agent (ignored with ``no_agent``). workdir: absolute - cwd for tools/scripts; AGENTS.md/CLAUDE.md/.cursorrules injected; unset = legacy behaviour. - monitor_script/monitor_url: cheap monitor source run FIRST each tick; output hashed as exact - bytes — unchanged suppresses the agent run, changed injects a MONITOR CHANGE DETECTED block - (diff + new output); mutually exclusive, incompatible with ``no_agent``. reasoning_effort: - per-job pin (canonical levels, case-insensitive) that wins over global/per-model config; - capability is NOT validated (provider transport clamps at send time); inert with ``no_agent``. + deliver defaults to "origin" when ``origin`` is given, else "local"; repeat None = forever. + script: stdout is injected as prompt context, or with ``no_agent=True`` IS the job (stdout + delivered verbatim, requires ``script``). context_from: job id(s) whose latest output is + injected. workdir: absolute cwd for tools/scripts. monitor_script/monitor_url: cheap monitor + source run FIRST each tick; unchanged output suppresses the agent run (mutually exclusive, + incompatible with ``no_agent``). reasoning_effort: per-job pin; capability NOT validated. """ parsed_schedule = parse_schedule(schedule) - repeat = normalize_repeat_value(repeat) if parsed_schedule["kind"] == "once" and repeat is None: repeat = 1 if deliver is None: deliver = "origin" if origin else "local" - job_id = uuid.uuid4().hex[:12] now = _hermes_now().isoformat() @@ -1806,11 +1663,9 @@ def create_job( normalized_reasoning_effort = _normalize_reasoning_effort(reasoning_effort) _validate_job_mode_invariants(f["monitor_script"], f["monitor_url"], f["no_agent"], f["script"]) - prompt_text = _coerce_job_text(prompt).strip() if not prompt_text and not f["script"] and not normalized_skills: raise ValueError(EMPTY_PAYLOAD_ERROR) - # Reject gateway-lifecycle commands (respawn loops) here, not just the CLI, to cover the cronjob tool. from cron.lifecycle_guard import check_gateway_lifecycle check_gateway_lifecycle(prompt_text, f["script"]) @@ -1822,19 +1677,9 @@ def create_job( or "cron job" ) name = name or label_source[:50].strip() - provider_snapshot, model_snapshot = _compute_provider_model_snapshots( - provider=f["provider"], model=f["model"], base_url=f["base_url"], no_agent=f["no_agent"] - ) - - next_run_at = compute_next_run(parsed_schedule) - if parsed_schedule.get("kind") == "once" and next_run_at is None: - run_at = parsed_schedule.get("run_at") or schedule - logger.warning( - "Rejecting one-shot cron job '%s': run_at %s is outside the %ss grace window", - name, run_at, ONESHOT_GRACE_SECONDS, - ) - raise _oneshot_past_grace_error(run_at) + provider=f["provider"], model=f["model"], base_url=f["base_url"], no_agent=f["no_agent"]) + next_run_at = _next_run_or_reject_past_oneshot(parsed_schedule, name, schedule, "") job = { "id": job_id, @@ -1874,14 +1719,15 @@ def create_job( "enabled_toolsets": f["enabled_toolsets"], "workdir": f["workdir"], } - # Persist only when explicitly set; absent key => fall back to global config resolution. - if normalized_attach is not None: - job["attach_to_session"] = normalized_attach - if normalized_reasoning_effort is not None: - job["reasoning_effort"] = normalized_reasoning_effort - # failure_deliver too: absent key = failures follow deliver, byte-identical to pre-feature jobs. - if f["failure_deliver"] is not None: - job["failure_deliver"] = f["failure_deliver"] + # Optional keys are persisted only when explicitly set: an absent key falls back to global config + # (attach/reasoning) or to ``deliver`` (failure_deliver), byte-identical to pre-feature jobs. + for key, value in ( + ("attach_to_session", normalized_attach), + ("reasoning_effort", normalized_reasoning_effort), + ("failure_deliver", f["failure_deliver"]), + ): + if value is not None: + job[key] = value with _jobs_lock(): save_jobs(load_jobs() + [job]) @@ -1903,8 +1749,7 @@ class AmbiguousJobReference(LookupError): ids = ", ".join(m["id"] for m in matches) super().__init__( f"Job name '{ref}' is ambiguous — matches {len(matches)} jobs: {ids}. " - f"Use the job ID instead." - ) + f"Use the job ID instead.") def resolve_job_ref(ref: str) -> Optional[Dict[str, Any]]: @@ -1921,9 +1766,7 @@ def resolve_job_ref(ref: str) -> Optional[Dict[str, Any]]: if not name_matches: return None if len(name_matches) > 1: - raise AmbiguousJobReference( - ref, [_normalize_job_record(j) for j in name_matches] - ) + raise AmbiguousJobReference(ref, [_normalize_job_record(j) for j in name_matches]) return _normalize_job_record(name_matches[0]) @@ -1956,8 +1799,7 @@ def _reject_terminal_activation(job: Dict[str, Any], updated: Dict[str, Any], jo ): raise ValueError( f"Cannot activate terminal cron job '{job.get('name', job_id)}' " - "through update_job; use cron resume --run-now or --at." - ) + "through update_job; use cron resume --run-now or --at.") def _normalize_job_updates(job: Dict[str, Any], updates: Dict[str, Any]) -> None: @@ -1966,8 +1808,6 @@ def _normalize_job_updates(job: Dict[str, Any], updates: Dict[str, Any]) -> None for key, norm in _UPDATE_FIELD_NORMALIZERS.items(): if key in updates: updates[key] = norm(updates[key]) - if "reasoning_effort" in updates: - updates["reasoning_effort"] = _normalize_reasoning_effort(updates["reasoning_effort"]) if "repeat" in updates: _rp = updates["repeat"] completed = (job.get("repeat") or {}).get("completed", 0) @@ -1980,30 +1820,51 @@ def _normalize_job_updates(job: Dict[str, Any], updates: Dict[str, Any]) -> None updates["repeat"] = {"times": normalize_repeat_value(_rp), "completed": completed} +def _apply_schedule_update(updated: Dict[str, Any], updates: Dict[str, Any], job_id: str) -> None: + """Parse a string schedule, refresh ``schedule_display`` and (unless paused) ``next_run_at``.""" + updated_schedule = updated["schedule"] + if isinstance(updated_schedule, str): + updated_schedule = parse_schedule(updated_schedule) + updated["schedule"] = updated_schedule + updated["schedule_display"] = updates.get( + "schedule_display", updated_schedule.get("display", updated.get("schedule_display"))) + if updated.get("state") != "paused": + updated["next_run_at"] = _next_run_or_reject_past_oneshot( + updated_schedule, updated.get("name", job_id), updated_schedule, "update ") + + +def _fill_missing_next_run(updated: Dict[str, Any]) -> None: + """An enabled, unpaused record must never persist without ``next_run_at`` (it would never fire).""" + if not updated.get("enabled", True) or updated.get("state") == "paused" or updated.get("next_run_at"): + return + next_run = compute_next_run(updated["schedule"]) + if next_run is None and updated["schedule"].get("kind") == "once": + run_at = updated["schedule"].get("run_at", "unknown") + raise ValueError( + f"Requested one-shot time {run_at} is in the past " + f"(grace window: {ONESHOT_GRACE_SECONDS}s) and cannot be scheduled.") + updated["next_run_at"] = next_run + + def update_job(job_id: str, updates: Dict[str, Any]) -> Optional[Dict[str, Any]]: """Update a job by ID, refreshing derived schedule fields when needed.""" # ``id`` is a path component under OUTPUT_DIR — changing it would leak path-escape values. bad_fields = _IMMUTABLE_JOB_FIELDS.intersection(updates or {}) if bad_fields: - raise ValueError( - f"Cron job field(s) cannot be updated: {', '.join(sorted(bad_fields))}" - ) + raise ValueError(f"Cron job field(s) cannot be updated: {', '.join(sorted(bad_fields))}") def apply(jobs, i, job): _normalize_job_updates(job, updates) previous_inference_axes = _normalized_inference_axes(job) updated = _apply_skill_fields({**job, **updates}) _reject_terminal_activation(job, updated, job_id) - # Re-check on the MERGED record; scoped to changed fields so legacy records keep loading. if {"monitor_script", "monitor_url", "no_agent", "script"}.intersection(updates): _validate_job_mode_invariants( updated.get("monitor_script") or None, updated.get("monitor_url") or None, bool(updated.get("no_agent")), - _normalize_job_optional_text(updated.get("script")), - ) - + _normalize_job_optional_text(updated.get("script"))) if any(k in updates for k in _PAYLOAD_FIELDS) and job_payload_is_empty(updated): raise ValueError(EMPTY_PAYLOAD_ERROR) inference_fields_changed = bool( @@ -2011,51 +1872,15 @@ def update_job(job_id: str, updates: Dict[str, Any]) -> Optional[Dict[str, Any]] ) and _normalized_inference_axes(updated) != previous_inference_axes if "schedule" in updates: - updated_schedule = updated["schedule"] - if isinstance(updated_schedule, str): - updated_schedule = parse_schedule(updated_schedule) - updated["schedule"] = updated_schedule - updated["schedule_display"] = updates.get( - "schedule_display", - updated_schedule.get("display", updated.get("schedule_display")), - ) - if updated.get("state") != "paused": - updated_next_run = compute_next_run(updated_schedule) - # Same guard as create_job: otherwise a ghost job (next_run_at=None) never fires. - if updated_next_run is None and updated_schedule.get("kind") == "once": - run_at = updated_schedule.get("run_at") or updated_schedule - logger.warning( - "Rejecting one-shot cron job update '%s': run_at %s " - "is outside the %ss grace window", - updated.get("name", job_id), - run_at, - ONESHOT_GRACE_SECONDS, - ) - raise _oneshot_past_grace_error(run_at) - updated["next_run_at"] = updated_next_run - + _apply_schedule_update(updated, updates, job_id) if inference_fields_changed: - provider_snapshot, model_snapshot = _compute_provider_model_snapshots( + updated["provider_snapshot"], updated["model_snapshot"] = _compute_provider_model_snapshots( provider=updated.get("provider"), model=updated.get("model"), base_url=updated.get("base_url"), - no_agent=updated.get("no_agent"), - ) - updated["provider_snapshot"] = provider_snapshot - updated["model_snapshot"] = model_snapshot - - if updated.get("enabled", True) and updated.get("state") != "paused" and not updated.get("next_run_at"): - next_run = compute_next_run(updated["schedule"]) - if next_run is None and updated["schedule"].get("kind") == "once": - run_at = updated["schedule"].get("run_at", "unknown") - raise ValueError( - f"Requested one-shot time {run_at} is in the past " - f"(grace window: {ONESHOT_GRACE_SECONDS}s) and cannot be scheduled." - ) - updated["next_run_at"] = next_run - + no_agent=updated.get("no_agent")) + _fill_missing_next_run(updated) _reject_terminal_activation(job, updated, job_id) - jobs[i] = updated save_jobs(jobs) return _normalize_job_record(updated) @@ -2068,15 +1893,12 @@ def pause_job(job_id: str, reason: Optional[str] = None) -> Optional[Dict[str, A job = resolve_job_ref(job_id) if not job: return None - return update_job( - job["id"], - { - "enabled": False, - "state": "paused", - "paused_at": _hermes_now().isoformat(), - "paused_reason": reason, - }, - ) + return update_job(job["id"], { + "enabled": False, + "state": "paused", + "paused_at": _hermes_now().isoformat(), + "paused_reason": reason, + }) def resume_job(job_id: str) -> Optional[Dict[str, Any]]: @@ -2084,29 +1906,22 @@ def resume_job(job_id: str) -> Optional[Dict[str, Any]]: job = resolve_job_ref(job_id) if not job: return None - next_run_at = compute_next_run(job["schedule"]) if next_run_at is None and job["schedule"].get("kind") == "once": run_at = job["schedule"].get("run_at", "unknown") raise ValueError( f"Cannot resume: one-shot time {run_at} is in the past " - f"(grace window: {ONESHOT_GRACE_SECONDS}s) and will never fire." - ) - return update_job( - job["id"], - { - "enabled": True, - "state": "scheduled", - "paused_at": None, - "paused_reason": None, - "next_run_at": next_run_at, - }, - ) + f"(grace window: {ONESHOT_GRACE_SECONDS}s) and will never fire.") + return update_job(job["id"], { + "enabled": True, + "state": "scheduled", + "paused_at": None, + "paused_reason": None, + "next_run_at": next_run_at, + }) -def trigger_job( - job_id: str, extra_prompt: Optional[str] = None -) -> Optional[Dict[str, Any]]: +def trigger_job(job_id: str, extra_prompt: Optional[str] = None) -> Optional[Dict[str, Any]]: """Schedule a job for the next tick (ID or name). ``extra_prompt`` is stamped as ``manual_run_prompt`` for that single fire only; ``mark_job_run`` clears it.""" job = resolve_job_ref(job_id) @@ -2118,31 +1933,34 @@ def trigger_job( raise ValueError( f"Cannot run: job '{name}' is {state} (terminal). " f"Create a new occurrence with 'hermes cron resume {name} " - "--run-now' or '--at '." - ) + "--run-now' or '--at '.") manual_run_at = _hermes_now().isoformat() - return update_job( - job["id"], - { - "enabled": True, - "state": "scheduled", - "paused_at": None, - "paused_reason": None, - "next_run_at": manual_run_at, - # Run-now intent, so cron expression/TZ repair guards don't treat it as stale state. - "manual_run_at": manual_run_at, - "manual_run_prompt": (extra_prompt or None), - }, - ) + return update_job(job["id"], { + "enabled": True, + "state": "scheduled", + "paused_at": None, + "paused_reason": None, + "next_run_at": manual_run_at, + # Run-now intent, so cron expression/TZ repair guards don't treat it as stale state. + "manual_run_at": manual_run_at, + "manual_run_prompt": (extra_prompt or None), + }) def _claim_is_live(claim: Any, now: datetime, ttl_seconds: float) -> bool: + """True for a well-formed claim aged within ``[0, ttl)``: future-dated (clock/TZ skew) or + malformed claims count as stale so they can never wedge a job.""" if not isinstance(claim, dict) or not claim.get("at"): return False claimed_at = _parse_aware(claim["at"]) return claimed_at is not None and 0 <= (now - claimed_at).total_seconds() < ttl_seconds +_REARM_RECURRING_ERROR = ( + "Cannot re-arm recurring jobs: re-arm is one-shot-only; use plain resume or cron run." +) + + def rearm_oneshot(job_id: str, run_at: Any) -> Optional[Dict[str, Any]]: """Re-arm a completed one-shot as an explicit new occurrence.""" job_ref = resolve_job_ref(job_id) @@ -2152,10 +1970,7 @@ def rearm_oneshot(job_id: str, run_at: Any) -> Optional[Dict[str, Any]]: run_at = run_at.isoformat() parsed_schedule = parse_schedule(str(run_at)) if parsed_schedule.get("kind") != "once": - raise ValueError( - "Cannot re-arm recurring jobs: re-arm is one-shot-only; " - "use plain resume or cron run." - ) + raise ValueError(_REARM_RECURRING_ERROR) next_run_at = compute_next_run(parsed_schedule) if next_run_at is None: raise _oneshot_past_grace_error(parsed_schedule.get("run_at") or run_at) @@ -2167,10 +1982,7 @@ def rearm_oneshot(job_id: str, run_at: Any) -> Optional[Dict[str, Any]]: if _claim_is_live(job.get("fire_claim"), now, 300): raise ValueError("Cannot re-arm one-shot over a live fire claim.") if job.get("schedule", {}).get("kind") != "once": - raise ValueError( - "Cannot re-arm recurring jobs: re-arm is one-shot-only; " - "use plain resume or cron run." - ) + raise ValueError(_REARM_RECURRING_ERROR) repeat = job.get("repeat") or {} repeat["completed"] = 0 job["schedule"] = parsed_schedule @@ -2196,48 +2008,23 @@ def remove_job(job_id: str) -> bool: jobs = load_jobs() original_len = len(jobs) jobs = [j for j in jobs if j["id"] != canonical_id] - if len(jobs) < original_len: - # Resolve BEFORE saving so a legacy unsafe ID fails closed without a half-applied removal. - job_output_dir = _job_output_dir(canonical_id) - save_jobs(jobs, removed_ids={canonical_id}) - if job_output_dir.exists(): - shutil.rmtree(job_output_dir) - try: - from cron.notepad import clear_notepad - clear_notepad(canonical_id) - except Exception: - logger.debug( - "Failed to clear notepad for removed job %s", - canonical_id, exc_info=True, - ) - # Prune the fire-fence lock entry so the registry doesn't grow monotonically. - _fence_key = f"{_current_cron_store().cron_dir.resolve()}::{canonical_id}" - with _fire_fence_locks_guard: - _fire_fence_locks.pop(_fence_key, None) - return True - return False - - -def mark_job_run( - job_id: str, - success: bool, - error: Optional[str] = None, - delivery_error: Optional[str] = None, - status: Optional[str] = None, - *, - expected_fire_owner: Optional[str] = None, -) -> bool: - with _fire_job_lock(job_id) as acquired: - if not acquired: + if len(jobs) == original_len: return False - return _mark_job_run_locked( - job_id, - success, - error, - delivery_error, - status=status, - expected_fire_owner=expected_fire_owner, - ) + # Resolve BEFORE saving so a legacy unsafe ID fails closed without a half-applied removal. + job_output_dir = _job_output_dir(canonical_id) + save_jobs(jobs, removed_ids={canonical_id}) + if job_output_dir.exists(): + shutil.rmtree(job_output_dir) + try: + from cron.notepad import clear_notepad + clear_notepad(canonical_id) + except Exception: + logger.debug("Failed to clear notepad for removed job %s", canonical_id, exc_info=True) + # Prune the fire-fence lock entry so the registry doesn't grow monotonically. + _fence_key = f"{_current_cron_store().cron_dir.resolve()}::{canonical_id}" + with _fire_fence_locks_guard: + _fire_fence_locks.pop(_fence_key, None) + return True def _set_alert_flag(job_id: str, field: str, value: bool) -> bool: @@ -2272,12 +2059,9 @@ def mark_drift_alerted(job_id: str) -> bool: def note_fire_forward_failure(job_id: str, detail: str) -> bool: - """Durably record that a scheduled fire could not be handed to the runner. - - Written by the dashboard fire webhook when the loopback forward to the gateway fails. Without - it the miss is invisible: no execution row exists and last_status/last_error only cover runs - that started. Stored as ``last_fire_error`` so list/CLI/dashboard surface it; mark_job_run clears it. - """ + """Durably record (as ``last_fire_error``) that a scheduled fire could not be handed to the + runner — written by the dashboard fire webhook when the loopback forward fails. Without it the + miss is invisible (no execution row, last_status only covers started runs); mark_job_run clears it.""" def apply(jobs, _i, job): job["last_fire_error"] = {"at": _hermes_now().isoformat(), "detail": str(detail or "")[:500]} save_jobs(jobs) @@ -2286,13 +2070,85 @@ def note_fire_forward_failure(job_id: str, detail: str) -> bool: return _with_job(job_id, apply, False) -def _mark_job_run_locked( +def _record_run_outcome( + job: Dict[str, Any], success: bool, error: Optional[str], delivery_error: Optional[str], + status: Optional[str], now: str, +) -> None: + """Stamp one completed run onto *job*: status fields, failure streak, alert markers, claims.""" + job["last_run_at"] = now + job.pop("manual_run_at", None) + # The transient manual-run context is single-fire: the run that just completed consumed it. + job.pop("manual_run_prompt", None) + delivery_failed = isinstance(delivery_error, str) and bool(delivery_error.strip()) + job["last_status"] = status or ( + "error" if not success else ("delivery_failed" if delivery_failed else "ok")) + job["last_error"] = error if not success else None + if success: + # Healthy run: drop the alert-once dedup markers so a FUTURE break re-alerts, and clear + # the forward-failure stamp so it only describes CURRENT auto-fire health. + job.pop("preflight_alerted", None) + job.pop("drift_alerted", None) + job.pop("last_fire_error", None) + job["failure_streak"] = 0 + else: + # Consecutive agent-failure streak; delivery failures do NOT count (scheduler._failure_streak_nudge). + job["failure_streak"] = int(job.get("failure_streak") or 0) + 1 + job["last_delivery_error"] = delivery_error + # Clear both claims: the run is over, so the job is claimable again. + job["fire_claim"] = None + if job.get("run_claim") is not None: # keep key absence for legacy records + job["run_claim"] = None + + +def _advance_after_run(job: Dict[str, Any], now: str) -> None: + """Bump ``repeat.completed`` and recompute ``next_run_at``; retire the record as a terminal + completion when the repeat limit is reached or a one-shot has no further run.""" + kind = job.get("schedule", {}).get("kind") + repeat = job.get("repeat") + if repeat: + times = repeat.get("times") + finite = times is not None and times > 0 + completed = repeat.get("completed", 0) + # Finite one-shots were pre-claimed by claim_dispatch() (completed already incremented) — + # do not double-count; recurring jobs and direct callers still get the increment. + if not (kind == "once" and finite and completed > 0): + completed += 1 + repeat["completed"] = completed + if finite and completed >= times: + # Limit reached: retain a terminal record instead of popping it, so the status just + # written stays inspectable in `cronjob list`; the retention sweep prunes it later. + _complete_job_record(job) + return + + job["next_run_at"] = compute_next_run(job["schedule"], now) + if job["next_run_at"] is not None: + if job.get("state") != "paused": + job["state"] = "scheduled" + elif kind in {"cron", "interval"}: + # Recurring: transient failure (e.g. croniter missing) — disabling it would turn a missing + # dep into "job completed" and silently drop the schedule. + job["state"] = "error" + if not job.get("last_error"): + job["last_error"] = ( + "Failed to compute next run for recurring " + "schedule (is the 'croniter' package " + "installed in the gateway's Python env?)") + logger.error( + "Job '%s' (%s) could not compute next_run_at; " + "leaving enabled and marking state=error so the " + "job is not silently disabled.", + job.get("name", job.get("id", "?")), kind) + else: + _complete_job_record(job) # one-shot: terminal completion + + +def mark_job_run( job_id: str, success: bool, error: Optional[str] = None, delivery_error: Optional[str] = None, - *, status: Optional[str] = None, + *, expected_fire_owner: Optional[str] = None, ) -> bool: """Mark a job as run: update last_run_at/last_status, bump completed, recompute next_run_at, @@ -2300,7 +2156,8 @@ def _mark_job_run_locked( ``delivery_error`` is separate from the agent error: agent succeeded but delivery failed records ``last_status = "delivery_failed"`` (never "ok") while ``failure_streak`` is left alone. An - explicit ``status`` (e.g. "blocked_config") overrides the derived value. + explicit ``status`` (e.g. "blocked_config") overrides the derived value. False when the fence + can't be taken, the job is missing, or ``expected_fire_owner`` no longer holds the fire claim. """ def apply(jobs, _i, job): if expected_fire_owner is not None: @@ -2309,85 +2166,22 @@ def _mark_job_run_locked( logger.warning( "mark_job_run: job_id %s fire claim owner changed; " "discarding stale completion", - job_id, - ) + job_id) return False now = _hermes_now().isoformat() - job["last_run_at"] = now - job.pop("manual_run_at", None) - # The transient manual-run context is single-fire: the run that just completed consumed it. - job.pop("manual_run_prompt", None) - delivery_failed = isinstance(delivery_error, str) and bool(delivery_error.strip()) - job["last_status"] = status or ( - "error" if not success else ("delivery_failed" if delivery_failed else "ok") - ) - job["last_error"] = error if not success else None - if success: - # Healthy run: drop the alert-once dedup markers so a FUTURE break re-alerts, and clear - # the forward-failure stamp so it only describes CURRENT auto-fire health. - job.pop("preflight_alerted", None) - job.pop("drift_alerted", None) - job.pop("last_fire_error", None) - job["failure_streak"] = 0 - else: - # Consecutive agent-failure streak; delivery failures do NOT count (scheduler._failure_streak_nudge). - job["failure_streak"] = int(job.get("failure_streak") or 0) + 1 - job["last_delivery_error"] = delivery_error - # Clear both claims: the run is over, so the job is claimable again. - job["fire_claim"] = None - if job.get("run_claim") is not None: # keep key absence for legacy records - job["run_claim"] = None - - kind = job.get("schedule", {}).get("kind") - repeat = job.get("repeat") - if repeat: - times = repeat.get("times") - finite = times is not None and times > 0 - completed = repeat.get("completed", 0) - # Finite one-shots were pre-claimed by claim_dispatch() (completed already incremented) — - # do not double-count; recurring jobs and direct callers still get the increment. - if not (kind == "once" and finite and completed > 0): - completed += 1 - repeat["completed"] = completed - if finite and completed >= times: - # Limit reached: retain a terminal record instead of popping it, so the status just - # written stays inspectable in `cronjob list`; the retention sweep prunes it later. - _complete_job_record(job) - save_jobs(jobs) - return True - - job["next_run_at"] = compute_next_run(job["schedule"], now) - if job["next_run_at"] is None: - # One-shot: terminal completion. Recurring: transient failure (e.g. croniter missing) — - # disabling it would turn a missing dep into "job completed" and silently drop the schedule. - if kind in {"cron", "interval"}: - job["state"] = "error" - if not job.get("last_error"): - job["last_error"] = ( - "Failed to compute next run for recurring " - "schedule (is the 'croniter' package " - "installed in the gateway's Python env?)" - ) - logger.error( - "Job '%s' (%s) could not compute next_run_at; " - "leaving enabled and marking state=error so the " - "job is not silently disabled.", - job.get("name", job.get("id", "?")), - kind, - ) - else: - _complete_job_record(job) - elif job.get("state") != "paused": - job["state"] = "scheduled" - + _record_run_outcome(job, success, error, delivery_error, status, now) + _advance_after_run(job, now) save_jobs(jobs) return True - found = _with_job(job_id, apply, missing=_MISSING) - if found is _MISSING: - logger.warning("mark_job_run: job_id %s not found, skipping save", job_id) - return False - return found + def locked(): + found = _with_job(job_id, apply, missing=_MISSING) + if found is _MISSING: + logger.warning("mark_job_run: job_id %s not found, skipping save", job_id) + return False + return found + + return _under_fire_fence(job_id, locked) def _write_oneshot_diagnostic(job: Dict[str, Any], text: str, what: str) -> bool: @@ -2401,11 +2195,8 @@ def _write_oneshot_diagnostic(job: Dict[str, Any], text: str, what: str) -> bool def _write_wedged_oneshot_diagnostic(job: Dict[str, Any]) -> None: - """Leave an operator-visible trace when a wedged one-shot is removed. - - A finite one-shot whose dispatch was claimed but which never reached mark_job_run was - interrupted mid-run; removing it silently would leave no output, error, or record. - """ + """Trace for a wedged one-shot removal: dispatch was claimed but mark_job_run never ran + (interrupted mid-run); removing it silently would leave no output, error, or record.""" if job.get("last_run_at") is not None: return # a prior run was recorded — normal completion race, not a wedge repeat = job.get("repeat") or {} @@ -2423,14 +2214,12 @@ def _write_wedged_oneshot_diagnostic(job: Dict[str, Any]) -> None: "process was most likely killed or restarted mid-execution. The " "job has been removed to stop it re-firing; recreate it to run " "again.\n", - "wedged-oneshot", - ) + "wedged-oneshot") if written: logger.warning( "Job '%s': removed without a completed run — diagnostic written to " "its output directory", - job.get("name", job.get("id", "?")), - ) + job.get("name", job.get("id", "?"))) def _write_missed_oneshot_diagnostic(job: Dict[str, Any], next_run: str) -> None: @@ -2449,8 +2238,9 @@ def _write_missed_oneshot_diagnostic(job: Dict[str, Any], next_run: str) -> None "enforced at create/update/resume time. The job was removed " "without running; recreate it (or use the Run button) to " "schedule it again.\n", - "missed-oneshot", - ) + "missed-oneshot") + + def claim_dispatch(job_id: str) -> bool: """Atomically claim a finite one-shot dispatch BEFORE execution. @@ -2493,28 +2283,11 @@ def claim_dispatch(job_id: str) -> bool: logger.debug( "claim_dispatch: job_id %s not in store — proceeding without claim " "(handed-in job dict; nothing to persist a claim against)", - job_id, - ) + job_id) return True return claimed -def heartbeat_run_claim(job_id: str, *, expected_owner: str) -> bool: - """Refresh a one-shot's ``run_claim`` timestamp while its run is alive. - - Called periodically by the run monitor so a long run keeps its claim fresh: an expired claim - then really means the claiming process died. Compare-and-refresh on ``expected_owner`` stops a - stale runner from extending a claim another process has since taken over. Returns False when - the job, claim, or ownership no longer matches. - """ - def apply(jobs, _i, job): - if job.get("schedule", {}).get("kind") != "once": - return False - return _refresh_claim(jobs, job.get("run_claim"), expected_owner) - - return _with_job(job_id, apply, False) - - def _refresh_claim(jobs: List[Dict[str, Any]], claim: Any, expected_owner: str) -> bool: """Compare-and-refresh a claim's ``at`` stamp; False unless *expected_owner* still holds it.""" if not isinstance(claim, dict) or claim.get("by") != expected_owner: @@ -2524,12 +2297,21 @@ def _refresh_claim(jobs: List[Dict[str, Any]], claim: Any, expected_owner: str) return True -def clear_run_claim(job_id: str) -> bool: - """Clear a one-shot's ``run_claim`` when dispatch itself fails. +def heartbeat_run_claim(job_id: str, *, expected_owner: str) -> bool: + """Refresh a one-shot's ``run_claim`` timestamp while its run is alive, so an expired claim + really means the claiming process died. Compare-and-refresh on ``expected_owner`` stops a stale + runner from extending a claim another process has since taken over.""" + def apply(jobs, _i, job): + if job.get("schedule", {}).get("kind") != "once": + return False + return _refresh_claim(jobs, job.get("run_claim"), expected_owner) - Such a job never reaches mark_job_run, so the stale claim would block re-dispatch until the TTL - expires; calling this on every early-exit path keeps "stays due, fires on the next healthy tick". - """ + return _with_job(job_id, apply, False) + + +def clear_run_claim(job_id: str) -> bool: + """Clear a one-shot's ``run_claim`` when dispatch itself fails: such a job never reaches + mark_job_run, so the stale claim would block re-dispatch until the TTL expires.""" def apply(jobs, _i, job): if job.get("schedule", {}).get("kind") != "once" or job.get("run_claim") is None: return False # recurring, or already cleared @@ -2573,8 +2355,7 @@ def advance_next_runs(job_ids) -> int: def advance_next_run(job_id: str) -> bool: """Advance a recurring job's next_run_at BEFORE run_job() so a mid-run crash cannot re-fire it on restart (at-most-once for recurring jobs — one missed run beats a crash-loop burst). - One-shots are left unchanged so they can retry. Returns True if next_run_at was advanced. - """ + One-shots are left unchanged so they can retry. Returns True if next_run_at was advanced.""" # >= 1 (not == 1): duplicate ids in a corrupted file all advance; still report the advance. return advance_next_runs([job_id]) >= 1 @@ -2599,31 +2380,13 @@ def claim_job_for_fire( claim_ttl_seconds: int = 300, force: bool = False, return_job: bool = False, -) -> Union[bool, Dict[str, Any]]: - with _fire_job_lock(job_id) as acquired: - if not acquired: - return False - return _claim_job_for_fire_locked( - job_id, - claim_ttl_seconds=claim_ttl_seconds, - force=force, - return_job=return_job, - ) - - -def _claim_job_for_fire_locked( - job_id: str, - *, - claim_ttl_seconds: int = 300, - force: bool = False, - return_job: bool = False, ) -> Union[bool, Dict[str, Any]]: """Atomically claim a job for one external 'fire' (multi-machine at-most-once); True iff THIS caller won. Used by ``CronScheduler.fire_due`` so exactly one of N replicas runs a job. - Under the file lock: reject missing/terminal/paused jobs unless ``force`` (an explicit manual - fire, which also enables/resumes the job atomically; external callbacks must leave it false so - a stale callback cannot resurrect a paused job). Lose if a fresh claim (younger than + Under the fence + file lock: reject missing/terminal/paused jobs unless ``force`` (an explicit + manual fire, which also enables/resumes the job atomically; external callbacks must leave it + false so a stale callback cannot resurrect a paused job). Lose if a fresh claim (younger than ``claim_ttl_seconds``) exists. Otherwise stamp ``fire_claim`` and, for recurring jobs, advance ``next_run_at`` so a stale re-delivery cannot re-fire; one-shots rely on the fresh claim alone. The TTL lets another fire reclaim a job whose claimant crashed; mark_job_run clears the claim. @@ -2636,8 +2399,6 @@ def _claim_job_for_fire_locked( if not force and not is_job_runnable(job): return False now = _hermes_now() - # _claim_is_live bounds the age on BOTH sides: a claim stamped in the future (clock/TZ skew, - # corrupted timestamp) must not stay "fresh" forever and wedge the job; malformed → overwrite. if _claim_is_live(job.get("fire_claim"), now, claim_ttl_seconds): return False # someone holds a fresh claim if force: @@ -2652,7 +2413,14 @@ def _claim_job_for_fire_locked( save_jobs(jobs) return copy.deepcopy(job) if return_job else True - return _with_job(job_id, apply, False) + return _under_fire_fence(job_id, lambda: _with_job(job_id, apply, False)) + + +def heartbeat_fire_claim(job_id: str, *, expected_owner: str) -> bool: + """Refresh an active ``fire_claim`` without extending another owner's lease: an execution may + outlive the TTL, and the owner check stops a stale runner from refreshing a recovered claim.""" + return _under_fire_fence(job_id, lambda: _with_job( + job_id, lambda jobs, _i, job: _refresh_claim(jobs, job.get("fire_claim"), expected_owner), False)) # Completed one-shots are retained in jobs.json (final status stays inspectable) and pruned by @@ -2677,10 +2445,7 @@ def _completed_oneshot_retention_days() -> float: def _sweep_completed_oneshots( - raw_jobs: List[Dict[str, Any]], - now: datetime, - *, - removed_ids: Optional[Set[str]] = None, + raw_jobs: List[Dict[str, Any]], now: datetime, *, removed_ids: Optional[Set[str]] = None, ) -> bool: """Prune completed one-shot records past retention (in place; True when anything was removed). @@ -2698,8 +2463,7 @@ def _sweep_completed_oneshots( if rj.get("state") != "completed": continue schedule = rj.get("schedule") - kind = schedule.get("kind") if isinstance(schedule, dict) else None - if kind != "once": + if (schedule.get("kind") if isinstance(schedule, dict) else None) != "once": continue last_run = rj.get("last_run_at") last_run_dt = _parse_aware(last_run) if isinstance(last_run, str) else None @@ -2713,36 +2477,14 @@ def _sweep_completed_oneshots( logger.info( "Job '%s': pruning completed one-shot record " "(finished %s, retention %.1f days)", - rj.get("name", rj.get("id", "?")), - last_run, - retention_days, - ) + rj.get("name", rj.get("id", "?")), last_run, retention_days) except Exception: logger.debug( - "Retention sweep skipped malformed job record %r", - rj.get("id", "?"), - exc_info=True, - ) + "Retention sweep skipped malformed job record %r", rj.get("id", "?"), exc_info=True) return removed -def heartbeat_fire_claim(job_id: str, *, expected_owner: str) -> bool: - with _fire_job_lock(job_id) as acquired: - if not acquired: - return False - return _heartbeat_fire_claim_locked( - job_id, - expected_owner=expected_owner, - ) - - -def _heartbeat_fire_claim_locked(job_id: str, *, expected_owner: str) -> bool: - """Refresh an active ``fire_claim`` without extending another owner's lease: an execution may - outlive the TTL, and the owner check stops a stale runner from refreshing a recovered claim.""" - return _with_job( - job_id, lambda jobs, _i, job: _refresh_claim(jobs, job.get("fire_claim"), expected_owner), False - ) - +# --- Due scan --- def get_due_jobs() -> List[Dict[str, Any]]: """Return all jobs due now. @@ -2785,13 +2527,10 @@ class _DueScan: def _normalize_due_scan_records(raw_jobs: List[Dict[str, Any]]) -> bool: - """Repair malformed store records in place BEFORE the due scan keys off them. - - A missing ``id`` (older writers used ``job_id``), non-dict ``schedule``, or non-ISO timestamp - used to raise mid-tick and abort the whole scan before save_jobs(), freezing the scheduler in a - fast-forward loop. Recover/synthesize the id, reset a bad schedule to ``{}``, and strip bad - timestamps so the "no next_run_at" path recomputes. Returns True when anything changed. - """ + """Repair malformed store records in place BEFORE the due scan keys off them: a missing ``id`` + (older writers used ``job_id``), non-dict ``schedule``, or non-ISO timestamp used to raise + mid-tick and abort the whole scan before save_jobs(), freezing the scheduler in a fast-forward + loop. Returns True when anything changed.""" changed = False for rj in raw_jobs: if not rj.get("id"): @@ -2803,7 +2542,7 @@ def _normalize_due_scan_records(raw_jobs: List[Dict[str, Any]]) -> bool: for key in ("next_run_at", "last_run_at"): value = rj.get(key) if value is not None and _parse_aware(value) is None: - rj.pop(key, None) + rj.pop(key, None) # the "no next_run_at" path recomputes changed = True return changed @@ -2811,14 +2550,12 @@ def _normalize_due_scan_records(raw_jobs: List[Dict[str, Any]]) -> bool: def _self_disable_half_paused(job: Dict[str, Any], scan: _DueScan) -> None: """Self-heal enabled=true with pause markers: the operator believes the job is frozen while the scheduler would still fire it. Force enabled=false so listings are honest; logged loudly since - pause_job now sets both fields atomically, so this should be rare.""" + pause_job sets both fields atomically, so this should be rare.""" jid = job.get("id") logger.error( "Job '%s' (%s) has pause markers while enabled=true; " "self-disabling so it cannot fire (pause must be authoritative).", - job.get("name", jid), - jid, - ) + job.get("name", jid), jid) rj = scan.find(jid) if rj is None: return @@ -2837,9 +2574,7 @@ def _recover_missing_next_run(job: Dict[str, Any], scan: _DueScan) -> Optional[s and would otherwise be silently skipped forever.""" schedule = job.get("schedule", {}) kind = schedule.get("kind") - recovered_next = _recoverable_oneshot_run_at( - schedule, scan.now, last_run_at=job.get("last_run_at") - ) + recovered_next = _recoverable_oneshot_run_at(schedule, scan.now, last_run_at=job.get("last_run_at")) recovery_kind = "one-shot" if recovered_next else None if not recovered_next and kind in {"cron", "interval"}: recovered_next = compute_next_run(schedule, scan.now.isoformat()) @@ -2850,20 +2585,14 @@ def _recover_missing_next_run(job: Dict[str, Any], scan: _DueScan) -> Optional[s job["next_run_at"] = recovered_next logger.info( "Job '%s' had no next_run_at; recovering %s run at %s", - job.get("name", job.get("id", "?")), - recovery_kind, - recovered_next, - ) + job.get("name", job.get("id", "?")), recovery_kind, recovered_next) scan.persist(job["id"], next_run_at=recovered_next) return recovered_next def _repair_timezone_shifted_cron( - job: Dict[str, Any], - schedule: Dict[str, Any], - raw_next_run_dt: datetime, - next_run_dt: datetime, - scan: _DueScan, + job: Dict[str, Any], schedule: Dict[str, Any], raw_next_run_dt: datetime, + next_run_dt: datetime, scan: _DueScan, ) -> bool: """Repair a cron job whose stored offset no longer matches now's (TZ migration). @@ -2884,22 +2613,14 @@ def _repair_timezone_shifted_cron( logger.info( "Job '%s' next_run_at offset changed (%s -> %s). " "Recomputing cron run to preserve local wall-clock intent: %s", - job.get("name", job.get("id", "?")), - raw_next_run_dt.utcoffset(), - scan.now.utcoffset(), - new_next, - ) + job.get("name", job.get("id", "?")), raw_next_run_dt.utcoffset(), scan.now.utcoffset(), new_next) scan.persist(job["id"], next_run_at=new_next) return True def _rearm_stale_error_recurring( - job: Dict[str, Any], - schedule: Dict[str, Any], - kind: Optional[str], - next_run: str, - next_run_dt: datetime, - scan: _DueScan, + job: Dict[str, Any], schedule: Dict[str, Any], kind: Optional[str], next_run: str, + next_run_dt: datetime, scan: _DueScan, ) -> datetime: """Re-arm a recurring job wedged in persisted last_status=error; returns the effective next_run_dt. @@ -2928,10 +2649,7 @@ def _rearm_stale_error_recurring( "job wedged in stale last_status=error without re-firing for " "a full cadence; re-arming next_run_at to %s so it " "re-dispatches without force-run/resume", - job.get("name", jid), - jid, - recovered_next, - ) + job.get("name", jid), jid, recovered_next) _record_persisted_error_recovery(job, next_run) job["next_run_at"] = recovered_next scan.persist(jid, next_run_at=recovered_next) @@ -2939,12 +2657,8 @@ def _rearm_stale_error_recurring( def _reanchor_stale_cron( - job: Dict[str, Any], - schedule: Dict[str, Any], - next_run: str, - raw_next_run_dt: datetime, - next_run_dt: datetime, - scan: _DueScan, + job: Dict[str, Any], schedule: Dict[str, Any], next_run: str, raw_next_run_dt: datetime, + next_run_dt: datetime, scan: _DueScan, ) -> bool: """Stale-schedule guard for a due cron instant; True when re-anchored without firing. @@ -2961,11 +2675,7 @@ def _reanchor_stale_cron( "Job '%s' next_run_at %s does not match its current " "cron expression %r (direct jobs.json edit?); " "re-anchoring to %s without firing.", - job.get("name", job.get("id", "?")), - next_run, - schedule.get("expr"), - new_next, - ) + job.get("name", job.get("id", "?")), next_run, schedule.get("expr"), new_next) if new_next: scan.persist(job["id"], next_run_at=new_next) return True @@ -2976,25 +2686,15 @@ def _reanchor_stale_cron( "pre-migration UTC offset (%s, now %s) and is a legal " "occurrence at its own wall clock; firing the due run " "instead of re-anchoring past it.", - job.get("name", job.get("id", "?")), - job.get("id"), - schedule.get("expr"), - next_run, - next_run_dt.isoformat(), - raw_next_run_dt.utcoffset(), - scan.now.utcoffset(), - ) + job.get("name", job.get("id", "?")), job.get("id"), schedule.get("expr"), next_run, + next_run_dt.isoformat(), raw_next_run_dt.utcoffset(), scan.now.utcoffset()) _record_timezone_migration_catchup(job, raw_next_run_dt, next_run_dt) return False def _fast_forward_missed_recurring( - job: Dict[str, Any], - schedule: Dict[str, Any], - next_run: str, - next_run_dt: datetime, - grace: int, - scan: _DueScan, + job: Dict[str, Any], schedule: Dict[str, Any], next_run: str, next_run_dt: datetime, + grace: int, scan: _DueScan, ) -> None: """Recurring job past its grace window: skip the accumulated misses, fire once now. @@ -3011,11 +2711,7 @@ def _fast_forward_missed_recurring( "Job '%s' missed its scheduled time (%s, grace=%ds). " "Running now; next run provisionally set to: %s " "(re-anchored on completion)", - job.get("name", job.get("id", "?")), - next_run, - grace, - new_next, - ) + job.get("name", job.get("id", "?")), next_run, grace, new_next) scan.persist(job["id"], next_run_at=new_next) record_catch_up_occurrence() @@ -3055,8 +2751,7 @@ def _oneshot_dispatch_limit_reached(job: Dict[str, Any], scan: _DueScan) -> bool logger.info( "Job '%s': dispatch limit reached (%d/%d) but its run is still in flight in this " "process — keeping entry", - name, completed, times, - ) + name, completed, times) return True if job.get("last_run_at") is not None: # A record with last_run_at completed a real run and was re-armed without a budget reset @@ -3066,13 +2761,11 @@ def _oneshot_dispatch_limit_reached(job: Dict[str, Any], scan: _DueScan) -> bool "a run (last_run_at=%s) — removing it WITHOUT firing. This record was re-armed " "without a budget reset (pre-#93615 store or hand edit); re-run it with " "'hermes cron resume --run-now' (#93524).", - name, completed, times, job.get("last_run_at"), - ) + name, completed, times, job.get("last_run_at")) else: logger.info( "Job '%s': one-shot dispatch limit reached (%d/%d) — removing stale due entry", - name, completed, times, - ) + name, completed, times) scan.retire(job["id"]) # The claimed run never completed here by definition — leave an operator-visible diagnostic. _write_wedged_oneshot_diagnostic(job) @@ -3097,6 +2790,7 @@ def _evaluate_due_job(job: Dict[str, Any], scan: _DueScan, run_claim_ttl: float) raw_next_run_dt = datetime.fromisoformat(next_run) schedule = job.get("schedule", {}) kind = schedule.get("kind") + recurring = kind in {"cron", "interval"} next_run_dt = _ensure_aware(raw_next_run_dt) # Intentionally string-exact on raw stored values: trigger_job stamps the SAME isoformat string # into both fields, and any rewrite of next_run_at (edit, re-anchor, fire-claim advance) must @@ -3119,7 +2813,7 @@ def _evaluate_due_job(job: Dict[str, Any], scan: _DueScan, run_claim_ttl: float) ): return False grace = _compute_grace_seconds(schedule) - if not manual_run and kind in {"cron", "interval"}: + if not manual_run and recurring: _fast_forward_missed_recurring(job, schedule, next_run, next_run_dt, grace, scan) if kind == "once": if _retire_expired_oneshot(job, next_run, next_run_dt, scan): @@ -3136,7 +2830,7 @@ def _evaluate_due_job(job: Dict[str, Any], scan: _DueScan, run_claim_ttl: float) # Missed-run visibility: persist scheduled-vs-actual timing so separate CLI processes can show # a late catch-up. Recurring only — expired one-shots were retired above; manual triggers aren't late. - if not manual_run and kind in {"cron", "interval"}: + if not manual_run and recurring: lateness = max(0.0, (now - next_run_dt).total_seconds()) dispatch_stamp = { "scheduled_at": next_run, @@ -3155,7 +2849,6 @@ def _get_due_jobs_locked() -> List[Dict[str, Any]]: scan = _DueScan(raw_jobs, _hermes_now()) scan.needs_save = _normalize_due_scan_records(raw_jobs) jobs = [_apply_skill_fields(j) for j in copy.deepcopy(raw_jobs)] - # One-shot run-claim TTL, resolved once per scan (see _oneshot_run_claim_ttl_seconds). run_claim_ttl = _oneshot_run_claim_ttl_seconds() @@ -3181,16 +2874,16 @@ def _get_due_jobs_locked() -> List[Dict[str, Any]]: except Exception: logger.exception( "Skipping malformed cron job %r during due scan", - job.get("name") or job.get("id") or "?", - ) - continue + job.get("name") or job.get("id") or "?") if scan.needs_save: save_jobs(raw_jobs, removed_ids=scan.removed or None) - return due -# Per-run output files (`cron/output//.md`) had no retention cap, so a frequent -# job could fill the disk over time. + + +# --- Run output --- + +# Per-run output files (`cron/output//.md`) are capped so a frequent job can't fill the disk. _CRON_OUTPUT_DEFAULT_KEEP = 50 @@ -3207,10 +2900,7 @@ def _prune_job_output(job_output_dir: Path, keep: int) -> int: return 0 try: files = sorted( - (f for f in job_output_dir.glob("*.md") if f.is_file()), - key=lambda f: f.name, - reverse=True, - ) + (f for f in job_output_dir.glob("*.md") if f.is_file()), key=lambda f: f.name, reverse=True) except OSError: return 0 deleted = 0 @@ -3229,15 +2919,10 @@ def save_job_output(job_id: str, output: str): job_output_dir = _job_output_dir(job_id) _ensure_cron_dir(job_output_dir) _secure_dir(job_output_dir) - - timestamp = _hermes_now().strftime("%Y-%m-%d_%H-%M-%S") - output_file = job_output_dir / f"{timestamp}.md" + output_file = job_output_dir / f"{_hermes_now().strftime('%Y-%m-%d_%H-%M-%S')}.md" atomic_write_text(output_file, output, tmp_prefix=".output_") _secure_file(output_file) - - # Bound per-job output growth so long-running deploys don't fill the disk. _prune_job_output(job_output_dir, _cron_output_keep()) - return output_file @@ -3255,10 +2940,7 @@ def _canonical_skill_ref(raw: Any) -> str: from agent.skill_utils import normalize_skill_lookup_name value = normalize_skill_lookup_name(value) or value except Exception: - logger.debug( - "referenced_skill_names: could not normalize skill ref %r", raw, - exc_info=True, - ) + logger.debug("referenced_skill_names: could not normalize skill ref %r", raw, exc_info=True) return value.strip().lstrip("/") @@ -3266,14 +2948,12 @@ def referenced_skill_names() -> Set[str]: """Skill names referenced by ANY cron job, deliberately including paused/disabled ones (a paused job never bumps its skills, yet resuming it must still find them). The curator uses this to protect referenced skills from inactivity archival. Names are canonicalized as the scheduler - does, so absolute paths are protected too. A corrupt store yields an empty set, never raises. - """ + does, so absolute paths are protected too. A corrupt store yields an empty set, never raises.""" try: jobs = load_jobs() except Exception: logger.debug("referenced_skill_names: failed to load cron jobs", exc_info=True) return set() - return { cleaned for job in jobs @@ -3284,8 +2964,7 @@ def referenced_skill_names() -> Set[str]: def rewrite_skill_refs( - consolidated: Optional[Dict[str, str]] = None, - pruned: Optional[List[str]] = None, + consolidated: Optional[Dict[str, str]] = None, pruned: Optional[List[str]] = None, ) -> Dict[str, Any]: """Rewrite cron job skill references after a curator consolidation pass. @@ -3293,14 +2972,11 @@ def rewrite_skill_refs( and skips missing skills). Consolidated names are replaced by their umbrella target without duplication, pruned names are dropped, ordering is preserved, and the legacy ``skill`` field is realigned. Returns ``{"rewrites": [{job_id, job_name, before, after, mapped, dropped}, ...], - "jobs_updated": N, "jobs_scanned": M}``. Load/save exceptions propagate so tests can assert - behaviour; the curator call site wraps this in try/except. + "jobs_updated": N, "jobs_scanned": M}``. Load/save exceptions propagate (the curator wraps). """ consolidated = dict(consolidated or {}) - pruned_set = set(pruned or []) # A skill listed in both wins as "consolidated" — it has a target, the more useful outcome. - pruned_set -= set(consolidated.keys()) - + pruned_set = set(pruned or []) - set(consolidated.keys()) if not consolidated and not pruned_set: return {"rewrites": [], "jobs_updated": 0, "jobs_scanned": 0} @@ -3311,11 +2987,9 @@ def rewrite_skill_refs( skills_before = _normalize_skill_list(job.get("skill"), job.get("skills")) if not skills_before: continue - mapped: Dict[str, str] = {} dropped: List[str] = [] new_skills: List[str] = [] - for name in skills_before: if name in consolidated: target = consolidated[name] @@ -3326,7 +3000,6 @@ def rewrite_skill_refs( dropped.append(name) elif name not in new_skills: new_skills.append(name) - if not mapped and not dropped: continue job["skills"] = new_skills @@ -3339,7 +3012,6 @@ def rewrite_skill_refs( "mapped": mapped, "dropped": dropped, }) - if rewrites: save_jobs(jobs) logger.info("Curator rewrote skill references in %d cron job(s)", len(rewrites))