# Conflicts: # AGENTS.md # acp_adapter/edit_approval.py # acp_adapter/server.py # agent/agent_init.py # agent/anthropic_adapter.py # agent/anthropic_credentials.py # agent/auxiliary_client.py # agent/azure_identity_adapter.py # agent/bedrock_adapter.py # agent/browser_registry.py # agent/chat_completion_helpers.py # agent/coding_context.py # agent/context_references.py # agent/conversation_loop.py # agent/copilot_acp_client.py # agent/credits_tracker.py # agent/curator.py # agent/curator_backup.py # agent/deadline.py # agent/display.py # agent/errors.py # agent/estop.py # agent/i18n.py # agent/image_gen_registry.py # agent/image_routing.py # agent/learning_graph.py # agent/learning_mutations.py # agent/lsp/servers.py # agent/model_metadata.py # agent/models_dev.py # agent/monitoring/gateway_health_export.py # agent/monitoring/otlp_exporter.py # agent/pet/store.py # agent/process_bootstrap.py # agent/prompt_builder.py # agent/proxy_sources/iron_proxy.py # agent/secret_sources/_cache.py # agent/secret_sources/bitwarden.py # agent/secret_sources/registry.py # agent/shell_hooks.py # agent/skill_bundles.py # agent/skill_commands.py # agent/skill_utils.py # agent/ssl_guard.py # agent/ssl_verify.py # agent/system_prompt.py # agent/terminal_env_registry.py # agent/trace_upload.py # agent/transcription_registry.py # agent/tts_registry.py # agent/verify/environment.py # agent/vertex_adapter.py # agent/video_gen_registry.py # agent/web_search_registry.py # cli.py # cron/jobs.py # cron/scheduler.py # gateway/agent_cache_pressure.py # gateway/cgroup_cleanup.py # gateway/channel_directory.py # gateway/config.py # gateway/control_socket.py # gateway/dead_targets.py # gateway/drain_control.py # gateway/hooks.py # gateway/kanban_watchers.py # gateway/lifecycle_ledger.py # gateway/mirror.py # gateway/pairing.py # gateway/platform_registry.py # gateway/platforms/helpers.py # gateway/platforms/weixin.py # gateway/readiness.py # gateway/restart_loop_guard.py # gateway/rich_sent_store.py # gateway/run.py # gateway/session.py # gateway/shutdown_flush.py # gateway/shutdown_forensics.py # gateway/slash_commands.py # gateway/status.py # gateway/sticker_cache.py # gateway/whatsapp_identity.py # hermes_bootstrap.py # hermes_cli/_early_recovery.py # hermes_cli/_install_repair.py # hermes_cli/_startup_fast.py # hermes_cli/_subprocess_compat.py # hermes_cli/agent_plugins.py # hermes_cli/auth.py # hermes_cli/backup.py # hermes_cli/banner.py # hermes_cli/browser_connect.py # hermes_cli/build_info.py # hermes_cli/cli_agent_setup_mixin.py # hermes_cli/cli_commands_mixin.py # hermes_cli/codex_models.py # hermes_cli/config.py # hermes_cli/config_defaults.py # hermes_cli/config_migrations.py # hermes_cli/container_boot.py # hermes_cli/dashboard_auth/registry.py # hermes_cli/debug.py # hermes_cli/dep_ensure.py # hermes_cli/doctor.py # hermes_cli/doctor_live.py # hermes_cli/dump.py # hermes_cli/env_loader.py # hermes_cli/foreign_sessions.py # hermes_cli/gateway.py # hermes_cli/gateway_windows.py # hermes_cli/gui_uninstall.py # hermes_cli/image_provenance.py # hermes_cli/install_identity.py # hermes_cli/kanban.py # hermes_cli/kanban_db.py # hermes_cli/linux_desktop_entry.py # hermes_cli/local_runtime/binaries.py # hermes_cli/local_runtime/endpoint.py # hermes_cli/local_runtime/growth.py # hermes_cli/local_runtime/supervisor.py # hermes_cli/logs.py # hermes_cli/macos_tcc_anchor.py # hermes_cli/main.py # hermes_cli/memory_setup.py # hermes_cli/model_catalog.py # hermes_cli/models.py # hermes_cli/nous_subscription.py # hermes_cli/npm_engine.py # hermes_cli/plugin_index.py # hermes_cli/plugins.py # hermes_cli/plugins_cmd.py # hermes_cli/profile_distribution.py # hermes_cli/profiles.py # hermes_cli/prompt_size.py # hermes_cli/psutil_android.py # hermes_cli/runtime_repair.py # hermes_cli/security_advisories.py # hermes_cli/security_audit.py # hermes_cli/security_audit_startup.py # hermes_cli/service_manager.py # hermes_cli/session_export_md.py # hermes_cli/setup.py # hermes_cli/skills_hub.py # hermes_cli/slack_cli.py # hermes_cli/status.py # hermes_cli/subcommands/gateway.py # hermes_cli/subcommands/uninstall.py # hermes_cli/tools_config.py # hermes_cli/uninstall.py # hermes_cli/update_cmd.py # hermes_cli/update_contract.py # hermes_cli/update_inventory.py # hermes_cli/update_lock.py # hermes_cli/update_receipt.py # hermes_cli/urllib_security.py # hermes_cli/web_routers/local_models.py # hermes_cli/web_routers/profiles.py # hermes_cli/web_routers/skills.py # hermes_cli/web_server.py # hermes_constants.py # hermes_state.py # plugins/disk-cleanup/__init__.py # plugins/disk-cleanup/disk_cleanup.py # plugins/google_meet/node/registry.py # plugins/google_meet/node/server.py # plugins/google_meet/process_manager.py # plugins/google_meet/realtime/openai_client.py # plugins/hermes-achievements/dashboard/plugin_api.py # plugins/memory/hindsight/__init__.py # plugins/memory/honcho/__init__.py # plugins/memory/honcho/cli.py # plugins/memory/honcho/client.py # plugins/memory/honcho/oauth.py # plugins/memory/honcho/session.py # plugins/memory/mem0/__init__.py # plugins/memory/mem0/_setup.py # plugins/memory/openviking/__init__.py # plugins/memory/retaindb/__init__.py # plugins/memory/supermemory/__init__.py # plugins/platforms/a2a/protocol.py # plugins/platforms/dingtalk/adapter.py # plugins/platforms/discord/adapter.py # plugins/platforms/feishu/adapter.py # plugins/platforms/google_chat/adapter.py # plugins/platforms/matrix/adapter.py # plugins/platforms/photon/adapter.py # plugins/platforms/photon/auth.py # plugins/platforms/photon/cli.py # plugins/platforms/slack/adapter.py # plugins/platforms/teams/adapter.py # plugins/platforms/telegram/adapter.py # plugins/platforms/wecom/callback_adapter.py # plugins/platforms/whatsapp/adapter.py # plugins/teams_pipeline/store.py # plugins/video_gen/fal/__init__.py # plugins/web/ddgs/provider.py # plugins/web/exa/provider.py # plugins/web/firecrawl/provider.py # plugins/web/parallel/provider.py # tests/agent/test_ssl_ca_guard.py # tests/hermes_cli/test_certifi_repair.py # tests/hermes_cli/test_cmd_update.py # tests/hermes_cli/test_cmd_update_apt.py # tests/hermes_cli/test_dashboard_unified_launch.py # tests/hermes_cli/test_dep_ensure.py # tests/hermes_cli/test_doctor.py # tests/hermes_cli/test_doctor_live.py # tests/hermes_cli/test_gui_command.py # tests/hermes_cli/test_kanban_boards.py # tests/hermes_cli/test_kanban_db.py # tests/hermes_cli/test_lazy_refresh_venv_repair.py # tests/hermes_cli/test_memory_setup_provider_arg.py # tests/hermes_cli/test_nous_subscription.py # tests/hermes_cli/test_pip_install_detection.py # tests/hermes_cli/test_profile_export_credentials.py # tests/hermes_cli/test_psutil_android_extract.py # tests/hermes_cli/test_status.py # tests/hermes_cli/test_tui_npm_install.py # tests/hermes_cli/test_update_fleet_restart_pending.py # tests/hermes_cli/test_update_head_moved_gate.py # tests/hermes_cli/test_update_interrupted_recovery.py # tests/hermes_cli/test_web_server.py # tests/hermes_cli/test_web_ui_build.py # tests/test_hermes_logging.py # tests/test_managed_runtime_resolution.py # tests/tools/test_browser_chromium_autoinstall.py # tests/tools/test_browser_chromium_check.py # tests/tools/test_browser_homebrew_paths.py # tests/tools/test_browser_lightpanda.py # tests/tools/test_browser_npx_warmup.py # tests/tools/test_browser_open_timeout.py # tests/tools/test_browser_orphan_reaper.py # tests/tools/test_browser_real_profile.py # tests/tools/test_browser_suspect_recycle.py # tests/tools/test_find_shell.py # tests/tools/test_local_env_blocklist.py # tests/tools/test_macos_protected_search.py # tests/tui_gateway/test_compute_host.py # tools/approval.py # tools/blueprints.py # tools/bot_mode_dm.py # tools/bot_mode_probe.py # tools/bot_relay.py # tools/browser_tool.py # tools/browser_use_cli.py # tools/checkpoint_manager.py # tools/code_execution_tool.py # tools/code_kernel.py # tools/computer_use/cua_backend.py # tools/cronjob_tools.py # tools/discord_tool.py # tools/environments/base.py # tools/environments/daytona.py # tools/environments/local.py # tools/environments/modal.py # tools/environments/vercel_sandbox.py # tools/fal_common.py # tools/file_operations.py # tools/lazy_deps.py # tools/mcp_tool.py # tools/neutts_synth.py # tools/process_registry.py # tools/read_extract.py # tools/registry.py # tools/skill_ledger.py # tools/skill_linter.py # tools/skill_manager_tool.py # tools/skill_usage.py # tools/skills_ast_audit.py # tools/skills_guard.py # tools/skills_hub.py # tools/skills_sync.py # tools/skills_sync_client.py # tools/skills_tool.py # tools/terminal_scope.py # tools/terminal_tool.py # tools/tirith_security.py # tools/transcription_tools.py # tools/tts_tool.py # tools/vision_tools.py # tools/voice_mode.py # tools/wake_word.py # tools/web_result_cache.py # tools/website_policy.py # tools/working_diff.py # tools/write_approval.py # tui_gateway/entry.py # tui_gateway/methods_tools.py # tui_gateway/server.py
186 lines
6.7 KiB
Python
186 lines
6.7 KiB
Python
"""Monitor-mode cron support — hash-suppressed change detection.
|
|
|
|
A monitor job attaches a cheap source (``monitor_script`` / ``monitor_url``) to an LLM cron job.
|
|
Each tick runs the source FIRST and hashes its EXACT output bytes (no timestamp/whitespace
|
|
normalization — scripts must emit stable output) against the hash from the last agent-triggering
|
|
tick: unchanged → agent run suppressed (silent ``no_change`` run); changed/first run → a "MONITOR
|
|
CHANGE DETECTED" block (capped unified diff + new output) is injected into the prompt; source
|
|
failure → an ERROR, never a change, and the stored hash is left untouched. State:
|
|
``job["monitor_state"]`` in jobs.json (hash + last_changed_at) and
|
|
``OUTPUT_DIR/<job_id>/monitor_last_output.txt`` (for the diff).
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import difflib
|
|
import hashlib
|
|
import logging
|
|
from dataclasses import dataclass
|
|
from typing import Optional
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# Prompt-injection caps: unified diff, and new-output block (mirrors the 8k context_from truncation
|
|
# in cron/scheduler.py). Then bounded-GET limits for monitor_url sources.
|
|
MAX_DIFF_CHARS = 4000
|
|
MAX_OUTPUT_CHARS = 8000
|
|
URL_TIMEOUT_SECONDS = 30
|
|
MAX_URL_BYTES = 262_144 # 256 KiB
|
|
|
|
_SNAPSHOT_FILENAME = "monitor_last_output.txt"
|
|
|
|
|
|
@dataclass
|
|
class MonitorOutcome:
|
|
"""Result of one monitor-source evaluation."""
|
|
|
|
ok: bool
|
|
changed: bool = False
|
|
first_run: bool = False
|
|
context_block: Optional[str] = None
|
|
error: Optional[str] = None
|
|
|
|
|
|
def hash_monitor_output(output: str) -> str:
|
|
"""Hash the monitor output as exact UTF-8 bytes (no normalization)."""
|
|
return hashlib.sha256(output.encode("utf-8", errors="replace")).hexdigest()
|
|
|
|
|
|
def build_monitor_diff(old: str, new: str) -> str:
|
|
"""Unified diff of old vs new monitor output, capped at MAX_DIFF_CHARS."""
|
|
diff = "\n".join(
|
|
difflib.unified_diff(
|
|
old.splitlines(), new.splitlines(), fromfile="previous", tofile="current", lineterm="",
|
|
)
|
|
)
|
|
if len(diff) > MAX_DIFF_CHARS:
|
|
diff = diff[:MAX_DIFF_CHARS] + "\n... [diff truncated]"
|
|
return diff
|
|
|
|
|
|
def _snapshot_path(job_id: str):
|
|
from cron.jobs import _job_output_dir
|
|
|
|
return _job_output_dir(job_id) / _SNAPSHOT_FILENAME
|
|
|
|
|
|
def _read_last_output(job_id: str) -> str:
|
|
try:
|
|
path = _snapshot_path(job_id)
|
|
if path.exists():
|
|
return path.read_text(encoding="utf-8-sig")
|
|
except Exception as exc:
|
|
logger.warning("Monitor: failed to read last output for %r: %s", job_id, exc)
|
|
return ""
|
|
|
|
|
|
def _write_last_output(job_id: str, output: str) -> None:
|
|
try:
|
|
from cron.jobs import _ensure_cron_dir
|
|
|
|
path = _snapshot_path(job_id)
|
|
_ensure_cron_dir(path.parent)
|
|
path.write_text(output, encoding="utf-8")
|
|
except Exception as exc:
|
|
logger.warning("Monitor: failed to persist last output for %r: %s", job_id, exc)
|
|
|
|
|
|
def _fetch_monitor_url(url: str) -> tuple[bool, str]:
|
|
"""Bounded GET of a monitor URL. Returns (ok, body-or-error)."""
|
|
import urllib.request
|
|
|
|
if not str(url).lower().startswith(("http://", "https://")):
|
|
return False, f"monitor_url must be http(s): {url!r}"
|
|
try:
|
|
req = urllib.request.Request(url, headers={"User-Agent": "hermes-cron-monitor"})
|
|
with urllib.request.urlopen(req, timeout=URL_TIMEOUT_SECONDS) as resp: # nosec B310 — scheme checked above
|
|
body = resp.read(MAX_URL_BYTES + 1)
|
|
return True, body[:MAX_URL_BYTES].decode("utf-8", errors="replace")
|
|
except Exception as exc:
|
|
return False, f"monitor_url fetch failed: {exc}"
|
|
|
|
|
|
def _field(job: dict, key: str) -> str:
|
|
return (job.get(key) or "").strip()
|
|
|
|
|
|
def _run_monitor_source(job: dict) -> tuple[bool, str]:
|
|
"""Run the job's monitor source (script or URL). Returns (ok, output)."""
|
|
monitor_script = _field(job, "monitor_script")
|
|
if monitor_script:
|
|
# Same containment + interpreter rules as the existing `script` field.
|
|
from cron.scheduler_script import _run_job_script
|
|
|
|
return _run_job_script(monitor_script, workdir=_field(job, "workdir") or None)
|
|
monitor_url = _field(job, "monitor_url")
|
|
if monitor_url:
|
|
return _fetch_monitor_url(monitor_url)
|
|
return False, "monitor job has neither monitor_script nor monitor_url"
|
|
|
|
|
|
def job_has_monitor(job: dict) -> bool:
|
|
return bool(_field(job, "monitor_script") or _field(job, "monitor_url"))
|
|
|
|
|
|
def check_monitor(job: dict) -> MonitorOutcome:
|
|
"""Run the monitor source and decide whether the agent should run.
|
|
|
|
On change (or first run) the new hash + snapshot are persisted BEFORE the agent runs — detection
|
|
time is the state boundary, so a failed agent run doesn't re-alert on the same content forever.
|
|
On failure nothing is persisted.
|
|
"""
|
|
job_id = str(job.get("id") or "")
|
|
ok, output = _run_monitor_source(job)
|
|
if not ok:
|
|
return MonitorOutcome(ok=False, error=output)
|
|
|
|
new_hash = hash_monitor_output(output)
|
|
raw_state = job.get("monitor_state")
|
|
last_hash = raw_state.get("last_output_hash") if isinstance(raw_state, dict) else None
|
|
|
|
if last_hash is not None and new_hash == last_hash:
|
|
return MonitorOutcome(ok=True, changed=False)
|
|
|
|
first_run = last_hash is None
|
|
old_output = "" if first_run else _read_last_output(job_id)
|
|
|
|
shown_output = output
|
|
if len(shown_output) > MAX_OUTPUT_CHARS:
|
|
shown_output = shown_output[:MAX_OUTPUT_CHARS] + "\n... [output truncated]"
|
|
|
|
current = f"### Current output\n\n```\n{shown_output}\n```"
|
|
if first_run:
|
|
context_block = (
|
|
"## Monitor Baseline (first run)\n\n"
|
|
"This is the first observation of the monitored source — there is "
|
|
"no previous output to diff against.\n\n" + current
|
|
)
|
|
else:
|
|
diff = build_monitor_diff(old_output, output)
|
|
context_block = (
|
|
"## MONITOR CHANGE DETECTED\n\n"
|
|
"The monitored source's output changed since the last run.\n\n"
|
|
f"### Diff (previous → current)\n\n```diff\n{diff}\n```\n\n" + current
|
|
)
|
|
|
|
_persist_monitor_state(job_id, new_hash, output)
|
|
return MonitorOutcome(ok=True, changed=True, first_run=first_run, context_block=context_block)
|
|
|
|
|
|
def _persist_monitor_state(job_id: str, new_hash: str, output: str) -> None:
|
|
from cron.jobs import _hermes_now, update_job
|
|
|
|
_write_last_output(job_id, output)
|
|
try:
|
|
update_job(
|
|
job_id,
|
|
{
|
|
"monitor_state": {
|
|
"last_output_hash": new_hash,
|
|
"last_changed_at": _hermes_now().isoformat(),
|
|
}
|
|
},
|
|
)
|
|
except Exception as exc:
|
|
logger.warning("Monitor: failed to persist state for %r: %s", job_id, exc)
|