#!/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)