feat(loop): add main-loop brain module for opm sidechat thinking and operator steering
This commit is contained in:
+413
@@ -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 <ts>] marker, NEVER a [JOB <id>] 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": <expires_epoch>}, "quiet_until": <epoch|0>,
|
||||||
|
"brain_watermark": <epoch>, "tick": <int>}}
|
||||||
|
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> <dur> (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 <agent>", 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> (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=<harvester reader>)
|
||||||
|
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
|
||||||
Reference in New Issue
Block a user