Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
225 changes: 225 additions & 0 deletions tests/test_bare_text_finalize_gate.py
Original file line number Diff line number Diff line change
@@ -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)
14 changes: 14 additions & 0 deletions workflows/agent_team/nodes/main_agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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.
Expand Down
12 changes: 12 additions & 0 deletions workflows/agent_team/nodes/reporter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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 ""),
}
50 changes: 50 additions & 0 deletions workflows/agent_team/observers/bare_text_finalize.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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(
Expand Down
4 changes: 4 additions & 0 deletions workflows/agent_team/spec.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand All @@ -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),
Expand All @@ -76,6 +78,8 @@
"final_answer_source",
"final_answer_rescued",
"final_answer_rescue_mode",
"finalize_gate_bypassed",
"finalize_gate_warning",
],
),
],
Expand Down
Loading
Loading