refactor(hermes_cli): cron/debug/dashboard_procs — inline single-use locals, collapse try/pass to suppress, compact rationale comments (every WHY kept)
This commit is contained in:
@@ -12,11 +12,9 @@ sys.path.insert(0, str(PROJECT_ROOT))
|
||||
|
||||
from hermes_cli.colors import Colors, color
|
||||
|
||||
# Gateway-lifecycle command detection lives in ``cron.lifecycle_guard`` (shared by every
|
||||
# job-creation path without a circular import). Re-exported here because
|
||||
# ``tools/terminal_tool.py`` imports it from this module to hard-block the same commands
|
||||
# at execution time when ``_HERMES_GATEWAY=1``.
|
||||
from cron.lifecycle_guard import ( # noqa: F401 (re-exported for terminal_tool)
|
||||
# Re-exported: ``tools/terminal_tool.py`` imports it from here to hard-block the same
|
||||
# gateway-lifecycle commands at execution time when ``_HERMES_GATEWAY=1``.
|
||||
from cron.lifecycle_guard import ( # noqa: F401
|
||||
contains_gateway_lifecycle_command as _contains_gateway_lifecycle_command,
|
||||
)
|
||||
|
||||
@@ -25,9 +23,8 @@ def _normalize_skills(single_skill=None, skills: Optional[Iterable[str]] = None)
|
||||
"""Deduped, stripped skill names; None when neither argument was given."""
|
||||
if skills is None and single_skill is None:
|
||||
return None
|
||||
raw_items = list(skills) if skills is not None else [single_skill]
|
||||
normalized: List[str] = []
|
||||
for item in raw_items:
|
||||
for item in list(skills) if skills is not None else [single_skill]:
|
||||
text = str(item or "").strip()
|
||||
if text and text not in normalized:
|
||||
normalized.append(text)
|
||||
@@ -40,11 +37,7 @@ def _cron_api(**kwargs):
|
||||
|
||||
|
||||
def _active_cron_provider_name() -> str:
|
||||
"""Resolved cron scheduler provider name ('builtin', 'chronos', …); 'builtin' on any failure.
|
||||
|
||||
Best-effort + offline (``resolve_cron_scheduler`` reads config; ``is_available()`` forbids
|
||||
network), so callers fall back to the historical ticker-based checks.
|
||||
"""
|
||||
"""Resolved cron scheduler provider name ('builtin', 'chronos', …); 'builtin' on failure."""
|
||||
try:
|
||||
from cron.scheduler_provider import resolve_cron_scheduler
|
||||
return resolve_cron_scheduler().name or "builtin"
|
||||
@@ -55,19 +48,15 @@ def _active_cron_provider_name() -> str:
|
||||
def _builtin_gateway_liveness() -> Optional[bool]:
|
||||
"""Tri-state liveness of the builtin cron scheduler's trigger (None = unknown).
|
||||
|
||||
Shared by the CLI and the ``cronjob`` model tool: the builtin ticker only runs inside the
|
||||
gateway process, so a scheduled job with no live gateway can never fire. Non-builtin
|
||||
providers fire jobs without the gateway.
|
||||
The builtin ticker only runs inside the gateway process, so a scheduled job with no live
|
||||
gateway can never fire; non-builtin providers fire jobs without the gateway.
|
||||
"""
|
||||
try:
|
||||
if _active_cron_provider_name() != "builtin":
|
||||
return True # external provider fires jobs without the gateway
|
||||
# The gateway runtime lock is held for exactly the gateway's lifetime — a more
|
||||
# reliable "ticker's process is alive" signal than PID scanning, and inside the
|
||||
# gateway process it short-circuits to True so the in-gateway cron tool never
|
||||
# emits a false "gateway not running" (find_gateway_pids can transiently miss the
|
||||
# gateway just after a restart).
|
||||
# A crashing lock probe is "unknown", not "dead" — let the pid scan decide.
|
||||
return True
|
||||
# The runtime lock is held for exactly the gateway's lifetime — more reliable than PID
|
||||
# scanning (find_gateway_pids transiently misses the gateway right after a restart, and
|
||||
# inside the gateway it must never say "not running"). A crashing probe is "unknown".
|
||||
with contextlib.suppress(Exception):
|
||||
from gateway.status import is_gateway_runtime_lock_active
|
||||
if is_gateway_runtime_lock_active():
|
||||
@@ -75,27 +64,19 @@ def _builtin_gateway_liveness() -> Optional[bool]:
|
||||
from hermes_cli.gateway import (
|
||||
find_gateway_pids, named_profile_served_by_running_multiplexer,
|
||||
)
|
||||
|
||||
if find_gateway_pids():
|
||||
return True
|
||||
# Satellite profile: no local gateway.pid, but the default multiplexer ticks this
|
||||
# profile's cron store.
|
||||
return named_profile_served_by_running_multiplexer()
|
||||
# Satellite profile: no local gateway.pid, but the default multiplexer ticks its store.
|
||||
return bool(find_gateway_pids()) or named_profile_served_by_running_multiplexer()
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
|
||||
def _warn_if_gateway_not_running() -> None:
|
||||
"""Warn that scheduled jobs won't fire unless the gateway is running.
|
||||
"""Warn that scheduled jobs won't fire unless the gateway is running (the #1 cron report).
|
||||
|
||||
The cron ticker only runs inside the gateway (no standalone daemon): without one,
|
||||
``next_run_at`` passes but jobs never fire and ``last_run_at`` stays null — the most
|
||||
common cron support report.
|
||||
False is the only warn-worthy liveness state (None = unknown).
|
||||
"""
|
||||
# _builtin_gateway_liveness never raises; False is the only warn-worthy state.
|
||||
if _builtin_gateway_liveness() is not False:
|
||||
return
|
||||
|
||||
print(color(" ⚠ Gateway is not running — jobs won't fire automatically.", Colors.YELLOW))
|
||||
print(color(" Start it with: hermes gateway install\n"
|
||||
" sudo hermes gateway install --system # Linux servers\n"
|
||||
@@ -125,9 +106,7 @@ def _dispatch_display(dispatch: dict) -> Optional[str]:
|
||||
"""
|
||||
if not isinstance(dispatch, dict):
|
||||
return None
|
||||
scheduled = dispatch.get("scheduled_at")
|
||||
actual = dispatch.get("dispatched_at")
|
||||
kind = dispatch.get("kind")
|
||||
scheduled, actual, kind = (dispatch.get(k) for k in ("scheduled_at", "dispatched_at", "kind"))
|
||||
if not scheduled or not actual or not kind:
|
||||
return None
|
||||
lateness = _format_lateness(dispatch.get("lateness_seconds", 0))
|
||||
@@ -148,9 +127,7 @@ def _print_banner(title: str) -> None:
|
||||
|
||||
|
||||
def _unverified_targets(unverified) -> str:
|
||||
if isinstance(unverified, list):
|
||||
return ", ".join(str(t) for t in unverified)
|
||||
return str(unverified)
|
||||
return ", ".join(map(str, unverified)) if isinstance(unverified, list) else str(unverified)
|
||||
|
||||
|
||||
_STATE_BADGES = {"paused": ("[paused]", Colors.YELLOW), "completed": ("[completed]", Colors.BLUE)}
|
||||
@@ -169,10 +146,9 @@ def cron_list(show_all: bool = False):
|
||||
_print_banner("Scheduled Jobs")
|
||||
|
||||
for job in jobs:
|
||||
# Scheduler-honoured flag — never show [paused] when enabled=true.
|
||||
# effective_job_state honours the scheduler flag — never [paused] when enabled=true.
|
||||
badge = _STATE_BADGES.get(effective_job_state(job)) or (
|
||||
("[active]", Colors.GREEN) if job.get("enabled", True) else ("[disabled]", Colors.RED)
|
||||
)
|
||||
("[active]", Colors.GREEN) if job.get("enabled", True) else ("[disabled]", Colors.RED))
|
||||
print(f" {color(job.get('id', '?'), Colors.YELLOW)} {color(*badge)}")
|
||||
for label, value in _job_rows(job):
|
||||
print(f" {label + ':':<11}{value}")
|
||||
@@ -188,8 +164,7 @@ def _last_run_display(job: Dict[str, Any]) -> str:
|
||||
if last_status == "ok":
|
||||
return color("ok", Colors.GREEN)
|
||||
if last_status == "delivery_failed":
|
||||
# Agent succeeded but the result never reached the user — not green; the detail
|
||||
# lives in last_delivery_error (last_error is None for these runs).
|
||||
# Agent succeeded but the result never reached the user — not green; last_error is None.
|
||||
return color(f"delivery_failed: {job.get('last_delivery_error') or '?'}", Colors.YELLOW)
|
||||
display = color(f"{last_status}: {job.get('last_error', '?')}", Colors.RED)
|
||||
streak = int(job.get("failure_streak") or 0)
|
||||
@@ -200,9 +175,7 @@ def _last_run_display(job: Dict[str, Any]) -> str:
|
||||
|
||||
def _job_rows(job: Dict[str, Any]) -> List[tuple[str, str]]:
|
||||
"""``(label, value)`` detail rows for one job in ``cron list``."""
|
||||
# `repeat` / `deliver` may be present-but-null in the record (a one-shot persisted with
|
||||
# "repeat": null), so coalesce rather than rely on the dict-default, which only applies
|
||||
# to a missing key — a null deliver would crash `", ".join(None)`.
|
||||
# `repeat` / `deliver` may be present-but-null (dict-default only covers a missing key).
|
||||
repeat_info = job.get("repeat") or {}
|
||||
repeat_times = repeat_info.get("times")
|
||||
deliver = job.get("deliver") or ["local"]
|
||||
@@ -257,15 +230,13 @@ def cron_tick():
|
||||
try:
|
||||
tick(verbose=True)
|
||||
except CronTickYielded as exc:
|
||||
# Not expected here (a one-shot CLI process has no boot fingerprint, so the yield
|
||||
# gate is inert) — report cleanly instead of a traceback if a future caller records one.
|
||||
# Inert for a one-shot CLI (no boot fingerprint); report cleanly rather than traceback.
|
||||
print(color(f"✗ {exc}", Colors.YELLOW))
|
||||
print(" A fresher gateway process owns the runtime lock and will fire due jobs; this "
|
||||
"stale process yielded its tick.")
|
||||
return 1
|
||||
except OSError as exc:
|
||||
# tick() propagates real lock-acquisition failures (EMFILE, EACCES on open, ...)
|
||||
# instead of swallowing them as contention; the gateway ticker loop handles its own retry.
|
||||
# Real lock-acquisition failures (EMFILE, EACCES) propagate; they are not contention.
|
||||
print(color(f"✗ Cron tick failed: {exc}", Colors.RED))
|
||||
print(" Check `hermes cron status` and the gateway log for details.")
|
||||
return 1
|
||||
@@ -293,8 +264,7 @@ _INCIDENT_STATE_COLORS = {"detected": Colors.RED, "alerted": Colors.YELLOW, "clo
|
||||
def cron_incidents(args) -> int:
|
||||
"""List (``[--state <s>]``) or ``ack <id>`` durable cron failure incidents.
|
||||
|
||||
The stored error is redacted and truncated at write time, safe for terminal display;
|
||||
acking closes an incident so its failure ping stays silent until the error signature changes.
|
||||
Acking closes an incident so its failure ping stays silent until the error signature changes.
|
||||
"""
|
||||
from cron.incidents import ack_incident, list_incidents
|
||||
action = getattr(args, "incident_action", "list")
|
||||
@@ -356,9 +326,7 @@ def _print_ticker_health(pids: list) -> None:
|
||||
TICKER_INTERVAL_SECONDS,
|
||||
)
|
||||
from cron.scheduler import _is_fd_exhaustion_text as _cron_is_fd_exhaustion_text
|
||||
# ~3 missed ticker iterations (+ slack) before declaring trouble; derived from the shared
|
||||
# interval so the threshold tracks the ticker cadence (= 200s at the 60s default).
|
||||
STALE_AFTER = TICKER_INTERVAL_SECONDS * 3 + 20
|
||||
STALE_AFTER = TICKER_INTERVAL_SECONDS * 3 + 20 # ~3 missed iterations + slack (200s @ 60s)
|
||||
hb_age = get_ticker_heartbeat_age()
|
||||
ok_age = get_ticker_success_age()
|
||||
pid_line = f" PID: {', '.join(map(str, pids))}" if pids else None
|
||||
@@ -369,25 +337,21 @@ def _print_ticker_health(pids: list) -> None:
|
||||
print(pid_line)
|
||||
|
||||
if hb_age is None:
|
||||
# No heartbeat file: ticker never started (non-cron profile, gateway started moments
|
||||
# ago, or a config issue blocking the ticker).
|
||||
# Ticker never started (non-cron profile, gateway just started, or a config issue).
|
||||
_warn("⚠ Gateway is running but the cron ticker has not reported a heartbeat.")
|
||||
print(" Cron jobs will NOT fire until the ticker writes its first heartbeat.\n"
|
||||
" If the gateway just started, wait ~60s and re-run `hermes cron status`.\n"
|
||||
" If heartbeat never appears, restart: hermes gateway restart")
|
||||
elif hb_age > STALE_AFTER:
|
||||
# No heartbeat at all → the ticker thread is gone.
|
||||
elif hb_age > STALE_AFTER: # ticker thread is gone
|
||||
_warn("⚠ Gateway is running but the cron ticker looks STALLED — "
|
||||
f"no heartbeat for {int(hb_age)}s (expected every ~60s).")
|
||||
print(" Cron jobs may NOT be firing. Restart: hermes gateway restart")
|
||||
elif ok_age is not None and ok_age > STALE_AFTER:
|
||||
# Loop alive (fresh heartbeat) but no tick SUCCEEDED in a long time → every tick fails.
|
||||
elif ok_age is not None and ok_age > STALE_AFTER: # loop alive but every tick fails
|
||||
_warn("⚠ Gateway and cron ticker are running, but no tick has "
|
||||
f"succeeded in {int(ok_age)}s — ticks may be failing.")
|
||||
last_error = get_ticker_last_error()
|
||||
if last_error:
|
||||
# Show WHY ticks fail — e.g. a root-rewritten jobs.json (PermissionError) that
|
||||
# silently locked out the ticker's uid, or fd exhaustion (EMFILE).
|
||||
# WHY ticks fail: root-rewritten jobs.json (PermissionError) or fd exhaustion.
|
||||
print(color(f" Last tick error: {last_error}", Colors.RED))
|
||||
if "Permission denied" in last_error:
|
||||
print(color(_PERMISSION_HINT, Colors.YELLOW))
|
||||
@@ -410,9 +374,8 @@ def cron_status():
|
||||
|
||||
provider = _active_cron_provider_name()
|
||||
if provider != "builtin":
|
||||
# An external provider (e.g. Chronos) arms one external one-shot per job, fired by a
|
||||
# NAS-mediated webhook: between fires there is intentionally NO ticker thread and NO
|
||||
# heartbeat file, so the ticker-liveness heuristics would always say "stalled".
|
||||
# External providers fire via webhook: no ticker thread / heartbeat file by design, so
|
||||
# the liveness heuristics would always say "stalled".
|
||||
print(color(f"✓ Cron provider: {provider} — jobs fire via the managed scheduler, "
|
||||
"not the in-process ticker.", Colors.GREEN))
|
||||
print(color(" (No ticker heartbeat is expected for an external provider; "
|
||||
@@ -421,9 +384,8 @@ def cron_status():
|
||||
pids = find_gateway_pids()
|
||||
gateway_alive_via_lock = False
|
||||
if not pids:
|
||||
# The pid scan can transiently miss a live gateway (just after a restart) while
|
||||
# the runtime lock — held for exactly the gateway's lifetime — proves the ticker's
|
||||
# process is alive. Only declare "not running" when both agree.
|
||||
# The pid scan transiently misses a live gateway right after a restart; the runtime
|
||||
# lock proves the process is alive. Declare "not running" only when both agree.
|
||||
with contextlib.suppress(Exception):
|
||||
from gateway.status import get_running_pid, is_gateway_runtime_lock_active
|
||||
gateway_alive_via_lock = is_gateway_runtime_lock_active()
|
||||
@@ -453,13 +415,9 @@ def _print_active_jobs_summary(jobs) -> None:
|
||||
print(f" {len(jobs)} active job(s)")
|
||||
if next_runs:
|
||||
print(f" Next run: {min(next_runs)}")
|
||||
# Missed-run visibility: call out jobs whose LAST dispatch was late or a catch-up so
|
||||
# post-downtime late fires show at status level, not just per-job in `cron list`.
|
||||
late = [
|
||||
j for j in jobs
|
||||
if isinstance(j.get("last_dispatch"), dict)
|
||||
and j["last_dispatch"].get("kind") in ("late", "catch_up")
|
||||
]
|
||||
# Post-downtime late fires show at status level, not just per-job in `cron list`.
|
||||
late = [j for j in jobs if isinstance(j.get("last_dispatch"), dict)
|
||||
and j["last_dispatch"].get("kind") in ("late", "catch_up")]
|
||||
if late:
|
||||
print()
|
||||
print(color(f" ⚠ {len(late)} job(s) last fired late (missed-fire catch-up):",
|
||||
@@ -473,11 +431,7 @@ def _print_active_jobs_summary(jobs) -> None:
|
||||
|
||||
|
||||
def _scripts_dir_for_cron() -> Path:
|
||||
"""Scripts directory used by cron jobs.
|
||||
|
||||
``cron.jobs.CRON_DIR.parent`` rather than a fresh ``get_hermes_home()`` so tests and
|
||||
profile-aware callers that monkeypatch cron storage inspect the same Hermes home.
|
||||
"""
|
||||
"""Scripts dir for cron jobs — via ``CRON_DIR`` so monkeypatched cron storage is honoured."""
|
||||
from cron.jobs import CRON_DIR
|
||||
return CRON_DIR.parent / "scripts"
|
||||
|
||||
@@ -487,12 +441,10 @@ def _script_health_issue(script: str) -> Optional[str]:
|
||||
scripts_dir = _scripts_dir_for_cron().resolve()
|
||||
raw = Path(script).expanduser()
|
||||
path = raw.resolve() if raw.is_absolute() else (scripts_dir / raw).resolve()
|
||||
|
||||
try:
|
||||
path.relative_to(scripts_dir)
|
||||
except ValueError:
|
||||
return f"script resolves outside HERMES_HOME/scripts: {script!r}"
|
||||
|
||||
if not path.exists():
|
||||
return f"script not found: {path}"
|
||||
if not path.is_file():
|
||||
@@ -500,8 +452,7 @@ def _script_health_issue(script: str) -> Optional[str]:
|
||||
return None
|
||||
|
||||
|
||||
# Grace before an overdue ``next_run_at`` is reported: the ticker runs once a minute and a
|
||||
# busy tick can push dispatch a few minutes late; only a next_run_at parked well in the past
|
||||
# A busy tick can push dispatch a few minutes late; only a next_run_at parked well in the past
|
||||
# means the job is silently not firing (ticker dead, gateway down, wedged fire-claim).
|
||||
_OVERDUE_GRACE_SECONDS = 15 * 60
|
||||
|
||||
@@ -524,39 +475,28 @@ def _next_run_overdue_issue(next_run: str) -> Optional[str]:
|
||||
|
||||
def _cron_doctor_issues_for_job(job: Dict[str, Any]) -> List[str]:
|
||||
issues: List[str] = []
|
||||
|
||||
last_status = str(job.get("last_status") or "").strip().lower()
|
||||
# "delivery_failed" means the agent run itself succeeded: the dedicated delivery issue
|
||||
# below reports it (and last_error is None, which would render as "unknown error" here).
|
||||
# "delivery_failed" = the agent run succeeded; the delivery issue below reports it.
|
||||
if last_status and last_status not in {"ok", "delivery_failed"}:
|
||||
err = str(job.get("last_error") or "unknown error").strip()
|
||||
issues.append(f"last run failed: {err}")
|
||||
|
||||
delivery_err = str(job.get("last_delivery_error") or "").strip()
|
||||
if delivery_err:
|
||||
issues.append(f"last run failed: {str(job.get('last_error') or 'unknown error').strip()}")
|
||||
if delivery_err := str(job.get("last_delivery_error") or "").strip():
|
||||
issues.append(f"last delivery failed: {delivery_err}")
|
||||
|
||||
unverified = job.get("last_delivery_unverified")
|
||||
if unverified:
|
||||
if unverified := job.get("last_delivery_unverified"):
|
||||
issues.append("last delivery unverified (adapter acked without evidence): "
|
||||
+ _unverified_targets(unverified))
|
||||
|
||||
if job.get("enabled", True) and job.get("state") not in {"paused", "completed"}:
|
||||
next_run = str(job.get("next_run_at") or "").strip()
|
||||
issue = _next_run_overdue_issue(next_run) if next_run else "active job has no next_run_at"
|
||||
if issue:
|
||||
issues.append(issue)
|
||||
|
||||
script = str(job.get("script") or "").strip()
|
||||
if job.get("no_agent") and not script:
|
||||
issues.append("no-agent job has no script")
|
||||
if script and (script_issue := _script_health_issue(script)):
|
||||
issues.append(script_issue)
|
||||
|
||||
workdir = str(job.get("workdir") or "").strip()
|
||||
if workdir and not Path(workdir).expanduser().exists():
|
||||
issues.append(f"workdir not found: {workdir}")
|
||||
|
||||
return issues
|
||||
|
||||
|
||||
@@ -565,13 +505,11 @@ def cron_doctor() -> int:
|
||||
from cron.jobs import list_jobs
|
||||
jobs = list_jobs(include_disabled=False)
|
||||
findings = [(job, issues) for job in jobs if (issues := _cron_doctor_issues_for_job(job))]
|
||||
|
||||
if not findings:
|
||||
print(color("✓ Cron doctor found no issues", Colors.GREEN))
|
||||
note = f" Checked {len(jobs)} active job(s)." if jobs else " No active jobs configured."
|
||||
print(color(note, Colors.DIM))
|
||||
return 0
|
||||
|
||||
issue_count = sum(len(issues) for _, issues in findings)
|
||||
print(color(f"Cron doctor found {issue_count} issue(s) across {len(findings)} job(s):", Colors.YELLOW))
|
||||
print()
|
||||
@@ -614,9 +552,8 @@ def _print_job_details(job_data: Dict[str, Any]) -> None:
|
||||
|
||||
|
||||
def cron_create(args):
|
||||
# The gateway-lifecycle guard lives in cron.jobs.create_job so it fires on every
|
||||
# job-creation path (CLI AND the agent's `cronjob` tool); a block surfaces as
|
||||
# result["error"] and is printed in red below.
|
||||
# The gateway-lifecycle guard lives in cron.jobs.create_job (every creation path); a block
|
||||
# surfaces as result["error"].
|
||||
result = _cron_api(
|
||||
action="create", schedule=args.schedule, prompt=args.prompt,
|
||||
skill=getattr(args, "skill", None),
|
||||
@@ -648,8 +585,7 @@ def cron_edit(args):
|
||||
if not job:
|
||||
print(color(f"Job not found: {args.job_id}", Colors.RED))
|
||||
return 1
|
||||
|
||||
existing_skills = list(job.get("skills") or ([] if not job.get("skill") else [job.get("skill")]))
|
||||
existing_skills = list(job.get("skills") or ([job["skill"]] if job.get("skill") else []))
|
||||
replacement_skills = _normalize_skills(getattr(args, "skill", None), getattr(args, "skills", None))
|
||||
add_skills = _normalize_skills(None, getattr(args, "add_skills", None)) or []
|
||||
remove_skills = set(_normalize_skills(None, getattr(args, "remove_skills", None)) or [])
|
||||
@@ -662,7 +598,6 @@ def cron_edit(args):
|
||||
elif add_skills or remove_skills:
|
||||
final_skills = [skill for skill in existing_skills if skill not in remove_skills]
|
||||
final_skills += [skill for skill in add_skills if skill not in final_skills]
|
||||
|
||||
result = _cron_api(action="update", job_id=args.job_id,
|
||||
schedule=getattr(args, "schedule", None),
|
||||
prompt=getattr(args, "prompt", None), skills=final_skills,
|
||||
@@ -682,13 +617,10 @@ def cron_edit(args):
|
||||
def _job_action(action: str, job_id: str, success_verb: str) -> int:
|
||||
_stateless_token = None
|
||||
if action == "run":
|
||||
# One-shot CLI: the process exits as soon as the command returns, so a
|
||||
# background-dispatched run (daemon thread of THIS process — triggered when the CLI
|
||||
# inherits a gateway/desktop session env, HERMES_SESSION_KEY) would be orphaned
|
||||
# mid-LLM-call, leaving the execution row stuck 'claimed'. Declare the channel
|
||||
# stateless so ``async_delivery_supported()`` gates it off and the run executes
|
||||
# synchronously. Scoped to this call (token reset in ``finally``) so in-process
|
||||
# callers (tests, embedding apps) are not tainted.
|
||||
# One-shot CLI: a background-dispatched run (daemon thread, triggered when the CLI
|
||||
# inherits HERMES_SESSION_KEY) would be orphaned mid-LLM-call, leaving the execution row
|
||||
# stuck 'claimed'. Declaring the channel stateless forces a synchronous run; scoped to
|
||||
# this call so in-process callers (tests, embedding apps) are not tainted.
|
||||
with contextlib.suppress(Exception):
|
||||
from gateway.session_context import _SESSION_ASYNC_DELIVERY
|
||||
_stateless_token = _SESSION_ASYNC_DELIVERY.set(False)
|
||||
@@ -712,9 +644,8 @@ def _job_action(action: str, job_id: str, success_verb: str) -> int:
|
||||
def _run_outcome(job: Dict[str, Any]) -> str:
|
||||
"""One-line verdict for a manual run.
|
||||
|
||||
A run may be dispatched to the gateway daemon's background delegation worker
|
||||
(execution_mode="background" and/or a delegation_id) and keeps running AFTER this CLI
|
||||
exits — a terminal success/failure verdict would be a lie, so report the dispatch.
|
||||
A background-dispatched run (execution_mode="background" / delegation_id) keeps running
|
||||
after this CLI exits, so report the dispatch rather than a success/failure verdict.
|
||||
"""
|
||||
if job.get("delegation_id"):
|
||||
return f"Running in background (delegation {job['delegation_id']})."
|
||||
@@ -749,22 +680,19 @@ def cron_resume(args) -> int:
|
||||
|
||||
|
||||
def cron_notepad(args) -> int:
|
||||
"""Handle ``hermes cron notepad <job_id> [get|set|delete|list]``.
|
||||
"""Handle ``hermes cron notepad <job_id> [get|set|delete|list]`` (per-job durable KV).
|
||||
|
||||
Write path for the per-job durable KV scratchpad (``cron/notepad.py``): a running cron
|
||||
agent updates its own notepad via its terminal tool, and the scheduler injects non-empty
|
||||
notepads into the job prompt on each run.
|
||||
A running cron agent updates its own notepad via its terminal tool; the scheduler injects
|
||||
non-empty notepads into the job prompt on each run.
|
||||
"""
|
||||
from cron import notepad
|
||||
job_id = str(getattr(args, "job_id", "") or "")
|
||||
action = getattr(args, "notepad_action", None) or "list"
|
||||
key = getattr(args, "key", None)
|
||||
value = getattr(args, "value", None)
|
||||
|
||||
if not job_id:
|
||||
print(color("A job ID is required.", Colors.RED))
|
||||
return 1
|
||||
|
||||
try:
|
||||
if action not in ("set", "get", "delete"): # list (default)
|
||||
notes = notepad.list_notes(job_id)
|
||||
@@ -797,8 +725,7 @@ def cron_notepad(args) -> int:
|
||||
return 1
|
||||
|
||||
|
||||
# Subcommand -> handler. Late-bound lambdas so module-level monkeypatching of the underlying
|
||||
# functions keeps working. cron_list/cron_status/cron_runs always return None -> exit 0.
|
||||
# Late-bound lambdas keep module-level monkeypatching working; list/status/runs return None -> 0.
|
||||
_CRON_SUBCOMMANDS = {
|
||||
"list": lambda a: cron_list(getattr(a, "all", False)) or 0,
|
||||
"status": lambda a: cron_status() or 0,
|
||||
@@ -825,7 +752,6 @@ def cron_command(args):
|
||||
handler = _CRON_SUBCOMMANDS.get("list" if subcmd is None else subcmd)
|
||||
if handler is not None:
|
||||
return handler(args)
|
||||
|
||||
print(f"Unknown cron command: {subcmd}\n"
|
||||
"Usage: hermes cron [list|create|edit|pause|resume|run|remove|status|runs|doctor|tick]")
|
||||
sys.exit(1)
|
||||
|
||||
@@ -10,8 +10,8 @@ import subprocess
|
||||
import sys
|
||||
from pathlib import Path
|
||||
|
||||
# Cmdline substrings identifying the long-lived server; ``hermes serve`` is the same server
|
||||
# under the headless name the desktop app spawns, reaped on update for the same reason.
|
||||
# Cmdline substrings identifying the long-lived server (``serve`` = the headless name Desktop
|
||||
# spawns; reaped on update for the same reason).
|
||||
_DASHBOARD_PATTERNS = tuple(
|
||||
f"{launcher} {cmd}"
|
||||
for cmd in ("dashboard", "serve")
|
||||
@@ -41,9 +41,9 @@ def _iter_process_table() -> list[tuple[int, str]]:
|
||||
"""``(pid, cmdline)`` for every process, via wmic (Windows) or ps. Raises on scan failure."""
|
||||
rows: list[tuple[int, str]] = []
|
||||
if sys.platform == "win32":
|
||||
# errors="ignore": wmic may emit the system code page (a decode error leaves stdout
|
||||
# None). bounded_probe_run, not run(): run()'s post-timeout cleanup joins pipe readers
|
||||
# unbounded and a conhost descendant holding duplicated handles wedges it forever.
|
||||
# errors="ignore": wmic may emit the system code page. bounded_probe_run, not run():
|
||||
# run()'s post-timeout cleanup joins pipe readers unbounded and a conhost descendant
|
||||
# holding duplicated handles wedges it forever.
|
||||
from hermes_cli._subprocess_compat import bounded_probe_run
|
||||
result = bounded_probe_run(
|
||||
["wmic", "process", "get", "ProcessId,CommandLine", "/FORMAT:LIST"],
|
||||
@@ -59,8 +59,7 @@ def _iter_process_table() -> list[tuple[int, str]]:
|
||||
elif line.startswith("ProcessId="):
|
||||
_append_row(rows, line[len("ProcessId=") :], current_cmd)
|
||||
return rows
|
||||
# ps (not `pgrep -f "hermes.*dashboard"`) keeps us consistent with gateway._scan_gateway_pids
|
||||
# and avoids a greedy regex matching unrelated cmdlines that merely contain both words.
|
||||
# ps, not `pgrep -f "hermes.*dashboard"` (greedy regex; consistent with gateway pid scan).
|
||||
result = subprocess.run(["ps", "-A", "-o", "pid=,command="], timeout=10, **_PS_RUN_KWARGS)
|
||||
if result.returncode == 0:
|
||||
for line in getattr(result, "stdout", "").split("\n"):
|
||||
@@ -73,34 +72,26 @@ def _iter_process_table() -> list[tuple[int, str]]:
|
||||
def _scan_dashboard_processes(*, exclude_pids: set[int] | None = None) -> list[tuple[int, str]]:
|
||||
"""``(pid, cmdline)`` of running ``dashboard``/``serve`` processes; empty on any scan error.
|
||||
|
||||
A forgotten dashboard keeps the old Python backend in memory against the new JS bundle
|
||||
after ``hermes update`` (new auth headers → every API call 401s). *exclude_pids* must never
|
||||
be returned: Desktop marks the backend it manages via ``HERMES_DESKTOP_CHILD_PID``.
|
||||
A forgotten dashboard keeps the old Python backend against the new JS bundle after
|
||||
``hermes update`` (every API call 401s). *exclude_pids* (Desktop's HERMES_DESKTOP_CHILD_PID
|
||||
backends) are never returned.
|
||||
"""
|
||||
self_pid = os.getpid()
|
||||
skip = {self_pid, *(exclude_pids or ())}
|
||||
skip = {os.getpid(), *(exclude_pids or ())}
|
||||
try:
|
||||
found = [
|
||||
(pid, cmd) for pid, cmd in _iter_process_table()
|
||||
if pid not in skip and any(p in cmd for p in _DASHBOARD_PATTERNS)
|
||||
]
|
||||
found = [(pid, cmd) for pid, cmd in _iter_process_table()
|
||||
if pid not in skip and any(p in cmd for p in _DASHBOARD_PATTERNS)]
|
||||
except (FileNotFoundError, subprocess.TimeoutExpired, OSError):
|
||||
return []
|
||||
|
||||
# Spawn-ledger augmentation: substring patterns miss profiled launches (`hermes --profile
|
||||
# p serve`); the ledger holds live-verified (pid, create_time) — positive identity.
|
||||
try:
|
||||
# Spawn-ledger augmentation: substring patterns miss profiled launches (`hermes --profile p
|
||||
# serve`); the ledger holds live-verified pids. Unavailable ledger → scan-only.
|
||||
with contextlib.suppress(Exception):
|
||||
from hermes_cli.process_identity import ledger_entries
|
||||
seen = {pid for pid, _ in found} | skip
|
||||
for entry in ledger_entries():
|
||||
pid = entry.get("pid")
|
||||
if entry.get("purpose") not in ("serve", "dashboard") or not isinstance(pid, int):
|
||||
continue
|
||||
if pid not in seen:
|
||||
if (entry.get("purpose") in ("serve", "dashboard") and isinstance(pid, int)
|
||||
and pid not in seen):
|
||||
found.append((pid, str(entry.get("argv") or "")))
|
||||
except Exception:
|
||||
pass # ledger unavailable → scan-only behavior
|
||||
|
||||
return found
|
||||
|
||||
|
||||
@@ -121,10 +112,7 @@ def _hermes_home_for_pid(pid: int) -> str | None:
|
||||
|
||||
|
||||
def _dashboard_subcommand_index(argv: list[str]) -> int | None:
|
||||
for i, tok in enumerate(argv):
|
||||
if tok in ("serve", "dashboard"):
|
||||
return i
|
||||
return None
|
||||
return next((i for i, tok in enumerate(argv) if tok in ("serve", "dashboard")), None)
|
||||
|
||||
|
||||
def _profile_flag_value(argv: list[str]) -> str | None:
|
||||
@@ -138,19 +126,13 @@ def _profile_flag_value(argv: list[str]) -> str | None:
|
||||
|
||||
|
||||
def _is_ephemeral_port_zero_backend(argv: list[str]) -> bool:
|
||||
"""True for Desktop-style ``serve|dashboard --port 0`` backends.
|
||||
|
||||
Owned by Hermes Desktop (or PPID-1 orphans of a prior respawn); replaying them after
|
||||
``hermes update`` multiplies listening backends because ``--port 0`` binds a fresh port.
|
||||
"""
|
||||
"""True for Desktop-style ``serve|dashboard --port 0`` backends — replaying them after
|
||||
``hermes update`` multiplies listening backends because ``--port 0`` binds a fresh port."""
|
||||
if _dashboard_subcommand_index(argv) is None:
|
||||
return False
|
||||
for i, tok in enumerate(argv):
|
||||
if tok == "--port" and i + 1 < len(argv) and str(argv[i + 1]) == "0":
|
||||
return True
|
||||
if tok.startswith("--port=") and tok.split("=", 1)[1].strip() == "0":
|
||||
return True
|
||||
return False
|
||||
return any((tok == "--port" and i + 1 < len(argv) and str(argv[i + 1]) == "0")
|
||||
or (tok.startswith("--port=") and tok.split("=", 1)[1].strip() == "0")
|
||||
for i, tok in enumerate(argv))
|
||||
|
||||
|
||||
def _normalize_dashboard_cmdline(argv: list[str]) -> tuple[str, ...]:
|
||||
@@ -187,9 +169,8 @@ def _normalized_home_for_compare(home: str) -> str:
|
||||
def _profile_key_for_respawn(argv: list[str], hermes_home: str | None = None) -> str:
|
||||
"""Stable owner key: ``HERMES_HOME`` when known, else ``--profile`` / ``-p``.
|
||||
|
||||
``HERMES_HOME`` ending in ``profiles/<name>`` → ``profile:<name>`` so it shares a cap with
|
||||
an explicit ``--profile``; other homes keep a resolved ``home:`` key so unrelated installs
|
||||
never collapse together.
|
||||
A home ending in ``profiles/<name>`` → ``profile:<name>`` (shares a cap with an explicit
|
||||
``--profile``); other homes keep a ``home:`` key so unrelated installs never collapse.
|
||||
"""
|
||||
if hermes_home:
|
||||
parts = _resolved_home(hermes_home).parts
|
||||
@@ -204,12 +185,11 @@ def _filter_dashboard_respawn_candidates(
|
||||
) -> list[list[str]]:
|
||||
"""Select which killed manual backends ``(pid, argv, hermes_home)`` to respawn after update.
|
||||
|
||||
Rules: never resurrect Desktop ``--port 0`` backends (Desktop owns them); never replay a
|
||||
backend from a **foreign** ``HERMES_HOME`` (the respawn is argv-only, so it would come back
|
||||
on the updating install's home and steal the foreign install's fixed port → its supervisor
|
||||
crash-loops on ``EADDRINUSE``; unreadable ``None`` stays eligible); dedupe by normalized
|
||||
cmdline; at most one backend per profile / home. PPID-1 is NOT skipped: a prior respawn
|
||||
detaches with ``start_new_session=True``, so fixed-port manual backends sit under init.
|
||||
Rules: never resurrect Desktop ``--port 0`` backends; never replay a backend from a
|
||||
**foreign** ``HERMES_HOME`` (the argv-only respawn would come back on this install's home
|
||||
and steal the foreign install's fixed port → EADDRINUSE crash-loop; unreadable ``None``
|
||||
stays eligible); dedupe by normalized cmdline; one backend per profile / home. PPID-1 is
|
||||
NOT skipped: a prior respawn detaches, so fixed-port manual backends sit under init.
|
||||
"""
|
||||
if own_home is None:
|
||||
try:
|
||||
@@ -218,11 +198,9 @@ def _filter_dashboard_respawn_candidates(
|
||||
except Exception:
|
||||
own_home = ""
|
||||
own_key = _normalized_home_for_compare(own_home) if own_home else ""
|
||||
|
||||
selected: list[list[str]] = []
|
||||
seen_cmdlines: set[tuple[str, ...]] = set()
|
||||
seen_profiles: set[str] = set()
|
||||
|
||||
for _pid, argv, hermes_home in candidates:
|
||||
if not argv or _is_ephemeral_port_zero_backend(argv):
|
||||
continue
|
||||
@@ -235,7 +213,6 @@ def _filter_dashboard_respawn_candidates(
|
||||
seen_cmdlines.add(norm)
|
||||
seen_profiles.add(profile_key)
|
||||
selected.append(list(argv))
|
||||
|
||||
return selected
|
||||
|
||||
|
||||
@@ -252,8 +229,7 @@ def _kill_pids_windows(pids: list[int], killed: list[int], failed: list[tuple[in
|
||||
"""``taskkill /F`` each PID after re-verifying its identity."""
|
||||
from gateway.status import get_process_start_time
|
||||
from hermes_cli._subprocess_compat import pid_is_hermes, windows_hide_flags
|
||||
# Capture identity immediately after discovery: a PID reused before the destructive
|
||||
# action fails the start-time check.
|
||||
# Identity captured right after discovery: a PID reused before the kill fails the check.
|
||||
pid_start_times = {pid: get_process_start_time(pid) for pid in pids}
|
||||
for pid in pids:
|
||||
try:
|
||||
@@ -294,16 +270,13 @@ def _kill_pids_posix(pids: list[int], killed: list[int], failed: list[tuple[int,
|
||||
|
||||
for pid in pids:
|
||||
_send(pid, _signal.SIGTERM)
|
||||
|
||||
deadline = _time.monotonic() + 3.0
|
||||
pending = [p for p in pids if p not in killed and p not in {f[0] for f in failed}]
|
||||
while pending and _time.monotonic() < deadline:
|
||||
_time.sleep(0.1)
|
||||
# os.kill(pid, 0) is NOT a no-op on Windows; use the portable check.
|
||||
alive = [p for p in pending if _pid_exists(p)]
|
||||
alive = [p for p in pending if _pid_exists(p)] # os.kill(pid, 0) breaks on Windows
|
||||
killed.extend(p for p in pending if p not in alive)
|
||||
pending = alive
|
||||
|
||||
for pid in pending:
|
||||
_send(pid, _signal.SIGKILL)
|
||||
|
||||
@@ -320,26 +293,20 @@ def _kill_stale_dashboard_processes(
|
||||
``.service`` suffix) are left untouched, not killed twice.
|
||||
"""
|
||||
if restart_managed and _m()._restart_managed_dashboard_service(reason):
|
||||
# The dashboard unit is handled but every OTHER backend is not (a host may also run
|
||||
# hermes-serve.service hosting tui_gateway): record the unit as handled (the filter
|
||||
# below drops PIDs it owns) and keep going.
|
||||
# The dashboard unit is handled but other backends (e.g. hermes-serve.service) are not:
|
||||
# mark the unit handled so the filter below drops its PIDs, and keep going.
|
||||
_dash_unit = getattr(_m(), "_DASHBOARD_SYSTEMD_UNIT", "hermes-dashboard.service")
|
||||
already_restarted_units = set(already_restarted_units or ()) | {
|
||||
str(_dash_unit).removesuffix(".service")}
|
||||
|
||||
exclude = _exclude_pids_from_env()
|
||||
if restart_managed:
|
||||
# An SSH-owned backend belongs to an attached Desktop client even when the updater
|
||||
# runs from an unrelated shell; killing it strands that client's fixed SSH
|
||||
# port-forward. Same ownership records as the reaper.
|
||||
# An SSH-owned backend belongs to an attached Desktop client; killing it strands that
|
||||
# client's fixed SSH port-forward. Same ownership records as the reaper.
|
||||
exclude |= _lock_owned_serve_pids()
|
||||
|
||||
pids = _m()._find_stale_dashboard_pids(exclude_pids=exclude or None)
|
||||
if not pids:
|
||||
return _empty_result()
|
||||
|
||||
# Snapshot systemd cgroup/unit and argv BEFORE killing (the cgroup disappears with the
|
||||
# process). Linux + update path only.
|
||||
# Snapshot systemd unit/cgroup and argv BEFORE killing (the cgroup dies with the process).
|
||||
pid_cgroup: dict[int, str | None] = {}
|
||||
pid_service: dict[int, str | None] = {}
|
||||
pid_cmdline: dict[int, list[str]] = {}
|
||||
@@ -349,36 +316,28 @@ def _kill_stale_dashboard_processes(
|
||||
pid_cgroup[pid] = _m()._get_pid_cgroup_path(pid)
|
||||
pid_service[pid] = _m()._get_systemd_service_for_pid(pid)
|
||||
if not pid_service[pid] and (cmdline := _m()._dashboard_cmdline_for_pid(pid)):
|
||||
# Manual process: keep exact argv + HERMES_HOME for the post-update respawn
|
||||
# and its per-profile cap.
|
||||
# Manual process: exact argv + HERMES_HOME for the respawn and its profile cap.
|
||||
pid_cmdline[pid] = cmdline
|
||||
pid_home[pid] = _hermes_home_for_pid(pid)
|
||||
if already_restarted_units:
|
||||
pids = [
|
||||
pid for pid in pids
|
||||
if (pid_service.get(pid) or "").removesuffix(".service") not in already_restarted_units
|
||||
]
|
||||
pids = [pid for pid in pids if (pid_service.get(pid) or "").removesuffix(".service")
|
||||
not in already_restarted_units]
|
||||
if not pids:
|
||||
return _empty_result()
|
||||
|
||||
print(f"\n⟲ Stopping {len(pids)} dashboard process(es) ({reason})")
|
||||
|
||||
killed: list[int] = []
|
||||
failed: list[tuple[int, str]] = []
|
||||
(_kill_pids_windows if sys.platform == "win32" else _kill_pids_posix)(pids, killed, failed)
|
||||
|
||||
for pid in killed:
|
||||
print(f" ✓ stopped PID {pid}")
|
||||
for pid, err_msg in failed:
|
||||
print(f" ✗ failed to stop PID {pid}: {err_msg}")
|
||||
|
||||
if killed and restart_managed:
|
||||
unrecovered = _restart_killed_backends(killed, pid_service, pid_cgroup, pid_cmdline, pid_home)
|
||||
else:
|
||||
unrecovered = list(killed)
|
||||
if killed:
|
||||
print(" Restart the dashboard when you're ready:\n hermes dashboard --port <port>")
|
||||
|
||||
return {"matched": list(pids), "killed": list(killed), "failed": list(failed),
|
||||
"unrecovered": list(unrecovered)}
|
||||
|
||||
@@ -387,11 +346,8 @@ def _restart_killed_backends(
|
||||
killed: list[int], pid_service: dict[int, str | None], pid_cgroup: dict[int, str | None],
|
||||
pid_cmdline: dict[int, list[str]], pid_home: dict[int, str | None],
|
||||
) -> list[int]:
|
||||
"""Update path: restart systemd-owned units, respawn manual argv. Returns PIDs not brought back.
|
||||
|
||||
Respawns are detached, headless, logged to logs/dashboard-restart.log; Desktop ``--port 0``
|
||||
backends are filtered out and duplicates collapse to one per profile.
|
||||
"""
|
||||
"""Update path: restart systemd units, respawn manual argv (detached, headless, logged to
|
||||
logs/dashboard-restart.log; one per profile, no ``--port 0``). Returns PIDs not brought back."""
|
||||
unrecovered: list[int] = []
|
||||
failed_restarts: list[tuple[str, str]] = []
|
||||
seen_services: set[str] = set()
|
||||
@@ -411,15 +367,12 @@ def _restart_killed_backends(
|
||||
respawn_candidates.append((pid, pid_cmdline[pid], pid_home.get(pid)))
|
||||
else:
|
||||
unrecovered.append(pid)
|
||||
|
||||
for svc, err in failed_restarts:
|
||||
print(f" ⚠ {svc}: {err}")
|
||||
|
||||
respawn_cmds = _filter_dashboard_respawn_candidates(respawn_candidates)
|
||||
failed_cmds = _m()._respawn_dashboard_processes(respawn_cmds) if respawn_cmds else None
|
||||
if failed_cmds:
|
||||
unrecovered.extend(p for p in killed if pid_cmdline.get(p) in failed_cmds)
|
||||
|
||||
if failed_restarts or unrecovered:
|
||||
print(" Restart anything not auto-restarted when you're ready:\n hermes dashboard --port <port>")
|
||||
return unrecovered
|
||||
@@ -440,32 +393,27 @@ def _detect_concurrent_hermes_instances(
|
||||
|
||||
Windows blocks DELETE/REPLACE on a running .exe, so a Desktop-spawned ``hermes.EXE`` makes
|
||||
the update's quarantine rename fail with ``[WinError 32]``. Excludes our PID and every
|
||||
*shim* ancestor (the setuptools launcher is a separate native process from the
|
||||
``python.exe`` it loads); ``proc.parents()`` at once because a per-hop loop bailed on the
|
||||
first AccessDenied. Empty off-Windows / without psutil. Never raises.
|
||||
*shim* ancestor (the setuptools launcher is a separate native process from its
|
||||
``python.exe``); ``proc.parents()`` at once because a per-hop loop bailed on the first
|
||||
AccessDenied. Empty off-Windows / without psutil. Never raises.
|
||||
"""
|
||||
if not _m()._is_windows():
|
||||
return []
|
||||
|
||||
try:
|
||||
import psutil
|
||||
except Exception:
|
||||
return []
|
||||
|
||||
shim_paths = {_norm_exe(shim) for shim in _m()._hermes_exe_shims(scripts_dir)}
|
||||
if not shim_paths:
|
||||
return []
|
||||
|
||||
seed = int(exclude_pid) if exclude_pid is not None else os.getpid()
|
||||
exclude_pids: set[int] = {seed}
|
||||
# Broad ``except Exception``: psutil may be partially stubbed in tests.
|
||||
with contextlib.suppress(Exception):
|
||||
with contextlib.suppress(Exception): # psutil may be partially stubbed in tests
|
||||
for ancestor in psutil.Process(seed).parents():
|
||||
with contextlib.suppress(Exception):
|
||||
anc_exe = ancestor.exe()
|
||||
if anc_exe and _norm_exe(anc_exe) in shim_paths:
|
||||
exclude_pids.add(int(ancestor.pid))
|
||||
|
||||
matches: list[tuple[int, str]] = []
|
||||
try:
|
||||
proc_iter = psutil.process_iter(["pid", "exe", "name"])
|
||||
@@ -479,7 +427,6 @@ def _detect_concurrent_hermes_instances(
|
||||
pid, exe = info.get("pid"), info.get("exe")
|
||||
if exe and pid is not None and pid not in exclude_pids and _norm_exe(exe) in shim_paths:
|
||||
matches.append((int(pid), str(info.get("name") or Path(exe).name)))
|
||||
|
||||
return matches
|
||||
|
||||
|
||||
@@ -492,12 +439,9 @@ def _is_desktop_local_serve_cmdline(command: str) -> bool:
|
||||
cmd = command.lower()
|
||||
if "serve" not in cmd or ("hermes" not in cmd and "hermes_cli" not in cmd):
|
||||
return False
|
||||
has_loopback = any(
|
||||
tok in cmd
|
||||
for tok in ("--host 127.0.0.1", "--host=127.0.0.1", "--host localhost", "--host=localhost")
|
||||
)
|
||||
has_ephemeral = "--port 0" in cmd or "--port=0" in cmd
|
||||
return has_loopback and has_ephemeral
|
||||
has_loopback = any(tok in cmd for tok in (
|
||||
"--host 127.0.0.1", "--host=127.0.0.1", "--host localhost", "--host=localhost"))
|
||||
return has_loopback and ("--port 0" in cmd or "--port=0" in cmd)
|
||||
|
||||
|
||||
def _process_ppid(pid: int) -> int | None:
|
||||
@@ -513,12 +457,10 @@ def _process_ppid(pid: int) -> int | None:
|
||||
return None
|
||||
|
||||
|
||||
# --- SSH remote-backend lock ownership -------------------------------------
|
||||
# ``backend.lock.json`` is written by the Desktop SSH runtime on the *remote* host for every
|
||||
# ``hermes serve`` it spawns (apps/desktop/electron/remote-lifecycle.ts). Such a backend is
|
||||
# legitimate even with no parent here (sshd exited → ppid 1); the reap must NEVER kill a PID a
|
||||
# valid lock claims — that once killed a production backend. Schema mirrors the writer; a
|
||||
# mismatched record is ignored (the reap only *spares*).
|
||||
# SSH remote-backend lock ownership: ``backend.lock.json`` is written by the Desktop SSH runtime
|
||||
# (apps/desktop/electron/remote-lifecycle.ts) for every ``hermes serve`` it spawns. Such a backend
|
||||
# is legitimate even at ppid 1 (sshd exited); the reap must NEVER kill a PID a valid lock claims
|
||||
# — that once killed a production backend. Schema mirrors the writer; mismatches are ignored.
|
||||
_LOCKFILE_SCHEMA_VERSION = 2
|
||||
_PROTOCOL_VERSION = 1
|
||||
_REMOTE_LOCK_SUBDIR = "desktop-ssh"
|
||||
@@ -547,24 +489,20 @@ def _valid_lockfile_payload(parsed: object, ownership_id: str) -> bool:
|
||||
):
|
||||
return False
|
||||
pid, port = parsed.get("pid"), parsed.get("port")
|
||||
if not (isinstance(pid, int) and 0 < pid <= 4194304):
|
||||
return False
|
||||
if not (isinstance(port, int) and 0 <= port <= 65535):
|
||||
if not (isinstance(pid, int) and 0 < pid <= 4194304 and isinstance(port, int)
|
||||
and 0 <= port <= 65535):
|
||||
return False
|
||||
# String fields must be present and bounded (the writer enforces <=1024).
|
||||
if any(not isinstance(parsed.get(f), str) or len(parsed[f]) > 1024
|
||||
for f in ("profile", "hermesPath", "hermesHome", "logPath", "startedAt")):
|
||||
return False
|
||||
# logPath is ``{lock_root}/{ownershipId}/{spawnNonce}.log``; only the suffix is checked so a
|
||||
# relocated HERMES_HOME can't falsely reject a legitimate backend (= re-introduce the kill).
|
||||
# Suffix-only check of logPath so a relocated HERMES_HOME can't reject a legitimate backend.
|
||||
return parsed["logPath"].endswith(f"/{ownership_id}/{parsed['spawnNonce']}.log")
|
||||
|
||||
|
||||
def _lock_owned_serve_pids(base_dir: Path | None = None) -> set[int]:
|
||||
"""PIDs claimed by valid ``{hermes_home}/desktop-ssh/<ownershipId>/backend.lock.json`` records.
|
||||
|
||||
Best-effort: a bad record contributes no PID; never raises.
|
||||
"""
|
||||
"""PIDs claimed by valid ``{hermes_home}/desktop-ssh/<ownershipId>/backend.lock.json`` records
|
||||
(best-effort: a bad record contributes no PID; never raises)."""
|
||||
import json
|
||||
root = base_dir if base_dir is not None else _hermes_home_dir() / _REMOTE_LOCK_SUBDIR
|
||||
owned: set[int] = set()
|
||||
@@ -575,8 +513,7 @@ def _lock_owned_serve_pids(base_dir: Path | None = None) -> set[int]:
|
||||
for entry in entries:
|
||||
ownership_id = entry.name
|
||||
lock_path = entry / "backend.lock.json"
|
||||
try:
|
||||
# Mirror validateOwnershipId(): exactly 32 lowercase hex chars.
|
||||
try: # validateOwnershipId(): exactly 32 lowercase hex chars
|
||||
if not entry.is_dir() or not _is_hex(ownership_id, 32) or not lock_path.is_file():
|
||||
continue
|
||||
data = lock_path.read_bytes()
|
||||
@@ -590,8 +527,7 @@ def _lock_owned_serve_pids(base_dir: Path | None = None) -> set[int]:
|
||||
return owned
|
||||
|
||||
|
||||
# Grace window before an orphaned-looking backend may be reaped: covers the gap between
|
||||
# process start and the Desktop client writing backend.lock.json.
|
||||
# Covers the gap between process start and the Desktop client writing backend.lock.json.
|
||||
_REAP_MIN_AGE_SECONDS = 180.0
|
||||
|
||||
|
||||
@@ -624,9 +560,7 @@ def _reap_orphaned_desktop_local_serves(
|
||||
sleep_fn = sleep_fn or _time.sleep
|
||||
lock_owned_pids_fn = lock_owned_pids_fn or _lock_owned_serve_pids
|
||||
process_age_seconds_fn = process_age_seconds_fn or _process_age_seconds
|
||||
|
||||
if sys.platform == "win32":
|
||||
# Windows desktop uses taskkill tree teardown; orphan scan is POSIX.
|
||||
if sys.platform == "win32": # Windows desktop uses taskkill tree teardown
|
||||
return _empty_result()
|
||||
|
||||
def _owned_pids() -> set[int]:
|
||||
@@ -635,34 +569,25 @@ def _reap_orphaned_desktop_local_serves(
|
||||
except Exception:
|
||||
return set() # never let lock scanning block or widen the reap
|
||||
|
||||
exclude = _exclude_pids_from_env() | {os.getpid()} | _owned_pids()
|
||||
try:
|
||||
exclude.add(os.getppid()) # the desktop / sshd wrapper
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
try:
|
||||
scanned = _scan_dashboard_processes(exclude_pids=exclude)
|
||||
except Exception:
|
||||
return _empty_result()
|
||||
|
||||
# Re-read ownership: a lock may have been written between scan and now.
|
||||
owned_now = _owned_pids()
|
||||
|
||||
def _is_stale_orphan(pid: int) -> bool:
|
||||
try: # never let a liveness probe failure widen the reap
|
||||
return process_age_seconds_fn(pid) >= _REAP_MIN_AGE_SECONDS
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
matched = [
|
||||
pid for pid, cmd in scanned
|
||||
if _is_desktop_local_serve_cmdline(cmd) and pid not in owned_now
|
||||
and _process_ppid(pid) in (0, 1) and _is_stale_orphan(pid)
|
||||
]
|
||||
exclude = _exclude_pids_from_env() | {os.getpid()} | _owned_pids()
|
||||
with contextlib.suppress(Exception):
|
||||
exclude.add(os.getppid()) # the desktop / sshd wrapper
|
||||
try:
|
||||
scanned = _scan_dashboard_processes(exclude_pids=exclude)
|
||||
except Exception:
|
||||
return _empty_result()
|
||||
owned_now = _owned_pids() # re-read: a lock may have been written since the scan
|
||||
matched = [pid for pid, cmd in scanned
|
||||
if _is_desktop_local_serve_cmdline(cmd) and pid not in owned_now
|
||||
and _process_ppid(pid) in (0, 1) and _is_stale_orphan(pid)]
|
||||
if not matched:
|
||||
return _empty_result()
|
||||
|
||||
killed: list[int] = []
|
||||
failed: list[int] = []
|
||||
for pid in matched:
|
||||
@@ -672,9 +597,7 @@ def _reap_orphaned_desktop_local_serves(
|
||||
continue
|
||||
except OSError:
|
||||
failed.append(pid)
|
||||
|
||||
# Brief grace, then SIGKILL survivors. psutil.pid_exists rather than os.kill(pid, 0),
|
||||
# which is a Windows footgun the linter blocks everywhere.
|
||||
# Brief grace, then SIGKILL survivors (psutil.pid_exists: os.kill(pid, 0) is a Windows footgun).
|
||||
sleep_fn(1.5)
|
||||
import psutil
|
||||
for pid in matched:
|
||||
@@ -688,7 +611,6 @@ def _reap_orphaned_desktop_local_serves(
|
||||
killed.append(pid)
|
||||
except OSError:
|
||||
failed.append(pid)
|
||||
|
||||
with contextlib.suppress(Exception):
|
||||
print(f"⟲ Reaped {len(killed)} orphaned desktop-local serve backend(s) ({reason}): {killed or matched}")
|
||||
return {"matched": matched, "killed": killed, "failed": failed}
|
||||
|
||||
@@ -20,36 +20,26 @@ from utils import atomic_replace
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Prepended to upload-bound content when redaction is enabled so reviewers of the public
|
||||
# paste know it was sanitized; trailing newline keeps it on its own line.
|
||||
# Prepended to upload-bound content when redaction is enabled so paste reviewers know.
|
||||
_REDACTION_BANNER = (
|
||||
"[hermes debug share: log content redacted at upload time. "
|
||||
"run with --no-redact to disable]\n"
|
||||
)
|
||||
|
||||
_EMAIL_ADDRESS_RE = re.compile(
|
||||
r"(?<![A-Za-z0-9._%+-])"
|
||||
r"[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.[A-Za-z]{2,}"
|
||||
r"(?![A-Za-z0-9._%+-])"
|
||||
)
|
||||
|
||||
# Paste services — paste.rs first, dpaste.com as fallback.
|
||||
_PASTE_RS_URL = "https://paste.rs/"
|
||||
_PASTE_RS_URL = "https://paste.rs/" # primary; dpaste.com is the fallback
|
||||
_DPASTE_COM_URL = "https://dpaste.com/api/"
|
||||
_USER_AGENT = "hermes-agent/debug-share"
|
||||
_MAX_LOG_BYTES = 512_000 # per log file for upload (paste.rs caps at ~1 MB)
|
||||
_AUTO_DELETE_SECONDS = 21600 # 6 hours
|
||||
|
||||
# Max bytes read from one log file for upload (paste.rs caps at ~1 MB; keep headroom).
|
||||
_MAX_LOG_BYTES = 512_000
|
||||
# Pending-deletion tracking: the gateway cron ticker calls ``_sweep_expired_pastes`` hourly and
|
||||
# ``hermes debug`` sweeps on entry (CLI-only users). Replaced a fork-and-sleep subprocess that
|
||||
# leaked ~20 MB per share.
|
||||
|
||||
# Auto-delete pastes after 6 hours.
|
||||
_AUTO_DELETE_SECONDS = 21600
|
||||
|
||||
|
||||
# ── Pending-deletion tracking ──────────────────────────────────────────────
|
||||
# Deletion is driven by the gateway's cron ticker (``gateway/run.py::_start_cron_ticker``
|
||||
# calls ``_sweep_expired_pastes`` hourly); ``hermes debug`` also sweeps opportunistically on
|
||||
# entry for CLI-only users who never start the gateway. Replaces the old fork-and-sleep
|
||||
# subprocess that leaked ~20 MB of resident interpreter per share.
|
||||
|
||||
def _pending_file() -> Path:
|
||||
return get_hermes_home() / "pastes" / "pending.json"
|
||||
@@ -85,11 +75,9 @@ def _sweep_expired_pastes(now: Optional[float] = None) -> tuple[int, int]:
|
||||
entries = _load_pending()
|
||||
if not entries:
|
||||
return (0, 0)
|
||||
|
||||
current = time.time() if now is None else now
|
||||
deleted = 0
|
||||
remaining: list[dict] = []
|
||||
|
||||
for entry in entries:
|
||||
try:
|
||||
expire_at = float(entry.get("expire_at", 0))
|
||||
@@ -106,23 +94,17 @@ def _sweep_expired_pastes(now: Optional[float] = None) -> tuple[int, int]:
|
||||
deleted += 1 # deleted, or given up on → count as reaped
|
||||
else:
|
||||
remaining.append(entry)
|
||||
|
||||
if deleted:
|
||||
_save_pending(remaining)
|
||||
|
||||
return (deleted, len(remaining))
|
||||
|
||||
|
||||
def _best_effort_sweep_expired_pastes() -> None:
|
||||
"""Pending-paste cleanup that never lets /debug fail offline."""
|
||||
try:
|
||||
with contextlib.suppress(Exception):
|
||||
_sweep_expired_pastes()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
# ── Privacy / delete helpers ───────────────────────────────────────────────
|
||||
|
||||
_PRIVACY_NOTICE = """\
|
||||
⚠️ This will upload system info + logs to a PUBLIC paste service.
|
||||
|
||||
@@ -162,7 +144,6 @@ def delete_paste(url: str) -> bool:
|
||||
paste_id = _extract_paste_id(url)
|
||||
if not paste_id:
|
||||
raise ValueError(f"Cannot delete: only paste.rs URLs are supported. Got: {url}")
|
||||
|
||||
req = urllib.request.Request(f"{_PASTE_RS_URL}{paste_id}", method="DELETE",
|
||||
headers={"User-Agent": _USER_AGENT})
|
||||
with urllib.request.urlopen(req, timeout=30) as resp:
|
||||
@@ -174,7 +155,6 @@ def _schedule_auto_delete(urls: list[str], delay_seconds: int = _AUTO_DELETE_SEC
|
||||
paste_rs_urls = [u for u in urls if _extract_paste_id(u)]
|
||||
if not paste_rs_urls:
|
||||
return
|
||||
|
||||
# Dedupe by URL, keeping the later expire_at.
|
||||
by_url: dict[str, float] = {e["url"]: float(e["expire_at"]) for e in _load_pending()}
|
||||
expire_at = time.time() + delay_seconds
|
||||
@@ -200,17 +180,11 @@ def _upload_paste_rs(content: str) -> str:
|
||||
|
||||
def _upload_dpaste_com(content: str, expiry_days: int = 7) -> str:
|
||||
boundary = "----HermesDebugBoundary9f3c"
|
||||
|
||||
def _field(name: str, value: str) -> str:
|
||||
return f'--{boundary}\r\nContent-Disposition: form-data; name="{name}"\r\n\r\n{value}\r\n'
|
||||
|
||||
body = (
|
||||
_field("content", content)
|
||||
+ _field("syntax", "text")
|
||||
+ _field("expiry_days", str(expiry_days))
|
||||
+ f"--{boundary}--\r\n"
|
||||
).encode("utf-8")
|
||||
return _post_paste("dpaste.com", _DPASTE_COM_URL, body, f"multipart/form-data; boundary={boundary}")
|
||||
fields = (("content", content), ("syntax", "text"), ("expiry_days", str(expiry_days)))
|
||||
body = ("".join(f'--{boundary}\r\nContent-Disposition: form-data; name="{n}"\r\n\r\n{v}\r\n'
|
||||
for n, v in fields) + f"--{boundary}--\r\n").encode("utf-8")
|
||||
return _post_paste("dpaste.com", _DPASTE_COM_URL, body,
|
||||
f"multipart/form-data; boundary={boundary}")
|
||||
|
||||
|
||||
def upload_to_pastebin(content: str, expiry_days: int = 7) -> str:
|
||||
@@ -227,12 +201,9 @@ def upload_to_pastebin(content: str, expiry_days: int = 7) -> str:
|
||||
raise RuntimeError("Failed to upload to any paste service:\n " + "\n ".join(errors))
|
||||
|
||||
|
||||
# ── Log file reading ───────────────────────────────────────────────────────
|
||||
|
||||
@dataclass
|
||||
class LogSnapshot:
|
||||
"""Single-read snapshot of a log file used by debug-share."""
|
||||
|
||||
path: Optional[Path]
|
||||
tail_text: str
|
||||
full_text: Optional[str]
|
||||
@@ -245,9 +216,8 @@ def _primary_log_path(log_name: str) -> Optional[Path]:
|
||||
return (get_hermes_home() / "logs" / filename) if filename else None
|
||||
|
||||
|
||||
# Logs written by a client process rather than this backend. When the desktop app talks to a
|
||||
# remote/docker/SSH backend, `hermes debug share` runs on the *backend* and can never see
|
||||
# them — a bare "(file not found)" would read as "the app logged nothing" and misdirect triage.
|
||||
# Logs written by a client process, invisible to a remote/docker/SSH backend running `debug
|
||||
# share`; a bare "(file not found)" would read as "the app logged nothing" and misdirect triage.
|
||||
_CLIENT_SIDE_LOGS = {
|
||||
"desktop": (
|
||||
"written by Hermes Desktop on the machine running the app, not by this "
|
||||
@@ -278,11 +248,8 @@ def _resolve_log_path(log_name: str) -> Optional[Path]:
|
||||
|
||||
|
||||
def _redact_log_text(text: str) -> str:
|
||||
"""Force-mode ``redact_sensitive_text`` (+ email scrub) over upload-bound text.
|
||||
|
||||
``force=True`` so redaction fires regardless of the operator's ``security.redact_secrets``
|
||||
setting; only the in-memory copy headed for the paste service is sanitized.
|
||||
"""
|
||||
"""``redact_sensitive_text(force=True)`` + email scrub — fires regardless of the operator's
|
||||
``security.redact_secrets`` setting; only the in-memory upload copy is sanitized."""
|
||||
if not text:
|
||||
return text
|
||||
from agent.redact import redact_sensitive_text
|
||||
@@ -293,11 +260,8 @@ def _redact_log_text(text: str) -> str:
|
||||
def _read_tail_bytes(
|
||||
log_path: Path, size: int, max_bytes: int, tail_lines: int,
|
||||
) -> tuple[bytes, bool]:
|
||||
"""Read the whole file, or enough of its tail for both views → (raw, truncated).
|
||||
|
||||
For oversized files, read backwards until we have ``max_bytes`` for the standalone upload
|
||||
AND enough newline context to render the summary tail from the same snapshot.
|
||||
"""
|
||||
"""Whole file, or (oversized) a backwards read holding ``max_bytes`` for the full upload AND
|
||||
enough newlines for the summary tail from the same snapshot → (raw, truncated)."""
|
||||
with open(log_path, "rb") as f:
|
||||
if size <= max_bytes:
|
||||
return f.read(), False
|
||||
@@ -322,12 +286,8 @@ def _read_tail_bytes(
|
||||
def _capture_log_snapshot(
|
||||
log_name: str, *, tail_lines: int, max_bytes: int = _MAX_LOG_BYTES, redact: bool = True,
|
||||
) -> LogSnapshot:
|
||||
"""Capture a log once and derive the summary tail and full-log views from it.
|
||||
|
||||
Both views must come from the same read: a rotation/truncate between reads would make the
|
||||
report look newer than the uploaded ``agent.log`` paste. With ``redact`` both texts are
|
||||
upload-safe; the on-disk file is never modified.
|
||||
"""
|
||||
"""Capture a log once and derive the summary tail and full-log views from that single read
|
||||
(a rotation between two reads would make the report look newer than the uploaded log)."""
|
||||
log_path = _resolve_log_path(log_name)
|
||||
if log_path is None:
|
||||
primary = _primary_log_path(log_name)
|
||||
@@ -336,33 +296,25 @@ def _capture_log_snapshot(
|
||||
|
||||
try:
|
||||
size = log_path.stat().st_size
|
||||
if size == 0:
|
||||
# race: file was truncated between _resolve_log_path and stat
|
||||
if size == 0: # truncated between _resolve_log_path and stat
|
||||
return LogSnapshot(path=log_path, tail_text="(file empty)", full_text=None)
|
||||
|
||||
raw, truncated = _read_tail_bytes(log_path, size, max_bytes, tail_lines)
|
||||
|
||||
full_raw = raw
|
||||
if truncated and len(full_raw) > max_bytes:
|
||||
cut = len(full_raw) - max_bytes
|
||||
# Only drop a partial first line when the cut lands genuinely mid-line (the byte
|
||||
# before the cut is not a newline).
|
||||
# Drop a partial first line only when the cut lands genuinely mid-line.
|
||||
on_boundary = cut > 0 and full_raw[cut - 1 : cut] == b"\n"
|
||||
full_raw = full_raw[cut:]
|
||||
if not on_boundary and b"\n" in full_raw:
|
||||
full_raw = full_raw.split(b"\n", 1)[1]
|
||||
|
||||
all_text = raw.decode("utf-8", errors="replace")
|
||||
tail_text = "".join(all_text.splitlines(keepends=True)[-tail_lines:]).rstrip("\n")
|
||||
|
||||
full_text = full_raw.decode("utf-8", errors="replace")
|
||||
if truncated:
|
||||
full_text = f"[... truncated — showing last ~{max_bytes // 1024}KB ...]\n{full_text}"
|
||||
|
||||
if redact:
|
||||
tail_text = _redact_log_text(tail_text)
|
||||
full_text = _redact_log_text(full_text)
|
||||
|
||||
return LogSnapshot(path=log_path, tail_text=tail_text, full_text=full_text)
|
||||
except Exception as exc:
|
||||
return LogSnapshot(path=log_path, tail_text=f"(error reading: {exc})", full_text=None)
|
||||
@@ -388,8 +340,6 @@ def _capture_default_log_snapshots(
|
||||
}
|
||||
|
||||
|
||||
# ── Debug report collection ────────────────────────────────────────────────
|
||||
|
||||
def _capture_dump() -> str:
|
||||
"""Run ``hermes dump`` and return its stdout as a string."""
|
||||
from hermes_cli.dump import run_dump
|
||||
@@ -408,42 +358,27 @@ def collect_debug_report(
|
||||
``dump_text`` is pre-captured dump output; when empty, ``hermes dump`` is run internally.
|
||||
"""
|
||||
buf = io.StringIO()
|
||||
|
||||
if not dump_text:
|
||||
dump_text = _capture_dump()
|
||||
buf.write(dump_text)
|
||||
|
||||
buf.write(dump_text or _capture_dump())
|
||||
if log_snapshots is None:
|
||||
log_snapshots = _capture_default_log_snapshots(log_lines)
|
||||
|
||||
# Sanitiser heal counters: in-process, in-memory — populated when this report is built
|
||||
# inside a process that ran agent turns (gateway /debug share); empty from a fresh CLI
|
||||
# process, where the errors.log tail below carries the same escalation lines instead.
|
||||
try:
|
||||
# In-process sanitiser heal counters: populated only inside a process that ran agent turns
|
||||
# (gateway /debug share); a fresh CLI's errors.log tail carries the same escalation lines.
|
||||
with contextlib.suppress(Exception):
|
||||
from agent.agent_runtime_helpers import get_sanitizer_heal_stats
|
||||
heal_stats = get_sanitizer_heal_stats()
|
||||
if heal_stats:
|
||||
buf.write("\n\n--- transcript sanitiser heal counters ---\n")
|
||||
for sess, st in sorted(heal_stats.items()):
|
||||
buf.write(
|
||||
f"session {sess}: {st['heal_events']} heal events, "
|
||||
f"{st['messages_healed']} messages healed, "
|
||||
f"escalated={st['escalated']}\n"
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
buf.write(f"session {sess}: {st['heal_events']} heal events, "
|
||||
f"{st['messages_healed']} messages healed, escalated={st['escalated']}\n")
|
||||
buf.write("\n")
|
||||
for name in _REPORT_LOGS:
|
||||
buf.write(f"\n--- {name}.log (last {_tail_budget(name, log_lines)} lines) ---\n"
|
||||
f"{log_snapshots[name].tail_text}\n")
|
||||
|
||||
return buf.getvalue()
|
||||
|
||||
|
||||
# ── Shared bundle collection (paste.rs and Nous-S3 paths) ──────────────────
|
||||
|
||||
# Bundle format identifier in the Nous-S3 JSON envelope; the discord-support viewer keys off it.
|
||||
# Nous-S3 envelope format id; the discord-support viewer keys off it.
|
||||
_NOUS_BUNDLE_FORMAT = "hermes-debug-share/1"
|
||||
|
||||
|
||||
@@ -466,23 +401,17 @@ def collect_share_bundle(log_lines: int = 200, redact: bool = True) -> dict[str,
|
||||
|
||||
|
||||
def build_nous_bundle(bundle: dict[str, str], redact: bool = True) -> bytes:
|
||||
"""Gzip a :func:`collect_share_bundle` mapping into the Nous envelope.
|
||||
|
||||
The JSON shape (``format``, ``redacted``, ``created``, ``files``) is what the
|
||||
discord-support viewer parses — keep it stable.
|
||||
"""
|
||||
"""Gzip a :func:`collect_share_bundle` mapping into the Nous envelope (shape parsed by the
|
||||
discord-support viewer — keep it stable)."""
|
||||
envelope = {"format": _NOUS_BUNDLE_FORMAT, "redacted": bool(redact),
|
||||
"created": datetime.datetime.now(datetime.timezone.utc).isoformat(),
|
||||
"files": bundle}
|
||||
return gzip.compress(json.dumps(envelope).encode("utf-8"))
|
||||
|
||||
|
||||
# ── CLI entry points ───────────────────────────────────────────────────────
|
||||
|
||||
@dataclass
|
||||
class DebugShareResult:
|
||||
"""Outcome of a ``debug share`` upload, so non-CLI callers can render real links."""
|
||||
|
||||
urls: dict # label -> paste URL (e.g. {"Report": "...", "agent.log": "..."})
|
||||
failures: list # human-readable "label: error" strings for optional uploads
|
||||
redacted: bool # whether force-mode redaction was applied before upload
|
||||
@@ -495,19 +424,17 @@ def build_debug_share(
|
||||
) -> DebugShareResult:
|
||||
"""Collect the debug report + full logs, upload each, return the URLs.
|
||||
|
||||
Shared core behind ``hermes debug share`` and the dashboard ``POST /api/ops/debug-share``.
|
||||
Blocking network I/O — callers inside an event loop must run it in a worker thread.
|
||||
Shared by ``hermes debug share`` and the dashboard ``POST /api/ops/debug-share``. Blocking
|
||||
network I/O — callers inside an event loop must run it in a worker thread.
|
||||
"""
|
||||
_best_effort_sweep_expired_pastes()
|
||||
bundle = collect_share_bundle(log_lines=log_lines, redact=redact)
|
||||
if redact:
|
||||
logger.info(
|
||||
"hermes debug share: applied force-mode redaction to log snapshots before upload"
|
||||
)
|
||||
"hermes debug share: applied force-mode redaction to log snapshots before upload")
|
||||
report = bundle["report"]
|
||||
failures: list[str] = []
|
||||
# Summary report is required — raises on failure so callers can fall back; full logs are
|
||||
# optional — failures are collected, not raised.
|
||||
# The summary report is required (raises so callers can fall back); full logs are optional.
|
||||
urls = {"Report": upload_to_pastebin(report, expiry_days=expiry)}
|
||||
for label, content in bundle.items():
|
||||
if label == "report":
|
||||
@@ -547,8 +474,7 @@ def run_debug_share(args):
|
||||
redact = not getattr(args, "no_redact", False)
|
||||
|
||||
if getattr(args, "local", False):
|
||||
# Never uploads — render via the shared collector so the output matches exactly what
|
||||
# would be uploaded, and bail before any network I/O.
|
||||
# Same collector as the upload path so the output matches exactly; no network I/O.
|
||||
_best_effort_sweep_expired_pastes()
|
||||
print("Collecting debug report...")
|
||||
bundle = collect_share_bundle(log_lines=log_lines, redact=redact)
|
||||
@@ -559,9 +485,7 @@ def run_debug_share(args):
|
||||
return
|
||||
|
||||
if getattr(args, "nous", False):
|
||||
_run_debug_share_nous(args, log_lines=log_lines, redact=redact)
|
||||
return
|
||||
|
||||
return _run_debug_share_nous(args, log_lines=log_lines, redact=redact)
|
||||
print(_PRIVACY_NOTICE)
|
||||
if not _confirm_upload(args):
|
||||
return
|
||||
@@ -572,7 +496,6 @@ def run_debug_share(args):
|
||||
print(f"\nUpload failed: {exc}", file=sys.stderr)
|
||||
print("\nRun `hermes debug share --local` to print the report instead.\n")
|
||||
sys.exit(1)
|
||||
|
||||
label_width = max(len(k) for k in result.urls)
|
||||
print("\nDebug report uploaded:")
|
||||
for label, url in result.urls.items():
|
||||
@@ -616,16 +539,11 @@ def _run_debug_share_nous(args, *, log_lines: int, redact: bool) -> None:
|
||||
try:
|
||||
res = share_to_nous(build_nous_bundle(bundle, redact=redact))
|
||||
except Exception as exc:
|
||||
print(
|
||||
f"\nNous upload failed: {exc}\n"
|
||||
"\nThe Nous diagnostics service may be unavailable or not yet "
|
||||
"provisioned.\n"
|
||||
"Run `hermes debug share --local` to print the report instead, "
|
||||
"or `hermes debug share` to upload to a public paste service.\n",
|
||||
file=sys.stderr,
|
||||
)
|
||||
print(f"\nNous upload failed: {exc}\n"
|
||||
"\nThe Nous diagnostics service may be unavailable or not yet provisioned.\n"
|
||||
"Run `hermes debug share --local` to print the report instead, "
|
||||
"or `hermes debug share` to upload to a public paste service.\n", file=sys.stderr)
|
||||
sys.exit(1)
|
||||
|
||||
view_url = res.get("viewUrl") or res.get("view_url")
|
||||
expires_at = res.get("expiresAt") or res.get("expires_at")
|
||||
print("\nDebug bundle uploaded to Nous (private):")
|
||||
@@ -633,7 +551,6 @@ def _run_debug_share_nous(args, *, log_lines: int, redact: bool) -> None:
|
||||
else f" (no view URL returned; upload id: {res.get('id', '?')})")
|
||||
print(f"\n⏱ Auto-deletes at {expires_at} (14-day retention)." if expires_at
|
||||
else "\n⏱ Auto-deletes after 14 days.")
|
||||
|
||||
print("\nShare this private link with the Nous team — only Nous staff "
|
||||
"(via Google login) can open it.\n"
|
||||
"\nPick up the discussion in:\n"
|
||||
@@ -649,7 +566,6 @@ def run_debug_delete(args):
|
||||
print("Usage: hermes debug delete <url> [<url> ...]\n"
|
||||
" Deletes paste.rs pastes uploaded by 'hermes debug share'.")
|
||||
return
|
||||
|
||||
for url in urls:
|
||||
try:
|
||||
if delete_paste(url):
|
||||
@@ -663,12 +579,10 @@ def run_debug_delete(args):
|
||||
|
||||
|
||||
def run_debug(args):
|
||||
"""Route debug subcommands."""
|
||||
# Opportunistic sweep of expired pastes on every ``hermes debug`` call.
|
||||
"""Route debug subcommands (sweeping expired pastes opportunistically on every call)."""
|
||||
_best_effort_sweep_expired_pastes()
|
||||
|
||||
handlers = {"share": run_debug_share, "delete": run_debug_delete}
|
||||
handler = handlers.get(getattr(args, "debug_command", None))
|
||||
handler = {"share": run_debug_share, "delete": run_debug_delete}.get(
|
||||
getattr(args, "debug_command", None))
|
||||
if handler is None:
|
||||
print(_DEBUG_USAGE)
|
||||
else:
|
||||
|
||||
Reference in New Issue
Block a user