fix(api): tolerate progressing tool input streams (#1271)

This commit is contained in:
程序员阿江(Relakkes)
2026-09-01 21:19:18 +08:00
parent c0f3318d9f
commit 5eecf5cbc5
6 changed files with 205 additions and 8 deletions
@@ -256,9 +256,8 @@ describe('ConversationService', () => {
// 240s apart keeps it alive forever. The overall-duration cap is NOT reset
// by chunks and is what actually frees that case (#766).
expect(env.CLAUDE_STREAM_MAX_DURATION_MS).toBe('600000')
// Tool JSON is not user-visible and should normally finish in seconds.
// Bound it separately so a model cannot spend the full response budget
// streaming a truncated Write payload (#1237).
// Tool JSON gets a shorter inactivity budget. Progress resets it, while
// the overall response cap still bounds a stream that trickles forever.
expect(env.CLAUDE_STREAM_TOOL_INPUT_MAX_DURATION_MS).toBe('120000')
// Non-streaming fallback stays off — its retry loop also hangs the UI (#766).
expect(env.CLAUDE_CODE_DISABLE_NONSTREAMING_FALLBACK).toBe('1')
+3 -3
View File
@@ -1640,9 +1640,9 @@ export class ConversationService {
// no completion (#766: "卡住" with slowly growing tokens). This independent
// cap frees such a stream after a fixed duration regardless of trickle.
CLAUDE_STREAM_MAX_DURATION_MS: cleanEnv.CLAUDE_STREAM_MAX_DURATION_MS || '600000',
// A local tool call should finish generating its JSON arguments quickly.
// Bound this separately from the full response so a truncated, continuously
// streaming Write payload cannot occupy the session for the full 10 minutes.
// Abort a local tool call when its JSON arguments stop making progress.
// Healthy input_json_delta events reset this budget; the independent full
// response cap above still bounds a stream that trickles forever.
CLAUDE_STREAM_TOOL_INPUT_MAX_DURATION_MS:
cleanEnv.CLAUDE_STREAM_TOOL_INPUT_MAX_DURATION_MS || '120000',
// Time-to-first-token budget: how long to wait for the FIRST streamed
+6
View File
@@ -2064,6 +2064,9 @@ async function* queryModel(
// 0 disables it (terminal CLI default); the desktop injects a value.
const STREAM_MAX_DURATION_MS =
parseInt(process.env.CLAUDE_STREAM_MAX_DURATION_MS || "", 10) || 0;
// Despite the legacy env name, this is an inactivity budget: each
// non-empty input_json_delta proves the tool call is still progressing.
// STREAM_MAX_DURATION_MS remains the independent total response cap.
const STREAM_TOOL_INPUT_MAX_DURATION_MS =
parseInt(
process.env.CLAUDE_STREAM_TOOL_INPUT_MAX_DURATION_MS || "",
@@ -2400,6 +2403,9 @@ async function* queryModel(
throw new Error("Content block input is not a string");
}
contentBlock.input += delta.partial_json;
if (delta.partial_json.length > 0) {
toolInputDurationGuard.progress(part.index);
}
break;
case "text_delta":
if (contentBlock.type !== "text") {
+165 -1
View File
@@ -143,6 +143,121 @@ function hangingToolResponse(model: string): Response {
})
}
function progressingToolResponse(model: string): Response {
const events = [
sseEvent('message_start', {
type: 'message_start',
message: {
id: 'msg_progressing_tool',
type: 'message',
role: 'assistant',
model,
content: [],
stop_reason: null,
stop_sequence: null,
usage: { input_tokens: 1, output_tokens: 0 },
},
}),
sseEvent('content_block_start', {
type: 'content_block_start',
index: 0,
content_block: {
type: 'tool_use',
id: 'tool_progressing_bash',
name: 'Bash',
input: {},
},
}),
...['{', '"command"', ':', '"echo OK"', '}'].map(partialJson =>
sseEvent('content_block_delta', {
type: 'content_block_delta',
index: 0,
delta: { type: 'input_json_delta', partial_json: partialJson },
}),
),
sseEvent('content_block_stop', { type: 'content_block_stop', index: 0 }),
sseEvent('message_delta', {
type: 'message_delta',
delta: { stop_reason: 'tool_use', stop_sequence: null },
usage: { output_tokens: 5 },
}),
sseEvent('message_stop', { type: 'message_stop' }),
]
let nextEvent = 0
let cancelled = false
return new Response(new ReadableStream({
async pull(controller) {
if (nextEvent > 0) await Bun.sleep(10)
if (cancelled) return
controller.enqueue(new TextEncoder().encode(events[nextEvent]))
nextEvent += 1
if (nextEvent === events.length) controller.close()
},
cancel() {
cancelled = true
},
}), {
headers: { 'content-type': 'text/event-stream' },
})
}
function tricklingToolResponse(model: string): Response {
const initialEvents = [
sseEvent('message_start', {
type: 'message_start',
message: {
id: 'msg_trickling_tool',
type: 'message',
role: 'assistant',
model,
content: [],
stop_reason: null,
stop_sequence: null,
usage: { input_tokens: 1, output_tokens: 0 },
},
}),
sseEvent('content_block_start', {
type: 'content_block_start',
index: 0,
content_block: {
type: 'tool_use',
id: 'tool_trickling_write',
name: 'Write',
input: {},
},
}),
sseEvent('content_block_delta', {
type: 'content_block_delta',
index: 0,
delta: { type: 'input_json_delta', partial_json: '{"content":"' },
}),
].join('')
const progressEvent = sseEvent('content_block_delta', {
type: 'content_block_delta',
index: 0,
delta: { type: 'input_json_delta', partial_json: 'x' },
})
let sentInitialEvents = false
let cancelled = false
return new Response(new ReadableStream({
async pull(controller) {
if (sentInitialEvents) await Bun.sleep(10)
if (cancelled) return
controller.enqueue(new TextEncoder().encode(
sentInitialEvents ? progressEvent : initialEvents,
))
sentInitialEvents = true
},
cancel() {
cancelled = true
},
}), {
headers: { 'content-type': 'text/event-stream' },
})
}
const ENV_KEYS = [
'NODE_ENV',
'CLAUDE_CONFIG_DIR',
@@ -371,7 +486,7 @@ test('drops a tool call truncated at the output-token boundary', async () => {
expect(error).toBe('max_output_tokens')
}, 10_000)
test('aborts a tool input that keeps streaming without completing', async () => {
test('aborts a tool input that stops progressing without completing', async () => {
const { content, error, requests } = await captureQueryRequest({
model: 'deepseek-v4-flash',
configureCapabilityOverrides: false,
@@ -395,6 +510,55 @@ test('aborts a tool input that keeps streaming without completing', async () =>
expect(error).toBe('server_error')
}, 10_000)
test('allows a progressing tool input to outlive its inactivity budget', async () => {
const { content, error, requests } = await captureQueryRequest({
model: 'deepseek-v4-flash',
configureCapabilityOverrides: false,
responseFactory: progressingToolResponse,
env: {
CLAUDE_CODE_DISABLE_NONSTREAMING_FALLBACK: '1',
CLAUDE_ENABLE_STREAM_WATCHDOG: '1',
CLAUDE_STREAM_IDLE_TIMEOUT_MS: '1000',
CLAUDE_STREAM_MAX_DURATION_MS: '1000',
CLAUDE_STREAM_TOOL_INPUT_MAX_DURATION_MS: '40',
},
})
expect(requests).toHaveLength(1)
expect(content).toEqual([
expect.objectContaining({
type: 'tool_use',
name: 'Bash',
input: { command: 'echo OK' },
}),
])
expect(error).toBeUndefined()
}, 10_000)
test('bounds a tool input that keeps progressing but never completes', async () => {
const { content, error, requests } = await captureQueryRequest({
model: 'deepseek-v4-flash',
configureCapabilityOverrides: false,
responseFactory: tricklingToolResponse,
env: {
CLAUDE_CODE_DISABLE_NONSTREAMING_FALLBACK: '1',
CLAUDE_ENABLE_STREAM_WATCHDOG: '1',
CLAUDE_STREAM_IDLE_TIMEOUT_MS: '1000',
CLAUDE_STREAM_MAX_DURATION_MS: '200',
CLAUDE_STREAM_TOOL_INPUT_MAX_DURATION_MS: '100',
},
})
expect(requests).toHaveLength(1)
expect(content).toEqual([
expect.objectContaining({
type: 'text',
text: expect.stringContaining('Stream max duration exceeded'),
}),
])
expect(error).toBe('server_error')
}, 10_000)
function clearCapabilityCache() {
;(get3PModelCapabilityOverride as typeof get3PModelCapabilityOverride & {
cache?: { clear?: () => void }
@@ -2,7 +2,7 @@ import { describe, expect, test } from 'bun:test'
import { StreamToolInputDurationGuard } from './streamToolInputDurationGuard.js'
describe('StreamToolInputDurationGuard', () => {
test('fires once when a tool input exceeds its generation budget', async () => {
test('fires once when a tool input makes no progress before its deadline', async () => {
const timedOut: number[] = []
const guard = new StreamToolInputDurationGuard({
enabled: true,
@@ -39,4 +39,27 @@ describe('StreamToolInputDurationGuard', () => {
guard.clear()
disabled.clear()
})
test('allows continued progress but times out after progress stops', async () => {
const timedOut: number[] = []
const guard = new StreamToolInputDurationGuard({
enabled: true,
timeoutMs: 100,
onTimeout: index => timedOut.push(index),
})
guard.start(4)
await Bun.sleep(30)
guard.progress(4)
await Bun.sleep(30)
guard.progress(4)
await Bun.sleep(50)
expect(timedOut).toEqual([])
await Bun.sleep(70)
expect(timedOut).toEqual([4])
guard.clear()
})
})
@@ -18,6 +18,11 @@ export class StreamToolInputDurationGuard {
}, this.options.timeoutMs))
}
progress(index: number): void {
if (!this.timers.has(index)) return
this.start(index)
}
stop(index: number): void {
const timer = this.timers.get(index)
if (timer === undefined) return