Files
6krrt/tests/test_watchdog.py

1389 lines
52 KiB
Python

"""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 <model>' 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 <model>' 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 <model>' 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)