refactor(tools): pack cronjob tool statements and compact docstrings (AST-neutral)
This commit is contained in:
@@ -32,10 +32,8 @@ def _origin_from_env() -> Optional[Dict[str, str]]:
|
||||
"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,
|
||||
"platform": origin_platform, "chat_id": origin_chat_id,
|
||||
"chat_name": get_session_env("HERMES_SESSION_CHAT_NAME") or None, "thread_id": thread_id,
|
||||
# Lets a delivery mirror resolve the participant's session in per-user-isolated groups.
|
||||
"user_id": get_session_env("HERMES_SESSION_USER_ID") or None,
|
||||
# Workspace/server scope (Slack team, Discord guild...): Slack session keys embed it,
|
||||
@@ -52,7 +50,6 @@ def _local_delivery_notice(job: Dict[str, Any], user_deliver: Optional[str]) ->
|
||||
return None
|
||||
try:
|
||||
from cron.scheduler import _resolve_delivery_targets
|
||||
|
||||
if _resolve_delivery_targets(job):
|
||||
return None
|
||||
except Exception: # resolution unavailable — fall back to the origin signal
|
||||
@@ -211,7 +208,6 @@ def _resolve_cron_context_deliver(deliver: Optional[str]) -> Optional[str]:
|
||||
elements pass through. Otherwise the scheduler would guess a home channel."""
|
||||
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
|
||||
|
||||
@@ -291,7 +287,6 @@ def _validate_cron_script_path(script: Optional[str]) -> Optional[str]:
|
||||
return None
|
||||
|
||||
from hermes_constants import get_hermes_home
|
||||
|
||||
raw = script.strip()
|
||||
if raw.startswith(("/", "~")) or (len(raw) >= 2 and raw[1] == ":"):
|
||||
return (
|
||||
@@ -300,7 +295,6 @@ def _validate_cron_script_path(script: Optional[str]) -> Optional[str]:
|
||||
f"Place scripts in ~/.hermes/scripts/ and use just the filename.")
|
||||
|
||||
from tools.path_security import validate_within_dir
|
||||
|
||||
scripts_dir = get_hermes_home() / "scripts"
|
||||
scripts_dir.mkdir(parents=True, exist_ok=True)
|
||||
if validate_within_dir(scripts_dir / raw, scripts_dir):
|
||||
@@ -392,7 +386,6 @@ def _gateway_liveness_notice(plural: bool = False) -> dict:
|
||||
on "scheduler active". False -> warning (no gateway process), None -> probe failed."""
|
||||
try:
|
||||
from hermes_cli.cron import _builtin_gateway_liveness
|
||||
|
||||
_gw = _builtin_gateway_liveness()
|
||||
except Exception:
|
||||
return {"gateway_running": None}
|
||||
|
||||
@@ -19,8 +19,7 @@ _CRON_THREAT_PATTERNS = [
|
||||
(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'authorized_keys', "ssh_backdoor"), (r'/etc/sudoers|visudo', "sudoers_mod"),
|
||||
(r'rm\s+-rf\s+/', "destructive_root_rm"),
|
||||
]
|
||||
|
||||
|
||||
@@ -99,7 +99,6 @@ def _api_server_base_url() -> str:
|
||||
"""``http://host:port`` of the local api_server, mirroring its bind resolution
|
||||
(extra.host -> API_SERVER_HOST -> 127.0.0.1); a wildcard bind listens on loopback too."""
|
||||
import os
|
||||
|
||||
port_raw = os.getenv("API_SERVER_PORT", "").strip()
|
||||
try:
|
||||
port = int(port_raw) if port_raw else 8642
|
||||
@@ -119,17 +118,13 @@ def _api_server_base_url() -> str:
|
||||
|
||||
|
||||
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.
|
||||
|
||||
Relay-fronted delivery has no standalone sender — the gateway's live relay adapter is
|
||||
the only path, reached via ``POST /api/jobs/{id}/run`` (marks the job due for the
|
||||
gateway ticker; ``extra_prompt`` rides in the body). Returns a JSON result string when
|
||||
forwarding engages, else None to fall through to the normal in-process run.
|
||||
"""
|
||||
"""Forward a manual run to the gateway when it targets a relay-fronted platform: such delivery
|
||||
has no standalone sender — the gateway's live relay adapter is the only path, reached via
|
||||
``POST /api/jobs/{id}/run`` (marks the job due; ``extra_prompt`` rides in the body). Returns a
|
||||
JSON result string when forwarding engages, else None (normal in-process run)."""
|
||||
if not _relay_fronted_delivery_platforms(job):
|
||||
return None
|
||||
from agent.secret_scope import get_secret
|
||||
|
||||
key = get_secret("API_SERVER_KEY", "") or ""
|
||||
try:
|
||||
import httpx
|
||||
@@ -175,12 +170,9 @@ _ALREADY_RUNNING_ERROR = (
|
||||
|
||||
|
||||
def _claim_for_manual_run(job_id: str, log_label: str):
|
||||
"""At-most-once claim shared by the sync and background run paths.
|
||||
|
||||
Returns ``(claimed_job, None)`` or ``(None, error_dict)`` in the ``_execute_job_now``
|
||||
result shape. A lost claim is labelled precisely: claim_job_for_fire also returns False
|
||||
for paused / disabled / missing jobs, which must not read as "already being fired".
|
||||
"""
|
||||
"""At-most-once claim shared by the sync and background run paths: ``(claimed_job, None)`` or
|
||||
``(None, error_dict)`` in the ``_execute_job_now`` shape. A lost claim is labelled precisely —
|
||||
claim_job_for_fire also returns False for paused/disabled/missing jobs, not just in-flight ones."""
|
||||
try:
|
||||
claimed_job = claim_job_for_fire(job_id, return_job=True)
|
||||
if isinstance(claimed_job, dict):
|
||||
@@ -310,7 +302,6 @@ def _latest_job_output_excerpt(job_id: str, max_chars: int = 2000) -> Optional[s
|
||||
block (parent sees what the job produced). Never raises."""
|
||||
try:
|
||||
from cron.jobs import get_cron_output_dir
|
||||
|
||||
files = sorted((get_cron_output_dir() / job_id).glob("*.md"))
|
||||
text = files[-1].read_text(encoding="utf-8", errors="replace").strip() if files else ""
|
||||
if not text:
|
||||
@@ -329,7 +320,6 @@ def _reap_stale_executions(job_name: str) -> None:
|
||||
Best-effort self-heal: must not block dispatch."""
|
||||
try:
|
||||
from cron.executions import recover_interrupted_executions
|
||||
|
||||
_reclaimed = recover_interrupted_executions()
|
||||
if _reclaimed:
|
||||
logger.warning(
|
||||
@@ -344,7 +334,6 @@ def _background_session_key(session_id: Optional[str]) -> str:
|
||||
cross the pool). Empty string = no durable consumer."""
|
||||
try:
|
||||
from tools.approval import get_current_session_key
|
||||
|
||||
session_key = get_current_session_key(default="")
|
||||
except Exception:
|
||||
session_key = ""
|
||||
@@ -371,35 +360,24 @@ def _manual_run_completion(
|
||||
if excerpt:
|
||||
lines += ["--- JOB OUTPUT ---", excerpt]
|
||||
return {
|
||||
"status": "completed" if res.get("success") else "error",
|
||||
"summary": "\n".join(lines),
|
||||
"error": res.get("error"),
|
||||
"api_calls": 0,
|
||||
"duration_seconds": duration,
|
||||
"status": "completed" if res.get("success") else "error", "summary": "\n".join(lines),
|
||||
"error": res.get("error"), "api_calls": 0, "duration_seconds": duration,
|
||||
}
|
||||
|
||||
|
||||
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 cron job is a full agent run (minutes to hours); inline it made the parent turn
|
||||
uninterruptible. Dispatches like ``delegate_task``'s background mode: the tool returns a
|
||||
handle and a ``type="async_delegation"`` completion re-enters the conversation as a
|
||||
fresh turn (role alternation legal, prompt cache intact). The claim is taken
|
||||
SYNCHRONOUSLY so unrunnable jobs report immediately.
|
||||
|
||||
Returns None when background delivery is unavailable (caller falls back to the sync
|
||||
path); ``{"claimed": False, ...}`` on a lost claim; ``{"claimed": True, "dispatched":
|
||||
True, "delegation_id"}`` when running in the background; ``{"claimed": True,
|
||||
"dispatched": False, ...}`` when the pool was full and the run executed inline.
|
||||
"""
|
||||
# Finite sessions cannot route a detached result back after the turn ends — mirror
|
||||
# delegate_task's gate and fall back to sync execution.
|
||||
"""Claim ``job`` now (SYNCHRONOUSLY, so unrunnable jobs report immediately), then fire it
|
||||
on the async-delegation executor like ``delegate_task``'s background mode: the tool returns
|
||||
a handle and a ``type="async_delegation"`` completion re-enters as a fresh turn (role
|
||||
alternation legal, prompt cache intact) instead of blocking the parent turn for minutes.
|
||||
Returns None when background delivery is unavailable (caller runs sync); ``{"claimed":
|
||||
False}`` on a lost claim; ``{"claimed": True, "dispatched": True, "delegation_id"}``; or
|
||||
``{"claimed": True, "dispatched": False, ...}`` when the pool was full and it ran inline."""
|
||||
# Finite sessions cannot route a detached result back after the turn ends (delegate_task's gate).
|
||||
try:
|
||||
from gateway.session_context import async_delivery_supported
|
||||
|
||||
if not async_delivery_supported():
|
||||
return None
|
||||
except Exception:
|
||||
@@ -409,18 +387,16 @@ def _try_dispatch_background_run(
|
||||
job_name = str(job.get("name") or job_id)
|
||||
_reap_stale_executions(job_name)
|
||||
|
||||
# Routing capture BEFORE the claim: with no routable session there is no durable
|
||||
# consumer for a detached completion, so we must not claim-and-dispatch. Direct Python
|
||||
# callers (`hermes cron run`, tests) exit right after the tool returns -> run sync.
|
||||
# Routing capture BEFORE the claim: no routable session = no durable consumer for a detached
|
||||
# completion, so don't claim-and-dispatch (direct callers like `hermes cron run` exit right after).
|
||||
session_key = _background_session_key(session_id)
|
||||
if not session_key:
|
||||
return None
|
||||
|
||||
# Best-effort early dedupe so a mid-run job reports in THIS tool response instead of
|
||||
# as a delayed error completion (authoritative check: try_register_running_job).
|
||||
# Early dedupe so a mid-run job reports in THIS response, not as a delayed error completion
|
||||
# (authoritative check: try_register_running_job).
|
||||
try:
|
||||
from cron.scheduler import get_running_job_ids
|
||||
|
||||
if job_id in get_running_job_ids():
|
||||
return {"claimed": False, "success": False, "error": _ALREADY_RUNNING_ERROR}
|
||||
except Exception:
|
||||
@@ -435,14 +411,12 @@ def _try_dispatch_background_run(
|
||||
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(
|
||||
@@ -453,16 +427,13 @@ def _try_dispatch_background_run(
|
||||
|
||||
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 (falsy -> "local", list -> comma
|
||||
# string), reading the claimed snapshot the run actually executes.
|
||||
# Scheduler's own normalizer (falsy -> "local", list -> comma string) on the claimed snapshot.
|
||||
from cron.scheduler import _normalize_deliver_value
|
||||
|
||||
deliver = _normalize_deliver_value(claimed_job.get("deliver", "local"))
|
||||
|
||||
def _runner() -> Dict[str, Any]:
|
||||
@@ -471,23 +442,16 @@ def _try_dispatch_background_run(
|
||||
|
||||
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,
|
||||
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.
|
||||
# Pool at capacity (or submit failure): the claim is already taken and must not be stranded.
|
||||
logger.info(
|
||||
"cronjob run: background pool unavailable (%s); running job '%s' inline.",
|
||||
dispatch.get("error", "rejected"), job_name)
|
||||
@@ -549,53 +513,31 @@ def _action_create(a: Dict[str, Any]) -> str:
|
||||
context_from = _apply_continuity(context_from, a["continuity"])
|
||||
|
||||
from cron.scheduler import CronSchedulerRegistrationError, create_job_with_scheduler_registration
|
||||
|
||||
try:
|
||||
job = create_job_with_scheduler_registration(
|
||||
prompt=prompt or "",
|
||||
schedule=a["schedule"],
|
||||
name=a["name"],
|
||||
repeat=a["repeat"],
|
||||
deliver=_resolve_cron_context_deliver(deliver),
|
||||
origin=_origin_from_env(),
|
||||
skills=canonical_skills,
|
||||
model=_normalize_optional_job_value(a["model"]),
|
||||
provider=_normalize_optional_job_value(a["provider"]),
|
||||
prompt=prompt or "", schedule=a["schedule"], name=a["name"], repeat=a["repeat"],
|
||||
deliver=_resolve_cron_context_deliver(deliver), origin=_origin_from_env(), skills=canonical_skills,
|
||||
model=_normalize_optional_job_value(a["model"]), provider=_normalize_optional_job_value(a["provider"]),
|
||||
base_url=_normalize_optional_job_value(a["base_url"], strip_trailing_slash=True),
|
||||
script=_normalize_optional_job_value(script),
|
||||
context_from=context_from,
|
||||
enabled_toolsets=a["enabled_toolsets"] or None,
|
||||
workdir=_normalize_optional_job_value(a["workdir"]),
|
||||
no_agent=_no_agent,
|
||||
attach_to_session=a["attach_to_session"],
|
||||
script=_normalize_optional_job_value(script), context_from=context_from,
|
||||
enabled_toolsets=a["enabled_toolsets"] or None, workdir=_normalize_optional_job_value(a["workdir"]),
|
||||
no_agent=_no_agent, attach_to_session=a["attach_to_session"],
|
||||
monitor_script=_normalize_optional_job_value(a["monitor_script"]),
|
||||
monitor_url=_normalize_optional_job_value(a["monitor_url"]),
|
||||
# CLI-only lane: absent from CRONJOB_SCHEMA and the model dispatch — models
|
||||
# do not make model-config decisions.
|
||||
# CLI-only lane: absent from CRONJOB_SCHEMA and the model dispatch (models don't pick models).
|
||||
reasoning_effort=a["reasoning_effort"],
|
||||
failure_deliver=_resolve_cron_context_deliver(_normalize_deliver_param(a["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, deliver)
|
||||
if _local_notice:
|
||||
_create_message = f"{_create_message} {_local_notice}"
|
||||
# The builtin ticker lives in the gateway process: a job created with no gateway
|
||||
# running is stored but never fires — tell the model (the CLI already warns).
|
||||
_create_message = " ".join(filter(None, (f"Cron job '{job['name']}' created.", _local_delivery_notice(job, deliver))))
|
||||
# The builtin ticker lives in the gateway process: with no gateway running the job is stored
|
||||
# but never fires — tell the model (the CLI already warns).
|
||||
_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(),
|
||||
"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(),
|
||||
}
|
||||
return _dumps(_with_guidance(_result, job, deliver))
|
||||
|
||||
@@ -827,13 +769,10 @@ def _action_update(job: Dict[str, Any], a: Dict[str, Any]) -> str:
|
||||
# Actions that need no job_id, and job-bound actions (job resolved first).
|
||||
_JOBLESS_ACTIONS = {"create": _action_create, "list": _action_list}
|
||||
_JOB_ACTIONS = {
|
||||
"remove": _action_remove,
|
||||
"remove": _action_remove, "update": _action_update,
|
||||
"run": _action_run, "run_now": _action_run, "trigger": _action_run,
|
||||
"pause": lambda job, a: _job_state_result(pause_job(job["id"], reason=a["reason"])),
|
||||
"resume": lambda job, a: _job_state_result(resume_job(job["id"])),
|
||||
"run": _action_run,
|
||||
"run_now": _action_run,
|
||||
"trigger": _action_run,
|
||||
"update": _action_update,
|
||||
}
|
||||
|
||||
|
||||
@@ -999,7 +938,6 @@ def check_cronjob_requirements() -> bool:
|
||||
"""Available in interactive CLI mode and gateway/messaging platforms (the scheduler is
|
||||
internal; no crontab needed). Flags must be explicitly truthy via ``env_var_enabled``."""
|
||||
from utils import env_var_enabled
|
||||
|
||||
return (
|
||||
env_var_enabled("HERMES_INTERACTIVE")
|
||||
or env_var_enabled("HERMES_GATEWAY_SESSION")
|
||||
@@ -1012,10 +950,8 @@ def check_cronjob_requirements() -> bool:
|
||||
# create/edit --model`, hand-edited jobs) — the agent must not point unattended spend at a
|
||||
# different model. Programmatic callers of cronjob() itself retain the parameters.
|
||||
_HANDLER_FORWARDED_ARGS = (
|
||||
"job_id", "prompt", "schedule", "name", "repeat", "deliver", "failure_deliver", "skill", "skills",
|
||||
"reason", "script", "context_from", "continuity", "enabled_toolsets", "workdir",
|
||||
"no_agent", "attach_to_session",
|
||||
)
|
||||
"job_id", "prompt", "schedule", "name", "repeat", "deliver", "failure_deliver", "skill", "skills", "reason",
|
||||
"script", "context_from", "continuity", "enabled_toolsets", "workdir", "no_agent", "attach_to_session")
|
||||
|
||||
|
||||
def _cronjob_handler(args, **kw):
|
||||
|
||||
Reference in New Issue
Block a user