diff --git a/cron/__init__.py b/cron/__init__.py index 4eae752136..9d12ef02d5 100644 --- a/cron/__init__.py +++ b/cron/__init__.py @@ -1,17 +1,6 @@ -"""Cron job scheduling system for Hermes Agent. - -This module provides scheduled task execution, allowing the agent to: -- Run automated tasks on schedules (cron expressions, intervals, one-shot) -- Self-schedule reminders and follow-up tasks -- Execute tasks in isolated sessions (no prior context) - -Cron jobs are executed automatically by the gateway daemon: - hermes gateway install # Install as a user service - sudo hermes gateway install --system # Linux servers: boot-time system service - hermes gateway # Or run in foreground - -The gateway ticks the scheduler every 60 seconds. A file lock prevents duplicate execution if -multiple processes overlap. +"""Cron job scheduling for Hermes Agent: scheduled tasks (cron expressions, intervals, one-shot), +self-scheduled reminders, isolated sessions. The gateway daemon (``hermes gateway [install]``) ticks +the scheduler every 60 seconds; a file lock prevents duplicate execution across processes. """ from cron.jobs import ( @@ -30,7 +19,7 @@ from cron.scheduler import tick __all__ = [ "create_job", - "get_job", + "get_job", "list_jobs", "remove_job", "update_job", diff --git a/cron/blueprint_catalog.py b/cron/blueprint_catalog.py index e6c34065d1..76a084d3f7 100644 --- a/cron/blueprint_catalog.py +++ b/cron/blueprint_catalog.py @@ -553,9 +553,17 @@ def get_blueprint(key: str) -> Optional[AutomationBlueprint]: return _CATALOG_BY_KEY.get(key) -# --------------------------------------------------------------------------- -# Renderers -# --------------------------------------------------------------------------- +def _slot(blueprint: AutomationBlueprint, *names: str, type: Optional[str] = None) -> Optional[BlueprintSlot]: + """First slot matching any of *names* (or *type*), else None.""" + return next((s for s in blueprint.slots if s.name in names or (type and s.type == type)), None) + + +def _slot_default(blueprint: AutomationBlueprint, *names: str) -> Any: + slot = _slot(blueprint, *names) + return slot.default if slot else None + + +# --- Renderers -------------------------------------------------------------------------------- def blueprint_form_schema(blueprint: AutomationBlueprint) -> Dict[str, Any]: """Emit the JSON a form renderer (dashboard / GUI) needs for this blueprint.""" @@ -582,11 +590,8 @@ def blueprint_form_schema(blueprint: AutomationBlueprint) -> Dict[str, Any]: def blueprint_slash_command(blueprint: AutomationBlueprint, values: Optional[Dict[str, Any]] = None) -> str: - """Build the flattened ``/blueprint slot=val …`` command string. - - Uses each slot's default when ``values`` is omitted, so the docs/dashboard - can show a ready-to-paste command. Free-text slots are quoted. - """ + """Build the flattened ``/blueprint slot=val …`` command string. Uses each slot's default + when ``values`` is omitted (ready-to-paste for docs/dashboard). Free-text slots are quoted.""" values = values or {} parts = [f"/blueprint {blueprint.key}"] for s in blueprint.slots: @@ -607,11 +612,11 @@ def blueprint_deeplink(blueprint: AutomationBlueprint, values: Optional[Dict[str from urllib.parse import quote, urlencode values = values or {} - query = {} - for s in blueprint.slots: - val = values.get(s.name, s.default) - if val not in (None, ""): - query[s.name] = str(val) + query = { + s.name: str(values.get(s.name, s.default)) + for s in blueprint.slots + if values.get(s.name, s.default) not in (None, "") + } qs = ("?" + urlencode(query)) if query else "" return f"hermes://blueprint/{quote(blueprint.key)}{qs}" @@ -620,34 +625,27 @@ def _humanize_schedule(blueprint: AutomationBlueprint) -> str: """A short human-readable description of when a blueprint runs (defaults).""" sched = blueprint.schedule_template if sched.startswith("*/"): - iv = next((s for s in blueprint.slots if s.name == "interval_min"), None) - every = (iv.default if iv else None) or sched.split("/")[1].split()[0] + every = _slot_default(blueprint, "interval_min") or sched.split("/")[1].split()[0] return f"every {every} minutes" if "{interval_hours}" in sched: - iv = next((s for s in blueprint.slots if s.name == "interval_hours"), None) - every = str((iv.default if iv else None) or "1") + every = str(_slot_default(blueprint, "interval_hours") or "1") scope = "weekdays, " if "* * 1-5" in sched else "" return f"{scope}every hour" if every == "1" else f"{scope}every {every} hours" - time_slot = next((s for s in blueprint.slots if s.type == "time"), None) + time_slot = _slot(blueprint, type="time") when = time_slot.default if time_slot else None if "* * 1-5" in sched: return f"weekdays at {when}" if when else "every weekday" if "{dow}" in sched: - day_slot = next((s for s in blueprint.slots if s.name in ("day", "recurrence")), None) - scope = (day_slot.default if day_slot else "") or "" + scope = _slot_default(blueprint, "day", "recurrence") or "" if scope and when: return f"{scope} at {when}" return f"at {when}" if when else "on a schedule" - if when: - return f"daily at {when}" - return "on a schedule" + return f"daily at {when}" if when else "on a schedule" def blueprint_catalog_entry(blueprint: AutomationBlueprint) -> Dict[str, Any]: - """Unified serializable shape for a blueprint — used by the docs generator - and the dashboard API. Combines the form schema, the ready-to-paste slash - command, the deep-link URL, and a human-readable schedule. - """ + """Unified serializable shape (docs generator + dashboard API): form schema + ready-to-paste slash + command + deep-link URL + human-readable schedule.""" return { **blueprint_form_schema(blueprint), "schedule": blueprint.schedule_template, @@ -657,9 +655,7 @@ def blueprint_catalog_entry(blueprint: AutomationBlueprint) -> Dict[str, Any]: } -# --------------------------------------------------------------------------- -# Fill + validate + translate to a create_job spec -# --------------------------------------------------------------------------- +# --- Fill + validate + translate to a create_job spec ----------------------------------------- _TIME_RE = re.compile(r"^([01]?\d|2[0-3]):([0-5]\d)$") _DAY_TO_DOW = { @@ -673,14 +669,13 @@ def _resolve_schedule(blueprint: AutomationBlueprint, values: Dict[str, Any]) -> sched = blueprint.schedule_template # A free-text `schedule` slot passes through verbatim (full flexibility). - if "schedule" in values and values["schedule"]: + if values.get("schedule"): return str(values["schedule"]) repl: Dict[str, str] = {} - # time -> minute/hour - time_val = values.get("time") if "{minute}" in sched or "{hour}" in sched: + time_val = values.get("time") if not time_val: raise BlueprintFillError("a time is required") m = _TIME_RE.match(str(time_val).strip()) @@ -689,7 +684,6 @@ def _resolve_schedule(blueprint: AutomationBlueprint, values: Dict[str, Any]) -> repl["hour"] = str(int(m.group(1))) repl["minute"] = str(int(m.group(2))) - # weekday set -> dow if "{dow}" in sched: if "recurrence" in values: preset = str(values.get("recurrence", "everyday")).lower() @@ -706,16 +700,14 @@ def _resolve_schedule(blueprint: AutomationBlueprint, values: Dict[str, Any]) -> else: repl["dow"] = "*" - # interval (minutes) for */N schedules if "{interval_min}" in sched: iv = str(values.get("interval_min", "")).strip() if not iv.isdigit() or int(iv) <= 0: raise BlueprintFillError(f"invalid interval {iv!r} — minutes as a positive integer") repl["interval_min"] = iv - # Any remaining {slot} placeholders are filled verbatim from validated - # enum/text slot values (e.g. an hour-range window). Enum options have - # already been checked in fill_blueprint, so these are safe to interpolate. + # Remaining {slot} placeholders are filled verbatim from validated enum/text slot values (enum + # options were already checked in fill_blueprint, so they are safe to interpolate). for name in re.findall(r"\{(\w+)\}", sched): if name not in repl and name in values: repl[name] = str(values[name]) @@ -734,10 +726,9 @@ def fill_blueprint( ) -> Dict[str, Any]: """Validate ``values`` and return ``cron.jobs.create_job`` kwargs. - Missing required (non-optional) slots raise BlueprintFillError naming the slot, so a form can - show field errors and the agent knows what to ask. Unknown slot names are rejected (a typo'd - ``tiem=07:15`` must not silently create a job with the default time). Enum values are checked - against their options. The result is passed straight to ``create_job`` — no second schema. + Missing required slots raise BlueprintFillError naming the slot (forms show field errors, the + agent knows what to ask). Unknown slot names are rejected (a typo'd ``tiem=07:15`` must not + silently create a job with the default time). Strict enums are checked against their options. """ known = {s.name for s in blueprint.slots} unknown = sorted(set(values) - known) @@ -761,7 +752,6 @@ def fill_blueprint( schedule = _resolve_schedule(blueprint, resolved) - # Render the prompt with whatever slots it references. try: prompt = blueprint.prompt_template.format(**resolved) except KeyError as e: diff --git a/cron/monitor.py b/cron/monitor.py index 9457e51863..9b82b6cc4c 100644 --- a/cron/monitor.py +++ b/cron/monitor.py @@ -19,12 +19,10 @@ from typing import Optional logger = logging.getLogger(__name__) -# Cap for the unified diff injected into the prompt. +# 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 -# Cap for the new-output block injected into the prompt (mirrors the 8k -# context_from truncation in cron/scheduler.py). MAX_OUTPUT_CHARS = 8000 -# Bounded GET limits for monitor_url sources. URL_TIMEOUT_SECONDS = 30 MAX_URL_BYTES = 262_144 # 256 KiB @@ -81,8 +79,9 @@ def _read_last_output(job_id: str) -> str: def _write_last_output(job_id: str, output: str) -> None: try: - path = _snapshot_path(job_id) 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: @@ -99,30 +98,31 @@ def _fetch_monitor_url(url: str) -> tuple[bool, str]: 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) - if len(body) > MAX_URL_BYTES: - body = body[:MAX_URL_BYTES] - return True, body.decode("utf-8", errors="replace") + 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 = (job.get("monitor_script") or "").strip() + monitor_script = _field(job, "monitor_script") if monitor_script: # Same containment + interpreter rules as the existing `script` field. from cron.scheduler import _run_job_script - workdir = (job.get("workdir") or "").strip() or None - return _run_job_script(monitor_script, workdir=workdir) - monitor_url = (job.get("monitor_url") or "").strip() + 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((job.get("monitor_script") or "").strip() or (job.get("monitor_url") or "").strip()) + return bool(_field(job, "monitor_script") or _field(job, "monitor_url")) def check_monitor(job: dict) -> MonitorOutcome: @@ -139,8 +139,7 @@ def check_monitor(job: dict) -> MonitorOutcome: new_hash = hash_monitor_output(output) raw_state = job.get("monitor_state") - state = raw_state if isinstance(raw_state, dict) else {} - last_hash = state.get("last_output_hash") + 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) @@ -152,26 +151,23 @@ def check_monitor(job: dict) -> MonitorOutcome: 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" - f"### Current output\n\n```\n{shown_output}\n```" + "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" - f"### Current output\n\n```\n{shown_output}\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 - ) + 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: diff --git a/cron/scripts/classify_items.py b/cron/scripts/classify_items.py index ba0cd42d4b..59c0db65cf 100644 --- a/cron/scripts/classify_items.py +++ b/cron/scripts/classify_items.py @@ -1,36 +1,14 @@ #!/usr/bin/env python3 """Classify candidate items by urgency/importance and emit only the urgent ones. -The proactive-monitor pattern: a fetch step (a watcher script, an inbox dump, a -feed) produces a list of candidate items; this script scores each with a cheap -LLM and prints ONLY the items at or above a threshold. Below-threshold runs -print nothing, so a cron job wrapping this stays silent unless something -actually matters -- the classic urgency-monitor pattern (fetch -> classify -urgency -> surface only what's above the bar). +The proactive-monitor pattern: a fetch step (watcher script, inbox dump, feed) produces a JSON list +of candidate items (stdin or --input-file); one call to the auxiliary ``monitor`` model scores the +whole batch and ONLY items at/above --threshold are printed. Empty stdout -> the cron job's +[SILENT]/empty-stdout path suppresses delivery, so quiet intervals never spam. A classifier +failure exits non-zero (never silently swallowed). Items are opaque objects; a +title/subject/summary/text field helps, and id/guid/message_id/url is echoed back for upstream dedup. -Design choices: - * Uses Hermes' auxiliary client with task="monitor", so the classifier model - is configured once in config.yaml (auxiliary.monitor.{provider,model}) and - can be a cheap fast model independent of the main chat model. - * Reads items as JSON (a list of objects) from stdin or --input-file. - * One LLM call scores the whole batch (cheap, single round-trip) and returns - structured scores; we filter locally. - * Empty result -> empty stdout -> the cron job's [SILENT]/empty-stdout path - suppresses delivery. No spam on quiet intervals. - -Usage (standalone): - cat items.json | python classify_items.py --threshold 7 \ - --criteria "Urgent if it needs a reply today or is from my manager/family" - -Usage (wired to a watcher via cron, agent mode): - Ask the agent: "Every 10 minutes, run watch_http_json.py for my inbox feed, - pipe its JSON into classify_items.py with my urgency criteria, and deliver - whatever it prints. Stay silent if it prints nothing." - -Item schema (flexible): each item is an object; the classifier sees the whole -object. A "title"/"subject"/"summary"/"text" field helps it judge. An "id" -field (any of id/guid/message_id/url) is echoed back so duplicates can be -deduped upstream. +Usage: cat items.json | python classify_items.py --threshold 7 --criteria "Urgent if ..." """ from __future__ import annotations @@ -40,13 +18,15 @@ import json import sys from typing import Any, Dict, List, Optional +_ID_KEYS = ("id", "guid", "message_id", "url", "link") +_VIEW_KEYS = ("title", "subject", "summary", "text", "body", "from", "sender", "url") + def _eprint(*args: Any) -> None: print(*args, file=sys.stderr) def _load_items(input_file: Optional[str]) -> List[Dict[str, Any]]: - raw = "" if input_file: with open(input_file, encoding="utf-8") as f: raw = f.read() @@ -72,39 +52,16 @@ def _load_items(input_file: Optional[str]) -> List[Dict[str, Any]]: def _item_id(item: Dict[str, Any], index: int) -> str: - for key in ("id", "guid", "message_id", "url", "link"): - val = item.get(key) - if val: - return str(val) - return f"item-{index}" - - -_CLASSIFY_INSTRUCTIONS = ( - "You are an urgency classifier for a proactive assistant. You will be given " - "a numbered list of items and the user's importance criteria. Score EACH " - "item from 0 (ignore entirely) to 10 (interrupt the user now). Return ONLY a " - "JSON array, one object per item, in the same order: " - '[{"index": , "score": , "reason": ""}]. ' - "No prose, no markdown fences. Be conservative: most items should score low. " - "Only score high when the item clearly meets the user's criteria." -) + return next((str(item[key]) for key in _ID_KEYS if item.get(key)), f"item-{index}") def _build_prompt(items: List[Dict[str, Any]], criteria: str) -> str: lines = [f"USER IMPORTANCE CRITERIA:\n{criteria}\n", "ITEMS:"] for i, item in enumerate(items): - # Show a compact view; the model sees the salient fields. - view = { - k: item[k] - for k in ("title", "subject", "summary", "text", "body", "from", "sender", "url") - if k in item - } - if not view: - view = item # fall back to the whole object + # Compact view of the salient fields; the whole object when none are present. + view = {k: item[k] for k in _VIEW_KEYS if k in item} or item lines.append(f"[{i}] {json.dumps(view, ensure_ascii=False)[:1200]}") - lines.append( - "\nReturn the JSON array of scores now (one object per item, same order)." - ) + lines.append("\nReturn the JSON array of scores now (one object per item, same order).") return "\n".join(lines) @@ -121,24 +78,34 @@ def _parse_scores(content: str, n_items: int) -> Dict[int, Dict[str, Any]]: # Last-ditch: find the first [...] block. start = text.find("[") end = text.rfind("]") - if start >= 0 and end > start: - try: - arr = json.loads(text[start : end + 1]) - except json.JSONDecodeError: - _eprint("classify_items: could not parse classifier output") - return {} - else: + if not (start >= 0 and end > start): _eprint("classify_items: classifier returned no JSON array") return {} - out: Dict[int, Dict[str, Any]] = {} - if isinstance(arr, list): - for obj in arr: - if not isinstance(obj, dict): - continue - idx = obj.get("index") - if isinstance(idx, int) and 0 <= idx < n_items: - out[idx] = obj - return out + try: + arr = json.loads(text[start : end + 1]) + except json.JSONDecodeError: + _eprint("classify_items: could not parse classifier output") + return {} + if not isinstance(arr, list): + return {} + return { + obj["index"]: obj + for obj in arr + if isinstance(obj, dict) and isinstance(obj.get("index"), int) and 0 <= obj["index"] < n_items + } + + +def _render_text(surfaced: list) -> str: + blocks = [] + for i, item, s in surfaced: + title = item.get("title") or item.get("subject") or item.get("summary") or _item_id(item, i) + block = f"## [{s.get('score')}/10] {title}" + if url := item.get("url") or item.get("link") or "": + block += f"\n{url}" + if reason := s.get("reason", ""): + block += f"\n_{reason}_" + blocks.append(block) + return "\n\n".join(blocks) def main() -> int: @@ -151,8 +118,7 @@ def main() -> int: items = _load_items(args.input_file) if not items: - # Nothing to classify -> silent. This is the common quiet-interval case. - return 0 + return 0 # nothing to classify -> silent (the common quiet-interval case) # Import here so --help works without the package importable. try: @@ -173,8 +139,7 @@ def main() -> int: if not isinstance(content, str): content = str(content) if content else "" except Exception as e: - # Classification failure is NOT silent -- surface it so a broken monitor - # doesn't quietly swallow important items. Non-zero exit -> cron alerts. + # A broken monitor must not quietly swallow important items: non-zero exit -> cron alerts. _eprint(f"classify_items: classifier call failed: {e}") return 4 @@ -187,38 +152,16 @@ def main() -> int: surfaced.append((i, item, s)) if not surfaced: - # Below threshold -> silent. Empty stdout; cron suppresses delivery. - return 0 + return 0 # below threshold -> silent; empty stdout suppresses delivery if args.format == "json": out = [ - { - "id": _item_id(item, i), - "score": s.get("score"), - "reason": s.get("reason", ""), - "item": item, - } + {"id": _item_id(item, i), "score": s.get("score"), "reason": s.get("reason", ""), "item": item} for (i, item, s) in surfaced ] print(json.dumps(out, ensure_ascii=False, indent=2)) else: - blocks = [] - for (i, item, s) in surfaced: - title = ( - item.get("title") - or item.get("subject") - or item.get("summary") - or _item_id(item, i) - ) - url = item.get("url") or item.get("link") or "" - reason = s.get("reason", "") - block = f"## [{s.get('score')}/10] {title}" - if url: - block += f"\n{url}" - if reason: - block += f"\n_{reason}_" - blocks.append(block) - print("\n\n".join(blocks)) + print(_render_text(surfaced)) return 0 diff --git a/cron/suggestion_catalog.py b/cron/suggestion_catalog.py index 39f77ba420..4cfeb8eea8 100644 --- a/cron/suggestion_catalog.py +++ b/cron/suggestion_catalog.py @@ -20,7 +20,7 @@ __all__ = ["CatalogEntry", "CATALOG", "seed_catalog_suggestions", "classify_item def classify_items_script_path() -> str: """Absolute path to the urgency classifier script shipped with cron/.""" - return str((Path(__file__).resolve().parent / "scripts" / "classify_items.py")) + return str(Path(__file__).resolve().parent / "scripts" / "classify_items.py") @dataclass(frozen=True) @@ -136,11 +136,8 @@ def seed_catalog_suggestions( if wanted is not None and entry.key not in wanted: continue rec = add_fn( - title=entry.title, - description=entry.description, - source="catalog", - job_spec=dict(entry.job_spec), - dedup_key=entry.key, + title=entry.title, description=entry.description, source="catalog", + job_spec=dict(entry.job_spec), dedup_key=entry.key, ) if rec is not None: created.append(rec)