Fix self_main_loop exit code, [!] false positive, state race
- Exit 0 on success (was 1 when prompts sent, confusing systemd) - [!] flag now uses word-boundary regex excluding hyphenated identities (operator-646 no longer triggers; 'operator needed' does) - State writes now hold fcntl exclusive lock with fresh reload, preventing timer check from clobbering enable/disable changes Session: sidechat/chromebox-ops
This commit is contained in:
+41
-10
@@ -21,7 +21,7 @@ First run per agent starts at the current newest message (no backfill spam).
|
||||
|
||||
Contract:
|
||||
- No new messages -> silent (no sidechat post), exit 0.
|
||||
- Digests sent -> exit 1.
|
||||
- Digests sent -> exit 0 (success; prompted count in JSON).
|
||||
- Read/send error -> logged to stderr, exit 2 (timer stays alive).
|
||||
- Overlap guard: fcntl LOCK_EX|LOCK_NB on a lockfile; a second
|
||||
concurrent run prints {"ok": false, "skipped": "already running"}
|
||||
@@ -34,8 +34,10 @@ stderr (systemd journal).
|
||||
import fcntl
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import subprocess
|
||||
import sys
|
||||
from contextlib import contextmanager
|
||||
from datetime import datetime, timezone
|
||||
|
||||
BASE = "/home/super/Projects/NetVM"
|
||||
@@ -44,6 +46,7 @@ BOX_CHAT = os.path.join(BIN, "box-chat.py")
|
||||
DM_PY = os.path.join(BIN, "dm.py")
|
||||
STATE_FILE = os.path.join(BASE, "self-main-loop-watermark.json")
|
||||
LOCK_FILE = os.path.join(BASE, "self-main-loop.lock")
|
||||
STATE_LOCK_FILE = os.path.join(BASE, "self-main-loop-state.lock")
|
||||
|
||||
DEFAULT_AGENTS = ["muse", "pip", "646", "opm"]
|
||||
# Prompting sidechat per agent — mirrors box-ctl.py NOTIFY_SIDECHATS.
|
||||
@@ -60,6 +63,28 @@ SEND_TIMEOUT = 180
|
||||
DIGEST_MAX = 600 # well under dm.py's 1000-char non-raw truncation
|
||||
PREVIEW_MAX = 120
|
||||
|
||||
# Word-boundary "operator" that does NOT match hyphenated identities like
|
||||
# "operator-646" (the "-" counts as a boundary for \b, so we exclude it
|
||||
# explicitly). "operator needed" matches; "operator-646" does not.
|
||||
OPERATOR_RE = re.compile(r"(?<![\w-])operator(?![\w-])", re.IGNORECASE)
|
||||
|
||||
|
||||
@contextmanager
|
||||
def state_locked():
|
||||
"""Exclusive lock for load-modify-save on the state file.
|
||||
|
||||
Prevents the timer's check (mid-run save) from clobbering an
|
||||
enable/disable change that landed during the slow chat reads.
|
||||
Held only for the brief critical section, never across network I/O.
|
||||
"""
|
||||
fh = open(STATE_LOCK_FILE, "w")
|
||||
try:
|
||||
fcntl.flock(fh, fcntl.LOCK_EX)
|
||||
yield
|
||||
finally:
|
||||
fcntl.flock(fh, fcntl.LOCK_UN)
|
||||
fh.close()
|
||||
|
||||
|
||||
def utcnow():
|
||||
return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
||||
@@ -152,7 +177,7 @@ def compose_digest(agent, new):
|
||||
for m in new[:5]:
|
||||
text = (m.get("text") or "").strip().replace("\n", " ")
|
||||
q = " [?]" if "?" in text else ""
|
||||
hot = " [!]" if ("operator" in text.lower() or "urgent" in text.lower()) else ""
|
||||
hot = " [!]" if (OPERATOR_RE.search(text) or "urgent" in text.lower()) else ""
|
||||
lines.append("- %s: %s%s%s" % (author_of(m), text[:PREVIEW_MAX], q, hot))
|
||||
if len(new) > 5:
|
||||
lines.append("(+%d more)" % (len(new) - 5))
|
||||
@@ -231,19 +256,24 @@ def do_check(only_agent=None):
|
||||
results[agent] = {"ok": False, "error": detail, "new": len(new)}
|
||||
errors += 1
|
||||
|
||||
st["agents"] = agents_state
|
||||
st["last_run"] = utcnow()
|
||||
st["last_result"] = {"prompted": prompted, "errors": errors,
|
||||
# Save under the state lock with a fresh reload: an enable/disable may
|
||||
# have landed during the slow chat reads; preserve its config changes
|
||||
# and only update the keys this run owns (watermarks, last_run/result).
|
||||
with state_locked():
|
||||
fresh = load_state()
|
||||
fresh["agents"] = agents_state
|
||||
fresh["last_run"] = utcnow()
|
||||
fresh["last_result"] = {"prompted": prompted, "errors": errors,
|
||||
"agents": results}
|
||||
# Preserve operator-edited config; seed defaults on first run.
|
||||
if "config" not in st:
|
||||
st["config"] = {"agents": cfg["agents"],
|
||||
# Seed defaults on first run; otherwise preserve operator-edited config.
|
||||
if "config" not in fresh:
|
||||
fresh["config"] = {"agents": cfg["agents"],
|
||||
"prompt_sidechat": cfg["prompt_sidechat"],
|
||||
"sender": cfg["sender"],
|
||||
"read_limit": cfg["read_limit"],
|
||||
"enabled": dict(cfg.get("enabled") or
|
||||
{a: True for a in cfg["agents"]})}
|
||||
save_state(st)
|
||||
save_state(fresh)
|
||||
return {"ok": errors == 0, "prompted": prompted, "errors": errors,
|
||||
"new_total": sum(r.get("new", 0) for r in results.values()),
|
||||
"agents": results}
|
||||
@@ -251,6 +281,7 @@ def do_check(only_agent=None):
|
||||
|
||||
def _set_enabled(agent, value):
|
||||
"""Enable/disable the loop for one agent (or all if agent is None)."""
|
||||
with state_locked():
|
||||
st = load_state()
|
||||
cfg = st.get("config") or {}
|
||||
agents = cfg.get("agents") or list(DEFAULT_AGENTS)
|
||||
@@ -320,7 +351,7 @@ def main(argv):
|
||||
print(json.dumps(result))
|
||||
if not result.get("ok"):
|
||||
return 2
|
||||
return 1 if result.get("prompted") else 0
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
Reference in New Issue
Block a user