refactor(cron): dedupe blueprint renderer slot lookups, compact classify_items/monitor/suggestion_catalog, drop dead _CLASSIFY_INSTRUCTIONS

This commit is contained in:
Teknium
2026-09-02 19:26:38 -07:00
parent b3d9b20a3b
commit f6460e681e
5 changed files with 103 additions and 188 deletions
+4 -15
View File
@@ -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",
+33 -43
View File
@@ -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 <key> 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 <key> 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:
+18 -22
View File
@@ -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:
+45 -102
View File
@@ -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": <int>, "score": <int 0-10>, "reason": "<short>"}]. '
"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
+3 -6
View File
@@ -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)