89 lines
2.5 KiB
Python
89 lines
2.5 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
Side-chat to main-chat work siphon — monitor loop.
|
|
|
|
Polls side chats for new messages, runs detection, siphons hits to main.
|
|
|
|
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
|
|
|
|
For the prototype, all three are injectable (see tests).
|
|
"""
|
|
|
|
import time
|
|
from typing import Callable, Dict, List
|
|
|
|
from detect import detect, is_opted_out
|
|
from siphon import siphon, RateLimiter
|
|
|
|
|
|
# Message shape: {"id": str, "text": str, "author": str, "ts": str}
|
|
Message = Dict[str, str]
|
|
|
|
|
|
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)
|
|
|
|
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
|
|
|
|
hit = detect(text, tid, mid, min_confidence)
|
|
if hit:
|
|
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)
|