From c1381653f638ce26df0e71685ec5001f844ca2b0 Mon Sep 17 00:00:00 2001 From: Rune Philosof Date: Mon, 28 Sep 2026 11:49:00 +0200 Subject: [PATCH 1/4] Fix crash on Azure stream chunks without delta Azure OpenAI sends stream chunks where choices[0] has no delta (e.g. content filter results). Destructuring delta threw inside an async Promise executor, which became an unhandled rejection and killed the process along with every in-flight chat. Read content with optional chaining and skip empty chunks. Move the stream loop into a method that never rejects: a stream failure errors the observable for that one answer, and completion callback failures are logged. Co-Authored-By: Claude Opus 5.5 --- src/openai/openai.service.spec.ts | 118 ++++++++++++++++++++++++++++++ src/openai/openai.service.ts | 85 +++++++++++++-------- 2 files changed, 174 insertions(+), 29 deletions(-) create mode 100644 src/openai/openai.service.spec.ts diff --git a/src/openai/openai.service.spec.ts b/src/openai/openai.service.spec.ts new file mode 100644 index 0000000..89de81d --- /dev/null +++ b/src/openai/openai.service.spec.ts @@ -0,0 +1,118 @@ +import { firstValueFrom, toArray } from 'rxjs'; +import { AzureOpenAI } from 'openai'; +import { OpenaiService } from './openai.service'; +import { AppConfigService } from '../common/config/appConfig.service'; + +jest.mock('openai', () => ({ + __esModule: true, + default: jest.fn(), + AzureOpenAI: jest.fn(), +})); + +const CONFIG: Record = { + aiProvider: 'openai-azure', + openaiAzureEndpoint: 'https://example.openai.azure.com', + openaiAzureKey: 'key', + openaiAzureVersion: '2024-06-01', +}; + +const contentChunk = (content: string) => ({ + choices: [{ index: 0, finish_reason: null, delta: { content } }], +}); + +async function* streamOf(chunks: unknown[], error?: Error) { + for (const chunk of chunks) yield chunk; + if (error) throw error; +} + +describe('OpenaiService', () => { + let service: OpenaiService; + let create: jest.Mock; + + beforeEach(() => { + create = jest.fn(); + (AzureOpenAI as unknown as jest.Mock).mockImplementation(() => ({ + apiKey: 'key', + chat: { completions: { create } }, + })); + + service = new OpenaiService({ + get: (key: string) => CONFIG[key], + } as unknown as AppConfigService); + + jest.spyOn(service, 'analyzeChatConversation').mockResolvedValue(''); + jest.spyOn(service, 'implementApiCalls').mockResolvedValue(undefined); + jest.spyOn(service['logger'], 'error').mockImplementation(() => undefined); + }); + + const requestStream = (completeCb?: jest.Mock) => + service.getChatGptCompletionStream( + { + messages: [{ role: 'user', content: 'Hi' }], + model: 'gpt-4o', + stream: true, + }, + completeCb, + ); + + describe('getChatGptCompletionStream', () => { + it('skips chunks without delta, such as Azure content filter chunks', async () => { + create.mockResolvedValue( + streamOf([ + { choices: [] }, + { + choices: [ + { index: 0, finish_reason: null, content_filter_results: {} }, + ], + }, + contentChunk('Hello'), + { choices: [{ index: 0, finish_reason: null, delta: {} }] }, + contentChunk(' world'), + { choices: [{ index: 0, finish_reason: 'stop', delta: {} }] }, + ]), + ); + const completeCb = jest.fn().mockResolvedValue(undefined); + + const observable = await requestStream(completeCb); + const values = await firstValueFrom(observable.pipe(toArray())); + + expect(values).toEqual([ + JSON.stringify({ content: 'Hello' }), + JSON.stringify({ content: ' world' }), + '[DONE]', + ]); + expect(completeCb).toHaveBeenCalledWith( + 'Hello world', + expect.objectContaining({ prompt: expect.any(Number) }), + ); + }); + + it('errors the observable instead of rejecting when the stream fails', async () => { + create.mockResolvedValue( + streamOf([contentChunk('Hel')], new Error('connection reset')), + ); + const completeCb = jest.fn(); + + const observable = await requestStream(completeCb); + + await expect(firstValueFrom(observable.pipe(toArray()))).rejects.toThrow( + 'Failed to generate answer', + ); + expect(completeCb).not.toHaveBeenCalled(); + }); + + it('logs instead of rejecting when the completion callback fails', async () => { + create.mockResolvedValue(streamOf([contentChunk('Hi')])); + const completeCb = jest.fn().mockRejectedValue(new Error('db down')); + + const observable = await requestStream(completeCb); + await firstValueFrom(observable.pipe(toArray())); + await new Promise(process.nextTick); + + expect(service['logger'].error).toHaveBeenCalledWith( + expect.stringContaining('callback'), + expect.any(Error), + ); + }); + }); +}); diff --git a/src/openai/openai.service.ts b/src/openai/openai.service.ts index 47ce7e8..510bba6 100644 --- a/src/openai/openai.service.ts +++ b/src/openai/openai.service.ts @@ -344,7 +344,7 @@ export class OpenaiService { // API Call try { const res = await openAiClient.chat.completions.create(data); - const chatResponse = res.choices[0].message.content; + const chatResponse = res.choices[0]?.message?.content; return { response: chatResponse, @@ -402,43 +402,70 @@ export class OpenaiService { data.messages.map((m) => m.content).join(' '), ); + let completionStream: AsyncIterable; try { - const completionStream = await openAiClient.chat.completions.create(data); + completionStream = await openAiClient.chat.completions.create(data); + } catch (error) { + if (APIError.isPrototypeOf(error)) { + this.logger.error('OpenAI ChatCompletion API error', error); + this.logger.error('Error response', error.data); + } + throw error; + } - let answer = ''; + // Not awaited: the caller needs the observable before chunks arrive. + // forwardCompletionStream never rejects. + void this.forwardCompletionStream( + completionStream, + observable, + promptTokens, + completeCb, + ); - const streamPromise = new Promise(async (res) => { - for await (const part of completionStream) { - if (part.choices.length === 0) continue; + return observable; + } + + private async forwardCompletionStream( + completionStream: AsyncIterable, + observable: Subject, + promptTokens: number, + completeCb?: ( + answer: string, + usage: ChatGTPResponse['tokenUsage'], + ) => Promise, + ) { + let answer = ''; - const { content } = part.choices[0].delta; + try { + for await (const part of completionStream) { + // Azure sends chunks without `delta` (e.g. content filter results) + const content = part.choices?.[0]?.delta?.content; + if (content == null) continue; - if (content !== undefined) { - observable.next(JSON.stringify({ content })); - answer += content; - } - } + observable.next(JSON.stringify({ content })); + answer += content; + } + } catch (error) { + this.logger.error('OpenAI ChatCompletion stream error', error); + observable.error(new Error('Failed to generate answer')); + return; + } - res(true); - }); + observable.next('[DONE]'); + observable.complete(); - streamPromise.then(() => { - observable.next('[DONE]'); - observable.complete(); - const completionTokens = this.getTokenCount(answer); - completeCb?.(answer, { - prompt: promptTokens, - completion: completionTokens, - total: promptTokens + completionTokens, - }); + const completionTokens = this.getTokenCount(answer); + try { + await completeCb?.(answer, { + prompt: promptTokens, + completion: completionTokens, + total: promptTokens + completionTokens, }); } catch (error) { - if (APIError.isPrototypeOf(error)) { - this.logger.error('OpenAI ChatCompletion API error', error); - this.logger.error('Error response', error.data); - } - throw error; + this.logger.error( + 'OpenAI ChatCompletion completion callback error', + error, + ); } - return observable; } } From 26a9351babcea633135e8f5faf9de4f4a6b9294b Mon Sep 17 00:00:00 2001 From: Rune Philosof Date: Mon, 28 Sep 2026 12:41:13 +0200 Subject: [PATCH 2/4] Handle token counting errors and empty completions Move getTokenCount inside the guarded block so a tokenizer failure cannot become an unhandled rejection, and default an empty non-streaming response to '' to match the streaming path. Co-Authored-By: Claude Opus 5.5 --- src/openai/openai.service.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/openai/openai.service.ts b/src/openai/openai.service.ts index 510bba6..679f2f9 100644 --- a/src/openai/openai.service.ts +++ b/src/openai/openai.service.ts @@ -344,7 +344,7 @@ export class OpenaiService { // API Call try { const res = await openAiClient.chat.completions.create(data); - const chatResponse = res.choices[0]?.message?.content; + const chatResponse = res.choices[0]?.message?.content ?? ''; return { response: chatResponse, @@ -454,8 +454,8 @@ export class OpenaiService { observable.next('[DONE]'); observable.complete(); - const completionTokens = this.getTokenCount(answer); try { + const completionTokens = this.getTokenCount(answer); await completeCb?.(answer, { prompt: promptTokens, completion: completionTokens, From e87b1ad0513c1d673f0873fd27ea3bea75d08f0f Mon Sep 17 00:00:00 2001 From: Rune Philosof Date: Mon, 28 Sep 2026 12:46:28 +0200 Subject: [PATCH 3/4] Fix APIError detection in stream error logging APIError.isPrototypeOf(error) is always false for instances, so API errors were never logged. Use instanceof and log the parsed error body, since APIError has no data field. Co-Authored-By: Claude Opus 5.5 --- src/openai/openai.service.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/openai/openai.service.ts b/src/openai/openai.service.ts index 679f2f9..ec052d8 100644 --- a/src/openai/openai.service.ts +++ b/src/openai/openai.service.ts @@ -406,9 +406,9 @@ export class OpenaiService { try { completionStream = await openAiClient.chat.completions.create(data); } catch (error) { - if (APIError.isPrototypeOf(error)) { + if (error instanceof APIError) { this.logger.error('OpenAI ChatCompletion API error', error); - this.logger.error('Error response', error.data); + this.logger.error('Error response', error.error); } throw error; } From 636c59bf237fd80fdf6c306e3af28fc5648570f2 Mon Sep 17 00:00:00 2001 From: Rune Philosof Date: Mon, 28 Sep 2026 12:51:34 +0200 Subject: [PATCH 4/4] Test empty non-streaming completion responses Co-Authored-By: Claude Opus 5.5 --- src/openai/openai.service.spec.ts | 23 +++++++++++++++++++++++ 1 file changed, 23 insertions(+) diff --git a/src/openai/openai.service.spec.ts b/src/openai/openai.service.spec.ts index 89de81d..d3b35a1 100644 --- a/src/openai/openai.service.spec.ts +++ b/src/openai/openai.service.spec.ts @@ -55,6 +55,29 @@ describe('OpenaiService', () => { completeCb, ); + describe('getChatGptCompletion', () => { + it.each([ + ['no choices', []], + [ + 'a filtered choice without message', + [{ index: 0, finish_reason: 'content_filter' }], + ], + [ + 'a null message content', + [{ index: 0, message: { role: 'assistant', content: null } }], + ], + ])('returns an empty response for %s', async (_, choices) => { + create.mockResolvedValue({ choices, usage: undefined }); + + const result = await service.getChatGptCompletion({ + messages: [{ role: 'user', content: 'Hi' }], + model: 'gpt-4o', + }); + + expect(result.response).toBe(''); + }); + }); + describe('getChatGptCompletionStream', () => { it('skips chunks without delta, such as Azure content filter chunks', async () => { create.mockResolvedValue(