"""Tests for watchdog.py — opencode loop detection orchestrator. Covers: resolve_roots, _read_rc_servers, _fire_alert, trigger idempotence, _make_key, detect_config_from_pydantic, calls_of (stubbed), and tick() state machine (lock, no_opencode, verdict, alert transition). """ import json import os import sqlite3 import sys import tempfile from datetime import datetime, timezone from unittest.mock import MagicMock, patch sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "src")) from config import NotificationsConfig from notifier import AlertEvent, Notifier from progress_detect import target_of from watchdog import ( _attribution, _fire_alert, _make_key, _read_rc_servers, calls_of, detect_config_from_pydantic, evaluate, resolve_roots, tick, ) # --------------------------------------------------------------------------- # helpers # --------------------------------------------------------------------------- def _stub_http(data): """Callable that always returns *data*.""" return lambda url, **kw: data def _make_conns(): """Temp SQLite with watchdog tables; returns (conn, tmp_path).""" import watchdog_store path = tempfile.mktemp(suffix=".db") conn = sqlite3.connect(path) conn.execute("PRAGMA foreign_keys=ON") # Use default sqlite3.Row (supports row["name"] and row[0]) conn.row_factory = sqlite3.Row watchdog_store.ensure_watchdog_tables(conn) conn.commit() return conn, path def _close_conn(conn, path): try: conn.close() except (sqlite3.Error, OSError): pass try: os.unlink(path) except OSError: pass def _make_cfg(**extra): """Minimal mock config with watchdog.detector attributes.""" cfg = MagicMock() cfg.watchdog = MagicMock() dc = cfg.watchdog.detector = MagicMock() dc.window = 600 dc.dup_min = 3 dc.top_min = 4 dc.top_min_ro = 5 dc.cum_min = 10 dc.cover_min = 0.5 dc.min_calls = 3 dc._ro_agents = None cfg.local_compute = MagicMock() cfg.local_compute.enabled = False for k, v in extra.items(): setattr(cfg.watchdog, k, v) return cfg def _sessions_tree(): """Root with two children.""" now = datetime.now(timezone.utc).timestamp() root = {"id": "ses_root123", "title": "refactor main.py", "parentID": None, "info": {"time": {"created": now}}} c1 = {"id": "ses_ch1", "title": "edit main", "parentID": "ses_root123", "info": {"time": {"created": now}}} c2 = {"id": "ses_ch2", "title": "edit utils", "parentID": "ses_root123", "info": {"time": {"created": now}}} return [root, c1, c2] # --------------------------------------------------------------------------- # resolve_roots # --------------------------------------------------------------------------- class TestResolveRoots: def test_three_depth_chain(self): now = datetime.now(timezone.utc).timestamp() sessions = [ {"id": "s1", "parentID": None, "info": {"time": {"created": now}}}, {"id": "s2", "parentID": "s1", "info": {"time": {"created": now}}}, {"id": "s3", "parentID": "s2", "info": {"time": {"created": now}}}, ] session_root, children = resolve_roots(sessions) assert session_root["s1"] == "s1" assert session_root["s2"] == "s1" assert session_root["s3"] == "s1" # children maps ALL nodes to their root (includes grandchildren) assert children["s1"] == ["s1", "s2", "s3"] def test_multiple_roots(self): now = datetime.now(timezone.utc).timestamp() sessions = [ {"id": "a", "parentID": None, "info": {"time": {"created": now}}}, {"id": "b", "parentID": None, "info": {"time": {"created": now}}}, {"id": "c", "parentID": "a", "info": {"time": {"created": now}}}, ] session_root, children = resolve_roots(sessions) assert session_root["a"] == "a" assert session_root["c"] == "a" # children["a"] has a AND c (grandchild c also maps to root a) assert children["a"] == ["a", "c"] assert children["b"] == ["b"] def test_empty(self): sr, ch = resolve_roots([]) assert sr == {} assert ch == {} def test_orphan_child_resolves_to_missing_parent_id(self): """orphan → _resolve("orphan") → pid="missing_parent" → _resolve("missing_parent") → not in id_to_s → returns "missing_parent".""" now = datetime.now(timezone.utc).timestamp() sessions = [ {"id": "solo", "parentID": None, "info": {"time": {"created": now}}}, {"id": "orphan", "parentID": "missing_parent", "info": {"time": {"created": now}}}, ] sr, _ch = resolve_roots(sessions) assert sr["orphan"] == "missing_parent" assert sr["solo"] == "solo" def test_children_sorted(self): now = datetime.now(timezone.utc).timestamp() sessions = [ {"id": "r", "parentID": None, "info": {"time": {"created": now}}}, {"id": "z", "parentID": "r", "info": {"time": {"created": now}}}, {"id": "a", "parentID": "r", "info": {"time": {"created": now}}}, {"id": "m", "parentID": "r", "info": {"time": {"created": now}}}, ] _, children = resolve_roots(sessions) assert children["r"] == ["a", "m", "r", "z"] def test_all_grandchildren_in_root_group(self): """Children dict maps every session to its root group, not just direct children.""" now = datetime.now(timezone.utc).timestamp() sessions = [ {"id": "gp", "parentID": None, "info": {"time": {"created": now}}}, {"id": "p", "parentID": "gp", "info": {"time": {"created": now}}}, {"id": "c", "parentID": "p", "info": {"time": {"created": now}}}, ] sr, ch = resolve_roots(sessions) assert sr["c"] == "gp" assert sr["p"] == "gp" # grandchild "c" is also in gp's children group assert ch["gp"] == ["c", "gp", "p"] # --------------------------------------------------------------------------- # _read_rc_servers # --------------------------------------------------------------------------- class TestReadRcServers: def _write_json(self, tmp_path, name, data): p = tmp_path / name p.write_text(json.dumps(data)) return p def test_nested_servers_format(self, tmp_path): path = self._write_json(tmp_path, "rc.json", { "servers": [ {"serverUrl": "http://localhost:4096"}, {"serverUrl": "http://localhost:4097"}, ] }) with patch("watchdog.RC_SERVERS_PATH", str(path)): urls = _read_rc_servers() assert set(urls) == {"http://localhost:4096", "http://localhost:4097"} assert urls[0] == "http://localhost:4096" def test_flat_keyed_format(self, tmp_path): path = self._write_json(tmp_path, "rc.json", { "/home/user/p1": {"serverUrl": "http://localhost:4096"}, "/home/user/p2": {"serverUrl": "http://localhost:4097"}, }) with patch("watchdog.RC_SERVERS_PATH", str(path)): urls = _read_rc_servers() assert set(urls) == {"http://localhost:4096", "http://localhost:4097"} def test_nonexistent_file_returns_empty(self): with patch("watchdog.os.path.exists", return_value=False): urls = _read_rc_servers() assert urls == [] def test_invalid_json_returns_empty(self, tmp_path): path = tmp_path / "rc.json" # Write raw bytes to ensure JSON decode fails path.write_bytes(b"not json {{{") with patch("watchdog.RC_SERVERS_PATH", str(path)): urls = _read_rc_servers() assert urls == [] def test_dedup_preserves_order(self, tmp_path): path = self._write_json(tmp_path, "rc.json", { "a": {"serverUrl": "http://localhost:4096"}, "b": {"serverUrl": "http://localhost:4096"}, "c": {"serverUrl": "http://localhost:4097"}, }) with patch("watchdog.RC_SERVERS_PATH", str(path)): urls = _read_rc_servers() assert urls == ["http://localhost:4096", "http://localhost:4097"] # --------------------------------------------------------------------------- # _make_key # --------------------------------------------------------------------------- class TestMakeKey: def test_format(self): assert _make_key("ses_root123") == "opencode-loop:ses_root123" def test_empty(self): assert _make_key("") == "opencode-loop:" # --------------------------------------------------------------------------- # _fire_alert — state machine # --------------------------------------------------------------------------- class TestFireAlert: def test_trigger_seeds_new_open_alert(self): conn, path = _make_conns() try: _fire_alert(conn, "k1", "trigger", "warning", "h1") row = conn.execute( "SELECT state, severity, dedup_key, resolved_at " "FROM watchdog_alerts" ).fetchone() assert row["state"] == "open" assert row["severity"] == "warning" assert row["dedup_key"] == "k1" assert row["resolved_at"] is None finally: _close_conn(conn, path) def test_escalate_updates_existing_open(self): conn, path = _make_conns() try: _fire_alert(conn, "dup", "trigger", "warning", "a") _fire_alert(conn, "dup", "escalate", "critical", "b") count = conn.execute( "SELECT count(*) as cnt FROM watchdog_alerts " "WHERE dedup_key='dup'" ).fetchone()["cnt"] assert count == 1 row = conn.execute( "SELECT state, severity FROM watchdog_alerts " "WHERE dedup_key='dup'" ).fetchone() assert row["state"] == "open" assert row["severity"] == "critical" finally: _close_conn(conn, path) def test_resolve_clears_open_alert(self): conn, path = _make_conns() try: _fire_alert(conn, "r1", "trigger", "warning", "alert") before = conn.execute( "SELECT resolved_at FROM watchdog_alerts " "WHERE dedup_key='r1'" ).fetchone()["resolved_at"] assert before is None _fire_alert(conn, "r1", "resolve", "info", "done") row = conn.execute( "SELECT resolved_at, state, severity " "FROM watchdog_alerts WHERE dedup_key='r1'" ).fetchone() assert row["resolved_at"] is not None assert row["state"] == "resolved" assert row["severity"] == "info" finally: _close_conn(conn, path) def test_resolve_no_op_when_already_resolved(self): conn, path = _make_conns() try: _fire_alert(conn, "r1", "trigger", "warning", "a") _fire_alert(conn, "r1", "resolve", "info", "done") _fire_alert(conn, "r1", "resolve", "info", "done2") row = conn.execute( "SELECT state FROM watchdog_alerts " "WHERE dedup_key='r1'" ).fetchone() assert row["state"] == "resolved" finally: _close_conn(conn, path) def test_trigger_reopens_with_escalate(self): """After resolve, escalate re-opens the alert with critical severity.""" conn, path = _make_conns() try: _fire_alert(conn, "r1", "trigger", "warning", "a") _fire_alert(conn, "r1", "resolve", "info", "done") row = conn.execute( "SELECT state, resolved_at FROM watchdog_alerts " "WHERE dedup_key='r1'" ).fetchone() assert row["state"] == "resolved" assert row["resolved_at"] is not None _fire_alert(conn, "r1", "escalate", "critical", "re") row = conn.execute( "SELECT state, severity, resolved_at " "FROM watchdog_alerts WHERE dedup_key='r1'" ).fetchone() assert row["state"] == "open" assert row["severity"] == "critical" assert row["resolved_at"] is not None finally: _close_conn(conn, path) def test_trigger_inserts_once(self): """_fire_alert with state='trigger' inserts at most one row per dedup_key.""" conn, path = _make_conns() try: _fire_alert(conn, "x", "trigger", "warning", "first") _fire_alert(conn, "x", "trigger", "warning", "second") count = conn.execute( "SELECT count(*) as cnt FROM watchdog_alerts" ).fetchone()["cnt"] assert count == 1 finally: _close_conn(conn, path) # --------------------------------------------------------------------------- # detect_config_from_pydantic # --------------------------------------------------------------------------- class TestDetectConfigFromPydantic: def test_fields_mapped(self): pc = MagicMock() pc.window = 300 pc.dup_min = 5 pc.top_min = 3 pc.top_min_ro = 4 pc.cum_min = 8 pc.cover_min = 0.3 pc.min_calls = 2 pc._ro_agents = ("planner", "architect") dc = detect_config_from_pydantic(pc) assert dc.window == 300 assert dc.dup_min == 5 assert dc.top_min == 3 assert dc.top_min_ro == 4 assert dc.cum_min == 8 assert dc.cover_min == 0.3 assert dc.min_calls == 2 assert dc.read_only_agent_keywords == ("planner", "architect") def test_none_ro_agents(self): pc = MagicMock() for a in ("window", "dup_min", "top_min", "top_min_ro", "cum_min", "cover_min", "min_calls"): setattr(pc, a, None) pc._ro_agents = None dc = detect_config_from_pydantic(pc) assert isinstance(dc.read_only_agent_keywords, tuple) # --------------------------------------------------------------------------- # calls_of (stubbed) # --------------------------------------------------------------------------- class TestCallsOf: def test_extracts_tool_calls(self): messages = [{ "parts": [{ "type": "tool", "tool": "edit", "state": { "input": {"filePath": "/main.py", "patches": "diff"}, "metadata": {"diff": "the diff"}, "time": {"start": 1000}, }, }], }] out = calls_of("http://localhost:4096", "ses_123", http_get=_stub_http(messages)) assert len(out) == 1 assert out[0][1] == "edit" assert out[0][3] is True # landed (diff present) def test_land_bash_git_commit(self): messages = [{ "parts": [{ "type": "tool", "tool": "bash", "state": { "input": {"command": "git commit -m 'fix'"}, "metadata": {"exit": 0}, "time": {"start": 2000}, }, }], }] out = calls_of("http://localhost:4096", "ses_123", http_get=_stub_http(messages)) assert out[0][1] == "bash" assert out[0][3] is True def test_land_false_bash_no_commit(self): messages = [{ "parts": [{ "type": "tool", "tool": "bash", "state": { "input": {"command": "ls -la"}, "metadata": {"exit": 0}, "time": {"start": 2000}, }, }], }] out = calls_of("http://localhost:4096", "ses_123", http_get=_stub_http(messages)) assert out[0][3] is False def test_land_false_no_exit(self): messages = [{ "parts": [{ "type": "tool", "tool": "bash", "state": { "input": {"command": "git commit"}, "metadata": {}, "time": {"start": 2000}, }, }], }] out = calls_of("http://localhost:4096", "ses_123", http_get=_stub_http(messages)) assert out[0][3] is False def test_non_tool_parts_ignored(self): messages = [{ "parts": [ {"type": "text", "text": "hello"}, {"type": "tool", "tool": "edit", "state": {}}, ], }] out = calls_of("http://localhost:4096", "ses_123", http_get=_stub_http(messages)) assert len(out) == 1 def test_empty_messages(self): out = calls_of("http://localhost:4096", "ses_123", http_get=_stub_http([])) assert out == [] def test_non_list_response_returns_empty(self): # Returning a dict: the for-loop iterates over keys, which are # strings. calls_of then calls "key".get("parts") which fails # — we expect this to be handled or return [] via the outer # try/except catch-all (since calls_of doesn't catch AttributeError). # Actually, dict iteration yields string keys, and calling .get() # on a string raises AttributeError. But our test checks that a # list-like response is required. Let's verify empty list first. messages = [] out = calls_of("http://localhost:4096", "ses_123", http_get=_stub_http(messages)) assert out == [] def test_http_error_returns_empty(self): out = calls_of("http://localhost:4096", "ses_123", http_get=lambda *a, **k: ( (_ for _ in ()).throw(OSError()) )) assert out == [] # --------------------------------------------------------------------------- # tick — state machine # --------------------------------------------------------------------------- class TestTick: def test_no_opencode_writes_no_opencode_tick(self): """When rc-servers is empty, writes 'no_opencode' tick row.""" conn, path = _make_conns() try: cfg = _make_cfg() cfg.database = MagicMock() cfg.database.path = path with patch("watchdog._read_rc_servers", return_value=[]): rc = tick(cfg, conn, http_get=lambda u, **kw: None) assert rc == 0 row = conn.execute( "SELECT outcome, sessions_seen " "FROM watchdog_ticks ORDER BY id DESC LIMIT 1" ).fetchone() assert row["outcome"] == "no_opencode" assert row["sessions_seen"] == 0 finally: _close_conn(conn, path) def test_servers_down_writes_no_opencode(self): """All servers unreachable → 'no_opencode' tick.""" conn, path = _make_conns() try: with patch("watchdog._read_rc_servers", return_value=["http://down:4096"]), \ patch("watchdog._fetch_sessions", return_value=None): rc = tick(cfg=_make_cfg(), conn=conn) assert rc == 0 row = conn.execute( "SELECT outcome FROM watchdog_ticks " "ORDER BY id DESC LIMIT 1" ).fetchone() assert row["outcome"] == "no_opencode" finally: _close_conn(conn, path) def test_tick_resolves_unflagged_open_alert(self): """Session that was flagged → now benign calls → resolve.""" conn, path = _make_conns() try: cfg = _make_cfg() cfg.database = MagicMock() cfg.database.path = path cfg.database.base_url = "http://localhost:4096" _fire_alert(conn, "opencode-loop:ses_root123", "trigger", "warning", "prev") conn.commit() now = datetime.now(timezone.utc) ts = now.timestamp() * 1000 # opencode call times are epoch ms messages = [{ "id": "m1", "parts": [{ "type": "tool", "tool": "read", # benign — should not flag "state": { "input": {"filePath": "/x.py"}, "metadata": {}, "time": {"start": ts}, }, }], "info": {"time": {"created": ts}}, }] sessions = _sessions_tree() # Patch route_decisions table creation to avoid OperationalError def _mock_ensure_tables(c): # Only ensure watchdog tables; route_decisions is optional pass def fake_http(url, **kw): if "message" in url: return messages return sessions with patch("watchdog._read_rc_servers", return_value=["http://localhost:4096"]), \ patch("watchdog._fetch_sessions", side_effect=lambda u, h: sessions), \ patch("watchdog._attribution", return_value=(None, None, 0.0, 0)): tick(cfg, conn, http_get=fake_http) alert = conn.execute( "SELECT state FROM watchdog_alerts " "WHERE dedup_key='opencode-loop:ses_root123'" ).fetchone() assert alert["state"] == "resolved", \ f"Expected resolved, got {alert['state']}" finally: _close_conn(conn, path) # --------------------------------------------------------------------------- # Rejection item tests — 9 new tests covering review rejection fixes # --------------------------------------------------------------------------- def _call(t, tool, args, landed=False): """Build a call tuple as progress_detect.evaluate expects.""" return (t, tool, json.dumps(args), landed) class TestFullHistoryFlag: """Rejection 1: full-history evaluation — 60 calls, only 3 recent → still flags.""" def test_full_history_flag(self): """evaluate() uses full call history, not only calls since last_tick.""" cfg = _make_cfg() dc = detect_config_from_pydantic(cfg.watchdog.detector) # 60 duplicate reads — all passed at once, no time-based filter calls = [_call(t, "read", {"filePath": "/x.py"}) for t in range(60)] flag, reason = evaluate(calls, dc, lambda _: 100, "test-session") assert flag is True, "Should flag with 60 duplicate calls in full history" assert reason["n"] >= 60 class TestPerSessionEvaluationWithRoChild: """Rejection 2: per-session evaluation — root flags normally, RO child evaluated separately.""" def test_root_flags_with_normal_thresholds(self): """Root session with non-RO title uses standard top_min threshold.""" cfg = _make_cfg() dc = detect_config_from_pydantic(cfg.watchdog.detector) calls = [_call(t, "read", {"filePath": "/x.py"}) for t in range(60)] flag, reason = evaluate(calls, dc, lambda _: 100, "normal agent") assert flag is True assert reason["ro"] is False def test_ro_child_evaluated_separately(self): """Read-only child title triggers is_ro with lower top_min_ro threshold.""" cfg = _make_cfg() cfg.watchdog.detector._ro_agents = ("explore", "librarian", "oracle") dc = detect_config_from_pydantic(cfg.watchdog.detector) calls = [_call(t, "read", {"filePath": "/x.py"}) for t in range(50)] _flag, reason = evaluate(calls, dc, lambda _: 100, "explore subagent") assert reason["ro"] is True, "RO title should set ro=True" assert reason["top"] >= dc.top_min_ro, ( f"Top count {reason['top']} should exceed RO threshold {dc.top_min_ro}" ) class TestAbortIdleSessionNotJudged: """Rejection 3: idle session not judged — session updated but no tool calls since last tick.""" def test_abort_idle_session_not_judged(self): """tick() skips sessions whose calls all predate last_tick.""" conn, path = _make_conns() try: cfg = _make_cfg() cfg.database = MagicMock() cfg.database.path = path cfg.database.base_url = "http://localhost:4096" now = datetime.now(timezone.utc) # An hour-old loop in opencode's real format (epoch MILLISECONDS): # 60 identical reads, which WOULD flag if judged. It must not be # judged, because none of its calls is newer than the last tick. # Regression: last_tick (seconds) compared against ms call times # made every historical call look recent. old_ts = (now.timestamp() - 3600) * 1000 # The tables a full tick reads, so a crash cannot pass as "no alert". conn.execute(""" CREATE TABLE IF NOT EXISTS route_decisions ( id INTEGER PRIMARY KEY AUTOINCREMENT, session_key TEXT, selected_model TEXT, selected_provider TEXT, est_cost_usd REAL DEFAULT 0, observed_at TEXT ) """) conn.execute( "CREATE TABLE IF NOT EXISTS models " "(model_id TEXT PRIMARY KEY, provider TEXT)" ) conn.commit() sessions = [{ "id": "ses_idle", "title": "idle session", "parentID": None, "time": {"created": old_ts, "updated": old_ts}, }] messages = [{ "id": f"m{i}", "parts": [{ "type": "tool", "tool": "read", "state": { "input": {"filePath": "/x.py"}, "metadata": {}, "time": {"start": old_ts + i}, }, }], "info": {"time": {"created": old_ts + i}}, } for i in range(60)] with patch("watchdog._read_rc_servers", return_value=["http://localhost:4096"]), \ patch("watchdog._fetch_sessions", return_value=sessions): tick(cfg, conn, http_get=lambda u, **kw: sessions if "message" not in u else messages) alerts = conn.execute( "SELECT count(*) as cnt FROM watchdog_alerts" ).fetchone()["cnt"] assert alerts == 0, "Idle session should not produce alerts" judged = conn.execute( "SELECT count(*) as cnt FROM watchdog_verdicts" ).fetchone()["cnt"] assert judged == 0, "Idle session should not be judged at all" ticks = conn.execute( "SELECT count(*) as cnt FROM watchdog_ticks" ).fetchone()["cnt"] assert ticks == 1, "The tick must complete and record itself" finally: _close_conn(conn, path) class TestEscalateFiresOnce: """Rejection 4: escalate fires exactly once — trigger once, escalate once, subsequent ticks produce no new rows. """ def test_escalate_fires_once(self): """_fire_alert: trigger at tick 1, escalate at tick 2; resolve and re-resolve fire once; subsequent resolve no-ops.""" conn, path = _make_conns() try: # Tick 1: trigger _fire_alert(conn, "escalation-test", "trigger", "warning", "first alert") assert conn.execute( "SELECT state FROM watchdog_alerts WHERE dedup_key='escalation-test'" ).fetchone()["state"] == "open" # Tick 2: escalate — updates severity to critical, bumps flagged_ticks _fire_alert(conn, "escalation-test", "escalate", "critical", "still flagged") rows = conn.execute( "SELECT count(*) as cnt FROM watchdog_alerts " "WHERE dedup_key='escalation-test'" ).fetchone()["cnt"] assert rows == 1, "Escalate should update existing row, not insert new" row = conn.execute( "SELECT state, severity FROM watchdog_alerts " "WHERE dedup_key='escalation-test'" ).fetchone() assert row["state"] == "open" assert row["severity"] == "critical" # Tick 3: resolve _fire_alert(conn, "escalation-test", "resolve", "info", "escalated then resolved") row = conn.execute( "SELECT state, resolved_at FROM watchdog_alerts " "WHERE dedup_key='escalation-test'" ).fetchone() assert row["state"] == "resolved" assert row["resolved_at"] is not None # Tick 4: another resolve on same key — must not change anything _fire_alert(conn, "escalation-test", "resolve", "info", "nothing to do") row = conn.execute( "SELECT state FROM watchdog_alerts " "WHERE dedup_key='escalation-test'" ).fetchone() assert row["state"] == "resolved" rows = conn.execute( "SELECT count(*) as cnt FROM watchdog_alerts " "WHERE dedup_key='escalation-test'" ).fetchone()["cnt"] assert rows == 1, "Resolve on already-resolved should not insert" finally: _close_conn(conn, path) class TestResolveFiresOnce: """Rejection 5: resolve fires at most once; historical flagged rows never re-fire resolve.""" def test_resolve_idempotent(self): """After resolve, further clean ticks fire nothing; historical resolved row stays resolved.""" conn, path = _make_conns() try: _fire_alert(conn, "resolve-once", "trigger", "warning", "alert") assert conn.execute( "SELECT state FROM watchdog_alerts WHERE dedup_key='resolve-once'" ).fetchone()["state"] == "open" # First resolve _fire_alert(conn, "resolve-once", "resolve", "info", "resolved") row = conn.execute( "SELECT state, resolved_at FROM watchdog_alerts " "WHERE dedup_key='resolve-once'" ).fetchone() assert row["state"] == "resolved" first_ts = row["resolved_at"] # Second resolve — no-op, timestamp unchanged _fire_alert(conn, "resolve-once", "resolve", "info", "resolved again") row = conn.execute( "SELECT resolved_at FROM watchdog_alerts " "WHERE dedup_key='resolve-once'" ).fetchone() assert row["resolved_at"] == first_ts, "Resolve on resolved should not update timestamp" _fire_alert(conn, "resolve-once", "escalate", "critical", "re-opened") row = conn.execute( "SELECT state, severity FROM watchdog_alerts " "WHERE dedup_key='resolve-once'" ).fetchone() assert row["state"] == "open" assert row["severity"] == "critical" finally: _close_conn(conn, path) class TestNotifierDeliverCalled: """Rejection 6: notifier.deliver called with AlertEvent containing #loops link in summary.""" def test_notifier_deliver_called_with_loops_link(self): """notifier.deliver() receives an AlertEvent whose summary includes a #loops anchor.""" event = AlertEvent( dedup_key="opencode-loop:ses_root", severity="warning", state="trigger", title="opencode loop: agent", summary="Loop detected — http://127.0.0.1:8080/admin#loops", ) notifier = Notifier(NotificationsConfig()) with patch.object(notifier, "_dispatch") as mock_dispatch: notifier.deliver(event) mock_dispatch.assert_called_once() assert "#loops" in event.summary, "Summary must contain #loops anchor" class TestAttributionUsesCKeys: """Rejection 7: attribution uses c: prefix keys, not ses_ prefix.""" def test_attribution_uses_c_keys(self): """_attribution queries by session_key with c: prefix, never ses_.""" conn, path = _make_conns() try: conn.execute(""" CREATE TABLE IF NOT EXISTS route_decisions ( id INTEGER PRIMARY KEY AUTOINCREMENT, session_key TEXT, selected_model TEXT, selected_provider TEXT, est_cost_usd REAL DEFAULT 0, observed_at TEXT ) """) conn.execute(""" CREATE TABLE IF NOT EXISTS models ( model_id TEXT PRIMARY KEY, provider TEXT ) """) conn.execute( "INSERT INTO models (model_id, provider) VALUES (?, ?)", ("gpt-4", "openai"), ) # Insert a row with c: prefix key conn.execute( "INSERT INTO route_decisions " "(session_key, selected_model, selected_provider, est_cost_usd, observed_at) " "VALUES (?, 'gpt-4', 'openai', 0.05, ?)", ("c:ses_123", "2026-01-01T00:00:00Z"), ) conn.commit() # Query using c: prefix — should match the inserted row model, prov, _cost, cnt = _attribution(conn, ["c:ses_123"], 0) assert model == "gpt-4" assert prov == "openai" assert cnt == 1 # Query using ses_ prefix — should NOT match _model2, _prov2, _cost2, cnt2 = _attribution(conn, ["ses_123"], 0) assert cnt2 == 0, "ses_ prefix should not match c: prefixed rows" finally: _close_conn(conn, path) class TestDetectorUsesInjectedFileLines: """Rejection 8: detector uses injected file_lines callable, no disk read.""" def test_detector_uses_injected_file_lines(self): """coverage() in detector uses the passed file_lines — no os.open / disk I/O.""" calls = [_call(t, "read", {"filePath": "/project/main.py", "limit": 100}) for t in range(60)] cfg = _make_cfg() dc = detect_config_from_pydantic(cfg.watchdog.detector) file_line_map = {"/project/main.py": 200} def injected_fl(fp): """Injected callable — no disk access, pure dict lookup.""" return file_line_map.get(fp, 0) or 0 flag, reason = evaluate(calls, dc, injected_fl, "test agent") assert flag is True assert reason["coverage"] > 0, ( "coverage should be > 0 when injected file_lines reports 200 lines" ) class TestBashCommentStrip: """Rejection 9: bash commands starting with # are treated as distinct targets.""" def test_bash_comment_strip(self): """bash commands with #-prefixed comments produce targets that include the comment prefix.""" cmd_comment = "# Check file permissions\nls -la" cmd_plain = "ls -la" result_comment = target_of("bash", json.dumps({"command": cmd_comment})) result_plain = target_of("bash", json.dumps({"command": cmd_plain})) assert result_comment[0] == "bash", "tool type should be bash" assert result_comment != result_plain, ( "Comment and non-comment versions of the same command " "must be distinct targets" ) def test_bash_comment_never_deduplicated_with_plain(self): """A #-prefixed bash command and a plain command on the same file are counted as different targets.""" cmd_comment = "# Check file permissions\nls -la" cmd_plain = "ls -la" result_comment = target_of("bash", json.dumps({"command": cmd_comment})) result_plain = target_of("bash", json.dumps({"command": cmd_plain})) assert result_comment != result_plain, ( "Comment and non-comment versions of the same command " "must be distinct targets" ) # --------------------------------------------------------------------------- # Descendant-landed scope — Bug 1 + Bug 2 fixes # --------------------------------------------------------------------------- class TestDescendantLandedScope: """Bug 1: descendant scope too wide (siblings masked looping worker). Bug 2: recompute block (flagged on healthy runs).""" def test_sibling_landing_does_not_mask_worker(self): """Worker with 15+ duplicate calls flags even when sibling lands commit.""" conn, path = _make_conns() try: cfg = _make_cfg() cfg.database = MagicMock() cfg.database.path = path cfg.database.base_url = "http://localhost:4096" now = datetime.now(timezone.utc) ts = now.timestamp() * 1000 # opencode call times are epoch ms # Create route_decisions table for cost_since_landed_map queries. conn.execute(""" CREATE TABLE IF NOT EXISTS route_decisions ( id INTEGER PRIMARY KEY AUTOINCREMENT, session_key TEXT, selected_model TEXT, selected_provider TEXT, est_cost_usd REAL DEFAULT 0, observed_at TEXT ) """) conn.commit() # Tree: root with two children (worker + sibling) sessions = [ {"id": "ses_p1", "title": "parent", "parentID": None, "info": {"time": {"created": ts}}}, {"id": "ses_worker", "title": "worker", "parentID": "ses_p1", "info": {"time": {"created": ts}}}, {"id": "ses_sibling", "title": "sibling", "parentID": "ses_p1", "info": {"time": {"created": ts}}}, ] # Worker: 15 identical reads (no landed) → will flag as high dup worker_msgs = [ { "id": f"wm{i}", "parts": [{ "type": "tool", "tool": "read", "state": { "input": {"filePath": "/x.py"}, "metadata": {}, "time": {"start": ts + i}, }, }], "info": {"time": {"created": ts + i}}, } for i in range(15) ] # Sibling: one landed edit → should NOT mask the worker sibling_msgs = [{ "id": "sm0", "parts": [{ "type": "tool", "tool": "edit", "state": { "input": {"filePath": "/y.py"}, "metadata": {"diff": "patch"}, "time": {"start": ts + 100}, }, }], "info": {"time": {"created": ts + 100}}, }] # Parent: one benign call so it shows as active parent_msgs = [{ "id": "pm0", "parts": [{ "type": "tool", "tool": "read", "state": { "input": {"filePath": "/z.py"}, "metadata": {}, "time": {"start": ts + 200}, }, }], "info": {"time": {"created": ts + 200}}, }] def fake_http(url, **kw): if "message" in url: if "ses_worker" in url: return worker_msgs if "ses_sibling" in url: return sibling_msgs if "ses_p1" in url: return parent_msgs return [] return sessions with patch("watchdog._read_rc_servers", return_value=["http://localhost:4096"]), \ patch("watchdog._attribution", return_value=(None, None, 0.0, 0)): tick(cfg, conn, http_get=fake_http) verdicts = conn.execute( "SELECT session_id, flagged FROM watchdog_verdicts" ).fetchall() worker_flagged = any( r["session_id"] == "ses_worker" and r["flagged"] == 1 for r in verdicts ) assert worker_flagged, ( "Worker should be flagged despite sibling landing" ) finally: _close_conn(conn, path) def test_parent_not_flagged_when_child_lands(self): """Parent with 20 identical git log calls NOT flagged when child lands.""" conn, path = _make_conns() try: cfg = _make_cfg() cfg.database = MagicMock() cfg.database.path = path cfg.database.base_url = "http://localhost:4096" now = datetime.now(timezone.utc) ts = now.timestamp() * 1000 # opencode call times are epoch ms # Create route_decisions table for cost_since_landed_map queries. conn.execute(""" CREATE TABLE IF NOT EXISTS route_decisions ( id INTEGER PRIMARY KEY AUTOINCREMENT, session_key TEXT, selected_model TEXT, selected_provider TEXT, est_cost_usd REAL DEFAULT 0, observed_at TEXT ) """) conn.commit() sessions = [ {"id": "ses_p2", "title": "parent s2", "parentID": None, "info": {"time": {"created": ts}}}, {"id": "ses_child", "title": "child", "parentID": "ses_p2", "info": {"time": {"created": ts}}}, ] # Parent: 20 identical git log calls → would flag as high dup # but child's landed should suppress the signal. parent_msgs = [ { "id": f"pm{i}", "parts": [{ "type": "tool", "tool": "bash", "state": { "input": {"command": "git log --oneline -5"}, "metadata": {"exit": 0}, "time": {"start": ts + i}, }, }], "info": {"time": {"created": ts + i}}, } for i in range(20) ] # Child: one landed edit (time must be within parent's window # [ts+0, ts+19] so extra_landed_times can suppress the flag). child_msgs = [{ "id": "cm0", "parts": [{ "type": "tool", "tool": "edit", "state": { "input": {"filePath": "/main.py"}, "metadata": {"diff": "changes"}, "time": {"start": ts + 5}, }, }], "info": {"time": {"created": ts + 5}}, }] def fake_http(url, **kw): if "message" in url: if "ses_p2" in url: return parent_msgs if "ses_child" in url: return child_msgs return [] return sessions with patch("watchdog._read_rc_servers", return_value=["http://localhost:4096"]), \ patch("watchdog._attribution", return_value=(None, None, 0.0, 0)): tick(cfg, conn, http_get=fake_http) verdicts = conn.execute( "SELECT session_id, flagged FROM watchdog_verdicts" ).fetchall() parent_flagged = any( r["session_id"] == "ses_p2" and r["flagged"] == 1 for r in verdicts ) assert not parent_flagged, ( "Parent should NOT be flagged when child lands within window" ) finally: _close_conn(conn, path) class TestMainExitCode: """`python -m watchdog --once` runs as a systemd oneshot: a tick that ran and fired alerts is a success, so the exit code must not carry the alert count (systemd would mark the unit FAILED on every real alert).""" def test_main_returns_zero_when_alerts_fire(self, tmp_path): import watchdog cfg = MagicMock() cfg.database.path = str(tmp_path / "router.db") cfg.notifications.channels = [] with patch("watchdog.load_config", return_value=cfg), \ patch("watchdog.tick", return_value=3) as fake_tick: rc = watchdog.main(["--once", "--config", "unused.yaml"]) assert fake_tick.called assert rc == 0 class TestAlertSummaryEvidence: """Trigger/escalate alert summaries carry calls, cost, and model evidence. The session loops (60 duplicate reads after a landed change), so the summary must report calls since the landed change, its cost, and the attributed model — and omit the 'on ' segment when no model is attributed. """ def _run_flagged_tick(self, conn, path, attribution, cost_after_landed): """Run one tick against a loop-flagged worker session; return events.""" cfg = _make_cfg() cfg.database = MagicMock() cfg.database.path = path cfg.database.base_url = "http://localhost:4096" cfg.watchdog.dashboard_base_url = "http://dash/" conn.execute(""" CREATE TABLE IF NOT EXISTS route_decisions ( id INTEGER PRIMARY KEY AUTOINCREMENT, session_key TEXT, selected_model TEXT, selected_provider TEXT, est_cost_usd REAL DEFAULT 0, observed_at TEXT ) """) conn.commit() now = datetime.now(timezone.utc) landed_ts = now.timestamp() * 1000 # ms # Worker: 60 identical reads (no landed) -> flags; each is after the # tree's last landed time so calls_since_landed is 60. worker_msgs = [{ "id": f"wm{i}", "parts": [{ "type": "tool", "tool": "read", "state": { "input": {"filePath": "/x.py"}, "metadata": {}, "time": {"start": landed_ts + 1 + i}, }, }], "info": {"time": {"created": landed_ts + 1 + i}}, } for i in range(60)] # Sibling: one landed edit at landed_ts -> tree_last_landed is set. sibling_msgs = [{ "id": "sm0", "parts": [{ "type": "tool", "tool": "edit", "state": { "input": {"filePath": "/y.py"}, "metadata": {"diff": "patch"}, "time": {"start": landed_ts}, }, }], "info": {"time": {"created": landed_ts}}, }] # Parent: one benign call so it shows as active (not flagged). parent_msgs = [{ "id": "pm0", "parts": [{ "type": "tool", "tool": "read", "state": { "input": {"filePath": "/z.py"}, "metadata": {}, "time": {"start": landed_ts + 200}, }, }], "info": {"time": {"created": landed_ts + 200}}, }] sessions = [ {"id": "ses_p", "title": "parent", "parentID": None, "time": {"created": landed_ts}, "info": {"time": {"created": landed_ts}}}, {"id": "ses_worker", "title": "worker", "parentID": "ses_p", "time": {"created": landed_ts}, "info": {"time": {"created": landed_ts}}}, {"id": "ses_sibling", "title": "sibling", "parentID": "ses_p", "time": {"created": landed_ts}, "info": {"time": {"created": landed_ts}}}, ] if cost_after_landed is not None: conn.execute( "INSERT INTO route_decisions (session_key, selected_model, " "selected_provider, est_cost_usd, observed_at) " "VALUES (?,?,?,?,?)", ("c:ses_worker", "gpt-4o", "neuralwatt", cost_after_landed, datetime.fromtimestamp((landed_ts + 1) / 1000, tz=timezone.utc).isoformat()), ) conn.commit() events = [] class _Recorder: def deliver(self, event): events.append(event) def fake_http(url, **kw): if "message" in url: if "ses_worker" in url: return worker_msgs if "ses_sibling" in url: return sibling_msgs if "ses_p" in url: return parent_msgs return [] return sessions with patch("watchdog._read_rc_servers", return_value=["http://localhost:4096"]), \ patch("watchdog._fetch_sessions", return_value=sessions), \ patch("watchdog._attribution", return_value=attribution): tick(cfg, conn, notifier=_Recorder(), http_get=fake_http) return events def _worker_event(self, events, state): """Return the event for the worker session of the given state.""" for e in events: if e.state == state and e.title == "opencode loop: worker": return e return None def test_alert_summary_trigger_includes_calls_cost_model(self): """Trigger summary lists calls, cost, and model since last landed change.""" conn, path = _make_conns() try: events = self._run_flagged_tick( conn, path, attribution=("gpt-4o", "neuralwatt", 0.5, 10), cost_after_landed=1.25, ) trigger = self._worker_event(events, "trigger") assert trigger, "Expected a worker trigger event" summary = trigger.summary assert "60 calls" in summary, f"calls missing: {summary}" assert "$1.25" in summary, f"cost missing: {summary}" assert "on gpt-4o" in summary, f"model missing: {summary}" assert "#loops" in summary, f"loops link missing: {summary}" finally: _close_conn(conn, path) def test_alert_summary_trigger_omits_model_when_null(self): """No 'on ' segment when model_id is NULL.""" conn, path = _make_conns() try: events = self._run_flagged_tick( conn, path, attribution=(None, None, 0.0, 0), cost_after_landed=1.25, ) trigger = self._worker_event(events, "trigger") assert trigger, "Expected a worker trigger event" summary = trigger.summary assert "60 calls" in summary, f"calls missing: {summary}" assert "on " not in summary, f"'on ' should be omitted: {summary}" finally: _close_conn(conn, path) def test_alert_summary_escalate_includes_calls_cost_model(self): """Escalate summary (3rd tick) carries the same calls/cost/model evidence.""" conn, path = _make_conns() try: events = [] for _ in range(3): events += self._run_flagged_tick( conn, path, attribution=("gpt-4o", "neuralwatt", 0.5, 10), cost_after_landed=1.25, ) escalate = self._worker_event(events, "escalate") assert escalate, "Expected a worker escalate event" summary = escalate.summary assert "60 calls" in summary, f"calls missing: {summary}" assert "$1.25" in summary, f"cost missing: {summary}" assert "on gpt-4o" in summary, f"model missing: {summary}" assert "#loops" in summary, f"loops link missing: {summary}" finally: _close_conn(conn, path)