2026-10-04 22:54:01 +00:00
|
|
|
#!/usr/bin/env python3
|
|
|
|
|
"""
|
2026-10-07 00:25:51 +00:00
|
|
|
Side-chat to main-chat siphon — monitor loop (INTEGRATED).
|
2026-10-04 22:54:01 +00:00
|
|
|
|
2026-10-07 00:25:51 +00:00
|
|
|
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.
|
2026-10-04 22:54:01 +00:00
|
|
|
|
|
|
|
|
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
|
2026-10-07 00:25:51 +00:00
|
|
|
- flush_digest() should be called on a schedule (e.g. every 30 min) and
|
|
|
|
|
its output posted to main chat once.
|
2026-10-04 22:54:01 +00:00
|
|
|
"""
|
|
|
|
|
|
|
|
|
|
import time
|
2026-10-07 00:25:51 +00:00
|
|
|
from datetime import datetime, timezone
|
|
|
|
|
from typing import Callable, Dict, List, Optional
|
2026-10-04 22:54:01 +00:00
|
|
|
|
|
|
|
|
from detect import detect, is_opted_out
|
2026-10-07 00:25:51 +00:00
|
|
|
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
|
2026-10-04 22:54:01 +00:00
|
|
|
|
|
|
|
|
|
|
|
|
|
# Message shape: {"id": str, "text": str, "author": str, "ts": str}
|
|
|
|
|
Message = Dict[str, str]
|
|
|
|
|
|
2026-10-07 00:25:51 +00:00
|
|
|
# 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
|
|
|
|
|
|
2026-10-04 22:54:01 +00:00
|
|
|
|
|
|
|
|
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)
|
2026-10-07 00:25:51 +00:00
|
|
|
digest = get_buffer() if get_buffer else None
|
2026-10-04 22:54:01 +00:00
|
|
|
|
|
|
|
|
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
|
|
|
|
|
|
2026-10-07 00:25:51 +00:00
|
|
|
# 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:
|
2026-10-04 22:54:01 +00:00
|
|
|
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)
|