Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 0 additions & 2 deletions src/__tests__/bugfixes.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -59,8 +59,6 @@ describe('Session timeout fix', () => {
const content = await file('services/api/openaiShim.ts').text()

expect(content).toContain('STREAM_IDLE_TIMEOUT_MS')
expect(content).toContain('readWithTimeout')
expect(content).toMatch(/readWithTimeout\(\)/)
})

test('codexShim has idle timeout for SSE streams', async () => {
Expand Down
289 changes: 289 additions & 0 deletions src/services/api/claude.lifecycle.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import {
acquireSharedMutationLock,
releaseSharedMutationLock,
} from '../../test/sharedMutationLock.js'
import { resetGrowthBook } from '../analytics/growthbook.js'
import { getEmptyToolPermissionContext } from '../../Tool.js'
import type { Message } from '../../types/message.js'
import { QueryLifecycleOperationTracker } from '../../utils/queryLifecycle.js'
Expand All @@ -36,6 +37,9 @@ const envKeys = [
'CLAUDE_CODE_USE_MISTRAL',
'CLAUDE_CODE_USE_OPENAI',
'CLAUDE_CODE_USE_VERTEX',
'CLAUDE_CODE_DISABLE_NONSTREAMING_FALLBACK',
'CLAUDE_FEATURE_FLAGS_FILE',
'CLAUDE_STREAM_IDLE_TIMEOUT_MS',
'GEMINI_API_KEY',
'OPENAI_API_KEY',
'OPENAI_BASE_URL',
Expand All @@ -51,6 +55,9 @@ let fixturesRoot: string | undefined

type FetchOverride = NonNullable<Options['fetchOverride']>
type LifecycleSnapshot = ReturnType<QueryLifecycleOperationTracker['snapshot']>
const TEST_STREAM_IDLE_TIMEOUT_MS = 25
const STREAM_IDLE_RECOVERY_ASSERTION_MS = 1_000
const STALLING_STREAM_CLEANUP_MS = 2_000
Comment thread
coderabbitai[bot] marked this conversation as resolved.

function makeJsonResponse(body: unknown, status = 200): Response {
return new Response(JSON.stringify(body), {
Expand Down Expand Up @@ -119,6 +126,98 @@ function makeOpenAIChatCompletionResponse(): Response {
})
}

function makeOpenAIStreamChunk(
delta: Record<string, unknown>,
finishReason: string | null = null,
): string {
return `data: ${JSON.stringify({
id: 'chatcmpl-lifecycle-stream',
object: 'chat.completion.chunk',
created: 1_771_264_800,
model: 'glm-5.2',
choices: [{ index: 0, delta, finish_reason: finishReason }],
})}\n\n`
}

function makeStallingOpenAIStreamResponse(
onCancel?: (reason: unknown) => void,
): Response {
const encoder = new TextEncoder()
let closeTimer: ReturnType<typeof setTimeout> | undefined

return new Response(
new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(
encoder.encode(
makeOpenAIStreamChunk({ role: 'assistant', content: 'partial' }),
),
)
// Bounded cleanup for current/baseline behavior: the idle-timeout
// assertions should fail before this close fires.
closeTimer = setTimeout(() => {
try {
controller.close()
} catch {
// stream may already be cancelled by the idle timeout path
}
}, STALLING_STREAM_CLEANUP_MS)
},
cancel(reason) {
if (closeTimer !== undefined) {
clearTimeout(closeTimer)
}
onCancel?.(reason)
},
}),
{
headers: {
'content-type': 'text/event-stream',
},
},
)
}

function makeRoleOnlyStallingOpenAIStreamResponse(
onInitialChunk: () => void,
onCancel?: (reason: unknown) => void,
): Response {
const encoder = new TextEncoder()
let closeTimer: ReturnType<typeof setTimeout> | undefined
let sentInitialChunk = false

return new Response(
new ReadableStream<Uint8Array>({
pull(controller) {
if (sentInitialChunk) return
sentInitialChunk = true
controller.enqueue(
encoder.encode(makeOpenAIStreamChunk({ role: 'assistant' })),
)
onInitialChunk()
closeTimer = setTimeout(() => {
try {
controller.close()
} catch {
// stream may already be cancelled by the abort path
}
}, 500)
},
cancel(reason) {
if (closeTimer !== undefined) {
clearTimeout(closeTimer)
}
onCancel?.(reason)
},
}),
{
headers: {
'content-type': 'text/event-stream',
},
},
)
}

function parseRequestBody(init: RequestInit | undefined): Record<string, unknown> {
if (typeof init?.body !== 'string') return {}
const parsed = JSON.parse(init.body) as unknown
Expand All @@ -136,6 +235,26 @@ async function drainGenerator<T>(
}
}

async function waitForPromise<T>(
promise: Promise<T>,
timeoutMs: number,
timeoutMessage: string,
): Promise<T> {
let timeoutId: ReturnType<typeof setTimeout> | undefined
try {
return await Promise.race([
promise,
new Promise<never>((_, reject) => {
timeoutId = setTimeout(() => {
reject(new Error(timeoutMessage))
}, timeoutMs)
}),
])
} finally {
if (timeoutId !== undefined) clearTimeout(timeoutId)
}
}

function makeParams(context: { model: string }): BetaMessageStreamParams {
return {
model: context.model,
Expand Down Expand Up @@ -178,7 +297,12 @@ function setClientTestEnv(): void {
}
process.env.ANTHROPIC_API_KEY = 'sk-test-lifecycle'
process.env.CLAUDE_CODE_TEST_FIXTURES_ROOT = fixturesRoot
process.env.CLAUDE_FEATURE_FLAGS_FILE = join(
fixturesRoot,
'feature-flags.json',
)
process.env.VCR_RECORD = '1'
resetGrowthBook()
}

beforeEach(async () => {
Expand All @@ -200,6 +324,7 @@ afterEach(() => {
delete (globalThis as Record<string, unknown>).MACRO
}
globalThis.fetch = originalFetch
resetGrowthBook()
if (fixturesRoot) {
rmSync(fixturesRoot, { force: true, recursive: true })
fixturesRoot = undefined
Expand Down Expand Up @@ -333,6 +458,170 @@ describe('Claude API lifecycle tracking', () => {
expect(queryLifecycle.snapshot().apiCalls).toEqual([])
})

test('parent abort during OpenAI-compatible stream does not start non-streaming fallback', async () => {
setClientTestEnv()
process.env.OPENCLAUDE_MAX_RETRIES = '0'
process.env.CLAUDE_STREAM_IDLE_TIMEOUT_MS = '1000'
const queryLifecycle = new QueryLifecycleOperationTracker()
const parent = new AbortController()
let fallbackRequests = 0
let fallbackNotifications = 0
let streamCancelled = false
const messages: unknown[] = []
let resolveStreamingRequestStarted!: () => void
const streamingRequestStarted = new Promise<void>(resolve => {
resolveStreamingRequestStarted = resolve
})
let resolveInitialStreamChunk!: () => void
const initialStreamChunk = new Promise<void>(resolve => {
resolveInitialStreamChunk = resolve
})

globalThis.fetch = (async (_input, init) => {
const body = parseRequestBody(init)
if (body.stream === true) {
resolveStreamingRequestStarted()
return makeRoleOnlyStallingOpenAIStreamResponse(
resolveInitialStreamChunk,
() => {
streamCancelled = true
},
)
}
fallbackRequests++
return makeOpenAIChatCompletionResponse()
}) as typeof fetch

let drainError: unknown
const drain = (async () => {
try {
const generator = queryModelWithStreaming({
messages: [
{
type: 'user',
uuid: '00000000-0000-0000-0000-000000000006',
timestamp: '2026-06-17T00:00:00.000Z',
message: { role: 'user', content: 'hello' },
} as Message,
],
systemPrompt: asSystemPrompt([]),
thinkingConfig: { type: 'disabled' },
tools: [],
signal: parent.signal,
options: {
...makeOptions(queryLifecycle),
providerOverride: {
model: 'glm-5.2',
baseURL: 'https://provider.example/v1',
apiKey: 'provider-test-key',
},
onStreamingFallback: () => {
fallbackNotifications++
},
},
})

for await (const message of generator) {
messages.push(message)
}
} catch (error) {
drainError = error
}
})()

await streamingRequestStarted
await initialStreamChunk
await Promise.resolve()
parent.abort()

await drain

expect(drainError).toBeUndefined()
expect(fallbackRequests).toBe(0)
expect(fallbackNotifications).toBe(0)
expect(streamCancelled).toBe(true)
expect(
messages.some(
message =>
typeof message === 'object' &&
message !== null &&
(message as { type?: unknown }).type === 'assistant',
),
).toBe(false)
})

test('stream idle timeout respects disabled non-streaming fallback guard', async () => {
setClientTestEnv()
process.env.OPENCLAUDE_MAX_RETRIES = '0'
process.env.CLAUDE_STREAM_IDLE_TIMEOUT_MS = String(TEST_STREAM_IDLE_TIMEOUT_MS)
process.env.CLAUDE_CODE_DISABLE_NONSTREAMING_FALLBACK = '1'
const queryLifecycle = new QueryLifecycleOperationTracker()
const parent = new AbortController()
let fallbackRequests = 0
const messages: unknown[] = []
const startedAt = Date.now()

globalThis.fetch = (async (_input, init) => {
const body = parseRequestBody(init)
if (body.stream === true) {
return makeStallingOpenAIStreamResponse()
}
fallbackRequests++
return makeOpenAIChatCompletionResponse()
}) as typeof fetch

let drainError: unknown
const drain = (async () => {
try {
const generator = queryModelWithStreaming({
messages: [
{
type: 'user',
uuid: '00000000-0000-0000-0000-000000000007',
timestamp: '2026-06-17T00:00:00.000Z',
message: { role: 'user', content: 'hello' },
} as Message,
],
systemPrompt: asSystemPrompt([]),
thinkingConfig: { type: 'disabled' },
tools: [],
signal: parent.signal,
options: {
...makeOptions(queryLifecycle),
providerOverride: {
model: 'glm-5.2',
baseURL: 'https://provider.example/v1',
apiKey: 'provider-test-key',
},
},
})

for await (const message of generator) {
messages.push(message)
}
} catch (error) {
drainError = error
}
})()

await drain
expect(Date.now() - startedAt).toBeLessThan(
TEST_STREAM_IDLE_TIMEOUT_MS + STREAM_IDLE_RECOVERY_ASSERTION_MS,
)

expect(drainError).toBeUndefined()
expect(fallbackRequests).toBe(0)
expect(
messages.some(
message =>
typeof message === 'object' &&
message !== null &&
(message as { type?: unknown }).type === 'assistant' &&
JSON.stringify((message as { message?: { content?: unknown } }).message?.content).includes('Stream idle timeout'),
),
).toBe(true)
})

test('tracks each non-streaming fallback request and clears it on success', async () => {
const queryLifecycle = new QueryLifecycleOperationTracker()
const requestSnapshots: ReturnType<
Expand Down
Loading