From 8209424a1fe8ffbdc739af82ba67b24185e8b57a Mon Sep 17 00:00:00 2001 From: Tingkai Ying Date: Fri, 25 Sep 2026 01:38:52 +0800 Subject: [PATCH 1/2] fix(agent-team): mark finalize-gate bypass on the final turn instead of delivering silently BareTextFinalizeObserver accepts the answer at turn >= max_turns - 1 even when finalize_gate blocks it (the never-lose-the-answer fallback) but latched it with no trace, so a run with open task-board items was delivered exactly like a clean success. Set metadata ["finalize_gate_bypassed"] and append a user-visible note naming the still-unresolved items. Mid-run rejection and clean-pass behavior are unchanged and locked by tests. Related to #19 --- CHANGELOG.md | 3 + tests/test_bare_text_finalize_gate.py | 103 ++++++++++++++++++ .../observers/bare_text_finalize.py | 30 +++++ 3 files changed, 136 insertions(+) create mode 100644 tests/test_bare_text_finalize_gate.py diff --git a/CHANGELOG.md b/CHANGELOG.md index a4ae801..5115e93 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -36,5 +36,8 @@ Initial open-source release of FrontierAgent. ### Fixed +- Surface finalize-gate bypasses on the final turn: an answer delivered despite + open task-board items now carries an unfinished-work note and a + `finalize_gate_bypassed` marker instead of reading as a clean success. - Apply benchmark question limits after seeded shuffling so repeated runs can sample different questions while `--no-shuffle` keeps canonical ordering. diff --git a/tests/test_bare_text_finalize_gate.py b/tests/test_bare_text_finalize_gate.py new file mode 100644 index 0000000..886fc5d --- /dev/null +++ b/tests/test_bare_text_finalize_gate.py @@ -0,0 +1,103 @@ +from __future__ import annotations + +import asyncio + +from frontier_agent.core.loop_types import TurnContext, notify_observers +from plugins.tools import task_board as tb +from workflows.agent_team.observers.bare_text_finalize import ( + BareTextFinalizeObserver, +) + +_TASK = "task" +_ANSWER = "Deployment complete; 100% of attacks blocked." + + +def _seed_board(resolutions: dict[str, str]) -> None: + tb._BOARDS[_TASK] = { + "seq": len(resolutions), + "tasks": { + tid: { + "description": f"work {tid}", + "resolution": resolution, + "owners": [], + } + for tid, resolution in resolutions.items() + }, + } + + +def _context(turn: int, max_turns: int = 20) -> TurnContext: + return TurnContext( + turn=turn, + max_turns=max_turns, + task_id=_TASK, + role_id="coordinator", + ai_text=_ANSWER, + thinking="", + tool_calls=[], + messages=[], + usage=None, + metadata={}, + ) + + +def _dispatch(ctx: TurnContext) -> list: + try: + return asyncio.run( + notify_observers( + [BareTextFinalizeObserver()], + "on_llm_response", + ctx, + ), + ) + finally: + tb._BOARDS.pop(_TASK, None) + + +def test_mid_run_unfinished_board_still_blocks() -> None: + """Away from the last turn the gate wins: rejected, never latched.""" + _seed_board({"t1": "open"}) + ctx = _context(turn=5) + + interventions = _dispatch(ctx) + + (intervention,) = interventions + assert intervention.continue_to_next_turn is True + assert intervention.stop_reason is None + assert "final_answer" not in ctx.metadata + + +def test_final_turn_bypass_marks_unfinished_work() -> None: + """The last-turn bypass must deliver WITH a visible unfinished-work note. + + ``finalize_gate`` blocks while board items are open, but the final turn + accepts the answer anyway (never lose it to max_turns). That bypass used + to be silent — an unfinished run looked like a clean success. It must + latch the answer, flag ``finalize_gate_bypassed``, and append a note + naming the still-open items (issue #19). + """ + _seed_board({"t1": "open", "t2": "open", "t3": "open"}) + ctx = _context(turn=19) + + interventions = _dispatch(ctx) + + (intervention,) = interventions + assert intervention.stop_reason == "final_answer" + latched = ctx.metadata["final_answer"] + assert "unfinished" in str(latched).casefold() + for task_id in ("t1", "t2", "t3"): + assert task_id in str(latched) + assert ctx.metadata.get("finalize_gate_bypassed") + + +def test_final_turn_clean_board_latches_untouched() -> None: + """Gate passes → the answer is latched verbatim, no note, no marker.""" + _seed_board({"t1": "resolved", "t2": "cancelled"}) + ctx = _context(turn=19) + + interventions = _dispatch(ctx) + + (intervention,) = interventions + assert intervention.stop_reason == "final_answer" + assert ctx.metadata["final_answer"] == _ANSWER + assert "finalize_gate_bypassed" not in ctx.metadata diff --git a/workflows/agent_team/observers/bare_text_finalize.py b/workflows/agent_team/observers/bare_text_finalize.py index 2164989..9e7ad10 100644 --- a/workflows/agent_team/observers/bare_text_finalize.py +++ b/workflows/agent_team/observers/bare_text_finalize.py @@ -7,10 +7,32 @@ from frontier_agent.core.loop_types import BaseObserver, Intervention, TurnContext from plugins.tools._bus_scope import resolve_bus_task_id from plugins.tools.finalize_answer import finalize_gate +from plugins.tools.task_board import unresolved_task_ids logger = logging.getLogger(__name__) +def _unfinished_note(task_id: str, text: str) -> str: + """Footnote an answer that is being delivered despite a blocked gate. + + The last turn must not lose the answer to ``max_turns`` (see the bypass + below), but an unfinished run must never read as a clean success: name the + board items that are still unresolved so the reader can tell which parts + of the answer were never corroborated. + """ + pending = unresolved_task_ids(task_id) + if pending: + what = f"task-board item(s) still unfinished: {', '.join(pending)}" + else: + what = "the finalize gate was still rejecting this submission" + return ( + f"{text}\n\n---\n\n" + f"> ⚠ Unfinished work at submission: {what}. This answer was " + "delivered on the final turn despite the gate — conclusions that " + "depend on that work are unverified." + ) + + class BareTextFinalizeObserver(BaseObserver): critical = True @@ -36,6 +58,14 @@ async def on_llm_response(self, ctx: TurnContext) -> Intervention | None: return Intervention(continue_to_next_turn=True, inject_messages=[err]) if isinstance(ctx.metadata, dict): + if err: + # Last-turn bypass: the gate says BLOCK, but the answer is + # delivered anyway rather than lost to max_turns. Keep that + # fallback — and make it visible (machine-readable marker + # plus a user-visible note), so an unfinished run is never + # presented as a clean success. + ctx.metadata["finalize_gate_bypassed"] = err + text = _unfinished_note(ctx.task_id, text) ctx.metadata["final_answer"] = text ctx.metadata["final_answer_confidence"] = 1.0 logger.info( From 4f4eb5fa74a2136a6993b2b4ab35192cafd09a97 Mon Sep 17 00:00:00 2001 From: Tingkai Ying Date: Fri, 25 Sep 2026 17:56:14 +0800 Subject: [PATCH 2/2] fix(agent-team): carry the bypass marker through workflow output and reporter finalization Review on #48: the marker died in loop-local metadata (consumers only saw answer_status="complete"), and the warning appended after a trailing References section was stripped by the reporter's citation cleanup. - publish finalize_gate_bypassed + finalize_gate_warning from main_agent_node and list them in both specs (main output_fields, reporter include_fields/output_fields) - the observer now stores a ready-to-append warning built while the task board is still live instead of mutating the latched answer - delivery nodes append the warning after finalization: main_agent for the coordinator's own answer, agent_team_reporter after _run_fast_reporter returns (References cleanup has already run) - append_bypass_warning is idempotent - five new tests cover both delivery boundaries and document the strip hazard, including strip_trailing_references' 30% heading guard --- tests/test_bare_text_finalize_gate.py | 148 ++++++++++++++++-- workflows/agent_team/nodes/main_agent.py | 14 ++ workflows/agent_team/nodes/reporter.py | 12 ++ .../observers/bare_text_finalize.py | 42 +++-- workflows/agent_team/spec.py | 4 + workflows/agent_team/spec_report.py | 4 + 6 files changed, 200 insertions(+), 24 deletions(-) diff --git a/tests/test_bare_text_finalize_gate.py b/tests/test_bare_text_finalize_gate.py index 886fc5d..d893292 100644 --- a/tests/test_bare_text_finalize_gate.py +++ b/tests/test_bare_text_finalize_gate.py @@ -1,15 +1,32 @@ from __future__ import annotations import asyncio +from collections.abc import Mapping +from typing import Any + +import pytest from frontier_agent.core.loop_types import TurnContext, notify_observers from plugins.tools import task_board as tb +from workflows._shared.citation_contract import ( + finalize_report_with_canonical_references, +) +from workflows.agent_team.nodes import fast_reporter_v1 +from workflows.agent_team.nodes.reporter import agent_team_reporter +from workflows.agent_team.observers import bare_text_finalize as btf from workflows.agent_team.observers.bare_text_finalize import ( BareTextFinalizeObserver, ) +from workflows.agent_team.spec import SWARM_SPEC +from workflows.agent_team.spec_report import AGENT_TEAM_REPORT_SPEC _TASK = "task" -_ANSWER = "Deployment complete; 100% of attacks blocked." +_ANSWER = "Deployment complete; 100% of attacks blocked [1]." +_REFERENCES = [{"url": "https://example.com/a", "title": "Example"}] +_BYPASS_ERR = ( + "Cannot finish: task board has unresolved item(s) ['t1']. " + "For each, call update_task(...)" +) def _seed_board(resolutions: dict[str, str]) -> None: @@ -67,31 +84,35 @@ def test_mid_run_unfinished_board_still_blocks() -> None: assert "final_answer" not in ctx.metadata -def test_final_turn_bypass_marks_unfinished_work() -> None: - """The last-turn bypass must deliver WITH a visible unfinished-work note. +def test_final_turn_bypass_stores_marker_and_warning() -> None: + """The bypass must leave machine- and human-readable traces in metadata. ``finalize_gate`` blocks while board items are open, but the final turn - accepts the answer anyway (never lose it to max_turns). That bypass used - to be silent — an unfinished run looked like a clean success. It must - latch the answer, flag ``finalize_gate_bypassed``, and append a note - naming the still-open items (issue #19). + accepts the answer anyway (never lose it to max_turns). The observer + stores the gate message (``finalize_gate_bypassed``) and a ready-to-append + warning built while the board is still live + (``finalize_gate_warning``); delivery nodes append the warning after any + finalization that would otherwise strip it (issue #19, review on #48). + The latched answer itself stays clean — appending here would be defeated + by the reporter's References cleanup. """ - _seed_board({"t1": "open", "t2": "open", "t3": "open"}) + _seed_board({"t1": "open", "t2": "open"}) ctx = _context(turn=19) interventions = _dispatch(ctx) (intervention,) = interventions assert intervention.stop_reason == "final_answer" - latched = ctx.metadata["final_answer"] - assert "unfinished" in str(latched).casefold() - for task_id in ("t1", "t2", "t3"): - assert task_id in str(latched) assert ctx.metadata.get("finalize_gate_bypassed") + warning = str(ctx.metadata.get("finalize_gate_warning") or "") + assert "unfinished" in warning.casefold() + assert "t1" in warning and "t2" in warning + # The answer is latched verbatim; the warning travels separately. + assert ctx.metadata["final_answer"] == _ANSWER def test_final_turn_clean_board_latches_untouched() -> None: - """Gate passes → the answer is latched verbatim, no note, no marker.""" + """Gate passes → the answer is latched verbatim, no marker, no warning.""" _seed_board({"t1": "resolved", "t2": "cancelled"}) ctx = _context(turn=19) @@ -101,3 +122,104 @@ def test_final_turn_clean_board_latches_untouched() -> None: assert intervention.stop_reason == "final_answer" assert ctx.metadata["final_answer"] == _ANSWER assert "finalize_gate_bypassed" not in ctx.metadata + assert "finalize_gate_warning" not in ctx.metadata + + +def test_references_finalizer_strips_a_trailing_warning() -> None: + """Documents WHY the warning must be re-attached after finalization. + + ``finalize_report_with_canonical_references`` drops everything from the + ``References`` heading to the end of the body — a warning appended after + that section does not survive the reporter's citation cleanup. + + ``strip_trailing_references`` refuses to cut when the heading starts + before 30% of the body (mid-body sections are legitimate), so the body + must be realistically long — as in a real report, where the reviewer's + reproduction showed the warning being stripped. + """ + body = ( + "Attackers probed hidden paths, enumerated backup archives, and " + "fuzzed administrative endpoints across several weeks of access " + "logs before the intrusion was detected. " * 3 + + "[1]\n\n" + "## References\n\n" + "[1] https://example.com/a\n" + ) + warning = "\n\n---\n\n> ⚠ Unfinished work at submission: task t1." + finalized = finalize_report_with_canonical_references( + body + warning, + references=_REFERENCES, + language="en", + ) + assert "Unfinished" not in finalized + assert "https://example.com/a" in finalized # canonical block re-appended + + +def test_append_bypass_warning_is_idempotent() -> None: + """Delivery boundaries append exactly once, even on repeated calls.""" + warning = "\n\n---\n\n> ⚠ Unfinished work at submission: task t1." + source: Mapping[str, Any] = {"finalize_gate_warning": warning} + + once = btf.append_bypass_warning("Report body.", source) + twice = btf.append_bypass_warning(once, source) + assert once.endswith(warning) + assert twice == once + # No stored warning → text untouched. + assert btf.append_bypass_warning("Report body.", {}) == "Report body." + + +@pytest.mark.asyncio +async def test_reporter_reappends_warning_after_finalization( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """The reporter's output must carry the warning exactly once (review #48). + + ``_run_fast_reporter`` is the post-finalization report: References + cleanup already ran inside the chain, so the node itself is the seam + that re-attaches the stored bypass warning. + """ + finalized_report = ( + "Rewritten report with citation [1].\n\n" + "## References\n\n- [1] https://example.com/a\n" + ) + + async def _stub_reporter( + _state: dict[str, Any], _ctx: Any, + ) -> str: + return finalized_report + + monkeypatch.setattr(fast_reporter_v1, "_run_fast_reporter", _stub_reporter) + warning = "\n\n---\n\n> ⚠ Unfinished work at submission: task t1." + state: dict[str, Any] = { + "reporter_backend": "fast", + "metadata": {}, + "finalize_gate_bypassed": _BYPASS_ERR, + "finalize_gate_warning": warning, + } + + out = await agent_team_reporter(state, None) # type: ignore[arg-type] + + final_answer = str(out.get("final_answer") or "") + assert final_answer.endswith(warning.rstrip("\n")) or warning in final_answer + assert final_answer.count("Unfinished work at submission") == 1 + # The marker must pass through so downstream consumers keep it. + assert out.get("finalize_gate_bypassed") == _BYPASS_ERR + + +def _node(spec: Any, node_id: str) -> Any: + return next(node for node in spec.nodes if node.node_id == node_id) + + +@pytest.mark.parametrize("spec", [SWARM_SPEC, AGENT_TEAM_REPORT_SPEC]) +def test_specs_publish_and_forward_the_bypass_marker(spec: Any) -> None: + """main_agent must publish the marker; the reporter must receive it. + + Without these list entries the marker dies in loop-local metadata and + consumers only ever see ``answer_status="complete"`` (review on #48). + """ + marker_keys = {"finalize_gate_bypassed", "finalize_gate_warning"} + main_fields = set(_node(spec, "main_agent").output_fields) + reporter = _node(spec, "agent_team_reporter") + assert marker_keys <= main_fields + assert marker_keys <= set(reporter.context_policy.include_fields) + assert marker_keys <= set(reporter.output_fields) diff --git a/workflows/agent_team/nodes/main_agent.py b/workflows/agent_team/nodes/main_agent.py index 0fceb8d..0ada121 100644 --- a/workflows/agent_team/nodes/main_agent.py +++ b/workflows/agent_team/nodes/main_agent.py @@ -104,6 +104,7 @@ from workflows.agent_team.observers.auto_fan_in import AutoFanInObserver from workflows.agent_team.observers.bare_text_finalize import ( BareTextFinalizeObserver, + append_bypass_warning, ) from workflows.agent_team.observers.console import RichConsoleObserver from workflows.agent_team.observers.no_progress_guard import NoProgressGuard @@ -1667,6 +1668,13 @@ async def _run_main_loop( url_repair_stats["unmatched"], url_repair_stats["checked"], ) + # Re-attach the finalize-gate bypass warning (if any) only now: the + # observer stores it while the board is live, and this is the delivery + # boundary for the coordinator's own answer. When the reporter runs it + # re-appends on its own freshly finalized text (References cleanup would + # strip an earlier append) — see review on #48. + final_text = append_bypass_warning(final_text, result.metadata) + if reporter_enabled and result.metadata.get("report_handoff"): logger.info( "agent_team: research stopped by %s; advancing to downstream reporter", @@ -1708,6 +1716,12 @@ async def _run_main_loop( result.metadata.get("final_answer_rescue_mode") or "", ), "final_answer_source": answer_source, + "finalize_gate_bypassed": str( + result.metadata.get("finalize_gate_bypassed") or "", + ), + "finalize_gate_warning": str( + result.metadata.get("finalize_gate_warning") or "", + ), # Conditional edge in both agent-team specs consumes this resolved # per-request value. False routes directly to END, preserving the # coordinator's answer as the protocol final. diff --git a/workflows/agent_team/nodes/reporter.py b/workflows/agent_team/nodes/reporter.py index 3eb220d..c12c3b6 100644 --- a/workflows/agent_team/nodes/reporter.py +++ b/workflows/agent_team/nodes/reporter.py @@ -337,6 +337,16 @@ async def agent_team_reporter( if not report_md.strip(): return {} + # Re-attach the finalize-gate bypass warning now: References cleanup + # inside the chain has already run, so this append survives it — the + # observer's stored warning is the same one main_agent would have + # appended had the reporter not replaced the answer (review on #48). + from workflows.agent_team.observers.bare_text_finalize import ( + append_bypass_warning, + ) + + report_md = append_bypass_warning(report_md, state) + try: _refresh_trace_terminal(state.get("metadata") or {}, report_md) except Exception as exc: @@ -355,4 +365,6 @@ async def agent_team_reporter( "final_answer_source": "reporter_llm", "final_answer_rescued": False, "final_answer_rescue_mode": "", + "finalize_gate_bypassed": str(state.get("finalize_gate_bypassed") or ""), + "finalize_gate_warning": str(state.get("finalize_gate_warning") or ""), } diff --git a/workflows/agent_team/observers/bare_text_finalize.py b/workflows/agent_team/observers/bare_text_finalize.py index 9e7ad10..0ee3db1 100644 --- a/workflows/agent_team/observers/bare_text_finalize.py +++ b/workflows/agent_team/observers/bare_text_finalize.py @@ -2,6 +2,8 @@ from __future__ import annotations import logging +from collections.abc import Mapping +from typing import Any from frontier_agent.core.execution_context import get_current_execution_scope from frontier_agent.core.loop_types import BaseObserver, Intervention, TurnContext @@ -12,13 +14,14 @@ logger = logging.getLogger(__name__) -def _unfinished_note(task_id: str, text: str) -> str: - """Footnote an answer that is being delivered despite a blocked gate. +def _build_bypass_warning(task_id: str) -> str: + """Build the standalone warning for an answer delivered past a blocked gate. - The last turn must not lose the answer to ``max_turns`` (see the bypass - below), but an unfinished run must never read as a clean success: name the - board items that are still unresolved so the reader can tell which parts - of the answer were never corroborated. + Built once while the task board is still live — the board is cleared + before the workflow output is assembled — and stored in + ``finalize_gate_warning`` for the delivery nodes to append after any + finalization that would otherwise strip it (the reporter's References + cleanup drops everything after that heading). """ pending = unresolved_task_ids(task_id) if pending: @@ -26,13 +29,27 @@ def _unfinished_note(task_id: str, text: str) -> str: else: what = "the finalize gate was still rejecting this submission" return ( - f"{text}\n\n---\n\n" + "\n\n---\n\n" f"> ⚠ Unfinished work at submission: {what}. This answer was " "delivered on the final turn despite the gate — conclusions that " "depend on that work are unverified." ) +def append_bypass_warning(text: str, source: Mapping[str, Any] | None) -> str: + """Re-attach the stored bypass warning at a delivery boundary, once. + + The observer stores the ready-made warning; the delivery nodes append + it *after* finalization (reporter References cleanup would strip an + earlier append). Idempotent: text already carrying the warning is + returned unchanged. + """ + warning = str((source or {}).get("finalize_gate_warning") or "") + if not warning or not text or warning in text: + return text + return f"{text.rstrip()}{warning}" + + class BareTextFinalizeObserver(BaseObserver): critical = True @@ -61,11 +78,14 @@ async def on_llm_response(self, ctx: TurnContext) -> Intervention | None: if err: # Last-turn bypass: the gate says BLOCK, but the answer is # delivered anyway rather than lost to max_turns. Keep that - # fallback — and make it visible (machine-readable marker - # plus a user-visible note), so an unfinished run is never - # presented as a clean success. + # fallback — and make it visible: store the gate message and + # a ready-to-append warning (built while the board is live). + # The visible text is appended by the delivery nodes AFTER + # reporter finalization, which would strip it otherwise. ctx.metadata["finalize_gate_bypassed"] = err - text = _unfinished_note(ctx.task_id, text) + ctx.metadata["finalize_gate_warning"] = _build_bypass_warning( + ctx.task_id, + ) ctx.metadata["final_answer"] = text ctx.metadata["final_answer_confidence"] = 1.0 logger.info( diff --git a/workflows/agent_team/spec.py b/workflows/agent_team/spec.py index 21a0da3..2fcc9b9 100644 --- a/workflows/agent_team/spec.py +++ b/workflows/agent_team/spec.py @@ -47,6 +47,7 @@ "answer_status", "answer_sentinel", "final_answer_rescued", "final_answer_rescue_mode", "final_answer_source", "stopped_by", + "finalize_gate_bypassed", "finalize_gate_warning", ], ), NodeDefinition( @@ -64,6 +65,7 @@ "live_followups", "effective_question", "reporter_backend", "reporter_wall_time_s", "reporter_deadline_monotonic_s", + "finalize_gate_bypassed", "finalize_gate_warning", ], ), compression=CompressionConfig(enabled=False), @@ -76,6 +78,8 @@ "final_answer_source", "final_answer_rescued", "final_answer_rescue_mode", + "finalize_gate_bypassed", + "finalize_gate_warning", ], ), ], diff --git a/workflows/agent_team/spec_report.py b/workflows/agent_team/spec_report.py index 87b67df..51576a7 100644 --- a/workflows/agent_team/spec_report.py +++ b/workflows/agent_team/spec_report.py @@ -49,6 +49,7 @@ "answer_status", "answer_sentinel", "final_answer_rescued", "final_answer_rescue_mode", "final_answer_source", "stopped_by", + "finalize_gate_bypassed", "finalize_gate_warning", ], ), NodeDefinition( @@ -66,6 +67,7 @@ "live_followups", "effective_question", "reporter_backend", "reporter_wall_time_s", "reporter_deadline_monotonic_s", + "finalize_gate_bypassed", "finalize_gate_warning", ], ), compression=CompressionConfig(enabled=False), @@ -78,6 +80,8 @@ "final_answer_source", "final_answer_rescued", "final_answer_rescue_mode", + "finalize_gate_bypassed", + "finalize_gate_warning", ], ), ],