diff --git a/CHANGELOG.md b/CHANGELOG.md index c96b5ee0..9d86a21e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/src/aipass/prax/.seedgo/bypass.json b/src/aipass/prax/.seedgo/bypass.json index fa5fd64e..30d57725 100644 --- a/src/aipass/prax/.seedgo/bypass.json +++ b/src/aipass/prax/.seedgo/bypass.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." } } -} \ No newline at end of file +} diff --git a/src/aipass/prax/apps/handlers/monitoring/commons_feed.py b/src/aipass/prax/apps/handlers/monitoring/commons_feed.py new file mode 100644 index 00000000..d46c5401 --- /dev/null +++ b/src/aipass/prax/apps/handlers/monitoring/commons_feed.py @@ -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 [,...] - 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) diff --git a/src/aipass/prax/apps/modules/monitor.py b/src/aipass/prax/apps/modules/monitor.py index d30a5b9e..a823a055 100755 --- a/src/aipass/prax/apps/modules/monitor.py +++ b/src/aipass/prax/apps/modules/monitor.py @@ -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: diff --git a/src/aipass/prax/tests/test_commons_feed.py b/src/aipass/prax/tests/test_commons_feed.py new file mode 100644 index 00000000..e4b5d3bc --- /dev/null +++ b/src/aipass/prax/tests/test_commons_feed.py @@ -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") diff --git a/src/aipass/prax/tests/test_monitor_module.py b/src/aipass/prax/tests/test_monitor_module.py index 6d8200b3..acda39ca 100644 --- a/src/aipass/prax/tests/test_monitor_module.py +++ b/src/aipass/prax/tests/test_monitor_module.py @@ -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 # ---------------------------------------------------------------------------