""" 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