diff --git a/CHANGELOG.md b/CHANGELOG.md index 0acb615..51bea5c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -62,6 +62,9 @@ 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. - Native mode puts the CLI's own Python environment ahead of the inherited `PATH`, so `read_file`, `download_file`, and `python3` inside `bash` use the interpreter the CLI was installed with rather than a system Python. diff --git a/tests/test_bare_text_finalize_gate.py b/tests/test_bare_text_finalize_gate.py new file mode 100644 index 0000000..d893292 --- /dev/null +++ b/tests/test_bare_text_finalize_gate.py @@ -0,0 +1,225 @@ +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 [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: + 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_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). 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"}) + ctx = _context(turn=19) + + interventions = _dispatch(ctx) + + (intervention,) = interventions + assert intervention.stop_reason == "final_answer" + 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 marker, no warning.""" + _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 + 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 2164989..0ee3db1 100644 --- a/workflows/agent_team/observers/bare_text_finalize.py +++ b/workflows/agent_team/observers/bare_text_finalize.py @@ -2,15 +2,54 @@ 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 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 _build_bypass_warning(task_id: str) -> str: + """Build the standalone warning for an answer delivered past a blocked gate. + + 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: + what = f"task-board item(s) still unfinished: {', '.join(pending)}" + else: + what = "the finalize gate was still rejecting this submission" + return ( + "\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 @@ -36,6 +75,17 @@ 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: 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 + 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", ], ), ],