Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions agent/estop.py
Original file line number Diff line number Diff line change
Expand Up @@ -57,11 +57,18 @@ def sentinel_path() -> Path:


def is_engaged() -> bool:
"""Cheap check (one stat): is the global emergency stop engaged?"""
"""Cheap check (one stat): is the global emergency stop engaged?

Fail SAFE on stat errors: if we cannot determine whether the sentinel
exists (permission error, transient I/O failure on HERMES_HOME), report
engaged. The module contract is that the pause must hold even when the
sentinel is unreadable — a fail-open here would silently lift an
operator's emergency stop exactly when the filesystem is misbehaving.
"""
try:
return sentinel_path().exists()
except OSError:
return False
return True


def engage(reason: Optional[str] = None) -> Path:
Expand Down
94 changes: 75 additions & 19 deletions cron/jobs.py
Original file line number Diff line number Diff line change
Expand Up @@ -1534,6 +1534,36 @@ def _normalized_inference_axes(job: Dict[str, Any]) -> Tuple[Optional[str], Opti
)


def _validate_job_mode_invariants(
monitor_script: Optional[str],
monitor_url: Optional[str],
no_agent: bool,
script: Optional[str],
) -> None:
"""Shared create/update validation for job execution-mode invariants.

ONE owner for the class: create_job and update_job both call this so an
invariant enforced at create time cannot be violated through the update
door (monitor jobs silently degrading when no_agent is flipped on, etc.).
"""
if monitor_script and monitor_url:
raise ValueError(
"monitor_script and monitor_url are mutually exclusive — a job "
"can only have one monitor source."
)
if (monitor_script or monitor_url) and no_agent:
raise ValueError(
"monitor_script/monitor_url cannot be combined with no_agent=True — "
"the whole point of a monitor job is to suppress or wake the AGENT "
"based on source changes. Use a plain no_agent script job instead."
)
if no_agent and not script:
raise ValueError(
"no_agent=True requires a script — with no agent and no script "
"there is nothing for the job to run."
)


