1389 lines
52 KiB
Python
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)
|