fix(api): explain interrupted provider streams with evidence

When a provider closes a 200 stream before message_stop and every stream
retry fails, the error now names the upstream model provider as the cause,
states it is not a context-limit error, and records what the attempt
received: event count, last event, stop_reason, whether message_stop was
missing, open blocks, elapsed time and retries spent. That separates a
reply cut mid-block from one that only dropped message_stop.

The error carries a new upstream_stream_interrupted business code. The
desktop shows a localized explanation in all five locales and keeps the
raw evidence line under it.
This commit is contained in:
程序员阿江(Relakkes)
2026-10-04 22:00:27 +08:00
parent 9623064b15
commit 36ba154a83
16 changed files with 237 additions and 4 deletions
@@ -8276,6 +8276,36 @@ describe('MessageList nested tool calls', () => {
expect(screen.queryByText(/This model does not support images/)).toBeNull()
})
it('blames the upstream provider for an interrupted stream and keeps the stream evidence', () => {
useSettingsStore.setState({ locale: 'zh' })
const raw =
'API Error: Provider stream ended before completing the response. The upstream model provider closed the stream before the reply finished; this is not a context-limit error. (retried 10 times · 412 events · last content_block_delta · stop_reason none · message_stop missing · 1 block open · 263s)'
useChatStore.setState({
sessions: {
[ACTIVE_TAB]: makeSessionState({
messages: [
{
id: 'error-1',
type: 'error',
code: 'unknown',
businessErrorCode: 'upstream_stream_interrupted',
message: raw,
timestamp: 1,
},
],
}),
},
})
render(<MessageList />)
expect(screen.getByText(/上游模型服务商在回复完成前断开了连接/)).toBeTruthy()
expect(screen.getByText(/不是上下文超限/)).toBeTruthy()
// The evidence is what tells a mid-block cut from a dropped message_stop
// when users send a screenshot; the translation alone would hide it.
expect(screen.getByText(raw)).toBeTruthy()
})
it.each([
['en', /too large for your provider or relay/, /removed automatically/],
['zh', /服务商或中转站允许的大小/, /自动移除/],
+8 -1
View File
@@ -152,6 +152,12 @@ type SelectionPointer = {
clientY: number
}
// Business errors whose raw message carries evidence worth keeping under the
// translated explanation (e.g. what the stream received before it was cut).
const BUSINESS_ERRORS_WITH_RAW_DETAIL: ReadonlySet<string> = new Set([
'upstream_stream_interrupted',
])
const CHAT_SELECTION_MENU_OFFSET = 10
const CHAT_SELECTION_MENU_WIDTH = 360
const CHAT_SELECTION_MENU_HEIGHT = 44
@@ -4103,7 +4109,8 @@ export const MessageBlock = memo(function MessageBlock({
? errorText
: message.message
const showRawDetail =
!message.businessErrorCode &&
(!message.businessErrorCode ||
BUSINESS_ERRORS_WITH_RAW_DETAIL.has(message.businessErrorCode)) &&
Boolean(message.message) &&
message.message.trim() !== '' &&
message.message !== displayMessage
+1
View File
@@ -3622,6 +3622,7 @@ Row 9, all 8 cells: continuing from straight down, turning left through lower-le
'businessError.request_too_large': 'The request is too large for your provider or relay. Earlier images and documents are removed automatically on your next message; if it still fails, compact the conversation or start a new session.',
'businessError.prompt_too_long': 'The prompt is too long for the selected model. Compact the conversation or retry with less context.',
'businessError.auto_mode_unavailable': 'Auto mode is unavailable for your current plan.',
'businessError.upstream_stream_interrupted': 'The upstream model provider closed the connection before the reply finished, and automatic retries did not recover it. This is usually a provider or network timeout or instability, not a context-limit error. Send "continue" to let the model pick up where it stopped, or retry later or switch providers.',
// ─── Server Status Verbs ──────────────────────────────────────
'serverVerb.Thinking': 'Thinking',
+1
View File
@@ -3623,6 +3623,7 @@ export const jp: Record<TranslationKey, string> = {
'businessError.request_too_large': 'リクエストがプロバイダーまたは中継サービスの許容サイズを超えました。次のメッセージで以前の画像とドキュメントが自動的に削除されます。それでも失敗する場合は、会話を圧縮するか、新しいセッションを開始してください。',
'businessError.prompt_too_long': 'プロンプトが選択したモデルには長すぎます。会話を圧縮するか、コンテキストを減らして再試行してください。',
'businessError.auto_mode_unavailable': '自動モードは現在のプランでは利用できません。',
'businessError.upstream_stream_interrupted': '上流のモデルプロバイダーが応答の完了前に接続を切断し、自動リトライでも回復しませんでした。通常はプロバイダーまたはネットワーク経路のタイムアウトや不安定さが原因で、コンテキスト上限のエラーではありません。「続けて」と送信して続きを依頼するか、時間をおいて再試行するか、プロバイダーを切り替えてください。',
// ─── Server Status Verbs ──────────────────────────────────────
'serverVerb.Thinking': '思考中',
+1
View File
@@ -3625,6 +3625,7 @@ export const kr: Record<TranslationKey, string> = {
'businessError.request_too_large': '요청이 공급자 또는 중계 서버가 허용하는 크기를 초과했습니다. 다음 메시지에서 이전 이미지와 문서가 자동으로 제거됩니다. 그래도 실패하면 대화를 압축하거나 새 세션을 시작하세요.',
'businessError.prompt_too_long': '프롬프트가 선택한 모델에 비해 너무 깁니다. 대화를 압축하거나 컨텍스트를 줄여 다시 시도하세요.',
'businessError.auto_mode_unavailable': '자동 모드는 현재 요금제에서 사용할 수 없습니다.',
'businessError.upstream_stream_interrupted': '상위 모델 제공업체가 응답이 끝나기 전에 연결을 끊었고, 자동 재시도로도 복구되지 않았습니다. 대개 제공업체나 네트워크 경로의 시간 초과 또는 불안정 때문이며, 컨텍스트 한도 오류가 아닙니다. "계속"을 보내 이어서 진행하게 하거나, 잠시 후 다시 시도하거나 제공업체를 바꿔 보세요.',
// ─── Server Status Verbs ──────────────────────────────────────
'serverVerb.Thinking': '사고 중',
+1
View File
@@ -3622,6 +3622,7 @@ export const zh: Record<TranslationKey, string> = {
'businessError.request_too_large': '這次請求超出了服務商或中轉站允許的大小。下一則訊息會自動移除較早的圖片和文件;如果仍然失敗,請壓縮會話或新開會話。',
'businessError.prompt_too_long': '當前上下文超出了模型限制。請先壓縮會話,或減少上下文後重試。',
'businessError.auto_mode_unavailable': '當前套餐不支援自動模式。',
'businessError.upstream_stream_interrupted': '上游模型服務商在回覆完成前中斷了連線,自動重試後仍未成功。這通常是服務商或網路線路逾時、不穩定造成的,不是上下文超限。可以傳送「繼續」讓模型接著做,或稍後重試、更換服務商。',
// ─── Server Status Verbs ──────────────────────────────────────
'serverVerb.Thinking': '思考中',
+1
View File
@@ -3621,6 +3621,7 @@ export const zh: Record<TranslationKey, string> = {
'businessError.request_too_large': '这次请求超出了服务商或中转站允许的大小。下一条消息会自动移除较早的图片和文档;如果仍然失败,请压缩会话或新开会话。',
'businessError.prompt_too_long': '当前上下文超出了模型限制。请先压缩会话,或减少上下文后重试。',
'businessError.auto_mode_unavailable': '当前套餐不支持自动模式。',
'businessError.upstream_stream_interrupted': '上游模型服务商在回复完成前断开了连接,自动重试后仍未成功。这通常是服务商或网络线路超时、不稳定造成的,不是上下文超限。可以发送「继续」让模型接着做,或稍后重试、更换服务商。',
// ─── Server Status Verbs ──────────────────────────────────────
'serverVerb.Thinking': '思考中',
+1
View File
@@ -7,6 +7,7 @@ export const BUSINESS_ERROR_CODES = {
REQUEST_TOO_LARGE: 'request_too_large',
PROMPT_TOO_LONG: 'prompt_too_long',
AUTO_MODE_UNAVAILABLE: 'auto_mode_unavailable',
UPSTREAM_STREAM_INTERRUPTED: 'upstream_stream_interrupted',
} as const
export type BusinessErrorCode =
@@ -1,7 +1,9 @@
import { describe, expect, spyOn, test } from 'bun:test'
import * as rateLimits from '../rateLimitMocking.js'
import { APIError } from '@anthropic-ai/sdk'
import { BUSINESS_ERROR_CODES } from '../../constants/businessErrors.js'
import { classifyAPIError, getAssistantMessageFromError } from './errors.js'
import { StreamEndedEarlyError } from './streamFallback.js'
function displayed(error: Error): string {
const result = getAssistantMessageFromError(error, 'claude-opus-5-5')
@@ -77,3 +79,31 @@ describe('specialized API error paths', () => {
expect(displayed(new APIError(400, { message: 'Bad field' }, undefined, undefined))).toBe('API Error: 400 Bad field')
})
})
describe('early stream end presentation', () => {
test('names the upstream provider, keeps the evidence and tags the business code', () => {
const error = new StreamEndedEarlyError('incomplete', {
eventCount: 412, lastEventType: 'content_block_delta', stopReason: null,
messageStopReceived: false, openBlockCount: 1, elapsedMs: 263_400,
})
const result = getAssistantMessageFromError(error, 'claude-opus-5-5', { streamRetries: 10 })
const text = result.message.content.map(block => block.type === 'text' ? block.text : '').join('')
expect(text.startsWith('API Error: Provider stream ended before completing the response. The upstream model provider')).toBe(true)
expect(text).toContain('retried 10 times')
expect(text).toContain('1 block open')
expect(result.isApiErrorMessage).toBe(true)
expect(result.error).toBe('unknown')
expect(result.businessErrorCode).toBe(BUSINESS_ERROR_CODES.UPSTREAM_STREAM_INTERRUPTED)
expect(JSON.parse(result.errorDetails!)).toEqual({
reason: 'incomplete', retries: 10, eventCount: 412, lastEventType: 'content_block_delta',
stopReason: null, messageStopReceived: false, openBlockCount: 1, elapsedMs: 263_400,
})
})
test('an empty stream without evidence is tagged the same way', () => {
const result = getAssistantMessageFromError(new StreamEndedEarlyError('no_events'), 'claude-opus-5-5')
expect(result.businessErrorCode).toBe(BUSINESS_ERROR_CODES.UPSTREAM_STREAM_INTERRUPTED)
expect(JSON.parse(result.errorDetails!)).toEqual({ reason: 'no_events', retries: 0 })
})
})
+21 -2
View File
@@ -230,6 +230,7 @@ import { isOpenAIPolicyError } from "../openaiAuth/policyError.js"
import {
shouldTriggerNonStreamingFallbackForEmptyStream,
StreamEndedEarlyError,
type StreamEndedEarlyEvidence,
} from "./streamFallback.js";
import { StreamAssistantCommitBuffer } from "./streamAssistantCommitBuffer.js";
import {
@@ -2784,6 +2785,24 @@ async function* queryModel(
);
}
// Recorded on an early-EOF error so the surfaced message shows whether
// the provider cut the reply mid-block or only dropped message_stop.
// Count open blocks before preservePartialText closes partial text.
const openBlockCountAtEof = contentBlocks.filter(
(block, index) => block && !completedBlockIndexes.has(index),
).length;
const streamEndEvidence = (): StreamEndedEarlyEvidence => {
const snapshot = streamWatchdogState.snapshot();
return {
eventCount: snapshot.eventCount,
lastEventType: snapshot.lastEventType ?? null,
stopReason,
messageStopReceived: snapshot.messageStopReceived,
openBlockCount: openBlockCountAtEof,
elapsedMs: Date.now() - start,
};
};
// Preserve visible text when the socket closes before its block_stop.
// Never synthesize a completed tool or unsigned thinking block from EOF.
preservePartialText();
@@ -2832,7 +2851,7 @@ async function* queryModel(
request_id: (streamRequestId ??
"unknown") as AnalyticsMetadata_I_VERIFIED_THIS_IS_NOT_CODE_OR_FILEPATHS,
});
throw new StreamEndedEarlyError("no_events");
throw new StreamEndedEarlyError("no_events", streamEndEvidence());
}
// A clean socket EOF is not a successful Anthropic response. Explicit
@@ -2843,7 +2862,7 @@ async function* queryModel(
(!streamWatchdogState.snapshot().messageStopReceived || stopReason === null ||
contentBlocks.some((block, index) => block && !completedBlockIndexes.has(index)))
) {
throw new StreamEndedEarlyError("incomplete");
throw new StreamEndedEarlyError("incomplete", streamEndEvidence());
}
// No tool boundary was crossed, so completed thinking/text blocks were
@@ -113,6 +113,21 @@ for (const block of ['text', 'tool_use']) {
expect(messages[0]?.isApiErrorMessage).not.toBe(true)
})
}
function upstreamInterruption(messages: AssistantMessage[]) {
const error = messages.find(message => message.isApiErrorMessage)
expect(error?.businessErrorCode).toBe('upstream_stream_interrupted')
return JSON.stringify(error?.message.content)
}
test('a reply cut mid-text reports the open block it was cut in', async () => {
const text = upstreamInterruption(await receiveResponse(fixtureEvents().slice(0, 3), true))
expect(text).toContain('The upstream model provider closed the stream')
expect(text).toContain('3 events · last content_block_delta · stop_reason none · message_stop missing · 1 block open')
})
test('a reply that only dropped message_stop reports the stop_reason it sent', async () => {
const text = upstreamInterruption(await receiveResponse([...fixtureEvents(), terminal()], true))
expect(text).toContain('5 events · last message_delta · stop_reason end_turn · message_stop missing')
expect(text).not.toContain('open')
})
for (const reason of ['max_tokens', 'model_context_window_exceeded', 'refusal']) {
test(`duplicate ${reason} terminal frame produces one terminal error`, async () => {
const messages = await receiveResponse([...fixtureEvents(), terminal(reason), terminal(reason), { type: 'message_stop' }])
@@ -295,6 +295,18 @@ test('retries are bounded, then the original error surfaces with the last attemp
expect(contentOf(errors)).toContain('ended without finish_reason')
}, 15_000)
test('a clean EOF that persists through every retry names the upstream and the retries spent', async () => {
const cutMidText = anthropicStream([messageStart, ...textBlock('cut-attempt').slice(0, 2)])
const { errors } = await run([cutMidText, cutMidText])
expect(requests).toHaveLength(2)
expect(errors).toHaveLength(1)
expect(errors[0]?.businessErrorCode).toBe('upstream_stream_interrupted')
expect(contentOf(errors)).toContain('The upstream model provider closed the stream')
expect(contentOf(errors)).toContain('retried 1 time · 3 events · last content_block_delta')
expect(JSON.parse(errors[0]!.errorDetails!)).toMatchObject({ reason: 'incomplete', retries: 1, openBlockCount: 1 })
}, 15_000)
for (const type of ['invalid_request_error', 'authentication_error', 'permission_error', 'billing_error']) {
test(`a mid-stream ${type} is not retried`, async () => {
const { errors } = await run([
+21
View File
@@ -51,6 +51,10 @@ import {
} from '../claudeAiLimits.js'
import { shouldProcessRateLimits } from '../rateLimitMocking.js' // Used for /mock-limits command
import { extractConnectionErrorDetails, formatAPIError } from './errorUtils.js'
import {
formatStreamEndedEarlyMessage,
StreamEndedEarlyError,
} from './streamFallback.js'
import { StreamWatchdogTimeoutError } from './streamWatchdog.js'
// Presentation only: classifiers, retries and diagnostic metadata keep the
@@ -659,6 +663,8 @@ export function getAssistantMessageFromError(
options?: {
messages?: Message[]
messagesForAPI?: (UserMessage | AssistantMessage)[]
/** Mid-stream re-sends already spent on this failure. */
streamRetries?: number
},
): AssistantMessage {
// Check for SDK timeout errors
@@ -1182,6 +1188,21 @@ export function getAssistantMessageFromError(
})
}
// The provider closed a 200 stream before finishing the reply. Name the
// upstream as the cause so it is not mistaken for a context-limit rejection.
if (error instanceof StreamEndedEarlyError) {
return createAssistantAPIErrorMessage({
content: `${API_ERROR_MESSAGE_PREFIX}: ${formatStreamEndedEarlyMessage(error, options?.streamRetries)}`,
error: 'unknown',
errorDetails: JSON.stringify({
reason: error.reason,
retries: options?.streamRetries ?? 0,
...error.evidence,
}),
businessErrorCode: BUSINESS_ERROR_CODES.UPSTREAM_STREAM_INTERRUPTED,
})
}
// Fallback for image rejections with unrecognized wording. The wording
// classifier above can't enumerate every provider/gateway phrasing, and a
// miss used to poison the session: the rejected image stayed in history and
+45
View File
@@ -1,5 +1,6 @@
import { describe, expect, test } from 'bun:test'
import {
formatStreamEndedEarlyMessage,
shouldTriggerNonStreamingFallbackForEmptyStream,
StreamEndedEarlyError,
} from './streamFallback.js'
@@ -14,6 +15,50 @@ describe('StreamEndedEarlyError', () => {
)
expect(new StreamEndedEarlyError('incomplete')).toBeInstanceOf(Error)
})
test('evidence does not change the message wording', () => {
const error = new StreamEndedEarlyError('incomplete', {
eventCount: 3, lastEventType: 'content_block_delta', stopReason: null,
messageStopReceived: false, openBlockCount: 1, elapsedMs: 1000,
})
expect(error.message).toBe('Provider stream ended before completing the response')
})
})
describe('formatStreamEndedEarlyMessage', () => {
const attribution =
'The upstream model provider closed the stream before the reply finished; this is not a context-limit error.'
test('a reply cut mid-block shows the open block and the missing terminal events', () => {
const error = new StreamEndedEarlyError('incomplete', {
eventCount: 412, lastEventType: 'content_block_delta', stopReason: null,
messageStopReceived: false, openBlockCount: 1, elapsedMs: 263_400,
})
expect(formatStreamEndedEarlyMessage(error, 10)).toBe(
`Provider stream ended before completing the response. ${attribution} ` +
'(retried 10 times · 412 events · last content_block_delta · stop_reason none · message_stop missing · 1 block open · 263s)',
)
})
test('a reply that only dropped message_stop shows the stop_reason it did send', () => {
const error = new StreamEndedEarlyError('incomplete', {
eventCount: 9, lastEventType: 'message_delta', stopReason: 'tool_use',
messageStopReceived: false, openBlockCount: 0, elapsedMs: 4_000,
})
expect(formatStreamEndedEarlyMessage(error, 1)).toBe(
`Provider stream ended before completing the response. ${attribution} ` +
'(retried 1 time · 9 events · last message_delta · stop_reason tool_use · message_stop missing · 4s)',
)
})
test('without evidence or retries only the wording and attribution remain', () => {
expect(formatStreamEndedEarlyMessage(new StreamEndedEarlyError('no_events'))).toBe(
`Stream ended without receiving any events. ${attribution}`,
)
expect(formatStreamEndedEarlyMessage(new StreamEndedEarlyError('no_events'), 0)).toBe(
`Stream ended without receiving any events. ${attribution}`,
)
})
})
describe('stream fallback policy', () => {
+48 -1
View File
@@ -14,6 +14,19 @@ export function shouldTriggerNonStreamingFallbackForEmptyStream({
export type StreamEndedEarlyReason = 'no_events' | 'incomplete'
/**
* What the attempt had received when the body closed. It tells a provider
* that cut the reply mid-block apart from one that only dropped message_stop.
*/
export type StreamEndedEarlyEvidence = {
eventCount: number
lastEventType: string | null
stopReason: string | null
messageStopReceived: boolean
openBlockCount: number
elapsedMs: number
}
/**
* A 200 SSE body that closed cleanly before the response completed: nothing
* usable arrived (`no_events`), or it stopped short of message_stop and a
@@ -21,7 +34,10 @@ export type StreamEndedEarlyReason = 'no_events' | 'incomplete'
* never ruled on the request, so — like a reset socket — a re-send can clear it.
*/
export class StreamEndedEarlyError extends Error {
constructor(readonly reason: StreamEndedEarlyReason) {
constructor(
readonly reason: StreamEndedEarlyReason,
readonly evidence?: StreamEndedEarlyEvidence,
) {
super(
reason === 'no_events'
? 'Stream ended without receiving any events'
@@ -30,3 +46,34 @@ export class StreamEndedEarlyError extends Error {
this.name = 'StreamEndedEarlyError'
}
}
/**
* User-facing text once the stream could not be recovered: the original
* wording, who closed the stream, and the evidence that shows how.
*/
export function formatStreamEndedEarlyMessage(
error: StreamEndedEarlyError,
retries?: number,
): string {
const details: string[] = []
if (retries !== undefined && retries > 0) {
details.push(`retried ${retries} ${retries === 1 ? 'time' : 'times'}`)
}
const evidence = error.evidence
if (evidence) {
details.push(
`${evidence.eventCount} ${evidence.eventCount === 1 ? 'event' : 'events'}`,
)
if (evidence.lastEventType) details.push(`last ${evidence.lastEventType}`)
details.push(`stop_reason ${evidence.stopReason ?? 'none'}`)
if (!evidence.messageStopReceived) details.push('message_stop missing')
if (evidence.openBlockCount > 0) {
details.push(
`${evidence.openBlockCount} ${evidence.openBlockCount === 1 ? 'block' : 'blocks'} open`,
)
}
details.push(`${Math.round(evidence.elapsedMs / 1000)}s`)
}
const summary = `${error.message}. The upstream model provider closed the stream before the reply finished; this is not a context-limit error.`
return details.length > 0 ? `${summary} (${details.join(' · ')})` : summary
}
+1
View File
@@ -109,6 +109,7 @@ export async function* withStreamRetry(
}
yield getAssistantMessageFromError(error.originalError, model, {
messages,
streamRetries: retries,
});
return;
}