From 5eecf5cbc5adf82f42b46a04f3ee9e4f24b9e507 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=A8=8B=E5=BA=8F=E5=91=98=E9=98=BF=E6=B1=9F=28Relakkes?= =?UTF-8?q?=29?= Date: Tue, 1 Sep 2026 21:19:18 +0800 Subject: [PATCH] fix(api): tolerate progressing tool input streams (#1271) --- .../__tests__/conversation-service.test.ts | 5 +- src/server/services/conversationService.ts | 6 +- src/services/api/claude.ts | 6 + .../api/claudeRequiredThinking.test.ts | 166 +++++++++++++++++- .../api/streamToolInputDurationGuard.test.ts | 25 ++- .../api/streamToolInputDurationGuard.ts | 5 + 6 files changed, 205 insertions(+), 8 deletions(-) diff --git a/src/server/__tests__/conversation-service.test.ts b/src/server/__tests__/conversation-service.test.ts index 4c13d9a8..9a709b25 100644 --- a/src/server/__tests__/conversation-service.test.ts +++ b/src/server/__tests__/conversation-service.test.ts @@ -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') diff --git a/src/server/services/conversationService.ts b/src/server/services/conversationService.ts index 7f04f7c1..cdfd27f0 100644 --- a/src/server/services/conversationService.ts +++ b/src/server/services/conversationService.ts @@ -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 diff --git a/src/services/api/claude.ts b/src/services/api/claude.ts index 8645a213..5f74d7d4 100644 --- a/src/services/api/claude.ts +++ b/src/services/api/claude.ts @@ -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") { diff --git a/src/services/api/claudeRequiredThinking.test.ts b/src/services/api/claudeRequiredThinking.test.ts index 8b508dfc..a15af10b 100644 --- a/src/services/api/claudeRequiredThinking.test.ts +++ b/src/services/api/claudeRequiredThinking.test.ts @@ -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 } diff --git a/src/services/api/streamToolInputDurationGuard.test.ts b/src/services/api/streamToolInputDurationGuard.test.ts index 3f1dd6cf..3de51066 100644 --- a/src/services/api/streamToolInputDurationGuard.test.ts +++ b/src/services/api/streamToolInputDurationGuard.test.ts @@ -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() + }) }) diff --git a/src/services/api/streamToolInputDurationGuard.ts b/src/services/api/streamToolInputDurationGuard.ts index bd39c705..752fa618 100644 --- a/src/services/api/streamToolInputDurationGuard.ts +++ b/src/services/api/streamToolInputDurationGuard.ts @@ -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