mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-06 00:02:30 +00:00
fix: stream the ai chat compaction request (#11514)
* fix: stream the ai chat compaction request instead of waiting on it Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix: reject a compaction stream stopped midway instead of keeping partial text Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * fix: reject a compaction stream that ends without a finish reason Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * test: cover a completed responses stream in the streamed completion helper Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
14d2995d76
commit
289d3c3941
@@ -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'
|
||||
}
|
||||
|
||||
@@ -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(
|
||||
'<analysis>scratchpad</analysis><summary>SUMMARY TEXT</summary>'
|
||||
)
|
||||
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(
|
||||
'<analysis>s</analysis><summary>SUM</summary>'
|
||||
)
|
||||
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>SUMMARY TEXT</summary>')
|
||||
mocks.getStreamedCompletionText.mockResolvedValue('<summary>SUMMARY TEXT</summary>')
|
||||
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('<summary>MANUAL SUMMARY</summary>')
|
||||
mocks.getStreamedCompletionText.mockResolvedValue('<summary>MANUAL SUMMARY</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('<summary>VIA COMMAND</summary>')
|
||||
mocks.getStreamedCompletionText.mockResolvedValue('<summary>VIA COMMAND</summary>')
|
||||
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('<summary>S</summary>')
|
||||
mocks.getStreamedCompletionText.mockResolvedValue('<summary>S</summary>')
|
||||
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', () => {
|
||||
|
||||
@@ -328,6 +328,7 @@ export async function* getOpenAIResponsesCompletionStream(
|
||||
forceModelProvider?: AIProviderModel
|
||||
openaiClient?: OpenAI
|
||||
reasoningEffort?: string
|
||||
maxTokensCap?: number
|
||||
}
|
||||
): AsyncGenerator<OpenAI.Chat.Completions.ChatCompletionChunk> {
|
||||
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
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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')
|
||||
|
||||
@@ -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<Stream<ChatCompletionChunk>> {
|
||||
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<string> {
|
||||
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
|
||||
|
||||
Reference in New Issue
Block a user