def create_job(
prompt: Optional[str],
schedule: str,
Expand Down Expand Up @@ -1650,26 +1680,16 @@ def create_job(

# Monitor-mode validation: exactly one source, and monitor mode only
# makes sense when there IS an agent to suppress/wake.
if normalized_monitor_script and normalized_monitor_url:
raise ValueError(
"monitor_script and monitor_url are mutually exclusive — a job "
"can only have one monitor source."
)
if (normalized_monitor_script or normalized_monitor_url) and normalized_no_agent:
raise ValueError(
"monitor_script/monitor_url cannot be combined with no_agent=True — "
"the whole point of a monitor job is to suppress or wake the AGENT "
"based on source changes. Use a plain no_agent script job instead."
)

# no_agent jobs are meaningless without a script — the script IS the job.
# Surface this as a clear ValueError at create time so bad configs never
# reach the scheduler.
if normalized_no_agent and not normalized_script:
raise ValueError(
"no_agent=True requires a script — with no agent and no script "
"there is nothing for the job to run."
)
# Surface these as clear ValueErrors at create time so bad configs never
# reach the scheduler (shared with update_job, see
# _validate_job_mode_invariants).
_validate_job_mode_invariants(
normalized_monitor_script,
normalized_monitor_url,
normalized_no_agent,
normalized_script,
)

# Normalize context_from: accept str or list of str, store as list or None
if isinstance(context_from, str):
Expand Down Expand Up @@ -1859,8 +1879,33 @@ def update_job(job_id: str, updates: Dict[str, Any]) -> Optional[Dict[str, Any]]
else:
updates["workdir"] = _normalize_workdir(_wd)

# Normalize monitor fields the same way create_job does (empty
# string clears the field).
for _mon_field in ("monitor_script", "monitor_url"):
if _mon_field in updates:
_mv = updates[_mon_field]
_mv = str(_mv).strip() if isinstance(_mv, str) else None
updates[_mon_field] = _mv or None

previous_inference_axes = _normalized_inference_axes(job)
updated = _apply_skill_fields({**job, **updates})

# Re-check execution-mode invariants on the MERGED record when
# any participating field changes, so create-time invariants
# can't be violated through the update door (e.g. flipping
# no_agent=True on a monitor job would silently disable the
# monitor: the scheduler's no_agent short-circuit runs before
# the monitor gate). Scoped to changed fields so legacy records
# untouched by this update keep loading.
if {"monitor_script", "monitor_url", "no_agent", "script"}.intersection(updates):
_upd_script = updated.get("script")
_upd_script = str(_upd_script).strip() if isinstance(_upd_script, str) else None
_validate_job_mode_invariants(
updated.get("monitor_script") or None,
updated.get("monitor_url") or None,
bool(updated.get("no_agent")),
_upd_script or None,
)
schedule_changed = "schedule" in updates
inference_fields_changed = bool(
{"provider", "model", "base_url", "no_agent"}.intersection(updates)
Expand Down Expand Up @@ -2011,6 +2056,17 @@ def remove_job(job_id: str) -> bool:
# Clean up output directory to prevent orphaned dirs accumulating
if job_output_dir.exists():
shutil.rmtree(job_output_dir)
# Clean up the job's durable notepad (cron/notepad.db) — without
# this, removed jobs orphan their KV rows forever. Best effort:
# a notepad failure must never block the removal itself.
try:
from cron.notepad import clear_notepad
clear_notepad(canonical_id)
except Exception:
logger.debug(
"Failed to clear notepad for removed job %s",
canonical_id, exc_info=True,
)
return True
return False

Expand Down
8 changes: 7 additions & 1 deletion cron/notepad.py
Original file line number Diff line number Diff line change
Expand Up @@ -150,7 +150,13 @@ def list_notes(job_id: str) -> List[Dict[str, Any]]:


def clear_notepad(job_id: str) -> int:
"""Delete every key for one job (e.g. on job removal). Returns row count."""
"""Delete every key for one job (e.g. on job removal). Returns row count.

Called from ``cron.jobs.remove_job`` so deleted jobs don't orphan their
rows. No-ops without creating the DB when no notepad file exists yet.
"""
if not NOTEPAD_FILE.exists():
return 0
with _transaction() as conn:
cur = conn.execute(
"DELETE FROM cron_notepad WHERE job_id=?", (str(job_id),)
Expand Down
97 changes: 91 additions & 6 deletions gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -14364,6 +14364,7 @@ async def _dispatch_busy_slash_command(
"restart": self._handle_restart_command,
"approve": self._handle_approve_command,
"deny": self._handle_deny_command,
"pause": self._handle_pause_command,
"agents": self._handle_agents_command,
"background": self._handle_background_command,
"kanban": self._handle_kanban_command,
Expand Down Expand Up @@ -14393,6 +14394,36 @@ async def _dispatch_busy_slash_command(
f"mid-turn. Wait for the current response or `/stop` first."
)

async def _handle_pause_command(self, event: MessageEvent):
"""`/pause [reason]` engages the global emergency stop; `/pause off`
(aliases: resume/stop) lifts it.

This is the in-band resume path for messaging-only operators — the
estop gate above deliberately lets recognized slash commands through
while paused so a user without host-shell access is never locked out.
"""
from agent import estop

args = (event.get_command_args() or "").strip()
if args.lower() in {"off", "resume", "stop", "disengage"}:
if estop.disengage():
return "▶️ Resumed — new work is accepted again."
return "Hermes wasn't paused."
state = estop.get_state()
if state is not None and not args:
reason = state.get("reason")
suffix = f" (reason: {reason})" if reason else ""
return (
f"⏸️ Hermes is already paused{suffix}. "
"Use `/pause off` to resume."
)
estop.engage(reason=args or None)
suffix = f" (reason: {args})" if args else ""
return (
f"⏸️ Paused{suffix}. New cron/kanban/gateway work is on hold; "
"in-flight work finishes normally. Use `/pause off` to resume."
)

async def _busy_start_command(self, event: MessageEvent, quick_key: str, source):
# Telegram sends /start for bot launches/deep-links. Treat it as a
# platform ping, not a user command: no help dump, no agent
Expand Down Expand Up @@ -14740,19 +14771,70 @@ async def _handle_message(self, event: MessageEvent) -> Optional[str]:
# gate — pause stops NEW work, it never kills or orphans running
# work. Placed after auth so unauthorized senders keep the normal
# silent/pairing behavior and can't probe pause state.
#
# Passthroughs (pause blocks new AGENT turns, not control traffic):
# * recognized slash commands — /status, /help, /new, /approve and
# friends must keep working while paused, and /pause off is the
# in-band resume path for messaging-only users;
# * replies owned by IN-FLIGHT work — a pending detached-update
# prompt, clarify, slash-confirm, or dangerous-command approval,
# plus any message steering a session whose agent is already
# running. Swallowing those would stall work the pause promised
# not to touch.
if not is_internal:
try:
from agent.estop import paused_reply as _estop_paused_reply
_paused_notice = _estop_paused_reply()
except ImportError:
_paused_notice = None
if _paused_notice is not None:
logger.info(
"Gateway turn paused by global emergency stop (platform=%s chat=%s)",
getattr(getattr(source, "platform", None), "value", "unknown"),
getattr(source, "chat_id", None) or "unknown",
)
return _paused_notice
_estop_allow = False
_estop_cmd = None
try:
_estop_cmd = event.get_command()
except Exception:
_estop_cmd = None
if _estop_cmd:
try:
from hermes_cli.commands import (
resolve_command as _resolve_estop_cmd,
)
_estop_allow = _resolve_estop_cmd(_estop_cmd) is not None
except Exception:
_estop_allow = False
if not _estop_allow:
try:
_estop_key = self._session_key_for_source(source)
_estop_state = self._peek_session_state(_estop_key)
if (
_estop_state is not None
and _estop_state.persistent.update_prompt_pending
):
_estop_allow = True
if not _estop_allow and self._is_session_running(_estop_key):
# Steering / interrupting in-flight work (which
# also covers pending clarify + tool approvals
# held by the running agent).
_estop_allow = True
if not _estop_allow:
from tools import slash_confirm as _estop_confirm_mod
if _estop_confirm_mod.get_pending(_estop_key):
_estop_allow = True
if not _estop_allow:
from tools.approval import (
has_blocking_approval as _estop_has_approval,
)
if _estop_has_approval(_estop_key):
_estop_allow = True
except Exception:
pass
if not _estop_allow:
logger.info(
"Gateway turn paused by global emergency stop (platform=%s chat=%s)",
getattr(getattr(source, "platform", None), "value", "unknown"),
getattr(source, "chat_id", None) or "unknown",
)
return _paused_notice

# Intercept messages that are responses to a pending /update prompt.
# The update process (detached) wrote .update_prompt.json; the watcher
Expand Down Expand Up @@ -15302,6 +15384,9 @@ async def _handle_message(self, event: MessageEvent) -> Optional[str]:
canonical = _cmd_def.name if _cmd_def else command
break

if canonical == "pause":
return await self._handle_pause_command(event)

if canonical == "new":
if await asyncio.to_thread(self._is_telegram_topic_root_lobby, source):
return self._telegram_topic_root_new_message()
Expand Down
7 changes: 6 additions & 1 deletion hermes_cli/commands.py
Original file line number Diff line number Diff line change
Expand Up @@ -139,6 +139,9 @@ class CommandDef:
cli_only=True, args_hint="<archive.tar.gz> [--name <name>]"),
CommandDef("stop", "Kill all running background processes", "Session",
busy_policy="interrupt_then_dispatch", busy_handler="stop"),
CommandDef("pause", "Pause new work globally (emergency stop); '/pause off' resumes", "Session",
gateway_only=True, args_hint="[reason | off]",
busy_policy="dispatch"),
CommandDef("approve", "Approve a pending dangerous command", "Session",
gateway_only=True, args_hint="[session|always]", busy_policy="dispatch"),
CommandDef("deny", "Deny a pending dangerous command (optionally with a reason)", "Session",
Expand Down Expand Up @@ -1272,7 +1275,9 @@ def discord_skill_commands_by_category(
# - refine: on-demand memory/skill review; reached via /hermes refine on
# Slack. Added at the 50-cap — a native slot would clamp an existing
# native slash.
_SLACK_VIA_HERMES_ONLY = frozenset({"topup", "moa", "debug", "egress", "init", "version", "diff", "update", "heartbeat", "refine"})
# - pause: global emergency stop; reached via /hermes pause [off] on
# Slack. Added at the 50-cap — a native slot would clamp /platform.
_SLACK_VIA_HERMES_ONLY = frozenset({"topup", "moa", "debug", "egress", "init", "version", "diff", "update", "heartbeat", "refine", "pause"})


def _sanitize_slack_name(raw: str) -> str:
Expand Down
Loading