refactor(agent/turn_tool_round,turn_stop_gates,turn_tool_validation,turn_final_response): dedupe nudge/error-result/partial-exit blocks, drop dead verdict scaffolding
This commit is contained in:
@@ -17,6 +17,12 @@ from agent.turn_stop_gates import apply_stop_gates
|
||||
|
||||
logger = logging.getLogger("agent.conversation_loop")
|
||||
|
||||
# Ephemeral retry scaffolding rows popped before the final answer becomes durable.
|
||||
_EPHEMERAL_SCAFFOLDING_FLAGS = (
|
||||
"_thinking_prefill", "_empty_recovery_synthetic", "_empty_terminal_sentinel",
|
||||
"_dropped_toolcall_nudge",
|
||||
)
|
||||
|
||||
|
||||
@dataclass
|
||||
class FinalResponseVerdict:
|
||||
@@ -136,11 +142,7 @@ def finish_text_response(
|
||||
interim_msg = agent._build_assistant_message(assistant_message, "incomplete")
|
||||
append_message(messages, interim_msg)
|
||||
agent._emit_interim_assistant_message(interim_msg)
|
||||
|
||||
continue_msg = {
|
||||
"role": "user", "content": _CODEX_ACK_CONTINUATION_NUDGE
|
||||
}
|
||||
append_message(messages, continue_msg)
|
||||
append_message(messages, {"role": "user", "content": _CODEX_ACK_CONTINUATION_NUDGE})
|
||||
agent._session_messages = messages
|
||||
# An acknowledgment is non-final: its text must not suppress
|
||||
# iteration-limit summarization if the continuation exhausts budget.
|
||||
@@ -204,12 +206,7 @@ def finish_text_response(
|
||||
while (
|
||||
messages
|
||||
and isinstance(messages[-1], dict)
|
||||
and (
|
||||
messages[-1].get("_thinking_prefill")
|
||||
or messages[-1].get("_empty_recovery_synthetic")
|
||||
or messages[-1].get("_empty_terminal_sentinel")
|
||||
or messages[-1].get("_dropped_toolcall_nudge")
|
||||
)
|
||||
and any(messages[-1].get(flag) for flag in _EPHEMERAL_SCAFFOLDING_FLAGS)
|
||||
):
|
||||
messages.pop()
|
||||
|
||||
@@ -243,4 +240,3 @@ def finish_text_response(
|
||||
if not agent.quiet_mode:
|
||||
agent._safe_print(f"🎉 Conversation completed after {api_call_count} OpenAI-compatible API call(s)")
|
||||
return _verdict("break")
|
||||
return _verdict("fallthrough")
|
||||
|
||||
@@ -1,12 +1,12 @@
|
||||
"""Text-response stop gates for the conversation turn loop.
|
||||
|
||||
Extracted from ``run_conversation``. When the model stops with a text answer, three
|
||||
gates may instead append the answer as an interim row plus a synthetic user-role nudge
|
||||
and continue the turn: verify-on-stop (#65919), the ``pre_verify`` plugin hook after code
|
||||
edits, and the kanban worker terminal-tool guard. Each keeps the candidate answer as a
|
||||
budget-exhaustion fallback (``pending_verification_response``) and clears
|
||||
``final_response`` so the finalizer can tell this gate from error exits (#61631).
|
||||
Nothing here imports ``agent.conversation_loop`` at module level (cycle).
|
||||
When the model stops with a text answer, three gates may instead append the answer as an
|
||||
interim row plus a synthetic user-role nudge and continue the turn: verify-on-stop (#65919),
|
||||
the ``pre_verify`` plugin hook after code edits, and the kanban worker terminal-tool guard.
|
||||
Each keeps the candidate answer as a budget-exhaustion fallback
|
||||
(``pending_verification_response``) and clears ``final_response`` so the finalizer can tell
|
||||
this gate from error exits (#61631). Nothing here imports ``agent.conversation_loop`` at
|
||||
module level (cycle).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -14,7 +14,7 @@ from __future__ import annotations
|
||||
import logging
|
||||
import os
|
||||
from dataclasses import dataclass
|
||||
from typing import Any, Dict, List
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from agent.message_metadata import append_message
|
||||
|
||||
@@ -32,6 +32,75 @@ class StopGateVerdict:
|
||||
pending_verification_response_previewed: Any
|
||||
|
||||
|
||||
def _verify_on_stop_nudge(agent) -> Optional[str]:
|
||||
try:
|
||||
from agent.verification_stop import (
|
||||
build_verify_on_stop_nudge, verify_on_stop_enabled
|
||||
)
|
||||
|
||||
if verify_on_stop_enabled():
|
||||
return build_verify_on_stop_nudge(
|
||||
session_id=getattr(agent, "session_id", None),
|
||||
changed_paths=getattr(agent, "_turn_file_mutation_paths", set()),
|
||||
attempts=getattr(agent, "_verification_stop_nudges", 0),
|
||||
)
|
||||
except Exception:
|
||||
logger.debug("verification stop-loop check failed", exc_info=True)
|
||||
return None
|
||||
|
||||
|
||||
def _pre_verify_nudge(agent, final_response, attempt: int) -> Optional[str]:
|
||||
"""After code edits a registered ``pre_verify`` hook may keep the agent going one
|
||||
more turn; no default continuation cost."""
|
||||
_edited = sorted(getattr(agent, "_turn_file_mutation_paths", set()) or [])
|
||||
try:
|
||||
from agent.verify_hooks import max_verify_nudges
|
||||
from hermes_cli.lifecycle import has_hook
|
||||
from hermes_cli.plugins import get_pre_verify_continue_message
|
||||
|
||||
if _edited and has_hook("pre_verify") and attempt < max_verify_nudges():
|
||||
# Posture is fixed for the session — resolve once + cache.
|
||||
coding = getattr(agent, "_resolved_is_coding", None)
|
||||
if coding is None:
|
||||
from agent.coding_context import is_coding_context
|
||||
coding = bool(is_coding_context(platform=getattr(agent, "platform", "") or ""))
|
||||
agent._resolved_is_coding = coding
|
||||
return get_pre_verify_continue_message(
|
||||
session_id=getattr(agent, "session_id", None) or "",
|
||||
platform=getattr(agent, "platform", "") or "",
|
||||
model=getattr(agent, "model", "") or "", coding=coding, attempt=attempt,
|
||||
final_response=final_response, changed_paths=_edited,
|
||||
)
|
||||
except Exception:
|
||||
logger.debug("pre_verify hook check failed", exc_info=True)
|
||||
return None
|
||||
|
||||
|
||||
def _kanban_stop_nudge(agent, messages) -> Optional[str]:
|
||||
"""Workers must end with kanban_complete / kanban_block; a narrated stop is recorded
|
||||
as protocol_violation, so nudge once or twice first."""
|
||||
try:
|
||||
from agent.kanban_stop import build_kanban_stop_nudge
|
||||
|
||||
return build_kanban_stop_nudge(
|
||||
messages=messages, attempts=getattr(agent, "_kanban_stop_nudges", 0)
|
||||
)
|
||||
except Exception:
|
||||
logger.debug("kanban stop-loop check failed", exc_info=True)
|
||||
return None
|
||||
|
||||
|
||||
def _append_interim_answer(agent, final_msg, messages, conversation_history, flush_fail_msg: str) -> None:
|
||||
"""Real content: persist and emit as interim so the user sees the attempted answer;
|
||||
only the nudge is flagged synthetic (#65919)."""
|
||||
agent._emit_interim_assistant_message(final_msg)
|
||||
append_message(messages, final_msg)
|
||||
try:
|
||||
agent._flush_messages_to_session_db(messages, conversation_history)
|
||||
except Exception:
|
||||
logger.debug(flush_fail_msg, exc_info=True)
|
||||
|
||||
|
||||
def apply_stop_gates(
|
||||
agent: Any, final_msg: Dict[str, Any], *, final_response: Any, messages: List[Dict[str, Any]],
|
||||
conversation_history: Any, pending_verification_response: Any,
|
||||
@@ -41,125 +110,50 @@ def apply_stop_gates(
|
||||
are user-role rows appended only after the assistant answer row, so role alternation
|
||||
holds. Hook lookups are imported lazily from their origin modules (tests patch them
|
||||
there)."""
|
||||
_pending_verification_response = pending_verification_response
|
||||
_pending_verification_response_previewed = pending_verification_response_previewed
|
||||
|
||||
def _verdict(continue_turn: bool) -> StopGateVerdict:
|
||||
def _continue(nudge: str, flag: str) -> StopGateVerdict:
|
||||
append_message(messages, {"role": "user", "content": nudge, flag: True})
|
||||
agent._session_messages = messages
|
||||
# Keep the answer only as a budget-exhaustion fallback; clear ``final_response`` so
|
||||
# the finalizer can tell this gate from error exits. Mark previewed only if the
|
||||
# candidate is reused (#61631).
|
||||
return StopGateVerdict(
|
||||
continue_turn=continue_turn, final_response=None if continue_turn else final_response,
|
||||
pending_verification_response=_pending_verification_response,
|
||||
pending_verification_response_previewed=_pending_verification_response_previewed,
|
||||
continue_turn=True, final_response=None,
|
||||
pending_verification_response=final_response,
|
||||
pending_verification_response_previewed=agent._interim_content_was_streamed(
|
||||
final_response or ""
|
||||
),
|
||||
)
|
||||
|
||||
try:
|
||||
from agent.verification_stop import (
|
||||
build_verify_on_stop_nudge, verify_on_stop_enabled
|
||||
)
|
||||
|
||||
if verify_on_stop_enabled():
|
||||
_verify_nudge = build_verify_on_stop_nudge(
|
||||
session_id=getattr(agent, "session_id", None),
|
||||
changed_paths=getattr(agent, "_turn_file_mutation_paths", set()),
|
||||
attempts=getattr(agent, "_verification_stop_nudges", 0),
|
||||
)
|
||||
else:
|
||||
_verify_nudge = None
|
||||
except Exception:
|
||||
logger.debug("verification stop-loop check failed", exc_info=True)
|
||||
_verify_nudge = None
|
||||
|
||||
_verify_nudge = _verify_on_stop_nudge(agent)
|
||||
if _verify_nudge:
|
||||
agent._verification_stop_nudges = (
|
||||
getattr(agent, "_verification_stop_nudges", 0) + 1
|
||||
)
|
||||
final_msg["finish_reason"] = "verification_required"
|
||||
# Real content: persist and emit as interim so the user sees the
|
||||
# attempted answer; only the nudge is flagged synthetic. (#65919)
|
||||
agent._emit_interim_assistant_message(final_msg)
|
||||
append_message(messages, final_msg)
|
||||
try:
|
||||
agent._flush_messages_to_session_db(messages, conversation_history)
|
||||
except Exception:
|
||||
logger.debug("verify-on-stop interim flush failed", exc_info=True)
|
||||
append_message(messages, {
|
||||
"role": "user", "content": _verify_nudge, "_verification_stop_synthetic": True
|
||||
})
|
||||
agent._session_messages = messages
|
||||
_append_interim_answer(
|
||||
agent, final_msg, messages, conversation_history, "verify-on-stop interim flush failed"
|
||||
)
|
||||
verdict = _continue(_verify_nudge, "_verification_stop_synthetic")
|
||||
# Internal nudge: stay silent on the terminal, debug-log only.
|
||||
logger.debug("verification stop-loop nudge issued (attempt %d)",
|
||||
agent._verification_stop_nudges)
|
||||
# Keep the answer only as a budget-exhaustion fallback; clear
|
||||
# ``final_response`` so the finalizer can tell this gate from error
|
||||
# exits. Mark previewed only if the candidate is reused. (#61631)
|
||||
_pending_verification_response = final_response
|
||||
_pending_verification_response_previewed = (
|
||||
agent._interim_content_was_streamed(final_response or "")
|
||||
)
|
||||
return _verdict(True)
|
||||
return verdict
|
||||
|
||||
# pre_verify hook gate: after code edits a registered hook may keep the
|
||||
# agent going one more turn; no default continuation cost.
|
||||
_verify_nudge2 = None
|
||||
_edited = sorted(getattr(agent, "_turn_file_mutation_paths", set()) or [])
|
||||
_attempt = getattr(agent, "_pre_verify_nudges", 0)
|
||||
try:
|
||||
from agent.verify_hooks import max_verify_nudges
|
||||
from hermes_cli.lifecycle import has_hook
|
||||
from hermes_cli.plugins import get_pre_verify_continue_message
|
||||
|
||||
if _edited and has_hook("pre_verify") and _attempt < max_verify_nudges():
|
||||
# Posture is fixed for the session — resolve once + cache.
|
||||
coding = getattr(agent, "_resolved_is_coding", None)
|
||||
if coding is None:
|
||||
from agent.coding_context import is_coding_context
|
||||
coding = bool(is_coding_context(platform=getattr(agent, "platform", "") or ""))
|
||||
agent._resolved_is_coding = coding
|
||||
_verify_nudge2 = get_pre_verify_continue_message(
|
||||
session_id=getattr(agent, "session_id", None) or "",
|
||||
platform=getattr(agent, "platform", "") or "",
|
||||
model=getattr(agent, "model", "") or "", coding=coding, attempt=_attempt,
|
||||
final_response=final_response, changed_paths=_edited,
|
||||
)
|
||||
except Exception:
|
||||
logger.debug("pre_verify hook check failed", exc_info=True)
|
||||
_verify_nudge2 = None
|
||||
|
||||
_verify_nudge2 = _pre_verify_nudge(agent, final_response, _attempt)
|
||||
if _verify_nudge2:
|
||||
agent._pre_verify_nudges = _attempt + 1
|
||||
final_msg["finish_reason"] = "verify_hook_continue"
|
||||
# Real content: persist and emit as interim so the user sees the
|
||||
# attempted answer; only the nudge is flagged synthetic. (#65919)
|
||||
agent._emit_interim_assistant_message(final_msg)
|
||||
append_message(messages, final_msg)
|
||||
try:
|
||||
agent._flush_messages_to_session_db(messages, conversation_history)
|
||||
except Exception:
|
||||
logger.debug("pre_verify interim flush failed", exc_info=True)
|
||||
append_message(messages, {
|
||||
"role": "user", "content": _verify_nudge2, "_pre_verify_synthetic": True
|
||||
})
|
||||
agent._session_messages = messages
|
||||
_append_interim_answer(
|
||||
agent, final_msg, messages, conversation_history, "pre_verify interim flush failed"
|
||||
)
|
||||
verdict = _continue(_verify_nudge2, "_pre_verify_synthetic")
|
||||
logger.debug("pre_verify nudge issued (attempt %d)",
|
||||
agent._pre_verify_nudges)
|
||||
_pending_verification_response = final_response
|
||||
_pending_verification_response_previewed = (
|
||||
agent._interim_content_was_streamed(final_response or "")
|
||||
)
|
||||
return _verdict(True)
|
||||
|
||||
# ── Kanban worker terminal-tool stop guard ─────────────
|
||||
# Workers must end with kanban_complete / kanban_block; a narrated stop
|
||||
# is recorded as protocol_violation, so nudge once or twice first.
|
||||
try:
|
||||
from agent.kanban_stop import build_kanban_stop_nudge
|
||||
|
||||
_kanban_nudge = build_kanban_stop_nudge(
|
||||
messages=messages, attempts=getattr(agent, "_kanban_stop_nudges", 0)
|
||||
)
|
||||
except Exception:
|
||||
logger.debug("kanban stop-loop check failed", exc_info=True)
|
||||
_kanban_nudge = None
|
||||
return verdict
|
||||
|
||||
_kanban_nudge = _kanban_stop_nudge(agent, messages)
|
||||
if _kanban_nudge:
|
||||
agent._kanban_stop_nudges = (
|
||||
getattr(agent, "_kanban_stop_nudges", 0) + 1
|
||||
@@ -167,10 +161,7 @@ def apply_stop_gates(
|
||||
final_msg["finish_reason"] = "kanban_terminal_required"
|
||||
final_msg["_kanban_stop_synthetic"] = True
|
||||
append_message(messages, final_msg)
|
||||
append_message(messages, {
|
||||
"role": "user", "content": _kanban_nudge, "_kanban_stop_synthetic": True
|
||||
})
|
||||
agent._session_messages = messages
|
||||
verdict = _continue(_kanban_nudge, "_kanban_stop_synthetic")
|
||||
logger.info(
|
||||
"kanban stop-loop nudge issued (attempt %d) task=%s",
|
||||
agent._kanban_stop_nudges,
|
||||
@@ -180,11 +171,9 @@ def apply_stop_gates(
|
||||
"⚠️ Kanban worker tried to exit without "
|
||||
"kanban_complete/kanban_block — nudging to finish"
|
||||
)
|
||||
# Same finalizer contract as verify-on-stop: clear final_response so
|
||||
# budget exhaustion doesn't treat the narrated stop as an answer.
|
||||
_pending_verification_response = final_response
|
||||
_pending_verification_response_previewed = (
|
||||
agent._interim_content_was_streamed(final_response or "")
|
||||
)
|
||||
return _verdict(True)
|
||||
return _verdict(False)
|
||||
return verdict
|
||||
return StopGateVerdict(
|
||||
continue_turn=False, final_response=final_response,
|
||||
pending_verification_response=pending_verification_response,
|
||||
pending_verification_response_previewed=pending_verification_response_previewed,
|
||||
)
|
||||
|
||||
@@ -1,15 +1,16 @@
|
||||
"""One tool-calling round of the conversation turn loop: validate/cap/dedupe the model's
|
||||
tool calls, persist the tool-call turn BEFORE any side effect, execute the tools, honour
|
||||
guardrail halts / persistence failures, then compress after tool results. Extracted from
|
||||
``run_conversation``'s ``if assistant_message.tool_calls:`` branch; nothing here imports
|
||||
``agent.conversation_loop`` at module level (cycle) — loop-internal helpers resolve lazily.
|
||||
guardrail halts / persistence failures, then compress after tool results. Nothing here
|
||||
imports ``agent.conversation_loop`` at module level (cycle) — loop-internal helpers resolve
|
||||
lazily.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from contextlib import suppress
|
||||
from dataclasses import dataclass
|
||||
import logging
|
||||
from typing import Any, Dict, Optional
|
||||
from typing import Any, Dict, Optional, Tuple
|
||||
|
||||
from agent.message_metadata import append_message
|
||||
from agent.message_sanitization import coalesce_tool_call_id
|
||||
@@ -18,6 +19,9 @@ from agent.turn_tool_validation import validate_tool_calls
|
||||
|
||||
logger = logging.getLogger("agent.conversation_loop")
|
||||
|
||||
# Post-response housekeeping tools: a round made only of these mutes tool progress.
|
||||
_HOUSEKEEPING_TOOLS = frozenset({"memory", "todo_list", "skill_manage", "session_search"})
|
||||
|
||||
|
||||
@dataclass
|
||||
class ToolRoundVerdict:
|
||||
@@ -49,9 +53,7 @@ def run_tool_round(
|
||||
durability invariant: resume must see the executed block if a destructive tool restarts
|
||||
Hermes; a failed canonical append ends the turn rather than running tools from
|
||||
process-only state."""
|
||||
from agent.conversation_loop import (
|
||||
_invalid_tool_name_error_content,
|
||||
)
|
||||
from agent.conversation_loop import _invalid_tool_name_error_content
|
||||
|
||||
def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> ToolRoundVerdict:
|
||||
return ToolRoundVerdict(
|
||||
@@ -75,39 +77,28 @@ def run_tool_round(
|
||||
conversation_history=conversation_history, api_call_count=api_call_count,
|
||||
effective_task_id=effective_task_id,
|
||||
)
|
||||
_mixed_invalid_batch = _tvv.mixed_invalid_batch
|
||||
if _tvv.action == "return":
|
||||
return _verdict("return", _tvv.result)
|
||||
if _tvv.action == "continue":
|
||||
return _verdict("continue")
|
||||
|
||||
# ── Post-call guardrails ──────────────────────────
|
||||
assistant_message.tool_calls = agent._cap_delegate_task_calls(
|
||||
assistant_message.tool_calls
|
||||
)
|
||||
assistant_message.tool_calls = agent._deduplicate_tool_calls(
|
||||
assistant_message.tool_calls
|
||||
agent._cap_delegate_task_calls(assistant_message.tool_calls)
|
||||
)
|
||||
|
||||
# Collect invalid calls so the assistant message keeps EVERY emitted
|
||||
# call (each tool_call needs a matching result) while only valid ones
|
||||
# dispatch.
|
||||
_invalid_batch_calls = []
|
||||
if _mixed_invalid_batch:
|
||||
_invalid_batch_calls = [
|
||||
tc for tc in assistant_message.tool_calls
|
||||
if tc.function.name not in agent.valid_tool_names
|
||||
]
|
||||
# Mixed batch: the assistant message keeps EVERY emitted call (each tool_call needs a
|
||||
# matching result) while only valid ones dispatch.
|
||||
_invalid_batch_calls = [
|
||||
tc for tc in assistant_message.tool_calls if tc.function.name not in agent.valid_tool_names
|
||||
] if _tvv.mixed_invalid_batch else []
|
||||
|
||||
_st = stage_tool_call_message(
|
||||
assistant_msg, duplicate_previous_interim = stage_tool_call_message(
|
||||
agent, assistant_message=assistant_message, finish_reason=finish_reason, messages=messages
|
||||
)
|
||||
assistant_msg = _st.assistant_msg
|
||||
duplicate_previous_interim = _st.duplicate_previous_interim
|
||||
append_message(messages, assistant_msg)
|
||||
|
||||
# Mixed batch: error-result invalid calls and drop them from execution.
|
||||
# The assistant message keeps all calls so tool_call/result pairs hold.
|
||||
if _invalid_batch_calls:
|
||||
for tc in _invalid_batch_calls:
|
||||
append_message(messages, {
|
||||
@@ -122,7 +113,6 @@ def run_tool_round(
|
||||
tc for tc in assistant_message.tool_calls if tc.function.name in agent.valid_tool_names
|
||||
]
|
||||
|
||||
_tool_turn_persisted = None
|
||||
try:
|
||||
# Persist the tool-call turn before any tool side effects so resume
|
||||
# sees the executed block if a destructive tool restarts Hermes.
|
||||
@@ -132,9 +122,7 @@ def run_tool_round(
|
||||
except Exception as exc:
|
||||
_tool_turn_persisted = False
|
||||
from hermes_state import classify_persistence_error
|
||||
agent._last_persistence_error_cause = (
|
||||
classify_persistence_error(exc)
|
||||
)
|
||||
agent._last_persistence_error_cause = classify_persistence_error(exc)
|
||||
logger.warning(
|
||||
"Incremental tool-call persistence failed before execution "
|
||||
"(session=%s): %s",
|
||||
@@ -162,10 +150,8 @@ def run_tool_round(
|
||||
# tool feed lines. Display callback only — TTS (_stream_callback) must
|
||||
# NOT receive None (its end-of-stream marker).
|
||||
if agent.stream_delta_callback:
|
||||
try:
|
||||
with suppress(Exception):
|
||||
agent.stream_delta_callback(None)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
agent._execute_tool_calls(assistant_message, messages, effective_task_id, api_call_count)
|
||||
|
||||
@@ -190,11 +176,9 @@ def run_tool_round(
|
||||
if final_response:
|
||||
agent._safe_print(f"\n{final_response}\n")
|
||||
if agent.stream_delta_callback:
|
||||
try:
|
||||
with suppress(Exception):
|
||||
agent.stream_delta_callback(final_response)
|
||||
agent.stream_delta_callback(None)
|
||||
except Exception:
|
||||
pass
|
||||
return _verdict("break")
|
||||
|
||||
# Reset per-turn retry counters so one truncation can't poison the turn.
|
||||
@@ -206,8 +190,7 @@ def run_tool_round(
|
||||
|
||||
# Refund the iteration when the ONLY tool was execute_code (programmatic
|
||||
# tool calling) — cheap RPC-style calls shouldn't eat the budget.
|
||||
_tc_names = {tc.function.name for tc in assistant_message.tool_calls}
|
||||
if _tc_names == {"execute_code"}:
|
||||
if {tc.function.name for tc in assistant_message.tool_calls} == {"execute_code"}:
|
||||
agent.iteration_budget.refund()
|
||||
|
||||
_ptc = compress_after_tool_results(
|
||||
@@ -232,40 +215,21 @@ def run_tool_round(
|
||||
# Touch activity so slow post-tool work plus a slow follow-up API call
|
||||
# can't exceed the gateway inactivity timeout (HERMES_AGENT_TIMEOUT).
|
||||
agent._touch_activity(f"tool results posted, continuing iteration #{api_call_count}")
|
||||
# Continue loop for next response
|
||||
return _verdict("continue")
|
||||
return _verdict("fallthrough")
|
||||
|
||||
|
||||
@dataclass
|
||||
class StagedToolCallMessage:
|
||||
"""Always ``action == "fallthrough"``. ``assistant_msg`` is the transcript row to append;
|
||||
``duplicate_previous_interim`` suppresses re-emitting interim commentary the previous
|
||||
``incomplete`` row already showed."""
|
||||
|
||||
action: str
|
||||
assistant_msg: Any
|
||||
duplicate_previous_interim: Any
|
||||
|
||||
|
||||
def stage_tool_call_message(
|
||||
agent: Any, *, assistant_message: Any, finish_reason: Any, messages: Any
|
||||
) -> StagedToolCallMessage:
|
||||
"""Build the assistant tool-call row and update the per-turn fallback/mute state: drop a bare
|
||||
bracketed marker beside a call (#78148), classify housekeeping-only rounds, keep visible
|
||||
content as the empty-follow-up fallback, pop thinking-only prefills (resetting their
|
||||
counters), re-arm the post-tool nudge and the dropped-tool-call stall budget."""
|
||||
from agent.conversation_loop import (
|
||||
_STALE_MARKER_RE,
|
||||
)
|
||||
) -> Tuple[Dict[str, Any], bool]:
|
||||
"""Build the assistant tool-call row and update the per-turn fallback/mute state.
|
||||
|
||||
def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> StagedToolCallMessage:
|
||||
return StagedToolCallMessage(
|
||||
action=action,
|
||||
assistant_msg=assistant_msg,
|
||||
duplicate_previous_interim=duplicate_previous_interim,
|
||||
|
||||
)
|
||||
Drops a bare bracketed marker beside a call (#78148), classifies housekeeping-only
|
||||
rounds, keeps visible content as the empty-follow-up fallback, pops thinking-only
|
||||
prefills (resetting their counters), re-arms the post-tool nudge and the
|
||||
dropped-tool-call stall budget. Returns ``(assistant_msg, duplicate_previous_interim)``;
|
||||
the flag suppresses re-emitting interim commentary the previous ``incomplete`` row
|
||||
already showed."""
|
||||
from agent.conversation_loop import _STALE_MARKER_RE
|
||||
|
||||
assistant_msg = agent._build_assistant_message(assistant_message, finish_reason)
|
||||
|
||||
@@ -285,9 +249,6 @@ def stage_tool_call_message(
|
||||
|
||||
# Classify tools regardless of visible content: a substantive tool-only
|
||||
# turn must invalidate any older housekeeping fallback.
|
||||
_HOUSEKEEPING_TOOLS = frozenset({
|
||||
"memory", "todo_list", "skill_manage", "session_search",
|
||||
})
|
||||
_all_housekeeping = all(
|
||||
tc.function.name in _HOUSEKEEPING_TOOLS for tc in assistant_message.tool_calls
|
||||
)
|
||||
@@ -316,17 +277,15 @@ def stage_tool_call_message(
|
||||
if clean:
|
||||
agent._vprint(f" ┊ 💬 {clean}")
|
||||
|
||||
# Pop thinking-only prefill message(s) before appending
|
||||
# (tool-call path — same rationale as the final-response path).
|
||||
# Pop thinking-only prefill message(s) before appending (same rationale as the
|
||||
# final-response path). Tool calls after a prefill recovery reset the prefill
|
||||
# counter, so each tool-call success is a fresh start, not a cumulative burn.
|
||||
_had_prefill = False
|
||||
while (
|
||||
messages and isinstance(messages[-1], dict) and messages[-1].get("_thinking_prefill")
|
||||
):
|
||||
messages.pop()
|
||||
_had_prefill = True
|
||||
|
||||
# Tool calls after a prefill recovery reset the prefill counter, so
|
||||
# each tool-call success is a fresh start, not a cumulative burn.
|
||||
if _had_prefill:
|
||||
agent._thinking_prefill_retries = 0
|
||||
agent._empty_content_retries = 0
|
||||
@@ -338,16 +297,11 @@ def stage_tool_call_message(
|
||||
|
||||
previous_msg = messages[-1] if messages else None
|
||||
current_interim_visible = agent._interim_assistant_visible_text(assistant_msg)
|
||||
previous_interim_visible = (
|
||||
agent._interim_assistant_visible_text(previous_msg)
|
||||
if isinstance(previous_msg, dict)
|
||||
else ""
|
||||
)
|
||||
duplicate_previous_interim = (
|
||||
bool(current_interim_visible)
|
||||
and isinstance(previous_msg, dict)
|
||||
and previous_msg.get("role") == "assistant"
|
||||
and previous_msg.get("finish_reason") == "incomplete"
|
||||
and previous_interim_visible == current_interim_visible
|
||||
and agent._interim_assistant_visible_text(previous_msg) == current_interim_visible
|
||||
)
|
||||
return _verdict("fallthrough")
|
||||
return assistant_msg, duplicate_previous_interim
|
||||
|
||||
@@ -2,10 +2,9 @@
|
||||
auto-repair and the 3-strike partial exit) and malformed JSON arguments (retry, then
|
||||
recovery tool results).
|
||||
|
||||
Extracted from ``run_conversation``. Role alternation is preserved on every path: an
|
||||
invalid batch is answered with tool-role error results (never a user message), and
|
||||
the exits close any open tool-result tail (#48879). Nothing here imports
|
||||
``agent.conversation_loop`` at module level (cycle); loop-internal helpers resolve lazily.
|
||||
Role alternation is preserved on every path: an invalid batch is answered with tool-role
|
||||
error results (never a user message), and the exits close any open tool-result tail
|
||||
(#48879). Nothing here imports ``agent.conversation_loop`` at module level (cycle).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -36,6 +35,37 @@ class ToolValidationVerdict:
|
||||
mixed_invalid_batch: bool
|
||||
|
||||
|
||||
def _preview_name(name: str) -> str:
|
||||
return name[:80] + "..." if len(name) > 80 else name
|
||||
|
||||
|
||||
def _append_tool_error_results(messages, tool_calls, content_for) -> None:
|
||||
"""One tool-role result per call so every tool_call keeps a matching result."""
|
||||
for tc in tool_calls:
|
||||
append_message(messages, {
|
||||
"role": "tool",
|
||||
"name": tc.function.name,
|
||||
"tool_call_id": coalesce_tool_call_id(tc),
|
||||
"content": content_for(tc),
|
||||
})
|
||||
|
||||
|
||||
def _partial_exit(agent, messages, conversation_history, api_call_count, final_response: str) -> Dict[str, Any]:
|
||||
"""Terminal partial result. Prior retries or an earlier tool batch leave a tool-result
|
||||
tail; close it as interrupt aborts do so the next turn is not tool→user (#48879).
|
||||
This path never reaches finalize_turn, so persist here."""
|
||||
close_interrupted_tool_sequence(messages, final_response)
|
||||
agent._persist_session(messages, conversation_history)
|
||||
return {
|
||||
"final_response": final_response,
|
||||
"messages": messages,
|
||||
"api_calls": api_call_count,
|
||||
"completed": False,
|
||||
"partial": True,
|
||||
"error": final_response,
|
||||
}
|
||||
|
||||
|
||||
def validate_tool_calls(
|
||||
agent: Any, assistant_message: Any, finish_reason: str, *, messages: List[Dict[str, Any]],
|
||||
conversation_history: Any, api_call_count: int, effective_task_id: Any,
|
||||
@@ -47,130 +77,95 @@ def validate_tool_calls(
|
||||
outright rather than retried."""
|
||||
from agent.conversation_loop import _invalid_tool_name_error_content
|
||||
|
||||
_mixed_invalid_batch = False
|
||||
tool_calls = assistant_message.tool_calls
|
||||
valid_names = agent.valid_tool_names
|
||||
|
||||
def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> ToolValidationVerdict:
|
||||
return ToolValidationVerdict(action=action, result=result, mixed_invalid_batch=_mixed_invalid_batch)
|
||||
|
||||
# Uniquify duplicate tool-call ids BEFORE any downstream consumer: the
|
||||
# pre-API sanitizer keeps only the first call/result per id. See
|
||||
# _uniquify_tool_call_ids.
|
||||
agent._uniquify_tool_call_ids(assistant_message.tool_calls)
|
||||
# pre-API sanitizer keeps only the first call/result per id.
|
||||
agent._uniquify_tool_call_ids(tool_calls)
|
||||
|
||||
# Validate tool call names - detect model hallucinations
|
||||
# Repair mismatched tool names before validating
|
||||
for tc in assistant_message.tool_calls:
|
||||
if tc.function.name not in agent.valid_tool_names:
|
||||
# Repair mismatched tool names before validating (model hallucinations).
|
||||
for tc in tool_calls:
|
||||
if tc.function.name not in valid_names:
|
||||
repaired = agent._repair_tool_call(tc.function.name)
|
||||
if repaired:
|
||||
print(f"{agent.log_prefix}🔧 Auto-repaired tool name: '{tc.function.name}' -> '{repaired}'")
|
||||
tc.function.name = repaired
|
||||
invalid_tool_calls = [
|
||||
tc.function.name for tc in assistant_message.tool_calls
|
||||
if tc.function.name not in agent.valid_tool_names
|
||||
]
|
||||
invalid_tool_calls = [tc.function.name for tc in tool_calls if tc.function.name not in valid_names]
|
||||
# Mixed batch: error-result ONLY the invalid calls and run the valid
|
||||
# ones; voiding the turn discards real work. Strikes advance only when a
|
||||
# turn has NO valid call, so a degenerate model still halts at 3.
|
||||
_mixed_invalid_batch = bool(invalid_tool_calls) and any(
|
||||
tc.function.name in agent.valid_tool_names for tc in assistant_message.tool_calls
|
||||
tc.function.name in valid_names for tc in tool_calls
|
||||
)
|
||||
if _mixed_invalid_batch:
|
||||
agent._invalid_tool_retries = 0
|
||||
invalid_name = invalid_tool_calls[0]
|
||||
invalid_preview = invalid_name[:80] + "..." if len(invalid_name) > 80 else invalid_name
|
||||
_n_valid = sum(
|
||||
1 for tc in assistant_message.tool_calls if tc.function.name in agent.valid_tool_names
|
||||
)
|
||||
_n_valid = sum(1 for tc in tool_calls if tc.function.name in valid_names)
|
||||
agent._buffer_vprint(
|
||||
f"⚠️ Unknown tool '{invalid_preview}' in batch — erroring that call, "
|
||||
f"⚠️ Unknown tool '{_preview_name(invalid_tool_calls[0])}' in batch — erroring that call, "
|
||||
f"executing {_n_valid} valid call(s)"
|
||||
)
|
||||
elif invalid_tool_calls:
|
||||
# Track retries for invalid tool calls
|
||||
agent._invalid_tool_retries += 1
|
||||
|
||||
# Return helpful error to model — model can agent-correct next turn
|
||||
invalid_name = invalid_tool_calls[0]
|
||||
invalid_preview = invalid_name[:80] + "..." if len(invalid_name) > 80 else invalid_name
|
||||
invalid_preview = _preview_name(invalid_tool_calls[0])
|
||||
agent._buffer_vprint(f"⚠️ Unknown tool '{invalid_preview}' — sending error to model for agent-correction ({agent._invalid_tool_retries}/3)")
|
||||
|
||||
if agent._invalid_tool_retries >= 3:
|
||||
agent._flush_status_buffer()
|
||||
agent._vprint(f"{agent.log_prefix}❌ Max retries (3) for invalid tool calls exceeded. Stopping as partial.", force=True)
|
||||
agent._invalid_tool_retries = 0
|
||||
_final_response = f"Model generated invalid tool call: {invalid_preview}"
|
||||
# Prior retries or an earlier tool batch leave a tool-result
|
||||
# tail; close it as interrupt aborts do so the next turn is not
|
||||
# tool→user. (#48879)
|
||||
close_interrupted_tool_sequence(messages, _final_response)
|
||||
agent._persist_session(messages, conversation_history)
|
||||
return _verdict("return", {
|
||||
"final_response": _final_response,
|
||||
"messages": messages,
|
||||
"api_calls": api_call_count,
|
||||
"completed": False,
|
||||
"partial": True,
|
||||
"error": _final_response
|
||||
})
|
||||
return _verdict("return", _partial_exit(
|
||||
agent, messages, conversation_history, api_call_count,
|
||||
f"Model generated invalid tool call: {invalid_preview}",
|
||||
))
|
||||
|
||||
assistant_msg = agent._build_assistant_message(assistant_message, finish_reason)
|
||||
append_message(messages, assistant_msg)
|
||||
for tc in assistant_message.tool_calls:
|
||||
_tc_name = tc.function.name
|
||||
if _tc_name not in agent.valid_tool_names:
|
||||
# See _invalid_tool_name_error_content for the
|
||||
# blank-name anti-priming rationale (#47967).
|
||||
content = _invalid_tool_name_error_content(
|
||||
_tc_name, agent.valid_tool_names
|
||||
)
|
||||
else:
|
||||
content = "Skipped: another tool call in this turn used an invalid name. Please retry this tool call."
|
||||
append_message(messages, {
|
||||
"role": "tool",
|
||||
"name": tc.function.name,
|
||||
"tool_call_id": coalesce_tool_call_id(tc),
|
||||
"content": content,
|
||||
})
|
||||
append_message(messages, agent._build_assistant_message(assistant_message, finish_reason))
|
||||
# See _invalid_tool_name_error_content for the blank-name anti-priming rationale (#47967).
|
||||
_append_tool_error_results(
|
||||
messages, tool_calls,
|
||||
lambda tc: (
|
||||
_invalid_tool_name_error_content(tc.function.name, valid_names)
|
||||
if tc.function.name not in valid_names
|
||||
else "Skipped: another tool call in this turn used an invalid name. Please retry this tool call."
|
||||
),
|
||||
)
|
||||
return _verdict("continue")
|
||||
# Reset retry counter on successful tool call validation
|
||||
agent._invalid_tool_retries = 0
|
||||
|
||||
# Validate tool call arguments are valid JSON
|
||||
# Handle empty strings as empty objects (common model quirk)
|
||||
# Validate tool call arguments are valid JSON; empty strings become empty
|
||||
# objects (common model quirk).
|
||||
invalid_json_args = []
|
||||
for tc in assistant_message.tool_calls:
|
||||
for tc in tool_calls:
|
||||
args = tc.function.arguments
|
||||
if isinstance(args, (dict, list)):
|
||||
tc.function.arguments = json.dumps(args)
|
||||
continue
|
||||
if args is not None and not isinstance(args, str):
|
||||
tc.function.arguments = str(args)
|
||||
args = tc.function.arguments
|
||||
# Treat empty/whitespace strings as empty object
|
||||
tc.function.arguments = args = str(args)
|
||||
if not args or not args.strip():
|
||||
tc.function.arguments = "{}"
|
||||
continue
|
||||
try:
|
||||
json.loads(args)
|
||||
except json.JSONDecodeError as e:
|
||||
if (
|
||||
_mixed_invalid_batch and tc.function.name not in agent.valid_tool_names
|
||||
):
|
||||
# This call never executes (invalid-name error result
|
||||
# below); don't let its broken args trigger the whole-turn
|
||||
# JSON retry.
|
||||
continue
|
||||
invalid_json_args.append((tc.function.name, str(e)))
|
||||
# A mixed-batch invalid-name call never executes (error result later);
|
||||
# don't let its broken args trigger the whole-turn JSON retry.
|
||||
if not (_mixed_invalid_batch and tc.function.name not in valid_names):
|
||||
invalid_json_args.append((tc.function.name, str(e)))
|
||||
|
||||
if invalid_json_args:
|
||||
invalid_names = {n for n, _ in invalid_json_args}
|
||||
# Routers may rewrite finish_reason "length" → "tool_calls", hiding
|
||||
# truncation; args not ending in } or ] (stripped) were cut off
|
||||
# mid-stream.
|
||||
_truncated = any(
|
||||
not (tc.function.arguments or "").rstrip().endswith(("}", "]"))
|
||||
for tc in assistant_message.tool_calls
|
||||
if tc.function.name in {n for n, _ in invalid_json_args}
|
||||
for tc in tool_calls if tc.function.name in invalid_names
|
||||
)
|
||||
if _truncated:
|
||||
agent._vprint(
|
||||
@@ -180,23 +175,12 @@ def validate_tool_calls(
|
||||
)
|
||||
agent._invalid_json_retries = 0
|
||||
agent._cleanup_task_resources(effective_task_id)
|
||||
_final_response = "Response truncated due to output length limit"
|
||||
# Same tool-tail close as interrupt / invalid-tool
|
||||
# exhaustion — this path never reaches finalize_turn.
|
||||
close_interrupted_tool_sequence(messages, _final_response)
|
||||
agent._persist_session(messages, conversation_history)
|
||||
return _verdict("return", {
|
||||
"final_response": _final_response,
|
||||
"messages": messages,
|
||||
"api_calls": api_call_count,
|
||||
"completed": False,
|
||||
"partial": True,
|
||||
"error": _final_response,
|
||||
})
|
||||
return _verdict("return", _partial_exit(
|
||||
agent, messages, conversation_history, api_call_count,
|
||||
"Response truncated due to output length limit",
|
||||
))
|
||||
|
||||
# Track retries for invalid JSON arguments
|
||||
agent._invalid_json_retries += 1
|
||||
|
||||
tool_name, error_msg = invalid_json_args[0]
|
||||
agent._buffer_vprint(f"⚠️ Invalid JSON in tool call arguments for '{tool_name}': {error_msg}")
|
||||
|
||||
@@ -204,35 +188,26 @@ def validate_tool_calls(
|
||||
agent._buffer_vprint(f"🔄 Retrying API call ({agent._invalid_json_retries}/3)...")
|
||||
# Don't add anything to messages, just retry the API call
|
||||
return _verdict("continue")
|
||||
else:
|
||||
# Instead of returning partial, inject tool error results so the model can recover.
|
||||
# Using tool results (not user messages) preserves role alternation.
|
||||
agent._buffer_vprint("⚠️ Injecting recovery tool results for invalid JSON...")
|
||||
agent._invalid_json_retries = 0 # Reset for next attempt
|
||||
# Instead of returning partial, inject tool error results so the model can recover.
|
||||
# Using tool results (not user messages) preserves role alternation.
|
||||
agent._buffer_vprint("⚠️ Injecting recovery tool results for invalid JSON...")
|
||||
agent._invalid_json_retries = 0 # Reset for next attempt
|
||||
# Append the assistant message with its (broken) tool_calls, then one
|
||||
# error result per call.
|
||||
append_message(messages, agent._build_assistant_message(assistant_message, finish_reason))
|
||||
|
||||
# Append the assistant message with its (broken) tool_calls
|
||||
recovery_assistant = agent._build_assistant_message(assistant_message, finish_reason)
|
||||
append_message(messages, recovery_assistant)
|
||||
def _json_error_result(tc) -> str:
|
||||
if tc.function.name not in invalid_names:
|
||||
return "Skipped: other tool call in this response had invalid JSON."
|
||||
err = next(e for n, e in invalid_json_args if n == tc.function.name)
|
||||
return (
|
||||
f"Error: Invalid JSON arguments. {err}. "
|
||||
f"For tools with no required parameters, use an empty object: {{}}. "
|
||||
f"Please retry with valid JSON."
|
||||
)
|
||||
|
||||
# Respond with tool error results for each tool call
|
||||
invalid_names = {name for name, _ in invalid_json_args}
|
||||
for tc in assistant_message.tool_calls:
|
||||
if tc.function.name in invalid_names:
|
||||
err = next(e for n, e in invalid_json_args if n == tc.function.name)
|
||||
tool_result = (
|
||||
f"Error: Invalid JSON arguments. {err}. "
|
||||
f"For tools with no required parameters, use an empty object: {{}}. "
|
||||
f"Please retry with valid JSON."
|
||||
)
|
||||
else:
|
||||
tool_result = "Skipped: other tool call in this response had invalid JSON."
|
||||
append_message(messages, {
|
||||
"role": "tool",
|
||||
"name": tc.function.name,
|
||||
"tool_call_id": coalesce_tool_call_id(tc),
|
||||
"content": tool_result,
|
||||
})
|
||||
return _verdict("continue")
|
||||
_append_tool_error_results(messages, tool_calls, _json_error_result)
|
||||
return _verdict("continue")
|
||||
|
||||
# Reset retry counter on successful JSON validation
|
||||
agent._invalid_json_retries = 0
|
||||
|
||||
Reference in New Issue
Block a user