Files
box/bin/monitor.py
T

155 lines
5.0 KiB
Python
Raw Normal View History

#!/usr/bin/env python3
"""
Side-chat to main-chat siphon — monitor loop (INTEGRATED).
Changes vs original (integrator):
1. Timestamp plumbing (agent 2's open item): message["ts"] is parsed to
epoch seconds and passed as message_ts to detect(), enabling the
15-minute stale-suppression for COMPLETED. Unparseable/missing ts →
backward-compatible (detect proceeds).
2. Author plumbing (agent 3 absent): message["author"] is attached to
the hit as hit.author, so relays attribute the real author instead
of the thread's registered agent.
3. Flood control (agent 4 absent): COMPLETED hits are routed to the
digest buffer instead of individual main-chat relays. ALERT, BLOCKER,
DECISION, MILESTONE still relay individually via siphon().
4. Persistent dedup: every processed hit is marked siphoned (including
digested ones) so a restart never re-relays or re-digests.
This is the integration point for bl. In production:
- list_sidechats() calls muse-chat-api.py or the sidechat manager
- get_messages() reads thread messages via CDP
- post_to_main() sends via muse-chat-api.py send to main chat
- flush_digest() should be called on a schedule (e.g. every 30 min) and
its output posted to main chat once.
"""
import time
from datetime import datetime, timezone
from typing import Callable, Dict, List, Optional
from detect import detect, is_opted_out
from siphon import siphon, mark_siphoned, already_siphoned, RateLimiter
try:
from digest import get_buffer, flush_digest # noqa: F401 (re-export)
except ImportError: # pragma: no cover — digest module optional
get_buffer = None
def flush_digest():
return None
# Message shape: {"id": str, "text": str, "author": str, "ts": str}
Message = Dict[str, str]
# Categories that batch into the digest instead of relaying individually.
DIGESTED_CATEGORIES = {"COMPLETED"}
def _parse_ts(ts) -> Optional[float]:
"""Parse a message timestamp to epoch seconds. None if unparseable."""
if ts is None:
return None
if isinstance(ts, (int, float)):
return float(ts)
s = str(ts).strip()
if not s:
return None
# Epoch as string?
try:
return float(s)
except ValueError:
pass
# ISO-8601 (with optional Z suffix)?
try:
iso = s.replace("Z", "+00:00")
dt = datetime.fromisoformat(iso)
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return dt.timestamp()
except ValueError:
return None
def monitor_once(
list_sidechats: Callable[[], List[Dict[str, str]]],
get_messages: Callable[[str, str], List[Message]],
post_to_main: Callable[[str], bool],
watermarks: Dict[str, str],
limiter: RateLimiter = None,
min_confidence: float = 0.6,
) -> Dict[str, str]:
"""
One poll cycle. Returns updated watermarks.
list_sidechats: () -> [{"id": thread_id, "name": str, "agent": str}]
get_messages: (thread_id, since_msg_id) -> [messages newer than watermark]
post_to_main: (text) -> True on success
watermarks: {thread_id: last_seen_message_id}
"""
lim = limiter or RateLimiter()
new_marks = dict(watermarks)
digest = get_buffer() if get_buffer else None
for chat in list_sidechats():
tid = chat["id"]
agent = chat.get("agent", "unknown")
if is_opted_out(tid):
continue
since = watermarks.get(tid, "")
try:
messages = get_messages(tid, since)
except Exception:
continue # don't let one bad chat kill the cycle
for msg in messages:
mid = msg.get("id", "")
text = msg.get("text", "")
if not mid or not text:
continue
# Update watermark to newest seen
new_marks[tid] = mid
# Persistent dedup first: never reprocess a seen message,
# even across restarts (marks are set for digested hits too).
if already_siphoned(mid):
continue
message_ts = _parse_ts(msg.get("ts"))
hit = detect(text, tid, mid, min_confidence,
message_ts=message_ts)
if hit is None:
continue
# Author plumbing: real author, never thread-owner-as-author.
hit.author = msg.get("author", "") or ""
if hit.category in DIGESTED_CATEGORIES and digest is not None:
digest.add(hit)
mark_siphoned(mid)
else:
siphon(hit, agent, post_to_main, lim)
return new_marks
def monitor_loop(
list_sidechats,
get_messages,
post_to_main,
poll_interval: int = 60,
watermarks: Dict[str, str] = None,
):
"""Run forever. For production use with systemd timer instead."""
marks = watermarks or {}
limiter = RateLimiter()
while True:
marks = monitor_once(
list_sidechats, get_messages, post_to_main, marks, limiter
)
time.sleep(poll_interval)