diff --git a/tools/cronjob_job_args.py b/tools/cronjob_job_args.py index 80ec1ef32c..ee7bdfa217 100644 --- a/tools/cronjob_job_args.py +++ b/tools/cronjob_job_args.py @@ -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} diff --git a/tools/cronjob_prompt_scan.py b/tools/cronjob_prompt_scan.py index 20ca81fd22..6f8522d8ee 100644 --- a/tools/cronjob_prompt_scan.py +++ b/tools/cronjob_prompt_scan.py @@ -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"), ] diff --git a/tools/cronjob_tools.py b/tools/cronjob_tools.py index 2efe325a64..3c5b6881d8 100644 --- a/tools/cronjob_tools.py +++ b/tools/cronjob_tools.py @@ -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):