diff --git a/agent/turn_final_response.py b/agent/turn_final_response.py index 64620e34bd..906308a1e5 100644 --- a/agent/turn_final_response.py +++ b/agent/turn_final_response.py @@ -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") diff --git a/agent/turn_stop_gates.py b/agent/turn_stop_gates.py index 5a44030174..2d92693905 100644 --- a/agent/turn_stop_gates.py +++ b/agent/turn_stop_gates.py @@ -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, + ) diff --git a/agent/turn_tool_round.py b/agent/turn_tool_round.py index 545dd03089..7c507e1911 100644 --- a/agent/turn_tool_round.py +++ b/agent/turn_tool_round.py @@ -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 diff --git a/agent/turn_tool_validation.py b/agent/turn_tool_validation.py index e71b1302aa..db3e89334e 100644 --- a/agent/turn_tool_validation.py +++ b/agent/turn_tool_validation.py @@ -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