From 23135c14a6ae09bf7ce341f769f4e502c46383da Mon Sep 17 00:00:00 2001 From: Adam Moussa Date: Thu, 18 Jun 2026 15:53:04 -0400 Subject: [PATCH] feat(agent-team): P5 checker-finding intake module + tests Add agent_team.transport.checker_intake: turn a confirmed, at/above-threshold Plane-1 checker FINDING into one Plane-2 pipeline remediation task via the committed coordinator intake entry (start_task), mirroring github_intake. - select_findings: status==confirmed AND severity>=threshold (default high); unverified/suppressed/below-threshold dropped; unknown threshold rejected. - finding_identity: stable de-dup key (finding id, else content-hash). In-memory set, best-effort, NOT durable across restart (ledger table is the follow-up). - finding_task_text/_sanitize: every repo-controlled field (title, proof, repo) is newline/control-char neutralised and length-bounded before it reaches the task text or operator log (log-injection hygiene). - load_report_findings/ingest_reports: read the exact checker report JSON shape (top-level object with findings[]; bare array and dir-of-*.json also accepted). 28 hermetic unit tests (stub coordinator, in-memory findings / temp reports). --- .../agent_team/transport/checker_intake.py | 404 ++++++++++++++++ agent-team/tests/test_checker_intake.py | 434 ++++++++++++++++++ 2 files changed, 838 insertions(+) create mode 100644 agent-team/agent_team/transport/checker_intake.py create mode 100644 agent-team/tests/test_checker_intake.py diff --git a/agent-team/agent_team/transport/checker_intake.py b/agent-team/agent_team/transport/checker_intake.py new file mode 100644 index 0000000..b1aad76 --- /dev/null +++ b/agent-team/agent_team/transport/checker_intake.py @@ -0,0 +1,404 @@ +"""Plane-1 checker-finding INTAKE: a confirmed finding becomes a pipeline task. + +This is the Plane-2 **P5 cross-plane loop**: it closes the gap between the +Plane-1 read-only *checkers* +(``security-review/checkers/compliance-drift.sh`` and +``security-review/checkers/dependency-cve.sh``) and the Plane-2 human-gated +SDLC *pipeline*. Where :mod:`agent_team.transport.github_intake` turns a +labeled GitHub issue into one pipeline task, this leaf turns a confirmed, +at-or-above-threshold checker *finding* into one pipeline remediation task by +calling the same committed coordinator intake entry, +:meth:`agent_team.coordinator.Coordinator.start_task` +(``task_text=``, ``transport_name=``). + +It deliberately mirrors the ``github_intake`` seam so the two front doors stay +consistent and equally testable: + +* Input is **plain data**, not network. The poller reads checker *report* JSON + (the exact shape the bash checkers already emit — a top-level object with a + ``findings`` array) from one or more files / a directory. No SDK, no socket, + no token: a checker run already wrote the report; this only reads it. +* ``coordinator`` is anything exposing ``start_task(task_text=..., + transport_name=...)``: the live :class:`~agent_team.coordinator.Coordinator` + in production, a stub in tests. No model or transport is touched here. + +Selection contract: + A finding is ingested only when it is both ``status == "confirmed"`` AND its + ``severity`` is at or above the configured threshold (default ``high``). + Unconfirmed / suppressed findings and below-threshold severities are + skipped. An ``unverified`` severity (the schema's auto-downgrade marker) is + treated as below every real threshold and never ingested. + +De-duplication (P5 scope note — matches the github_intake discipline): + The poller tracks already-ingested findings in an **in-memory** set keyed by + a stable content identity (the finding ``id`` when present, else a + content-hash of checker+repo+title+severity). So re-reading the same nightly + report — or two reports that both carry the same finding — does not start a + second task within one process. This is deliberately simple and mirrors + ``github_intake``: it does NOT survive a process restart. Durable de-dup (a + ledger table of ingested finding ids, mirroring the ``pending_questions`` + discipline) is the known FOLLOW-UP and is intentionally not shipped here. + After a restart an already-ingested finding still present in a fresh report + would be re-ingested; treat the in-memory set as a best-effort guard, not a + durable contract. + +Untrusted-input hygiene: + Checker findings carry **repo-controlled strings** (titles, proofs) — a + repo name, a dependency advisory summary, a PR title. Those must never reach + operator logs or the rendered task text raw, or a forged multi-line value + could spoof log lines / pipeline-task framing (log injection). Every such + string is sanitized via :func:`_sanitize` (newline/control-char neutralised, + length-bounded) before it is logged or rendered, mirroring the coordinator's + ``start_task`` log-injection hardening. + +Design constraints (pre-deployment scaffolding): + * **No live infrastructure.** Nothing is provisioned or called at import. + The poller only reads local JSON and calls the injected coordinator. + * **P2 stays the production default; this is OPT-IN and INERT.** It is wired + ONLY behind the ``intake-checker`` run-team subcommand, never into the + always-on ``serve`` path. It does no CI, OIDC, git/patch apply, or + network; it only reads a report file and calls the existing intake entry. +""" + +from __future__ import annotations + +import hashlib +import json +import logging +from pathlib import Path +from typing import Any + +__all__ = [ + "CHECKER_TRANSPORT_NAME", + "DEFAULT_SEVERITY_THRESHOLD", + "SEVERITY_RANK", + "CheckerFindingIntake", + "finding_identity", + "finding_task_text", + "load_report_findings", + "select_findings", +] + +_LOG = logging.getLogger(__name__) + +# Transport name handed to the coordinator's intake entry so the resulting +# remediation task's clarifier question-sets route over a real channel. Findings +# are about org repos, so GitHub is the natural default (mirrors github_intake). +CHECKER_TRANSPORT_NAME = "github" + +# Severity ordering, matching the finding.schema.json enum. Higher rank = more +# severe. ``info`` / ``unverified`` sit below every real remediation threshold. +SEVERITY_RANK: dict[str, int] = { + "unverified": -1, + "info": 0, + "low": 1, + "medium": 2, + "high": 3, + "critical": 4, +} + +# Default minimum severity a confirmed finding must reach to spawn a task. +DEFAULT_SEVERITY_THRESHOLD = "high" + +# Hard cap on any single rendered/logged untrusted string. Generous enough for a +# real title/proof, tight enough that a forged megastring cannot flood the log or +# the task text. +_MAX_FIELD_LEN = 500 + + +def _sanitize(value: Any, *, max_len: int = _MAX_FIELD_LEN) -> str: + """Neutralise an untrusted, repo-controlled string for logs / task text. + + Coerces ``value`` to ``str``, strips surrounding whitespace, replaces every + control character (newlines, carriage returns, tabs, and other C0/C1 + controls) with a visible escape so a forged multi-line value cannot spoof a + log line or the framing of the rendered task text, and bounds the length so a + megastring cannot flood the sink. Mirrors the coordinator's ``start_task`` + log-injection hardening, generalised to every untrusted field. + """ + text = str(value if value is not None else "").strip() + if len(text) > max_len: + text = text[:max_len] + "…(truncated)" + out: list[str] = [] + for ch in text: + if ch == "\n": + out.append("\\n") + elif ch == "\r": + out.append("\\r") + elif ch == "\t": + out.append("\\t") + elif ord(ch) < 0x20 or ord(ch) == 0x7F: + # Any other C0 control (and DEL) -> visible escape. + out.append(f"\\x{ord(ch):02x}") + else: + out.append(ch) + return "".join(out) + + +def _proof_hint(proof: Any) -> str: + """Render the remediation hint from a finding's ``proof`` object, sanitized. + + Both checkers nest the actionable detail under ``proof``: + + * compliance-drift: ``{"outcome": ""}`` + * dependency-cve: ``{"package", "version", "advisory_id", "summary", + "fixed_version", ...}`` + + so this renders whichever keys are present into one compact, sanitized line. + A non-mapping or empty ``proof`` yields an empty hint (the caller omits the + line). Every value is run through :func:`_sanitize` because proofs are + repo-controlled (e.g. an advisory summary copied from an upstream feed). + """ + if not isinstance(proof, dict): + return "" + # Order keys for a stable, readable hint; unknown keys are appended after. + preferred = ( + "outcome", + "summary", + "package", + "version", + "fixed_version", + "advisory_id", + ) + parts: list[str] = [] + seen: set[str] = set() + for key in preferred: + if key in proof and proof[key] not in (None, ""): + parts.append(f"{key}={_sanitize(proof[key])}") + seen.add(key) + for key, val in proof.items(): + if key in seen or val in (None, ""): + continue + parts.append(f"{_sanitize(key, max_len=80)}={_sanitize(val)}") + return "; ".join(parts) + + +def finding_identity(finding: dict[str, Any]) -> str: + """Return the stable de-dup identity for ``finding`` as a string. + + Prefers the finding ``id`` (the checkers mint a stable ``-``), + which keeps the same finding from spawning two tasks across nightly runs. + When ``id`` is absent (a malformed/partial report), falls back to a content + hash of checker+repo+title+severity so two structurally identical findings + still collapse to one task rather than slipping the de-dup. + """ + raw_id = finding.get("id") + if raw_id not in (None, ""): + return str(raw_id) + payload = "\x1f".join( + str(finding.get(key) or "") for key in ("check", "repo", "title", "severity") + ) + return "sha256:" + hashlib.sha256(payload.encode("utf-8")).hexdigest() + + +def finding_task_text(finding: dict[str, Any], *, checker: str = "") -> str: + """Render one confirmed finding into the pipeline task's ``task_text``. + + Produces a compact, fully-sanitized remediation brief: a headline line with + the checker, repo, and severity; the finding title; and a remediation hint + drawn from ``proof``. Every interpolated value is repo-controlled and so is + passed through :func:`_sanitize` first (no raw newline / control char reaches + the task text or, downstream, the operator log). ``checker`` falls back to + the finding's own ``check`` field when not supplied by the report header. + """ + repo = _sanitize(finding.get("repo") or "", max_len=120) + severity = _sanitize(finding.get("severity") or "", max_len=40) + title = _sanitize(finding.get("title") or "") + checker_name = _sanitize(checker or finding.get("check") or "checker", max_len=80) + + lines = [ + f"[Plane-1 {checker_name}] remediation for {repo} (severity={severity})", + f"Finding: {title}", + ] + hint = _proof_hint(finding.get("proof")) + if hint: + lines.append(f"Remediation hint: {hint}") + return "\n".join(lines) + + +def _meets_threshold(severity: Any, *, threshold_rank: int) -> bool: + """True when ``severity`` is a known level at or above ``threshold_rank``. + + Unknown / missing severities (and the schema's ``unverified`` downgrade + marker, ranked below zero) never meet a real threshold, so a malformed + finding can never sneak past the gate. + """ + rank = SEVERITY_RANK.get(str(severity).strip().lower(), -99) + return rank >= threshold_rank + + +def select_findings( + findings: list[dict[str, Any]], + *, + threshold: str = DEFAULT_SEVERITY_THRESHOLD, +) -> list[dict[str, Any]]: + """Filter ``findings`` to those eligible to spawn a remediation task. + + A finding is selected only when BOTH: + + * ``status == "confirmed"`` (unverified / suppressed are dropped), and + * its ``severity`` is at or above ``threshold`` (default ``high``). + + Order is preserved. ``threshold`` must be one of the schema severities; + an unknown threshold is rejected so a typo cannot silently widen the gate. + """ + key = threshold.strip().lower() + if key not in SEVERITY_RANK or key in ("unverified", "info"): + raise ValueError( + f"invalid severity threshold {threshold!r}; expected one of " + "low/medium/high/critical" + ) + threshold_rank = SEVERITY_RANK[key] + selected: list[dict[str, Any]] = [] + for finding in findings: + if str(finding.get("status")).strip().lower() != "confirmed": + continue + if not _meets_threshold(finding.get("severity"), threshold_rank=threshold_rank): + continue + selected.append(finding) + return selected + + +def load_report_findings(path: Path) -> list[dict[str, Any]]: + """Read a checker report file (or every ``*.json`` in a dir) into findings. + + Accepts the exact shape the bash checkers emit: a top-level object with a + ``findings`` array. A bare JSON array is also accepted (a caller that has + already extracted ``.findings``). When ``path`` is a directory, every + ``*.json`` file directly inside it is read and the findings concatenated + (a malformed file raises, surfacing the bad report rather than silently + skipping it). Non-mapping finding entries are ignored defensively. + """ + if path.is_dir(): + findings: list[dict[str, Any]] = [] + for report in sorted(path.glob("*.json")): + findings.extend(load_report_findings(report)) + return findings + + raw = path.read_text(encoding="utf-8") + data = json.loads(raw) if raw.strip() else {} + if isinstance(data, list): + items = data + elif isinstance(data, dict): + items = data.get("findings") or [] + else: + items = [] + return [item for item in items if isinstance(item, dict)] + + +class CheckerFindingIntake: + """Turn confirmed at/above-threshold checker findings into pipeline tasks. + + Construct with an injected ``coordinator`` (anything exposing + ``start_task(task_text=..., transport_name=...)``), an optional severity + ``threshold`` (default ``high``), and the ``transport_name`` the resulting + remediation tasks should deliver clarifier questions over. Then call + :meth:`ingest_findings` (in-memory findings) or :meth:`ingest_reports` + (report files / a directory). + + De-dup is in-memory only (see the module docstring): the set of ingested + finding identities lives on the instance, so re-reading the same report never + double-ingests within one process, but a restart loses the set. Durable + de-dup is the known follow-up. + + Nothing here touches the network or any SDK: it reads local JSON and calls + the injected coordinator, so the whole intake is unit-testable with a stub + coordinator and in-memory findings. + """ + + def __init__( + self, + *, + coordinator: Any, + threshold: str = DEFAULT_SEVERITY_THRESHOLD, + transport_name: str = CHECKER_TRANSPORT_NAME, + ) -> None: + """Bind the intake to one coordinator, severity threshold, and transport. + + Args: + coordinator: The intake target. Must expose + ``start_task(task_text=..., transport_name=...)``: the live + :class:`~agent_team.coordinator.Coordinator` in production. + threshold: Minimum severity a confirmed finding must reach to spawn a + task (``low``/``medium``/``high``/``critical``; default ``high``). + An invalid value is rejected up-front. + transport_name: Channel the resulting task's clarifier question-sets + route over (default ``github``). + """ + # Validate the threshold eagerly (reuses select_findings' guard). + select_findings([], threshold=threshold) + self._coordinator = coordinator + self._threshold = threshold.strip().lower() + self._transport_name = transport_name + # In-memory de-dup set (P5 scope: best-effort, NOT durable across a + # restart; see the module docstring). Tracks finding identities already + # turned into tasks so re-reading a report does not double-ingest. + self._ingested: set[str] = set() + + @property + def threshold(self) -> str: + """The configured severity threshold (read-only).""" + return self._threshold + + @property + def ingested_ids(self) -> frozenset[str]: + """A snapshot of finding identities ingested this process (read-only).""" + return frozenset(self._ingested) + + def ingest_findings(self, findings: list[dict[str, Any]]) -> list[str]: + """Select, de-dup, and start one task per unique eligible finding. + + For each finding that passes :func:`select_findings` and is not already + ingested this process, calls ``coordinator.start_task(task_text=, transport_name=)`` and records its identity so a + subsequent pass does not re-ingest it. The identity is recorded ONLY + after ``start_task`` returns, so a failing intake leaves the finding + eligible for retry rather than silently dropping it (mirrors + github_intake). + + Returns the list of finding identities ingested on THIS pass (empty when + nothing new), so an operator loop can meter intake volume. + """ + ingested_now: list[str] = [] + for finding in select_findings(findings, threshold=self._threshold): + identity = finding_identity(finding) + if identity in self._ingested: + _LOG.debug( + "checker-intake: finding %s already ingested; skip", + _sanitize(identity, max_len=120), + ) + continue + + checker = str(finding.get("check") or "") + task_text = finding_task_text(finding, checker=checker) + # task_text is already sanitized field-by-field; log a sanitized + # summary (never the raw repo/title) to avoid log injection. + _LOG.info( + "checker-intake: starting remediation task for %s " + "(repo=%s severity=%s checker=%s)", + _sanitize(identity, max_len=120), + _sanitize(finding.get("repo"), max_len=120), + _sanitize(finding.get("severity"), max_len=40), + _sanitize(checker, max_len=80), + ) + self._coordinator.start_task( + task_text=task_text, + transport_name=self._transport_name, + ) + self._ingested.add(identity) + ingested_now.append(identity) + + return ingested_now + + def ingest_reports(self, paths: list[Path]) -> list[str]: + """Load checker report files / dirs and ingest their eligible findings. + + Each entry in ``paths`` may be a report file or a directory of ``*.json`` + reports (see :func:`load_report_findings`). All loaded findings are + concatenated, then handed to :meth:`ingest_findings` (so de-dup spans the + whole batch). Returns the finding identities ingested on this call. + """ + findings: list[dict[str, Any]] = [] + for path in paths: + findings.extend(load_report_findings(path)) + return self.ingest_findings(findings) diff --git a/agent-team/tests/test_checker_intake.py b/agent-team/tests/test_checker_intake.py new file mode 100644 index 0000000..fa0e2b6 --- /dev/null +++ b/agent-team/tests/test_checker_intake.py @@ -0,0 +1,434 @@ +"""Unit tests for agent_team.transport.checker_intake (P5 cross-plane loop). + +Fully hermetic: the coordinator is an injected in-memory fake and findings are +plain dicts (or temp report files), so no network call, token, GitHub SDK, model +or live checker is exercised. The tests pin the P5 contract: + +* a CONFIRMED finding at/above the threshold creates exactly one task, +* unconfirmed / suppressed / below-threshold findings are ignored, +* re-reading the same finding does not double-ingest (de-dup), +* the rendered task text is safe (no log/text injection from repo-controlled + fields), and +* the loop is OPT-IN — it is NOT wired into the default coordinator serve path. +""" + +from __future__ import annotations + +import importlib.util +import json +from pathlib import Path +from typing import Any + +import pytest + +from agent_team.transport.checker_intake import ( + CHECKER_TRANSPORT_NAME, + CheckerFindingIntake, + finding_identity, + finding_task_text, + load_report_findings, + select_findings, +) + +# --------------------------------------------------------------------------- # +# Fakes / builders +# --------------------------------------------------------------------------- # + + +class FakeCoordinator: + """In-memory coordinator double recording every ``start_task`` call. + + Mirrors the real coordinator's intake entry signature and captures each + call's keyword args so a test can assert exactly what was ingested. + """ + + def __init__(self) -> None: + self.calls: list[dict[str, Any]] = [] + + def start_task(self, *, task_text: str, transport_name: str) -> str: + self.calls.append({"task_text": task_text, "transport_name": transport_name}) + return f"thread-{len(self.calls)}" + + +def _finding( + *, + fid: str | None = "repo-a-branchprot-no-pr", + repo: str = "repo-a", + title: str = "main does not require a PR for merge", + severity: str = "high", + status: str = "confirmed", + check: str = "branch-protection", + proof: dict[str, Any] | None = None, +) -> dict[str, Any]: + """Build a checker-finding-shaped mapping matching the bash checker emit.""" + finding: dict[str, Any] = { + "repo": repo, + "title": title, + "severity": severity, + "category": "other", + "check": check, + "status": status, + "proof": proof if proof is not None else {"outcome": "require a PR for merges"}, + } + if fid is not None: + finding["id"] = fid + return finding + + +def _report( + findings: list[dict[str, Any]], *, checker: str = "compliance-drift" +) -> dict[str, Any]: + """Wrap findings in the top-level report object the checkers emit.""" + return { + "checker": checker, + "generated": "2026-06-18T00:00:00Z", + "org": "Sea-Haven-Industries", + "findings": findings, + } + + +# --------------------------------------------------------------------------- # +# select_findings — the gate +# --------------------------------------------------------------------------- # + + +def test_selects_confirmed_at_or_above_threshold() -> None: + findings = [ + _finding(fid="r-high", severity="high"), + _finding(fid="r-crit", severity="critical"), + ] + selected = select_findings(findings, threshold="high") + assert [f["id"] for f in selected] == ["r-high", "r-crit"] + + +def test_ignores_below_threshold() -> None: + findings = [ + _finding(fid="r-low", severity="low"), + _finding(fid="r-med", severity="medium"), + _finding(fid="r-high", severity="high"), + ] + selected = select_findings(findings, threshold="high") + assert [f["id"] for f in selected] == ["r-high"] + + +def test_ignores_unconfirmed_and_suppressed() -> None: + findings = [ + _finding(fid="r-unver", status="unverified", severity="critical"), + _finding(fid="r-supp", status="suppressed", severity="critical"), + _finding(fid="r-ok", status="confirmed", severity="high"), + ] + selected = select_findings(findings, threshold="high") + assert [f["id"] for f in selected] == ["r-ok"] + + +def test_unverified_severity_never_meets_threshold() -> None: + # status confirmed but severity is the schema's downgrade marker. + findings = [_finding(fid="r-x", status="confirmed", severity="unverified")] + assert select_findings(findings, threshold="low") == [] + + +def test_threshold_lower_admits_more() -> None: + findings = [ + _finding(fid="r-low", severity="low"), + _finding(fid="r-med", severity="medium"), + ] + selected = select_findings(findings, threshold="low") + assert [f["id"] for f in selected] == ["r-low", "r-med"] + + +def test_invalid_threshold_rejected() -> None: + for bad in ("info", "unverified", "nope", ""): + with pytest.raises(ValueError): + select_findings([], threshold=bad) + + +# --------------------------------------------------------------------------- # +# finding_task_text / sanitization — untrusted-input hygiene +# --------------------------------------------------------------------------- # + + +def test_task_text_renders_checker_repo_severity_title_and_hint() -> None: + text = finding_task_text( + _finding( + repo="orchestrator", + title="vulnerable dependency", + severity="critical", + check="vulnerable-dependency", + proof={ + "package": "requests", + "version": "2.0.0", + "fixed_version": "2.32.0", + "advisory_id": "GHSA-xxxx", + "summary": "RCE in requests", + }, + ), + checker="dependency-cve", + ) + assert "[Plane-1 dependency-cve]" in text + assert "orchestrator" in text + assert "severity=critical" in text + assert "Finding: vulnerable dependency" in text + assert "package=requests" in text + assert "fixed_version=2.32.0" in text + + +def test_task_text_neutralises_newline_injection_in_title() -> None: + # A forged title that tries to inject a fake task line / log line. + evil = "real title\nFinding: SPOOFED\nadmin=true" + text = finding_task_text(_finding(title=evil)) + # No raw newline from the untrusted field reaches the task text body: the + # title occupies exactly one line, so the only real newlines are the ones + # WE add between the headline / finding / hint lines. + lines = text.split("\n") + finding_lines = [ln for ln in lines if ln.startswith("Finding: ")] + assert finding_lines == ["Finding: real title\\nFinding: SPOOFED\\nadmin=true"] + + +def test_task_text_neutralises_carriage_return_and_controls() -> None: + evil = "x\r\ty\x00z\x1b[31m" + text = finding_task_text(_finding(title=evil)) + assert "\r" not in text + assert "\x00" not in text + assert "\x1b" not in text + assert "\\r" in text and "\\t" in text and "\\x00" in text + + +def test_task_text_bounds_megastring() -> None: + text = finding_task_text(_finding(title="A" * 10_000)) + assert "…(truncated)" in text + assert len(text) < 1_000 + + +def test_proof_hint_sanitises_repo_controlled_summary() -> None: + text = finding_task_text( + _finding(proof={"outcome": "line1\nline2\rline3"}), + ) + assert "outcome=line1\\nline2\\rline3" in text + assert "\n line2" not in text + + +# --------------------------------------------------------------------------- # +# finding_identity / de-dup +# --------------------------------------------------------------------------- # + + +def test_identity_prefers_id() -> None: + assert finding_identity(_finding(fid="repo-a-x")) == "repo-a-x" + + +def test_identity_falls_back_to_content_hash_without_id() -> None: + f = _finding(fid=None) + ident = finding_identity(f) + assert ident.startswith("sha256:") + # Stable for identical content. + assert finding_identity(_finding(fid=None)) == ident + + +# --------------------------------------------------------------------------- # +# CheckerFindingIntake.ingest_findings — the core contract +# --------------------------------------------------------------------------- # + + +def test_confirmed_finding_creates_exactly_one_task() -> None: + coordinator = FakeCoordinator() + intake = CheckerFindingIntake(coordinator=coordinator) + ingested = intake.ingest_findings([_finding(fid="repo-a-x")]) + + assert ingested == ["repo-a-x"] + assert len(coordinator.calls) == 1 + call = coordinator.calls[0] + assert call["transport_name"] == CHECKER_TRANSPORT_NAME + assert "repo-a" in call["task_text"] + + +def test_below_threshold_and_unconfirmed_ignored_by_intake() -> None: + coordinator = FakeCoordinator() + intake = CheckerFindingIntake(coordinator=coordinator) + ingested = intake.ingest_findings( + [ + _finding(fid="r-low", severity="low"), + _finding(fid="r-unver", status="unverified", severity="critical"), + _finding(fid="r-ok", severity="high"), + ] + ) + assert ingested == ["r-ok"] + assert len(coordinator.calls) == 1 + + +def test_repeated_finding_does_not_double_ingest() -> None: + coordinator = FakeCoordinator() + intake = CheckerFindingIntake(coordinator=coordinator) + + first = intake.ingest_findings([_finding(fid="repo-a-x")]) + second = intake.ingest_findings([_finding(fid="repo-a-x")]) + + assert first == ["repo-a-x"] + assert second == [] # already ingested -> no second task + assert len(coordinator.calls) == 1 + assert intake.ingested_ids == frozenset({"repo-a-x"}) + + +def test_dedup_within_single_batch() -> None: + coordinator = FakeCoordinator() + intake = CheckerFindingIntake(coordinator=coordinator) + # Same finding present twice in one report batch. + ingested = intake.ingest_findings( + [_finding(fid="repo-a-x"), _finding(fid="repo-a-x")] + ) + assert ingested == ["repo-a-x"] + assert len(coordinator.calls) == 1 + + +def test_new_finding_on_second_pass_is_ingested() -> None: + coordinator = FakeCoordinator() + intake = CheckerFindingIntake(coordinator=coordinator) + + first = intake.ingest_findings([_finding(fid="repo-a-x")]) + second = intake.ingest_findings( + [_finding(fid="repo-a-x"), _finding(fid="repo-b-y", repo="repo-b")] + ) + assert first == ["repo-a-x"] + assert second == ["repo-b-y"] + assert len(coordinator.calls) == 2 + + +def test_custom_threshold_admits_medium() -> None: + coordinator = FakeCoordinator() + intake = CheckerFindingIntake(coordinator=coordinator, threshold="medium") + ingested = intake.ingest_findings([_finding(fid="r-med", severity="medium")]) + assert ingested == ["r-med"] + + +def test_failed_start_task_leaves_finding_eligible_for_retry() -> None: + class FlakyCoordinator: + def __init__(self) -> None: + self.attempts = 0 + + def start_task(self, *, task_text: str, transport_name: str) -> str: + self.attempts += 1 + if self.attempts == 1: + raise RuntimeError("transient intake failure") + return "thread-ok" + + coordinator = FlakyCoordinator() + intake = CheckerFindingIntake(coordinator=coordinator) + + with pytest.raises(RuntimeError): + intake.ingest_findings([_finding(fid="repo-a-x")]) + assert intake.ingested_ids == frozenset() # not recorded -> retryable + + ingested = intake.ingest_findings([_finding(fid="repo-a-x")]) + assert ingested == ["repo-a-x"] + assert coordinator.attempts == 2 + + +def test_invalid_threshold_rejected_at_construction() -> None: + with pytest.raises(ValueError): + CheckerFindingIntake(coordinator=FakeCoordinator(), threshold="bogus") + + +# --------------------------------------------------------------------------- # +# load_report_findings / ingest_reports — report-file path +# --------------------------------------------------------------------------- # + + +def test_load_report_findings_reads_top_level_object(tmp_path: Path) -> None: + report = tmp_path / "compliance-drift.json" + report.write_text(json.dumps(_report([_finding(fid="repo-a-x")])), encoding="utf-8") + findings = load_report_findings(report) + assert [f["id"] for f in findings] == ["repo-a-x"] + + +def test_load_report_findings_accepts_bare_array(tmp_path: Path) -> None: + report = tmp_path / "bare.json" + report.write_text(json.dumps([_finding(fid="repo-a-x")]), encoding="utf-8") + assert [f["id"] for f in load_report_findings(report)] == ["repo-a-x"] + + +def test_load_report_findings_reads_directory_of_reports(tmp_path: Path) -> None: + (tmp_path / "a.json").write_text( + json.dumps(_report([_finding(fid="repo-a-x")])), encoding="utf-8" + ) + (tmp_path / "b.json").write_text( + json.dumps(_report([_finding(fid="repo-b-y", repo="repo-b")])), + encoding="utf-8", + ) + ids = sorted(f["id"] for f in load_report_findings(tmp_path)) + assert ids == ["repo-a-x", "repo-b-y"] + + +def test_ingest_reports_selects_dedups_across_files(tmp_path: Path) -> None: + # Two reports both carrying the same high finding plus a unique one each. + (tmp_path / "a.json").write_text( + json.dumps( + _report( + [ + _finding(fid="shared", severity="high"), + _finding(fid="only-a", repo="repo-a"), + _finding(fid="low-a", severity="low"), + ] + ) + ), + encoding="utf-8", + ) + (tmp_path / "b.json").write_text( + json.dumps( + _report([_finding(fid="shared", severity="high"), _finding(fid="only-b")]) + ), + encoding="utf-8", + ) + coordinator = FakeCoordinator() + intake = CheckerFindingIntake(coordinator=coordinator) + ingested = intake.ingest_reports([tmp_path]) + + assert sorted(ingested) == ["only-a", "only-b", "shared"] + assert len(coordinator.calls) == 3 # low-a dropped, shared once + + +def test_empty_report_file_yields_nothing(tmp_path: Path) -> None: + report = tmp_path / "empty.json" + report.write_text("", encoding="utf-8") + assert load_report_findings(report) == [] + + +# --------------------------------------------------------------------------- # +# OPT-IN: not wired into the default serve path +# --------------------------------------------------------------------------- # + + +def _load_run_team(): + """Import run-team.py (hyphenated, so loaded by path) as a module.""" + cli_path = Path(__file__).resolve().parent.parent / "run-team.py" + 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 + + +def test_intake_checker_is_a_subcommand_not_the_default() -> None: + cli = _load_run_team() + parser = cli.build_parser() + # The subcommand exists... + args = parser.parse_args( + ["intake-checker", "--report", "/tmp/x.json", "--threshold", "high"] + ) + assert args.func is cli._cmd_intake_checker + assert args.command == "intake-checker" + + +def test_serve_path_does_not_invoke_checker_intake() -> None: + """The always-on serve loop must not reference checker intake (opt-in only).""" + cli = _load_run_team() + import inspect + + serve_src = inspect.getsource(cli._cmd_serve) + assert "checker" not in serve_src.lower() + assert "CheckerFindingIntake" not in serve_src + + # And coordinator.serve itself does not pull in the checker intake module. + from agent_team import coordinator as coordinator_mod + + coord_src = inspect.getsource(coordinator_mod.Coordinator.serve) + assert "checker_intake" not in coord_src + assert "CheckerFindingIntake" not in coord_src