2191 lines
101 KiB
Python
2191 lines
101 KiB
Python
"""
|
|
Cron job management tools for Hermes Agent.
|
|
|
|
Expose a single compressed action-oriented tool to avoid schema/context bloat.
|
|
Compatibility wrappers remain for direct Python callers and legacy tests.
|
|
"""
|
|
|
|
import json
|
|
import logging
|
|
import re
|
|
import sys
|
|
import threading
|
|
import time
|
|
from pathlib import Path
|
|
from typing import Any, Dict, List, Optional, Union
|
|
|
|
from hermes_constants import display_hermes_home
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# Cadence for the heartbeat that keeps the calling agent's inactivity watchdog
|
|
# at bay while a manual `cronjob(action="run")` executes the job synchronously
|
|
# in-process (#76502). Mirrors the 10s cadence of
|
|
# tools/environments/base.py::touch_activity_if_due (delegate_task's heartbeat
|
|
# uses 30s) — comfortably below the 1800s default HERMES_AGENT_TIMEOUT.
|
|
_CRON_RUN_HEARTBEAT_INTERVAL = 10.0
|
|
|
|
# Hard ceiling on how long the heartbeat keeps the parent watchdog at bay.
|
|
# The child cron run has its own inactivity watchdog (HERMES_CRON_TIMEOUT,
|
|
# default 600s) that bounds a wedged job, but with HERMES_CRON_TIMEOUT=0
|
|
# (explicit "unlimited") a truly hung run_one_job would otherwise mask the
|
|
# gateway watchdog forever — pre-#76502 the parent was at least reaped at
|
|
# ~1800s. After this ceiling the heartbeat stops and the gateway watchdog
|
|
# regains authority over the turn.
|
|
_CRON_RUN_HEARTBEAT_CEILING = 6 * 3600.0
|
|
|
|
# Import from cron module (will be available when properly installed)
|
|
sys.path.insert(0, str(Path(__file__).parent.parent))
|
|
|
|
from cron.jobs import (
|
|
AmbiguousJobReference,
|
|
claim_job_for_fire,
|
|
effective_job_state,
|
|
get_job,
|
|
is_job_runnable,
|
|
list_jobs,
|
|
mark_job_run,
|
|
parse_schedule,
|
|
pause_job,
|
|
remove_job,
|
|
resolve_job_ref,
|
|
resume_job,
|
|
update_job,
|
|
)
|
|
|
|
|
|
def _notify_provider_jobs_changed_safe() -> None:
|
|
"""Tell the active cron scheduler provider the job set changed (no-op for
|
|
the built-in). Best-effort — never lets a provider error break the tool."""
|
|
try:
|
|
from cron.scheduler import _notify_provider_jobs_changed
|
|
_notify_provider_jobs_changed()
|
|
except Exception:
|
|
pass
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Cron prompt scanning
|
|
# ---------------------------------------------------------------------------
|
|
#
|
|
# Two threat surfaces, two scanners:
|
|
#
|
|
# 1. User-supplied cron prompt (small, written as a directive).
|
|
# Strict scanning is appropriate — a legit cron prompt has no business
|
|
# saying "cat ~/.hermes/.env" or "rm -rf /". `_scan_cron_prompt()` runs
|
|
# against this at create/update time and as a runtime defense-in-depth.
|
|
#
|
|
# 2. Assembled prompt that includes loaded skill content (large markdown
|
|
# bodies, often security docs, postmortems, runbooks discussing attack
|
|
# patterns in PROSE). Reusing the strict patterns here false-positives
|
|
# every time a skill *describes* a command — see #3968 follow-up: the
|
|
# `hermes-agent-dev` skill contains a security postmortem mentioning
|
|
# `cat ~/.hermes/.env`, which tripped `read_secrets` and silently
|
|
# killed all PR-scout jobs.
|
|
#
|
|
# Skill bodies are user-curated and scanned at install time by
|
|
# `skills_guard.py`. The runtime cron scan only needs to catch the
|
|
# patterns whose phrasing does NOT survive normal English prose:
|
|
# classic prompt-injection directives ("ignore previous instructions",
|
|
# "disregard your rules"), deception directives, and invisible
|
|
# unicode. `_scan_cron_skill_assembled()` runs against the assembled
|
|
# prompt with this tighter pattern set.
|
|
#
|
|
# Both scanners share the invisible-unicode check and the GitHub Authorization
|
|
# header exemption.
|
|
|
|
# Strict patterns — applied to the user prompt only.
|
|
_CRON_THREAT_PATTERNS = [
|
|
(r'ignore\s+(?:\w+\s+)*(?:previous|all|above|prior)\s+(?:\w+\s+)*instructions', "prompt_injection"),
|
|
(r'do\s+not\s+tell\s+the\s+user', "deception_hide"),
|
|
(r'system\s+prompt\s+override', "sys_prompt_override"),
|
|
(r'disregard\s+(your|all|any)\s+(instructions|rules|guidelines)', "disregard_rules"),
|
|
(r'cat\s+[^\n]*(\.env|credentials|\.netrc|\.pgpass|id_rsa|id_ed25519|id_ecdsa)', "read_secrets"),
|
|
(r'authorized_keys', "ssh_backdoor"),
|
|
(r'/etc/sudoers|visudo', "sudoers_mod"),
|
|
(r'rm\s+-rf\s+/', "destructive_root_rm"),
|
|
]
|
|
|
|
# Looser pattern set — applied to the assembled prompt when skills are
|
|
# attached. Only patterns whose phrasing is unambiguous in any context;
|
|
# command-shape patterns are dropped because they false-positive on prose
|
|
# in security docs / postmortems. Skill bodies are scanned at install time
|
|
# by `skills_guard.py`, so the runtime cron scan is purely a tripwire for
|
|
# obvious injection directives surviving a malicious skill that slipped
|
|
# through install.
|
|
_CRON_SKILL_ASSEMBLED_PATTERNS = [
|
|
(r'ignore\s+(?:\w+\s+)*(?:previous|all|above|prior)\s+(?:\w+\s+)*instructions', "prompt_injection"),
|
|
(r'do\s+not\s+tell\s+the\s+user', "deception_hide"),
|
|
(r'system\s+prompt\s+override', "sys_prompt_override"),
|
|
(r'disregard\s+(your|all|any)\s+(instructions|rules|guidelines)', "disregard_rules"),
|
|
]
|
|
|
|
_CRON_SECRET_VAR_RE = r'\$\{?\w*(?:KEY|TOKEN|SECRET|PASSWORD|CREDENTIAL|API)\w*\}?'
|
|
_CRON_EXFIL_COMMAND_PATTERNS = [
|
|
# Tighten exfil detection to obvious leak paths: embedding a secret
|
|
# directly in the destination URL, sending it in POST/FORM payloads,
|
|
# or shipping it via Authorization headers to arbitrary hosts. The
|
|
# only intended allowlist exception today is the bundled GitHub skill
|
|
# pattern that talks to api.github.com.
|
|
(rf'curl\s+[^\n]*https?://[^\s"\'`]*{_CRON_SECRET_VAR_RE}', "exfil_curl_url"),
|
|
(rf'wget\s+[^\n]*https?://[^\s"\'`]*{_CRON_SECRET_VAR_RE}', "exfil_wget_url"),
|
|
(rf'curl\s+[^\n]*(?:--data(?:-raw|-binary|-urlencode)?|-d|--form|-F)\s+[^\n]*{_CRON_SECRET_VAR_RE}', "exfil_curl_data"),
|
|
(rf'wget\s+[^\n]*--post-(?:data|file)=[^\n]*{_CRON_SECRET_VAR_RE}', "exfil_wget_post"),
|
|
(rf'curl\s+[^\n]*(?:-H|--header)\s+["\']Authorization:\s*(?:Bearer|token)\s+{_CRON_SECRET_VAR_RE}["\']', "exfil_curl_auth_header"),
|
|
]
|
|
|
|
# Single source of truth, shared with the install-time scanner
|
|
# (threat_patterns.INVISIBLE_CHARS / skills_guard). Keeping a separate, narrower
|
|
# copy here let an obfuscated injection directive slip past this runtime cron
|
|
# tripwire while being caught at install time (or vice versa): U+2062-U+2064
|
|
# (invisible math operators) and U+2066-U+2069 (directional isolates) are real
|
|
# attack tools and were missing from the cron-local set. Importing the canonical
|
|
# set keeps the cron tripwire and the install scanner from drifting apart.
|
|
from tools.threat_patterns import INVISIBLE_CHARS as _CRON_INVISIBLE_CHARS
|
|
|
|
# U+200D Zero-Width Joiner is also a legitimate, required part of many
|
|
# Unicode emoji sequences (for example 👨👩👧, 🏳️🌈, ❤️🩹, 🧑💻).
|
|
# We should still block ZWJ when it is hiding between plain text characters,
|
|
# but not when it is clearly part of an emoji grapheme cluster.
|
|
_EMOJI_NEIGHBOUR_CP_RANGES = (
|
|
(0x1F000, 0x1FFFF),
|
|
(0x2600, 0x27BF),
|
|
(0x2300, 0x23FF),
|
|
(0x1F1E6, 0x1F1FF),
|
|
(0x20E3, 0x20E3),
|
|
)
|
|
_VARIATION_SELECTOR_CP = 0xFE0F
|
|
|
|
|
|
def _is_emoji_cp(cp: int) -> bool:
|
|
return any(lo <= cp <= hi for lo, hi in _EMOJI_NEIGHBOUR_CP_RANGES)
|
|
|
|
|
|
def _zwj_has_emoji_neighbour(text: str, idx: int) -> bool:
|
|
"""Return True when the ZWJ at text[idx] appears inside an emoji sequence."""
|
|
left = idx - 1
|
|
while left >= 0 and ord(text[left]) == _VARIATION_SELECTOR_CP:
|
|
left -= 1
|
|
right = idx + 1
|
|
while right < len(text) and ord(text[right]) == _VARIATION_SELECTOR_CP:
|
|
right += 1
|
|
return (
|
|
left >= 0 and right < len(text)
|
|
and _is_emoji_cp(ord(text[left]))
|
|
and _is_emoji_cp(ord(text[right]))
|
|
)
|
|
|
|
|
|
def _strip_legitimate_emoji_zwj(prompt: str) -> str:
|
|
if '\u200d' not in prompt:
|
|
return prompt
|
|
cleaned: list[str] = []
|
|
for idx, ch in enumerate(prompt):
|
|
if ch == '\u200d' and _zwj_has_emoji_neighbour(prompt, idx):
|
|
continue
|
|
cleaned.append(ch)
|
|
return ''.join(cleaned)
|
|
|
|
|
|
def _strip_cron_safe_constructs(prompt: str) -> str:
|
|
"""Strip the GitHub `Authorization: token $GITHUB_TOKEN` auth-header
|
|
pattern so it doesn't trip the broader curl-auth-header exfil rule.
|
|
|
|
Allows the bundled GitHub skill fallback without opening a blanket
|
|
exemption for arbitrary Authorization-header exfiltration.
|
|
|
|
Uses ``re.sub`` so EVERY occurrence is scrubbed, not just the first — a
|
|
cron job that loads 2+ GitHub skills (e.g. github-issues +
|
|
github-pr-workflow + github-code-review) contains several such blocks,
|
|
and the old ``re.search`` + single ``str.replace`` left the rest to trip
|
|
the exfil_curl_auth_header detector on every run. The trailing
|
|
``[^\\s;&|$`]*`` consumes only the URL path — never whitespace, command
|
|
separators, or subshell openers — so a payload smuggled onto the same
|
|
line (``;``, ``&&``, ``|``, ``$(...)``, backticks) survives the strip
|
|
and is still scanned. The host must be exactly ``api.github.com``
|
|
followed by ``/``, whitespace, quote, or end: lookalike authorities
|
|
(``api.github.com.evil.com``, ``api.github.com@evil.com``) are not the
|
|
trusted construct and fall through to the exfil detectors, while
|
|
legitimately quoted bare-host URLs stay exempt.
|
|
"""
|
|
return re.sub(
|
|
rf'curl\s+[^\n;&|$`]*(?:-H|--header)\s+["\']Authorization:\s*token\s+{_CRON_SECRET_VAR_RE}["\']'
|
|
r'\s+["\']?https://api\.github\.com(?::\d+)?(?:/|\s|$|["\'])[^\s;&|$`]*',
|
|
'curl https://api.github.com/user',
|
|
prompt,
|
|
flags=re.IGNORECASE,
|
|
)
|
|
|
|
|
|
def _check_invisible_unicode(prompt: str) -> str:
|
|
"""Return an error string if the prompt contains invisible-unicode
|
|
injection markers (ZWJ inside legitimate emoji sequences is allowed).
|
|
"""
|
|
prompt_for_invisible_scan = _strip_legitimate_emoji_zwj(prompt)
|
|
for char in _CRON_INVISIBLE_CHARS:
|
|
if char in prompt_for_invisible_scan:
|
|
return f"Blocked: prompt contains invisible unicode U+{ord(char):04X} (possible injection)."
|
|
return ""
|
|
|
|
|
|
def _strip_invisible_unicode(prompt: str) -> tuple[str, list[str]]:
|
|
"""Strip invisible-unicode characters from *prompt*, preserving the ZWJ
|
|
that lives inside legitimate emoji sequences.
|
|
|
|
Returns ``(cleaned_prompt, removed_codepoints)`` where ``removed_codepoints``
|
|
is the sorted list of ``U+XXXX`` labels that were stripped (empty when the
|
|
prompt was already clean). Used by the skills-attached cron path, where the
|
|
skill body is already vetted at install time by ``skills_guard.py`` — a
|
|
stray zero-width space in a code example should be sanitized, not turned
|
|
into a hard block that permanently kills the job.
|
|
"""
|
|
if not prompt:
|
|
return prompt, []
|
|
# Keep emoji-ZWJ: temporarily remove the legitimate joiners, scan/strip the
|
|
# rest, then the legitimate joiners survive because we operate on the
|
|
# original string and only drop chars that are NOT part of an emoji cluster.
|
|
removed: set[str] = set()
|
|
cleaned: list[str] = []
|
|
for idx, ch in enumerate(prompt):
|
|
if ch in _CRON_INVISIBLE_CHARS:
|
|
if ch == '\u200d' and _zwj_has_emoji_neighbour(prompt, idx):
|
|
cleaned.append(ch) # legitimate emoji joiner — keep
|
|
continue
|
|
removed.add(f"U+{ord(ch):04X}")
|
|
continue
|
|
cleaned.append(ch)
|
|
return ''.join(cleaned), sorted(removed)
|
|
|
|
|
|
def _scan_cron_prompt(prompt: str) -> str:
|
|
"""Scan the USER-SUPPLIED cron prompt for critical threats.
|
|
|
|
Strict pattern set — used at job create/update time and as a runtime
|
|
defense-in-depth for prompts authored before the scanner existed.
|
|
The user prompt is small and directive; bare `cat .env` or `rm -rf /`
|
|
there is a smoking gun, not prose. Returns an error string when
|
|
blocked, else empty string.
|
|
"""
|
|
prompt_to_scan = _strip_cron_safe_constructs(prompt)
|
|
invisible_err = _check_invisible_unicode(prompt_to_scan)
|
|
if invisible_err:
|
|
return invisible_err
|
|
for pattern, pid in _CRON_THREAT_PATTERNS:
|
|
if re.search(pattern, prompt_to_scan, re.IGNORECASE):
|
|
return f"Blocked: prompt matches threat pattern '{pid}'. Cron prompts must not contain injection or exfiltration payloads."
|
|
for pattern, pid in _CRON_EXFIL_COMMAND_PATTERNS:
|
|
if re.search(pattern, prompt_to_scan, re.IGNORECASE):
|
|
return f"Blocked: prompt matches threat pattern '{pid}'. Cron prompts must not contain injection or exfiltration payloads."
|
|
return ""
|
|
|
|
|
|
def _scan_cron_skill_assembled(assembled: str) -> tuple[str, str]:
|
|
"""Scan an ASSEMBLED cron prompt that includes loaded skill content.
|
|
|
|
Looser pattern set — only catches unambiguous prompt-injection
|
|
directives. Drops command-shape patterns (cat .env, rm -rf /,
|
|
authorized_keys, /etc/sudoers) because they false-positive on
|
|
legitimate skill markdown that *describes* attack commands in
|
|
security postmortems and runbooks.
|
|
|
|
Invisible unicode is SANITIZED, not blocked. Skill bodies are
|
|
user-curated and already scanned at install time by
|
|
``skills_guard.py``; a stray zero-width space in a code example
|
|
(common in copy-pasted unicode docs) should not permanently kill the
|
|
job. The offending codepoints are stripped and logged, the cleaned
|
|
prompt is returned. The hard block remains for raw user prompts via
|
|
``_scan_cron_prompt`` — that path is the actual injection surface.
|
|
|
|
Returns ``(cleaned_prompt, error)``; ``error`` is empty when the
|
|
prompt passed (after sanitization).
|
|
"""
|
|
cleaned, removed = _strip_invisible_unicode(assembled)
|
|
if removed:
|
|
logger.warning(
|
|
"Cron skill-assembled prompt: stripped %d invisible-unicode "
|
|
"char(s) (%s) from vetted skill content",
|
|
len(removed), ", ".join(removed),
|
|
)
|
|
prompt_to_scan = _strip_cron_safe_constructs(cleaned)
|
|
for pattern, pid in _CRON_SKILL_ASSEMBLED_PATTERNS:
|
|
if re.search(pattern, prompt_to_scan, re.IGNORECASE):
|
|
return cleaned, f"Blocked: prompt matches threat pattern '{pid}'. Cron prompts must not contain injection or exfiltration payloads."
|
|
return cleaned, ""
|
|
|
|
|
|
def _origin_from_env() -> Optional[Dict[str, str]]:
|
|
from gateway.session_context import get_session_env
|
|
origin_platform = get_session_env("HERMES_SESSION_PLATFORM")
|
|
origin_chat_id = get_session_env("HERMES_SESSION_CHAT_ID")
|
|
if origin_platform and origin_chat_id:
|
|
thread_id = get_session_env("HERMES_SESSION_THREAD_ID") or None
|
|
# Slack thread-per-message session keying (native parity: thread_ts =
|
|
# event.thread_ts or ts) stamps every TOP-LEVEL message's own id as
|
|
# the session thread. That stamp is a per-message session KEY, not a
|
|
# durable conversation location — persisting it as origin routing
|
|
# pins every future delivery inside the ephemeral thread spawned
|
|
# around the creation message. Recognize it at the source: a Slack
|
|
# thread id equal to the triggering message's own id is synthetic.
|
|
# A genuine in-thread creation (thread == the parent's id != this
|
|
# message's id) keeps its thread.
|
|
if thread_id and origin_platform == "slack":
|
|
message_id = get_session_env("HERMES_SESSION_MESSAGE_ID") or None
|
|
if message_id and str(thread_id) == str(message_id):
|
|
logger.debug(
|
|
"Cron origin: dropping synthetic per-message Slack "
|
|
"thread_id=%s (== creation message id)", thread_id,
|
|
)
|
|
thread_id = None
|
|
if thread_id:
|
|
logger.debug(
|
|
"Cron origin captured thread_id=%s for %s:%s",
|
|
thread_id, origin_platform, origin_chat_id,
|
|
)
|
|
return {
|
|
"platform": origin_platform,
|
|
"chat_id": origin_chat_id,
|
|
"chat_name": get_session_env("HERMES_SESSION_CHAT_NAME") or None,
|
|
"thread_id": thread_id,
|
|
# Captured so an opt-in delivery mirror (cron.mirror_delivery /
|
|
# attach_to_session) can resolve the exact participant's session in
|
|
# per-user-isolated group chats — parity with interactive
|
|
# send_message, which passes HERMES_SESSION_USER_ID to
|
|
# gateway.mirror.mirror_to_session. Harmless for DMs/shared sessions.
|
|
"user_id": get_session_env("HERMES_SESSION_USER_ID") or None,
|
|
# Workspace/server scope (Slack team, Discord guild, Matrix
|
|
# server). build_session_key embeds it in every Slack session key
|
|
# (dm/group/thread alike), so a continuable cron seed built
|
|
# WITHOUT it creates a row no scoped reply ever resolves to —
|
|
# the seeded key is agent:main:slack:dm:<chat>:<thread> while the
|
|
# reply keys agent:main:slack:dm:<team>:<chat>:<thread>. Captured
|
|
# here so the scheduler's seed helpers can reproduce the reply's
|
|
# exact key. Same session-context var async_delegation already
|
|
# snapshots; None for platforms without scope.
|
|
"scope_id": get_session_env("HERMES_SESSION_SCOPE_ID") or None,
|
|
}
|
|
return None
|
|
|
|
|
|
def _local_delivery_notice(job: Dict[str, Any], user_deliver: Optional[str]) -> Optional[str]:
|
|
"""Return an informational notice when a created job won't deliver anywhere.
|
|
|
|
TUI/CLI sessions cannot be captured as a cron ``origin`` (no
|
|
``HERMES_SESSION_PLATFORM``/``CHAT_ID`` is set for them), so a
|
|
``deliver="origin"`` request — or an omitted ``deliver`` that defaults to
|
|
origin-or-local — produces a job that runs and saves output to
|
|
``last_output`` but is never delivered back into the session. This is by
|
|
design (there is no live-delivery channel for local sessions), but silently
|
|
dropping the user's "tell me when it runs" intent is the trap reported in
|
|
#51568. Surface it at create time so the agent can relay it instead of
|
|
promising a delivery that never happens.
|
|
|
|
Returns ``None`` when the user explicitly asked for ``local`` (no surprise),
|
|
or when the job resolves to a real delivery target.
|
|
"""
|
|
# An explicit local request is exactly what the user asked for — no notice.
|
|
if (user_deliver or "").strip().lower() == "local":
|
|
return None
|
|
try:
|
|
from cron.scheduler import _resolve_delivery_targets
|
|
|
|
if _resolve_delivery_targets(job):
|
|
return None # Will actually deliver somewhere — nothing to flag.
|
|
except Exception:
|
|
# If resolution can't be evaluated, fall back to the origin signal.
|
|
if job.get("origin"):
|
|
return None
|
|
return (
|
|
"This is a local-only cron job: its output is saved (view it with "
|
|
"cronjob(action='list')) but will NOT be delivered back into this "
|
|
"session — CLI/TUI sessions have no live-delivery channel. To be "
|
|
"notified when it runs, recreate or update the job with deliver set to "
|
|
"a gateway-connected platform, e.g. deliver='telegram' or deliver='all'."
|
|
)
|
|
|
|
|
|
def _mode_guidance_notes(job: Dict[str, Any], user_deliver: Optional[str]) -> List[str]:
|
|
"""Mode-specific guidance echoed in the create/update response.
|
|
|
|
The teaching that used to live in CRONJOB_SCHEMA parameter descriptions
|
|
(paid for on every API call of every session) is delivered here instead —
|
|
once, in the tool result, at the moment the model actually created a job
|
|
in that mode. Keep each note short and actionable; only fire notes for
|
|
modes the job actually uses.
|
|
"""
|
|
notes: List[str] = []
|
|
if job.get("monitor_script") or job.get("monitor_url"):
|
|
notes.append(
|
|
"Monitor mode: the source runs first each tick and its output is "
|
|
"hashed as exact bytes — unchanged output suppresses the agent run "
|
|
"(silent no_change tick), changed output injects a MONITOR CHANGE "
|
|
"DETECTED diff into the prompt. The first tick always runs as "
|
|
"baseline. The source must emit STABLE output (no timestamps, no "
|
|
"random ordering) or every tick will look changed."
|
|
)
|
|
if job.get("no_agent"):
|
|
notes.append(
|
|
"no_agent mode: stdout is delivered verbatim; EMPTY stdout sends "
|
|
"nothing at all (watchdog pattern — script should stay quiet when "
|
|
"there is nothing to report). Non-zero exit or timeout sends an "
|
|
"error alert. prompt/skills are ignored."
|
|
)
|
|
_deliver = (user_deliver or "").strip().lower()
|
|
if _deliver:
|
|
if "all" in _deliver.split(","):
|
|
notes.append(
|
|
"deliver='all' resolves at fire time and never includes "
|
|
"bot-chat targets — channels connected later are picked up "
|
|
"automatically."
|
|
)
|
|
if _deliver.startswith("bot-chat:"):
|
|
notes.append(
|
|
"Targeting another profile's Bot Chat costs that bot an agent "
|
|
"turn per run."
|
|
)
|
|
# platform:chat_id with no thread segment loses topic targeting —
|
|
# warn once here instead of carrying the warning in the schema.
|
|
for target in _deliver.split(","):
|
|
parts = target.strip().split(":")
|
|
if (
|
|
len(parts) == 2
|
|
and parts[0] not in ("bot-chat", "sms")
|
|
and parts[1]
|
|
and not parts[1].startswith("#")
|
|
):
|
|
notes.append(
|
|
f"deliver target '{target.strip()}' has no :thread_id "
|
|
"segment — on thread/topic platforms the delivery lands in "
|
|
"the main chat, not a topic."
|
|
)
|
|
break
|
|
return notes
|
|
|
|
|
|
def _split_monitor_arg(
|
|
monitor: Optional[str],
|
|
monitor_script: Optional[str],
|
|
monitor_url: Optional[str],
|
|
) -> tuple:
|
|
"""Resolve the model-facing ``monitor`` field into the stored pair.
|
|
|
|
The schema advertises ONE ``monitor`` field; the value's shape decides the
|
|
transport: ``http(s)://...`` is a URL source, anything else is a script
|
|
path (a legal script path can never start with a URL scheme). Jobs keep
|
|
storing ``monitor_script``/``monitor_url`` separately — this is an
|
|
interface merge, not a storage migration — and the legacy field names are
|
|
still accepted as aliases so older transcripts/replays keep working.
|
|
|
|
Returns ``(monitor_script, monitor_url)`` with update semantics:
|
|
``None`` = leave unchanged, ``''`` = clear. Setting one source via
|
|
``monitor`` clears the other, so switching transports in one call never
|
|
trips the mutual-exclusion invariant. An explicit ``monitor`` wins over
|
|
the legacy aliases.
|
|
"""
|
|
if monitor is None:
|
|
return monitor_script, monitor_url
|
|
value = monitor.strip()
|
|
if not value:
|
|
return "", "" # clear both sources
|
|
if value.lower().startswith(("http://", "https://")):
|
|
return "", value
|
|
return value, ""
|
|
|
|
|
|
def _repeat_display(job: Dict[str, Any]) -> str:
|
|
times = (job.get("repeat") or {}).get("times")
|
|
completed = (job.get("repeat") or {}).get("completed", 0)
|
|
if times is None:
|
|
return "forever"
|
|
if times == 1:
|
|
return "once" if completed == 0 else "1/1"
|
|
return f"{completed}/{times}" if completed else f"{times} times"
|
|
|
|
|
|
def _canonical_skills(skill: Optional[str] = None, skills: Optional[Any] = None) -> List[str]:
|
|
if skills is None:
|
|
raw_items = [skill] if skill else []
|
|
elif isinstance(skills, str):
|
|
raw_items = [skills]
|
|
else:
|
|
raw_items = list(skills)
|
|
|
|
normalized: List[str] = []
|
|
for item in raw_items:
|
|
text = str(item or "").strip()
|
|
if text and text not in normalized:
|
|
normalized.append(text)
|
|
return normalized
|
|
|
|
|
|
|
|
|
|
def _normalize_optional_job_value(value: Optional[Any], *, strip_trailing_slash: bool = False) -> Optional[str]:
|
|
if value is None:
|
|
return None
|
|
text = str(value).strip()
|
|
if strip_trailing_slash:
|
|
text = text.rstrip("/")
|
|
return text or None
|
|
|
|
|
|
def _normalize_deliver_param(value: Any) -> Optional[str]:
|
|
"""Normalize a user-supplied ``deliver`` value to the canonical string form.
|
|
|
|
The cron schema documents ``deliver`` as a string (``"local"``, ``"origin"``,
|
|
``"telegram"``, ``"telegram:chat_id[:thread_id]"``, or comma-separated combos).
|
|
Some callers — MCP clients passing arrays, scripts building the payload as a
|
|
list — supply ``["telegram"]``. ``create_job``/``update_job`` store it as-is,
|
|
and the scheduler's ``str(deliver).split(",")`` then serializes the list to
|
|
the literal ``"['telegram']"`` which is not a known platform. Flatten lists
|
|
/ tuples at the API boundary so storage is always a string. Returns ``None``
|
|
for ``None``/empty so callers can treat it as "not supplied".
|
|
"""
|
|
if value is None:
|
|
return None
|
|
if isinstance(value, (list, tuple)):
|
|
parts = [str(p).strip() for p in value if str(p).strip()]
|
|
return ",".join(parts) if parts else None
|
|
text = str(value).strip()
|
|
return text or None
|
|
|
|
|
|
def _validate_bot_chat_deliver(deliver: Optional[str]) -> Optional[str]:
|
|
"""Validate any ``bot-chat[:<profile>]`` deliver elements at create time.
|
|
|
|
Bot Chat delivery is machine-local: the named profile must exist on THIS
|
|
machine (the one whose scheduler will fire the job). Failing loudly here
|
|
beats a per-run ``last_delivery_error`` at 3am — especially for Desktop
|
|
clients whose merged multi-gateway rosters may show same-named profiles
|
|
from other machines. Returns an error string or None.
|
|
"""
|
|
if not deliver:
|
|
return None
|
|
try:
|
|
from cron.scheduler import parse_bot_chat_deliver_token
|
|
from hermes_cli.profiles import normalize_profile_name, profile_exists
|
|
except Exception:
|
|
return None # validation is best-effort; resolution re-checks at fire time
|
|
for part in str(deliver).split(","):
|
|
profile_arg = parse_bot_chat_deliver_token(part.strip())
|
|
if profile_arg is None or not profile_arg:
|
|
continue # not a bot-chat token, or bare token (own profile)
|
|
try:
|
|
canon = normalize_profile_name(profile_arg)
|
|
except Exception:
|
|
return f"invalid bot-chat profile name '{profile_arg}'"
|
|
if not profile_exists(canon):
|
|
return (
|
|
f"bot-chat delivery profile '{profile_arg}' not found on this "
|
|
"gateway's machine. Bot Chat delivery is machine-local — use a "
|
|
"profile that exists here (hermes profile list), or omit the "
|
|
"name (deliver='bot-chat') for the job's own profile."
|
|
)
|
|
return None
|
|
|
|
|
|
def _resolve_cron_context_deliver(deliver: Optional[str]) -> Optional[str]:
|
|
"""Resolve ``origin`` to a concrete target for cron-context creates.
|
|
|
|
A job created FROM a cron run must never store the literal ``origin``:
|
|
the creating session is ephemeral, so by fire time there is no origin to
|
|
resolve and the scheduler would fall back to guessing a home channel.
|
|
Resolve at create time instead, using the creating run's own concrete
|
|
delivery target — the ``HERMES_CRON_AUTO_DELIVER_*`` contextvars that
|
|
``run_job`` publishes per run (already per-job-safe under the parallel
|
|
pool). Rules:
|
|
|
|
* Not a cron-context session → returned unchanged (chat/CLI creates keep
|
|
today's fire-time ``origin`` semantics, byte-identical).
|
|
* ``origin`` element (or an omitted value, which the scheduler treats as
|
|
origin) → replaced with ``platform:chat_id[:thread_id]`` from the
|
|
creating run's target; ``local`` when the creating run has no concrete
|
|
target (e.g. its own deliver is ``local``).
|
|
* Every other element (``local``, ``all``, explicit ``platform:...``)
|
|
passes through verbatim, including inside comma lists.
|
|
"""
|
|
from gateway.session_context import get_session_env
|
|
from utils import is_truthy_value
|
|
|
|
if not is_truthy_value(get_session_env("HERMES_CRON_SESSION", "")):
|
|
return deliver
|
|
|
|
def _creator_target() -> str:
|
|
platform = get_session_env("HERMES_CRON_AUTO_DELIVER_PLATFORM", "").strip()
|
|
chat_id = get_session_env("HERMES_CRON_AUTO_DELIVER_CHAT_ID", "").strip()
|
|
if not platform or not chat_id:
|
|
return "local"
|
|
thread_id = get_session_env("HERMES_CRON_AUTO_DELIVER_THREAD_ID", "").strip()
|
|
if thread_id:
|
|
return f"{platform}:{chat_id}:{thread_id}"
|
|
return f"{platform}:{chat_id}"
|
|
|
|
if deliver is None:
|
|
return _creator_target()
|
|
parts = [p.strip() for p in str(deliver).split(",") if p.strip()]
|
|
resolved = [_creator_target() if p.lower() == "origin" else p for p in parts]
|
|
# De-dup while preserving order: 'origin,local' with a local-target
|
|
# creator would otherwise store 'local,local'.
|
|
seen: set = set()
|
|
unique = [p for p in resolved if not (p in seen or seen.add(p))]
|
|
return ",".join(unique) if unique else None
|
|
|
|
|
|
def _validate_cron_base_url(
|
|
provider: Optional[Any], base_url: Optional[Any]
|
|
) -> Optional[str]:
|
|
"""Reject pairing a named provider's stored credential with an off-host base_url.
|
|
|
|
The cron tool is model-callable, so a prompt-injected job could set a real
|
|
provider plus an attacker ``base_url``; on fire the scheduler resolves that
|
|
provider's stored API key and sends it to the URL, exfiltrating the
|
|
credential (CWE-200/CWE-522). Allow a ``base_url`` override only when it
|
|
cannot leak a stored secret: no override at all, a configured custom/byok
|
|
provider that carries its own endpoint+key, or an override whose host
|
|
matches the named provider's own endpoint.
|
|
|
|
Returns an error string if blocked, else None (valid).
|
|
"""
|
|
bu = _normalize_optional_job_value(base_url, strip_trailing_slash=True)
|
|
if not bu:
|
|
return None
|
|
prov = _normalize_optional_job_value(provider)
|
|
if not prov:
|
|
# A base_url with no explicit provider inherits the default/session
|
|
# provider's stored key — the same exfil primitive without naming a
|
|
# provider. Require an explicit (custom) provider for custom endpoints.
|
|
return (
|
|
"base_url override requires an explicit provider. Set provider to a "
|
|
"configured custom provider to use a custom endpoint."
|
|
)
|
|
try:
|
|
from hermes_cli.runtime_provider import (
|
|
has_named_custom_provider,
|
|
resolve_requested_provider,
|
|
_get_named_custom_provider,
|
|
)
|
|
from hermes_cli.auth import PROVIDER_REGISTRY
|
|
from utils import base_url_host_matches, base_url_hostname
|
|
except Exception:
|
|
# Can't resolve provider metadata -> fail closed.
|
|
return f"Unable to validate base_url override for provider {prov!r}; refused."
|
|
|
|
if prov.lower() == "custom":
|
|
# Bare/inline 'custom' (and aliases that resolve to it) is pure BYOK: the
|
|
# runtime derives the key from a pool keyed by THIS base_url or from
|
|
# host-gated env vars, never an arbitrary stored secret. Safe to allow.
|
|
return None
|
|
if has_named_custom_provider(prov):
|
|
# A NAMED custom provider carries a STORED key, and
|
|
# _resolve_named_custom_runtime prefers the override base_url while still
|
|
# sending that stored key — so an off-host override exfiltrates it.
|
|
# Require the override host to match the provider's CONFIGURED endpoint.
|
|
try:
|
|
cp = _get_named_custom_provider(prov)
|
|
except Exception:
|
|
cp = None
|
|
cfg_host = base_url_hostname((cp or {}).get("base_url", "")) if cp else ""
|
|
if cfg_host and base_url_host_matches(bu, cfg_host):
|
|
return None
|
|
return (
|
|
f"base_url {bu!r} is not allowed for provider {prov!r}. A named "
|
|
f"custom provider's stored credential may only be sent to its own "
|
|
f"configured endpoint ({cfg_host or 'unknown'})."
|
|
)
|
|
try:
|
|
resolved = resolve_requested_provider(prov)
|
|
except Exception:
|
|
resolved = prov
|
|
pconfig = PROVIDER_REGISTRY.get(resolved) if isinstance(resolved, str) else None
|
|
known_host = base_url_hostname(getattr(pconfig, "inference_base_url", "") if pconfig else "")
|
|
if known_host and base_url_host_matches(bu, known_host):
|
|
return None
|
|
# Fail closed: any non-custom provider we cannot host-match to its own
|
|
# endpoint is refused. This covers named providers with a stored credential
|
|
# AND aliases/unknown names we can't resolve to a known host (e.g. "openai",
|
|
# "google"), which would otherwise pair a stored key with the override URL.
|
|
return (
|
|
f"base_url {bu!r} is not allowed for provider {prov!r}. A named "
|
|
f"provider's stored credential may only be sent to its own endpoint; "
|
|
f'use a configured custom provider (provider="custom") for a custom base_url.'
|
|
)
|
|
|
|
|
|
def _validate_cron_script_path(script: Optional[str]) -> Optional[str]:
|
|
"""Validate a cron job script path at the API boundary.
|
|
|
|
Scripts must be relative paths that resolve within HERMES_HOME/scripts/.
|
|
Absolute paths and ~ expansion are rejected to prevent arbitrary script
|
|
execution via prompt injection.
|
|
|
|
Returns an error string if blocked, else None (valid).
|
|
"""
|
|
if not script or not script.strip():
|
|
return None # empty/None = clearing the field, always OK
|
|
|
|
from hermes_constants import get_hermes_home
|
|
|
|
raw = script.strip()
|
|
|
|
# Reject absolute paths and ~ expansion at the API boundary.
|
|
# Only relative paths within ~/.hermes/scripts/ are allowed.
|
|
if raw.startswith(("/", "~")) or (len(raw) >= 2 and raw[1] == ":"):
|
|
return (
|
|
f"Script path must be relative to ~/.hermes/scripts/. "
|
|
f"Got absolute or home-relative path: {raw!r}. "
|
|
f"Place scripts in ~/.hermes/scripts/ and use just the filename."
|
|
)
|
|
|
|
# Validate containment after resolution
|
|
from tools.path_security import validate_within_dir
|
|
|
|
scripts_dir = get_hermes_home() / "scripts"
|
|
scripts_dir.mkdir(parents=True, exist_ok=True)
|
|
containment_error = validate_within_dir(scripts_dir / raw, scripts_dir)
|
|
if containment_error:
|
|
return (
|
|
f"Script path escapes the scripts directory via traversal: {raw!r}"
|
|
)
|
|
|
|
return None
|
|
|
|
|
|
def _format_job(job: Dict[str, Any]) -> Dict[str, Any]:
|
|
prompt = str(job.get("prompt") or "")
|
|
skills = _canonical_skills(job.get("skill"), job.get("skills"))
|
|
job_id = str(job.get("id") or "unknown")
|
|
name = str(job.get("name") or prompt[:50] or (skills[0] if skills else "") or job_id or "cron job")
|
|
result = {
|
|
"job_id": job_id,
|
|
"name": name,
|
|
"skill": skills[0] if skills else None,
|
|
"skills": skills,
|
|
"prompt_preview": prompt[:100] + "..." if len(prompt) > 100 else prompt,
|
|
"model": job.get("model"),
|
|
"provider": job.get("provider"),
|
|
"base_url": job.get("base_url"),
|
|
"schedule": job.get("schedule_display") or "?",
|
|
"repeat": _repeat_display(job),
|
|
"deliver": job.get("deliver", "local"),
|
|
"next_run_at": job.get("next_run_at"),
|
|
"last_run_at": job.get("last_run_at"),
|
|
"last_status": job.get("last_status"),
|
|
"last_delivery_error": job.get("last_delivery_error"),
|
|
"last_delivery_unverified": job.get("last_delivery_unverified"),
|
|
"last_fire_error": job.get("last_fire_error"),
|
|
"enabled": job.get("enabled", True),
|
|
# Derive from enabled so half-paused records never render as paused.
|
|
"state": effective_job_state(job),
|
|
"paused_at": job.get("paused_at"),
|
|
"paused_reason": job.get("paused_reason"),
|
|
}
|
|
if job.get("script"):
|
|
result["script"] = job["script"]
|
|
if job.get("reasoning_effort"):
|
|
result["reasoning_effort"] = job["reasoning_effort"]
|
|
if job.get("monitor_script"):
|
|
result["monitor_script"] = job["monitor_script"]
|
|
if job.get("monitor_url"):
|
|
result["monitor_url"] = job["monitor_url"]
|
|
if job.get("monitor_state"):
|
|
result["monitor_state"] = job["monitor_state"]
|
|
if job.get("no_agent"):
|
|
result["no_agent"] = True
|
|
if job.get("enabled_toolsets"):
|
|
result["enabled_toolsets"] = job["enabled_toolsets"]
|
|
if job.get("workdir"):
|
|
result["workdir"] = job["workdir"]
|
|
stored_refs = job.get("context_from") or []
|
|
if isinstance(stored_refs, str):
|
|
stored_refs = [stored_refs]
|
|
if any(str(r).strip().lower() == "self" or r == job.get("id") for r in stored_refs):
|
|
result["continuity"] = True
|
|
external_refs = [
|
|
r for r in stored_refs
|
|
if str(r).strip().lower() != "self" and r != job.get("id")
|
|
]
|
|
if external_refs:
|
|
result["context_from"] = external_refs
|
|
if isinstance(job.get("attach_to_session"), bool):
|
|
result["attach_to_session"] = job["attach_to_session"]
|
|
return result
|
|
|
|
|
|
def _relay_fronted_delivery_platforms(job: Dict[str, Any]) -> set:
|
|
"""Delivery-platform names for this job that the relay connector fronts."""
|
|
try:
|
|
from gateway.relay import relay_fronted_platforms
|
|
except Exception:
|
|
return set()
|
|
fronted = relay_fronted_platforms()
|
|
if not fronted:
|
|
return set()
|
|
try:
|
|
from cron.scheduler import _resolve_delivery_targets
|
|
|
|
targets = _resolve_delivery_targets(job) or []
|
|
except Exception:
|
|
return set()
|
|
theirs = {t.get("platform") for t in targets if t.get("platform")}
|
|
return theirs & fronted
|
|
|
|
|
|
def _forward_relay_fronted_run(
|
|
job: Dict[str, Any], extra_prompt: Optional[str] = None
|
|
) -> Optional[str]:
|
|
"""Forward a manual run to the gateway when it targets a relay-fronted
|
|
platform and this process has no live relay adapter.
|
|
|
|
Relay-fronted delivery has no standalone sender: the connector owns the
|
|
credential and the gateway's live relay adapter is the only path. The
|
|
gateway api_server's ``POST /api/jobs/{id}/run`` marks the job due for its
|
|
own ticker, which fires it with the live adapter. ``extra_prompt``
|
|
(transient per-run context) rides in the request body so the forwarded
|
|
fire keeps it. Returns a JSON result string when forwarding engages
|
|
(dispatch or the accurate error), else None to fall through to the normal
|
|
in-process run.
|
|
"""
|
|
if not _relay_fronted_delivery_platforms(job):
|
|
return None
|
|
job_id = job["id"]
|
|
import os
|
|
|
|
port_raw = os.getenv("API_SERVER_PORT", "").strip()
|
|
try:
|
|
port = int(port_raw) if port_raw else 8642
|
|
except ValueError:
|
|
port = 8642
|
|
# Mirror the api_server's own bind resolution (adapter reads
|
|
# extra.host -> API_SERVER_HOST -> 127.0.0.1). A wildcard bind
|
|
# (0.0.0.0/::) listens on loopback too, so dial loopback for those.
|
|
host = ""
|
|
try:
|
|
from hermes_cli.config import cfg_get, load_config_readonly
|
|
|
|
host = str(
|
|
cfg_get(
|
|
load_config_readonly(), "platforms", "api_server", "extra", "host",
|
|
default="",
|
|
)
|
|
or ""
|
|
).strip()
|
|
except Exception:
|
|
host = ""
|
|
if not host:
|
|
host = os.getenv("API_SERVER_HOST", "").strip()
|
|
if not host or host in ("0.0.0.0", "::", "*"):
|
|
host = "127.0.0.1"
|
|
if ":" in host and not host.startswith("["):
|
|
host = f"[{host}]" # bare IPv6 literal
|
|
url = f"http://{host}:{port}/api/jobs/{job_id}/run"
|
|
|
|
from agent.secret_scope import get_secret
|
|
|
|
key = get_secret("API_SERVER_KEY", "") or ""
|
|
|
|
resp = None
|
|
try:
|
|
import httpx
|
|
|
|
resp = httpx.post(
|
|
url,
|
|
headers={"Authorization": f"Bearer {key}"},
|
|
json=({"prompt": extra_prompt} if extra_prompt else {}),
|
|
timeout=10.0,
|
|
)
|
|
except Exception:
|
|
resp = None
|
|
|
|
if resp is not None and resp.status_code < 300:
|
|
return json.dumps(
|
|
{
|
|
"success": True,
|
|
"forwarded_to_gateway": True,
|
|
"note": (
|
|
"This job targets a relay-fronted platform; it was dispatched "
|
|
"to the running gateway, whose live relay adapter owns that "
|
|
"delivery."
|
|
),
|
|
},
|
|
indent=2,
|
|
)
|
|
return json.dumps(
|
|
{
|
|
"success": False,
|
|
"error": (
|
|
"This job targets a relay-fronted platform, which has no "
|
|
"standalone sender. Start the gateway — its ticker will "
|
|
"deliver the job on schedule via the live relay adapter."
|
|
),
|
|
},
|
|
indent=2,
|
|
)
|
|
|
|
|
|
def _manual_run_delivery_note(deliver: str, refreshed: Dict[str, Any]) -> str:
|
|
"""Parenthetical delivery note for a manual run's completion summary.
|
|
|
|
Follows the refreshed job record (#83993): ``run_one_job`` writes
|
|
``last_delivery_error`` via ``mark_job_run`` when the post-run delivery
|
|
(telegram/discord/…) failed, and the summary must not claim success over
|
|
that record — the calling agent relays this line to the user. Local jobs
|
|
never deliver; an empty/missing error keeps the legacy wording
|
|
byte-for-byte.
|
|
"""
|
|
# Falsy deliver ("", stored JSON null) means no delivery target — the
|
|
# fire-time path normalizes it to "local" (no delivery, output persisted
|
|
# in last_output, no delivery error), so it must read as saved-locally,
|
|
# not as a delivered remote target. Whitespace-only values are NOT folded
|
|
# in here: they keep falling through to the error check, where the
|
|
# fire-time "no delivery target resolved" error gets surfaced.
|
|
if not deliver or deliver == "local":
|
|
return " (output saved locally only)"
|
|
err = str(refreshed.get("last_delivery_error") or "").strip()
|
|
if not err:
|
|
return " (output was delivered there by the job itself)"
|
|
return f" (⚠ delivery FAILED: {err[:200]})"
|
|
|
|
|
|
def _execute_job_now(
|
|
job: Dict[str, Any], extra_prompt: Optional[str] = None
|
|
) -> Dict[str, Any]:
|
|
"""Execute a cron job immediately, outside the scheduler tick.
|
|
|
|
Atomically claims the job first via ``claim_job_for_fire`` — the same
|
|
at-most-once CAS the scheduler/external-provider fire path uses — so a
|
|
concurrently-running gateway ticker cannot also fire it (the claim both
|
|
blocks a duplicate fire and advances ``next_run_at`` for recurring jobs).
|
|
If the claim is lost (another fire is in flight), this is a no-op.
|
|
|
|
The actual firing is delegated to ``run_one_job`` — the single shared
|
|
execute→save→deliver→mark body the ticker and external providers use — so
|
|
failure delivery, ``[SILENT]`` handling, and live-adapter delivery stay
|
|
identical across paths and can't drift.
|
|
|
|
Returns {"claimed": bool, "success": bool, "error": str|None}.
|
|
"""
|
|
job_id = job["id"]
|
|
claimed_job = None
|
|
try:
|
|
# At-most-once claim: bail without running if a tick/other fire owns it.
|
|
claimed_job = claim_job_for_fire(job_id, return_job=True)
|
|
if not isinstance(claimed_job, dict):
|
|
# claim_job_for_fire returns False for paused/disabled/missing
|
|
# jobs too — don't mislabel those as "already being fired"
|
|
# (#60703): that message sends the user chasing a phantom
|
|
# in-flight run when the job simply isn't runnable.
|
|
refreshed = get_job(job_id)
|
|
if refreshed is None:
|
|
reason = "Job no longer exists; nothing to run."
|
|
elif not is_job_runnable(refreshed):
|
|
reason = "Job is paused/disabled; resume it before running."
|
|
else:
|
|
reason = "Job is already being fired by the scheduler; not run again."
|
|
return {"claimed": False, "success": False, "error": reason}
|
|
except Exception as e:
|
|
logger.error("Failed to claim cron job %s for immediate run: %s", job_id, e)
|
|
try:
|
|
mark_job_run(job_id, False, str(e))
|
|
except Exception:
|
|
pass
|
|
return {"claimed": True, "success": False, "error": str(e)}
|
|
|
|
return _run_claimed_job(claimed_job, extra_prompt=extra_prompt)
|
|
|
|
|
|
def _run_claimed_job(
|
|
job: Dict[str, Any], extra_prompt: Optional[str] = None
|
|
) -> Dict[str, Any]:
|
|
"""Fire an already-claimed job through the shared ``run_one_job`` body.
|
|
|
|
Split out of ``_execute_job_now`` so the background dispatch path
|
|
(``_try_dispatch_background_run``) can take the claim synchronously — so
|
|
the tool response can report "paused"/"already firing" immediately — and
|
|
hand the actual run to a daemon worker.
|
|
|
|
Returns {"claimed": True, "success": bool, "error": str|None}.
|
|
"""
|
|
job_id = job["id"]
|
|
_registered = False
|
|
fire_owner = None
|
|
try:
|
|
from cron.scheduler import (
|
|
release_running_job,
|
|
run_one_job,
|
|
try_register_running_job,
|
|
)
|
|
|
|
# In-flight dedupe (idea from #53395 by @izumi0uu): the fire claim's
|
|
# TTL (300s) is routinely outlived by real jobs, so it alone cannot
|
|
# stop a manual run from double-firing a job the ticker (or another
|
|
# manual run) is still executing. Register in the scheduler's shared
|
|
# running set — the same guard _submit_with_guard uses — which also
|
|
# makes this run visible to the gateway shutdown drain
|
|
# (get_running_job_ids, #60432) and mark_running_jobs_interrupted.
|
|
if not try_register_running_job(job_id):
|
|
return {
|
|
"claimed": True,
|
|
"success": False,
|
|
"error": (
|
|
"Job is already running (a scheduler tick or another "
|
|
"manual run is executing it); not started again."
|
|
),
|
|
}
|
|
_registered = True
|
|
|
|
claim = job.get("fire_claim")
|
|
fire_owner = str(claim.get("by") or "") if isinstance(claim, dict) else None
|
|
|
|
# run_one_job records last_run_at/last_status via mark_job_run (which
|
|
# also clears the fire claim) and returns True iff it processed the job.
|
|
# ``job`` here is the exact claimed snapshot (owner-bearing), so the
|
|
# shared body fences every terminal write by that owner.
|
|
#
|
|
# A manual `run` executes the job synchronously on the caller's thread,
|
|
# and a cron job is itself a full agent run that routinely takes
|
|
# minutes. The calling turn emits no tool activity for that entire
|
|
# window, so the gateway inactivity watchdog concludes the agent is
|
|
# hung and kills the parent turn (#76502). Fire a heartbeat into the
|
|
# caller's activity tracker (the same signal tool progress uses) while
|
|
# the job runs, so the watchdog sees a working tool instead of a
|
|
# silent one — mirrors the delegate_task heartbeat pattern. Best-effort:
|
|
# if no activity callback is registered (direct Python callers, tests),
|
|
# behavior is unchanged.
|
|
try:
|
|
from tools.environments.base import get_activity_callback
|
|
|
|
# Capture on THIS thread: the callback is thread-local (installed
|
|
# by the tool executor as the calling agent's _touch_activity), so
|
|
# a freshly spawned thread cannot read it back.
|
|
activity_cb = get_activity_callback()
|
|
except Exception:
|
|
activity_cb = None
|
|
|
|
_heartbeat_stop = threading.Event()
|
|
_heartbeat_thread = None
|
|
|
|
if activity_cb is not None:
|
|
job_name = str(job.get("name") or job_id)
|
|
|
|
def _heartbeat_loop() -> None:
|
|
started = time.monotonic()
|
|
while not _heartbeat_stop.wait(_CRON_RUN_HEARTBEAT_INTERVAL):
|
|
elapsed = time.monotonic() - started
|
|
if elapsed > _CRON_RUN_HEARTBEAT_CEILING:
|
|
# Stop masking the gateway watchdog — a run this long
|
|
# with an unlimited child watchdog is likely wedged.
|
|
logger.warning(
|
|
"cronjob run heartbeat ceiling reached for job "
|
|
"'%s' (%.0fs) — stopping heartbeat; gateway "
|
|
"watchdog regains authority",
|
|
job_name, elapsed,
|
|
)
|
|
return
|
|
try:
|
|
activity_cb(
|
|
f"cronjob: running job '{job_name}' ({int(elapsed)}s elapsed)"
|
|
)
|
|
except Exception:
|
|
# Never break the job run; keep heartbeating — one
|
|
# transient callback error must not silently drop
|
|
# watchdog protection for the rest of a long job.
|
|
continue
|
|
|
|
_heartbeat_thread = threading.Thread(
|
|
target=_heartbeat_loop,
|
|
daemon=True,
|
|
name="cronjob-run-heartbeat",
|
|
)
|
|
_heartbeat_thread.start()
|
|
|
|
# Manual runs invoked from a gateway agent execute outside the scheduler
|
|
# ticker, but they still share the process with the live platform
|
|
# adapters. Pass the gateway-owned adapter map and event loop through
|
|
# to run_one_job so delivery is scheduled on the loop that owns clients
|
|
# such as Matrix/aiohttp. Calling those clients from run_one_job's
|
|
# standalone asyncio.run() loop raises errors like "Timeout context
|
|
# manager should be used inside a task" and can break encrypted Matrix
|
|
# delivery (#61495 — salvaged from #63586 by @Fly-onlyone).
|
|
gateway_module = sys.modules.get("gateway.run")
|
|
runner_ref = getattr(gateway_module, "_gateway_runner_ref", None)
|
|
runner = runner_ref() if callable(runner_ref) else None
|
|
adapters = getattr(runner, "adapters", None) if runner is not None else None
|
|
gateway_loop = getattr(runner, "_gateway_loop", None) if runner is not None else None
|
|
|
|
try:
|
|
try:
|
|
processed = run_one_job(
|
|
job, adapters=adapters, loop=gateway_loop,
|
|
extra_prompt=extra_prompt,
|
|
)
|
|
finally:
|
|
_heartbeat_stop.set()
|
|
if _heartbeat_thread is not None:
|
|
_heartbeat_thread.join(timeout=_CRON_RUN_HEARTBEAT_INTERVAL + 1)
|
|
finally:
|
|
_registered = False
|
|
release_running_job(job_id)
|
|
refreshed = get_job(job_id) or {}
|
|
execution = None
|
|
execution_id = job.get("execution_id")
|
|
if execution_id:
|
|
from cron.executions import get_execution
|
|
|
|
execution = get_execution(str(execution_id))
|
|
last_status = refreshed.get("last_status")
|
|
# "delivery_failed" (#83993): the agent run itself succeeded but the
|
|
# output never reached the user. That is NOT a success for the caller
|
|
# — the calling agent relays this result — so report it as failed
|
|
# and surface the delivery error, which lives in last_delivery_error
|
|
# (last_error is None for these runs, and a bare success=False with
|
|
# error=None reads as an unexplained failure).
|
|
ok = last_status == "ok"
|
|
run_error = refreshed.get("last_error")
|
|
if last_status == "delivery_failed" and not run_error:
|
|
run_error = refreshed.get("last_delivery_error")
|
|
if execution is not None and execution.get("status") != "completed":
|
|
ok = False
|
|
run_error = (
|
|
execution.get("error")
|
|
or f"execution ended in {execution.get('status') or 'unknown'} state"
|
|
)
|
|
return {
|
|
"claimed": True,
|
|
"success": bool(processed and ok),
|
|
"error": run_error,
|
|
}
|
|
|
|
except Exception as e:
|
|
logger.error("Failed to execute cron job %s immediately: %s", job_id, e)
|
|
if _registered:
|
|
# Registration succeeded but we raised before the run's own
|
|
# release ran (e.g. heartbeat setup) — don't leave the job
|
|
# permanently marked in-flight. Only release registrations WE
|
|
# took: a bare discard here could erase a ticker-owned entry.
|
|
try:
|
|
from cron.scheduler import release_running_job as _release
|
|
|
|
_release(job_id)
|
|
except Exception:
|
|
pass
|
|
try:
|
|
mark_job_run(
|
|
job_id,
|
|
False,
|
|
str(e),
|
|
expected_fire_owner=fire_owner,
|
|
)
|
|
except Exception:
|
|
pass
|
|
return {
|
|
"claimed": True,
|
|
"success": False,
|
|
"error": str(e),
|
|
}
|
|
|
|
|
|
def _latest_job_output_excerpt(job_id: str, max_chars: int = 2000) -> Optional[str]:
|
|
"""Best-effort excerpt of the job's most recent saved output file.
|
|
|
|
Included in the background-run completion block so the parent agent sees
|
|
what the job actually produced without having to dig through
|
|
``~/.hermes/cron/output/``. Never raises.
|
|
"""
|
|
try:
|
|
from cron.jobs import get_cron_output_dir
|
|
|
|
out_dir = get_cron_output_dir() / job_id
|
|
files = sorted(out_dir.glob("*.md"))
|
|
if not files:
|
|
return None
|
|
text = files[-1].read_text(encoding="utf-8", errors="replace").strip()
|
|
if not text:
|
|
return None
|
|
if len(text) > max_chars:
|
|
text = text[:max_chars] + f"\n… (truncated; full output: {files[-1]})"
|
|
return text
|
|
except Exception:
|
|
return None
|
|
|
|
|
|
def _try_dispatch_background_run(
|
|
job: Dict[str, Any], session_id: Optional[str] = None,
|
|
extra_prompt: Optional[str] = None,
|
|
) -> Optional[Dict[str, Any]]:
|
|
"""Claim ``job`` now, then fire it on the async-delegation daemon executor.
|
|
|
|
A manual ``cronjob(action='run')`` used to execute the job synchronously
|
|
on the calling agent's tool thread. A cron job is a full agent run that
|
|
routinely takes minutes-to-hours, so the parent turn sat inside ONE tool
|
|
call the whole time: uninterruptible (the interrupt flag is only checked
|
|
between loop iterations) and serial (a batch of runs executed one by one).
|
|
|
|
This dispatches the run like ``delegate_task``'s background mode: the tool
|
|
returns immediately with a handle, the run executes on the shared async
|
|
daemon executor, and a ``type="async_delegation"`` completion event
|
|
re-enters the conversation as a fresh turn when the job finishes — riding
|
|
the existing completion-queue rail (CLI drain + gateway watcher), which
|
|
keeps message-role alternation legal and the prompt cache intact.
|
|
|
|
The at-most-once claim is taken SYNCHRONOUSLY before dispatch so
|
|
unrunnable jobs (paused / missing / already firing) report in the tool
|
|
response immediately instead of as a delayed completion event.
|
|
|
|
Returns
|
|
-------
|
|
None
|
|
Background delivery unavailable on this session runtime (one-shot
|
|
``hermes -z``, stateless HTTP, Kanban worker, nested cron run).
|
|
Caller falls back to the synchronous path unchanged.
|
|
dict
|
|
``{"claimed": False, "success": False, "error": ...}`` — claim lost;
|
|
same shape as ``_execute_job_now`` so the caller's existing response
|
|
formatting applies.
|
|
``{"claimed": True, "dispatched": True, "delegation_id": ...}`` —
|
|
run is executing in the background.
|
|
``{"claimed": True, "dispatched": False, "success": ..., "error": ...}``
|
|
— dispatch pool was at capacity; the run executed inline (the claim
|
|
was already taken and must not be stranded).
|
|
"""
|
|
# Finite sessions cannot route a detached result back after the turn
|
|
# ends — mirror delegate_task's gate and fall back to sync execution.
|
|
try:
|
|
from gateway.session_context import async_delivery_supported
|
|
|
|
if not async_delivery_supported():
|
|
return None
|
|
except Exception:
|
|
pass
|
|
|
|
job_id = job["id"]
|
|
job_name = str(job.get("name") or job_id)
|
|
|
|
# Reap any execution row this job (or any job) left stranded 'claimed'/
|
|
# 'running' by a dead owner process -- e.g. a PRIOR one-shot `hermes
|
|
# cron run` invocation whose dispatched runner died with the exiting
|
|
# process before writing a terminal status (issue #86721). The
|
|
# long-lived scheduler ticker already does this once at its own
|
|
# startup (cron/scheduler.py's self.recover_interrupted()); a one-shot
|
|
# CLI invocation has no equivalent "startup" moment of its own, so it
|
|
# never got this self-heal -- leaving a permanently-stale claim that
|
|
# blocked every subsequent manual run on the same job. Safe and cheap:
|
|
# only provably-dead owners (PID gone, or PID reused by a different
|
|
# process per its start time) are reaped; a genuinely live owner's row
|
|
# is left untouched.
|
|
try:
|
|
from cron.executions import recover_interrupted_executions
|
|
|
|
_reclaimed = recover_interrupted_executions()
|
|
if _reclaimed:
|
|
logger.warning(
|
|
"Reclaimed %d stale cron execution(s) from dead owner(s) "
|
|
"before dispatching job '%s'",
|
|
_reclaimed,
|
|
job_name,
|
|
)
|
|
except Exception as _reap_exc:
|
|
# Best-effort self-heal; a failure here must not block dispatch —
|
|
# but stay diagnosable (mirrors the scheduler tick's reap handling).
|
|
logger.debug("Stale execution reclaim failed: %s", _reap_exc)
|
|
|
|
# ---- routing capture (on THIS thread; contextvars don't cross the pool) ----
|
|
# Resolved BEFORE the claim: with no routable session there is no durable
|
|
# consumer for a detached completion, so we must not claim-and-dispatch.
|
|
try:
|
|
from tools.approval import get_current_session_key
|
|
|
|
session_key = get_current_session_key(default="")
|
|
except Exception:
|
|
session_key = ""
|
|
if not session_key and session_id:
|
|
# CLI path: the approval contextvar is only bound during gateway/TUI
|
|
# turns. The CLI drain filters completions by the durable agent
|
|
# session id (#64240), so stamp it as the key — an empty key would
|
|
# fail closed and the completion could never be claimed.
|
|
session_key = str(session_id)
|
|
if not session_key:
|
|
# Direct Python callers (`hermes cron run`, tests) have no agent
|
|
# session to deliver a completion to — the process exits right after
|
|
# the tool returns. Run synchronously.
|
|
return None
|
|
|
|
# ---- synchronous claim (same semantics as _execute_job_now) ----
|
|
try:
|
|
# Best-effort early dedupe so a mid-run job reports in THIS tool
|
|
# response instead of as a delayed error completion event. The
|
|
# authoritative (atomic) check is try_register_running_job inside
|
|
# _run_claimed_job on the worker.
|
|
try:
|
|
from cron.scheduler import get_running_job_ids
|
|
|
|
if job_id in get_running_job_ids():
|
|
return {
|
|
"claimed": False,
|
|
"success": False,
|
|
"error": (
|
|
"Job is already running (a scheduler tick or another "
|
|
"manual run is executing it); not started again."
|
|
),
|
|
}
|
|
except Exception:
|
|
pass
|
|
|
|
# Same snapshot claim as _execute_job_now: carry the owner-bearing
|
|
# record into the run so terminal writes stay fenced by this owner.
|
|
claimed_job = claim_job_for_fire(job_id, return_job=True)
|
|
if not isinstance(claimed_job, dict):
|
|
refreshed = get_job(job_id)
|
|
if refreshed is None:
|
|
reason = "Job no longer exists; nothing to run."
|
|
elif not is_job_runnable(refreshed):
|
|
reason = "Job is paused/disabled; resume it before running."
|
|
else:
|
|
reason = "Job is already being fired by the scheduler; not run again."
|
|
return {"claimed": False, "success": False, "error": reason}
|
|
except Exception as e:
|
|
logger.error("Failed to claim cron job %s for background run: %s", job_id, e)
|
|
try:
|
|
mark_job_run(job_id, False, str(e))
|
|
except Exception:
|
|
pass
|
|
return {"claimed": True, "dispatched": False, "success": False, "error": str(e)}
|
|
|
|
origin_ui_session_id = ""
|
|
try:
|
|
from gateway.session_context import get_session_env
|
|
|
|
origin_ui_session_id = get_session_env("HERMES_UI_SESSION_ID", "") or ""
|
|
except Exception:
|
|
pass
|
|
|
|
try:
|
|
from tools.async_delegation import (
|
|
_current_origin_session_id,
|
|
dispatch_async_delegation,
|
|
)
|
|
|
|
origin_session_id = _current_origin_session_id()
|
|
except Exception as e:
|
|
logger.warning(
|
|
"cronjob run: async delegation registry unavailable (%s); "
|
|
"running job '%s' inline.", e, job_name,
|
|
)
|
|
result = _run_claimed_job(claimed_job, extra_prompt=extra_prompt)
|
|
result["dispatched"] = False
|
|
return result
|
|
|
|
try:
|
|
from tools.delegate_tool import _get_max_async_children
|
|
|
|
max_async = _get_max_async_children()
|
|
except Exception:
|
|
max_async = 3
|
|
|
|
started_at = time.time()
|
|
# Canonicalize with the scheduler's own normalizer so the summary states
|
|
# the same target fire time will use: falsy ("", stored JSON null) reads
|
|
# "local", legacy list-form deliver flattens to its comma string. Read
|
|
# from the claimed snapshot — the owner-bearing record the run actually
|
|
# executes — not the pre-claim `job` the tool loaded.
|
|
from cron.scheduler import _normalize_deliver_value
|
|
|
|
deliver = _normalize_deliver_value(claimed_job.get("deliver", "local"))
|
|
|
|
def _runner() -> Dict[str, Any]:
|
|
res = _run_claimed_job(claimed_job, extra_prompt=extra_prompt)
|
|
duration = round(time.time() - started_at, 2)
|
|
refreshed = get_job(job_id) or {}
|
|
lines = [
|
|
f"Cron job '{job_name}' ({job_id}) finished its manual run.",
|
|
f"Result: {'ok' if res.get('success') else 'FAILED'}"
|
|
+ (f" — {res.get('error')}" if res.get("error") else ""),
|
|
f"Delivery target: {deliver}"
|
|
+ _manual_run_delivery_note(deliver, refreshed),
|
|
]
|
|
if refreshed.get("next_run_at"):
|
|
lines.append(f"Next scheduled run: {refreshed['next_run_at']}")
|
|
excerpt = _latest_job_output_excerpt(job_id)
|
|
if excerpt:
|
|
lines.append("--- JOB OUTPUT ---")
|
|
lines.append(excerpt)
|
|
return {
|
|
"status": "completed" if res.get("success") else "error",
|
|
"summary": "\n".join(lines),
|
|
"error": res.get("error"),
|
|
"api_calls": 0,
|
|
"duration_seconds": duration,
|
|
}
|
|
|
|
dispatch = dispatch_async_delegation(
|
|
goal=f"Manual run of cron job '{job_name}' ({job_id})",
|
|
context=(
|
|
"Triggered via cronjob(action='run'). The job executed in its own "
|
|
"fresh cron session; this block reports its outcome."
|
|
),
|
|
toolsets=None,
|
|
role="cron_run",
|
|
model=job.get("model"),
|
|
session_key=session_key,
|
|
parent_session_id=str(session_id) if session_id else None,
|
|
runner=_runner,
|
|
origin_ui_session_id=origin_ui_session_id,
|
|
origin_session_id=origin_session_id,
|
|
max_async_children=max_async,
|
|
)
|
|
|
|
if dispatch.get("status") == "dispatched":
|
|
return {
|
|
"claimed": True,
|
|
"dispatched": True,
|
|
"delegation_id": dispatch.get("delegation_id"),
|
|
}
|
|
|
|
# Pool at capacity (or submit failure): the claim is already taken and
|
|
# must not be stranded — run inline exactly as the legacy path did.
|
|
logger.info(
|
|
"cronjob run: background pool unavailable (%s); running job '%s' inline.",
|
|
dispatch.get("error", "rejected"), job_name,
|
|
)
|
|
result = _run_claimed_job(job, extra_prompt=extra_prompt)
|
|
result["dispatched"] = False
|
|
return result
|
|
|
|
|
|
def _apply_continuity(
|
|
context_from: Optional[Union[str, List[str]]],
|
|
continuity: bool,
|
|
) -> Optional[List[str]]:
|
|
"""Translate the ``continuity`` flag into the ``context_from`` list.
|
|
|
|
``continuity=True`` ensures ``"self"`` is present (the job's own previous
|
|
output is injected each run); ``continuity=False`` removes any
|
|
``"self"``/own-id entry. Other entries are preserved untouched.
|
|
"""
|
|
if isinstance(context_from, str):
|
|
refs = [context_from.strip()] if context_from.strip() else []
|
|
elif context_from:
|
|
refs = [str(j).strip() for j in context_from if str(j).strip()]
|
|
else:
|
|
refs = []
|
|
has_self = any(r.lower() == "self" for r in refs)
|
|
if continuity and not has_self:
|
|
refs.append("self")
|
|
elif not continuity and has_self:
|
|
refs = [r for r in refs if r.lower() != "self"]
|
|
return refs or None
|
|
|
|
|
|
def _gateway_liveness_notice(plural: bool = False) -> dict:
|
|
"""Build the ``gateway_running``/``warning`` payload for tool results.
|
|
|
|
Thin adapter over the shared CLI helper ``hermes_cli.cron._builtin_gateway_liveness``
|
|
(#87033) so the CLI and this tool can never disagree about what "scheduler
|
|
active" means. Returns ``{"gateway_running": False, "warning": ...}`` when
|
|
the builtin ticker has no gateway process to run it, ``{"gateway_running":
|
|
None}`` when the probe failed, and ``{"gateway_running": True}`` when the
|
|
scheduler is active. ``plural`` rewords the warning for multi-job results
|
|
(the ``list`` action).
|
|
"""
|
|
try:
|
|
from hermes_cli.cron import _builtin_gateway_liveness
|
|
|
|
_gw = _builtin_gateway_liveness()
|
|
except Exception:
|
|
return {"gateway_running": None}
|
|
subject = "these jobs are saved" if plural else "this job is saved"
|
|
if _gw is False:
|
|
return {
|
|
"gateway_running": False,
|
|
"warning": (
|
|
f"The Hermes gateway is not running — {subject} "
|
|
"but will NOT fire until the gateway is started "
|
|
"(hermes gateway install / hermes gateway start). "
|
|
"Tell the user the task is scheduled but not active yet."
|
|
),
|
|
}
|
|
if _gw is None:
|
|
return {"gateway_running": None}
|
|
return {"gateway_running": True}
|
|
|
|
|
|
def cronjob(
|
|
action: str,
|
|
job_id: Optional[str] = None,
|
|
prompt: Optional[str] = None,
|
|
schedule: Optional[str] = None,
|
|
name: Optional[str] = None,
|
|
repeat: Optional[int] = None,
|
|
deliver: Optional[str] = None,
|
|
include_disabled: bool = False,
|
|
skill: Optional[str] = None,
|
|
skills: Optional[List[str]] = None,
|
|
model: Optional[str] = None,
|
|
provider: Optional[str] = None,
|
|
base_url: Optional[str] = None,
|
|
reason: Optional[str] = None,
|
|
script: Optional[str] = None,
|
|
context_from: Optional[Union[str, List[str]]] = None,
|
|
continuity: Optional[bool] = None,
|
|
enabled_toolsets: Optional[List[str]] = None,
|
|
workdir: Optional[str] = None,
|
|
no_agent: Optional[bool] = None,
|
|
attach_to_session: Optional[bool] = None,
|
|
monitor_script: Optional[str] = None,
|
|
monitor_url: Optional[str] = None,
|
|
reasoning_effort: Optional[str] = None,
|
|
failure_deliver: Optional[Union[str, List[str]]] = None,
|
|
task_id: str = None,
|
|
session_id: Optional[str] = None,
|
|
) -> str:
|
|
"""Unified cron job management tool."""
|
|
del task_id # unused but kept for handler signature compatibility
|
|
|
|
try:
|
|
normalized = (action or "").strip().lower()
|
|
|
|
if normalized == "create":
|
|
if not schedule:
|
|
return tool_error("schedule is required for create", success=False)
|
|
canonical_skills = _canonical_skills(skill, skills)
|
|
_no_agent = bool(no_agent)
|
|
# Job-shape validation differs by mode:
|
|
# - no_agent=True → script is the job; prompt/skills are optional
|
|
# (and irrelevant to execution).
|
|
# - no_agent=False (default) → at least one of prompt/skills must
|
|
# be set, same as before.
|
|
if _no_agent:
|
|
if not script:
|
|
return tool_error(
|
|
"create with no_agent=True requires a script — "
|
|
"the script is the job. In no_agent mode the LLM is "
|
|
"skipped entirely: prompt and skills are ignored, "
|
|
"non-empty stdout is delivered verbatim, empty stdout "
|
|
"sends nothing (watchdog pattern), and a non-zero "
|
|
"exit or timeout sends an error alert.",
|
|
success=False,
|
|
)
|
|
elif not prompt and not canonical_skills:
|
|
return tool_error("create requires either prompt or at least one skill", success=False)
|
|
if prompt:
|
|
scan_error = _scan_cron_prompt(prompt)
|
|
if scan_error:
|
|
return tool_error(scan_error, success=False)
|
|
|
|
# Validate script path before storing
|
|
if script:
|
|
script_error = _validate_cron_script_path(script)
|
|
if script_error:
|
|
return tool_error(script_error, success=False)
|
|
|
|
# Validate monitor source (same containment rules as script).
|
|
if monitor_script:
|
|
monitor_error = _validate_cron_script_path(monitor_script)
|
|
if monitor_error:
|
|
return tool_error(monitor_error, success=False)
|
|
|
|
# Reject a model-supplied base_url that would route a named
|
|
# provider's stored credential to an attacker endpoint (F8).
|
|
base_url_error = _validate_cron_base_url(provider, base_url)
|
|
if base_url_error:
|
|
return tool_error(base_url_error, success=False)
|
|
|
|
# bot-chat deliver targets are machine-local: named profiles must
|
|
# exist here, and a bad name should fail the CREATE, not the run.
|
|
bot_chat_error = _validate_bot_chat_deliver(_normalize_deliver_param(deliver))
|
|
if bot_chat_error:
|
|
return tool_error(bot_chat_error, success=False)
|
|
# failure_deliver shares deliver's grammar and validators (NS-788).
|
|
bot_chat_error = _validate_bot_chat_deliver(
|
|
_normalize_deliver_param(failure_deliver)
|
|
)
|
|
if bot_chat_error:
|
|
return tool_error(bot_chat_error, success=False)
|
|
|
|
# Validate context_from references existing jobs
|
|
if context_from:
|
|
from cron.jobs import get_job as _get_job
|
|
refs = [context_from] if isinstance(context_from, str) else context_from
|
|
for ref_id in refs:
|
|
# "self" is resolved to the job's own id at run time —
|
|
# it can't be validated against the store (the job does
|
|
# not exist yet at create time).
|
|
if isinstance(ref_id, str) and ref_id.strip().lower() == "self":
|
|
continue
|
|
if not _get_job(ref_id):
|
|
return tool_error(
|
|
f"context_from job '{ref_id}' not found. "
|
|
"Use cronjob(action='list') to see available jobs.",
|
|
success=False,
|
|
)
|
|
|
|
# continuity=True is sugar for context_from including "self":
|
|
# the job wakes up with its own previous run's output injected.
|
|
if continuity is not None:
|
|
context_from = _apply_continuity(context_from, continuity)
|
|
|
|
from cron.scheduler import (
|
|
CronSchedulerRegistrationError,
|
|
create_job_with_scheduler_registration,
|
|
)
|
|
|
|
try:
|
|
job = create_job_with_scheduler_registration(
|
|
prompt=prompt or "",
|
|
schedule=schedule,
|
|
name=name,
|
|
repeat=repeat,
|
|
deliver=_resolve_cron_context_deliver(
|
|
_normalize_deliver_param(deliver)
|
|
),
|
|
origin=_origin_from_env(),
|
|
skills=canonical_skills,
|
|
model=_normalize_optional_job_value(model),
|
|
provider=_normalize_optional_job_value(provider),
|
|
base_url=_normalize_optional_job_value(base_url, strip_trailing_slash=True),
|
|
script=_normalize_optional_job_value(script),
|
|
context_from=context_from,
|
|
enabled_toolsets=enabled_toolsets or None,
|
|
workdir=_normalize_optional_job_value(workdir),
|
|
no_agent=_no_agent,
|
|
attach_to_session=attach_to_session,
|
|
monitor_script=_normalize_optional_job_value(monitor_script),
|
|
monitor_url=_normalize_optional_job_value(monitor_url),
|
|
# reasoning_effort reaches here from the CLI
|
|
# (hermes cron create --reasoning-effort) ONLY — it is
|
|
# deliberately absent from CRONJOB_SCHEMA and the model
|
|
# dispatch below: models do not make model-config
|
|
# decisions (standing policy).
|
|
reasoning_effort=reasoning_effort,
|
|
failure_deliver=_resolve_cron_context_deliver(
|
|
_normalize_deliver_param(failure_deliver)
|
|
),
|
|
)
|
|
except CronSchedulerRegistrationError as exc:
|
|
_partial = exc.to_dict()
|
|
return tool_error(_partial.pop("error"), success=False, **_partial)
|
|
_create_message = f"Cron job '{job['name']}' created."
|
|
_local_notice = _local_delivery_notice(job, _normalize_deliver_param(deliver))
|
|
if _local_notice:
|
|
_create_message = f"{_create_message} {_local_notice}"
|
|
# Gateway liveness surfacing (#87033): the builtin scheduler's
|
|
# ticker lives in the gateway process, so a job created with no
|
|
# gateway running is stored but will never fire. Tell the model
|
|
# here — the CLI already warns, but the agent path saw only a
|
|
# clean success and confidently told the user it was scheduled.
|
|
_result = {
|
|
"success": True,
|
|
"job_id": job["id"],
|
|
"name": job["name"],
|
|
"skill": job.get("skill"),
|
|
"skills": job.get("skills", []),
|
|
"schedule": job["schedule_display"],
|
|
"repeat": _repeat_display(job),
|
|
"deliver": job.get("deliver", "local"),
|
|
"next_run_at": job["next_run_at"],
|
|
"job": _format_job(job),
|
|
"message": _create_message,
|
|
**_gateway_liveness_notice(),
|
|
}
|
|
# Mode-specific guidance rides in the create response (once, when
|
|
# relevant) instead of in the schema (every API call). See
|
|
# _mode_guidance_notes.
|
|
_notes = _mode_guidance_notes(job, _normalize_deliver_param(deliver))
|
|
if _notes:
|
|
_result["guidance"] = _notes
|
|
return json.dumps(_result, indent=2)
|
|
|
|
if normalized == "list":
|
|
jobs = [_format_job(job) for job in list_jobs(include_disabled=include_disabled)]
|
|
_result = {"success": True, "count": len(jobs), "jobs": jobs}
|
|
# Same silent-inert-job class as create (#87033): an agent
|
|
# inspecting existing jobs in a gateway-less environment must
|
|
# learn they are not firing, not just see a clean list. An empty
|
|
# list has nothing inert — stay quiet (and skip the probe).
|
|
if jobs:
|
|
_result.update(_gateway_liveness_notice(plural=True))
|
|
return json.dumps(_result, indent=2)
|
|
|
|
if not job_id:
|
|
return tool_error(f"job_id is required for action '{normalized}'", success=False)
|
|
|
|
try:
|
|
job = resolve_job_ref(job_id)
|
|
except AmbiguousJobReference as exc:
|
|
return json.dumps(
|
|
{
|
|
"success": False,
|
|
"error": str(exc),
|
|
"matches": [
|
|
{
|
|
"id": m["id"],
|
|
"name": m.get("name"),
|
|
"schedule": m.get("schedule_display"),
|
|
"next_run_at": m.get("next_run_at"),
|
|
}
|
|
for m in exc.matches
|
|
],
|
|
},
|
|
indent=2,
|
|
)
|
|
if not job:
|
|
return json.dumps(
|
|
{"success": False, "error": f"Job with ID or name '{job_id}' not found. Use cronjob(action='list') to inspect jobs."},
|
|
indent=2,
|
|
)
|
|
# Resolve to canonical ID (supports name-based lookup)
|
|
job_id = job["id"]
|
|
|
|
if normalized == "remove":
|
|
removed = remove_job(job_id)
|
|
if not removed:
|
|
return tool_error(f"Failed to remove job '{job_id}'", success=False)
|
|
_notify_provider_jobs_changed_safe()
|
|
return json.dumps(
|
|
{
|
|
"success": True,
|
|
"message": f"Cron job '{job['name']}' removed.",
|
|
"removed_job": {
|
|
"id": job_id,
|
|
"name": job["name"],
|
|
"schedule": job.get("schedule_display"),
|
|
},
|
|
},
|
|
indent=2,
|
|
)
|
|
|
|
if normalized == "pause":
|
|
updated = pause_job(job_id, reason=reason)
|
|
_notify_provider_jobs_changed_safe()
|
|
return json.dumps({"success": True, "job": _format_job(updated)}, indent=2)
|
|
|
|
if normalized == "resume":
|
|
updated = resume_job(job_id)
|
|
_notify_provider_jobs_changed_safe()
|
|
return json.dumps({"success": True, "job": _format_job(updated)}, indent=2)
|
|
|
|
if normalized in {"run", "run_now", "trigger"}:
|
|
# Per-run context (#57331, salvaged from #57342/@liuhao1024 and
|
|
# #57360/@ghedeselmabot): `prompt` on the run action is transient
|
|
# context appended to the stored prompt for THIS fire only, never
|
|
# persisted. It goes through the same strict injection scan as
|
|
# stored prompts before firing.
|
|
extra_prompt = prompt or None
|
|
if extra_prompt:
|
|
scan_error = _scan_cron_prompt(extra_prompt)
|
|
if scan_error:
|
|
return tool_error(scan_error, success=False)
|
|
# Execute the job immediately rather than only scheduling it for the
|
|
# next scheduler tick — a manual `run` should actually run, even when
|
|
# no gateway/ticker is active (the #41037 case). The claim (taken
|
|
# inside both paths below) advances next_run_at and blocks a
|
|
# concurrent tick from double-firing.
|
|
#
|
|
# Preferred path: dispatch the run to the background like
|
|
# delegate_task — the tool returns a handle immediately and the
|
|
# job's outcome re-enters the conversation as a completion event.
|
|
# A cron job is a full agent run (minutes to hours); executing it
|
|
# inline made the parent turn uninterruptible and serialized
|
|
# batches of manual runs (#80xxx — the "stuck Telegram session"
|
|
# incident). Falls back to inline execution when the session
|
|
# runtime can't receive detached completions.
|
|
bg = _try_dispatch_background_run(
|
|
job, session_id=session_id, extra_prompt=extra_prompt
|
|
)
|
|
if bg is not None and bg.get("dispatched"):
|
|
_notify_provider_jobs_changed_safe()
|
|
result = _format_job(get_job(job_id) or {"id": job_id})
|
|
result["executed"] = True
|
|
result["execution_mode"] = "background"
|
|
result["delegation_id"] = bg.get("delegation_id")
|
|
return json.dumps(
|
|
{
|
|
"success": True,
|
|
"job": result,
|
|
"note": (
|
|
"The job is running in the background. You and the "
|
|
"user can keep working; its outcome re-enters the "
|
|
"conversation as a new message when it finishes. "
|
|
"Do not wait or poll — just continue."
|
|
),
|
|
},
|
|
indent=2,
|
|
)
|
|
# bg carries a terminal result (claim lost, or inline fallback
|
|
# after pool rejection); None means background delivery is
|
|
# unsupported here — run synchronously as before.
|
|
if bg is not None:
|
|
exec_result = bg
|
|
else:
|
|
# Relay-fronted manual run: a standalone process has no live
|
|
# relay adapter and no standalone sender, so forward to the
|
|
# running gateway (its live adapter owns that delivery).
|
|
forwarded = _forward_relay_fronted_run(job, extra_prompt=extra_prompt)
|
|
if forwarded is not None:
|
|
return forwarded
|
|
exec_result = _execute_job_now(job, extra_prompt=extra_prompt)
|
|
# A claimed direct run advances next_run_at and may race the
|
|
# external one-shot for the same occurrence. If Chronos loses that
|
|
# claim, its consumed fire cannot re-arm itself; reconcile from the
|
|
# winning direct path after the run has persisted its final state.
|
|
if exec_result.get("claimed", False):
|
|
_notify_provider_jobs_changed_safe()
|
|
# Re-read so the response reflects the post-run last_run_at/last_status.
|
|
result = _format_job(get_job(job_id) or {"id": job_id})
|
|
result["executed"] = exec_result.get("claimed", False)
|
|
result["execution_success"] = exec_result.get("success", False)
|
|
if not exec_result.get("claimed", False):
|
|
result["execution_skipped"] = exec_result.get("error") or (
|
|
"Already being fired by the scheduler; not run again."
|
|
)
|
|
elif exec_result.get("error"):
|
|
result["execution_error"] = exec_result["error"]
|
|
return json.dumps({"success": True, "job": result}, indent=2)
|
|
|
|
if normalized == "update":
|
|
updates: Dict[str, Any] = {}
|
|
if prompt is not None:
|
|
scan_error = _scan_cron_prompt(prompt)
|
|
if scan_error:
|
|
return tool_error(scan_error, success=False)
|
|
updates["prompt"] = prompt
|
|
if name is not None and name.strip():
|
|
# Blank name is a no-op, not a clear. The `is not None` sentinel
|
|
# treats every supplied field as an explicit edit, and a model
|
|
# that re-sends the whole schema with type-default empties ("", [], 0)
|
|
# then wipes fields it never meant to touch.
|
|
updates["name"] = name
|
|
if deliver is not None:
|
|
bot_chat_error = _validate_bot_chat_deliver(_normalize_deliver_param(deliver))
|
|
if bot_chat_error:
|
|
return tool_error(bot_chat_error, success=False)
|
|
updates["deliver"] = _resolve_cron_context_deliver(
|
|
_normalize_deliver_param(deliver)
|
|
)
|
|
if failure_deliver is not None:
|
|
# '' clears the override (job falls back to deliver on
|
|
# failures); non-empty values share deliver's validation
|
|
# AND its cron-context origin resolution (a job created
|
|
# from inside a cron run must never store literal
|
|
# 'origin' — same rule as deliver).
|
|
_norm_fd = _normalize_deliver_param(failure_deliver)
|
|
if _norm_fd:
|
|
bot_chat_error = _validate_bot_chat_deliver(_norm_fd)
|
|
if bot_chat_error:
|
|
return tool_error(bot_chat_error, success=False)
|
|
_norm_fd = _resolve_cron_context_deliver(_norm_fd)
|
|
updates["failure_deliver"] = _norm_fd
|
|
if skills is not None or skill is not None:
|
|
canonical_skills = _canonical_skills(skill, skills)
|
|
updates["skills"] = canonical_skills
|
|
updates["skill"] = canonical_skills[0] if canonical_skills else None
|
|
if model is not None:
|
|
updates["model"] = _normalize_optional_job_value(model)
|
|
if provider is not None:
|
|
updates["provider"] = _normalize_optional_job_value(provider)
|
|
if base_url is not None:
|
|
updates["base_url"] = _normalize_optional_job_value(base_url, strip_trailing_slash=True)
|
|
if reasoning_effort is not None:
|
|
# CLI-only lane (see create above): update_job validates
|
|
# against the canonical grammar; empty string clears the pin.
|
|
updates["reasoning_effort"] = reasoning_effort
|
|
# Re-validate the EFFECTIVE provider/base_url on EVERY update, not
|
|
# only when this update supplies provider/base_url. A job persisted
|
|
# before this guard (or written directly to the jobs store) may
|
|
# already hold an unsafe named-provider + off-host base_url pair;
|
|
# if we only checked when the update touches those axes, editing any
|
|
# unrelated field (name, schedule, ...) would succeed and leave that
|
|
# exfil-capable pair active and schedulable (F8). The effective pair
|
|
# merges this update's normalized values over the stored job; an
|
|
# operator can still remediate in the same update by clearing
|
|
# base_url or pointing provider/base_url at a safe pair.
|
|
eff_provider = (
|
|
updates["provider"] if "provider" in updates else job.get("provider")
|
|
)
|
|
eff_base_url = (
|
|
updates["base_url"] if "base_url" in updates else job.get("base_url")
|
|
)
|
|
base_url_error = _validate_cron_base_url(eff_provider, eff_base_url)
|
|
if base_url_error:
|
|
return tool_error(base_url_error, success=False)
|
|
if script is not None:
|
|
# Pass empty string to clear an existing script
|
|
if script:
|
|
script_error = _validate_cron_script_path(script)
|
|
if script_error:
|
|
return tool_error(script_error, success=False)
|
|
updates["script"] = _normalize_optional_job_value(script) if script else None
|
|
if monitor_script is not None:
|
|
# Pass empty string to clear an existing monitor_script
|
|
if monitor_script:
|
|
monitor_error = _validate_cron_script_path(monitor_script)
|
|
if monitor_error:
|
|
return tool_error(monitor_error, success=False)
|
|
updates["monitor_script"] = (
|
|
_normalize_optional_job_value(monitor_script) if monitor_script else None
|
|
)
|
|
if monitor_url is not None:
|
|
# Pass empty string to clear an existing monitor_url
|
|
updates["monitor_url"] = (
|
|
_normalize_optional_job_value(monitor_url) if monitor_url else None
|
|
)
|
|
if monitor_script is not None or monitor_url is not None:
|
|
eff_mon_script = (
|
|
updates["monitor_script"] if "monitor_script" in updates else job.get("monitor_script")
|
|
)
|
|
eff_mon_url = (
|
|
updates["monitor_url"] if "monitor_url" in updates else job.get("monitor_url")
|
|
)
|
|
if eff_mon_script and eff_mon_url:
|
|
return tool_error(
|
|
"monitor_script and monitor_url are mutually exclusive — "
|
|
"clear one before setting the other.",
|
|
success=False,
|
|
)
|
|
if context_from is not None or continuity is not None:
|
|
# Empty string / empty list clears the field; otherwise validate
|
|
# each referenced job exists before storing. Normalized to a list
|
|
# (or None) to match the shape stored by create_job().
|
|
if context_from is None:
|
|
# continuity-only update: start from the job's stored refs.
|
|
existing = job.get("context_from") or []
|
|
refs = [str(j).strip() for j in existing if str(j).strip()]
|
|
elif isinstance(context_from, str):
|
|
refs = [context_from.strip()] if context_from.strip() else []
|
|
else:
|
|
refs = [str(j).strip() for j in context_from if str(j).strip()]
|
|
if continuity is not None:
|
|
refs = _apply_continuity(refs, continuity) or []
|
|
if refs:
|
|
from cron.jobs import get_job as _get_job
|
|
for ref_id in refs:
|
|
# "self" resolves to the job's own id at run time.
|
|
if ref_id.lower() == "self":
|
|
continue
|
|
if not _get_job(ref_id):
|
|
return tool_error(
|
|
f"context_from job '{ref_id}' not found. "
|
|
"Use cronjob(action='list') to see available jobs.",
|
|
success=False,
|
|
)
|
|
updates["context_from"] = refs or None
|
|
if enabled_toolsets is not None:
|
|
updates["enabled_toolsets"] = enabled_toolsets or None
|
|
if attach_to_session is not None:
|
|
updates["attach_to_session"] = bool(attach_to_session)
|
|
if workdir is not None:
|
|
# Empty string clears the field (restores old behaviour);
|
|
# otherwise pass raw — update_job() validates / normalizes.
|
|
updates["workdir"] = _normalize_optional_job_value(workdir) or None
|
|
if no_agent is not None:
|
|
# Toggling no_agent on/off at update time. If flipping to True,
|
|
# we need a script to already exist on the job (or be part of
|
|
# the same update) — otherwise the next tick would error out.
|
|
target_no_agent = bool(no_agent)
|
|
if target_no_agent:
|
|
effective_script = updates.get("script") if "script" in updates else job.get("script")
|
|
if not effective_script:
|
|
return tool_error(
|
|
"Cannot set no_agent=True on a job without a script. "
|
|
"Set `script` in the same update, or on the job first.",
|
|
success=False,
|
|
)
|
|
updates["no_agent"] = target_no_agent
|
|
if repeat is not None:
|
|
# Coerce string forms ('forever'/'once'/'3') and 0/negative
|
|
# via the shared chokepoint — a bare `repeat <= 0` here
|
|
# raised TypeError for string repeats on the UPDATE path
|
|
# (create was fixed first; same class).
|
|
from cron.jobs import normalize_repeat_value
|
|
normalized_repeat = normalize_repeat_value(repeat)
|
|
repeat_state = dict(job.get("repeat") or {})
|
|
repeat_state["times"] = normalized_repeat
|
|
updates["repeat"] = repeat_state
|
|
if schedule is not None:
|
|
parsed_schedule = parse_schedule(schedule)
|
|
updates["schedule"] = parsed_schedule
|
|
updates["schedule_display"] = parsed_schedule.get("display", schedule)
|
|
if job.get("state") != "paused":
|
|
updates["state"] = "scheduled"
|
|
updates["enabled"] = True
|
|
if not updates:
|
|
return tool_error("No updates provided.", success=False)
|
|
updated = update_job(job_id, updates)
|
|
_notify_provider_jobs_changed_safe()
|
|
_upd_result: Dict[str, Any] = {"success": True, "job": _format_job(updated)}
|
|
# An update can switch a job into monitor / no_agent mode or
|
|
# change its delivery — echo the same mode guidance as create.
|
|
_upd_notes = _mode_guidance_notes(updated, _normalize_deliver_param(deliver))
|
|
if _upd_notes:
|
|
_upd_result["guidance"] = _upd_notes
|
|
return json.dumps(_upd_result, indent=2)
|
|
|
|
return tool_error(f"Unknown cron action '{action}'", success=False)
|
|
|
|
except Exception as e:
|
|
return tool_error(str(e), success=False)
|
|
|
|
|
|
|
|
CRONJOB_SCHEMA = {
|
|
"name": "cronjob_manage",
|
|
"description": """Manage scheduled cron jobs: action='create' schedules a job from a prompt and/or skills; 'list' inspects jobs; 'update'/'pause'/'resume'/'remove' manage one by job_id (always list first — never guess job IDs); 'run' fires a job immediately in the BACKGROUND (returns a handle at once, outcome re-enters the conversation when done — do not wait or poll; optional 'prompt' adds transient context for that fire only).
|
|
|
|
Jobs run in a fresh session with no current-chat context, so prompts must be self-contained, and the agent's FINAL RESPONSE is what gets delivered — cron runs are autonomous and cannot ask questions. Prefer updating an existing job over creating near-duplicates.""",
|
|
"parameters": {
|
|
"type": "object",
|
|
"properties": {
|
|
"action": {
|
|
"type": "string",
|
|
"description": "One of: create, list, update, pause, resume, remove, run. When action=create, the 'schedule' and 'prompt' fields are REQUIRED."
|
|
},
|
|
"job_id": {
|
|
"type": "string",
|
|
"description": "Required for update/pause/resume/remove/run"
|
|
},
|
|
"prompt": {
|
|
"type": "string",
|
|
"description": "For create: the full self-contained prompt (paired with any skills as the task instruction). For run: optional transient context for that single fire (never persisted)."
|
|
},
|
|
"schedule": {
|
|
"type": "string",
|
|
"type": "string",
|
|
"description": "REQUIRED for create. Schedule forms: (1) recurring interval — '30m', 'every 2h', 'every hour' (EVERY 30 minutes / 2 hours / hour, forever by default); (2) explicit one-shot by duration — 'in 30m', 'in 2h' (fires ONCE that far from now; use this for 'remind me in N minutes' — do NOT hand-compute an absolute timestamp); (3) natural day/time — 'every monday 9am', 'weekdays at 9am', 'every day at 9am' (recurring weekly/daily); (4) cron syntax — '0 9 * * *' (daily 9am); (5) absolute one-shot — ISO timestamp '2026-06-01T09:00:00'."
|
|
},
|
|
"name": {
|
|
"type": "string",
|
|
"description": "Optional human-friendly name"
|
|
},
|
|
"repeat": {
|
|
"type": "integer",
|
|
"description": "Optional repeat count. Omit for defaults (once for one-shot, forever for recurring)."
|
|
},
|
|
"deliver": {
|
|
"type": "string",
|
|
"description": "Where the job's output is POSTED as a one-way message (the job itself always runs in a fresh session with no chat context). Omit to address the chat/topic this job was created from. Otherwise: 'local' (save only, no delivery), 'all' (every connected home channel, resolved at fire time), 'bot-chat' or 'bot-chat:<profile>' (inject into a Bot Chat as a real message), or platform:chat_id:thread_id (e.g. 'telegram:-1001234567890:17585'). Comma-combine like 'origin,all'."
|
|
},
|
|
"failure_deliver": {
|
|
"type": "string",
|
|
"description": "Optional override target for FAILURE notices only (same grammar as deliver). When set, engine failure/interruption notices go here instead of the deliver target; 'local' suppresses them entirely (state still recorded in cron list/run history). Use for jobs delivering into shared channels where failure noise is unwanted. Omit = failures follow deliver (default). On update, '' clears."
|
|
},
|
|
"skills": {
|
|
"type": "array",
|
|
"items": {"type": "string"},
|
|
"description": "Optional ordered skill names loaded before the cron prompt. On update, [] clears."
|
|
},
|
|
"script": {
|
|
"type": "string",
|
|
"description": f"Optional script run each tick; stdout is injected into the agent's prompt as context (with no_agent=True the script IS the job). Relative paths resolve under {display_hermes_home()}/scripts/; .sh/.bash via bash, else Python. On update, '' clears."
|
|
},
|
|
"monitor": {
|
|
"type": "string",
|
|
"description": "Optional change-detector that gates the agent: an http(s) URL (fetched each tick) or a script path (same rules as `script`, run each tick) — cheap, no LLM. Output identical to the previous tick skips the agent run entirely; changed output wakes the agent with a diff injected into the prompt. First tick always runs (baseline). Output must be deterministic (no timestamps) or every tick looks changed. Incompatible with no_agent. On update, '' clears."
|
|
},
|
|
"no_agent": {
|
|
"type": "boolean",
|
|
"default": False,
|
|
"description": "True = no LLM: the scheduler runs `script` (required) on schedule and delivers its stdout verbatim; empty stdout sends nothing (watchdog pattern). Use for script-only pings with fixed output; keep False for anything needing reasoning."
|
|
},
|
|
"context_from": {
|
|
"type": "array",
|
|
"items": {"type": "string"},
|
|
"description": "Optional job ID(s) whose most recent completed output is injected as context each run — chains jobs (A collects, B processes). For a job's OWN previous output prefer `continuity`. On update, [] clears."
|
|
},
|
|
"continuity": {
|
|
"type": "boolean",
|
|
"description": "True = each run sees the job's own previous output, so it can dedupe and continue where it left off (scouts, monitors, incremental digests). Default false. On update, false turns it off."
|
|
},
|
|
"enabled_toolsets": {
|
|
"type": "array",
|
|
"items": {"type": "string"},
|
|
"description": "Optional toolset names to restrict the job's agent to (e.g. [\"web\", \"terminal\"]) — cuts token overhead. Infer from the prompt. Omit for all default tools. On update, [] clears."
|
|
},
|
|
"workdir": {
|
|
"type": "string",
|
|
"description": "Optional absolute existing path to run the job from: injects that directory's AGENTS.md/context files and anchors terminal/file tools there. On update, '' clears."
|
|
},
|
|
"attach_to_session": {
|
|
"type": "boolean",
|
|
"description": "True = the job's delivery is CONTINUABLE — the user can reply and the agent has the brief in context (threads on thread-capable platforms, mirrored into the DM elsewhere). Use for conversational recurring jobs (briefings); leave unset for fire-and-forget alerts. Scope: the job's own conversation only — the origin chat, the home-channel fallback when deliver='origin' captured no origin (script-created jobs), or the job's single explicit platform:chat target (this flag is the only way to attach an explicit target). Broadcast targets are never attached; no effect when deliver='local'."
|
|
},
|
|
},
|
|
"required": ["action"]
|
|
}
|
|
}
|
|
|
|
|
|
def check_cronjob_requirements() -> bool:
|
|
"""
|
|
Check if cronjob tools can be used.
|
|
|
|
Available in interactive CLI mode and gateway/messaging platforms.
|
|
The cron system is internal (JSON file-based scheduler ticked by the gateway),
|
|
so no external crontab executable is required.
|
|
|
|
Session env vars must hold an explicit truthy string (``1``, ``true``,
|
|
``yes``, ``on``) — false-like values (``0``, ``false``, ``no``, ``off``)
|
|
leave the tool disabled. Uses the shared ``env_var_enabled`` helper so
|
|
every consumer of these flags agrees on the truthy set.
|
|
"""
|
|
from utils import env_var_enabled
|
|
|
|
return (
|
|
env_var_enabled("HERMES_INTERACTIVE")
|
|
or env_var_enabled("HERMES_GATEWAY_SESSION")
|
|
or env_var_enabled("HERMES_EXEC_ASK")
|
|
)
|
|
|
|
|
|
# --- Registry ---
|
|
from tools.registry import registry, tool_error
|
|
|
|
|
|
def _cronjob_handler(args, **kw):
|
|
"""Model-tool dispatch for ``cronjob``.
|
|
|
|
Resolves the one model-facing ``monitor`` field into the stored
|
|
``monitor_script``/``monitor_url`` pair (legacy field names still accepted
|
|
as aliases so older transcripts/replays keep working).
|
|
"""
|
|
_mon_script, _mon_url = _split_monitor_arg(
|
|
args.get("monitor"), args.get("monitor_script"), args.get("monitor_url")
|
|
)
|
|
return cronjob(
|
|
action=args.get("action", ""),
|
|
job_id=args.get("job_id"),
|
|
prompt=args.get("prompt"),
|
|
schedule=args.get("schedule"),
|
|
name=args.get("name"),
|
|
repeat=args.get("repeat"),
|
|
deliver=args.get("deliver"),
|
|
failure_deliver=args.get("failure_deliver"),
|
|
include_disabled=args.get("include_disabled", True),
|
|
skill=args.get("skill"),
|
|
skills=args.get("skills"),
|
|
# model / provider / base_url are intentionally NOT read from the
|
|
# agent's arguments: per-job inference pins are user-owned (dashboard,
|
|
# `hermes cron create/edit --model`, or hand-edited jobs). The agent
|
|
# must not be able to point unattended spend at a different model.
|
|
# Programmatic callers of cronjob() itself retain the parameters.
|
|
reason=args.get("reason"),
|
|
script=args.get("script"),
|
|
context_from=args.get("context_from"),
|
|
continuity=args.get("continuity"),
|
|
enabled_toolsets=args.get("enabled_toolsets"),
|
|
workdir=args.get("workdir"),
|
|
no_agent=args.get("no_agent"),
|
|
attach_to_session=args.get("attach_to_session"),
|
|
monitor_script=_mon_script,
|
|
monitor_url=_mon_url,
|
|
task_id=kw.get("task_id"),
|
|
session_id=kw.get("session_id"),
|
|
)
|
|
|
|
|
|
registry.register(
|
|
name="cronjob_manage",
|
|
toolset="cronjob",
|
|
schema=CRONJOB_SCHEMA,
|
|
handler=_cronjob_handler,
|
|
check_fn=check_cronjob_requirements,
|
|
emoji="⏰",
|
|
)
|