From 67ef2e50fe50b5f4cd76ed4f386c0a9aa6a05d86 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 15:37:47 -0700 Subject: [PATCH] refactor(turn): extract response intake (normalize/hooks/scratchpad/codex-incomplete) into agent/turn_response_intake.py --- agent/conversation_loop.py | 179 +++------------------------ agent/turn_response_intake.py | 223 ++++++++++++++++++++++++++++++++++ 2 files changed, 242 insertions(+), 160 deletions(-) create mode 100644 agent/turn_response_intake.py diff --git a/agent/conversation_loop.py b/agent/conversation_loop.py index 23db1632a5..23b908f744 100644 --- a/agent/conversation_loop.py +++ b/agent/conversation_loop.py @@ -37,7 +37,6 @@ from agent.turn_empty_response import recover_empty_response from agent.turn_stop_gates import apply_stop_gates from agent.turn_tool_validation import validate_tool_calls from agent.turn_truncation import ( - continue_codex_incomplete, handle_content_policy_refusal, recover_from_truncation, ) @@ -82,15 +81,14 @@ from agent.prompt_caching import ( strip_anthropic_cache_control, strip_anthropic_tool_cache_control, ) -from agent.provider_projection import splice_provider_projection from agent.retry_utils import ( adaptive_rate_limit_backoff, # noqa: F401 — resolved lazily by agent.turn_recovery (tests patch it here) jittered_backoff, ) -from agent.trajectory import has_incomplete_scratchpad # Bind before the turn starts so a source-tree swap cannot load a skewed # finalizer at turn end. from agent.turn_finalizer import finalize_turn +from agent.turn_response_intake import normalize_model_response from agent.turn_loop_errors import handle_outer_loop_error from hermes_logging import set_session_context from tools.skill_provenance import set_current_write_origin @@ -3350,164 +3348,25 @@ def run_conversation( break try: - _transport = agent._get_transport() - _normalize_kwargs = {} - if agent.api_mode == "anthropic_messages": - _normalize_kwargs["strip_tool_prefix"] = agent._is_anthropic_oauth - normalized = _transport.normalize_response(response, **_normalize_kwargs) - assistant_message = normalized - finish_reason = normalized.finish_reason - - # Some OpenAI-compatible servers (llama-server) return content as dict/list, - # which crashes downstream .strip(); normalize to str. - if assistant_message.content is not None and not isinstance(assistant_message.content, str): - raw = assistant_message.content - if isinstance(raw, dict): - assistant_message.content = raw.get("text", "") or raw.get("content", "") or json.dumps(raw) - elif isinstance(raw, list): - # Multimodal content list — extract text parts - parts = [] - for part in raw: - if isinstance(part, str): - parts.append(part) - elif isinstance(part, dict) and part.get("type") == "text": - parts.append(part.get("text", "")) - elif isinstance(part, dict) and "text" in part: - parts.append(str(part["text"])) - assistant_message.content = "\n".join(parts) - else: - assistant_message.content = str(raw) - - # ── Agent-as-provider projection ────────────────────────────── - # Splice the provider-agent's own tool work in as call/result rows before - # this turn's assistant message; no-op for ordinary providers. - splice_provider_projection(agent, response, messages) - - try: - from hermes_cli.lifecycle import ( - has_hook, - invoke_hook as _invoke_hook, - ) - if has_hook("post_api_request"): - _assistant_tool_calls = ( - getattr(assistant_message, "tool_calls", None) or [] - ) - _assistant_text = assistant_message.content or "" - _api_ended_at = api_start_time + api_duration - _invoke_hook( - "post_api_request", - task_id=effective_task_id, - turn_id=turn_id, - api_request_id=api_request_id, - session_id=agent.session_id or "", - platform=agent.platform or "", - model=agent.model, - provider=agent.provider, - base_url=agent.base_url, - api_mode=agent.api_mode, - api_call_count=api_call_count, - api_duration=api_duration, - started_at=api_start_time, - ended_at=_api_ended_at, - # First stream chunk time (epoch s) from - # interruptible_streaming_api_call; None if not streamed / no - # chunk. TTFB = first_chunk_at - started_at. - first_chunk_at=getattr( - agent, "_last_api_first_chunk_at", None - ), - finish_reason=finish_reason, - message_count=len(api_messages), - response_model=getattr(response, "model", None), - response=agent._api_response_payload_for_hook( - response, - assistant_message, - finish_reason=finish_reason, - ), - usage=agent._usage_summary_for_api_request_hook(response), - assistant_message=assistant_message, - assistant_content_chars=len(_assistant_text), - assistant_tool_call_count=len(_assistant_tool_calls), - moa_references=_moa_reference_metrics_for_hook(agent), - ) - except Exception: - pass - - # Handle assistant response - if assistant_message.content and not agent.quiet_mode: - if agent.verbose_logging: - agent._vprint(f"{agent.log_prefix}🤖 Assistant: {assistant_message.content}") - else: - agent._vprint(f"{agent.log_prefix}🤖 Assistant: {assistant_message.content[:100]}{'...' if len(assistant_message.content) > 100 else ''}") - - # Notify progress callback of model's thinking (used by subagent - # delegation to relay the child's reasoning to the parent display). - if (assistant_message.content and agent.tool_progress_callback): - _think_text = assistant_message.content.strip() - # Strip reasoning XML tags that shouldn't leak to parent display - _think_text = re.sub( - r'', '', _think_text - ).strip() - # For subagents: relay first line to parent display (existing behaviour). - # For all agents with a structured callback: emit reasoning.available event. - first_line = _think_text.split('\n')[0][:80] if _think_text else "" - if first_line and getattr(agent, '_delegate_depth', 0) > 0: - try: - agent.tool_progress_callback("_thinking", first_line) - except Exception: - pass - elif _think_text: - try: - agent.tool_progress_callback("reasoning.available", "_thinking", _think_text[:500], None) - except Exception: - pass - - # Check for incomplete (opened but never closed) - # This means the model ran out of output tokens mid-reasoning — retry up to 2 times - if has_incomplete_scratchpad(assistant_message.content or ""): - agent._incomplete_scratchpad_retries += 1 - - agent._buffer_vprint("⚠️ Incomplete detected (opened but never closed)") - - if agent._incomplete_scratchpad_retries <= 2: - agent._buffer_vprint(f"🔄 Retrying API call ({agent._incomplete_scratchpad_retries}/2)...") - # Don't add the broken message, just retry - continue - else: - # Max retries - discard this turn and save as partial - agent._flush_status_buffer() - agent._vprint(f"{agent.log_prefix}❌ Max retries (2) for incomplete scratchpad. Saving as partial.", force=True) - agent._incomplete_scratchpad_retries = 0 - - rolled_back_messages = agent._get_messages_up_to_last_assistant(messages) - agent._cleanup_task_resources(effective_task_id) - agent._persist_session(messages, conversation_history) - - return { - "final_response": "Incomplete REASONING_SCRATCHPAD after 2 retries", - "messages": rolled_back_messages, - "api_calls": api_call_count, - "completed": False, - "partial": True, - "error": "Incomplete REASONING_SCRATCHPAD after 2 retries" - } - - # Reset incomplete scratchpad counter on clean response - agent._incomplete_scratchpad_retries = 0 - - if agent.api_mode == "codex_responses" and finish_reason == "incomplete": - _codex_result = continue_codex_incomplete( - agent, - assistant_message, - finish_reason, - messages=messages, - conversation_history=conversation_history, - api_call_count=api_call_count, - ) - if _codex_result is not None: - return _codex_result + _ri = normalize_model_response( + agent, + response=response, + messages=messages, + api_messages=api_messages, + conversation_history=conversation_history, + api_call_count=api_call_count, + api_duration=api_duration, + api_start_time=api_start_time, + api_request_id=api_request_id, + effective_task_id=effective_task_id, + turn_id=turn_id, + ) + assistant_message = _ri.assistant_message + finish_reason = _ri.finish_reason + if _ri.action == "return": + return _ri.result + if _ri.action == "continue": continue - elif hasattr(agent, "_codex_incomplete_retries"): - agent._codex_incomplete_retries = 0 # Check for tool calls if assistant_message.tool_calls: diff --git a/agent/turn_response_intake.py b/agent/turn_response_intake.py new file mode 100644 index 0000000000..bd3fe722b7 --- /dev/null +++ b/agent/turn_response_intake.py @@ -0,0 +1,223 @@ +"""Response intake for the conversation turn loop: normalize the raw provider response into +the assistant message, splice agent-as-provider projections, fire ``post_api_request``, relay +reasoning to the progress callback, and apply the incomplete-scratchpad / Codex-incomplete +continuation guards. Extracted from ``run_conversation``; nothing here imports +``agent.conversation_loop`` at module level (cycle). +""" + +from __future__ import annotations + +import logging +from dataclasses import dataclass +from typing import Any, Dict, Optional +import json +import re + +from agent.provider_projection import splice_provider_projection +from agent.trajectory import has_incomplete_scratchpad +from agent.turn_truncation import continue_codex_incomplete + +logger = logging.getLogger("agent.conversation_loop") + + +@dataclass +class ResponseIntakeVerdict: + """``action``: ``"fallthrough"`` (process ``assistant_message``), ``"continue"`` (retry the + iteration: incomplete scratchpad / Codex continuation) or ``"return"`` (``result`` is the + turn's result dict). ``assistant_message``/``finish_reason`` are the normalized outputs.""" + + action: str + assistant_message: Any + finish_reason: Any + result: Optional[Dict[str, Any]] = None + + +def normalize_model_response( + agent: Any, + *, + response: Any, + messages: Any, + api_messages: Any, + conversation_history: Any, + api_call_count: Any, + api_duration: Any, + api_start_time: Any, + api_request_id: Any, + effective_task_id: Any, + turn_id: Any, +) -> ResponseIntakeVerdict: + """Normalize ``response`` into ``assistant_message`` (str content, never dict/list) and run + the post-response hooks and continuation guards, in the original order.""" + from agent.conversation_loop import ( + _moa_reference_metrics_for_hook, + ) + assistant_message = None + finish_reason = None + + def _verdict(action: str, result: Optional[Dict[str, Any]] = None) -> ResponseIntakeVerdict: + return ResponseIntakeVerdict( + action=action, + assistant_message=assistant_message, + finish_reason=finish_reason, + result=result, + ) + + _transport = agent._get_transport() + _normalize_kwargs = {} + if agent.api_mode == "anthropic_messages": + _normalize_kwargs["strip_tool_prefix"] = agent._is_anthropic_oauth + normalized = _transport.normalize_response(response, **_normalize_kwargs) + assistant_message = normalized + finish_reason = normalized.finish_reason + + # Some OpenAI-compatible servers (llama-server) return content as dict/list, + # which crashes downstream .strip(); normalize to str. + if assistant_message.content is not None and not isinstance(assistant_message.content, str): + raw = assistant_message.content + if isinstance(raw, dict): + assistant_message.content = raw.get("text", "") or raw.get("content", "") or json.dumps(raw) + elif isinstance(raw, list): + # Multimodal content list — extract text parts + parts = [] + for part in raw: + if isinstance(part, str): + parts.append(part) + elif isinstance(part, dict) and part.get("type") == "text": + parts.append(part.get("text", "")) + elif isinstance(part, dict) and "text" in part: + parts.append(str(part["text"])) + assistant_message.content = "\n".join(parts) + else: + assistant_message.content = str(raw) + + # ── Agent-as-provider projection ────────────────────────────── + # Splice the provider-agent's own tool work in as call/result rows before + # this turn's assistant message; no-op for ordinary providers. + splice_provider_projection(agent, response, messages) + + try: + from hermes_cli.lifecycle import ( + has_hook, + invoke_hook as _invoke_hook, + ) + if has_hook("post_api_request"): + _assistant_tool_calls = ( + getattr(assistant_message, "tool_calls", None) or [] + ) + _assistant_text = assistant_message.content or "" + _api_ended_at = api_start_time + api_duration + _invoke_hook( + "post_api_request", + task_id=effective_task_id, + turn_id=turn_id, + api_request_id=api_request_id, + session_id=agent.session_id or "", + platform=agent.platform or "", + model=agent.model, + provider=agent.provider, + base_url=agent.base_url, + api_mode=agent.api_mode, + api_call_count=api_call_count, + api_duration=api_duration, + started_at=api_start_time, + ended_at=_api_ended_at, + # First stream chunk time (epoch s) from + # interruptible_streaming_api_call; None if not streamed / no + # chunk. TTFB = first_chunk_at - started_at. + first_chunk_at=getattr( + agent, "_last_api_first_chunk_at", None + ), + finish_reason=finish_reason, + message_count=len(api_messages), + response_model=getattr(response, "model", None), + response=agent._api_response_payload_for_hook( + response, + assistant_message, + finish_reason=finish_reason, + ), + usage=agent._usage_summary_for_api_request_hook(response), + assistant_message=assistant_message, + assistant_content_chars=len(_assistant_text), + assistant_tool_call_count=len(_assistant_tool_calls), + moa_references=_moa_reference_metrics_for_hook(agent), + ) + except Exception: + pass + + # Handle assistant response + if assistant_message.content and not agent.quiet_mode: + if agent.verbose_logging: + agent._vprint(f"{agent.log_prefix}🤖 Assistant: {assistant_message.content}") + else: + agent._vprint(f"{agent.log_prefix}🤖 Assistant: {assistant_message.content[:100]}{'...' if len(assistant_message.content) > 100 else ''}") + + # Notify progress callback of model's thinking (used by subagent + # delegation to relay the child's reasoning to the parent display). + if (assistant_message.content and agent.tool_progress_callback): + _think_text = assistant_message.content.strip() + # Strip reasoning XML tags that shouldn't leak to parent display + _think_text = re.sub( + r'', '', _think_text + ).strip() + # For subagents: relay first line to parent display (existing behaviour). + # For all agents with a structured callback: emit reasoning.available event. + first_line = _think_text.split('\n')[0][:80] if _think_text else "" + if first_line and getattr(agent, '_delegate_depth', 0) > 0: + try: + agent.tool_progress_callback("_thinking", first_line) + except Exception: + pass + elif _think_text: + try: + agent.tool_progress_callback("reasoning.available", "_thinking", _think_text[:500], None) + except Exception: + pass + + # Check for incomplete (opened but never closed) + # This means the model ran out of output tokens mid-reasoning — retry up to 2 times + if has_incomplete_scratchpad(assistant_message.content or ""): + agent._incomplete_scratchpad_retries += 1 + + agent._buffer_vprint("⚠️ Incomplete detected (opened but never closed)") + + if agent._incomplete_scratchpad_retries <= 2: + agent._buffer_vprint(f"🔄 Retrying API call ({agent._incomplete_scratchpad_retries}/2)...") + # Don't add the broken message, just retry + return _verdict("continue") + else: + # Max retries - discard this turn and save as partial + agent._flush_status_buffer() + agent._vprint(f"{agent.log_prefix}❌ Max retries (2) for incomplete scratchpad. Saving as partial.", force=True) + agent._incomplete_scratchpad_retries = 0 + + rolled_back_messages = agent._get_messages_up_to_last_assistant(messages) + agent._cleanup_task_resources(effective_task_id) + agent._persist_session(messages, conversation_history) + + return _verdict("return", { + "final_response": "Incomplete REASONING_SCRATCHPAD after 2 retries", + "messages": rolled_back_messages, + "api_calls": api_call_count, + "completed": False, + "partial": True, + "error": "Incomplete REASONING_SCRATCHPAD after 2 retries" + }) + + # Reset incomplete scratchpad counter on clean response + agent._incomplete_scratchpad_retries = 0 + + if agent.api_mode == "codex_responses" and finish_reason == "incomplete": + _codex_result = continue_codex_incomplete( + agent, + assistant_message, + finish_reason, + messages=messages, + conversation_history=conversation_history, + api_call_count=api_call_count, + ) + if _codex_result is not None: + return _verdict("return", _codex_result) + return _verdict("continue") + elif hasattr(agent, "_codex_incomplete_retries"): + agent._codex_incomplete_retries = 0 + return _verdict("fallthrough")