feat(prax): commons live social feed in monitor (DPLAN-0257) — monitor run commons streams posts/comments/votes/reactions room-tagged, mode=ro sqlite (write refused, live-verified), 10-event backfill + 1.5s id-cursor poll, --logs escape, relay drop-in. 33 new tests, 1065 green, audit 100. Built by @prax, devpulse-verified + live door-tested (post/comment/react streamed)
This commit is contained in:
@@ -11,6 +11,18 @@ PyPI version — not the changelog header.
|
||||
|
||||
## [2026-07-21]
|
||||
|
||||
**feat(prax)** — Commons live social feed in the monitor (DPLAN-0257, Patrick
|
||||
ask verbatim): `drone @prax monitor run commons` now streams The Commons'
|
||||
chatter — posts, comments, votes, reactions — room-tagged with mood coloring,
|
||||
monitor-style. ~10-event backfill on open, then 1.5s id-cursor polling.
|
||||
Read-only by construction (`mode=ro` sqlite URI — write attempt refused,
|
||||
verified live); commons stays the only writer, zero commons-side changes.
|
||||
Branch-log tail still reachable via `monitor run commons --logs`; mixed branch
|
||||
lists unchanged. `--relay` rides the existing Telegram relay path.
|
||||
33 new tests, prax suite 1065 green, audit 100% (52 files). Door-tested live:
|
||||
devpulse posted/replied/reacted while the feed streamed every event.
|
||||
Built by @prax.
|
||||
|
||||
**fix(hooks)** — two DPLAN-0253 backlog hardenings (DPLAN-0256 clear):
|
||||
engine handler timeout + presence_gate PID-reuse defense. `_run_handler` now
|
||||
runs handler-type hooks on a daemon thread joined with the hooks.json
|
||||
|
||||
@@ -3,7 +3,7 @@
|
||||
"version": "2.0.0",
|
||||
"created": "2026-03-07T22:43:24.315842",
|
||||
"description": "Standards bypass configuration for prax branch",
|
||||
"last_updated": "2026-04-26T00:06:19.911705"
|
||||
"last_updated": "2026-07-21T00:00:00.000000"
|
||||
},
|
||||
"bypass": [
|
||||
{
|
||||
@@ -511,6 +511,12 @@
|
||||
"file": "tests/test_logging.py",
|
||||
"standard": "log_structure",
|
||||
"reason": "Test assertions reference _AIPASS_PKG_ROOT which resolves to /home/ paths at runtime. Not log config \u2014 test path constants."
|
||||
},
|
||||
{
|
||||
"file": "apps/handlers/monitoring/commons_feed.py",
|
||||
"standard": "cli",
|
||||
"pattern": "console.print(",
|
||||
"reason": "This IS the display handler for the commons live feed (DPLAN-0257). commons_feed.py is a self-contained view \u2014 polling, formatting, and terminal rendering together, same shape as unified_stream.py. console.print() is its designated purpose, not a violation."
|
||||
}
|
||||
],
|
||||
"notes": {
|
||||
@@ -533,4 +539,4 @@
|
||||
"trigger_integration": "Several handlers optionally import trigger.modules.core for event firing. All have graceful ImportError fallbacks. This is cross-branch integration, not an architectural violation."
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,474 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: commons_feed.py
|
||||
# Description: Commons Live Social Feed (read-only monitor view)
|
||||
# Version: 1.0.0
|
||||
# Created: 2026-07-21
|
||||
# Modified: 2026-07-21
|
||||
# =============================================
|
||||
|
||||
"""
|
||||
PRAX Monitor — Commons Live Feed
|
||||
|
||||
`drone @prax monitor run commons` streams The Commons' social chatter
|
||||
(posts, comments, votes, reactions) room-tagged, monitor-style, instead of
|
||||
tailing commons' technical prax logs (still reachable via `--logs`).
|
||||
|
||||
Read-only by construction: connects to commons.db via a `mode=ro` URI, so
|
||||
the feed can never write — commons stays the only writer. Polls with per-table
|
||||
id cursors, showing only genuinely new rows each cycle.
|
||||
"""
|
||||
|
||||
import os
|
||||
import sqlite3
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from typing import Dict, List, Optional, Tuple
|
||||
|
||||
from aipass.prax.apps.modules.logger import system_logger as logger
|
||||
from aipass.cli.apps.modules import console, header, error, warning
|
||||
from aipass.prax.apps.handlers.json import json_handler
|
||||
from aipass.prax.apps.handlers.monitoring.event_queue import MonitoringEvent
|
||||
from aipass.prax.apps.handlers.monitoring.telegram_relay import (
|
||||
init_relay,
|
||||
relay_event,
|
||||
stop_relay,
|
||||
is_relay_enabled_by_env,
|
||||
)
|
||||
|
||||
# =============================================================================
|
||||
# CONFIGURATION
|
||||
# =============================================================================
|
||||
|
||||
POLL_INTERVAL = 1.5
|
||||
BACKFILL_LIMIT = 10
|
||||
|
||||
_PRAX_ROOT = Path(__file__).resolve().parents[3] # monitoring/ -> handlers/ -> apps/ -> prax/
|
||||
_ECOSYSTEM_ROOT = _PRAX_ROOT.parent
|
||||
|
||||
CURSOR_TABLES = ("posts", "comments", "votes", "reactions")
|
||||
|
||||
_CURSOR_TABLE_FOR_KIND = {
|
||||
"post": "posts",
|
||||
"comment": "comments",
|
||||
"vote": "votes",
|
||||
"reaction": "reactions",
|
||||
}
|
||||
|
||||
MOOD_COLORS = {
|
||||
"welcoming": "green",
|
||||
"focused": "cyan",
|
||||
"relaxed": "yellow",
|
||||
"formal": "blue",
|
||||
"creative": "magenta",
|
||||
"neutral": "white",
|
||||
}
|
||||
DEFAULT_ROOM_COLOR = "white"
|
||||
ROOM_LABEL_WIDTH = 18
|
||||
|
||||
_stop_event = threading.Event()
|
||||
|
||||
# =============================================================================
|
||||
# DB ACCESS (read-only)
|
||||
# =============================================================================
|
||||
|
||||
|
||||
def _get_commons_db_path() -> Path:
|
||||
"""Resolve commons.db — sibling branch under the same ecosystem root.
|
||||
|
||||
AIPASS_COMMONS_DB_PATH overrides for tests, mirroring AIPASS_TEST_LOG_DIR.
|
||||
"""
|
||||
override = os.environ.get("AIPASS_COMMONS_DB_PATH")
|
||||
if override:
|
||||
return Path(override)
|
||||
return _ECOSYSTEM_ROOT / "commons" / "commons.db"
|
||||
|
||||
|
||||
def connect_readonly(db_path: Path) -> sqlite3.Connection:
|
||||
"""Open commons.db strictly read-only via a mode=ro URI — writes always fail."""
|
||||
conn = sqlite3.connect(f"file:{db_path}?mode=ro", uri=True, check_same_thread=False)
|
||||
conn.row_factory = sqlite3.Row
|
||||
return conn
|
||||
|
||||
|
||||
def _load_room_moods(conn: sqlite3.Connection) -> Dict[str, str]:
|
||||
"""Load room name -> mood, for cheap color tagging."""
|
||||
rows = conn.execute("SELECT name, mood FROM rooms").fetchall()
|
||||
return {row["name"]: row["mood"] for row in rows}
|
||||
|
||||
|
||||
def initial_cursors(conn: sqlite3.Connection) -> Dict[str, int]:
|
||||
"""Snapshot the current max id per event table — the poll starting line."""
|
||||
cursors = {}
|
||||
for table in CURSOR_TABLES:
|
||||
row = conn.execute(f"SELECT COALESCE(MAX(id), 0) AS m FROM {table}").fetchone()
|
||||
cursors[table] = row["m"]
|
||||
return cursors
|
||||
|
||||
|
||||
def _fetch_new_posts(conn: sqlite3.Connection, since_id: int) -> List[dict]:
|
||||
rows = conn.execute(
|
||||
"SELECT id, room_name, author, title, content, created_at FROM posts WHERE id > ? ORDER BY id ASC",
|
||||
(since_id,),
|
||||
).fetchall()
|
||||
return [dict(r) for r in rows]
|
||||
|
||||
|
||||
def _fetch_new_comments(conn: sqlite3.Connection, since_id: int) -> List[dict]:
|
||||
rows = conn.execute(
|
||||
"SELECT c.id, c.post_id, c.parent_id, c.author, c.content, c.created_at, "
|
||||
"p.room_name AS room_name, p.author AS post_author "
|
||||
"FROM comments c JOIN posts p ON c.post_id = p.id "
|
||||
"WHERE c.id > ? ORDER BY c.id ASC",
|
||||
(since_id,),
|
||||
).fetchall()
|
||||
return [dict(r) for r in rows]
|
||||
|
||||
|
||||
def _fetch_new_votes(conn: sqlite3.Connection, since_id: int) -> List[dict]:
|
||||
rows = conn.execute(
|
||||
"SELECT v.id, v.agent_name, v.target_id, v.target_type, v.direction, v.created_at, "
|
||||
"COALESCE(p1.room_name, p2.room_name) AS room_name "
|
||||
"FROM votes v "
|
||||
"LEFT JOIN posts p1 ON v.target_type = 'post' AND v.target_id = p1.id "
|
||||
"LEFT JOIN comments c ON v.target_type = 'comment' AND v.target_id = c.id "
|
||||
"LEFT JOIN posts p2 ON c.post_id = p2.id "
|
||||
"WHERE v.id > ? ORDER BY v.id ASC",
|
||||
(since_id,),
|
||||
).fetchall()
|
||||
return [dict(r) for r in rows]
|
||||
|
||||
|
||||
def _fetch_new_reactions(conn: sqlite3.Connection, since_id: int) -> List[dict]:
|
||||
rows = conn.execute(
|
||||
"SELECT r.id, r.agent_name, r.post_id, r.comment_id, r.reaction, r.created_at, "
|
||||
"COALESCE(p1.room_name, p2.room_name) AS room_name "
|
||||
"FROM reactions r "
|
||||
"LEFT JOIN posts p1 ON r.post_id = p1.id "
|
||||
"LEFT JOIN comments c ON r.comment_id = c.id "
|
||||
"LEFT JOIN posts p2 ON c.post_id = p2.id "
|
||||
"WHERE r.id > ? ORDER BY r.id ASC",
|
||||
(since_id,),
|
||||
).fetchall()
|
||||
return [dict(r) for r in rows]
|
||||
|
||||
|
||||
_NEW_FETCHERS = {
|
||||
"post": _fetch_new_posts,
|
||||
"comment": _fetch_new_comments,
|
||||
"vote": _fetch_new_votes,
|
||||
"reaction": _fetch_new_reactions,
|
||||
}
|
||||
|
||||
|
||||
def fetch_new_events(conn: sqlite3.Connection, cursors: Dict[str, int]) -> Tuple[List[dict], Dict[str, int]]:
|
||||
"""Fetch rows newer than each cursor, tag with kind, return sorted events + advanced cursors."""
|
||||
events: List[dict] = []
|
||||
new_cursors = dict(cursors)
|
||||
|
||||
for kind, table in _CURSOR_TABLE_FOR_KIND.items():
|
||||
rows = _NEW_FETCHERS[kind](conn, cursors[table])
|
||||
if rows:
|
||||
events.extend({**row, "kind": kind} for row in rows)
|
||||
new_cursors[table] = rows[-1]["id"]
|
||||
|
||||
events.sort(key=lambda e: (e["created_at"], e["id"]))
|
||||
return events, new_cursors
|
||||
|
||||
|
||||
def fetch_backfill(conn: sqlite3.Connection, cursors: Dict[str, int], limit: int = BACKFILL_LIMIT) -> List[dict]:
|
||||
"""Fetch the last `limit` events at/under the given cursors — startup context."""
|
||||
events: List[dict] = []
|
||||
queries = {
|
||||
"post": "SELECT id, room_name, author, title, content, created_at FROM posts "
|
||||
"WHERE id <= ? ORDER BY id DESC LIMIT ?",
|
||||
"comment": "SELECT c.id, c.post_id, c.parent_id, c.author, c.content, c.created_at, "
|
||||
"p.room_name AS room_name, p.author AS post_author FROM comments c "
|
||||
"JOIN posts p ON c.post_id = p.id WHERE c.id <= ? ORDER BY c.id DESC LIMIT ?",
|
||||
"vote": "SELECT v.id, v.agent_name, v.target_id, v.target_type, v.direction, v.created_at, "
|
||||
"COALESCE(p1.room_name, p2.room_name) AS room_name FROM votes v "
|
||||
"LEFT JOIN posts p1 ON v.target_type = 'post' AND v.target_id = p1.id "
|
||||
"LEFT JOIN comments c ON v.target_type = 'comment' AND v.target_id = c.id "
|
||||
"LEFT JOIN posts p2 ON c.post_id = p2.id WHERE v.id <= ? ORDER BY v.id DESC LIMIT ?",
|
||||
"reaction": "SELECT r.id, r.agent_name, r.post_id, r.comment_id, r.reaction, r.created_at, "
|
||||
"COALESCE(p1.room_name, p2.room_name) AS room_name FROM reactions r "
|
||||
"LEFT JOIN posts p1 ON r.post_id = p1.id "
|
||||
"LEFT JOIN comments c ON r.comment_id = c.id "
|
||||
"LEFT JOIN posts p2 ON c.post_id = p2.id WHERE r.id <= ? ORDER BY r.id DESC LIMIT ?",
|
||||
}
|
||||
|
||||
for kind, sql in queries.items():
|
||||
table = _CURSOR_TABLE_FOR_KIND[kind]
|
||||
rows = conn.execute(sql, (cursors[table], limit)).fetchall()
|
||||
events.extend({**dict(row), "kind": kind} for row in rows)
|
||||
|
||||
events.sort(key=lambda e: (e["created_at"], e["id"]))
|
||||
return events[-limit:]
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# DISPLAY
|
||||
# =============================================================================
|
||||
|
||||
|
||||
def _snippet(text: Optional[str], length: int = 100) -> str:
|
||||
"""Collapse whitespace and truncate to ~length chars with an ellipsis."""
|
||||
if not text:
|
||||
return ""
|
||||
collapsed = " ".join(text.split())
|
||||
if len(collapsed) <= length:
|
||||
return collapsed
|
||||
return collapsed[: length - 1].rstrip() + "…"
|
||||
|
||||
|
||||
def event_room(event: dict) -> str:
|
||||
"""Room name for an event, or 'commons' when it can't be resolved (orphaned row)."""
|
||||
return event.get("room_name") or "commons"
|
||||
|
||||
|
||||
def format_event(event: dict) -> str:
|
||||
"""Plain-text (no Rich markup) description of an event — shared by console + relay."""
|
||||
kind = event["kind"]
|
||||
|
||||
if kind == "post":
|
||||
return f'{event["author"]} posted: "{event["title"]}" — {_snippet(event["content"])}'
|
||||
if kind == "comment":
|
||||
target = event.get("post_author") or "?"
|
||||
return f"{event['author']} replied to {target}: {_snippet(event['content'])}"
|
||||
if kind == "vote":
|
||||
direction = "up" if (event.get("direction") or 0) > 0 else "down"
|
||||
return f"{event['agent_name']} voted {direction} on {event['target_type']} #{event['target_id']}"
|
||||
if kind == "reaction":
|
||||
target_kind = "post" if event.get("post_id") else "comment"
|
||||
reaction = event.get("reaction") or ""
|
||||
return f"{event['agent_name']} reacted {reaction} to a {target_kind}"
|
||||
return str(event)
|
||||
|
||||
|
||||
def _print_feed_event(room: str, message: str, mood: Optional[str]) -> None:
|
||||
"""Print one feed line — room tag colored by mood, monitor-style timestamp."""
|
||||
ts = datetime.now().strftime("%H:%M:%S")
|
||||
color = MOOD_COLORS.get(mood or "", DEFAULT_ROOM_COLOR)
|
||||
label = f"[{room}]"
|
||||
console.print(f"[dim]{ts}[/dim] [{color}]{label:<{ROOM_LABEL_WIDTH}}[/{color}] {message}")
|
||||
|
||||
|
||||
def _emit(event: dict, room_moods: Dict[str, str]) -> None:
|
||||
"""Render an event to the console and relay it (relay is a no-op when inactive)."""
|
||||
room = event_room(event)
|
||||
message = format_event(event)
|
||||
_print_feed_event(room, message, room_moods.get(room))
|
||||
|
||||
relay_event(
|
||||
MonitoringEvent(
|
||||
priority=3,
|
||||
event_type="log",
|
||||
branch=room.upper(),
|
||||
message=message,
|
||||
level="info",
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# FEED STATE + INTERACTIVE COMMANDS
|
||||
# =============================================================================
|
||||
|
||||
|
||||
class FeedState:
|
||||
"""Tracks rooms/agents/events seen and the active room filter."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.rooms_seen: set = set()
|
||||
self.agents_seen: set = set()
|
||||
self.events_count: int = 0
|
||||
self.room_filter: Optional[set] = None
|
||||
|
||||
def record(self, event: dict) -> None:
|
||||
"""Count an event toward status stats — independent of the display filter."""
|
||||
self.rooms_seen.add(event_room(event))
|
||||
agent = event.get("author") or event.get("agent_name")
|
||||
if agent:
|
||||
self.agents_seen.add(agent)
|
||||
self.events_count += 1
|
||||
|
||||
def visible(self, event: dict) -> bool:
|
||||
"""Whether the active room filter allows this event to display."""
|
||||
if not self.room_filter:
|
||||
return True
|
||||
return event_room(event) in self.room_filter
|
||||
|
||||
|
||||
def _feed_help_text() -> str:
|
||||
return """
|
||||
Available Commands:
|
||||
filter <room>[,<room>...] - Show only these rooms
|
||||
filter clear - Remove the room filter
|
||||
status - Show feed stats (rooms/agents/events)
|
||||
help - Show this help
|
||||
quit/exit - Stop the feed
|
||||
"""
|
||||
|
||||
|
||||
def _print_feed_status(state: FeedState) -> None:
|
||||
console.print()
|
||||
console.print("[bold cyan]Commons Feed Status:[/bold cyan]")
|
||||
console.print(f" [green]Rooms seen:[/green] {len(state.rooms_seen)}")
|
||||
console.print(f" [green]Agents seen:[/green] {len(state.agents_seen)}")
|
||||
console.print(f" [green]Events streamed:[/green] {state.events_count}")
|
||||
if state.room_filter:
|
||||
console.print(f" [yellow]Filter:[/yellow] {', '.join(sorted(state.room_filter))}")
|
||||
console.print()
|
||||
|
||||
|
||||
def _handle_feed_cmd(cmd: str, cmd_args: List[str], state: FeedState) -> None:
|
||||
"""Dispatch an interactive feed command."""
|
||||
if cmd == "help":
|
||||
console.print(_feed_help_text())
|
||||
return
|
||||
if cmd == "status":
|
||||
_print_feed_status(state)
|
||||
return
|
||||
if cmd == "filter":
|
||||
if not cmd_args or cmd_args[0] in ("clear", "all"):
|
||||
state.room_filter = None
|
||||
warning("Filter cleared — showing all rooms")
|
||||
return
|
||||
rooms = {room.strip() for token in cmd_args for room in token.split(",") if room.strip()}
|
||||
state.room_filter = rooms
|
||||
warning(f"Filtering to rooms: {', '.join(sorted(rooms))}")
|
||||
return
|
||||
error(f"Unknown command: {cmd}")
|
||||
console.print("[dim]Type 'help' for available commands[/dim]")
|
||||
|
||||
|
||||
def _interactive_feed_loop(state: FeedState) -> None:
|
||||
"""Interactive command loop for the feed — mirrors the branch monitor's loop shape."""
|
||||
if not sys.stdin.isatty():
|
||||
logger.info("[commons_feed] No TTY detected - passive mode (Ctrl+C to stop)")
|
||||
try:
|
||||
while not _stop_event.is_set():
|
||||
time.sleep(0.5)
|
||||
except KeyboardInterrupt:
|
||||
logger.info("[commons_feed] Stopped by user (passive mode)")
|
||||
console.print("\n[yellow]Stopping feed...[/yellow]")
|
||||
return
|
||||
|
||||
from aipass.prax.apps.handlers.monitoring.interactive_filter import parse_command
|
||||
|
||||
while not _stop_event.is_set():
|
||||
try:
|
||||
user_input = input().strip()
|
||||
if not user_input:
|
||||
continue
|
||||
|
||||
cmd, cmd_args = parse_command(user_input)
|
||||
if not cmd:
|
||||
continue
|
||||
|
||||
if cmd in ("quit", "exit", "q"):
|
||||
console.print("[yellow]Stopping feed...[/yellow]")
|
||||
break
|
||||
|
||||
_handle_feed_cmd(cmd, cmd_args, state)
|
||||
|
||||
except KeyboardInterrupt:
|
||||
logger.info("[commons_feed] Stopped by user")
|
||||
console.print("\n[yellow]Stopping feed...[/yellow]")
|
||||
break
|
||||
except EOFError:
|
||||
logger.info("[commons_feed] EOF received, stopping interactive loop")
|
||||
break
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# POLL WORKER
|
||||
# =============================================================================
|
||||
|
||||
|
||||
def _poll_worker(
|
||||
conn: sqlite3.Connection, cursors: Dict[str, int], state: FeedState, room_moods: Dict[str, str]
|
||||
) -> None:
|
||||
"""Background thread: poll for new rows, display + relay the visible ones."""
|
||||
while not _stop_event.is_set():
|
||||
try:
|
||||
events, new_cursors = fetch_new_events(conn, cursors)
|
||||
cursors.update(new_cursors)
|
||||
for event in events:
|
||||
state.record(event)
|
||||
if state.visible(event):
|
||||
_emit(event, room_moods)
|
||||
except sqlite3.Error as exc:
|
||||
logger.warning("[commons_feed] Poll error: %s", exc)
|
||||
_stop_event.wait(POLL_INTERVAL)
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# ENTRY POINT
|
||||
# =============================================================================
|
||||
|
||||
|
||||
def run_commons_feed(args: List[str], relay_config: Optional[dict] = None) -> bool:
|
||||
"""Launch the commons live social feed — read-only, room-tagged, monitor-style.
|
||||
|
||||
relay_config is loaded by monitor.py (the module layer, which owns the
|
||||
cross-branch @api secrets lookup) and passed in — this handler never
|
||||
reaches outside its own branch.
|
||||
"""
|
||||
_stop_event.clear()
|
||||
json_handler.log_operation("commons_feed_started", {"args": args})
|
||||
|
||||
relay_enabled = "--relay" in args or is_relay_enabled_by_env()
|
||||
init_relay(relay_enabled, relay_config if relay_enabled else None)
|
||||
if relay_enabled:
|
||||
console.print("[green]monitor → Telegram relay ON (prax_monitor)[/green]")
|
||||
|
||||
db_path = _get_commons_db_path()
|
||||
if not db_path.exists():
|
||||
error(f"commons.db not found at {db_path}")
|
||||
stop_relay()
|
||||
return True
|
||||
|
||||
try:
|
||||
conn = connect_readonly(db_path)
|
||||
cursors = initial_cursors(conn)
|
||||
except sqlite3.Error as exc:
|
||||
logger.warning("[commons_feed] Could not open commons.db read-only: %s", exc)
|
||||
error(f"Could not open commons.db read-only: {exc}")
|
||||
stop_relay()
|
||||
return True
|
||||
|
||||
state = FeedState()
|
||||
room_moods = _load_room_moods(conn)
|
||||
|
||||
console.print()
|
||||
header("PRAX Mission Control - Commons Live Feed")
|
||||
console.print()
|
||||
console.print("[green]Live — read-only feed of The Commons (posts, comments, votes, reactions)[/green]")
|
||||
is_tty = sys.stdin.isatty()
|
||||
console.print("[dim]Type 'help' for commands[/dim]" if is_tty else "[dim]Ctrl+C to stop[/dim]")
|
||||
console.print()
|
||||
|
||||
for event in fetch_backfill(conn, cursors):
|
||||
state.record(event)
|
||||
_emit(event, room_moods)
|
||||
|
||||
poll_thread = threading.Thread(target=_poll_worker, args=(conn, cursors, state, room_moods), daemon=True)
|
||||
poll_thread.start()
|
||||
|
||||
try:
|
||||
_interactive_feed_loop(state)
|
||||
except KeyboardInterrupt:
|
||||
logger.info("[commons_feed] KeyboardInterrupt escaped interactive loop")
|
||||
console.print("\n[yellow]Feed stopped.[/yellow]")
|
||||
|
||||
_stop_event.set()
|
||||
poll_thread.join(timeout=POLL_INTERVAL + 2.0)
|
||||
stop_relay()
|
||||
conn.close()
|
||||
|
||||
# sys.exit(0) prevents drone's post-execution json_handler from running
|
||||
# after the feed exits, avoiding a json.load crash on Ctrl+C (matches monitor.py).
|
||||
sys.exit(0)
|
||||
@@ -108,6 +108,14 @@ def print_help():
|
||||
"drone @prax monitor run [branches]",
|
||||
"Monitor specific branches (comma-separated)\n Example: drone @prax monitor run seedgo,cli,flow",
|
||||
),
|
||||
(
|
||||
"drone @prax monitor run commons",
|
||||
"Live social feed of The Commons (posts, comments, votes, reactions), read-only",
|
||||
),
|
||||
(
|
||||
"drone @prax monitor run commons --logs",
|
||||
"Old behavior — tail commons branch's technical prax logs instead of the feed",
|
||||
),
|
||||
(
|
||||
"drone @prax monitor run --relay",
|
||||
"Enable Telegram relay (mirrors feed to prax_monitor bot)"
|
||||
@@ -153,7 +161,7 @@ def handle_command(command: str, args: List[str]) -> bool:
|
||||
# Subcommand routing
|
||||
subcmd = args[0]
|
||||
if subcmd == "run":
|
||||
return _run_monitor(args[1:])
|
||||
return _dispatch_run(args[1:])
|
||||
|
||||
# Unknown subcommand
|
||||
error(f"Unknown monitor subcommand: {subcmd}")
|
||||
@@ -161,6 +169,28 @@ def handle_command(command: str, args: List[str]) -> bool:
|
||||
return True
|
||||
|
||||
|
||||
def _dispatch_run(run_args: List[str]) -> bool:
|
||||
"""Route 'monitor run' — a bare 'commons' target opens the live social feed.
|
||||
|
||||
'commons --logs' escapes back to the branch-log tail, and mixed lists
|
||||
(e.g. 'seedgo,commons') keep treating commons as a branch (feed is
|
||||
standalone-only in v1).
|
||||
"""
|
||||
positional = [a for a in run_args if not a.startswith("--")]
|
||||
target = positional[0] if positional else None
|
||||
logs_escape = "--logs" in run_args
|
||||
|
||||
if target == "commons" and not logs_escape:
|
||||
from aipass.prax.apps.handlers.monitoring.commons_feed import run_commons_feed
|
||||
|
||||
feed_args = [a for a in run_args if a != target]
|
||||
relay_enabled = "--relay" in feed_args or is_relay_enabled_by_env()
|
||||
relay_config = _load_relay_config() if relay_enabled else None
|
||||
return run_commons_feed(feed_args, relay_config=relay_config)
|
||||
|
||||
return _run_monitor([a for a in run_args if a != "--logs"])
|
||||
|
||||
|
||||
def _load_relay_config() -> Optional[dict]:
|
||||
"""Load Telegram relay config from @api secrets."""
|
||||
try:
|
||||
|
||||
@@ -0,0 +1,366 @@
|
||||
# =================== AIPass ====================
|
||||
# Name: test_commons_feed.py
|
||||
# Description: Tests for the commons live feed handler
|
||||
# Version: 1.0.0
|
||||
# Created: 2026-07-21
|
||||
# Modified: 2026-07-21
|
||||
# =============================================
|
||||
|
||||
"""Tests for apps/handlers/monitoring/commons_feed.py (DPLAN-0257)
|
||||
|
||||
Covers:
|
||||
- connect_readonly: mode=ro connection actually refuses writes
|
||||
- initial_cursors / fetch_new_events: only-new-row cursor semantics
|
||||
- fetch_backfill: last-N-events context on start, bounded by cursor
|
||||
- display formatting: format_event, event_room, _snippet
|
||||
- FeedState: record/visible room filtering
|
||||
- _get_commons_db_path: env override
|
||||
"""
|
||||
|
||||
import importlib
|
||||
import sqlite3
|
||||
|
||||
import pytest
|
||||
|
||||
feed = importlib.import_module("aipass.prax.apps.handlers.monitoring.commons_feed")
|
||||
|
||||
|
||||
SCHEMA = """
|
||||
CREATE TABLE rooms (
|
||||
name TEXT PRIMARY KEY,
|
||||
mood TEXT DEFAULT 'neutral'
|
||||
);
|
||||
CREATE TABLE posts (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
room_name TEXT NOT NULL,
|
||||
author TEXT NOT NULL,
|
||||
title TEXT NOT NULL,
|
||||
content TEXT DEFAULT '',
|
||||
created_at TEXT NOT NULL
|
||||
);
|
||||
CREATE TABLE comments (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
post_id INTEGER NOT NULL,
|
||||
parent_id INTEGER,
|
||||
author TEXT NOT NULL,
|
||||
content TEXT NOT NULL,
|
||||
created_at TEXT NOT NULL
|
||||
);
|
||||
CREATE TABLE votes (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
agent_name TEXT NOT NULL,
|
||||
target_id INTEGER NOT NULL,
|
||||
target_type TEXT NOT NULL,
|
||||
direction INTEGER NOT NULL,
|
||||
created_at TEXT NOT NULL
|
||||
);
|
||||
CREATE TABLE reactions (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
agent_name TEXT NOT NULL,
|
||||
post_id INTEGER,
|
||||
comment_id INTEGER,
|
||||
reaction TEXT NOT NULL,
|
||||
created_at TEXT NOT NULL
|
||||
);
|
||||
"""
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def db_path(tmp_path):
|
||||
"""A throwaway commons.db with the tables the feed queries."""
|
||||
path = tmp_path / "commons.db"
|
||||
conn = sqlite3.connect(str(path))
|
||||
conn.executescript(SCHEMA)
|
||||
conn.execute("INSERT INTO rooms (name, mood) VALUES ('general', 'welcoming')")
|
||||
conn.commit()
|
||||
conn.close()
|
||||
return path
|
||||
|
||||
|
||||
def _seed_post(path, room="general", author="alice", title="Hello", content="world", ts="2026-07-21T10:00:00Z"):
|
||||
conn = sqlite3.connect(str(path))
|
||||
cur = conn.execute(
|
||||
"INSERT INTO posts (room_name, author, title, content, created_at) VALUES (?, ?, ?, ?, ?)",
|
||||
(room, author, title, content, ts),
|
||||
)
|
||||
conn.commit()
|
||||
post_id = cur.lastrowid
|
||||
conn.close()
|
||||
return post_id
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# connect_readonly — write refusal
|
||||
# =============================================================================
|
||||
|
||||
|
||||
class TestConnectReadonly:
|
||||
"""The feed's connection must never be able to write to commons.db."""
|
||||
|
||||
def test_read_succeeds(self, db_path):
|
||||
"""A mode=ro connection can still read existing rows."""
|
||||
_seed_post(db_path)
|
||||
conn = feed.connect_readonly(db_path)
|
||||
try:
|
||||
rows = conn.execute("SELECT * FROM posts").fetchall()
|
||||
assert len(rows) == 1
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
def test_write_attempt_raises_operational_error(self, db_path):
|
||||
"""INSERT on a mode=ro connection raises instead of succeeding."""
|
||||
conn = feed.connect_readonly(db_path)
|
||||
try:
|
||||
with pytest.raises(sqlite3.OperationalError):
|
||||
conn.execute(
|
||||
"INSERT INTO posts (room_name, author, title, created_at) VALUES ('general', 'bob', 'x', 'y')"
|
||||
)
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
def test_write_attempt_leaves_db_unchanged(self, db_path):
|
||||
"""A refused DELETE leaves the underlying rows untouched."""
|
||||
conn = feed.connect_readonly(db_path)
|
||||
try:
|
||||
with pytest.raises(sqlite3.OperationalError):
|
||||
conn.execute("DELETE FROM posts")
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
verify = sqlite3.connect(str(db_path))
|
||||
count = verify.execute("SELECT COUNT(*) FROM rooms").fetchone()[0]
|
||||
verify.close()
|
||||
assert count == 1
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# Cursor logic — only new rows
|
||||
# =============================================================================
|
||||
|
||||
|
||||
class TestCursorLogic:
|
||||
"""Cursors should only ever surface rows newer than the snapshot."""
|
||||
|
||||
def test_initial_cursors_snapshot_current_max(self, db_path):
|
||||
"""initial_cursors captures the current max id per table."""
|
||||
_seed_post(db_path)
|
||||
_seed_post(db_path, title="Second")
|
||||
conn = feed.connect_readonly(db_path)
|
||||
try:
|
||||
cursors = feed.initial_cursors(conn)
|
||||
finally:
|
||||
conn.close()
|
||||
assert cursors == {"posts": 2, "comments": 0, "votes": 0, "reactions": 0}
|
||||
|
||||
def test_fetch_new_events_empty_when_nothing_past_cursor(self, db_path):
|
||||
"""No rows inserted since the cursor snapshot means no events."""
|
||||
_seed_post(db_path)
|
||||
conn = feed.connect_readonly(db_path)
|
||||
try:
|
||||
cursors = feed.initial_cursors(conn)
|
||||
events, new_cursors = feed.fetch_new_events(conn, cursors)
|
||||
finally:
|
||||
conn.close()
|
||||
assert events == []
|
||||
assert new_cursors == cursors
|
||||
|
||||
def test_fetch_new_events_returns_only_rows_after_cursor(self, db_path):
|
||||
"""Rows inserted before the cursor snapshot are excluded."""
|
||||
_seed_post(db_path, title="Before cursor")
|
||||
conn = feed.connect_readonly(db_path)
|
||||
cursors = feed.initial_cursors(conn)
|
||||
|
||||
_seed_post(db_path, title="After cursor", ts="2026-07-21T11:00:00Z")
|
||||
|
||||
events, new_cursors = feed.fetch_new_events(conn, cursors)
|
||||
conn.close()
|
||||
|
||||
assert len(events) == 1
|
||||
assert events[0]["title"] == "After cursor"
|
||||
assert new_cursors["posts"] == 2
|
||||
|
||||
def test_fetch_new_events_advances_cursor_only_for_kinds_with_new_rows(self, db_path):
|
||||
"""Cursors for tables with no new rows stay put; posts' cursor advances."""
|
||||
_seed_post(db_path)
|
||||
conn = feed.connect_readonly(db_path)
|
||||
cursors = feed.initial_cursors(conn)
|
||||
|
||||
_seed_post(db_path, title="Another", ts="2026-07-21T11:00:00Z")
|
||||
|
||||
_, new_cursors = feed.fetch_new_events(conn, cursors)
|
||||
conn.close()
|
||||
|
||||
assert new_cursors["posts"] == 2
|
||||
assert new_cursors["comments"] == 0
|
||||
assert new_cursors["votes"] == 0
|
||||
assert new_cursors["reactions"] == 0
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# Backfill — startup context
|
||||
# =============================================================================
|
||||
|
||||
|
||||
class TestFetchBackfill:
|
||||
"""On start the feed shows recent history for context, bounded by cursor + limit."""
|
||||
|
||||
def test_backfill_returns_events_up_to_limit(self, db_path):
|
||||
"""Backfill returns only the most recent `limit` events, oldest first."""
|
||||
for i in range(15):
|
||||
_seed_post(db_path, title=f"Post {i}", ts=f"2026-07-21T10:{i:02d}:00Z")
|
||||
conn = feed.connect_readonly(db_path)
|
||||
cursors = feed.initial_cursors(conn)
|
||||
events = feed.fetch_backfill(conn, cursors, limit=10)
|
||||
conn.close()
|
||||
|
||||
assert len(events) == 10
|
||||
assert [e["title"] for e in events] == [f"Post {i}" for i in range(5, 15)]
|
||||
|
||||
def test_backfill_respects_cursor_bound(self, db_path):
|
||||
"""Backfill never includes rows that came in after the cursor snapshot.
|
||||
|
||||
Rows inserted after the cursor is taken belong to the live poll, not
|
||||
the startup backfill — including them here would double-emit.
|
||||
"""
|
||||
_seed_post(db_path, title="Old")
|
||||
conn = feed.connect_readonly(db_path)
|
||||
cursors = feed.initial_cursors(conn)
|
||||
|
||||
_seed_post(db_path, title="New", ts="2026-07-21T12:00:00Z")
|
||||
|
||||
events = feed.fetch_backfill(conn, cursors, limit=10)
|
||||
conn.close()
|
||||
|
||||
assert [e["title"] for e in events] == ["Old"]
|
||||
|
||||
def test_backfill_empty_db_returns_empty(self, db_path):
|
||||
"""An empty database backfills to an empty list, not an error."""
|
||||
conn = feed.connect_readonly(db_path)
|
||||
cursors = feed.initial_cursors(conn)
|
||||
events = feed.fetch_backfill(conn, cursors)
|
||||
conn.close()
|
||||
assert events == []
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# Display formatting
|
||||
# =============================================================================
|
||||
|
||||
|
||||
class TestFormatEvent:
|
||||
def test_post_event(self):
|
||||
"""A post event formats as author + quoted title + content snippet."""
|
||||
event = {"kind": "post", "author": "alice", "title": "Hello world", "content": "some content here"}
|
||||
line = feed.format_event(event)
|
||||
assert line == 'alice posted: "Hello world" — some content here'
|
||||
|
||||
def test_comment_event(self):
|
||||
"""A comment event formats as author replying to the post's author."""
|
||||
event = {"kind": "comment", "author": "bob", "content": "nice one", "post_author": "alice"}
|
||||
line = feed.format_event(event)
|
||||
assert line == "bob replied to alice: nice one"
|
||||
|
||||
def test_comment_event_missing_post_author(self):
|
||||
"""A missing post_author falls back to '?' rather than crashing."""
|
||||
event = {"kind": "comment", "author": "bob", "content": "nice one", "post_author": None}
|
||||
line = feed.format_event(event)
|
||||
assert line == "bob replied to ?: nice one"
|
||||
|
||||
def test_vote_event_up(self):
|
||||
"""A positive direction renders as 'voted up'."""
|
||||
event = {"kind": "vote", "agent_name": "bob", "direction": 1, "target_type": "post", "target_id": 5}
|
||||
assert feed.format_event(event) == "bob voted up on post #5"
|
||||
|
||||
def test_vote_event_down(self):
|
||||
"""A negative direction renders as 'voted down'."""
|
||||
event = {"kind": "vote", "agent_name": "bob", "direction": -1, "target_type": "comment", "target_id": 7}
|
||||
assert feed.format_event(event) == "bob voted down on comment #7"
|
||||
|
||||
def test_reaction_event_on_post(self):
|
||||
"""A reaction with post_id set targets 'a post'."""
|
||||
event = {"kind": "reaction", "agent_name": "carol", "post_id": 3, "comment_id": None, "reaction": "agree"}
|
||||
assert feed.format_event(event) == "carol reacted agree to a post"
|
||||
|
||||
def test_reaction_event_on_comment(self):
|
||||
"""A reaction with comment_id set targets 'a comment'."""
|
||||
event = {"kind": "reaction", "agent_name": "carol", "post_id": None, "comment_id": 9, "reaction": "thinking"}
|
||||
assert feed.format_event(event) == "carol reacted thinking to a comment"
|
||||
|
||||
|
||||
class TestEventRoom:
|
||||
def test_returns_room_name_when_present(self):
|
||||
"""A populated room_name passes through unchanged."""
|
||||
assert feed.event_room({"room_name": "dev"}) == "dev"
|
||||
|
||||
def test_falls_back_to_commons_when_missing(self):
|
||||
"""A None or absent room_name falls back to 'commons'."""
|
||||
assert feed.event_room({"room_name": None}) == "commons"
|
||||
assert feed.event_room({}) == "commons"
|
||||
|
||||
|
||||
class TestSnippet:
|
||||
def test_short_text_unchanged(self):
|
||||
"""Text under the length limit passes through unchanged."""
|
||||
assert feed._snippet("short text") == "short text"
|
||||
|
||||
def test_collapses_whitespace(self):
|
||||
"""Runs of whitespace (including newlines) collapse to single spaces."""
|
||||
assert feed._snippet("a b\n\nc") == "a b c"
|
||||
|
||||
def test_truncates_long_text_with_ellipsis(self):
|
||||
"""Text over the length limit is truncated with a trailing ellipsis."""
|
||||
text = "x" * 150
|
||||
result = feed._snippet(text, length=100)
|
||||
assert len(result) == 100
|
||||
assert result.endswith("…")
|
||||
|
||||
def test_none_returns_empty_string(self):
|
||||
"""None input returns an empty string rather than raising."""
|
||||
assert feed._snippet(None) == ""
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# FeedState
|
||||
# =============================================================================
|
||||
|
||||
|
||||
class TestFeedState:
|
||||
def test_record_tracks_rooms_agents_and_count(self):
|
||||
"""record() accumulates distinct rooms, distinct agents, and a running count."""
|
||||
state = feed.FeedState()
|
||||
state.record({"room_name": "dev", "author": "alice"})
|
||||
state.record({"room_name": "dev", "agent_name": "bob"})
|
||||
assert state.rooms_seen == {"dev"}
|
||||
assert state.agents_seen == {"alice", "bob"}
|
||||
assert state.events_count == 2
|
||||
|
||||
def test_visible_true_when_no_filter(self):
|
||||
"""With no active room filter, every event is visible."""
|
||||
state = feed.FeedState()
|
||||
assert state.visible({"room_name": "dev"}) is True
|
||||
|
||||
def test_visible_respects_room_filter(self):
|
||||
"""With an active room filter, only matching-room events are visible."""
|
||||
state = feed.FeedState()
|
||||
state.room_filter = {"dev"}
|
||||
assert state.visible({"room_name": "dev"}) is True
|
||||
assert state.visible({"room_name": "general"}) is False
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# _get_commons_db_path — env override
|
||||
# =============================================================================
|
||||
|
||||
|
||||
class TestGetCommonsDbPath:
|
||||
def test_env_override(self, monkeypatch, tmp_path):
|
||||
"""AIPASS_COMMONS_DB_PATH overrides the default sibling-branch path."""
|
||||
override = tmp_path / "custom.db"
|
||||
monkeypatch.setenv("AIPASS_COMMONS_DB_PATH", str(override))
|
||||
assert feed._get_commons_db_path() == override
|
||||
|
||||
def test_default_is_sibling_commons_branch(self, monkeypatch):
|
||||
"""Without an env override, the path resolves to the commons branch's db."""
|
||||
monkeypatch.delenv("AIPASS_COMMONS_DB_PATH", raising=False)
|
||||
path = feed._get_commons_db_path()
|
||||
assert path.parts[-2:] == ("commons", "commons.db")
|
||||
@@ -116,6 +116,78 @@ class TestHandleCommand:
|
||||
mod.error.assert_called()
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# _dispatch_run routing tests (DPLAN-0257 — commons feed vs. branch log tail)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
_COMMONS_FEED_TARGET = "aipass.prax.apps.handlers.monitoring.commons_feed.run_commons_feed"
|
||||
|
||||
|
||||
class TestDispatchRun:
|
||||
"""'run commons' opens the live social feed; --logs and mixed lists stay branch tail."""
|
||||
|
||||
def test_bare_commons_routes_to_feed(self):
|
||||
"""A standalone 'commons' target opens the feed, not the branch monitor."""
|
||||
mod = _import_monitor()
|
||||
mod.is_relay_enabled_by_env.return_value = False
|
||||
with (
|
||||
patch(_COMMONS_FEED_TARGET, return_value=True) as mock_feed,
|
||||
patch.object(mod, "_run_monitor") as mock_branch,
|
||||
):
|
||||
result = mod.handle_command("monitor", ["run", "commons"])
|
||||
assert result is True
|
||||
mock_feed.assert_called_once_with([], relay_config=None)
|
||||
mock_branch.assert_not_called()
|
||||
|
||||
def test_commons_logs_escape_routes_to_branch_tail(self):
|
||||
"""'commons --logs' is the escape hatch back to the old branch-log tail."""
|
||||
mod = _import_monitor()
|
||||
with (
|
||||
patch(_COMMONS_FEED_TARGET) as mock_feed,
|
||||
patch.object(mod, "_run_monitor", return_value=True) as mock_branch,
|
||||
):
|
||||
result = mod.handle_command("monitor", ["run", "commons", "--logs"])
|
||||
assert result is True
|
||||
mock_branch.assert_called_once_with(["commons"])
|
||||
mock_feed.assert_not_called()
|
||||
|
||||
def test_mixed_branch_list_with_commons_stays_branch_tail(self):
|
||||
"""A mixed list (e.g. 'seedgo,commons') treats commons as a branch — feed is standalone-only."""
|
||||
mod = _import_monitor()
|
||||
with (
|
||||
patch(_COMMONS_FEED_TARGET) as mock_feed,
|
||||
patch.object(mod, "_run_monitor", return_value=True) as mock_branch,
|
||||
):
|
||||
result = mod.handle_command("monitor", ["run", "seedgo,commons"])
|
||||
assert result is True
|
||||
mock_branch.assert_called_once_with(["seedgo,commons"])
|
||||
mock_feed.assert_not_called()
|
||||
|
||||
def test_commons_feed_loads_relay_config_when_relay_flag_set(self):
|
||||
"""--relay on the feed path loads config via monitor.py's own _load_relay_config."""
|
||||
mod = _import_monitor()
|
||||
mod.is_relay_enabled_by_env.return_value = False
|
||||
with (
|
||||
patch(_COMMONS_FEED_TARGET, return_value=True) as mock_feed,
|
||||
patch.object(mod, "_load_relay_config", return_value={"bot_token": "x", "chat_id": 1}) as mock_load,
|
||||
):
|
||||
mod.handle_command("monitor", ["run", "commons", "--relay"])
|
||||
mock_load.assert_called_once()
|
||||
mock_feed.assert_called_once_with(["--relay"], relay_config={"bot_token": "x", "chat_id": 1})
|
||||
|
||||
def test_commons_feed_no_relay_config_when_relay_disabled(self):
|
||||
"""Without --relay (and no env flag), relay_config stays None."""
|
||||
mod = _import_monitor()
|
||||
mod.is_relay_enabled_by_env.return_value = False
|
||||
with (
|
||||
patch(_COMMONS_FEED_TARGET, return_value=True) as mock_feed,
|
||||
patch.object(mod, "_load_relay_config") as mock_load,
|
||||
):
|
||||
mod.handle_command("monitor", ["run", "commons"])
|
||||
mock_load.assert_not_called()
|
||||
mock_feed.assert_called_once_with([], relay_config=None)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# _get_watch_directories tests
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
Reference in New Issue
Block a user