From c74ebaa65837add65a06c1b4c2d2d7532544322c Mon Sep 17 00:00:00 2001 From: operator Date: Mon, 5 Oct 2026 16:02:58 +0000 Subject: [PATCH] feat(loop): add main-loop brain module for opm sidechat thinking and operator steering --- bin/brain.py | 413 +++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 413 insertions(+) create mode 100644 bin/brain.py diff --git a/bin/brain.py b/bin/brain.py new file mode 100644 index 0000000..c66a5f1 --- /dev/null +++ b/bin/brain.py @@ -0,0 +1,413 @@ +""" +brain.py — Main-loop brain workspace module. + +The brain sidechat ("main-loop brain" on opm's account) is where the main loop +OPERATES: it posts its thinking there, reads operator instructions from there, +and keeps its working state visible there. + +Deploy to: ~/Projects/NetVM/bin/brain.py on bl (alongside self_main_loop.py). + +Design source: ~/workspace/main-loop-brain-design.md +User directive 2026-10-04: "we need main loop to operate its brains in side chat" + +SAFETY CONTRACT (do not weaken): + - Brain posts carry a [BRAIN ] marker, NEVER a [JOB ] marker. + The response-harvester keys off [JOB ...]; a thinking note must never + look actionable. + - Brain posts NEVER use dm.py --expect-reply. Thinking creates no followup + records, no nudges, no escalations. + - When reading the brain, the loop skips its own messages (sender check). + The loop must never digest its own thinking as agent activity. + - !loop commands are honored ONLY from AUTHORIZED_SENDERS. Everything else + is read as context, never as instruction. + - Malformed commands get a one-line correction posted to the brain. + Never silent, never a crash. + - Cap: MAX_BRAIN_POSTS_PER_TICK posts per tick. The brain must not amplify. + +State lives in the existing watermark JSON file under the "brain" key: + {"brain": {"ignores": {"646": }, "quiet_until": , + "brain_watermark": , "tick": }} +Absolute expiries so a dead loop cannot leave an agent ignored forever. +""" + +import json +import os +import re +import subprocess +import time + +# --------------------------------------------------------------------------- +# Constants +# --------------------------------------------------------------------------- + +BRAIN_AGENT = "opm" +BRAIN_SIDECHAT_NAME = "main-loop brain" + +# Operator identities allowed to issue !loop commands. These must match the +# DM-signer / board identity strings; do not invent new ones here. +AUTHORIZED_SENDERS = {"super", "operator-646", "operator-main"} + +# Sender identities the loop itself posts under (skipped on read-back). +OWN_SENDERS = {"main-loop", "self_main_loop", "operator-main-loop", "opm"} + +BRAIN_MARKER_RE = re.compile(r"\[BRAIN\s+([^\]]+)\]") +JOB_MARKER_RE = re.compile(r"\[JOB\s+([^\]]+)\]") +LOOP_CMD_RE = re.compile(r"^\s*!loop\s+(\S+)(.*)$", re.IGNORECASE) + +MAX_BRAIN_POSTS_PER_TICK = 4 +DEFAULT_BIN_DIR = os.path.expanduser("~/Projects/NetVM/bin") + +VALID_AGENTS = {"muse", "pip", "646", "opm"} + +DUR_RE = re.compile(r"^\s*(\d+)\s*([smh])\s*$", re.IGNORECASE) + + +# --------------------------------------------------------------------------- +# Pure logic — fully unit-testable, no I/O +# --------------------------------------------------------------------------- + +def parse_duration(text): + """'30m' -> 1800.0, '1h' -> 3600.0, '90s' -> 90.0. None if malformed.""" + m = DUR_RE.match(text or "") + if not m: + return None + n, unit = int(m.group(1)), m.group(2).lower() + return float(n * {"s": 1, "m": 60, "h": 3600}[unit]) + + +def is_own_message(sender): + s = (sender or "").strip().lower() + return s in {x.lower() for x in OWN_SENDERS} + + +def is_authorized(sender): + return (sender or "").strip() in AUTHORIZED_SENDERS + + +def extract_loop_commands(messages): + """messages: list of {"sender": str, "text": str, "ts": float}. + Returns [(sender, verb, args, msg)] for !loop lines from any sender + (authorization is applied by the caller so corrections can name names).""" + cmds = [] + for msg in messages: + text = msg.get("text") or "" + for line in text.splitlines(): + m = LOOP_CMD_RE.match(line) + if m: + cmds.append((msg.get("sender", "?"), + m.group(1).lower(), + m.group(2).strip(), + msg)) + return cmds + + +def prune_expired(state, now=None): + """Drop expired ignores / quiet. Returns (notes, changed).""" + now = now if now is not None else time.time() + notes, changed = [], False + ignores = state.setdefault("ignores", {}) + for agent in list(ignores): + if ignores[agent] <= now: + del ignores[agent] + notes.append("resuming %s (ignore expired)" % agent) + changed = True + if state.get("quiet_until", 0) and state["quiet_until"] <= now: + state["quiet_until"] = 0 + notes.append("quiet period ended, prompts resumed") + changed = True + return notes, changed + + +def is_ignored(state, agent, now=None): + now = now if now is not None else time.time() + return state.get("ignores", {}).get(agent, 0) > now + + +def is_quiet(state, now=None): + now = now if now is not None else time.time() + return (state.get("quiet_until", 0) or 0) > now + + +def apply_command(sender, verb, args, state, now=None): + """Apply one !loop command. Returns (ack_text, changed).""" + now = now if now is not None else time.time() + changed = False + + if verb == "ignore": + parts = args.split() + if len(parts) != 2 or parts[0] not in VALID_AGENTS: + return ("usage: !loop ignore (agent: %s, dur like 30m/1h)" + % "/".join(sorted(VALID_AGENTS)), False) + dur = parse_duration(parts[1]) + if dur is None or dur <= 0: + return ("bad duration %r, try 30m or 1h" % parts[1], False) + state.setdefault("ignores", {})[parts[0]] = now + dur + return ("ignoring %s for %s (until %s)" + % (parts[0], parts[1], + time.strftime("%H:%M UTC", time.gmtime(now + dur))), True) + + if verb == "unignore": + agent = args.strip() + if agent not in VALID_AGENTS: + return ("usage: !loop unignore ", False) + if state.get("ignores", {}).pop(agent, None) is not None: + changed = True + return ("resuming %s now" % agent, True) + return ("%s was not ignored" % agent, False) + + if verb == "quiet": + dur = parse_duration(args) + if dur is None or dur <= 0: + return ("usage: !loop quiet (dur like 30m/1h)", False) + state["quiet_until"] = now + dur + changed = True + return ("quiet for %s: reads continue, digests land here, " + "no per-agent escalation" % args.strip(), True) + + if verb == "unquiet": + if state.get("quiet_until"): + state["quiet_until"] = 0 + changed = True + return ("prompts resumed", True) + return ("was not quiet", False) + + if verb == "check": + # The tick loop honors this by running the agent reads immediately + # rather than waiting for the next timer fire. Caller sets the flag. + return ("CHECK_REQUESTED", True) + + if verb == "status": + return ("STATUS_REQUESTED", False) + + return ("unknown command %r, try: ignore, unignore, quiet, unquiet, " + "check, status" % verb, False) + + +def format_thinking(tick, summary, state, now=None): + """summary: dict with keys seen{agent:(new,q,urgent)}, escalated[JOB ids], + skipped_info[int], errors[int], closure_rate[float|None].""" + now = now if now is not None else time.time() + ts = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime(now)) + seen_bits = [] + for agent in ("646", "pip", "muse", "opm"): + new, q, u = summary.get("seen", {}).get(agent, (0, 0, 0)) + bit = "%s:%dnew" % (agent, new) + if q: + bit += "(%d?)" % q + if u: + bit += "(%d!)" % u + seen_bits.append(bit) + esc = summary.get("escalated", []) + lines = [ + "[BRAIN %s tick=%d]" % (ts, tick), + "seen: " + ", ".join(seen_bits), + "decided: escalated %d%s | skipped %d info | ignored: %s" % ( + len(esc), + (" (%s)" % ", ".join(esc[:3])) if esc else "", + summary.get("skipped_info", 0), + ", ".join(sorted(state.get("ignores", {}))) or "none"), + "errors: %d | quiet: %s | closure: %s" % ( + summary.get("errors", 0), + "yes" if is_quiet(state, now) else "no", + ("%.2f" % summary["closure_rate"]) + if summary.get("closure_rate") is not None else "n/a"), + ] + return "\n".join(lines) + + +def format_status(tick, cfg, state, watermarks, health, now=None): + now = now if now is not None else time.time() + ts = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime(now)) + lines = ["[BRAIN %s tick=%d] status" % (ts, tick)] + lines.append("agents: " + ", ".join( + "%s=%s" % (a, "on" if cfg.get(a) else "off") + for a in ("muse", "pip", "646", "opm"))) + ign = state.get("ignores", {}) + lines.append("ignores: " + (", ".join( + "%s until %s" % (a, time.strftime("%H:%M UTC", time.gmtime(e))) + for a, e in sorted(ign.items())) or "none")) + q = state.get("quiet_until", 0) + lines.append("quiet: " + ("until %s" % time.strftime("%H:%M UTC", time.gmtime(q)) + if q and q > now else "no")) + ages = [] + for a in ("muse", "pip", "646", "opm"): + wm = (watermarks or {}).get(a, 0) + ages.append("%s:%dm" % (a, int((now - wm) / 60)) if wm else "%s:never" % a) + lines.append("watermark age: " + ", ".join(ages)) + if health: + lines.append("digest health: delivered=%d acked=%d closed=%d stale=%d " + "closure=%.2f" % ( + health.get("delivered", 0), health.get("acked", 0), + health.get("closed", 0), health.get("stale", 0), + health.get("closure_rate", 0.0))) + return "\n".join(lines) + + +# --------------------------------------------------------------------------- +# I/O adapters — thin shells over the bl tools. Verify paths on bl. +# --------------------------------------------------------------------------- + +class BrainIO: + """Send/receive against the brain sidechat. + + send: dm.py, WITHOUT --expect-reply (thinking is never actionable). + read: pluggable read_fn(messages-since-ts); defaults to None and must be + wired to the same chat-read primitive the response-harvester uses. + """ + + def __init__(self, bin_dir=DEFAULT_BIN_DIR, read_fn=None, + sidechat_map_path=None): + self.bin_dir = bin_dir + self.read_fn = read_fn + self.sidechat_map_path = (sidechat_map_path or + os.path.join(bin_dir, "..", + "job-sidechats.json")) + self._posts_this_tick = 0 + + # -- name resolution: never hardcode the thread UUID ------------------- + def brain_thread_uuid(self): + """Resolve 'main-loop brain' -> UUID via job-sidechats.json.""" + try: + with open(os.path.normpath(self.sidechat_map_path)) as f: + data = json.load(f) + except (OSError, ValueError): + return None + # schema: {"sidechats": {"main-loop brain": {"uuid": ...}}} or flat + node = data.get("sidechats", data).get(BRAIN_SIDECHAT_NAME) + if isinstance(node, dict): + return node.get("uuid") or node.get("thread_uuid") + return node if isinstance(node, str) else None + + # -- write -------------------------------------------------------------- + def reset_tick_budget(self): + self._posts_this_tick = 0 + + def post(self, text): + """Post thinking/acks to the brain. Returns True on VERIFIED send.""" + if self._posts_this_tick >= MAX_BRAIN_POSTS_PER_TICK: + return False + if JOB_MARKER_RE.search(text): + raise ValueError("refusing to post [JOB ...] to the brain") + dm = os.path.join(self.bin_dir, "dm.py") + # Same path the loop uses for opm digests; NO --expect-reply: + # thinking is never actionable and must not create followups. + cmd = [dm, "send", "--agent", BRAIN_AGENT, + "--to", BRAIN_AGENT, "--target", BRAIN_SIDECHAT_NAME, + "--message", text] + try: + p = subprocess.run(cmd, capture_output=True, text=True, + timeout=120) + except (OSError, subprocess.TimeoutExpired): + return False + ok = "VERIFIED" in (p.stdout or "") + if ok: + self._posts_this_tick += 1 + return ok + + # -- read --------------------------------------------------------------- + def read_new(self, since_ts): + """Return [{"sender","text","ts"}] newer than since_ts. Skips own.""" + if self.read_fn is None: + return [] + try: + msgs = self.read_fn(BRAIN_AGENT, BRAIN_SIDECHAT_NAME, since_ts) or [] + except Exception: + return [] + return [m for m in msgs + if (m.get("ts", 0) or 0) > since_ts + and not is_own_message(m.get("sender"))] + + +# --------------------------------------------------------------------------- +# Workspace — one object per tick +# --------------------------------------------------------------------------- + +class BrainWorkspace: + """Owns the brain side of a main-loop tick. + + Usage in self_main_loop.py tick(): + brain = BrainWorkspace(watermark_path, read_fn=) + cmds_outcome = brain.intake() # read, parse, apply, ack + ... existing per-agent reads, skipping brain.ignored(agent) ... + brain.post_thinking(summary) # the loop's reasoning, visible + brain.save() + """ + + def __init__(self, watermark_path, read_fn=None, bin_dir=DEFAULT_BIN_DIR): + self.watermark_path = watermark_path + self.io = BrainIO(bin_dir=bin_dir, read_fn=read_fn) + self._data = self._load() + self.state = self._data.setdefault("brain", {}) + self.io.reset_tick_budget() + self.check_requested = False + self.status_requested = False + + # -- persistence --------------------------------------------------------- + def _load(self): + try: + with open(self.watermark_path) as f: + return json.load(f) + except (OSError, ValueError): + return {} + + def save(self): + tmp = self.watermark_path + ".tmp" + with open(tmp, "w") as f: + json.dump(self._data, f, indent=2) + os.replace(tmp, self.watermark_path) + + # -- tick intake: read -> parse -> apply -> ack -------------------------- + def intake(self): + """Process new brain messages. Returns dict of what happened.""" + now = time.time() + outcome = {"commands": 0, "acks": 0, "expired_notes": 0, + "ignored_senders": 0} + since = float(self.state.get("brain_watermark", 0)) + msgs = self.io.read_new(since) + newest = since + for m in msgs: + newest = max(newest, float(m.get("ts", 0) or 0)) + if msgs: + self.state["brain_watermark"] = newest + + notes, _ = prune_expired(self.state, now) + for n in notes: + if self.io.post("[BRAIN %s] %s" % ( + time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime(now)), n)): + outcome["expired_notes"] += 1 + + for sender, verb, args, _msg in extract_loop_commands(msgs): + if not is_authorized(sender): + outcome["ignored_senders"] += 1 + continue + outcome["commands"] += 1 + ack, _changed = apply_command(sender, verb, args, self.state, now) + if ack == "CHECK_REQUESTED": + self.check_requested = True + ack = "out-of-cycle check armed for this tick" + elif ack == "STATUS_REQUESTED": + self.status_requested = True + continue # status posts at end of tick with full context + if self.io.post("[BRAIN %s] @%s %s" % ( + time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime(now)), + sender, ack)): + outcome["acks"] += 1 + return outcome + + # -- state queries for the tick loop -------------------------------------- + def ignored(self, agent): + return is_ignored(self.state, agent) + + def quiet(self): + return is_quiet(self.state) + + # -- end of tick ----------------------------------------------------------- + def post_thinking(self, summary, cfg=None, watermarks=None, health=None): + tick = int(self.state.get("tick", 0)) + 1 + self.state["tick"] = tick + ok = self.io.post(format_thinking(tick, summary, self.state)) + if self.status_requested and cfg is not None: + self.io.post(format_status(tick, cfg, self.state, + watermarks, health)) + self.status_requested = False + return ok