diff --git a/tests/auto_agents/controller/test_metadata_hygiene.py b/tests/auto_agents/controller/test_metadata_hygiene.py new file mode 100644 index 000000000..e8a193921 --- /dev/null +++ b/tests/auto_agents/controller/test_metadata_hygiene.py @@ -0,0 +1,1559 @@ +"""Phase 4 — metadata-hygiene tests. + +Covers the five Phase 4 checks + their config: + +- ``completed_not_closed`` — close linked issue on MERGED +- ``closing_keyword_fixup`` — add ``Closes #N`` to PR body +- ``label_sync_from_issue`` — copy Priority/Type/MoSCoW labels +- ``state_label_inference`` — sync State/* labels to current_state +- ``milestone_assignment`` — copy milestone from linked issue +- ``MetadataHygieneConfig`` — kill switch + per-check flags + +Each check has: happy-path execute, idempotency skip, kill-switch +respect, dry-run no-write, and a check-specific edge case. +""" + +from __future__ import annotations + +from dataclasses import replace + +import pytest +from sqlalchemy import create_engine, text +from sqlalchemy.orm import sessionmaker + +from tools.controller.db.models import Workflow +from tools.controller.db.session import create_all +from tools.controller.master.metadata_hygiene import ( + MetadataHygieneCallbacks, + _STATE_TO_LABEL, + run_closing_keyword_fixup, + run_completed_not_closed, + run_label_sync_from_issue, + run_metadata_hygiene_tick, + run_milestone_assignment, + run_state_label_inference, +) +from tools.controller.master.metadata_hygiene_config import ( + MetadataHygieneConfig, + get_metadata_hygiene_config, +) + + +# ─── fixtures ───────────────────────────────────────────────────────── + + +@pytest.fixture +def engine(): + eng = create_engine( + "sqlite:///:memory:", + connect_args={"check_same_thread": False}, + ) + create_all(eng) + yield eng + eng.dispose() + + +@pytest.fixture +def session_factory(engine): + return sessionmaker(engine, expire_on_commit=False) + + +def _seed_workflow( + session_factory, + *, + state: str = "ANALYZING", + pr_number: int = 401, + kind: str = "pr", +) -> int: + with session_factory.begin() as s: + wf = Workflow( + kind=kind, + owner="drew", + repo="cleveragents-core", + entity_number=pr_number, + current_state=state, + ) + s.add(wf) + s.flush() + return wf.workflow_id + + +class _RecordingForgejo: + """Recording fake; lets tests assert what Forgejo was asked to do. + Default canned responses are HTTP 200 + a benign body.""" + + def __init__(self): + self.calls: list[tuple[str, tuple]] = [] + self.pr_details: dict[tuple[str, str, int], dict] = {} + # issue_details stores issue-endpoint responses used by check #5 + # (milestone_assignment). Keyed separately from pr_details so + # tests can distinguish "this is an issue" from "this is a PR". + self.issue_details: dict[tuple[str, str, int], dict] = {} + self.labels: dict[tuple[str, str, int], list[dict]] = {} + self.add_label_ok = True + self.remove_label_ok = True + self.patch_status = 200 + + def get_pr_details(self, owner, repo, n): + self.calls.append(("get_pr_details", (owner, repo, n))) + return self.pr_details.get((owner, repo, n)) + + def get_issue_state(self, owner, repo, n): + self.calls.append(("get_issue_state", (owner, repo, n))) + return self.issue_details.get((owner, repo, n)) + + def get_labels(self, owner, repo, n): + self.calls.append(("get_labels", (owner, repo, n))) + return list(self.labels.get((owner, repo, n), [])) + + def add_label(self, owner, repo, n, lbl): + self.calls.append(("add_label", (owner, repo, n, lbl))) + if self.add_label_ok: + # Mutate state so subsequent get_labels returns it. + self.labels.setdefault((owner, repo, n), []).append({"name": lbl}) + return self.add_label_ok + + def remove_label(self, owner, repo, n, lbl): + self.calls.append(("remove_label", (owner, repo, n, lbl))) + existing = self.labels.get((owner, repo, n), []) + if self.remove_label_ok: + self.labels[(owner, repo, n)] = [ + x for x in existing if x.get("name") != lbl + ] + return self.remove_label_ok + + def patch_pr_state(self, owner, repo, n, state): + self.calls.append(("patch_pr_state", (owner, repo, n, state))) + return {"status": self.patch_status, "body": {"state": state}} + + def patch_pr_body(self, owner, repo, n, body): + self.calls.append(("patch_pr_body", (owner, repo, n, body))) + return {"status": self.patch_status, "body": {"body": body}} + + def patch_pr_milestone(self, owner, repo, n, mid): + self.calls.append(("patch_pr_milestone", (owner, repo, n, mid))) + return {"status": self.patch_status, "body": {"milestone": mid}} + + +def _callbacks( + fake: _RecordingForgejo, + *, + completed_not_closed=True, + closing_keyword_fixup=True, + label_sync=True, + state_label_inference=True, + milestone_assignment=True, + dry_run=False, +) -> MetadataHygieneCallbacks: + return MetadataHygieneCallbacks( + get_pr_details=fake.get_pr_details, + get_issue_state=fake.get_issue_state, + get_labels=fake.get_labels, + add_label=fake.add_label, + remove_label=fake.remove_label, + patch_pr_state=fake.patch_pr_state, + patch_pr_body=fake.patch_pr_body, + patch_pr_milestone=fake.patch_pr_milestone, + completed_not_closed_enabled=completed_not_closed, + closing_keyword_fixup_enabled=closing_keyword_fixup, + label_sync_enabled=label_sync, + state_label_inference_enabled=state_label_inference, + milestone_assignment_enabled=milestone_assignment, + dry_run=dry_run, + ) + + +# ─── config ─────────────────────────────────────────────────────────── + + +class TestMetadataHygieneConfig: + def test_defaults_all_disabled(self, monkeypatch): + for v in ( + "CONTROLLER_METADATA_HYGIENE_ENABLED", + "CONTROLLER_METADATA_HYGIENE_COMPLETED_NOT_CLOSED", + "CONTROLLER_METADATA_HYGIENE_CLOSING_KEYWORD_FIXUP", + "CONTROLLER_METADATA_HYGIENE_LABEL_SYNC", + "CONTROLLER_METADATA_HYGIENE_STATE_LABEL_INFERENCE", + "CONTROLLER_METADATA_HYGIENE_MILESTONE_ASSIGNMENT", + ): + monkeypatch.delenv(v, raising=False) + cfg = get_metadata_hygiene_config() + assert cfg == MetadataHygieneConfig() # all-False default + assert cfg.enabled is False + + def test_per_check_flags_independent(self, monkeypatch): + monkeypatch.setenv("CONTROLLER_METADATA_HYGIENE_ENABLED", "true") + monkeypatch.setenv("CONTROLLER_METADATA_HYGIENE_LABEL_SYNC", "true") + monkeypatch.delenv( + "CONTROLLER_METADATA_HYGIENE_CLOSING_KEYWORD_FIXUP", raising=False + ) + cfg = get_metadata_hygiene_config() + assert cfg.enabled is True + assert cfg.label_sync_from_issue is True + # Other flags untouched stay False. + assert cfg.completed_not_closed is False + assert cfg.closing_keyword_fixup is False + + +# ─── check #1: completed_not_closed ────────────────────────────────── + + +class TestCompletedNotClosed: + def test_closes_linked_issue_on_merged(self, engine, session_factory): + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + wf_id = _seed_workflow(session_factory, state="MERGED") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Implements thing.\n\nCloses #42", + } + report = MetadataHygieneReport() + # Function returns None — mutates the passed-in report. + run_completed_not_closed( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + assert report.completed_not_closed_executed == 1 + # The PATCH for issue #42: + assert ("patch_pr_state", ("drew", "cleveragents-core", 42, "closed")) in fake.calls + # Audit row landed. + with engine.connect() as c: + row = c.execute( + text( + "SELECT executed, reason_category, stage, verdict " + "FROM grooming_decisions " + "WHERE workflow_id = :w AND check_name = 'completed_not_closed'" + ), + {"w": wf_id}, + ).first() + assert row.executed == 1 + assert row.reason_category == "closed_1_of_1" + # Phase-4 audit shape — distinguishes from Phase-1/2/3 rows. + assert row.stage == "metadata_hygiene" + assert row.verdict == "metadata_fixup" + + def test_skips_when_no_linked_issues(self, engine, session_factory): + wf_id = _seed_workflow(session_factory, state="MERGED") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Cleans up tests.", # no Closes #N + } + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_completed_not_closed( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + assert report.completed_not_closed_skipped == 1 + # No PATCH on any issue. + assert not any(c[0] == "patch_pr_state" for c in fake.calls) + # Audit row recorded the skip (so we don't re-check). + with engine.connect() as c: + row = c.execute( + text( + "SELECT reason_category FROM grooming_decisions " + "WHERE workflow_id = :w AND check_name = 'completed_not_closed'" + ), + {"w": wf_id}, + ).first() + assert row.reason_category == "no_linked_issues" + + def test_idempotent_skip(self, engine, session_factory): + """After a successful close, the LEFT-JOIN candidate query + excludes the resolved workflow entirely — the second sweep + finds zero candidates and does nothing. Round-1 architect- + review refactor: pre-refactor the per-row ``_already_done`` + SELECT incremented ``skipped`` counters; the new candidate + SELECT skips at the SQL level so no per-row processing + happens AND no skip counter increments.""" + wf_id = _seed_workflow(session_factory, state="MERGED") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Closes #42", + } + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + r1 = MetadataHygieneReport() + run_completed_not_closed(engine=engine, callbacks=_callbacks(fake), report=r1) + assert r1.completed_not_closed_executed == 1 + # Second run is a no-op — the workflow is excluded from + # candidates by the LEFT JOIN on executed=1. + fake2 = _RecordingForgejo() + fake2.pr_details = fake.pr_details + r2 = MetadataHygieneReport() + run_completed_not_closed(engine=engine, callbacks=_callbacks(fake2), report=r2) + assert r2.completed_not_closed_executed == 0 + assert fake2.calls == [] # No Forgejo activity at all. + + def test_dry_run_skips_forgejo_writes(self, engine, session_factory): + _seed_workflow(session_factory, state="MERGED") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = {"body": "Closes #42"} + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_completed_not_closed( + engine=engine, callbacks=_callbacks(fake, dry_run=True), report=report, + ) + assert not any(c[0] == "patch_pr_state" for c in fake.calls) + + def test_kill_switch_short_circuit(self, engine, session_factory): + _seed_workflow(session_factory, state="MERGED") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = {"body": "Closes #42"} + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_completed_not_closed( + engine=engine, + callbacks=_callbacks(fake, completed_not_closed=False), + report=report, + ) + assert fake.calls == [] + assert report.completed_not_closed_executed == 0 + + def test_only_runs_on_merged_workflows(self, engine, session_factory): + # An ANALYZING workflow with Closes #42 must NOT be acted on. + _seed_workflow(session_factory, state="ANALYZING", pr_number=402) + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 402)] = {"body": "Closes #42"} + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_completed_not_closed(engine=engine, callbacks=_callbacks(fake), report=report) + assert fake.calls == [] + + +# ─── check #2: closing_keyword_fixup ────────────────────────────────── + + +class TestClosingKeywordFixup: + def test_adds_closes_when_bare_ref_exists(self, engine, session_factory): + wf_id = _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Fixes the bug in #42 — refactor.", + } + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_closing_keyword_fixup(engine=engine, callbacks=_callbacks(fake), report=report) + # Should PATCH the body adding 'Closes #42'. + patch_calls = [c for c in fake.calls if c[0] == "patch_pr_body"] + assert len(patch_calls) == 1 + assert "Closes #42" in patch_calls[0][1][3] + assert report.closing_keyword_fixup_executed == 1 + + def test_skip_when_closes_already_present(self, engine, session_factory): + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Closes #42 — already covered.", + } + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_closing_keyword_fixup(engine=engine, callbacks=_callbacks(fake), report=report) + assert not any(c[0] == "patch_pr_body" for c in fake.calls) + assert report.closing_keyword_fixup_skipped == 1 + + def test_skip_when_multiple_bare_refs(self, engine, session_factory): + """Multiple bare refs is ambiguous — auto-fix would guess + wrong. Test the safety perimeter.""" + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Related to #42 and also #43.", + } + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_closing_keyword_fixup(engine=engine, callbacks=_callbacks(fake), report=report) + assert not any(c[0] == "patch_pr_body" for c in fake.calls) + + def test_skip_when_no_bare_ref(self, engine, session_factory): + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Just a refactor — no issue ref.", + } + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_closing_keyword_fixup(engine=engine, callbacks=_callbacks(fake), report=report) + assert not any(c[0] == "patch_pr_body" for c in fake.calls) + + def test_skip_on_terminal_workflow(self, engine, session_factory): + # Should not act on MERGED / ABANDONED workflows. + _seed_workflow(session_factory, state="MERGED") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = {"body": "Fixes #42"} + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_closing_keyword_fixup(engine=engine, callbacks=_callbacks(fake), report=report) + assert fake.calls == [] + + +# ─── check #3: label_sync_from_issue ────────────────────────────────── + + +class TestLabelSyncFromIssue: + def test_copies_matching_labels(self, engine, session_factory): + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = {"body": "Closes #42"} + # Issue #42 has Priority/High + Type/Bug + Random/Other. + fake.labels[("drew", "cleveragents-core", 42)] = [ + {"name": "Priority/High"}, + {"name": "Type/Bug"}, + {"name": "Random/Other"}, + ] + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_label_sync_from_issue(engine=engine, callbacks=_callbacks(fake), report=report) + # add_label called for Priority/High + Type/Bug, NOT Random/Other. + added = {c[1][3] for c in fake.calls if c[0] == "add_label"} + assert "Priority/High" in added + assert "Type/Bug" in added + assert "Random/Other" not in added + assert report.label_sync_executed == 1 + + def test_no_op_when_no_matching_labels(self, engine, session_factory): + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = {"body": "Closes #42"} + fake.labels[("drew", "cleveragents-core", 42)] = [ + {"name": "Random/Other"}, + ] + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_label_sync_from_issue(engine=engine, callbacks=_callbacks(fake), report=report) + assert not any(c[0] == "add_label" for c in fake.calls) + + def test_no_op_when_pr_has_no_linked_issue(self, engine, session_factory): + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = {"body": "no issue"} + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_label_sync_from_issue(engine=engine, callbacks=_callbacks(fake), report=report) + assert not any(c[0] == "add_label" for c in fake.calls) + + def test_idempotent_skip(self, engine, session_factory): + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = {"body": "Closes #42"} + fake.labels[("drew", "cleveragents-core", 42)] = [{"name": "Priority/High"}] + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + r1 = MetadataHygieneReport() + run_label_sync_from_issue( + engine=engine, callbacks=_callbacks(fake), report=r1, + ) + assert r1.label_sync_executed == 1 + # Second run is a no-op — workflow excluded by LEFT JOIN on + # the executed=1 audit row. + fake2 = _RecordingForgejo() + fake2.pr_details = fake.pr_details + fake2.labels = fake.labels + r2 = MetadataHygieneReport() + run_label_sync_from_issue(engine=engine, callbacks=_callbacks(fake2), report=r2) + assert r2.label_sync_executed == 0 + assert fake2.calls == [] + + +# ─── check #4: state_label_inference ────────────────────────────────── + + +class TestStateLabelInference: + def test_adds_state_label_matching_current_state(self, engine, session_factory): + wf_id = _seed_workflow(session_factory, state="REVIEWING") + fake = _RecordingForgejo() + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_state_label_inference(engine=engine, callbacks=_callbacks(fake), report=report) + added = [c[1][3] for c in fake.calls if c[0] == "add_label"] + assert "State/Reviewing" in added + + def test_removes_stale_state_labels(self, engine, session_factory): + _seed_workflow(session_factory, state="REVIEWING") + fake = _RecordingForgejo() + # Existing label State/Implementing should get removed. + fake.labels[("drew", "cleveragents-core", 401)] = [ + {"name": "State/Implementing"}, + ] + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_state_label_inference(engine=engine, callbacks=_callbacks(fake), report=report) + removed = [c[1][3] for c in fake.calls if c[0] == "remove_label"] + assert "State/Implementing" in removed + + def test_no_op_when_state_already_synced(self, engine, session_factory): + wf_id = _seed_workflow(session_factory, state="REVIEWING") + fake = _RecordingForgejo() + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + # First run syncs. + run_state_label_inference( + engine=engine, callbacks=_callbacks(fake), report=MetadataHygieneReport(), + ) + # Second run on the same state must skip. + fake2 = _RecordingForgejo() + fake2.labels = fake.labels + r2 = MetadataHygieneReport() + run_state_label_inference(engine=engine, callbacks=_callbacks(fake2), report=r2) + assert r2.state_label_skipped == 1 + assert not any(c[0] == "add_label" for c in fake2.calls) + + def test_re_syncs_on_state_change(self, engine, session_factory): + wf_id = _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + # First sync = Analyzing. + run_state_label_inference( + engine=engine, callbacks=_callbacks(fake), report=MetadataHygieneReport(), + ) + # Now transition to REVIEWING. + with session_factory.begin() as s: + wf = s.get(Workflow, wf_id) + wf.current_state = "REVIEWING" + # Second run picks up the change. + fake2 = _RecordingForgejo() + fake2.labels = fake.labels + r2 = MetadataHygieneReport() + run_state_label_inference(engine=engine, callbacks=_callbacks(fake2), report=r2) + assert r2.state_label_executed == 1 + added = [c[1][3] for c in fake2.calls if c[0] == "add_label"] + assert "State/Reviewing" in added + + def test_state_label_vocab_covers_all_known_states(self): + """The _STATE_TO_LABEL map MUST cover every KNOWN_STATES name — + otherwise workflows in an uncovered state silently skip without + a State/* label. A new state added to the state machine without + an update here would silently regress the operator's Forgejo + board view.""" + from tools.controller.state_machine import KNOWN_STATES + + missing = KNOWN_STATES - set(_STATE_TO_LABEL) + assert not missing, ( + f"_STATE_TO_LABEL is missing entries for {sorted(missing)} — " + f"workflows in those states won't get a Forgejo State/* label" + ) + + +# ─── check #5: milestone_assignment ────────────────────────────────── + + +class TestMilestoneAssignment: + def test_copies_milestone_from_linked_issue(self, engine, session_factory): + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Closes #42", "milestone": None, + } + # Issue #42's milestone is fetched via get_issue_state (/issues/{n}) + # not get_pr_details (/pulls/{n}) — round-2 principal-review fix. + fake.issue_details[("drew", "cleveragents-core", 42)] = { + "milestone": {"id": 7, "title": "v1.0"}, + } + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_milestone_assignment(engine=engine, callbacks=_callbacks(fake), report=report) + patch_calls = [c for c in fake.calls if c[0] == "patch_pr_milestone"] + assert len(patch_calls) == 1 + assert patch_calls[0][1][3] == 7 + assert report.milestone_executed == 1 + + def test_skips_when_pr_already_has_milestone(self, engine, session_factory): + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Closes #42", + "milestone": {"id": 99, "title": "Existing"}, + } + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_milestone_assignment(engine=engine, callbacks=_callbacks(fake), report=report) + assert not any(c[0] == "patch_pr_milestone" for c in fake.calls) + assert report.milestone_skipped == 1 + + def test_skips_when_issue_has_no_milestone(self, engine, session_factory): + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Closes #42", "milestone": None, + } + # Issue fetched via get_issue_state; milestone=None → no_issue_milestone. + fake.issue_details[("drew", "cleveragents-core", 42)] = {"milestone": None} + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_milestone_assignment(engine=engine, callbacks=_callbacks(fake), report=report) + assert not any(c[0] == "patch_pr_milestone" for c in fake.calls) + + +# ─── dispatcher + cross-tick ───────────────────────────────────────── + + +class TestDispatcher: + def test_dispatcher_invokes_all_enabled_checks(self, engine, session_factory): + _seed_workflow(session_factory, state="MERGED", pr_number=401) + _seed_workflow(session_factory, state="REVIEWING", pr_number=402) + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = {"body": "Closes #42"} + fake.pr_details[("drew", "cleveragents-core", 402)] = { + "body": "Closes #43", "milestone": None, + } + # Issue #43's milestone is fetched via get_issue_state, not + # get_pr_details (round-2 fix). + fake.issue_details[("drew", "cleveragents-core", 43)] = { + "milestone": {"id": 5, "title": "v1.0"}, + } + fake.labels[("drew", "cleveragents-core", 43)] = [{"name": "Priority/High"}] + report = run_metadata_hygiene_tick(engine=engine, callbacks=_callbacks(fake)) + # All five checks fired at least once. + assert report.completed_not_closed_executed == 1 # PR #401 + assert report.state_label_executed == 2 # both PRs + assert report.milestone_executed == 1 # PR #402 + assert report.label_sync_executed == 1 # PR #402 + # closing_keyword_fixup: only PR #402 is a candidate (the other + # is MERGED → excluded from non-terminal-only checks). #402's + # body has 'Closes #43' already — skipped. + assert report.closing_keyword_fixup_skipped == 1 + + def test_master_kill_switch_skips_everything(self, engine, session_factory): + """When all per-check flags are False, no Forgejo calls fire + even with candidate workflows.""" + _seed_workflow(session_factory, state="MERGED", pr_number=401) + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = {"body": "Closes #42"} + cbs = _callbacks( + fake, + completed_not_closed=False, closing_keyword_fixup=False, + label_sync=False, state_label_inference=False, + milestone_assignment=False, + ) + run_metadata_hygiene_tick(engine=engine, callbacks=cbs) + assert fake.calls == [] + + +# ─── round-1 review additions ──────────────────────────────────────── + + +class TestNoDialectSpecificSQL: + """Phase 2/3 grep-test pattern — Phase 4's inline ``text()`` + queries on ``workflows`` + ``grooming_decisions`` MUST be dialect- + portable (no SQLite-specific ``json_extract`` or PostgreSQL- + specific ``->>`` operators). Round-1 test-engineer-review fix. + """ + + def test_metadata_hygiene_uses_no_json_extract(self): + from pathlib import Path + + path = ( + Path(__file__).resolve().parents[3] + / "tools" / "controller" / "master" / "metadata_hygiene.py" + ) + source = path.read_text() + assert "json_extract(" not in source, ( + "metadata_hygiene.py contains a hardcoded json_extract( — " + "use a dialect-aware helper instead so the query works on " + "PostgreSQL" + ) + assert "->>" not in source + + +class TestKillSwitchWiringContract: + """Phase 3 grep-test pattern — verify the master kill switch + actually gates the wiring. Without these tests a refactor that + always-wires the callbacks bypasses the kill switch silently. + Round-1 test-engineer-review fix. + """ + + def test_loop_guards_metadata_hygiene_tick_on_callbacks_none(self): + from pathlib import Path + + loop_source = ( + Path(__file__).resolve().parents[3] + / "tools" / "controller" / "master" / "loop.py" + ).read_text() + assert "if metadata_hygiene_callbacks is not None:" in loop_source, ( + "loop.py is missing the None-guard for " + "metadata_hygiene_callbacks — the kill switch's audit-only " + "mode depends on it" + ) + + def test_main_wires_metadata_hygiene_only_when_enabled(self): + from pathlib import Path + + main_source = ( + Path(__file__).resolve().parents[3] + / "tools" / "controller" / "master" / "__main__.py" + ).read_text() + assert "metadata_hygiene_callbacks=" in main_source + assert "if _metadata_hygiene_cfg.enabled" in main_source, ( + "__main__.py is missing the _metadata_hygiene_cfg.enabled " + "gate on metadata_hygiene_callbacks; the kill switch is " + "bypassed and every deploy goes full-write" + ) + + +class TestErrorPathHandling: + """Round-1 test-engineer-review fix: each tick has try/except + around Forgejo callbacks; the raise path → report.errors was + untested for any of the 5 checks.""" + + def test_completed_not_closed_forgejo_raises_recorded( + self, engine, session_factory + ): + _seed_workflow(session_factory, state="MERGED") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = {"body": "Closes #42"} + + def boom(*a, **k): + raise RuntimeError("forgejo down") + fake.patch_pr_state = boom # type: ignore[assignment] + + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_completed_not_closed( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + assert any( + check == "completed_not_closed" and "forgejo down" in msg + for (check, _wf, msg) in report.errors + ) + + def test_closing_keyword_fixup_forgejo_raises_recorded( + self, engine, session_factory + ): + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Mentions #42.", + } + + def boom(*a, **k): + raise RuntimeError("forgejo down") + fake.patch_pr_body = boom # type: ignore[assignment] + + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_closing_keyword_fixup( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + assert any( + check == "closing_keyword_fixup" and "forgejo down" in msg + for (check, _wf, msg) in report.errors + ) + + def test_milestone_assignment_forgejo_raises_recorded( + self, engine, session_factory + ): + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Closes #42", "milestone": None, + } + # Issue #42 fetched via get_issue_state (round-2 fix). + fake.issue_details[("drew", "cleveragents-core", 42)] = { + "milestone": {"id": 7, "title": "v1.0"}, + } + + def boom(*a, **k): + raise RuntimeError("forgejo down") + fake.patch_pr_milestone = boom # type: ignore[assignment] + + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_milestone_assignment( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + assert any( + check == "milestone_assignment" and "forgejo down" in msg + for (check, _wf, msg) in report.errors + ) + + +class TestPRNotFoundHandling: + """Round-1 test-engineer-review fix: the ``pr_not_found`` branch + (Forgejo returns nothing for the PR — deleted / 404) writes an + executed=True audit row for permanent skip. Untested previously + for checks #1, #2, #3, #5.""" + + @pytest.mark.parametrize("check_fn_name,state", [ + ("run_completed_not_closed", "MERGED"), + ("run_closing_keyword_fixup", "ANALYZING"), + ("run_label_sync_from_issue", "ANALYZING"), + ("run_milestone_assignment", "ANALYZING"), + ]) + def test_pr_not_found_records_permanent_skip( + self, engine, session_factory, check_fn_name, state, + ): + from tools.controller.master import metadata_hygiene + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + wf_id = _seed_workflow(session_factory, state=state) + fake = _RecordingForgejo() + # PR not seeded → get_pr_details returns None → pr_not_found. + report = MetadataHygieneReport() + getattr(metadata_hygiene, check_fn_name)( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + with engine.connect() as c: + row = c.execute( + text( + "SELECT executed, reason_category FROM grooming_decisions " + "WHERE workflow_id = :w" + ), + {"w": wf_id}, + ).first() + assert row is not None + assert row.executed == 1, "pr_not_found is a PERMANENT skip" + assert row.reason_category == "pr_not_found" + + +class TestClosingKeywordFixupMarkdownEdgeCases: + """Round-1 principal-review + test-engineer-review fix: the + ``_BARE_REF_RE`` regex must not misread markdown links + ``[#42](url)`` or URL fragments ``host#42`` as bare refs (would + cause spurious PATCHing of the PR body).""" + + def test_markdown_link_not_misread(self, engine, session_factory): + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Related to [#42](https://example/42).", + } + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_closing_keyword_fixup( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + # Must NOT PATCH the body — the [#42](...) link is intentional + # and not a bare ref the fixup should auto-close. + assert not any(c[0] == "patch_pr_body" for c in fake.calls) + + def test_url_fragment_not_misread(self, engine, session_factory): + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "See https://example.com/page#42 for context.", + } + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + report = MetadataHygieneReport() + run_closing_keyword_fixup( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + assert not any(c[0] == "patch_pr_body" for c in fake.calls) + + def test_no_bare_ref_writes_executed_zero_for_re_check( + self, engine, session_factory, + ): + """Round-1 architect-review fix: skips for ``no_bare_ref`` / + ``multiple_bare_refs`` write executed=False so the PR body + can evolve into a single bare ref later and be re-evaluated. + Pre-fix executed=True permanently locked future re-evaluation. + """ + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Mentions #42 and also #43.", # ambiguous: 2 bare refs + } + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + run_closing_keyword_fixup( + engine=engine, callbacks=_callbacks(fake), report=MetadataHygieneReport(), + ) + with engine.connect() as c: + row = c.execute( + text( + "SELECT executed, reason_category FROM grooming_decisions " + "WHERE check_name = 'closing_keyword_fixup'" + ), + ).first() + assert row.executed == 0, ( + "transient skip — must allow re-evaluation if the PR body " + "is edited later to disambiguate" + ) + assert row.reason_category == "multiple_bare_refs" + + def test_executed_audit_preserves_pre_body( + self, engine, session_factory, + ): + """Round-1 architect-review fix: a successful body PATCH + persists the original body to llm_reasoning so the mutation + is recoverable if it ever corrupts a PR body in production.""" + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + original_body = "Mentions #42." + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": original_body, + } + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + run_closing_keyword_fixup( + engine=engine, callbacks=_callbacks(fake), report=MetadataHygieneReport(), + ) + with engine.connect() as c: + row = c.execute( + text( + "SELECT llm_reasoning FROM grooming_decisions " + "WHERE check_name = 'closing_keyword_fixup'" + ), + ).first() + assert row.llm_reasoning is not None + assert original_body in row.llm_reasoning + + +class TestForgejoStatusClassification: + """Round-1 test-engineer-review fix: the 4xx/5xx/429 classification + branches feed report.errors paths but were untested. Phase 4 now + imports the helper from forgejo_writes; this test pins the + contract by exercising the helper directly.""" + + @pytest.mark.parametrize("status,expected_category", [ + (200, "ok"), + (201, "ok"), + (404, "ok-no-op"), + (429, "pending-retry"), + (500, "pending-retry"), + (502, "pending-retry"), + (400, "failed"), + (401, "failed"), + (403, "failed"), + (422, "failed"), + ]) + def test_classify_status(self, status, expected_category): + from tools.controller.master.forgejo_writes import ( + _classify_forgejo_status, + ) + cat, _err = _classify_forgejo_status(status) + assert cat == expected_category + + +class TestCandidatesForCheck: + """Round-1 architect-review fix: the new LEFT-JOIN candidate query + is the load-bearing piece for both (a) idempotency (excludes + workflows with executed=1 audit rows) and (b) scale (workflows + with only executed=0 transient skips stay in the candidate set + for re-evaluation).""" + + def test_executed_one_excludes_workflow(self, engine, session_factory): + from tools.controller.db.session import session_scope + from tools.controller.master.metadata_hygiene import ( + _candidates_for_check, _record_audit, + ) + + wf_id = _seed_workflow(session_factory, state="MERGED") + with session_scope(engine) as s: + _record_audit( + s, workflow_id=wf_id, check_name="completed_not_closed", + reason_category="closed_1_of_1", executed=True, + ) + with session_scope(engine) as s: + rows = _candidates_for_check( + s, "completed_not_closed", state_filter="merged", + ) + assert rows == [] + + def test_executed_zero_keeps_workflow_as_candidate( + self, engine, session_factory, + ): + from tools.controller.db.session import session_scope + from tools.controller.master.metadata_hygiene import ( + _candidates_for_check, _record_audit, + ) + + wf_id = _seed_workflow(session_factory, state="ANALYZING") + with session_scope(engine) as s: + _record_audit( + s, workflow_id=wf_id, check_name="label_sync_from_issue", + reason_category="no_linked_issues", executed=False, + ) + with session_scope(engine) as s: + rows = _candidates_for_check( + s, "label_sync_from_issue", state_filter="non_terminal", + ) + assert len(rows) == 1 + assert rows[0].workflow_id == wf_id + + def test_unknown_state_filter_raises(self, engine, session_factory): + from tools.controller.db.session import session_scope + from tools.controller.master.metadata_hygiene import _candidates_for_check + + with session_scope(engine) as s: + with pytest.raises(ValueError, match="unknown state_filter"): + _candidates_for_check(s, "x", state_filter="garbage") + + +# ─── round-2 review additions ──────────────────────────────────────── +# Covers seven gaps left after round-1: +# Gap 1 — completed_not_closed partial-success (N refs, K < N succeed) +# Gap 2 — adjust_labels failure in label_sync → no executed=1 row → retry +# Gap 3 — adjust_labels failure in state_label_inference → no executed=1 row +# Gap 4 — state_label_inference runs on MERGED / ABANDONED workflows +# Gap 5 — _last_synced_state returns None when only executed=0 rows exist +# (dry-run then real-run re-fires correctly) +# Gap 6 — closing_keyword_fixup body has bare ref AND matching Closes #N +# for the SAME issue → candidates == {} → skips +# Gap 7 — error path coverage for label_sync + state_label_inference +# (TestErrorPathHandling was missing these two checks) + + +class TestCompletedNotClosedPartialSuccess: + """Gap 1 — partial-success when N refs, K < N close OK. + + Round-2 architect fix: completed_not_closed now writes executed=True + only when ALL refs close (``all_closed = closed_count == len(refs)``). + Partial success writes executed=False so the workflow stays in + candidates and retries next tick; already-closed issues hit 200/ok-no-op + and do not double-close. + """ + + def test_partial_success_one_of_two_refs_fails( + self, engine, session_factory + ): + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + wf_id = _seed_workflow(session_factory, state="MERGED") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Closes #42\n\nCloses #55", + } + + # Issue #42 → 200 OK; issue #55 → 429 (rate-limited, not closed). + def patch_state_mixed(owner, repo, n, state): + fake.calls.append(("patch_pr_state", (owner, repo, n, state))) + if n == 42: + return {"status": 200, "body": {"state": "closed"}} + return {"status": 429, "body": {}} # #55 fails + + fake.patch_pr_state = patch_state_mixed # type: ignore[assignment] + + report = MetadataHygieneReport() + run_completed_not_closed( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + + # Both issues were attempted (loop does NOT short-circuit on error). + attempted = [c for c in fake.calls if c[0] == "patch_pr_state"] + assert len(attempted) == 2 + + # Partial success: closed_1_of_2 — executed=False (retry next tick). + with engine.connect() as c: + row = c.execute( + text( + "SELECT reason_category, executed " + "FROM grooming_decisions " + "WHERE workflow_id = :w AND check_name = 'completed_not_closed'" + ), + {"w": wf_id}, + ).first() + assert row is not None + assert row.reason_category == "closed_1_of_2" + # Round-2 fix: executed=False — partial success must NOT permanently + # lock out the remaining unprocessed issue. + assert row.executed == 0 + + # 429 path lands in report.errors. + assert any( + check == "completed_not_closed" and "55" in msg + for (check, _wf, msg) in report.errors + ) + + def test_zero_of_two_refs_succeed_writes_executed_false( + self, engine, session_factory + ): + """When ALL close calls fail (both 429), audit row is written with + closed_0_of_2 and executed=False so the workflow retries. + Previously (pre round-2 fix) executed=True permanently locked + out both issues even though neither closed.""" + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + wf_id = _seed_workflow(session_factory, state="MERGED") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Closes #42\n\nCloses #55", + } + fake.patch_status = 429 # all PATCH calls fail + + report = MetadataHygieneReport() + run_completed_not_closed( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + + with engine.connect() as c: + row = c.execute( + text( + "SELECT reason_category, executed " + "FROM grooming_decisions " + "WHERE workflow_id = :w AND check_name = 'completed_not_closed'" + ), + {"w": wf_id}, + ).first() + assert row is not None + assert row.reason_category == "closed_0_of_2" + # Round-2 fix: executed=False — zero success must not lock out + # either issue; the workflow must retry next tick. + assert row.executed == 0 + assert len(report.errors) == 2 # one per failed issue + + def test_full_success_writes_executed_true(self, engine, session_factory): + """Baseline: all refs close OK → executed=True (permanent lock, + check #1 is done for this workflow).""" + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + wf_id = _seed_workflow(session_factory, state="MERGED") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Closes #42\n\nCloses #55", + } + # Default patch_status = 200 → all succeed. + report = MetadataHygieneReport() + run_completed_not_closed( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + with engine.connect() as c: + row = c.execute( + text( + "SELECT reason_category, executed FROM grooming_decisions " + "WHERE workflow_id = :w AND check_name = 'completed_not_closed'" + ), + {"w": wf_id}, + ).first() + assert row.reason_category == "closed_2_of_2" + assert row.executed == 1 + + +class TestLabelSyncAdjustLabelsFailure: + """Gap 2 — when adjust_labels returns a 'failed' result for label_sync, + an executed=False audit row is written (for observability) but NO + executed=True row is written → workflow stays in candidates → retry + next tick. + + Round-2 architect fix: pre-fix wrote no audit row at all, leaving + no log-grep anchor if a label sync silently failed on every tick. + """ + + def test_adjust_labels_failure_writes_no_executed_one_row( + self, engine, session_factory + ): + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + wf_id = _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = {"body": "Closes #42"} + fake.labels[("drew", "cleveragents-core", 42)] = [{"name": "Priority/High"}] + # Simulate add_label returning False → adjust_labels action='failed'. + fake.add_label_ok = False + + report = MetadataHygieneReport() + run_label_sync_from_issue( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + + # Error recorded in report. + assert any( + check == "label_sync_from_issue" + for (check, _wf, _msg) in report.errors + ) + + # No executed=1 audit row — workflow stays in candidates for retry. + with engine.connect() as c: + row = c.execute( + text( + "SELECT executed, reason_category FROM grooming_decisions " + "WHERE workflow_id = :w " + "AND check_name = 'label_sync_from_issue'" + ), + {"w": wf_id}, + ).first() + # Round-2 fix: an executed=False row IS written for observability. + assert row is not None + assert row.executed == 0 + assert row.reason_category == "label_apply_failed" + + def test_adjust_labels_failure_workflow_retried_next_tick( + self, engine, session_factory + ): + """After a failed adjust_labels pass, the next tick (with + add_label_ok restored) succeeds and writes executed=1.""" + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + wf_id = _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = {"body": "Closes #42"} + fake.labels[("drew", "cleveragents-core", 42)] = [{"name": "Priority/High"}] + fake.add_label_ok = False + + # Tick 1 — fails. + run_label_sync_from_issue( + engine=engine, callbacks=_callbacks(fake), report=MetadataHygieneReport(), + ) + + # Tick 2 — Forgejo recovers. + fake.add_label_ok = True + r2 = MetadataHygieneReport() + run_label_sync_from_issue( + engine=engine, callbacks=_callbacks(fake), report=r2, + ) + assert r2.label_sync_executed == 1 + + with engine.connect() as c: + row = c.execute( + text( + "SELECT executed FROM grooming_decisions " + "WHERE workflow_id = :w " + "AND check_name = 'label_sync_from_issue' " + "AND executed = 1" + ), + {"w": wf_id}, + ).first() + assert row is not None + + +class TestStateLabelAdjustLabelsFailure: + """Gap 3 — when adjust_labels returns 'failed' in state_label_inference, + NO executed=1 audit row is written → ``_last_synced_state`` still + returns the prior state (or None) → workflow retries next tick. + + Latent bug risk: an errant executed=True write would permanently + anchor the State/* label to the wrong state. + """ + + def test_adjust_labels_failure_writes_no_executed_one_row( + self, engine, session_factory + ): + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + wf_id = _seed_workflow(session_factory, state="REVIEWING") + fake = _RecordingForgejo() + fake.add_label_ok = False # simulate Forgejo label-add failure + + report = MetadataHygieneReport() + run_state_label_inference( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + + assert any( + check == "state_label_inference" + for (check, _wf, _msg) in report.errors + ) + + with engine.connect() as c: + row = c.execute( + text( + "SELECT executed FROM grooming_decisions " + "WHERE workflow_id = :w " + "AND check_name = 'state_label_inference' " + "AND executed = 1" + ), + {"w": wf_id}, + ).first() + assert row is None, ( + "state_label_inference adjust_labels failure must NOT write " + "executed=1 — _last_synced_state must return None so the " + "label syncs on the next tick" + ) + + def test_adjust_labels_failure_workflow_retried_next_tick( + self, engine, session_factory + ): + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + wf_id = _seed_workflow(session_factory, state="REVIEWING") + fake = _RecordingForgejo() + fake.add_label_ok = False + + # Tick 1 — fails. + run_state_label_inference( + engine=engine, callbacks=_callbacks(fake), report=MetadataHygieneReport(), + ) + + # Tick 2 — Forgejo recovers. + fake.add_label_ok = True + fake2 = _RecordingForgejo() + fake2.labels = fake.labels + r2 = MetadataHygieneReport() + run_state_label_inference( + engine=engine, callbacks=_callbacks(fake2), report=r2, + ) + assert r2.state_label_executed == 1 + + +class TestStateLabelInferenceTerminalWorkflows: + """Gap 4 — state_label_inference must sync terminal workflows + (MERGED, ABANDONED) since _STATE_TO_LABEL covers them. + + Latent bug: if the query were tightened to 'non_terminal', the + Forgejo board view for closed PRs would freeze at their last + non-terminal State/* label (e.g. State/Reviewing). + """ + + @pytest.mark.parametrize("terminal_state,expected_label", [ + ("MERGED", "State/Merged"), + ("ABANDONED", "State/Abandoned"), + ]) + def test_syncs_state_label_for_terminal_workflow( + self, engine, session_factory, terminal_state, expected_label + ): + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + _seed_workflow(session_factory, state=terminal_state) + fake = _RecordingForgejo() + + report = MetadataHygieneReport() + run_state_label_inference( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + + added = [c[1][3] for c in fake.calls if c[0] == "add_label"] + assert expected_label in added, ( + f"state_label_inference did not sync {expected_label} for " + f"a {terminal_state} workflow" + ) + assert report.state_label_executed == 1 + + def test_terminal_workflow_audit_row_has_executed_true( + self, engine, session_factory + ): + """Audit row for a terminal workflow: executed=1, reason=state.""" + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + wf_id = _seed_workflow(session_factory, state="MERGED") + fake = _RecordingForgejo() + + run_state_label_inference( + engine=engine, + callbacks=_callbacks(fake), + report=MetadataHygieneReport(), + ) + + with engine.connect() as c: + row = c.execute( + text( + "SELECT executed, reason_category FROM grooming_decisions " + "WHERE workflow_id = :w AND check_name = 'state_label_inference'" + ), + {"w": wf_id}, + ).first() + assert row is not None + assert row.executed == 1 + assert row.reason_category == "MERGED" + + +class TestLastSyncedStateDryRunThenReal: + """Gap 5 — _last_synced_state returns None when only executed=0 rows + exist (WHERE requires executed=1). + + Dry-run pass writes executed=False, so _last_synced_state still + returns None on the following real pass → real pass correctly fires + adjust_labels and writes executed=1. + + Latent bug: if _last_synced_state read executed=0 rows, the dry-run + would permanently suppress the real label sync. + """ + + def test_dry_run_does_not_satisfy_last_synced_state( + self, engine, session_factory + ): + from tools.controller.master.metadata_hygiene import ( + MetadataHygieneReport, + _last_synced_state, + ) + from tools.controller.db.session import session_scope + + wf_id = _seed_workflow(session_factory, state="REVIEWING") + fake_dry = _RecordingForgejo() + + # Dry-run pass — writes executed=False audit row. + run_state_label_inference( + engine=engine, + callbacks=_callbacks(fake_dry, dry_run=True), + report=MetadataHygieneReport(), + ) + + # Confirm dry-run wrote executed=0 row. + with engine.connect() as c: + dry_row = c.execute( + text( + "SELECT executed FROM grooming_decisions " + "WHERE workflow_id = :w AND check_name = 'state_label_inference'" + ), + {"w": wf_id}, + ).first() + assert dry_row is not None + assert dry_row.executed == 0 + + # _last_synced_state must still return None — only executed=1 counts. + with session_scope(engine) as s: + last = _last_synced_state(s, wf_id) + assert last is None, ( + "_last_synced_state must ignore executed=0 rows; dry-run pass " + "must not suppress the subsequent real sync" + ) + + def test_real_run_after_dry_run_syncs_forgejo_label( + self, engine, session_factory + ): + """Integration: dry-run pass → real pass → adjust_labels fires.""" + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + wf_id = _seed_workflow(session_factory, state="REVIEWING") + fake_dry = _RecordingForgejo() + + # Dry-run pass. + run_state_label_inference( + engine=engine, + callbacks=_callbacks(fake_dry, dry_run=True), + report=MetadataHygieneReport(), + ) + assert not any(c[0] == "add_label" for c in fake_dry.calls) + + # Real pass — must fire adjust_labels. + fake_real = _RecordingForgejo() + fake_real.labels = fake_dry.labels + r2 = MetadataHygieneReport() + run_state_label_inference( + engine=engine, + callbacks=_callbacks(fake_real, dry_run=False), + report=r2, + ) + assert r2.state_label_executed == 1 + added = [c[1][3] for c in fake_real.calls if c[0] == "add_label"] + assert "State/Reviewing" in added + + +class TestClosingKeywordFixupBareRefAlreadyCovered: + """Gap 6 — bare #N ref already covered by a 'Closes #N' for the SAME + issue: candidates = bare_refs - existing = {} → skips as 'no_bare_ref'. + + Latent bug: without the subtraction the fixup would append a duplicate + 'Closes #42', corrupting the PR body. + """ + + def test_bare_ref_already_covered_skips_no_bare_ref( + self, engine, session_factory + ): + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + wf_id = _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + # Body has bare #42 AND "Closes #42" → candidates after subtraction = {}. + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Related to #42.\n\nCloses #42", + } + + report = MetadataHygieneReport() + run_closing_keyword_fixup( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + + # No PATCH — the bare ref is already covered. + assert not any(c[0] == "patch_pr_body" for c in fake.calls), ( + "closing_keyword_fixup must not PATCH when the bare ref is " + "already covered by an existing Closes keyword" + ) + assert report.closing_keyword_fixup_skipped == 1 + + # Transient skip — executed=False so the workflow can be re-checked. + with engine.connect() as c: + row = c.execute( + text( + "SELECT executed, reason_category FROM grooming_decisions " + "WHERE workflow_id = :w AND check_name = 'closing_keyword_fixup'" + ), + {"w": wf_id}, + ).first() + assert row is not None + assert row.executed == 0 + assert row.reason_category == "no_bare_ref" + + def test_bare_ref_different_issue_than_closes_still_fires( + self, engine, session_factory + ): + """Negative control: bare #43 is not covered by Closes #42 → + candidates = {43} → fixup proceeds.""" + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = { + "body": "Fixes the thing in #43.\n\nCloses #42", + } + + report = MetadataHygieneReport() + run_closing_keyword_fixup( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + + patch_calls = [c for c in fake.calls if c[0] == "patch_pr_body"] + assert len(patch_calls) == 1 + assert "Closes #43" in patch_calls[0][1][3] + assert report.closing_keyword_fixup_executed == 1 + + +class TestErrorPathHandlingRound2: + """Gap 7 — error path coverage for label_sync (#3) and state_label + (#4), which TestErrorPathHandling from round-1 left dark. + + These two have a different failure shape: they fail via adjust_labels + returning action='failed', not via an exception from a Forgejo + callback. Both failure shapes (exception + failed-result) are covered. + """ + + def test_label_sync_get_labels_raises_recorded( + self, engine, session_factory + ): + """get_labels raising inside _process_label_sync is caught + per-issue; loop continues; error appended to report.errors.""" + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = {"body": "Closes #42"} + + def boom_get_labels(owner, repo, n): + fake.calls.append(("get_labels", (owner, repo, n))) + raise RuntimeError("labels endpoint down") + + fake.get_labels = boom_get_labels # type: ignore[assignment] + + report = MetadataHygieneReport() + run_label_sync_from_issue( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + assert any( + check == "label_sync_from_issue" and "labels endpoint down" in msg + for (check, _wf, msg) in report.errors + ) + + def test_label_sync_adjust_labels_failed_recorded( + self, engine, session_factory + ): + """adjust_labels action='failed' → error in report + executed=False.""" + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + wf_id = _seed_workflow(session_factory, state="ANALYZING") + fake = _RecordingForgejo() + fake.pr_details[("drew", "cleveragents-core", 401)] = {"body": "Closes #42"} + fake.labels[("drew", "cleveragents-core", 42)] = [{"name": "Priority/High"}] + fake.add_label_ok = False + + report = MetadataHygieneReport() + run_label_sync_from_issue( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + assert any( + check == "label_sync_from_issue" + for (check, _wf, _msg) in report.errors + ), "adjust_labels failure must be recorded in report.errors" + + def test_state_label_inference_adjust_labels_failed_recorded( + self, engine, session_factory + ): + """adjust_labels action='failed' inside _process_state_label_sync + → error in report.errors; skipped counter NOT incremented (this is + a hard failure, not a benign transient skip).""" + from tools.controller.master.metadata_hygiene import MetadataHygieneReport + + _seed_workflow(session_factory, state="REVIEWING") + fake = _RecordingForgejo() + fake.add_label_ok = False + + report = MetadataHygieneReport() + run_state_label_inference( + engine=engine, callbacks=_callbacks(fake), report=report, + ) + assert any( + check == "state_label_inference" + for (check, _wf, _msg) in report.errors + ), "adjust_labels failure must be recorded in report.errors" + # Hard failure ≠ benign skip — skipped counter must stay zero. + assert report.state_label_skipped == 0 diff --git a/tools/controller/master/__main__.py b/tools/controller/master/__main__.py index bdb9a00e4..9bedbb72f 100644 --- a/tools/controller/master/__main__.py +++ b/tools/controller/master/__main__.py @@ -39,6 +39,11 @@ from .gate3_abandon_config import ( ) from .grooming_side_effects import GroomingCallbacks from .loop import MasterConfig, master_main_loop +from .metadata_hygiene import MetadataHygieneCallbacks +from .metadata_hygiene_config import ( + get_metadata_hygiene_config as _get_metadata_hygiene_cfg, + log_effective_config as _log_metadata_hygiene_cfg, +) from .reviewer_abandon_side_effects import ReviewerAbandonCallbacks from .prefetch import PrefetchDataCallbacks, make_prefetch_callback @@ -282,6 +287,11 @@ def main(argv: list[str] | None = None) -> int: # incident happened?". _log_gate3_abandon_cfg() _gate3_abandon_cfg = _get_gate3_abandon_cfg() + # Phase 4 (metadata-hygiene, 2026-05-25): log + load the config. + # Same startup-forensics rationale — operators get a log-grep + # anchor for "what was set when this incident happened?". + _log_metadata_hygiene_cfg() + _metadata_hygiene_cfg = _get_metadata_hygiene_cfg() # Phase 1k++ (N6): parser-coverage check runs BEFORE backfill + # main loop so strict-mode failure exits 2 without wasting a @@ -500,6 +510,43 @@ def main(argv: list[str] | None = None) -> int: if _gate3_abandon_cfg.enabled else None ), + # Phase 4 (2026-05-25): metadata-hygiene tick. Gated on the + # master CONTROLLER_METADATA_HYGIENE_ENABLED flag (default + # False) — each of the five sub-checks has its own per-check + # flag. ``dry_run`` shares the grooming flag for symmetry with + # Gates 1/2/3 (one safe-rollout toggle covers all metadata + # write paths). + metadata_hygiene_callbacks=( + MetadataHygieneCallbacks( + get_pr_details=callbacks.get_pr_details, + # get_issue_state hits /issues/{n} — needed by check #5 + # to read an issue's milestone. Using get_pr_details + # (/pulls/{n}) would 404 on plain issues. (Round-2 fix.) + get_issue_state=callbacks.get_issue_state, + get_labels=callbacks.get_labels, + add_label=callbacks.add_label, + remove_label=callbacks.remove_label, + patch_pr_state=callbacks.patch_pr_state, + patch_pr_body=callbacks.patch_pr_body, + patch_pr_milestone=callbacks.patch_pr_milestone, + completed_not_closed_enabled=( + _metadata_hygiene_cfg.completed_not_closed + ), + closing_keyword_fixup_enabled=( + _metadata_hygiene_cfg.closing_keyword_fixup + ), + label_sync_enabled=_metadata_hygiene_cfg.label_sync_from_issue, + state_label_inference_enabled=( + _metadata_hygiene_cfg.state_label_inference + ), + milestone_assignment_enabled=( + _metadata_hygiene_cfg.milestone_assignment + ), + dry_run=_grooming_cfg.dry_run, + ) + if _metadata_hygiene_cfg.enabled + else None + ), # RUN_CI_LOCAL: skip ci_poll_exhaustion while local CI is busy # (None — a no-op — under remote CI). local_ci_in_flight=local_ci_in_flight, diff --git a/tools/controller/master/forgejo_http.py b/tools/controller/master/forgejo_http.py index 2a98edc43..447058704 100644 --- a/tools/controller/master/forgejo_http.py +++ b/tools/controller/master/forgejo_http.py @@ -61,6 +61,19 @@ GetCIStatusCallback = _Callable[[str, str, str], dict | None] # text for the failing jobs (empty string when none / unreachable). GetFailureLogsCallback = _Callable[[str, str, str], str] +# Phase 4 metadata-hygiene PATCH callbacks (2026-05-25). Both PATCH +# Forgejo's unified issues endpoint: +# PATCH /repos/{owner}/{repo}/issues/{n} body={"": } +# Return the raw ``{"status": int, "body": ...}`` shape so the per- +# check orchestrator can dispatch on the error-handling matrix (200 = +# success, 404 = treat as no-op, 422 = stuck, 5xx/429 = retry). +# patch_pr_body(owner, repo, pr_number, body) → dict +PatchPRBodyCallback = _Callable[[str, str, int, str], dict] +# patch_pr_milestone(owner, repo, pr_number, milestone_id) → dict +# ``milestone_id`` may be ``None`` to clear the assignment, or an int +# to set it. +PatchPRMilestoneCallback = _Callable[[str, str, int, "int | None"], dict] + @dataclass class ForgejoCallbacks: @@ -78,6 +91,12 @@ class ForgejoCallbacks: # PATCH /issues/{n} {"state": ...} — used by grooming's close path # (Phase 0 grooming plan; orchestration lives in forgejo_writes.close_issue). patch_pr_state: fw.PatchPRStateCallback + # PATCH /issues/{n} {"body": ...} — Phase 4 metadata-hygiene + # (closing-keyword fixup adds ``Closes #N`` to the PR body). + patch_pr_body: "PatchPRBodyCallback" + # PATCH /issues/{n} {"milestone": ...} — Phase 4 metadata-hygiene + # (milestone assignment copies milestone from linked issue). + patch_pr_milestone: "PatchPRMilestoneCallback" merge_pr: mg.MergeCallback # Reconciliation callbacks (Phase 1g): get_pr_state: rec.GetPRStateCallback @@ -136,6 +155,8 @@ def build_callbacks( add_label=_make_add_label(cfg, runtime), remove_label=_make_remove_label(cfg, runtime), patch_pr_state=_make_patch_pr_state(cfg, runtime), + patch_pr_body=_make_patch_pr_body(cfg, runtime), + patch_pr_milestone=_make_patch_pr_milestone(cfg, runtime), merge_pr=_make_merge_pr(cfg, runtime), get_pr_state=_make_get_pr_state(cfg, runtime), get_issue_state=_make_get_issue_state(cfg, runtime), @@ -296,6 +317,53 @@ def _make_patch_pr_state(cfg, runtime): return patch_pr_state +def _make_patch_pr_body(cfg, runtime): + """Build the PATCH-PR-body closure (Phase 4 metadata-hygiene). + + Forgejo's unified issues endpoint accepts ``body`` mutations for + both issues and PRs: + PATCH /repos/{owner}/{repo}/issues/{n} body={"body": "..."} + + Used by the closing-keyword-fixup tick to add ``Closes #N`` to a + PR body that references issue N without the closing keyword. + """ + + def patch_pr_body( + owner: str, + repo: str, + pr_number: int, + body: str, + ) -> dict: + path = f"/repos/{owner}/{repo}/issues/{int(pr_number)}" + return runtime.patch(path, cfg, {"body": body}) + + return patch_pr_body + + +def _make_patch_pr_milestone(cfg, runtime): + """Build the PATCH-PR-milestone closure (Phase 4 metadata-hygiene). + + Forgejo's unified issues endpoint accepts ``milestone`` mutations + for both issues and PRs: + PATCH /repos/{owner}/{repo}/issues/{n} body={"milestone": } + + Pass ``None`` to clear, int milestone-id to set. Used by the + milestone-assignment tick to copy the milestone from a linked + issue onto its PR when the PR has none. + """ + + def patch_pr_milestone( + owner: str, + repo: str, + pr_number: int, + milestone_id: int | None, + ) -> dict: + path = f"/repos/{owner}/{repo}/issues/{int(pr_number)}" + return runtime.patch(path, cfg, {"milestone": milestone_id}) + + return patch_pr_milestone + + def _make_remove_label(cfg, runtime): def remove_label( owner: str, diff --git a/tools/controller/master/loop.py b/tools/controller/master/loop.py index c3e0d7c2d..ac9ffbd3f 100644 --- a/tools/controller/master/loop.py +++ b/tools/controller/master/loop.py @@ -58,6 +58,11 @@ from .grooming_side_effects import ( run_grooming_side_effects_tick, ) from .merging import MergeCallback, MergingHandlerReport, run_merging_tick +from .metadata_hygiene import ( + MetadataHygieneCallbacks, + MetadataHygieneReport, + run_metadata_hygiene_tick, +) from .reviewer_abandon_side_effects import ( ReviewerAbandonCallbacks, ReviewerAbandonSideEffectReport, @@ -237,6 +242,13 @@ def master_main_loop( # and performs the Forgejo close via ``forgejo_writes.close_act`` # with ``cause=Cause.REVIEWER_ABANDON``. Same shape + dry_run # source as the estimator-abandon tick — Phase 3 (2026-05-25). + metadata_hygiene_callbacks: MetadataHygieneCallbacks | None = None, + # Phase 4 (2026-05-25): when set, the metadata-hygiene tick runs + # every iteration; it invokes each enabled check (completed-not- + # closed, closing-keyword fixup, label sync, state-label inference, + # milestone assignment). Each check is independently gated via + # per-check flags on the callbacks dataclass. None disables the + # entire phase regardless of per-check flags. grooming_callbacks: GroomingCallbacks | None = None, # When set, the grooming side-effect tick fires every iteration: # it finds workflows whose state-machine just transitioned via @@ -474,6 +486,25 @@ def master_main_loop( "reviewer_abandon_side_effects tick raised; continuing" ) + # Phase 4 (2026-05-25) — metadata-hygiene dispatcher. + # Runs every iteration; each of the five checks + # (completed-not-closed, closing-keyword fixup, label sync, + # state-label inference, milestone assignment) is gated by + # its own per-check flag on ``metadata_hygiene_callbacks``. + # ``None`` disables the whole phase regardless of per-check + # flags. Cheap when no candidates exist per check. + metadata_hygiene_report: MetadataHygieneReport | None = None + if metadata_hygiene_callbacks is not None: + try: + metadata_hygiene_report = run_metadata_hygiene_tick( + engine=engine, + callbacks=metadata_hygiene_callbacks, + ) + except Exception: + logger.exception( + "metadata_hygiene tick raised; continuing" + ) + # Phase 1k+++ (real-run): MERGING handler. For workflows # in MERGING state, call the Forgejo merge endpoint via # the injected callback. Without this, workflows that diff --git a/tools/controller/master/metadata_hygiene.py b/tools/controller/master/metadata_hygiene.py new file mode 100644 index 000000000..e2ffdcada --- /dev/null +++ b/tools/controller/master/metadata_hygiene.py @@ -0,0 +1,1136 @@ +"""Phase 4 — metadata-hygiene side-effect ticks. + +Five deterministic, idempotent checks the controller performs on +Forgejo PRs / issues. None of these affect the state machine; they +mutate Forgejo-side metadata to keep PR/issue state consistent with +the workflow + each other. + +The five checks: + +1. **completed_not_closed** — When a workflow reaches MERGED, close + any open issue the PR's body links via ``Closes #N`` / ``Fixes #N``. + Common case: implementer landed the work, reviewer merged, but the + linked issue is still open because Forgejo's auto-close-on-merge + only fires when the keywords are present BEFORE merge (and Phase + 4's own closing-keyword fixup may have added them after — too + late for Forgejo's hook). + +2. **closing_keyword_fixup** — When the PR body contains a bare + ``#N`` reference (no closing verb) but issue N exists + is open, + PATCH the body to add ``Closes #N`` so Forgejo auto-closes the + issue at merge time. Conservative: only fixes when there is + exactly ONE bare ``#N`` reference and it isn't already covered + by a closing keyword. + +3. **label_sync_from_issue** — At pickup, copy ``Priority/*``, + ``Type/*``, ``MoSCoW/*`` labels from the linked issue to the PR + so the PR carries the same metadata the issue does. Idempotent: + no-op if the PR already has the label. + +4. **state_label_inference** — Periodically sync a ``State/*`` label + to mirror the workflow's ``current_state``. Operators using the + Forgejo board view see at-a-glance pipeline state without + needing the controller's CLI. + +5. **milestone_assignment** — At pickup, copy the milestone from the + linked issue to the PR if the PR has none. Idempotent — never + overrides an explicit milestone. + +Architecture +------------ +Each check has its own selection query + idempotency guard + +Forgejo write. Audit rows land in ``grooming_decisions`` with +``verdict='metadata_fixup'`` and ``check_name`` discriminating the +specific check. The ``executed=1`` flag is the per-(workflow, +check_name) idempotency marker — once a check runs successfully, +it never re-runs against the same workflow. + +Exception: ``state_label_inference`` runs on EVERY iteration for +every non-terminal workflow (it has to, since the workflow's +``current_state`` changes over time). It writes an audit row only on +state CHANGES (the previous audit row's ``reason_category`` records +the last-synced state; a tick that finds it matches the workflow's +current state is a no-op without writing). + +Trigger placement +----------------- +- All five run from ``master_main_loop`` per-iteration via + ``run_metadata_hygiene_tick`` — the dispatcher reads the config + and only invokes the enabled checks. +- Each check is internally cheap when there's nothing to do (zero + candidates → zero work). + +Safety +------ +Default-off via the master ``CONTROLLER_METADATA_HYGIENE_ENABLED`` +kill switch + per-check granular flags. Audit-only mode runs the +checks but skips the Forgejo writes (``dry_run=True``). +""" + +from __future__ import annotations + +import json +import logging +import re +from dataclasses import dataclass, field +from typing import Any + +from sqlalchemy import text +from sqlalchemy.engine import Engine + +from ..db.session import session_scope +from .forgejo_writes import ( + AddLabelCallback, + GetLabelsCallback, + PatchPRStateCallback, + RemoveLabelCallback, + _classify_forgejo_status, + adjust_labels, +) +from .forgejo_http import PatchPRBodyCallback, PatchPRMilestoneCallback +from .grooming import extract_closes_refs + +logger = logging.getLogger(__name__) + + +# ─── shared audit-row helpers ───────────────────────────────────────── + + +_VERDICT = "metadata_fixup" + + +def _candidates_for_check( + session, + check_name: str, + *, + state_filter: str, +) -> list: + """Return workflows that need ``check_name`` evaluated this sweep. + + Selection uses a LEFT JOIN against ``grooming_decisions`` to + EXCLUDE workflows that already have an ``executed=1`` audit row + for this check (the permanent-resolution case). Workflows with + only ``executed=0`` rows (transient skips — "no linked issue + yet", "no bare ref yet", etc.) STAY in the candidate list and + are re-checked each tick because the underlying condition may + have changed (PR body edited, issue gained a label, etc.). + + Round-1 architect-review fixes: + - Replaces the prior per-row ``_already_done`` SELECT (one extra + query per candidate per tick → O(N) extra SELECTs on a busy + repo) with a single LEFT JOIN in the main candidate query. + - Eliminates the unbounded-MERGED-scan complaint: a repo with 500 + historical MERGED workflows now returns ZERO candidates after + the first sweep instead of re-scanning all 500 every tick. + + ``state_filter`` is one of: + - ``"merged"`` — PR-kind workflows in MERGED (check #1 only). + - ``"non_terminal"`` — PR-kind workflows in any non-terminal + state (checks #2, #3, #5). + - ``"all_pr"`` — every PR workflow (check #4 — state-label sync + runs on every state including terminal). + """ + if state_filter == "merged": + state_clause = "AND w.kind = 'pr' AND w.current_state = 'MERGED'" + elif state_filter == "non_terminal": + state_clause = ( + "AND w.kind = 'pr' " + "AND w.current_state NOT IN " + "('MERGED', 'ABANDONED', 'STUCK', 'CREATED_PR')" + ) + elif state_filter == "all_pr": + state_clause = "AND w.kind = 'pr'" + else: + raise ValueError( + f"unknown state_filter {state_filter!r}; expected one of " + "'merged', 'non_terminal', 'all_pr'" + ) + + sql = f""" + SELECT w.workflow_id, w.owner, w.repo, w.entity_number, w.current_state + FROM workflows w + LEFT JOIN grooming_decisions gd + ON gd.workflow_id = w.workflow_id + AND gd.check_name = :check + AND gd.executed = 1 + WHERE gd.decision_id IS NULL + {state_clause} + """ + return session.execute(text(sql), {"check": check_name}).fetchall() + + +def _last_synced_state(session, workflow_id: int) -> str | None: + """Used only by state_label_inference: returns the last state we + synced a Forgejo label for, or None if we've never synced this + workflow.""" + row = session.execute( + text( + """ + SELECT reason_category + FROM grooming_decisions + WHERE workflow_id = :wf_id + AND verdict = :v + AND check_name = 'state_label_inference' + AND executed = 1 + ORDER BY decided_at DESC, decision_id DESC + LIMIT 1 + """ + ), + {"wf_id": workflow_id, "v": _VERDICT}, + ).first() + return row.reason_category if row else None + + +def _record_audit( + session, + *, + workflow_id: int, + check_name: str, + reason_category: str, + executed: bool, + llm_reasoning: str | None = None, + target_workflow_id: int | None = None, +) -> None: + """Write the audit row. Mirrors close_act's grooming_decisions + insert shape. ``executed=False`` is the dry-run / skipped path.""" + session.execute( + text( + """ + INSERT INTO grooming_decisions + (workflow_id, decided_at, check_name, stage, verdict, + reason_category, llm_reasoning, target_workflow_id, + executed) + VALUES + (:wf_id, CURRENT_TIMESTAMP, :check, 'metadata_hygiene', + :v, :rc, :lr, :tw, :exec) + """ + ), + { + "wf_id": workflow_id, + "check": check_name, + "v": _VERDICT, + "rc": reason_category, + "lr": llm_reasoning, + "tw": target_workflow_id, + "exec": 1 if executed else 0, + }, + ) + + +# Note: ``_classify_forgejo_status`` is imported from forgejo_writes so +# Phase 4 reuses the same status-code → (result, error) mapping as +# Phases 0-3. Returns ``"pending-retry"`` (not ``"retry"``) for 429/5xx +# — the literal string matters for log-grep + incident response. + + +# ─── report ─────────────────────────────────────────────────────────── + + +@dataclass +class MetadataHygieneReport: + """Per-sweep summary — counts per check.""" + + completed_not_closed_executed: int = 0 + completed_not_closed_skipped: int = 0 + closing_keyword_fixup_executed: int = 0 + closing_keyword_fixup_skipped: int = 0 + label_sync_executed: int = 0 + label_sync_skipped: int = 0 + state_label_executed: int = 0 + state_label_skipped: int = 0 + milestone_executed: int = 0 + milestone_skipped: int = 0 + errors: list[tuple[str, int, str]] = field(default_factory=list) + + +# ─── callbacks bundle ───────────────────────────────────────────────── + + +# Type alias for the Forgejo "get PR/issue details" callback. The +# existing PrefetchDataCallbacks.get_pr_details has the same shape; +# metadata-hygiene re-types it here for documentation clarity. +GetPRDetailsCallback = Any # Callable[[str, str, int], dict | None] +# get_issue_state(owner, repo, issue_n) → dict | None. Used by +# milestone_assignment to fetch the linked issue's milestone — hits +# GET /repos/{o}/{r}/issues/{n} (the issues endpoint, NOT /pulls/). +# Critical distinction: plain issues return 404 from /pulls/; using +# the issues endpoint ensures milestone data is always present. +# (Round-2 principal-review fix.) +GetIssueStateCallback = Any # Callable[[str, str, int], dict | None] + + +@dataclass(frozen=True) +class MetadataHygieneCallbacks: + """Bundle of Forgejo callbacks + per-check enable flags for the + Phase 4 metadata-hygiene tick. + + The five per-check enable flags map 1:1 to + ``MetadataHygieneConfig`` (the tick reads them off this dataclass + rather than re-importing the config so tests can flip individual + checks without monkey-patching env vars). + """ + + # Forgejo callbacks (all from forgejo_http.ForgejoCallbacks). + get_pr_details: GetPRDetailsCallback + # get_issue_state hits /issues/{n} (not /pulls/{n}) — required for + # milestone_assignment check #5 so plain issues (non-PR) return their + # milestone correctly. (Round-2 principal-review fix replacing the + # dead list_issues field.) + get_issue_state: GetIssueStateCallback + get_labels: GetLabelsCallback + add_label: AddLabelCallback + remove_label: RemoveLabelCallback + patch_pr_state: PatchPRStateCallback + patch_pr_body: PatchPRBodyCallback + patch_pr_milestone: PatchPRMilestoneCallback + + # Per-check enable flags. + completed_not_closed_enabled: bool = False + closing_keyword_fixup_enabled: bool = False + label_sync_enabled: bool = False + state_label_inference_enabled: bool = False + milestone_assignment_enabled: bool = False + + # Shared safety toggle (sources from CONTROLLER_GROOMING_DRY_RUN + # for symmetry with Gates 1/2/3). + dry_run: bool = False + + +# ─── label vocabularies ─────────────────────────────────────────────── + + +# Labels copied from issue to PR by check #3 (label_sync_from_issue). +# Operators can extend by editing this regex. +_LABEL_SYNC_REGEX = re.compile(r"^(Priority|Type|MoSCoW)/.+$") + +# Controller state → Forgejo State/* label. Used by check #4 +# (state_label_inference). The reverse — Forgejo label → state — is +# NOT used; the controller drives, Forgejo reflects. +_STATE_TO_LABEL: dict[str, str] = { + "DISCOVERED": "State/Discovered", + "GROOMING": "State/Grooming", + "ANALYZING": "State/Analyzing", + "IMPLEMENTING": "State/Implementing", + "AWAITING_CI": "State/Awaiting-CI", + "CONFLICT_RESOLVING": "State/Conflict-Resolving", + "ESCALATING": "State/Escalating", + "REVIEWING": "State/Reviewing", + "APPROVED": "State/Approved", + "MERGING": "State/Merging", + "PAUSED": "State/Paused", + "OPERATOR_ATTENTION": "State/Operator-Attention", + "MERGED": "State/Merged", + "ABANDONED": "State/Abandoned", + "STUCK": "State/Stuck", + "CREATED_PR": "State/Created-PR", +} + + +# Bare ``#N`` reference matcher used by check #2 (closing_keyword_fixup). +# Conservative: negative lookbehind excludes: +# - ``\w`` — `foo#42` cross-references (Python `#42` comments, etc.) +# - ``/`` — `owner/repo#42` cross-repo refs +# - ``[`` — markdown links `[#42](url)` (round-1 fix) +# - ``(`` — markdown link destinations `[text](#42)` or HTML hrefs +# `href="(#42)"` — without this the URL-fragment `(#42)` was +# misread as a bare ref and the fixup would PATCH a spurious +# ``Closes #42`` into PR bodies that merely anchor-link to a +# section heading. (Round-2 principal-review fix.) +# The closing-keyword extractor handles ``Closes #N`` etc. separately. +_BARE_REF_RE = re.compile(r"(? None: + """When a workflow reaches MERGED, close any open issue the PR + body links via ``Closes #N`` / ``Fixes #N`` / ``Resolves #N``. + + Candidate selection: all PR-kind workflows in MERGED that haven't + been checked yet (no executed=1 audit row). For each: fetch PR + body, extract closing refs, fetch each ref's issue state, close + any that are still open. + """ + if not callbacks.completed_not_closed_enabled: + return + with session_scope(engine) as session: + rows = _candidates_for_check( + session, "completed_not_closed", state_filter="merged", + ) + + for row in rows: + wf_id = row.workflow_id + try: + with session_scope(engine) as session: + _process_completed_not_closed( + session=session, + workflow_id=wf_id, + owner=row.owner, + repo=row.repo, + pr_number=row.entity_number, + callbacks=callbacks, + report=report, + ) + except Exception as exc: # noqa: BLE001 + logger.exception( + "completed_not_closed: workflow_id=%s raised; continuing", wf_id, + ) + report.errors.append(("completed_not_closed", wf_id, str(exc))) + + +def _process_completed_not_closed( + *, + session, + workflow_id: int, + owner: str, + repo: str, + pr_number: int, + callbacks: MetadataHygieneCallbacks, + report: MetadataHygieneReport, +) -> None: + pr = callbacks.get_pr_details(owner, repo, pr_number) + if not pr: + # PR vanished from Forgejo (404 / deletion). PERMANENT skip — + # PRs don't come back from the dead. + _record_audit( + session, + workflow_id=workflow_id, + check_name="completed_not_closed", + reason_category="pr_not_found", + executed=True, + ) + report.completed_not_closed_skipped += 1 + return + body = pr.get("body") if isinstance(pr, dict) else None + refs = extract_closes_refs(body) + if not refs: + # No linked issues — nothing to close. PERMANENT skip: at + # MERGED state the PR body is stable for closing-keyword + # purposes (operators don't add new Closes refs post-merge — + # if they do, the closing-keyword-fixup check handles them). + _record_audit( + session, + workflow_id=workflow_id, + check_name="completed_not_closed", + reason_category="no_linked_issues", + executed=True, + ) + report.completed_not_closed_skipped += 1 + return + + # Dry-run: log each issue we WOULD close, write one audit row, and + # return early — no Forgejo writes, no increment of _executed counter. + # (Round-2 architect-review fix: pre-fix the dry-run path ran through + # the same loop, inflated ``closed_count``, and incremented + # ``completed_not_closed_executed`` even though no write occurred.) + if callbacks.dry_run: + for issue_n in sorted(refs): + logger.info( + "completed_not_closed[dry-run]: wf=%s would close %s/%s#%d", + workflow_id, owner, repo, issue_n, + ) + _record_audit( + session, + workflow_id=workflow_id, + check_name="completed_not_closed", + reason_category=f"would_close_{len(refs)}_of_{len(refs)}", + executed=False, + llm_reasoning=f"refs={sorted(refs)}", + ) + report.completed_not_closed_skipped += 1 + return + + closed_count = 0 + for issue_n in sorted(refs): + try: + resp = callbacks.patch_pr_state(owner, repo, issue_n, "closed") + except Exception as exc: # noqa: BLE001 + report.errors.append( + ("completed_not_closed", workflow_id, + f"patch_pr_state(#{issue_n}) raised: {exc}") + ) + continue + status = int((resp or {}).get("status") or 0) + category, err = _classify_forgejo_status(status) + if category in ("ok", "ok-no-op"): + closed_count += 1 + else: + report.errors.append( + ("completed_not_closed", workflow_id, + f"patch_pr_state(#{issue_n}) status={status}: {err}") + ) + + # Only write executed=True when ALL linked issues closed successfully. + # Partial success (e.g. some issues hit 429) writes executed=False so + # the workflow stays in candidates and the tick retries on the next + # sweep. Already-closed issues return 200/ok-no-op and do not + # double-close. (Round-2 architect-review fix: pre-fix always wrote + # executed=True, permanently locking out unprocessed issues.) + all_closed = closed_count == len(refs) + _record_audit( + session, + workflow_id=workflow_id, + check_name="completed_not_closed", + reason_category=f"closed_{closed_count}_of_{len(refs)}", + executed=all_closed, + llm_reasoning=f"refs={sorted(refs)}", + ) + report.completed_not_closed_executed += 1 + + +# ─── check #2: closing_keyword_fixup ────────────────────────────────── + + +def run_closing_keyword_fixup( + *, + engine: Engine, + callbacks: MetadataHygieneCallbacks, + report: MetadataHygieneReport, +) -> None: + """When the PR body has a bare ``#N`` reference but no closing + verb, PATCH the body to add ``Closes #N``. + + Conservative: only fires when there's exactly ONE bare ref AND it + isn't already covered by a closing keyword. The ``one bare ref`` + rule is the safety perimeter — multiple bare refs is ambiguous + and a human should disambiguate. The ``not already covered`` rule + prevents double-fixups when ``Closes #N`` already exists but the + PR body redundantly mentions ``#N`` elsewhere. + + Runs on any non-terminal workflow that hasn't been fixed up yet. + """ + if not callbacks.closing_keyword_fixup_enabled: + return + with session_scope(engine) as session: + rows = _candidates_for_check( + session, "closing_keyword_fixup", state_filter="non_terminal", + ) + + for row in rows: + wf_id = row.workflow_id + try: + with session_scope(engine) as session: + _process_closing_keyword_fixup( + session=session, + workflow_id=wf_id, + owner=row.owner, + repo=row.repo, + pr_number=row.entity_number, + callbacks=callbacks, + report=report, + ) + except Exception as exc: # noqa: BLE001 + logger.exception( + "closing_keyword_fixup: workflow_id=%s raised; continuing", wf_id, + ) + report.errors.append(("closing_keyword_fixup", wf_id, str(exc))) + + +def _process_closing_keyword_fixup( + *, + session, + workflow_id: int, + owner: str, + repo: str, + pr_number: int, + callbacks: MetadataHygieneCallbacks, + report: MetadataHygieneReport, +) -> None: + pr = callbacks.get_pr_details(owner, repo, pr_number) + if not pr: + # PERMANENT skip — PR vanished from Forgejo. + _record_audit( + session, workflow_id=workflow_id, + check_name="closing_keyword_fixup", + reason_category="pr_not_found", executed=True, + ) + report.closing_keyword_fixup_skipped += 1 + return + body = (pr.get("body") if isinstance(pr, dict) else None) or "" + existing = extract_closes_refs(body) + bare_refs = {int(m.group(1)) for m in _BARE_REF_RE.finditer(body)} + # Only candidate refs that AREN'T already covered by a closing kw. + candidates = bare_refs - existing + if len(candidates) != 1: + # Zero or multiple bare refs — too ambiguous to auto-fix. + # TRANSIENT skip (executed=False): the PR body can be edited + # later to add the single bare ref this check needs. Round-1 + # architect-review fix: pre-fix we wrote executed=True here, + # which permanently locked future re-evaluation. + _record_audit( + session, workflow_id=workflow_id, + check_name="closing_keyword_fixup", + reason_category=( + "no_bare_ref" if not candidates else "multiple_bare_refs" + ), + executed=False, + llm_reasoning=f"bare={sorted(bare_refs)} existing={sorted(existing)}", + ) + report.closing_keyword_fixup_skipped += 1 + return + + issue_n = next(iter(candidates)) + new_body = body.rstrip() + f"\n\nCloses #{issue_n}\n" + + if callbacks.dry_run: + logger.info( + "closing_keyword_fixup[dry-run]: wf=%s would add 'Closes #%d' to PR body", + workflow_id, issue_n, + ) + _record_audit( + session, workflow_id=workflow_id, + check_name="closing_keyword_fixup", + reason_category=f"would_add_closes_{issue_n}", + executed=False, + ) + report.closing_keyword_fixup_skipped += 1 + return + + try: + resp = callbacks.patch_pr_body(owner, repo, pr_number, new_body) + except Exception as exc: # noqa: BLE001 + report.errors.append( + ("closing_keyword_fixup", workflow_id, f"patch_pr_body raised: {exc}") + ) + return + status = int((resp or {}).get("status") or 0) + category, err = _classify_forgejo_status(status) + if category in ("ok", "ok-no-op"): + # Persist the pre-mutation body to the audit row (round-1 + # architect-review fix). PATCH-ing user content needs a + # recovery path: if the fixup ever corrupts a body (encoding + # bug, escape mismatch, race with an operator edit), the + # original is recoverable from grooming_decisions.llm_reasoning. + # The body is capped at ~64KB by Forgejo so this fits inside + # llm_reasoning's TEXT column without trimming. + _record_audit( + session, workflow_id=workflow_id, + check_name="closing_keyword_fixup", + reason_category=f"added_closes_{issue_n}", + executed=True, + llm_reasoning=f"pre_body={body!r}", + target_workflow_id=None, + ) + report.closing_keyword_fixup_executed += 1 + else: + report.errors.append( + ("closing_keyword_fixup", workflow_id, + f"patch_pr_body status={status}: {err}") + ) + + +# ─── check #3: label_sync_from_issue ────────────────────────────────── + + +def run_label_sync_from_issue( + *, + engine: Engine, + callbacks: MetadataHygieneCallbacks, + report: MetadataHygieneReport, +) -> None: + """At pickup, copy ``Priority/*`` / ``Type/*`` / ``MoSCoW/*`` + labels from the linked issue to the PR. + + Only PRs with at least one ``Closes #N`` ref are candidates; + standalone PRs are no-op. Idempotency via per-workflow audit row. + """ + if not callbacks.label_sync_enabled: + return + with session_scope(engine) as session: + rows = _candidates_for_check( + session, "label_sync_from_issue", state_filter="non_terminal", + ) + + for row in rows: + wf_id = row.workflow_id + try: + with session_scope(engine) as session: + _process_label_sync( + session=session, + workflow_id=wf_id, + owner=row.owner, + repo=row.repo, + pr_number=row.entity_number, + callbacks=callbacks, + report=report, + ) + except Exception as exc: # noqa: BLE001 + logger.exception( + "label_sync_from_issue: workflow_id=%s raised; continuing", wf_id, + ) + report.errors.append(("label_sync_from_issue", wf_id, str(exc))) + + +def _process_label_sync( + *, + session, + workflow_id: int, + owner: str, + repo: str, + pr_number: int, + callbacks: MetadataHygieneCallbacks, + report: MetadataHygieneReport, +) -> None: + pr = callbacks.get_pr_details(owner, repo, pr_number) + if not pr: + # PERMANENT skip — PR vanished from Forgejo. + _record_audit( + session, workflow_id=workflow_id, + check_name="label_sync_from_issue", + reason_category="pr_not_found", executed=True, + ) + report.label_sync_skipped += 1 + return + refs = extract_closes_refs((pr.get("body") if isinstance(pr, dict) else None)) + if not refs: + # TRANSIENT skip — the operator might add a `Closes #N` later. + # Round-1 architect-review fix: pre-fix this wrote executed=True + # which permanently locked future re-evaluation if the PR body + # gained a `Closes #N` later. + _record_audit( + session, workflow_id=workflow_id, + check_name="label_sync_from_issue", + reason_category="no_linked_issues", executed=False, + ) + report.label_sync_skipped += 1 + return + + # For each linked issue, fetch its labels + collect Priority/Type/MoSCoW ones. + to_add: set[str] = set() + for issue_n in sorted(refs): + try: + issue_labels = callbacks.get_labels(owner, repo, issue_n) or [] + except Exception as exc: # noqa: BLE001 + report.errors.append( + ("label_sync_from_issue", workflow_id, + f"get_labels(issue #{issue_n}) raised: {exc}") + ) + continue + for lbl in issue_labels: + name = (lbl or {}).get("name") if isinstance(lbl, dict) else None + if name and _LABEL_SYNC_REGEX.match(name): + to_add.add(name) + + if not to_add: + # TRANSIENT skip — operator might add a Priority/Type/MoSCoW + # label to the linked issue later. + _record_audit( + session, workflow_id=workflow_id, + check_name="label_sync_from_issue", + reason_category="no_matching_labels", executed=False, + ) + report.label_sync_skipped += 1 + return + + if callbacks.dry_run: + logger.info( + "label_sync_from_issue[dry-run]: wf=%s would add labels %s", + workflow_id, sorted(to_add), + ) + _record_audit( + session, workflow_id=workflow_id, + check_name="label_sync_from_issue", + reason_category=f"would_add_{len(to_add)}", + executed=False, + llm_reasoning=f"labels={sorted(to_add)}", + ) + report.label_sync_skipped += 1 + return + + results = adjust_labels( + owner=owner, repo=repo, pr_number=pr_number, + add=sorted(to_add), remove=[], + get_labels=callbacks.get_labels, + add_label=callbacks.add_label, + remove_label=callbacks.remove_label, + ) + failed = [r for r in results if r.action == "failed"] + if failed: + err_detail = "; ".join(f"{r.label}:{r.error}" for r in failed) + report.errors.append(("label_sync_from_issue", workflow_id, err_detail)) + # Transient failure — write executed=False so the workflow stays + # in candidates and retries next tick. The audit row provides + # an observability trail for persistent failures (e.g. a label + # id deleted on Forgejo). (Round-2 architect-review fix: pre-fix + # wrote no audit row at all, giving operators no log-grep anchor + # when the label never appeared.) + _record_audit( + session, workflow_id=workflow_id, + check_name="label_sync_from_issue", + reason_category="label_apply_failed", + executed=False, + llm_reasoning=err_detail, + ) + return + _record_audit( + session, workflow_id=workflow_id, + check_name="label_sync_from_issue", + reason_category=f"added_{sum(1 for r in results if r.action == 'added')}", + executed=True, + llm_reasoning=f"labels={sorted(to_add)}", + ) + report.label_sync_executed += 1 + + +# ─── check #4: state_label_inference ────────────────────────────────── + + +def run_state_label_inference( + *, + engine: Engine, + callbacks: MetadataHygieneCallbacks, + report: MetadataHygieneReport, +) -> None: + """Periodically sync a ``State/*`` label to mirror the workflow's + ``current_state``. Runs on every non-terminal workflow on every + iteration; cheap because most workflows haven't changed state since + the last sweep (the ``_last_synced_state`` audit row trims those). + """ + if not callbacks.state_label_inference_enabled: + return + # state_label_inference does NOT use the executed=1 LEFT JOIN + # idempotency pattern — its idempotency is per-(workflow, state), + # checked via ``_last_synced_state`` matching the workflow's + # current ``current_state``. So we always scan every PR-kind + # workflow regardless of past audit rows; the per-row state-match + # guard short-circuits when nothing changed. + with session_scope(engine) as session: + # Fetch the list of PR-kind workflow IDs to process this sweep. + # We do NOT fetch current_state here — we re-read it inside the + # per-workflow inner session to avoid acting on a stale snapshot. + # (Round-2 architect-review fix: the state machine can transition + # a workflow between the outer SELECT and the inner session, + # causing a spurious label write with the old target_label.) + rows = session.execute( + text( + """ + SELECT workflow_id, owner, repo, entity_number + FROM workflows + WHERE kind = 'pr' + """ + ) + ).fetchall() + + for row in rows: + wf_id = row.workflow_id + try: + with session_scope(engine) as session: + # Re-read current_state from DB to get the current truth. + fresh = session.execute( + text( + "SELECT current_state FROM workflows " + "WHERE workflow_id = :wf_id" + ), + {"wf_id": wf_id}, + ).first() + if not fresh: + # Workflow was deleted between the outer SELECT and + # this inner session — skip. + continue + state = fresh.current_state + target_label = _STATE_TO_LABEL.get(state) + if target_label is None: + # Unknown state — skip rather than guess a label name. + continue + last = _last_synced_state(session, wf_id) + if last == state: + report.state_label_skipped += 1 + continue + _process_state_label_sync( + session=session, + workflow_id=wf_id, + owner=row.owner, + repo=row.repo, + pr_number=row.entity_number, + state=state, + target_label=target_label, + callbacks=callbacks, + report=report, + ) + except Exception as exc: # noqa: BLE001 + logger.exception( + "state_label_inference: workflow_id=%s raised; continuing", wf_id, + ) + report.errors.append(("state_label_inference", wf_id, str(exc))) + + +def _process_state_label_sync( + *, + session, + workflow_id: int, + owner: str, + repo: str, + pr_number: int, + state: str, + target_label: str, + callbacks: MetadataHygieneCallbacks, + report: MetadataHygieneReport, +) -> None: + # Compute add + remove sets: add target_label, remove any OTHER + # State/* labels (so a stale State/Implementing doesn't co-exist + # with the new State/Reviewing). + to_remove = [lbl for lbl in _STATE_TO_LABEL.values() if lbl != target_label] + + if callbacks.dry_run: + logger.info( + "state_label_inference[dry-run]: wf=%s would sync %s → %s", + workflow_id, state, target_label, + ) + _record_audit( + session, workflow_id=workflow_id, + check_name="state_label_inference", + reason_category=state, executed=False, + ) + report.state_label_skipped += 1 + return + + results = adjust_labels( + owner=owner, repo=repo, pr_number=pr_number, + add=[target_label], remove=to_remove, + get_labels=callbacks.get_labels, + add_label=callbacks.add_label, + remove_label=callbacks.remove_label, + ) + failed = [r for r in results if r.action == "failed"] + if failed: + report.errors.append( + ("state_label_inference", workflow_id, + "; ".join(f"{r.label}:{r.error}" for r in failed)) + ) + return + _record_audit( + session, workflow_id=workflow_id, + check_name="state_label_inference", + reason_category=state, executed=True, + ) + report.state_label_executed += 1 + + +# ─── check #5: milestone_assignment ─────────────────────────────────── + + +def run_milestone_assignment( + *, + engine: Engine, + callbacks: MetadataHygieneCallbacks, + report: MetadataHygieneReport, +) -> None: + """At pickup, copy the milestone from the linked issue to the PR + if the PR has none. + + Idempotency: per-workflow audit row. Never overrides an explicit + milestone — only sets when the PR's current milestone is None. + """ + if not callbacks.milestone_assignment_enabled: + return + with session_scope(engine) as session: + rows = _candidates_for_check( + session, "milestone_assignment", state_filter="non_terminal", + ) + + for row in rows: + wf_id = row.workflow_id + try: + with session_scope(engine) as session: + _process_milestone_assignment( + session=session, + workflow_id=wf_id, + owner=row.owner, + repo=row.repo, + pr_number=row.entity_number, + callbacks=callbacks, + report=report, + ) + except Exception as exc: # noqa: BLE001 + logger.exception( + "milestone_assignment: workflow_id=%s raised; continuing", wf_id, + ) + report.errors.append(("milestone_assignment", wf_id, str(exc))) + + +def _milestone_id(milestone: Any) -> int | None: + """Forgejo PR/issue ``milestone`` field is either None or a dict + with an ``id`` field. Return the int id or None.""" + if isinstance(milestone, dict): + mid = milestone.get("id") + if isinstance(mid, int): + return mid + return None + + +def _process_milestone_assignment( + *, + session, + workflow_id: int, + owner: str, + repo: str, + pr_number: int, + callbacks: MetadataHygieneCallbacks, + report: MetadataHygieneReport, +) -> None: + pr = callbacks.get_pr_details(owner, repo, pr_number) + if not pr: + # PERMANENT skip — PR vanished from Forgejo. + _record_audit( + session, workflow_id=workflow_id, + check_name="milestone_assignment", + reason_category="pr_not_found", executed=True, + ) + report.milestone_skipped += 1 + return + # Never override an explicit milestone. + if _milestone_id(pr.get("milestone") if isinstance(pr, dict) else None): + # PERMANENT skip — once the PR has a milestone, the check is + # done. We never replace an explicit assignment. + _record_audit( + session, workflow_id=workflow_id, + check_name="milestone_assignment", + reason_category="pr_already_has_milestone", executed=True, + ) + report.milestone_skipped += 1 + return + refs = extract_closes_refs((pr.get("body") if isinstance(pr, dict) else None)) + if not refs: + # TRANSIENT skip — operator might add a `Closes #N` later. + _record_audit( + session, workflow_id=workflow_id, + check_name="milestone_assignment", + reason_category="no_linked_issues", executed=False, + ) + report.milestone_skipped += 1 + return + # Use the FIRST linked issue's milestone (sorted for determinism). + # If multiple linked issues have different milestones, taking the + # smallest issue # is the conservative choice — the same one + # GitHub/Forgejo close-on-merge would auto-link to first. + # + # IMPORTANT: use get_issue_state (hits /issues/{n}), NOT + # get_pr_details (hits /pulls/{n}). Plain issues return 404 from + # the /pulls/ endpoint because Forgejo filters that endpoint to PRs + # only; get_pr_details would silently return None and the loop would + # exhaust all refs producing no_issue_milestone even when the issue + # has a milestone. (Round-2 principal-review fix.) + target_mid: int | None = None + chosen_issue: int | None = None + for issue_n in sorted(refs): + issue_pr = callbacks.get_issue_state(owner, repo, issue_n) + if not isinstance(issue_pr, dict): + continue + mid = _milestone_id(issue_pr.get("milestone")) + if mid: + target_mid = mid + chosen_issue = issue_n + break + if not target_mid: + # TRANSIENT skip — the linked issue might gain a milestone + # later. Round-1 architect-review fix: pre-fix executed=True + # permanently locked future re-evaluation. + _record_audit( + session, workflow_id=workflow_id, + check_name="milestone_assignment", + reason_category="no_issue_milestone", executed=False, + ) + report.milestone_skipped += 1 + return + + if callbacks.dry_run: + logger.info( + "milestone_assignment[dry-run]: wf=%s would set milestone %d " + "(from issue #%d)", + workflow_id, target_mid, chosen_issue, + ) + _record_audit( + session, workflow_id=workflow_id, + check_name="milestone_assignment", + reason_category=f"would_set_milestone_{target_mid}", + executed=False, + ) + report.milestone_skipped += 1 + return + + try: + resp = callbacks.patch_pr_milestone(owner, repo, pr_number, target_mid) + except Exception as exc: # noqa: BLE001 + report.errors.append( + ("milestone_assignment", workflow_id, + f"patch_pr_milestone raised: {exc}") + ) + return + status = int((resp or {}).get("status") or 0) + category, err = _classify_forgejo_status(status) + if category in ("ok", "ok-no-op"): + _record_audit( + session, workflow_id=workflow_id, + check_name="milestone_assignment", + reason_category=f"set_milestone_{target_mid}", + executed=True, + llm_reasoning=f"from_issue=#{chosen_issue}", + ) + report.milestone_executed += 1 + else: + report.errors.append( + ("milestone_assignment", workflow_id, + f"patch_pr_milestone status={status}: {err}") + ) + + +# ─── dispatcher ─────────────────────────────────────────────────────── + + +def run_metadata_hygiene_tick( + *, + engine: Engine, + callbacks: MetadataHygieneCallbacks, +) -> MetadataHygieneReport: + """One sweep: invoke each enabled metadata-hygiene check in turn. + + Order doesn't matter (checks are independent + idempotent), but + pickup-time checks (#2, #3, #5) are grouped before the periodic + (#4) for log-readability. Completed-not-closed (#1) runs first + because it's the only one keyed off MERGED workflows. + """ + report = MetadataHygieneReport() + + run_completed_not_closed(engine=engine, callbacks=callbacks, report=report) + run_closing_keyword_fixup(engine=engine, callbacks=callbacks, report=report) + run_label_sync_from_issue(engine=engine, callbacks=callbacks, report=report) + run_milestone_assignment(engine=engine, callbacks=callbacks, report=report) + run_state_label_inference(engine=engine, callbacks=callbacks, report=report) + + if ( + report.completed_not_closed_executed + or report.closing_keyword_fixup_executed + or report.label_sync_executed + or report.state_label_executed + or report.milestone_executed + or report.errors + ): + logger.info( + "metadata_hygiene: completed_not_closed=%d closing_kw_fixup=%d " + "label_sync=%d state_label=%d milestone=%d errors=%d", + report.completed_not_closed_executed, + report.closing_keyword_fixup_executed, + report.label_sync_executed, + report.state_label_executed, + report.milestone_executed, + len(report.errors), + ) + return report + + +__all__ = [ + "MetadataHygieneCallbacks", + "MetadataHygieneReport", + "run_metadata_hygiene_tick", + "run_completed_not_closed", + "run_closing_keyword_fixup", + "run_label_sync_from_issue", + "run_state_label_inference", + "run_milestone_assignment", +] diff --git a/tools/controller/master/metadata_hygiene_config.py b/tools/controller/master/metadata_hygiene_config.py new file mode 100644 index 000000000..112535a13 --- /dev/null +++ b/tools/controller/master/metadata_hygiene_config.py @@ -0,0 +1,124 @@ +"""Phase 4 — metadata-hygiene configuration. + +Phase 4 adds five deterministic, idempotent metadata-hygiene checks +the controller can perform on Forgejo PRs / issues. Each check is +behind an INDIVIDUAL feature flag (so operators can roll out one at +a time + revert a problem check without disabling the whole phase) +PLUS a master enable flag (so operators can disable all five during +an incident without touching individual flags). + +Standard rollout sequence per check: +- master ENABLED=false (default) — none of the five run; zero risk. +- master ENABLED=true + per-check ENABLED=false + DRY_RUN=true — the + tick computes what it WOULD do, writes an audit row with + ``executed=0``, and emits a log line. No Forgejo writes. +- master ENABLED=true + per-check ENABLED=true + DRY_RUN=false — + full path active. + +The kill switches matter because each check writes to Forgejo: +- completed-not-closed: closes a (likely operator-watched) issue +- closing-keyword fixup: mutates the PR body +- label sync: adds labels (operator-visible) +- state-label inference: adds/removes labels on every iteration +- milestone assignment: assigns a milestone + +A mis-fire of any of these is operator-visible noise at best, +operator-visible damage at worst. Default-off + per-check granular +control is the right safety posture. + +Pairs with ``CONTROLLER_GROOMING_DRY_RUN`` (which metadata-hygiene +shares as its safe-rollout layer for symmetry with Gate 1/2/3) + +this module's own ``CONTROLLER_METADATA_HYGIENE_*`` toggles. +""" + +from __future__ import annotations + +import json +import logging +import os +from dataclasses import asdict, dataclass + +logger = logging.getLogger(__name__) + + +@dataclass(frozen=True) +class MetadataHygieneConfig: + """Frozen config for the Phase 4 metadata-hygiene checks. + + Each flag defaults to False so a fresh deploy of Phase 4 code is + audit-only until the operator explicitly enables each path. The + master ``enabled`` flag is the kill switch — when False, the loop + skips the entire metadata-hygiene block regardless of per-check + flags. + """ + + # Master switch. When False, ``__main__.py`` passes + # ``metadata_hygiene_callbacks=None`` to ``master_main_loop``; + # all five checks are skipped. + enabled: bool = False + + # Per-check feature flags. Each can be flipped independently of + # the others (e.g. ship label-sync first, hold the closing- + # keyword-fixup until operators have watched a few days of audit + # rows). + completed_not_closed: bool = False + closing_keyword_fixup: bool = False + label_sync_from_issue: bool = False + state_label_inference: bool = False + milestone_assignment: bool = False + + +def _bool(env_name: str, default: bool) -> bool: + raw = os.environ.get(env_name) + if raw is None: + return default + return raw.strip().lower() in {"1", "true", "yes", "on"} + + +def get_metadata_hygiene_config() -> MetadataHygieneConfig: + """Read the effective Phase 4 config from environment variables. + + Pure function; safe to call repeatedly. Operators typically log + the result via ``log_effective_config`` at controller startup so + incident-response can recover "what was set when this happened?" + from the log stream. + """ + return MetadataHygieneConfig( + enabled=_bool("CONTROLLER_METADATA_HYGIENE_ENABLED", False), + completed_not_closed=_bool( + "CONTROLLER_METADATA_HYGIENE_COMPLETED_NOT_CLOSED", False + ), + closing_keyword_fixup=_bool( + "CONTROLLER_METADATA_HYGIENE_CLOSING_KEYWORD_FIXUP", False + ), + label_sync_from_issue=_bool( + "CONTROLLER_METADATA_HYGIENE_LABEL_SYNC", False + ), + state_label_inference=_bool( + "CONTROLLER_METADATA_HYGIENE_STATE_LABEL_INFERENCE", False + ), + milestone_assignment=_bool( + "CONTROLLER_METADATA_HYGIENE_MILESTONE_ASSIGNMENT", False + ), + ) + + +def log_effective_config() -> None: + """Emit one INFO-level line with the effective Phase 4 config. + + Called from controller startup so operators have a log-grep anchor + for "what was set when this incident happened?" without needing to + reconstruct env-var state from systemd / shell history. + """ + cfg = get_metadata_hygiene_config() + logger.info( + "metadata_hygiene config: %s", + json.dumps(asdict(cfg), sort_keys=True), + ) + + +__all__ = [ + "MetadataHygieneConfig", + "get_metadata_hygiene_config", + "log_effective_config", +]