refactor(agent): one home for the stream parse-error markers

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.
This commit is contained in:
kshitijk4poor
2026-09-13 20:44:13 +05:30
committed by kshitij
parent 982e504262
commit 32a0d2d2e1
4 changed files with 11 additions and 8 deletions
+3 -1
View File
@@ -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):
+6 -3
View File
@@ -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)
-1
View File
@@ -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")
+2 -3
View File
@@ -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")