131 lines
3.6 KiB
Python
131 lines
3.6 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
Side-chat to main-chat work siphon — siphon action (INTEGRATED).
|
|
|
|
Changes vs original (integrator; agent 3 of 5 never delivered, so the
|
|
minimal reversible versions below stand in for its authorship/dedup work):
|
|
- format_siphon imported from detect (single definition; honest labeling
|
|
+ honest authorship live there).
|
|
- Deduplication is PERSISTENT: siphoned message IDs are stored as JSON
|
|
on disk (SIPHON_STATE_DIR or ~/.siphon-state/siphoned_ids.json) so a
|
|
restart can never re-relay. In-memory set kept as a fast path.
|
|
- mark_siphoned() writes through to disk on every call.
|
|
|
|
Safety (unchanged):
|
|
- Rate limited (max N siphons per hour per thread)
|
|
- Never posts full message content
|
|
- Respects opt-out registry
|
|
- Deduplicates (same message_id never siphoned twice, even across restarts)
|
|
"""
|
|
|
|
import json
|
|
import os
|
|
import time
|
|
from dataclasses import dataclass, field
|
|
from typing import Callable, Optional
|
|
|
|
from detect import SiphonHit, is_opted_out, format_siphon # noqa: F401 (re-export)
|
|
|
|
|
|
# --- Rate limiting ---
|
|
|
|
@dataclass
|
|
class RateLimiter:
|
|
max_per_hour: int = 5
|
|
_timestamps: dict = field(default_factory=dict) # thread_id -> [ts, ...]
|
|
|
|
def allow(self, thread_id: str) -> bool:
|
|
now = time.time()
|
|
stamps = self._timestamps.get(thread_id, [])
|
|
# Prune older than 1 hour
|
|
stamps = [s for s in stamps if now - s < 3600]
|
|
if len(stamps) >= self.max_per_hour:
|
|
return False
|
|
stamps.append(now)
|
|
self._timestamps[thread_id] = stamps
|
|
return True
|
|
|
|
|
|
# --- Deduplication (persistent) ---
|
|
|
|
_STATE_DIR = os.environ.get(
|
|
"SIPHON_STATE_DIR", os.path.expanduser("~/.siphon-state"))
|
|
DEDUP_FILE = os.path.join(_STATE_DIR, "siphoned_ids.json")
|
|
_DEDUP_MAX_IDS = 5000 # bound disk growth; oldest evicted first
|
|
|
|
|
|
def _load_siphoned() -> set:
|
|
try:
|
|
with open(DEDUP_FILE) as f:
|
|
data = json.load(f)
|
|
ids = data.get("ids", []) if isinstance(data, dict) else []
|
|
return set(ids)
|
|
except (OSError, ValueError):
|
|
return set()
|
|
|
|
|
|
def _save_siphoned(ids: set) -> None:
|
|
try:
|
|
os.makedirs(_STATE_DIR, exist_ok=True)
|
|
trimmed = sorted(ids)[-_DEDUP_MAX_IDS:]
|
|
tmp = DEDUP_FILE + ".tmp"
|
|
with open(tmp, "w") as f:
|
|
json.dump({"ids": trimmed, "updated": time.time()}, f)
|
|
os.replace(tmp, DEDUP_FILE)
|
|
except OSError:
|
|
pass # dedup degrades to in-memory; never crash the relay on IO
|
|
|
|
|
|
_siphoned_ids: set = _load_siphoned()
|
|
|
|
|
|
def already_siphoned(message_id: str) -> bool:
|
|
return message_id in _siphoned_ids
|
|
|
|
|
|
def mark_siphoned(message_id: str):
|
|
_siphoned_ids.add(message_id)
|
|
_save_siphoned(_siphoned_ids)
|
|
|
|
|
|
def siphoned_count() -> int:
|
|
return len(_siphoned_ids)
|
|
|
|
|
|
# --- Siphon action ---
|
|
|
|
CATEGORY_EMOJI = {
|
|
"COMPLETED": "✅",
|
|
"BLOCKER": "🚧",
|
|
"DECISION": "❓",
|
|
"ALERT": "🚨",
|
|
"MILESTONE": "🎯",
|
|
}
|
|
|
|
|
|
def siphon(hit: SiphonHit,
|
|
agent_name: str,
|
|
post_to_main: Callable[[str], bool],
|
|
limiter: Optional[RateLimiter] = None) -> bool:
|
|
"""
|
|
Execute the siphon: post summary to main chat.
|
|
|
|
post_to_main: callable that posts text to main chat, returns True on success.
|
|
Returns True if siphoned, False if suppressed.
|
|
"""
|
|
# Safety checks
|
|
if is_opted_out(hit.thread_id):
|
|
return False
|
|
if already_siphoned(hit.message_id):
|
|
return False
|
|
|
|
lim = limiter or RateLimiter()
|
|
if not lim.allow(hit.thread_id):
|
|
return False
|
|
|
|
text = format_siphon(hit, agent_name)
|
|
ok = post_to_main(text)
|
|
if ok:
|
|
mark_siphoned(hit.message_id)
|
|
return ok
|