diff --git a/frontend/src/lib/components/copilot/chat/AIChatManager.svelte.ts b/frontend/src/lib/components/copilot/chat/AIChatManager.svelte.ts index 7bf9685e87..8b0c346244 100644 --- a/frontend/src/lib/components/copilot/chat/AIChatManager.svelte.ts +++ b/frontend/src/lib/components/copilot/chat/AIChatManager.svelte.ts @@ -52,7 +52,7 @@ import { loadApiTools } from './api/apiTools' import { prepareScriptUserMessage } from './script/core' import { prepareNavigatorUserMessage } from './navigator/core' import { sendUserToast } from '$lib/toast' -import { workspaceAIClients, getNonStreamingCompletion } from '../lib' +import { workspaceAIClients, getStreamedCompletionText } from '../lib' import { logFeatureUsage } from '$lib/utils/featureUsage' import { modelSupportsVision } from '../modelConfig' import { getEffectiveModelContextWindow } from '../modelConfig' @@ -1587,10 +1587,8 @@ export class AIChatManager implements ChatViewHost { this.compacting = true try { // Cap the summarizer's output at the budget already reserved for the - // summary. Without a cap the model's default max_tokens applies, and the - // Anthropic SDK rejects non-streaming requests whose max_tokens implies - // >10 minutes of generation (~21k tokens) before anything is sent. - const raw = await getNonStreamingCompletion( + // summary: without a cap the model's default max_tokens applies. + const raw = await getStreamedCompletionText( [ // Strip image blobs from the summarizer input — the summary text stands in // for them, so re-sending base64 to the summarizer only wastes tokens. @@ -1600,7 +1598,7 @@ export class AIChatManager implements ChatViewHost { abortController, { maxTokensCap: SUMMARY_OUTPUT_RESERVE_TOKENS } ) - const formatted = formatCompactSummary(raw ?? '') + const formatted = formatCompactSummary(raw) if (!formatted) { return 'empty' } diff --git a/frontend/src/lib/components/copilot/chat/AIChatManager.test.ts b/frontend/src/lib/components/copilot/chat/AIChatManager.test.ts index 97a1afb2ff..dbb3ac51e3 100644 --- a/frontend/src/lib/components/copilot/chat/AIChatManager.test.ts +++ b/frontend/src/lib/components/copilot/chat/AIChatManager.test.ts @@ -33,7 +33,7 @@ const mocks = vi.hoisted(() => ({ sendUserToast: vi.fn(), getOpenaiClient: vi.fn(), getAnthropicClient: vi.fn(), - getNonStreamingCompletion: vi.fn(), + getStreamedCompletionText: vi.fn(), runChatLoop: vi.fn(), listResource: vi.fn(), getJob: vi.fn(), @@ -143,7 +143,7 @@ vi.mock('../lib', () => ({ getOpenaiClient: mocks.getOpenaiClient, getAnthropicClient: mocks.getAnthropicClient }, - getNonStreamingCompletion: mocks.getNonStreamingCompletion, + getStreamedCompletionText: mocks.getStreamedCompletionText, providerSupportsWebSearch: (provider: string) => provider === 'openai' || provider === 'anthropic' })) @@ -3350,7 +3350,7 @@ describe('AIChatManager context compaction', () => { it('summarizes the older prefix and keeps the recent tail verbatim', async () => { mocks.getCurrentModel.mockReturnValue(gpt4oModel) mocks.tryGetCurrentModel.mockReturnValue(gpt4oModel) - mocks.getNonStreamingCompletion.mockResolvedValue( + mocks.getStreamedCompletionText.mockResolvedValue( 'scratchpadSUMMARY TEXT' ) const manager = new AIChatManager() @@ -3360,8 +3360,8 @@ describe('AIChatManager context compaction', () => { // The prefix (the four OLD messages) was sent to the summarizer, followed // by the summary-instruction user message. - expect(mocks.getNonStreamingCompletion).toHaveBeenCalledTimes(1) - const summaryReq = mocks.getNonStreamingCompletion.mock.calls[0][0] + expect(mocks.getStreamedCompletionText).toHaveBeenCalledTimes(1) + const summaryReq = mocks.getStreamedCompletionText.mock.calls[0][0] expect(summaryReq).toHaveLength(5) expect(summaryReq[0].content).toContain('OLD1') expect(summaryReq[3].content).toContain('OLD4') @@ -3392,7 +3392,7 @@ describe('AIChatManager context compaction', () => { it('carries folded-away message files on the summary', async () => { mocks.getCurrentModel.mockReturnValue(gpt4oModel) mocks.tryGetCurrentModel.mockReturnValue(gpt4oModel) - mocks.getNonStreamingCompletion.mockResolvedValue( + mocks.getStreamedCompletionText.mockResolvedValue( 'sSUM' ) const manager = new AIChatManager() @@ -3423,7 +3423,7 @@ describe('AIChatManager context compaction', () => { it('never lands the tail boundary on a screenshot follow-up that has no display counterpart', async () => { mocks.getCurrentModel.mockReturnValue(gpt4oModel) mocks.tryGetCurrentModel.mockReturnValue(gpt4oModel) - mocks.getNonStreamingCompletion.mockResolvedValue('SUMMARY TEXT') + mocks.getStreamedCompletionText.mockResolvedValue('SUMMARY TEXT') const manager = new AIChatManager() manager.messages = [ { role: 'user', content: 'OLD1' + 'a'.repeat(100_000) }, @@ -3471,7 +3471,7 @@ describe('AIChatManager context compaction', () => { it('falls back to drop-oldest when summarization fails', async () => { mocks.getCurrentModel.mockReturnValue(gpt4oModel) mocks.tryGetCurrentModel.mockReturnValue(gpt4oModel) - mocks.getNonStreamingCompletion.mockRejectedValue(new Error('summary boom')) + mocks.getStreamedCompletionText.mockRejectedValue(new Error('summary boom')) const manager = new AIChatManager() seedForSummary(manager) @@ -3479,7 +3479,7 @@ describe('AIChatManager context compaction', () => { // Summarization was attempted, then the request still went out — via // drop-oldest, so no summary boundary anywhere. - expect(mocks.getNonStreamingCompletion).toHaveBeenCalledTimes(1) + expect(mocks.getStreamedCompletionText).toHaveBeenCalledTimes(1) expect(mocks.runChatLoop).toHaveBeenCalledTimes(1) const sent = mocks.runChatLoop.mock.calls[0][0].messages expect(sent[0].content).not.toContain('continued from a previous conversation') @@ -3501,7 +3501,7 @@ describe('AIChatManager context compaction', () => { await manager.sendRequest() // A two-message prefix isn't worth a summary round-trip. - expect(mocks.getNonStreamingCompletion).not.toHaveBeenCalled() + expect(mocks.getStreamedCompletionText).not.toHaveBeenCalled() expect(manager.displayMessages.some((m) => m.role === 'summary')).toBe(false) }) @@ -3510,7 +3510,7 @@ describe('AIChatManager context compaction', () => { mocks.tryGetCurrentModel.mockReturnValue(gpt4oModel) // The user hits Stop while the summary request is in flight: it aborts the // turn's controller and rejects. - mocks.getNonStreamingCompletion.mockImplementation(async (_msgs: any, ac: AbortController) => { + mocks.getStreamedCompletionText.mockImplementation(async (_msgs: any, ac: AbortController) => { ac.abort('user_cancelled') throw new Error('aborted') }) @@ -3530,7 +3530,7 @@ describe('AIChatManager context compaction', () => { // Summarization was attempted and aborted, but the abort must NOT trigger a // destructive drop-oldest fallback: the full prefix survives and the unsent // turn is rolled back to the pre-send history (the head pair is still there). - expect(mocks.getNonStreamingCompletion).toHaveBeenCalledTimes(1) + expect(mocks.getStreamedCompletionText).toHaveBeenCalledTimes(1) expect(manager.messages).toHaveLength(6) expect(manager.messages[0].content).toContain('OLD1') expect(manager.displayMessages.some((m) => m.role === 'summary')).toBe(false) @@ -3565,7 +3565,7 @@ describe('AIChatManager manual compaction', () => { } it('folds the whole history into a single summary boundary, keeping nothing verbatim', async () => { - mocks.getNonStreamingCompletion.mockResolvedValue('MANUAL SUMMARY') + mocks.getStreamedCompletionText.mockResolvedValue('MANUAL SUMMARY') const manager = new AIChatManager() seedExchange(manager) manager.contextUsage = 123 @@ -3574,15 +3574,15 @@ describe('AIChatManager manual compaction', () => { await manager.compactManually() // The summarizer saw the entire history, then the summary instruction. - expect(mocks.getNonStreamingCompletion).toHaveBeenCalledTimes(1) - const summaryReq = mocks.getNonStreamingCompletion.mock.calls[0][0] + expect(mocks.getStreamedCompletionText).toHaveBeenCalledTimes(1) + const summaryReq = mocks.getStreamedCompletionText.mock.calls[0][0] expect(summaryReq).toHaveLength(5) expect(summaryReq[0].content).toBe('q1') expect(summaryReq[3].content).toBe('a2') expect(summaryReq[4].content).toContain('detailed summary') - // The summarizer's output must stay capped: without it the model default - // applies and the Anthropic SDK rejects the non-streaming call pre-flight. - expect(mocks.getNonStreamingCompletion.mock.calls[0][2]).toEqual({ maxTokensCap: 8000 }) + // The summarizer's output must stay capped at the budget reserved for the + // summary: without it the model default applies. + expect(mocks.getStreamedCompletionText.mock.calls[0][2]).toEqual({ maxTokensCap: 8000 }) // Nothing kept verbatim: messages collapse to just the summary user message. expect(manager.messages).toHaveLength(1) @@ -3608,13 +3608,13 @@ describe('AIChatManager manual compaction', () => { await manager.compactManually() - expect(mocks.getNonStreamingCompletion).not.toHaveBeenCalled() + expect(mocks.getStreamedCompletionText).not.toHaveBeenCalled() expect(mocks.sendUserToast).toHaveBeenCalledWith('Nothing to compact yet.') expect(manager.messages).toHaveLength(1) }) it('leaves history untouched when the user stops mid-summary', async () => { - mocks.getNonStreamingCompletion.mockImplementation(async (_msgs: any, ac: AbortController) => { + mocks.getStreamedCompletionText.mockImplementation(async (_msgs: any, ac: AbortController) => { ac.abort('user_cancelled') throw new Error('aborted') }) @@ -3631,7 +3631,7 @@ describe('AIChatManager manual compaction', () => { }) it('routes the /compact session command to manual compaction instead of the model', async () => { - mocks.getNonStreamingCompletion.mockResolvedValue('VIA COMMAND') + mocks.getStreamedCompletionText.mockResolvedValue('VIA COMMAND') const manager = new AIChatManager() manager.isSessionChat = true seedExchange(manager) @@ -3643,13 +3643,13 @@ describe('AIChatManager manual compaction', () => { expect(sent).toBe(true) expect(mocks.runChatLoop).not.toHaveBeenCalled() // ...it ran the summarizer and compacted in place, clearing the composer. - expect(mocks.getNonStreamingCompletion).toHaveBeenCalledTimes(1) + expect(mocks.getStreamedCompletionText).toHaveBeenCalledTimes(1) expect(manager.displayMessages[0]).toMatchObject({ role: 'summary', content: 'VIA COMMAND' }) expect(manager.instructions).toBe('') }) it('auto-sends a message queued while compaction was running', async () => { - mocks.getNonStreamingCompletion.mockResolvedValue('S') + mocks.getStreamedCompletionText.mockResolvedValue('S') mocks.runChatLoop.mockImplementation(async (config: any) => { const message = { role: 'assistant' as const, content: 'done' } config.addedMessages?.push(message) @@ -3669,7 +3669,7 @@ describe('AIChatManager manual compaction', () => { await manager.compactManually() // Compaction ran once, then the queued message went out as a real turn. - expect(mocks.getNonStreamingCompletion).toHaveBeenCalledTimes(1) + expect(mocks.getStreamedCompletionText).toHaveBeenCalledTimes(1) expect(mocks.runChatLoop).toHaveBeenCalledTimes(1) const sent = mocks.runChatLoop.mock.calls[0][0].messages expect(sent[sent.length - 1].content).toContain('follow-up question') @@ -3687,7 +3687,7 @@ describe('AIChatManager manual compaction', () => { // without ever reaching the model... expect(sent).toBe(true) expect(mocks.runChatLoop).not.toHaveBeenCalled() - expect(mocks.getNonStreamingCompletion).not.toHaveBeenCalled() + expect(mocks.getStreamedCompletionText).not.toHaveBeenCalled() // ...it reset the conversation and cleared the composer. expect(manager.displayMessages).toEqual([]) expect(manager.messages).toEqual([]) @@ -3758,7 +3758,7 @@ describe('AIChatManager manual compaction', () => { // Without the session-chat command surface, /compact is a normal message. expect(mocks.runChatLoop).toHaveBeenCalledTimes(1) - expect(mocks.getNonStreamingCompletion).not.toHaveBeenCalled() + expect(mocks.getStreamedCompletionText).not.toHaveBeenCalled() }) it('shadows a selected skill that collides with a built-in command', () => { diff --git a/frontend/src/lib/components/copilot/chat/openai-responses.ts b/frontend/src/lib/components/copilot/chat/openai-responses.ts index ebaa5b1c38..3a839f4bad 100644 --- a/frontend/src/lib/components/copilot/chat/openai-responses.ts +++ b/frontend/src/lib/components/copilot/chat/openai-responses.ts @@ -328,6 +328,7 @@ export async function* getOpenAIResponsesCompletionStream( forceModelProvider?: AIProviderModel openaiClient?: OpenAI reasoningEffort?: string + maxTokensCap?: number } ): AsyncGenerator { const { provider, config } = getProviderAndCompletionConfig({ @@ -335,6 +336,7 @@ export async function* getOpenAIResponsesCompletionStream( stream: true, tools, forceModelProvider: options?.forceModelProvider, + maxTokensCap: options?.maxTokensCap, reasoningEffort: options?.reasoningEffort }) const { instructions, input } = convertMessagesToResponsesInput(messages) @@ -383,6 +385,20 @@ export async function* getOpenAIResponsesCompletionStream( } ] } as OpenAI.Chat.Completions.ChatCompletionChunk + } else if (event.type === 'response.completed' || event.type === 'response.incomplete') { + yield { + id: 'chatcmpl-' + Date.now(), + object: 'chat.completion.chunk', + created: Date.now(), + model: responsesConfig.model, + choices: [ + { + index: 0, + delta: {}, + finish_reason: event.type === 'response.completed' ? 'stop' : 'length' + } + ] + } as OpenAI.Chat.Completions.ChatCompletionChunk } } } diff --git a/frontend/src/lib/components/copilot/lib.anthropicRouting.test.ts b/frontend/src/lib/components/copilot/lib.anthropicRouting.test.ts index 4bffcd80ac..52d58863c6 100644 --- a/frontend/src/lib/components/copilot/lib.anthropicRouting.test.ts +++ b/frontend/src/lib/components/copilot/lib.anthropicRouting.test.ts @@ -99,6 +99,7 @@ async function setupClients() { textDelta('Hel'), { type: 'content_block_delta', delta: { type: 'input_json_delta', partial_json: '{' } }, textDelta('lo'), + { type: 'message_delta', delta: { stop_reason: 'end_turn' } }, { type: 'message_stop' } ]) ) @@ -175,8 +176,9 @@ describe('Anthropic Messages API routing', () => { } expect(anthropicStream).toHaveBeenCalledTimes(1) - // only the two text deltas surface; message_start/stop and input_json are dropped - expect(chunks).toBe(2) + // the two text deltas and the stop reason surface; message_start/stop and + // input_json are dropped + expect(chunks).toBe(3) expect(text).toBe('Hello') }) @@ -228,6 +230,87 @@ describe('Anthropic Messages API routing', () => { expect(METADATA_MAX_TOKENS).toBeLessThanOrEqual(21333) }) + it('getStreamedCompletionText streams the capped request and falls back off the Responses API', async () => { + const { getStreamedCompletionText, workspaceAIClients } = await import('./lib') + + h.currentModel = { provider: 'anthropic', model: 'claude-sonnet-4-6' } + expect( + await getStreamedCompletionText(messages, new AbortController(), { maxTokensCap: 8000 }) + ).toBe('Hello') + expect(anthropicCreate).not.toHaveBeenCalled() + expect(anthropicStream.mock.calls[0][0].max_tokens).toBe(8000) + + h.currentModel = { provider: 'openai', model: 'gpt-4o' } + const completedResponses = vi + .fn() + .mockReturnValue( + streamOf([ + { type: 'response.created' }, + { type: 'response.output_text.delta', delta: 'responses ' }, + { type: 'response.output_text.delta', delta: 'text' }, + { type: 'response.completed' } + ]) + ) + vi.spyOn(workspaceAIClients, 'getOpenaiClient').mockReturnValue({ + chat: { completions: { create: openaiCreate } }, + responses: { stream: completedResponses } + } as any) + expect(await getStreamedCompletionText(messages, new AbortController())).toBe('responses text') + expect(openaiCreate).not.toHaveBeenCalled() + + // The Responses stream fails on iteration, not on creation. + const responsesStream = vi.fn().mockReturnValue( + (async function* () { + throw new Error('responses api not served') + })() + ) + openaiCreate.mockResolvedValue( + streamOf([ + { choices: [{ delta: { content: 'chat ' } }] }, + { choices: [{ delta: { content: 'text' } }] }, + { choices: [{ delta: {}, finish_reason: 'stop' }] }, + { choices: [], usage: {} } + ]) + ) + vi.spyOn(workspaceAIClients, 'getOpenaiClient').mockReturnValue({ + chat: { completions: { create: openaiCreate } }, + responses: { stream: responsesStream } + } as any) + vi.spyOn(console, 'error').mockImplementation(() => {}) + h.currentModel = { provider: 'openai', model: 'gpt-4o' } + + expect( + await getStreamedCompletionText(messages, new AbortController(), { maxTokensCap: 8000 }) + ).toBe('chat text') + expect(responsesStream.mock.calls[0][0].max_output_tokens).toBe(8000) + expect(openaiCreate.mock.calls[0][0]).toMatchObject({ stream: true, max_tokens: 8000 }) + }) + + it('getStreamedCompletionText rejects when the stream is cut or stopped midway', async () => { + const { getStreamedCompletionText } = await import('./lib') + h.currentModel = { provider: 'deepseek', model: 'deepseek-chat' } + + // A hop closing the response early ends the iteration cleanly, with no + // finish reason. + openaiCreate.mockResolvedValue(streamOf([{ choices: [{ delta: { content: 'partial' } }] }])) + await expect(getStreamedCompletionText(messages, new AbortController())).rejects.toThrow( + 'ended before the model finished' + ) + + const abortController = new AbortController() + // The OpenAI SDK swallows the abort and just ends the iteration. + openaiCreate.mockResolvedValue( + (async function* () { + yield { choices: [{ delta: { content: 'partial' } }] } + abortController.abort('user_cancelled') + })() + ) + + await expect(getStreamedCompletionText(messages, abortController)).rejects.toBe( + 'user_cancelled' + ) + }) + it('caps max_output_tokens for metadata completions on the OpenAI Responses path', async () => { const { getNonStreamingCompletion, getNonStreamingMetadataCompletion, METADATA_MAX_TOKENS } = await import('./lib') diff --git a/frontend/src/lib/components/copilot/lib.ts b/frontend/src/lib/components/copilot/lib.ts index 40f1dee6a6..dc420be2cc 100644 --- a/frontend/src/lib/components/copilot/lib.ts +++ b/frontend/src/lib/components/copilot/lib.ts @@ -746,6 +746,20 @@ function getAnthropicStreamingCompletion({ model: params.modelProvider.model, choices: [{ index: 0, delta: { content: event.delta.text }, finish_reason: null }] } + } else if (event.type === 'message_delta' && event.delta.stop_reason) { + yield { + id: '', + object: 'chat.completion.chunk', + created: 0, + model: params.modelProvider.model, + choices: [ + { + index: 0, + delta: {}, + finish_reason: event.delta.stop_reason === 'max_tokens' ? 'length' : 'stop' + } + ] + } } } } @@ -1158,12 +1172,18 @@ export async function getCompletion( openaiClient?: OpenAI reasoningEffort?: string promptCaching?: boolean + maxTokensCap?: number } ): Promise> { const modelProvider = options?.forceModelProvider ?? getCurrentModel() if (usesAnthropicMessagesApi(modelProvider.provider, modelProvider.model)) { - return getAnthropicStreamingCompletion({ messages, modelProvider, abortController }) + return getAnthropicStreamingCompletion({ + messages, + modelProvider, + abortController, + maxTokensCap: options?.maxTokensCap + }) } const { provider, config } = getProviderAndCompletionConfig({ @@ -1171,6 +1191,7 @@ export async function getCompletion( stream: true, tools, forceModelProvider: options?.forceModelProvider, + maxTokensCap: options?.maxTokensCap, promptCaching: options?.promptCaching, reasoningEffort: options?.reasoningEffort }) @@ -1181,7 +1202,8 @@ export async function getCompletion( const stream = getOpenAIResponsesCompletionStream(messages, abortController, tools, { forceModelProvider: options?.forceModelProvider, openaiClient: options?.openaiClient, - reasoningEffort: options?.reasoningEffort + reasoningEffort: options?.reasoningEffort, + maxTokensCap: options?.maxTokensCap }) as any return stream } catch (error) { @@ -1223,6 +1245,51 @@ export async function getCompletion( return completion } +/** + * Streams a completion and returns its whole text. For a generation that can run + * for minutes: a non-streaming request sends no byte until the model is done, so + * an idle timeout on any hop between the browser and the provider cuts it. + */ +export async function getStreamedCompletionText( + messages: ChatCompletionMessageParam[], + abortController: AbortController, + options?: { maxTokensCap?: number } +): Promise { + const drain = async (forceCompletions: boolean) => { + const stream = await getCompletion(messages, abortController, undefined, { + forceCompletions, + maxTokensCap: options?.maxTokensCap + }) + let text = '' + let finished = false + for await (const chunk of stream) { + text += getResponseFromEvent(chunk) + finished ||= !!chunk.choices?.[0]?.finish_reason + } + // The OpenAI SDK ends an aborted stream, and one closed early by a hop in + // between, without throwing: partial text must not pass for the whole + // completion. + abortController.signal.throwIfAborted() + if (!finished) { + throw new Error('The completion stream ended before the model finished') + } + return text + } + + try { + return await drain(false) + } catch (error) { + // The Responses API stream only fails once iterated, so the fallback to + // chat completions for a deployment that doesn't serve it lives here. + const { provider } = getCurrentModel() + if (abortController.signal.aborted || (provider !== 'openai' && provider !== 'azure_openai')) { + throw error + } + console.error('Error using Responses API:', error) + return drain(true) + } +} + function extractFirstJSON(str: string) { let depth = 0, i = 0