- run-team.py: add the 'dispatch <thread_id>' operator command (P3 option-b). The read-only box parks at DISPATCH; this completes it with a just-in-time WRITE token: reads candidate_diff + scope from the checkpoint (or --diff/--scope files), pushes the head branch + fires workflow_dispatch via dispatch_apply_verify, prints the located run_id, and (--write-back) writes it into the task checkpoint so VERIFY binds. +2 tests. - OPERATOR-RUNBOOK: fix the misleading 'systemctl show -p Environment' check (it does NOT show EnvironmentFile= vars) -> use /proc/<MainPID>/environ + _p3_env_is_configured(); document the operator-initiated dispatch flow + the fine-grained-token write-probe caveat. Suite green, ruff clean. Branch only; not merged.
1361 lines
48 KiB
Python
1361 lines
48 KiB
Python
"""Unit tests for the ``run-team.py`` operator CLI (design §3.3.1, §7.1 P1).
|
|
|
|
``run-team.py`` is a hyphenated entry script (per the design's "entry CLI
|
|
``run-team.py``"), so it cannot be imported by normal ``import`` syntax. These
|
|
tests load it via :mod:`importlib` from its file path and exercise the manual
|
|
ledger path against the FOUNDATION ``agent_team.db.schema`` ledger.
|
|
|
|
The tests assert the §3.3.1 manual-path contract: list open/parked questions,
|
|
answer-on-behalf / force-expire / supersede gated behind ``--confirm`` and
|
|
audit-logged, first-answer-wins semantics inherited from the foundation
|
|
compare-and-set, and read-only commands needing no confirmation.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import importlib.util
|
|
import io
|
|
import json
|
|
from datetime import timedelta
|
|
from pathlib import Path
|
|
from types import ModuleType
|
|
from typing import Any
|
|
|
|
import pytest
|
|
|
|
from agent_team.db.schema import QUESTION_STATES, connect, init_db
|
|
|
|
# Path to the hyphenated entry CLI (sibling of the agent_team package).
|
|
_CLI_PATH = Path(__file__).resolve().parents[1] / "run-team.py"
|
|
|
|
|
|
def _load_cli() -> ModuleType:
|
|
"""Import ``run-team.py`` from its file path as a module."""
|
|
spec = importlib.util.spec_from_file_location("run_team_cli", _CLI_PATH)
|
|
assert spec is not None and spec.loader is not None
|
|
module = importlib.util.module_from_spec(spec)
|
|
spec.loader.exec_module(module)
|
|
return module
|
|
|
|
|
|
@pytest.fixture(scope="module")
|
|
def cli() -> ModuleType:
|
|
"""The loaded run-team CLI module (loaded once per test module)."""
|
|
return _load_cli()
|
|
|
|
|
|
@pytest.fixture()
|
|
def db_path(tmp_path: Path) -> Path:
|
|
"""A fresh, initialized ledger DB for each test."""
|
|
path = tmp_path / "state" / "agent_team.sqlite"
|
|
init_db(path)
|
|
return path
|
|
|
|
|
|
@pytest.fixture()
|
|
def audit_log(tmp_path: Path) -> Path:
|
|
"""Path to a per-test audit log (not created until first destructive op)."""
|
|
return tmp_path / "state" / "audit.log.jsonl"
|
|
|
|
|
|
def _insert_question(
|
|
db_path: Path,
|
|
*,
|
|
question_id: str,
|
|
thread_id: str = "thread-a",
|
|
turn: int = 0,
|
|
status: str = "open",
|
|
transport: str = "slack",
|
|
deadline_at: str | None = None,
|
|
) -> None:
|
|
"""Insert a pending_questions row directly for test setup."""
|
|
conn = connect(db_path)
|
|
try:
|
|
conn.execute(
|
|
"INSERT INTO pending_questions "
|
|
"(question_id, thread_id, turn, status, transport, posted_at, "
|
|
"deadline_at) VALUES (?, ?, ?, ?, ?, ?, ?)",
|
|
(
|
|
question_id,
|
|
thread_id,
|
|
turn,
|
|
status,
|
|
transport,
|
|
"2026-06-17T00:00:00+00:00",
|
|
deadline_at,
|
|
),
|
|
)
|
|
finally:
|
|
conn.close()
|
|
|
|
|
|
def _status_of(db_path: Path, question_id: str) -> str | None:
|
|
conn = connect(db_path)
|
|
try:
|
|
row = conn.execute(
|
|
"SELECT status FROM pending_questions WHERE question_id = ?",
|
|
(question_id,),
|
|
).fetchone()
|
|
finally:
|
|
conn.close()
|
|
return None if row is None else row["status"]
|
|
|
|
|
|
def _run(
|
|
cli: ModuleType,
|
|
db_path: Path,
|
|
audit_log: Path,
|
|
*args: str,
|
|
) -> tuple[int, str]:
|
|
"""Invoke ``main`` with the standard global flags, capturing stdout."""
|
|
out = io.StringIO()
|
|
argv = ["--db", str(db_path), "--audit-log", str(audit_log), *args]
|
|
code = cli.main(argv, out=out)
|
|
return code, out.getvalue()
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# Foundation-import / structural assertions
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_cli_file_exists_and_is_hyphenated() -> None:
|
|
assert _CLI_PATH.name == "run-team.py"
|
|
assert _CLI_PATH.is_file()
|
|
|
|
|
|
def test_cli_imports_foundation_contracts_verbatim(cli: ModuleType) -> None:
|
|
# The CLI must import the foundation, not redefine it.
|
|
from agent_team.db import schema as foundation_schema
|
|
|
|
assert cli.answer_question is foundation_schema.answer_question
|
|
assert cli.expire_question is foundation_schema.expire_question
|
|
assert cli.supersede_question is foundation_schema.supersede_question
|
|
assert cli.connect is foundation_schema.connect
|
|
assert cli.init_db is foundation_schema.init_db
|
|
|
|
|
|
def test_build_parser_has_no_side_effects(cli: ModuleType) -> None:
|
|
parser = cli.build_parser()
|
|
assert parser.prog == "run-team.py"
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# init-db
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_init_db_creates_ledger_tables(cli: ModuleType, tmp_path: Path) -> None:
|
|
db_path = tmp_path / "state" / "fresh.sqlite"
|
|
audit_log = tmp_path / "audit.jsonl"
|
|
code, out = _run(cli, db_path, audit_log, "init-db")
|
|
assert code == 0
|
|
assert db_path.exists()
|
|
conn = connect(db_path)
|
|
try:
|
|
names = {
|
|
r["name"]
|
|
for r in conn.execute(
|
|
"SELECT name FROM sqlite_master WHERE type='table'"
|
|
).fetchall()
|
|
}
|
|
finally:
|
|
conn.close()
|
|
assert "pending_questions" in names
|
|
assert "budget_ledger" in names
|
|
|
|
|
|
def test_init_db_is_idempotent(cli: ModuleType, tmp_path: Path) -> None:
|
|
db_path = tmp_path / "state" / "fresh.sqlite"
|
|
audit_log = tmp_path / "audit.jsonl"
|
|
assert _run(cli, db_path, audit_log, "init-db")[0] == 0
|
|
assert _run(cli, db_path, audit_log, "init-db")[0] == 0
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# list / show (read-only, no confirmation)
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_list_open_default(cli: ModuleType, db_path: Path, audit_log: Path) -> None:
|
|
_insert_question(db_path, question_id="q-open", status="open")
|
|
_insert_question(db_path, question_id="q-exp", status="expired")
|
|
code, out = _run(cli, db_path, audit_log, "list")
|
|
assert code == 0
|
|
payload = json.loads(out)
|
|
ids = {row["question_id"] for row in payload}
|
|
assert ids == {"q-open"}
|
|
|
|
|
|
def test_list_all(cli: ModuleType, db_path: Path, audit_log: Path) -> None:
|
|
_insert_question(db_path, question_id="q-open", status="open")
|
|
_insert_question(db_path, question_id="q-exp", status="expired")
|
|
code, out = _run(cli, db_path, audit_log, "list", "--all")
|
|
assert code == 0
|
|
ids = {row["question_id"] for row in json.loads(out)}
|
|
assert ids == {"q-open", "q-exp"}
|
|
|
|
|
|
def test_list_parked_excludes_open(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q-open", status="open")
|
|
_insert_question(db_path, question_id="q-ans", status="answered")
|
|
_insert_question(db_path, question_id="q-exp", status="expired")
|
|
_insert_question(db_path, question_id="q-sup", status="superseded")
|
|
code, out = _run(cli, db_path, audit_log, "list", "--parked")
|
|
assert code == 0
|
|
ids = {row["question_id"] for row in json.loads(out)}
|
|
assert ids == {"q-ans", "q-exp", "q-sup"}
|
|
assert "q-open" not in ids
|
|
|
|
|
|
def test_list_status_filter(cli: ModuleType, db_path: Path, audit_log: Path) -> None:
|
|
_insert_question(db_path, question_id="q-open", status="open")
|
|
_insert_question(db_path, question_id="q-exp", status="expired")
|
|
code, out = _run(cli, db_path, audit_log, "list", "--status", "expired")
|
|
assert code == 0
|
|
ids = {row["question_id"] for row in json.loads(out)}
|
|
assert ids == {"q-exp"}
|
|
|
|
|
|
def test_list_empty_returns_empty_array(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
code, out = _run(cli, db_path, audit_log, "list")
|
|
assert code == 0
|
|
assert json.loads(out) == []
|
|
|
|
|
|
def test_show_existing(cli: ModuleType, db_path: Path, audit_log: Path) -> None:
|
|
_insert_question(db_path, question_id="q1", thread_id="t1", turn=3)
|
|
code, out = _run(cli, db_path, audit_log, "show", "q1")
|
|
assert code == 0
|
|
row = json.loads(out)
|
|
assert row["question_id"] == "q1"
|
|
assert row["thread_id"] == "t1"
|
|
assert row["turn"] == 3
|
|
|
|
|
|
def test_show_missing_returns_1(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
code, _ = _run(cli, db_path, audit_log, "show", "nope")
|
|
assert code == 1
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# Destructive actions require --confirm and are audit-logged
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_expire_without_confirm_refuses_and_does_not_mutate(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
code, _ = _run(cli, db_path, audit_log, "expire", "q1")
|
|
assert code == 1
|
|
# Unchanged: the guard fired before touching the ledger.
|
|
assert _status_of(db_path, "q1") == "open"
|
|
assert not audit_log.exists()
|
|
|
|
|
|
def test_answer_without_confirm_refuses(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
code, _ = _run(cli, db_path, audit_log, "answer", "q1", "--answer", "yes")
|
|
assert code == 1
|
|
assert _status_of(db_path, "q1") == "open"
|
|
|
|
|
|
def test_supersede_without_confirm_refuses(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
code, _ = _run(cli, db_path, audit_log, "supersede", "q1")
|
|
assert code == 1
|
|
assert _status_of(db_path, "q1") == "open"
|
|
|
|
|
|
def test_expire_with_confirm_flips_status_and_audits(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
code, out = _run(
|
|
cli, db_path, audit_log, "--operator", "adam", "expire", "q1", "--confirm"
|
|
)
|
|
assert code == 0
|
|
assert _status_of(db_path, "q1") == "expired"
|
|
entries = [json.loads(line) for line in audit_log.read_text().splitlines()]
|
|
# Attempt is recorded BEFORE the mutation, outcome after, so a mutation can
|
|
# never land without a trail (§3.3.1).
|
|
assert len(entries) == 2
|
|
assert entries[0]["phase"] == "attempt"
|
|
assert "applied" not in entries[0]
|
|
assert entries[-1]["phase"] == "outcome"
|
|
assert entries[-1]["action"] == "expire"
|
|
assert entries[-1]["question_id"] == "q1"
|
|
assert entries[-1]["operator"] == "adam"
|
|
assert entries[-1]["applied"] is True
|
|
assert "ts" in entries[-1]
|
|
|
|
|
|
def test_answer_with_confirm_flips_status_records_via(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
code, out = _run(
|
|
cli,
|
|
db_path,
|
|
audit_log,
|
|
"--operator",
|
|
"adam",
|
|
"answer",
|
|
"q1",
|
|
"--answer",
|
|
'{"choice": "B"}',
|
|
"--confirm",
|
|
)
|
|
assert code == 0
|
|
assert _status_of(db_path, "q1") == "answered"
|
|
conn = connect(db_path)
|
|
try:
|
|
row = conn.execute(
|
|
"SELECT answer_json, answered_via FROM pending_questions "
|
|
"WHERE question_id = ?",
|
|
("q1",),
|
|
).fetchone()
|
|
finally:
|
|
conn.close()
|
|
assert row["answer_json"] == '{"choice": "B"}'
|
|
assert row["answered_via"] == "cli:adam"
|
|
|
|
|
|
def test_answer_explicit_via_overrides_default(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
code, _ = _run(
|
|
cli,
|
|
db_path,
|
|
audit_log,
|
|
"answer",
|
|
"q1",
|
|
"--answer",
|
|
"x",
|
|
"--via",
|
|
"slack:U123",
|
|
"--confirm",
|
|
)
|
|
assert code == 0
|
|
conn = connect(db_path)
|
|
try:
|
|
row = conn.execute(
|
|
"SELECT answered_via FROM pending_questions WHERE question_id = ?",
|
|
("q1",),
|
|
).fetchone()
|
|
finally:
|
|
conn.close()
|
|
assert row["answered_via"] == "slack:U123"
|
|
|
|
|
|
def test_supersede_with_confirm_flips_status(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
code, _ = _run(cli, db_path, audit_log, "supersede", "q1", "--confirm")
|
|
assert code == 0
|
|
assert _status_of(db_path, "q1") == "superseded"
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# First-answer-wins / no-op semantics inherited from the foundation
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_answer_already_expired_is_noop_returns_1_but_audits(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="expired")
|
|
code, _ = _run(
|
|
cli, db_path, audit_log, "answer", "q1", "--answer", "x", "--confirm"
|
|
)
|
|
assert code == 1
|
|
# Status unchanged (compare-and-set lost), but the attempt is audited.
|
|
assert _status_of(db_path, "q1") == "expired"
|
|
entries = [json.loads(line) for line in audit_log.read_text().splitlines()]
|
|
assert entries[-1]["action"] == "answer"
|
|
assert entries[-1]["applied"] is False
|
|
|
|
|
|
def test_expire_missing_question_is_noop_returns_1(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
code, _ = _run(cli, db_path, audit_log, "expire", "ghost", "--confirm")
|
|
assert code == 1
|
|
entries = [json.loads(line) for line in audit_log.read_text().splitlines()]
|
|
assert entries[-1]["applied"] is False
|
|
|
|
|
|
def test_double_answer_second_is_noop(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
first, _ = _run(
|
|
cli, db_path, audit_log, "answer", "q1", "--answer", "a", "--confirm"
|
|
)
|
|
second, _ = _run(
|
|
cli, db_path, audit_log, "answer", "q1", "--answer", "b", "--confirm"
|
|
)
|
|
assert first == 0
|
|
assert second == 1 # first-answer-wins; second is a no-op
|
|
conn = connect(db_path)
|
|
try:
|
|
row = conn.execute(
|
|
"SELECT answer_json FROM pending_questions WHERE question_id = ?",
|
|
("q1",),
|
|
).fetchone()
|
|
finally:
|
|
conn.close()
|
|
assert row["answer_json"] == "a" # original answer preserved
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# Audit log durability (append-only, multiple actions)
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_audit_log_appends_across_actions(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
_insert_question(db_path, question_id="q2", status="open")
|
|
_run(cli, db_path, audit_log, "expire", "q1", "--confirm")
|
|
_run(cli, db_path, audit_log, "answer", "q2", "--answer", "y", "--confirm")
|
|
entries = [json.loads(line) for line in audit_log.read_text().splitlines()]
|
|
# Each destructive action writes an attempt + an outcome record (append-only).
|
|
assert len(entries) == 4
|
|
actions = [e["action"] for e in entries]
|
|
assert actions == ["expire", "expire", "answer", "answer"]
|
|
outcomes = [e["action"] for e in entries if e["phase"] == "outcome"]
|
|
assert outcomes == ["expire", "answer"]
|
|
|
|
|
|
def test_audit_entries_are_valid_json_lines(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
_run(cli, db_path, audit_log, "expire", "q1", "--confirm")
|
|
content = audit_log.read_text()
|
|
assert content.endswith("\n")
|
|
for line in content.splitlines():
|
|
json.loads(line) # raises if any line is not valid JSON
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# argparse-level usage errors
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_no_subcommand_is_usage_error(cli: ModuleType) -> None:
|
|
with pytest.raises(SystemExit) as exc:
|
|
cli.main([])
|
|
assert exc.value.code == 2
|
|
|
|
|
|
def test_unknown_status_choice_is_usage_error(cli: ModuleType) -> None:
|
|
with pytest.raises(SystemExit) as exc:
|
|
cli.main(["list", "--status", "bogus"])
|
|
assert exc.value.code == 2
|
|
|
|
|
|
def test_answer_requires_answer_flag(cli: ModuleType) -> None:
|
|
with pytest.raises(SystemExit) as exc:
|
|
cli.main(["answer", "q1", "--confirm"])
|
|
assert exc.value.code == 2
|
|
|
|
|
|
def test_parked_states_derived_from_foundation(cli: ModuleType) -> None:
|
|
# The parked-context states are exactly the non-open foundation states.
|
|
assert set(cli._PARKED_STATES) == set(QUESTION_STATES) - {"open"}
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# re-deliver + force-resume (design-named operator verbs, §3.3.1 / §6.6)
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def _set_channel_ref(db_path: Path, question_id: str, ref: str) -> None:
|
|
conn = connect(db_path)
|
|
try:
|
|
conn.execute(
|
|
"UPDATE pending_questions SET channel_ref = ? WHERE question_id = ?",
|
|
(ref, question_id),
|
|
)
|
|
finally:
|
|
conn.close()
|
|
|
|
|
|
def _channel_ref_of(db_path: Path, question_id: str) -> str | None:
|
|
conn = connect(db_path)
|
|
try:
|
|
row = conn.execute(
|
|
"SELECT channel_ref FROM pending_questions WHERE question_id = ?",
|
|
(question_id,),
|
|
).fetchone()
|
|
finally:
|
|
conn.close()
|
|
return None if row is None else row["channel_ref"]
|
|
|
|
|
|
def test_redeliver_clears_channel_ref_and_audits(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
_set_channel_ref(db_path, "q1", "slack:123.456")
|
|
code, out = _run(cli, db_path, audit_log, "redeliver", "q1")
|
|
assert code == 0
|
|
assert _channel_ref_of(db_path, "q1") is None
|
|
entries = [json.loads(line) for line in audit_log.read_text().splitlines()]
|
|
# Non-destructive but audited: attempt + outcome, no --confirm needed.
|
|
assert [e["phase"] for e in entries] == ["attempt", "outcome"]
|
|
assert entries[-1]["action"] == "redeliver"
|
|
assert entries[-1]["applied"] is True
|
|
assert entries[-1]["prior_channel_ref"] == "slack:123.456"
|
|
|
|
|
|
def test_redeliver_needs_no_confirm(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
# redeliver is not in the destructive set, so it runs without --confirm.
|
|
assert "redeliver" not in cli._DESTRUCTIVE_ACTIONS
|
|
|
|
|
|
def test_redeliver_non_open_is_noop_returns_1(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="answered")
|
|
code, _ = _run(cli, db_path, audit_log, "redeliver", "q1")
|
|
assert code == 1
|
|
entries = [json.loads(line) for line in audit_log.read_text().splitlines()]
|
|
assert entries[-1]["applied"] is False
|
|
|
|
|
|
def test_force_resume_reopens_expired_parked_question(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
# The parked case: an expired question is RE-OPENED so it can be answered,
|
|
# NOT superseded (superseding would make it permanently un-resumable).
|
|
_insert_question(db_path, question_id="q1", status="expired")
|
|
code, out = _run(
|
|
cli, db_path, audit_log, "--operator", "adam", "force-resume", "q1", "--confirm"
|
|
)
|
|
assert code == 0
|
|
assert _status_of(db_path, "q1") == "open" # reopened, not superseded
|
|
entries = [json.loads(line) for line in audit_log.read_text().splitlines()]
|
|
assert [e["phase"] for e in entries] == ["attempt", "outcome"]
|
|
assert entries[-1]["action"] == "force-resume"
|
|
assert entries[-1]["resume_requested"] is True
|
|
assert entries[-1]["applied"] is True
|
|
assert entries[-1]["operator"] == "adam"
|
|
|
|
|
|
def test_force_resume_does_not_supersede_answered_row(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
# Regression: force-resume must NOT flip an answered-but-unresumed row out of
|
|
# the state the recovery sweep resumes from. It stays 'answered'.
|
|
_insert_question(db_path, question_id="q1", status="answered")
|
|
code, _ = _run(cli, db_path, audit_log, "force-resume", "q1", "--confirm")
|
|
assert code == 0
|
|
assert _status_of(db_path, "q1") == "answered" # untouched, still resumable
|
|
|
|
|
|
def test_force_resume_on_open_is_noop(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
code, _ = _run(cli, db_path, audit_log, "force-resume", "q1", "--confirm")
|
|
assert code == 1 # an open (not parked) question has nothing to force
|
|
assert _status_of(db_path, "q1") == "open"
|
|
|
|
|
|
def test_force_resume_without_confirm_refuses(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
code, _ = _run(cli, db_path, audit_log, "force-resume", "q1")
|
|
assert code == 1
|
|
assert _status_of(db_path, "q1") == "open" # unmutated
|
|
assert not audit_log.exists() # refused before any audit (confirm-check first)
|
|
|
|
|
|
def test_force_resume_is_in_destructive_set(cli: ModuleType) -> None:
|
|
assert "force-resume" in cli._DESTRUCTIVE_ACTIONS
|
|
|
|
|
|
def test_operator_defaults_to_os_login_not_empty(
|
|
cli: ModuleType, db_path: Path, audit_log: Path
|
|
) -> None:
|
|
# AUTHZ regression: --operator defaulted to "" → non-attributable audit.
|
|
# Omitting it must record a real (non-empty) operator identity.
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
code, _ = _run(cli, db_path, audit_log, "expire", "q1", "--confirm")
|
|
assert code == 0
|
|
entries = [json.loads(line) for line in audit_log.read_text().splitlines()]
|
|
assert entries[-1]["operator"] # non-empty
|
|
assert cli._default_operator() # helper never returns empty
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# Audit-before-mutate: an unwritable audit path aborts BEFORE the ledger mutates
|
|
# (regression: previously the row was mutated, then the audit append crashed,
|
|
# leaving a mutation with no record and an uncaught traceback)
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_unwritable_audit_path_aborts_before_mutation(
|
|
cli: ModuleType, db_path: Path, tmp_path: Path
|
|
) -> None:
|
|
_insert_question(db_path, question_id="q1", status="open")
|
|
# Point the audit log at a path whose parent is a FILE, so the atomic write
|
|
# of the attempt record fails with OSError before the mutation runs.
|
|
blocker = tmp_path / "not-a-dir"
|
|
blocker.write_text("x")
|
|
bad_audit = blocker / "audit.jsonl"
|
|
code, _ = _run(cli, db_path, bad_audit, "expire", "q1", "--confirm")
|
|
assert code == 1 # clean failure, not an uncaught traceback
|
|
assert _status_of(db_path, "q1") == "open" # NOT mutated — no trail, no change
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# start / serve coordinator commands + transport factory (lazy, token-tolerant)
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
class _FakeCoordinator:
|
|
"""Records setup()/start_task() so the ``start`` CLI boundary is testable."""
|
|
|
|
instances: list[_FakeCoordinator] = []
|
|
|
|
def __init__(
|
|
self,
|
|
*,
|
|
db_path: Any,
|
|
transport: Any,
|
|
build_clarify_node: Any = None,
|
|
build_plan_node: Any = None,
|
|
review_wiring: Any = None,
|
|
build_verify_wiring: Any = None,
|
|
dispatch_node_wiring: Any = None,
|
|
notify: Any = None,
|
|
alarm_hook: Any = None,
|
|
ci_poller: Any = None,
|
|
ci_timeout: Any = None,
|
|
) -> None:
|
|
self.db_path = db_path
|
|
self.transport = transport
|
|
self.notify = notify
|
|
self.alarm_hook = alarm_hook
|
|
# The production CLI opts the coordinator into the P2 graph by injecting
|
|
# these factories; record them so the wiring is asserted, not ignored.
|
|
self.build_clarify_node = build_clarify_node
|
|
self.build_plan_node = build_plan_node
|
|
self.review_wiring = review_wiring
|
|
# P3 fail-safe serve default (Decision 5): only the ``serve`` command
|
|
# auto-binds these; start/intake leave them None.
|
|
self.build_verify_wiring = build_verify_wiring
|
|
self.dispatch_node_wiring = dispatch_node_wiring
|
|
# CI-watcher seams (§4 Decision 2): the live serve path binds the poller +
|
|
# timeout here and the provider post-construction; inert leaves all None.
|
|
self.ci_poller = ci_poller
|
|
self.ci_timeout = ci_timeout
|
|
self._ci_pending_provider: Any = None
|
|
# Draft-PR runaway/stale monitor provider (P3 A4): the live serve path
|
|
# binds a read-only enumerator here post-construction; inert leaves None.
|
|
self._draft_pr_provider: Any = None
|
|
self.setup_called = False
|
|
self.start_kwargs: dict[str, Any] | None = None
|
|
self.new_task_callback: Any = None
|
|
_FakeCoordinator.instances.append(self)
|
|
|
|
def _enumerate_ci_pending(self) -> list[Any]:
|
|
# Stand-in for the durable enumerator the live path binds as the
|
|
# ci_pending_provider; identity is what the wiring test asserts.
|
|
return []
|
|
|
|
def setup(self) -> None:
|
|
self.setup_called = True
|
|
|
|
def set_new_task_callback(self, callback: Any) -> None:
|
|
self.new_task_callback = callback
|
|
|
|
def start_task(self, *, task_text: str, transport_name: str) -> str:
|
|
self.start_kwargs = {"task_text": task_text, "transport_name": transport_name}
|
|
return "thread-minted-42"
|
|
|
|
|
|
def test_start_runs_setup_and_start_task_and_prints_thread_id(
|
|
cli: ModuleType,
|
|
db_path: Path,
|
|
audit_log: Path,
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
"""``start --dry-run`` builds a coordinator, runs setup + start_task, prints id.
|
|
|
|
``_build_coordinator`` imports ``Coordinator`` lazily from
|
|
``agent_team.coordinator``, so patching it there intercepts construction.
|
|
``--dry-run`` means no Slack token is required (the real
|
|
``_build_transport`` returns a ``_DryRunTransport``).
|
|
"""
|
|
_FakeCoordinator.instances.clear()
|
|
monkeypatch.setattr(
|
|
"agent_team.coordinator.Coordinator", _FakeCoordinator, raising=True
|
|
)
|
|
|
|
code, out = _run(
|
|
cli, db_path, audit_log, "start", "--dry-run", "--task", "do the thing"
|
|
)
|
|
|
|
assert code == 0
|
|
assert len(_FakeCoordinator.instances) == 1
|
|
coord = _FakeCoordinator.instances[0]
|
|
assert coord.setup_called is True
|
|
assert coord.start_kwargs == {
|
|
"task_text": "do the thing",
|
|
"transport_name": "slack",
|
|
}
|
|
# The minted thread_id is printed to the captured stdout.
|
|
assert out.strip() == "thread-minted-42"
|
|
# --dry-run substitutes the non-posting transport (no token needed).
|
|
assert isinstance(coord.transport, cli._DryRunTransport)
|
|
# The production CLI opts into the full P2 graph: planner + review factories
|
|
# are injected (callables), not left at the P1-stub default of None.
|
|
assert callable(coord.build_plan_node)
|
|
assert callable(coord.review_wiring)
|
|
|
|
|
|
def test_build_transport_live_github_builds_github_transport(
|
|
cli: ModuleType, monkeypatch: pytest.MonkeyPatch
|
|
) -> None:
|
|
"""``--transport github`` builds the live GitHub transport from the env thread.
|
|
|
|
Token resolution is deferred to call time, so with GITHUB_TOKEN set and the
|
|
issue thread configured the builder constructs a GitHubTransport bound to
|
|
that ``owner/repo#issue`` (no network — only the issue thread is wired).
|
|
"""
|
|
monkeypatch.setenv("GITHUB_TOKEN", "ghp_test")
|
|
monkeypatch.setenv("GITHUB_OWNER", "Sea-Haven-Industries")
|
|
monkeypatch.setenv("GITHUB_REPO", "orchestrator")
|
|
monkeypatch.setenv("GITHUB_ISSUE_NUMBER", "42")
|
|
|
|
from agent_team.transport.github_adapter import GitHubTransport
|
|
|
|
args = argparse.Namespace(dry_run=False, transport="github")
|
|
transport = cli._build_transport(args)
|
|
|
|
assert isinstance(transport, GitHubTransport)
|
|
assert transport.owner == "Sea-Haven-Industries"
|
|
assert transport.repo == "orchestrator"
|
|
assert transport.issue_number == 42
|
|
|
|
|
|
def test_build_transport_live_github_missing_thread_raises_system_exit(
|
|
cli: ModuleType, monkeypatch: pytest.MonkeyPatch
|
|
) -> None:
|
|
"""A github transport with no issue thread configured fails loudly."""
|
|
monkeypatch.delenv("GITHUB_OWNER", raising=False)
|
|
monkeypatch.delenv("GITHUB_REPO", raising=False)
|
|
monkeypatch.delenv("GITHUB_ISSUE_NUMBER", raising=False)
|
|
|
|
args = argparse.Namespace(dry_run=False, transport="github")
|
|
with pytest.raises(SystemExit, match="GITHUB_OWNER"):
|
|
cli._build_transport(args)
|
|
|
|
|
|
def test_build_transport_live_github_non_integer_issue_raises_system_exit(
|
|
cli: ModuleType, monkeypatch: pytest.MonkeyPatch
|
|
) -> None:
|
|
"""A non-integer GITHUB_ISSUE_NUMBER fails loudly rather than half-building."""
|
|
monkeypatch.setenv("GITHUB_OWNER", "o")
|
|
monkeypatch.setenv("GITHUB_REPO", "r")
|
|
monkeypatch.setenv("GITHUB_ISSUE_NUMBER", "not-a-number")
|
|
|
|
args = argparse.Namespace(dry_run=False, transport="github")
|
|
with pytest.raises(SystemExit, match="must be an integer"):
|
|
cli._build_transport(args)
|
|
|
|
|
|
def test_build_transport_live_claude_code_builds_file_drop_transport(
|
|
cli: ModuleType, monkeypatch: pytest.MonkeyPatch, tmp_path: Path
|
|
) -> None:
|
|
"""``--transport claude_code`` builds the live file-drop transport from env."""
|
|
monkeypatch.setenv("CLAUDE_CODE_DROP_DIR", str(tmp_path / "drops"))
|
|
|
|
from agent_team.transport.claude_code_adapter import ClaudeCodeAdapter
|
|
|
|
args = argparse.Namespace(dry_run=False, transport="claude_code")
|
|
transport = cli._build_transport(args)
|
|
|
|
assert isinstance(transport, ClaudeCodeAdapter)
|
|
|
|
|
|
def test_build_transport_live_claude_code_missing_drop_dir_raises_system_exit(
|
|
cli: ModuleType, monkeypatch: pytest.MonkeyPatch
|
|
) -> None:
|
|
"""claude_code with no CLAUDE_CODE_DROP_DIR fails loudly."""
|
|
monkeypatch.delenv("CLAUDE_CODE_DROP_DIR", raising=False)
|
|
|
|
args = argparse.Namespace(dry_run=False, transport="claude_code")
|
|
with pytest.raises(SystemExit, match="CLAUDE_CODE_DROP_DIR"):
|
|
cli._build_transport(args)
|
|
|
|
|
|
def test_intake_github_polls_and_starts_tasks_dry_run(
|
|
cli: ModuleType,
|
|
db_path: Path,
|
|
audit_log: Path,
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
"""``intake-github --dry-run`` builds a coordinator, polls, and ingests issues.
|
|
|
|
The lazily-imported ``Coordinator`` is replaced with the fake (so no graph /
|
|
SDK is built), and ``build_default_issue_client`` is patched to return an
|
|
in-memory fake lister (so no GitHub network). The fake coordinator records
|
|
the ``start_task`` calls the intake leaf makes — one per labeled issue.
|
|
"""
|
|
_FakeCoordinator.instances.clear()
|
|
monkeypatch.setattr(
|
|
"agent_team.coordinator.Coordinator", _FakeCoordinator, raising=True
|
|
)
|
|
|
|
class _FakeIssueClient:
|
|
def list_open_issues(self, *, label: str) -> list[dict[str, Any]]:
|
|
assert label == "agent-team"
|
|
return [
|
|
{"id": 1001, "title": "Do thing A", "body": "details A"},
|
|
{"id": 1002, "title": "Do thing B", "body": ""},
|
|
]
|
|
|
|
def _fake_build_client(*, owner: str, repo: str, **_kw: Any) -> Any:
|
|
assert owner == "Sea-Haven-Industries"
|
|
assert repo == "orchestrator"
|
|
return _FakeIssueClient()
|
|
|
|
monkeypatch.setattr(
|
|
"agent_team.transport.github_intake.build_default_issue_client",
|
|
_fake_build_client,
|
|
raising=True,
|
|
)
|
|
|
|
code, out = _run(
|
|
cli,
|
|
db_path,
|
|
audit_log,
|
|
"intake-github",
|
|
"--dry-run",
|
|
"--owner",
|
|
"Sea-Haven-Industries",
|
|
"--repo",
|
|
"orchestrator",
|
|
"--label",
|
|
"agent-team",
|
|
)
|
|
|
|
assert code == 0
|
|
assert len(_FakeCoordinator.instances) == 1
|
|
coord = _FakeCoordinator.instances[0]
|
|
assert coord.setup_called is True
|
|
# Both labeled issues were ingested; the leaf passes title+body and the
|
|
# GitHub transport name. (The fake records only the LAST call's kwargs.)
|
|
assert coord.start_kwargs == {
|
|
"task_text": "Do thing B",
|
|
"transport_name": "github",
|
|
}
|
|
# The ingested issue ids are printed (one per line).
|
|
assert out.split() == ["1001", "1002"]
|
|
|
|
|
|
def test_intake_github_label_required(cli: ModuleType) -> None:
|
|
"""``intake-github`` requires --owner/--repo/--label (argparse usage error)."""
|
|
parser = cli.build_parser()
|
|
with pytest.raises(SystemExit):
|
|
parser.parse_args(["intake-github", "--owner", "o", "--repo", "r"])
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# serve fail-safe P3 wiring default (design Decision 5; UNIT 0e)
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def _serve_args(db_path: Path, command: str) -> argparse.Namespace:
|
|
"""A minimal args namespace for ``_build_coordinator`` (dry-run, no token)."""
|
|
return argparse.Namespace(
|
|
command=command,
|
|
db=db_path,
|
|
transport="slack",
|
|
dry_run=True,
|
|
)
|
|
|
|
|
|
def test_serve_binds_inert_p3_wiring_when_env_unset(
|
|
cli: ModuleType,
|
|
db_path: Path,
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
"""``serve`` is the new P3 default but degrades to inert (None) when the P3
|
|
env is unset — never crashing serve-start."""
|
|
for name in (
|
|
"AGENT_TEAM_REPO_OWNER",
|
|
"AGENT_TEAM_REPO_NAME",
|
|
"AGENT_TEAM_CI_READ_TOKEN",
|
|
"GITHUB_TOKEN",
|
|
):
|
|
monkeypatch.delenv(name, raising=False)
|
|
_FakeCoordinator.instances.clear()
|
|
monkeypatch.setattr(
|
|
"agent_team.coordinator.Coordinator", _FakeCoordinator, raising=True
|
|
)
|
|
|
|
coord = cli._build_coordinator(_serve_args(db_path, "serve"))
|
|
|
|
assert coord.build_verify_wiring is None
|
|
assert coord.dispatch_node_wiring is None
|
|
# Inert box: no CI-watcher seams, so the tick() CI sweep is a NO-OP.
|
|
assert coord.ci_poller is None
|
|
assert coord.ci_timeout is None
|
|
assert coord._ci_pending_provider is None
|
|
# A4 draft-PR monitor stays inert too, gated on the same signal as ci_poller:
|
|
# the sweep is a NO-OP, so no production runaway/stale sweep ever fires.
|
|
assert coord._draft_pr_provider is None
|
|
|
|
|
|
def test_serve_binds_live_p3_wiring_when_env_set(
|
|
cli: ModuleType,
|
|
db_path: Path,
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
"""With the P3 env provisioned, ``serve`` binds the live wiring pair."""
|
|
monkeypatch.setenv("AGENT_TEAM_REPO_OWNER", "Sea-Haven-Industries")
|
|
monkeypatch.setenv("AGENT_TEAM_REPO_NAME", "orchestrator")
|
|
monkeypatch.setenv("AGENT_TEAM_CI_READ_TOKEN", "ghp_test")
|
|
_FakeCoordinator.instances.clear()
|
|
monkeypatch.setattr(
|
|
"agent_team.coordinator.Coordinator", _FakeCoordinator, raising=True
|
|
)
|
|
|
|
coord = cli._build_coordinator(_serve_args(db_path, "serve"))
|
|
|
|
assert callable(coord.build_verify_wiring)
|
|
assert callable(coord.dispatch_node_wiring)
|
|
# Live box: the CI-watcher seams are wired so a task suspended at VERIFY
|
|
# awaiting CI gets resumed/parked rather than waiting forever. The poller is
|
|
# the read-only default; the provider is the coordinator's durable
|
|
# enumerator (bound post-construction); the timeout is the 30-min default.
|
|
assert callable(coord.ci_poller)
|
|
assert coord.ci_timeout == timedelta(minutes=30)
|
|
assert coord._ci_pending_provider == coord._enumerate_ci_pending
|
|
# A4 draft-PR monitor: the live serve path binds a read-only provider so the
|
|
# sweep has a real snapshot to ALARM / remind on (gated on the same live-pair
|
|
# signal as ci_poller). Without this the wired sweep would always see no PRs.
|
|
assert callable(coord._draft_pr_provider)
|
|
|
|
|
|
def test_serve_draft_pr_provider_reads_open_apply_draft_prs(
|
|
cli: ModuleType,
|
|
db_path: Path,
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
"""The bound draft-PR provider is READ-ONLY: it shells a scoped ``gh pr list``
|
|
(no write/close/dispatch) and maps the JSON into ``DraftPr`` snapshots."""
|
|
from agent_team.draft_pr_monitor import DraftPr
|
|
|
|
monkeypatch.setenv("AGENT_TEAM_REPO_OWNER", "Sea-Haven-Industries")
|
|
monkeypatch.setenv("AGENT_TEAM_REPO_NAME", "orchestrator")
|
|
monkeypatch.setenv("AGENT_TEAM_CI_READ_TOKEN", "ghp_test")
|
|
_FakeCoordinator.instances.clear()
|
|
monkeypatch.setattr(
|
|
"agent_team.coordinator.Coordinator", _FakeCoordinator, raising=True
|
|
)
|
|
|
|
captured: dict[str, Any] = {}
|
|
|
|
class _FakeProc:
|
|
stdout = json.dumps(
|
|
[
|
|
{
|
|
"number": 7,
|
|
"createdAt": "2026-06-23T10:00:00Z",
|
|
"updatedAt": "2026-06-23T10:05:00Z",
|
|
},
|
|
{
|
|
"number": 9,
|
|
"createdAt": "2026-06-10T09:00:00Z",
|
|
"updatedAt": "2026-06-10T09:00:00Z",
|
|
},
|
|
]
|
|
)
|
|
|
|
def _fake_run(args: Any, **kwargs: Any) -> _FakeProc:
|
|
# The provider must shell a READ-ONLY, namespace-scoped enumeration —
|
|
# never a write/close/merge subcommand.
|
|
captured["args"] = args
|
|
captured["kwargs"] = kwargs
|
|
return _FakeProc()
|
|
|
|
monkeypatch.setattr("subprocess.run", _fake_run, raising=True)
|
|
|
|
coord = cli._build_coordinator(_serve_args(db_path, "serve"))
|
|
assert callable(coord._draft_pr_provider)
|
|
|
|
prs = coord._draft_pr_provider()
|
|
|
|
# READ-ONLY + scoped: it is a `gh pr list` over the apply/ head namespace, not
|
|
# a mutating subcommand, and never carries --shell.
|
|
assert captured["args"][:3] == ["gh", "pr", "list"]
|
|
assert "--draft" in captured["args"]
|
|
assert f"head:{cli._DRAFT_PR_HEAD_PREFIX}" in captured["args"]
|
|
assert "Sea-Haven-Industries/orchestrator" in captured["args"]
|
|
assert not any(
|
|
tok in captured["args"] for tok in ("close", "merge", "edit", "ready")
|
|
)
|
|
# The JSON rows map field-for-field into the snapshot the monitor expects.
|
|
assert prs == [
|
|
DraftPr(
|
|
number=7,
|
|
opened_at="2026-06-23T10:00:00Z",
|
|
updated_at="2026-06-23T10:05:00Z",
|
|
),
|
|
DraftPr(
|
|
number=9,
|
|
opened_at="2026-06-10T09:00:00Z",
|
|
updated_at="2026-06-10T09:00:00Z",
|
|
),
|
|
]
|
|
|
|
|
|
def test_serve_draft_pr_provider_fails_soft_on_gh_error(
|
|
cli: ModuleType,
|
|
db_path: Path,
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
"""A non-zero ``gh`` exit yields an empty snapshot (the sweep no-ops) rather
|
|
than raising and breaking the tick loop."""
|
|
import subprocess
|
|
|
|
monkeypatch.setenv("AGENT_TEAM_REPO_OWNER", "Sea-Haven-Industries")
|
|
monkeypatch.setenv("AGENT_TEAM_REPO_NAME", "orchestrator")
|
|
monkeypatch.setenv("AGENT_TEAM_CI_READ_TOKEN", "ghp_test")
|
|
_FakeCoordinator.instances.clear()
|
|
monkeypatch.setattr(
|
|
"agent_team.coordinator.Coordinator", _FakeCoordinator, raising=True
|
|
)
|
|
|
|
def _boom(args: Any, **kwargs: Any) -> Any:
|
|
raise subprocess.CalledProcessError(returncode=1, cmd=args)
|
|
|
|
monkeypatch.setattr("subprocess.run", _boom, raising=True)
|
|
|
|
coord = cli._build_coordinator(_serve_args(db_path, "serve"))
|
|
assert coord._draft_pr_provider() == []
|
|
|
|
|
|
def test_start_does_not_bind_p3_wiring_even_when_env_set(
|
|
cli: ModuleType,
|
|
db_path: Path,
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
"""Only ``serve`` (the daemon) auto-binds P3; the one-shot ``start`` path runs
|
|
to the first human gate and never wires build/verify/dispatch."""
|
|
monkeypatch.setenv("AGENT_TEAM_REPO_OWNER", "Sea-Haven-Industries")
|
|
monkeypatch.setenv("AGENT_TEAM_REPO_NAME", "orchestrator")
|
|
monkeypatch.setenv("AGENT_TEAM_CI_READ_TOKEN", "ghp_test")
|
|
_FakeCoordinator.instances.clear()
|
|
monkeypatch.setattr(
|
|
"agent_team.coordinator.Coordinator", _FakeCoordinator, raising=True
|
|
)
|
|
|
|
coord = cli._build_coordinator(_serve_args(db_path, "start"))
|
|
|
|
assert coord.build_verify_wiring is None
|
|
assert coord.dispatch_node_wiring is None
|
|
|
|
|
|
def test_build_transport_dry_run_returns_dry_run_transport(cli: ModuleType) -> None:
|
|
"""``dry_run=True`` yields a _DryRunTransport whose post returns a synthetic ref."""
|
|
args = argparse.Namespace(dry_run=True, transport="slack")
|
|
transport = cli._build_transport(args)
|
|
|
|
assert isinstance(transport, cli._DryRunTransport)
|
|
ref = transport.post_question(
|
|
thread_id="t1",
|
|
question_id="q1",
|
|
turn=0,
|
|
question_set=None,
|
|
deadline="2026-06-18T00:00:00+00:00",
|
|
)
|
|
assert ref == "dry-run:q1"
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# fix subcommand (Plane-1 Tier-3 fixer dry-run; §7 Phase 5)
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def _write_dep_report(path: Path, finding_id: str = "r-vuln-1") -> Path:
|
|
"""Write a minimal dependency-cve.json report with one confirmed finding."""
|
|
report = {
|
|
"checker": "dependency-cve",
|
|
"findings": [
|
|
{
|
|
"repo": "r",
|
|
"id": finding_id,
|
|
"title": "requests 2.19.0 is vulnerable (CVE-2018-18074)",
|
|
"severity": "high",
|
|
"category": "other",
|
|
"check": "vulnerable-dependency",
|
|
"status": "confirmed",
|
|
"proof": {
|
|
"package": "requests",
|
|
"version": "2.19.0",
|
|
"advisory_id": "CVE-2018-18074",
|
|
"summary": "leaks auth on redirect",
|
|
"fixed_version": "2.20.0",
|
|
},
|
|
}
|
|
],
|
|
}
|
|
path.write_text(json.dumps(report), encoding="utf-8")
|
|
return path
|
|
|
|
|
|
def test_fix_dry_run_prints_plan_and_dispatches_nothing(
|
|
cli: ModuleType, tmp_path: Path, monkeypatch: pytest.MonkeyPatch
|
|
) -> None:
|
|
"""``fix --dry-run`` plans a fix (fake model seams) and prints what it WOULD
|
|
dispatch — without firing a workflow."""
|
|
from agent_team.nodes import fixer as fixer_mod
|
|
|
|
report = _write_dep_report(tmp_path / "dependency-cve.json")
|
|
|
|
# Inject fake spec + build seams so no Claude/DeepSeek/network is touched.
|
|
bump = (
|
|
"diff --git a/requirements.txt b/requirements.txt\n"
|
|
"--- a/requirements.txt\n"
|
|
"+++ b/requirements.txt\n"
|
|
"@@ -1,1 +1,1 @@\n"
|
|
"-requests==2.19.0\n"
|
|
"+requests==2.20.0\n"
|
|
)
|
|
real_plan_fix = fixer_mod.plan_fix
|
|
|
|
def _patched_plan_fix(finding: Any, **kw: Any) -> Any:
|
|
kw.setdefault("spec", lambda _f: "bump it")
|
|
kw.setdefault("build", lambda _i: bump)
|
|
return real_plan_fix(finding, **kw)
|
|
|
|
monkeypatch.setattr(fixer_mod, "plan_fix", _patched_plan_fix)
|
|
|
|
out = io.StringIO()
|
|
rc = cli.main(
|
|
[
|
|
"fix",
|
|
"--report",
|
|
str(report),
|
|
"--finding-id",
|
|
"r-vuln-1",
|
|
"--task-id",
|
|
"t-cli-1",
|
|
"--dry-run",
|
|
],
|
|
out=out,
|
|
)
|
|
assert rc == 0
|
|
text = out.getvalue()
|
|
assert "workflow_dispatch inputs" in text
|
|
assert "candidate-diff-t-cli-1" in text
|
|
assert "DRY-RUN: nothing dispatched" in text
|
|
|
|
|
|
def test_fix_requires_dry_run(cli: ModuleType, tmp_path: Path) -> None:
|
|
"""Without --dry-run the fix command refuses (live dispatch is gated)."""
|
|
report = _write_dep_report(tmp_path / "dependency-cve.json")
|
|
with pytest.raises(SystemExit):
|
|
cli.main(
|
|
[
|
|
"fix",
|
|
"--report",
|
|
str(report),
|
|
"--finding-id",
|
|
"r-vuln-1",
|
|
"--task-id",
|
|
"t",
|
|
],
|
|
out=io.StringIO(),
|
|
)
|
|
|
|
|
|
def test_fix_unknown_finding_id_exits(cli: ModuleType, tmp_path: Path) -> None:
|
|
report = _write_dep_report(tmp_path / "dependency-cve.json")
|
|
with pytest.raises(SystemExit):
|
|
cli.main(
|
|
[
|
|
"fix",
|
|
"--report",
|
|
str(report),
|
|
"--finding-id",
|
|
"does-not-exist",
|
|
"--task-id",
|
|
"t",
|
|
"--dry-run",
|
|
],
|
|
out=io.StringIO(),
|
|
)
|
|
|
|
|
|
def test_fix_non_fixable_finding_returns_one(
|
|
cli: ModuleType, tmp_path: Path, monkeypatch: pytest.MonkeyPatch
|
|
) -> None:
|
|
"""A finding that is not a confirmed dependency-cve yields a fail-safe plan
|
|
(rc=1) and dispatches nothing."""
|
|
path = tmp_path / "dependency-cve.json"
|
|
report = {
|
|
"findings": [
|
|
{
|
|
"repo": "r",
|
|
"id": "r-not-dep",
|
|
"title": "x",
|
|
"severity": "high",
|
|
"category": "injection",
|
|
"check": "sqli",
|
|
"status": "confirmed",
|
|
"proof": {"input": "x", "outcome": "y"},
|
|
}
|
|
]
|
|
}
|
|
path.write_text(json.dumps(report), encoding="utf-8")
|
|
|
|
out = io.StringIO()
|
|
rc = cli.main(
|
|
[
|
|
"fix",
|
|
"--report",
|
|
str(path),
|
|
"--finding-id",
|
|
"r-not-dep",
|
|
"--task-id",
|
|
"t",
|
|
"--dry-run",
|
|
],
|
|
out=out,
|
|
)
|
|
assert rc == 1
|
|
assert "FIX NOT PLANNED" in out.getvalue()
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# _build_notifiers — the Slack lifecycle-milestone sink (one-thread-per-task)
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
|
|
def test_notify_sink_forwards_thread_ts(
|
|
cli: ModuleType, monkeypatch: pytest.MonkeyPatch
|
|
) -> None:
|
|
"""The notify sink must pass `thread_ts` so milestones thread under the task.
|
|
|
|
Regression: the sink was `def notify(message)` with no `thread_ts`, so the
|
|
coordinator's `_emit(message, thread_ts=root)` hit a TypeError and silently
|
|
fell back to a TOP-LEVEL post — every plan-ready/parked/failed milestone
|
|
landed unthreaded. build_slack_poster already forwards a `thread_ts` payload
|
|
key, so the only gap was this wrapper.
|
|
"""
|
|
captured: list[dict] = []
|
|
monkeypatch.setenv("SLACK_CHANNEL_ID", "C123")
|
|
# _build_notifiers imports build_slack_poster from slack_live at call time.
|
|
import agent_team.transport.slack_live as slack_live
|
|
|
|
monkeypatch.setattr(
|
|
slack_live,
|
|
"build_slack_poster",
|
|
lambda: lambda payload: captured.append(payload),
|
|
)
|
|
|
|
args = argparse.Namespace(dry_run=False, transport="slack")
|
|
notify, _alarm = cli._build_notifiers(args)
|
|
assert notify is not None
|
|
|
|
notify("threaded milestone", thread_ts="ROOT.TS")
|
|
notify("top-level milestone")
|
|
|
|
assert captured[0] == {
|
|
"channel": "C123",
|
|
"text": "threaded milestone",
|
|
"thread_ts": "ROOT.TS",
|
|
}
|
|
# No thread_ts when none is given (top-level post, not a broken key).
|
|
assert "thread_ts" not in captured[1]
|
|
assert captured[1] == {"channel": "C123", "text": "top-level milestone"}
|
|
|
|
|
|
def test_dispatch_from_files_invokes_dispatcher(
|
|
cli: ModuleType, db_path: Path, audit_log: Path, tmp_path: Path, monkeypatch
|
|
) -> None:
|
|
"""`dispatch` reads a diff/scope file and fires dispatch_apply_verify (P3 op-b)."""
|
|
from agent_team import dispatcher as d
|
|
|
|
diff_f = tmp_path / "d.diff"
|
|
diff_f.write_text("diff --git a/x b/x\n@@ -1 +1 @@\n-a\n+b\n", encoding="utf-8")
|
|
scope_f = tmp_path / "s.txt"
|
|
scope_f.write_text("agent_team/\n", encoding="utf-8")
|
|
|
|
calls: dict = {}
|
|
|
|
def _fake_dispatch(*, owner, repo, task_id, diff_text, declared_scope, base):
|
|
calls.update(
|
|
owner=owner, repo=repo, task_id=task_id, scope=declared_scope, base=base
|
|
)
|
|
return d.DispatchResult(
|
|
inputs=d.build_dispatch_inputs(
|
|
task_id=task_id, diff_text=diff_text, declared_scope=declared_scope
|
|
),
|
|
run_id="27990718108",
|
|
dispatched_at="2026-06-24T00:00:00Z",
|
|
correlation_tag=task_id,
|
|
)
|
|
|
|
monkeypatch.setattr(d, "dispatch_apply_verify", _fake_dispatch)
|
|
code, out = _run(
|
|
cli,
|
|
db_path,
|
|
audit_log,
|
|
"dispatch",
|
|
"task-xyz",
|
|
"--owner",
|
|
"Sea-Haven-Industries",
|
|
"--repo",
|
|
"orchestrator",
|
|
"--diff",
|
|
str(diff_f),
|
|
"--scope",
|
|
str(scope_f),
|
|
)
|
|
assert code == 0, out
|
|
assert calls["owner"] == "Sea-Haven-Industries"
|
|
assert calls["repo"] == "orchestrator"
|
|
assert calls["task_id"] == "task-xyz"
|
|
assert "agent_team/" in calls["scope"]
|
|
assert "27990718108" in out # the located run_id is reported
|
|
|
|
|
|
def test_dispatch_requires_owner_repo(
|
|
cli: ModuleType, db_path: Path, audit_log: Path, tmp_path: Path, monkeypatch
|
|
) -> None:
|
|
"""Without owner/repo (args or env) dispatch refuses with exit 2, no dispatch."""
|
|
monkeypatch.delenv("AGENT_TEAM_REPO_OWNER", raising=False)
|
|
monkeypatch.delenv("AGENT_TEAM_REPO_NAME", raising=False)
|
|
diff_f = tmp_path / "d.diff"
|
|
diff_f.write_text("diff --git a/x b/x\n", encoding="utf-8")
|
|
code, _ = _run(cli, db_path, audit_log, "dispatch", "t1", "--diff", str(diff_f))
|
|
assert code == 2
|