refactor(hermes_cli): cron.py — fold subcommand wrappers into dispatch table, dedupe unverified-target rendering, compact narrative comments
This commit is contained in:
@@ -11,12 +11,10 @@ sys.path.insert(0, str(PROJECT_ROOT))
|
||||
|
||||
from hermes_cli.colors import Colors, color
|
||||
|
||||
# Gateway-lifecycle command detection lives in ``cron.lifecycle_guard`` so it
|
||||
# can be shared across every job-creation path (CLI + the agent's ``cronjob``
|
||||
# model tool via ``cron.jobs.create_job``) without a circular import. Re-export
|
||||
# ``_contains_gateway_lifecycle_command`` here for back-compat: ``tools/
|
||||
# terminal_tool.py`` imports it from this module to hard-block the same
|
||||
# commands at execution time when ``_HERMES_GATEWAY=1``.
|
||||
# Gateway-lifecycle command detection lives in ``cron.lifecycle_guard`` (shared by every
|
||||
# job-creation path without a circular import). Re-exported here because
|
||||
# ``tools/terminal_tool.py`` imports it from this module to hard-block the same commands
|
||||
# at execution time when ``_HERMES_GATEWAY=1``.
|
||||
from cron.lifecycle_guard import ( # noqa: F401 (re-exported for terminal_tool)
|
||||
contains_gateway_lifecycle_command as _contains_gateway_lifecycle_command,
|
||||
)
|
||||
@@ -45,11 +43,10 @@ def _cron_api(**kwargs):
|
||||
|
||||
|
||||
def _active_cron_provider_name() -> str:
|
||||
"""Name of the resolved cron scheduler provider ('builtin', 'chronos', …).
|
||||
"""Resolved cron scheduler provider name ('builtin', 'chronos', …); 'builtin' on any failure.
|
||||
|
||||
Best-effort + offline (``resolve_cron_scheduler`` reads config and the provider's
|
||||
``is_available()`` contract forbids network). Returns 'builtin' on any failure so callers fall
|
||||
back to the historical ticker-based checks.
|
||||
Best-effort + offline (``resolve_cron_scheduler`` reads config; ``is_available()`` forbids
|
||||
network), so callers fall back to the historical ticker-based checks.
|
||||
"""
|
||||
try:
|
||||
from cron.scheduler_provider import resolve_cron_scheduler
|
||||
@@ -60,39 +57,33 @@ def _active_cron_provider_name() -> str:
|
||||
|
||||
|
||||
def _builtin_gateway_liveness() -> Optional[bool]:
|
||||
"""Tri-state liveness of the builtin cron scheduler's trigger.
|
||||
"""Tri-state liveness of the builtin cron scheduler's trigger (None = unknown).
|
||||
|
||||
Single source of truth shared by the CLI (``_warn_if_gateway_not_running``) and the ``cronjob``
|
||||
model tool (#87033): the builtin ticker only runs inside the gateway process, so a scheduled job
|
||||
with no live gateway can never fire. Non-builtin providers (e.g.
|
||||
Shared by the CLI and the ``cronjob`` model tool: the builtin ticker only runs inside the
|
||||
gateway process, so a scheduled job with no live gateway can never fire. Non-builtin
|
||||
providers fire jobs without the gateway.
|
||||
"""
|
||||
try:
|
||||
if _active_cron_provider_name() != "builtin":
|
||||
return True # external provider fires jobs without the gateway
|
||||
# The gateway runtime lock is held for exactly the gateway's lifetime, so it
|
||||
# is a more reliable "is the ticker's process alive" signal than PID scanning
|
||||
# — and inside the gateway process it short-circuits to True, so the in-gateway
|
||||
# cron tool never emits a false "gateway not running" (find_gateway_pids can
|
||||
# transiently miss the gateway just after a restart).
|
||||
# The gateway runtime lock is held for exactly the gateway's lifetime — a more
|
||||
# reliable "ticker's process is alive" signal than PID scanning, and inside the
|
||||
# gateway process it short-circuits to True so the in-gateway cron tool never
|
||||
# emits a false "gateway not running" (find_gateway_pids can transiently miss the
|
||||
# gateway just after a restart).
|
||||
try:
|
||||
from gateway.status import is_gateway_runtime_lock_active
|
||||
|
||||
if is_gateway_runtime_lock_active():
|
||||
return True
|
||||
except Exception:
|
||||
# A crashing lock probe is "unknown", not "dead" — let the pid
|
||||
# scan below still decide instead of collapsing the whole
|
||||
# tri-state to None.
|
||||
pass
|
||||
from hermes_cli.gateway import (
|
||||
find_gateway_pids,
|
||||
named_profile_served_by_running_multiplexer,
|
||||
)
|
||||
pass # a crashing lock probe is "unknown", not "dead" — let the pid scan decide
|
||||
from hermes_cli.gateway import find_gateway_pids, named_profile_served_by_running_multiplexer
|
||||
|
||||
if find_gateway_pids():
|
||||
return True
|
||||
# Satellite profile: no local gateway.pid, but the default multiplexer
|
||||
# ticks this profile's cron store (#97120).
|
||||
# Satellite profile: no local gateway.pid, but the default multiplexer ticks this
|
||||
# profile's cron store.
|
||||
return named_profile_served_by_running_multiplexer()
|
||||
except Exception:
|
||||
return None
|
||||
@@ -101,12 +92,11 @@ def _builtin_gateway_liveness() -> Optional[bool]:
|
||||
def _warn_if_gateway_not_running() -> None:
|
||||
"""Warn that scheduled jobs won't fire unless the gateway is running.
|
||||
|
||||
The cron ticker only runs inside the gateway (``_start_cron_ticker`` in gateway/run.py); there
|
||||
is no standalone cron daemon. Without a running gateway, ``next_run_at`` passes but jobs never
|
||||
fire and ``last_run_at`` stays null — the most common cron support report (#51038).
|
||||
The cron ticker only runs inside the gateway (no standalone daemon): without one,
|
||||
``next_run_at`` passes but jobs never fire and ``last_run_at`` stays null — the most
|
||||
common cron support report.
|
||||
"""
|
||||
# _builtin_gateway_liveness never raises (it maps probe failures to None),
|
||||
# so no guard is needed here — False is the only warn-worthy state.
|
||||
# _builtin_gateway_liveness never raises; False is the only warn-worthy state.
|
||||
if _builtin_gateway_liveness() is not False:
|
||||
return
|
||||
|
||||
@@ -138,10 +128,10 @@ def _format_lateness(seconds: float) -> str:
|
||||
|
||||
|
||||
def _dispatch_display(dispatch: dict) -> Optional[str]:
|
||||
"""One-line scheduled-vs-actual dispatch summary for a job (#99879).
|
||||
"""One-line scheduled-vs-actual dispatch summary; None when the stamp is malformed.
|
||||
|
||||
None when the stamp is malformed. On-time dispatches render dim; late/catch-up dispatches render
|
||||
loudly so a run fired long after gateway downtime doesn't look like an ordinary success.
|
||||
On-time dispatches render dim; late/catch-up dispatches render loudly so a run fired long
|
||||
after gateway downtime doesn't look like an ordinary success.
|
||||
"""
|
||||
if not isinstance(dispatch, dict):
|
||||
return None
|
||||
@@ -170,6 +160,10 @@ def _print_banner(title: str) -> None:
|
||||
print()
|
||||
|
||||
|
||||
def _unverified_targets(unverified) -> str:
|
||||
return ", ".join(str(t) for t in unverified) if isinstance(unverified, list) else str(unverified)
|
||||
|
||||
|
||||
def cron_list(show_all: bool = False):
|
||||
"""List all scheduled jobs."""
|
||||
from cron.jobs import effective_job_state, list_jobs
|
||||
@@ -187,23 +181,18 @@ def cron_list(show_all: bool = False):
|
||||
job_id = job.get("id", "?")
|
||||
name = job.get("name", "(unnamed)")
|
||||
schedule = job.get("schedule_display", job.get("schedule", {}).get("value", "?"))
|
||||
# Derive from the scheduler-honoured flag — never show [paused] when
|
||||
# enabled=true (half-paused contradiction must not look frozen).
|
||||
# Scheduler-honoured flag — never show [paused] when enabled=true.
|
||||
state = effective_job_state(job)
|
||||
next_run = job.get("next_run_at", "?")
|
||||
|
||||
# `repeat` may be present-but-null in the job record (e.g. a one-shot
|
||||
# job persisted with "repeat": null), so coalesce to {} rather than
|
||||
# relying on the dict-default, which only applies to a missing key.
|
||||
# `repeat` / `deliver` may be present-but-null in the record (a one-shot persisted
|
||||
# with "repeat": null), so coalesce rather than rely on the dict-default, which only
|
||||
# applies to a missing key — a null deliver would crash `", ".join(None)`.
|
||||
repeat_info = job.get("repeat") or {}
|
||||
repeat_times = repeat_info.get("times")
|
||||
repeat_completed = repeat_info.get("completed", 0)
|
||||
repeat_str = f"{repeat_completed}/{repeat_times}" if repeat_times else "∞"
|
||||
|
||||
# `deliver` may be present-but-null in the job record (same pitfall as
|
||||
# `repeat` above), so coalesce to the default rather than relying on the
|
||||
# dict-default, which only applies to a missing key. A null value would
|
||||
# otherwise reach `", ".join(None)` and crash the whole listing (#32896).
|
||||
deliver = job.get("deliver") or ["local"]
|
||||
if isinstance(deliver, str):
|
||||
deliver = [deliver]
|
||||
@@ -242,16 +231,14 @@ def cron_list(show_all: bool = False):
|
||||
if workdir:
|
||||
print(f" Workdir: {workdir}")
|
||||
|
||||
# Execution history
|
||||
last_status = job.get("last_status")
|
||||
if last_status:
|
||||
last_run = job.get("last_run_at", "?")
|
||||
if last_status == "ok":
|
||||
status_display = color("ok", Colors.GREEN)
|
||||
elif last_status == "delivery_failed":
|
||||
# The agent succeeded but the result never reached the user —
|
||||
# not green, and the detail lives in last_delivery_error
|
||||
# (last_error is None for these runs).
|
||||
# Agent succeeded but the result never reached the user — not green; the
|
||||
# detail lives in last_delivery_error (last_error is None for these runs).
|
||||
detail = job.get("last_delivery_error") or "?"
|
||||
status_display = color(f"delivery_failed: {detail}", Colors.YELLOW)
|
||||
else:
|
||||
@@ -276,15 +263,13 @@ def cron_list(show_all: bool = False):
|
||||
if delivery_err:
|
||||
print(f" {color('⚠ Delivery failed:', Colors.YELLOW)} {delivery_err}")
|
||||
|
||||
# A live adapter acked the last send but returned no message_id /
|
||||
# raw_response (Slack/Matrix/Mattermost shape): accepted as delivered,
|
||||
# but say so here rather than only in a WARNING log line.
|
||||
# A live adapter acked the last send but returned no message_id / raw_response
|
||||
# (Slack/Matrix/Mattermost shape): accepted as delivered, but say so here.
|
||||
unverified = job.get("last_delivery_unverified")
|
||||
if unverified:
|
||||
targets = ", ".join(str(t) for t in unverified) if isinstance(unverified, list) else str(unverified)
|
||||
print(
|
||||
f" {color('⚠ Delivery UNVERIFIED:', Colors.YELLOW)} "
|
||||
f"adapter acked {targets} without message_id/raw_response"
|
||||
f"adapter acked {_unverified_targets(unverified)} without message_id/raw_response"
|
||||
)
|
||||
|
||||
fire_err = job.get("last_fire_error")
|
||||
@@ -305,9 +290,8 @@ def cron_tick():
|
||||
try:
|
||||
tick(verbose=True)
|
||||
except CronTickYielded as exc:
|
||||
# Not expected on this surface (a one-shot CLI process has no boot
|
||||
# fingerprint, so the yield gate is inert) — but if a future caller
|
||||
# records one, report cleanly instead of a traceback.
|
||||
# Not expected here (a one-shot CLI process has no boot fingerprint, so the yield
|
||||
# gate is inert) — report cleanly instead of a traceback if a future caller records one.
|
||||
print(color(f"✗ {exc}", Colors.YELLOW))
|
||||
print(
|
||||
" A fresher gateway process owns the runtime lock and will fire "
|
||||
@@ -315,10 +299,8 @@ def cron_tick():
|
||||
)
|
||||
return 1
|
||||
except OSError as exc:
|
||||
# tick() now propagates real lock-acquisition failures (EMFILE,
|
||||
# EACCES on open, ...) instead of swallowing them as contention
|
||||
# (#87644). For the one-shot CLI surface, report cleanly instead of
|
||||
# dumping a traceback; the gateway ticker loop handles its own retry.
|
||||
# tick() propagates real lock-acquisition failures (EMFILE, EACCES on open, ...)
|
||||
# instead of swallowing them as contention; the gateway ticker loop handles its own retry.
|
||||
print(color(f"✗ Cron tick failed: {exc}", Colors.RED))
|
||||
print(" Check `hermes cron status` and the gateway log for details.")
|
||||
return 1
|
||||
@@ -343,19 +325,14 @@ def cron_runs(job_id: Optional[str] = None, limit: int = 20):
|
||||
print(f" {record['error']}")
|
||||
|
||||
|
||||
_INCIDENT_STATE_COLORS = {
|
||||
"detected": Colors.RED,
|
||||
"alerted": Colors.YELLOW,
|
||||
"closed": Colors.GREEN,
|
||||
}
|
||||
_INCIDENT_STATE_COLORS = {"detected": Colors.RED, "alerted": Colors.YELLOW, "closed": Colors.GREEN}
|
||||
|
||||
|
||||
def cron_incidents(args) -> int:
|
||||
"""List or acknowledge durable cron failure incidents.
|
||||
"""List (``[--state <s>]``) or ``ack <id>`` durable cron failure incidents.
|
||||
|
||||
``hermes cron incidents [--state <s>]`` lists incidents (the stored error is redacted and
|
||||
truncated at write time, safe for terminal display); ``hermes cron incidents ack <id>`` closes
|
||||
one so its failure ping stays silent until the error signature changes.
|
||||
The stored error is redacted and truncated at write time, safe for terminal display;
|
||||
acking closes an incident so its failure ping stays silent until the error signature changes.
|
||||
"""
|
||||
from cron.incidents import ack_incident, list_incidents
|
||||
|
||||
@@ -394,13 +371,11 @@ def cron_incidents(args) -> int:
|
||||
if inc.get("output_file"):
|
||||
print(f" Output: {inc['output_file']}")
|
||||
print()
|
||||
print(
|
||||
color(
|
||||
f" {len(incidents)} incident(s) | ack one with: "
|
||||
"hermes cron incidents ack <id>",
|
||||
Colors.DIM,
|
||||
)
|
||||
)
|
||||
print(color(
|
||||
f" {len(incidents)} incident(s) | ack one with: "
|
||||
"hermes cron incidents ack <id>",
|
||||
Colors.DIM,
|
||||
))
|
||||
return 0
|
||||
|
||||
|
||||
@@ -421,9 +396,8 @@ _FD_EXHAUSTION_HINT = (
|
||||
def _print_ticker_health(pids: list) -> None:
|
||||
"""Report builtin-ticker liveness for a gateway process known to be alive.
|
||||
|
||||
The gateway PROCESS is alive — but the cron ticker THREAD inside it can die silently, or stay
|
||||
alive while every tick fails. Check both the liveness heartbeat and the last-successful-tick
|
||||
marker so we don't report "will fire" when the ticker is dead or failing (#32612, #32895).
|
||||
The ticker THREAD can die silently or stay alive while every tick fails, so check both
|
||||
the liveness heartbeat and the last-successful-tick marker before saying "will fire".
|
||||
"""
|
||||
from cron.jobs import (
|
||||
get_ticker_heartbeat_age,
|
||||
@@ -433,10 +407,9 @@ def _print_ticker_health(pids: list) -> None:
|
||||
)
|
||||
from cron.scheduler import _is_fd_exhaustion_text as _cron_is_fd_exhaustion_text
|
||||
|
||||
# Allow ~3 missed ticker iterations (+ a little slack) before declaring
|
||||
# trouble. Derived from the shared interval constant so this threshold
|
||||
# tracks the ticker cadence instead of assuming a hardcoded 60s.
|
||||
STALE_AFTER = TICKER_INTERVAL_SECONDS * 3 + 20 # = 200s at the 60s default
|
||||
# ~3 missed ticker iterations (+ slack) before declaring trouble; derived from the shared
|
||||
# interval so the threshold tracks the ticker cadence (= 200s at the 60s default).
|
||||
STALE_AFTER = TICKER_INTERVAL_SECONDS * 3 + 20
|
||||
hb_age = get_ticker_heartbeat_age()
|
||||
ok_age = get_ticker_success_age()
|
||||
pid_line = f" PID: {', '.join(map(str, pids))}" if pids else None
|
||||
@@ -447,9 +420,8 @@ def _print_ticker_health(pids: list) -> None:
|
||||
print(pid_line)
|
||||
|
||||
if hb_age is None:
|
||||
# No heartbeat file means the ticker thread has never started: gateway
|
||||
# not in a cron-enabled profile, started moments ago (heartbeat is
|
||||
# written after startup), or a config issue blocking the ticker.
|
||||
# No heartbeat file: ticker never started (non-cron profile, gateway started moments
|
||||
# ago, or a config issue blocking the ticker).
|
||||
_warn("⚠ Gateway is running but the cron ticker has not reported a heartbeat.")
|
||||
print(" Cron jobs will NOT fire until the ticker writes its first heartbeat.")
|
||||
print(" If the gateway just started, wait ~60s and re-run `hermes cron status`.")
|
||||
@@ -462,18 +434,15 @@ def _print_ticker_health(pids: list) -> None:
|
||||
)
|
||||
print(" Cron jobs may NOT be firing. Restart: hermes gateway restart")
|
||||
elif ok_age is not None and ok_age > STALE_AFTER:
|
||||
# Loop is alive (fresh heartbeat) but no tick has SUCCEEDED in a
|
||||
# long time → ticks are failing every iteration.
|
||||
# Loop alive (fresh heartbeat) but no tick SUCCEEDED in a long time → failing every iteration.
|
||||
_warn(
|
||||
"⚠ Gateway and cron ticker are running, but no tick has "
|
||||
f"succeeded in {int(ok_age)}s — ticks may be failing."
|
||||
)
|
||||
last_error = get_ticker_last_error()
|
||||
if last_error:
|
||||
# Show WHY ticks fail — e.g. a root-rewritten jobs.json
|
||||
# (PermissionError) that silently locked out the ticker's uid for
|
||||
# ~14h in the field (#68483), or fd exhaustion (EMFILE) that used
|
||||
# to stall the scheduler invisibly (#87644).
|
||||
# Show WHY ticks fail — e.g. a root-rewritten jobs.json (PermissionError) that
|
||||
# silently locked out the ticker's uid, or fd exhaustion (EMFILE).
|
||||
print(color(f" Last tick error: {last_error}", Colors.RED))
|
||||
if "Permission denied" in last_error:
|
||||
print(color(_PERMISSION_HINT, Colors.YELLOW))
|
||||
@@ -497,13 +466,9 @@ def cron_status():
|
||||
|
||||
provider = _active_cron_provider_name()
|
||||
if provider != "builtin":
|
||||
# An external provider (e.g. Chronos) does NOT run the in-process 60s
|
||||
# ticker — it arms one external one-shot per job and is fired by a
|
||||
# NAS-mediated webhook, so between fires there is intentionally NO
|
||||
# ticker thread and NO heartbeat file. Reporting the ticker-heartbeat
|
||||
# staleness here would always say "stalled / not firing" on a perfectly
|
||||
# healthy Chronos instance. Report the provider instead and skip the
|
||||
# ticker-liveness heuristics entirely.
|
||||
# An external provider (e.g. Chronos) arms one external one-shot per job, fired by a
|
||||
# NAS-mediated webhook: between fires there is intentionally NO ticker thread and NO
|
||||
# heartbeat file, so the ticker-liveness heuristics would always say "stalled".
|
||||
print(color(
|
||||
f"✓ Cron provider: {provider} — jobs fire via the managed scheduler, "
|
||||
"not the in-process ticker.",
|
||||
@@ -518,11 +483,9 @@ def cron_status():
|
||||
pids = find_gateway_pids()
|
||||
gateway_alive_via_lock = False
|
||||
if not pids:
|
||||
# Same false-alarm class the cronjob tool fixed (#95947): the pid scan
|
||||
# can transiently miss a live gateway (just after a restart) while the
|
||||
# runtime lock — held for exactly the gateway's lifetime — proves the
|
||||
# ticker's process is alive. Only declare "not running" when both the
|
||||
# scan AND the lock say so.
|
||||
# The pid scan can transiently miss a live gateway (just after a restart) while
|
||||
# the runtime lock — held for exactly the gateway's lifetime — proves the ticker's
|
||||
# process is alive. Only declare "not running" when both agree.
|
||||
try:
|
||||
from gateway.status import get_running_pid, is_gateway_runtime_lock_active
|
||||
|
||||
@@ -548,50 +511,40 @@ def cron_status():
|
||||
print()
|
||||
|
||||
|
||||
|
||||
def _print_active_jobs_summary(jobs) -> None:
|
||||
"""Print the '<N> active job(s)' + next-run line shared by every status
|
||||
path (built-in ticker AND external provider)."""
|
||||
if jobs:
|
||||
next_runs = [j.get("next_run_at") for j in jobs if j.get("next_run_at")]
|
||||
print(f" {len(jobs)} active job(s)")
|
||||
if next_runs:
|
||||
print(f" Next run: {min(next_runs)}")
|
||||
# Missed-run visibility (#99879): call out jobs whose LAST dispatch
|
||||
# was late or a catch-up so post-downtime late fires are visible at
|
||||
# status level, not just buried per-job in `hermes cron list`.
|
||||
late = [
|
||||
j for j in jobs
|
||||
if isinstance(j.get("last_dispatch"), dict)
|
||||
and j["last_dispatch"].get("kind") in ("late", "catch_up")
|
||||
]
|
||||
if late:
|
||||
print()
|
||||
print(color(
|
||||
f" ⚠ {len(late)} job(s) last fired late (missed-fire catch-up):",
|
||||
Colors.YELLOW,
|
||||
))
|
||||
for j in late:
|
||||
d = j["last_dispatch"]
|
||||
print(
|
||||
f" {j.get('id', '?')} {j.get('name', '(unnamed)')}: "
|
||||
f"scheduled {d.get('scheduled_at', '?')}, "
|
||||
f"ran {d.get('dispatched_at', '?')} "
|
||||
+ color(
|
||||
f"({_format_lateness(d.get('lateness_seconds', 0))} late)",
|
||||
Colors.YELLOW,
|
||||
)
|
||||
)
|
||||
else:
|
||||
"""Print the '<N> active job(s)' + next-run line shared by every status path."""
|
||||
if not jobs:
|
||||
print(" No active jobs")
|
||||
return
|
||||
next_runs = [j.get("next_run_at") for j in jobs if j.get("next_run_at")]
|
||||
print(f" {len(jobs)} active job(s)")
|
||||
if next_runs:
|
||||
print(f" Next run: {min(next_runs)}")
|
||||
# Missed-run visibility: call out jobs whose LAST dispatch was late or a catch-up so
|
||||
# post-downtime late fires show at status level, not just per-job in `cron list`.
|
||||
late = [
|
||||
j for j in jobs
|
||||
if isinstance(j.get("last_dispatch"), dict)
|
||||
and j["last_dispatch"].get("kind") in ("late", "catch_up")
|
||||
]
|
||||
if late:
|
||||
print()
|
||||
print(color(f" ⚠ {len(late)} job(s) last fired late (missed-fire catch-up):", Colors.YELLOW))
|
||||
for j in late:
|
||||
d = j["last_dispatch"]
|
||||
print(
|
||||
f" {j.get('id', '?')} {j.get('name', '(unnamed)')}: "
|
||||
f"scheduled {d.get('scheduled_at', '?')}, "
|
||||
f"ran {d.get('dispatched_at', '?')} "
|
||||
+ color(f"({_format_lateness(d.get('lateness_seconds', 0))} late)", Colors.YELLOW)
|
||||
)
|
||||
|
||||
|
||||
def _scripts_dir_for_cron() -> Path:
|
||||
"""Return the scripts directory used by cron jobs.
|
||||
"""Scripts directory used by cron jobs.
|
||||
|
||||
Prefer ``cron.jobs.CRON_DIR.parent`` over a fresh ``get_hermes_home()`` call so tests and
|
||||
profile-aware callers that monkeypatch cron storage inspect the same Hermes home the jobs were
|
||||
loaded from.
|
||||
``cron.jobs.CRON_DIR.parent`` rather than a fresh ``get_hermes_home()`` so tests and
|
||||
profile-aware callers that monkeypatch cron storage inspect the same Hermes home.
|
||||
"""
|
||||
from cron.jobs import CRON_DIR
|
||||
|
||||
@@ -599,7 +552,7 @@ def _scripts_dir_for_cron() -> Path:
|
||||
|
||||
|
||||
def _script_health_issue(script: str) -> Optional[str]:
|
||||
"""Return a human-readable script issue, or ``None`` when the path is OK."""
|
||||
"""Human-readable script issue, or ``None`` when the path is OK."""
|
||||
scripts_dir = _scripts_dir_for_cron().resolve()
|
||||
raw = Path(script).expanduser()
|
||||
path = raw.resolve() if raw.is_absolute() else (scripts_dir / raw).resolve()
|
||||
@@ -616,15 +569,14 @@ def _script_health_issue(script: str) -> Optional[str]:
|
||||
return None
|
||||
|
||||
|
||||
# Grace period before an overdue ``next_run_at`` is reported. The ticker runs
|
||||
# once a minute and a busy tick can push dispatch a few minutes late; only a
|
||||
# next_run_at parked well in the past means the job is silently not firing
|
||||
# (ticker dead, gateway down, or a wedged fire-claim).
|
||||
# Grace before an overdue ``next_run_at`` is reported: the ticker runs once a minute and a
|
||||
# busy tick can push dispatch a few minutes late; only a next_run_at parked well in the past
|
||||
# means the job is silently not firing (ticker dead, gateway down, wedged fire-claim).
|
||||
_OVERDUE_GRACE_SECONDS = 15 * 60
|
||||
|
||||
|
||||
def _next_run_overdue_issue(next_run: str) -> Optional[str]:
|
||||
"""Return an issue string when ``next_run_at`` is parked in the past."""
|
||||
"""Issue string when ``next_run_at`` is parked in the past."""
|
||||
from datetime import datetime, timezone
|
||||
|
||||
try:
|
||||
@@ -646,9 +598,8 @@ def _cron_doctor_issues_for_job(job: Dict[str, Any]) -> List[str]:
|
||||
issues: List[str] = []
|
||||
|
||||
last_status = str(job.get("last_status") or "").strip().lower()
|
||||
# "delivery_failed" means the agent run itself succeeded, so it is not a
|
||||
# failed last run — the dedicated delivery issue below reports it (and
|
||||
# last_error is None, which would render as "unknown error" here).
|
||||
# "delivery_failed" means the agent run itself succeeded: the dedicated delivery issue
|
||||
# below reports it (and last_error is None, which would render as "unknown error" here).
|
||||
if last_status and last_status not in {"ok", "delivery_failed"}:
|
||||
err = str(job.get("last_error") or "unknown error").strip()
|
||||
issues.append(f"last run failed: {err}")
|
||||
@@ -659,8 +610,9 @@ def _cron_doctor_issues_for_job(job: Dict[str, Any]) -> List[str]:
|
||||
|
||||
unverified = job.get("last_delivery_unverified")
|
||||
if unverified:
|
||||
targets = ", ".join(str(t) for t in unverified) if isinstance(unverified, list) else str(unverified)
|
||||
issues.append(f"last delivery unverified (adapter acked without evidence): {targets}")
|
||||
issues.append(
|
||||
f"last delivery unverified (adapter acked without evidence): {_unverified_targets(unverified)}"
|
||||
)
|
||||
|
||||
if job.get("enabled", True) and job.get("state") not in {"paused", "completed"}:
|
||||
next_run = str(job.get("next_run_at") or "").strip()
|
||||
@@ -720,17 +672,10 @@ def cron_doctor() -> int:
|
||||
|
||||
|
||||
_JOB_ARG_FIELDS = (
|
||||
("name", "name"),
|
||||
("deliver", "deliver"),
|
||||
("failure_deliver", "failure_deliver"),
|
||||
("repeat", "repeat"),
|
||||
("script", "script"),
|
||||
("workdir", "workdir"),
|
||||
("model", "model"),
|
||||
("provider", "model_provider"),
|
||||
("monitor_script", "monitor_script"),
|
||||
("monitor_url", "monitor_url"),
|
||||
("continuity", "continuity"),
|
||||
("name", "name"), ("deliver", "deliver"), ("failure_deliver", "failure_deliver"),
|
||||
("repeat", "repeat"), ("script", "script"), ("workdir", "workdir"), ("model", "model"),
|
||||
("provider", "model_provider"), ("monitor_script", "monitor_script"),
|
||||
("monitor_url", "monitor_url"), ("continuity", "continuity"),
|
||||
("reasoning_effort", "reasoning_effort"),
|
||||
)
|
||||
|
||||
@@ -757,12 +702,9 @@ def _print_job_details(job_data: Dict[str, Any]) -> None:
|
||||
|
||||
|
||||
def cron_create(args):
|
||||
# The gateway-lifecycle guard lives in cron.jobs.create_job so it fires on
|
||||
# every job-creation path (this CLI subcommand AND the agent's `cronjob`
|
||||
# model tool, which calls create_job directly). When it blocks, create_job
|
||||
# raises GatewayLifecycleBlocked, the `cronjob` tool wrapper catches it and
|
||||
# returns it as result["error"], and the `if not result.get("success")`
|
||||
# branch below prints it in red and exits 1 — same UX as before.
|
||||
# The gateway-lifecycle guard lives in cron.jobs.create_job so it fires on every
|
||||
# job-creation path (CLI AND the agent's `cronjob` tool); a block surfaces as
|
||||
# result["error"] and is printed in red below.
|
||||
result = _cron_api(
|
||||
action="create",
|
||||
schedule=args.schedule,
|
||||
@@ -844,16 +786,13 @@ def cron_edit(args):
|
||||
def _job_action(action: str, job_id: str, success_verb: str) -> int:
|
||||
_stateless_token = None
|
||||
if action == "run":
|
||||
# One-shot CLI: this process exits as soon as the command returns, so
|
||||
# a background-dispatched run (daemon thread of THIS process) would be
|
||||
# orphaned mid-LLM-call — the delegation dies 'unknown' and the job's
|
||||
# execution row is stuck 'claimed', blocking future runs (#86721).
|
||||
# The background path in ``_try_dispatch_background_run`` triggers when
|
||||
# the CLI inherits a gateway/desktop session env (HERMES_SESSION_KEY);
|
||||
# declare the channel stateless so ``async_delivery_supported()`` gates
|
||||
# it off and the run executes synchronously to completion instead.
|
||||
# The declaration is scoped to this call (token reset in ``finally``)
|
||||
# so in-process callers (tests, embedding apps) are not tainted.
|
||||
# One-shot CLI: the process exits as soon as the command returns, so a
|
||||
# background-dispatched run (daemon thread of THIS process — triggered when the CLI
|
||||
# inherits a gateway/desktop session env, HERMES_SESSION_KEY) would be orphaned
|
||||
# mid-LLM-call, leaving the execution row stuck 'claimed'. Declare the channel
|
||||
# stateless so ``async_delivery_supported()`` gates it off and the run executes
|
||||
# synchronously. Scoped to this call (token reset in ``finally``) so in-process
|
||||
# callers (tests, embedding apps) are not tainted.
|
||||
try:
|
||||
from gateway.session_context import _SESSION_ASYNC_DELIVERY
|
||||
|
||||
@@ -874,14 +813,9 @@ def _job_action(action: str, job_id: str, success_verb: str) -> int:
|
||||
print(f" Next run: {result['job']['next_run_at']}")
|
||||
if action == "run":
|
||||
job = result.get("job", {})
|
||||
# A manual run can be dispatched to the gateway daemon's background
|
||||
# delegation worker instead of executing inline (e.g. when the CLI
|
||||
# process inherits a gateway/desktop session env and the run
|
||||
# resolves a session key). Such responses carry
|
||||
# execution_mode="background" and/or a delegation_id, and the job
|
||||
# keeps running AFTER this CLI process exits — a terminal
|
||||
# success/failure verdict would be a lie (#83340). Report the
|
||||
# background dispatch instead of claiming the run failed.
|
||||
# A manual run may be dispatched to the gateway daemon's background delegation worker
|
||||
# (execution_mode="background" and/or a delegation_id) and keeps running AFTER this
|
||||
# CLI exits — a terminal success/failure verdict would be a lie, so report the dispatch.
|
||||
delegation_id = job.get("delegation_id")
|
||||
if job.get("execution_mode") == "background" or delegation_id:
|
||||
if delegation_id:
|
||||
@@ -927,9 +861,9 @@ def cron_resume(args) -> int:
|
||||
def cron_notepad(args) -> int:
|
||||
"""Handle ``hermes cron notepad <job_id> [get|set|delete|list]``.
|
||||
|
||||
This CLI is the write path for the per-job durable KV scratchpad (``cron/notepad.py``): a
|
||||
running cron agent updates its own notepad via its terminal tool, and the scheduler injects
|
||||
non-empty notepads into the job prompt on each run.
|
||||
Write path for the per-job durable KV scratchpad (``cron/notepad.py``): a running cron
|
||||
agent updates its own notepad via its terminal tool, and the scheduler injects non-empty
|
||||
notepads into the job prompt on each run.
|
||||
"""
|
||||
from cron import notepad
|
||||
|
||||
@@ -977,46 +911,26 @@ def cron_notepad(args) -> int:
|
||||
return 1
|
||||
|
||||
|
||||
def _cron_list_cmd(args) -> int:
|
||||
cron_list(getattr(args, "all", False))
|
||||
return 0
|
||||
|
||||
|
||||
def _cron_status_cmd(args) -> int:
|
||||
cron_status()
|
||||
return 0
|
||||
|
||||
|
||||
def _cron_runs_cmd(args) -> int:
|
||||
cron_runs(getattr(args, "job_id", None), getattr(args, "limit", 20))
|
||||
return 0
|
||||
|
||||
|
||||
def _cron_remove_cmd(args) -> int:
|
||||
return _job_action("remove", args.job_id, "Removed")
|
||||
|
||||
|
||||
# Subcommand -> handler. Late-bound lambdas so module-level monkeypatching of
|
||||
# the underlying functions keeps working.
|
||||
# Subcommand -> handler. Late-bound lambdas so module-level monkeypatching of the underlying
|
||||
# functions keeps working. cron_list/cron_status/cron_runs always return None -> exit 0.
|
||||
_CRON_SUBCOMMANDS = {
|
||||
"list": _cron_list_cmd,
|
||||
"status": _cron_status_cmd,
|
||||
"doctor": lambda args: cron_doctor(),
|
||||
"tick": lambda args: cron_tick(),
|
||||
"runs": _cron_runs_cmd,
|
||||
"history": _cron_runs_cmd,
|
||||
"incidents": lambda args: cron_incidents(args),
|
||||
"notepad": lambda args: cron_notepad(args),
|
||||
"create": lambda args: cron_create(args),
|
||||
"add": lambda args: cron_create(args),
|
||||
"edit": lambda args: cron_edit(args),
|
||||
"pause": lambda args: _job_action("pause", args.job_id, "Paused"),
|
||||
"resume": lambda args: cron_resume(args),
|
||||
"run": lambda args: _job_action("run", args.job_id, "Triggered"),
|
||||
"remove": _cron_remove_cmd,
|
||||
"rm": _cron_remove_cmd,
|
||||
"delete": _cron_remove_cmd,
|
||||
"list": lambda a: cron_list(getattr(a, "all", False)) or 0,
|
||||
"status": lambda a: cron_status() or 0,
|
||||
"doctor": lambda a: cron_doctor(),
|
||||
"tick": lambda a: cron_tick(),
|
||||
"runs": lambda a: cron_runs(getattr(a, "job_id", None), getattr(a, "limit", 20)) or 0,
|
||||
"incidents": lambda a: cron_incidents(a),
|
||||
"notepad": lambda a: cron_notepad(a),
|
||||
"create": lambda a: cron_create(a),
|
||||
"edit": lambda a: cron_edit(a),
|
||||
"pause": lambda a: _job_action("pause", a.job_id, "Paused"),
|
||||
"resume": lambda a: cron_resume(a),
|
||||
"run": lambda a: _job_action("run", a.job_id, "Triggered"),
|
||||
"remove": lambda a: _job_action("remove", a.job_id, "Removed"),
|
||||
}
|
||||
_CRON_SUBCOMMANDS["history"] = _CRON_SUBCOMMANDS["runs"]
|
||||
_CRON_SUBCOMMANDS["add"] = _CRON_SUBCOMMANDS["create"]
|
||||
_CRON_SUBCOMMANDS["rm"] = _CRON_SUBCOMMANDS["delete"] = _CRON_SUBCOMMANDS["remove"]
|
||||
|
||||
|
||||
def cron_command(args):
|
||||
|
||||
Reference in New Issue
Block a user