From 32a0d2d2e1dd45019743e04468f126dfcddd386d Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Sun, 13 Sep 2026 20:44:13 +0530 Subject: [PATCH] refactor(agent): one home for the stream parse-error markers MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The SDK's malformed-frame ValueError substrings were hand-copied into the retry classifier and the error summary; a third SDK message would need both remembered. The mid-tool retry branch no longer clears partial_tool_names itself — _start_stream_attempt resets it for every attempt. buffer_anthropic_tool_input's docstring states the trade: the knob stays on the turn's kwargs and the changed tools block costs one prompt-cache miss. --- agent/anthropic_adapter.py | 4 +++- agent/api_error_summary.py | 9 ++++++--- agent/chat_completion_helpers.py | 1 - run_agent.py | 5 ++--- 4 files changed, 11 insertions(+), 8 deletions(-) diff --git a/agent/anthropic_adapter.py b/agent/anthropic_adapter.py index 82cc9d5643..f118e2ddb6 100644 --- a/agent/anthropic_adapter.py +++ b/agent/anthropic_adapter.py @@ -628,7 +628,9 @@ def buffer_anthropic_tool_input(api_kwargs: dict[str, Any], base_url: str | None """Retry knob for a malformed fine-grained tool-JSON stream (#107830): the beta streams tool args unvalidated, so a model that emits ``{"names": cronjob_manage}`` breaks the SDK parser and an identical retry breaks identically. ``eager_input_streaming: false`` per tool restores - Anthropic's buffered, validated args for THIS request only. Off the happy path on purpose: + Anthropic's buffered, validated args for the rest of this turn (the flag lives on the turn's + kwargs, so a later retry of the same turn keeps it; the changed ``tools`` block costs one + prompt-cache miss, cheaper than a dead turn). Off the happy path on purpose: buffering a large payload is a zero-event gap the stale-stream detector kills. No-op on endpoints that never get the beta (MiniMax) rather than sending them an unknown field.""" if _TOOL_STREAMING_BETA not in _common_betas_for_base_url(base_url): diff --git a/agent/api_error_summary.py b/agent/api_error_summary.py index 5a1e884c27..1217be328c 100644 --- a/agent/api_error_summary.py +++ b/agent/api_error_summary.py @@ -10,6 +10,11 @@ from typing import Any, Dict, Optional from agent.redact import redact_sensitive_text +# Substrings of the plain ``ValueError`` the Anthropic SDK raises for a malformed event-stream +# frame (wire trouble, not local validation). Read by ``AIAgent._is_provider_stream_parse_error``. +PROVIDER_STREAM_PARSE_MARKERS = ("expected ident at line", "expected value at line") + + # Offline DNS failures are wrapped in a generic "Connection error" by SDKs — inspect the chain. _NETWORK_RESOLUTION_MARKERS = ( "temporary failure in name resolution", @@ -142,9 +147,7 @@ class ApiErrorSummaryMixin: ) current = current.__cause__ or current.__context__ - if isinstance(error, ValueError) and any( - marker in raw.lower() for marker in ("expected ident at line", "expected value at line") - ): + if isinstance(error, ValueError) and any(marker in raw.lower() for marker in PROVIDER_STREAM_PARSE_MARKERS): return f"Malformed provider streaming response: {raw[:300]}" prefix = _http_prefix(error) diff --git a/agent/chat_completion_helpers.py b/agent/chat_completion_helpers.py index 6c33b7cf4b..e59e238891 100644 --- a/agent/chat_completion_helpers.py +++ b/agent/chat_completion_helpers.py @@ -3171,7 +3171,6 @@ class _StreamingCall(StreamingWaitMonitor): # reset the streamed-text buffer so it isn't double-recorded; fresh accumulators. self._quiet(self.agent._fire_stream_delta, "\n\n⚠ Connection dropped mid tool-call; reconnecting…\n\n") self._quiet(self.agent._reset_stream_delivery_tracking) - self.result["partial_tool_names"] = [] self.deltas_were_sent["yes"] = False self.first_delta_fired["done"] = False self._retry_after_drop(e, attempt, max_retries, mid_tool_call=True, reason="stream_mid_tool_retry_cleanup") diff --git a/run_agent.py b/run_agent.py index 3174c413ed..5744825cfa 100644 --- a/run_agent.py +++ b/run_agent.py @@ -109,7 +109,7 @@ from agent.client_lifecycle import ClientLifecycleMixin from agent.stream_delivery import StreamDeliveryMixin from agent.status_output import StatusOutputMixin from agent.api_request_hooks import ApiRequestHooksMixin -from agent.api_error_summary import ApiErrorSummaryMixin +from agent.api_error_summary import PROVIDER_STREAM_PARSE_MARKERS, ApiErrorSummaryMixin from agent.interrupt_control import InterruptControlMixin from agent.turn_explainers import TurnExplainersMixin from agent.activity_tracking import ActivityTrackingMixin @@ -485,8 +485,7 @@ class AIAgent( that is wire trouble, not local validation, so it follows the truncated-JSON retry path.""" return (getattr(self, "api_mode", None) == "anthropic_messages" and isinstance(error, ValueError) and not isinstance(error, (UnicodeEncodeError, json.JSONDecodeError)) - and any(marker in str(error).strip().lower() for marker in ( - "expected ident at line", "expected value at line"))) + and any(marker in str(error).strip().lower() for marker in PROVIDER_STREAM_PARSE_MARKERS)) _log_stream_retry = _forward("agent.stream_diag", "log_stream_retry") _emit_stream_drop = _forward("agent.stream_diag", "emit_stream_drop")