diff --git a/packages/backend/src/ai/ai.service.spec.ts b/packages/backend/src/ai/ai.service.spec.ts index c1a0aeb8..05b12232 100644 --- a/packages/backend/src/ai/ai.service.spec.ts +++ b/packages/backend/src/ai/ai.service.spec.ts @@ -136,6 +136,23 @@ describe('AIService', () => { expect(aiService.redis.removeInflight).toHaveBeenCalledWith('U1', 'T1'); expect(aiService.redis.decrementDailyRequests).toHaveBeenCalledWith('U1', 'T1'); }); + + it('alerts #muzzlefeedback when OpenAI returns a 429 error', async () => { + const createSpy = aiService.openAi.responses.create as Mock; + createSpy.mockRejectedValue( + Object.assign(new Error('Rate limit exceeded'), { + status: 429, + error: { message: 'Please slow down.' }, + }), + ); + + await expect(aiService.generateText('U1', 'T1', 'C1', 'hello')).rejects.toThrow('Rate limit exceeded'); + + expect(aiService.webService.sendMessage).toHaveBeenCalledWith( + '#muzzlefeedback', + 'OpenAI 429 during generateText: Please slow down.', + ); + }); }); describe('generateImage', () => { @@ -517,6 +534,88 @@ describe('AIService', () => { }); }); + describe('alertOnOpenAiRateLimit', () => { + it('alerts #muzzlefeedback for 429 errors with the OpenAI message', async () => { + await aiService.alertOnOpenAiRateLimit( + Object.assign(new Error('Rate limit exceeded'), { + status: 429, + error: { message: 'Please slow down.' }, + }), + 'generateText', + ); + + expect(aiService.webService.sendMessage).toHaveBeenCalledWith( + '#muzzlefeedback', + 'OpenAI 429 during generateText: Please slow down.', + ); + }); + + it('does not alert for non-429 errors', async () => { + await aiService.alertOnOpenAiRateLimit(new Error('API error'), 'generateText'); + + expect(aiService.webService.sendMessage).not.toHaveBeenCalled(); + }); + + it('falls back to the top-level error message when the nested OpenAI message is blank', async () => { + await aiService.alertOnOpenAiRateLimit( + { + status: 429, + error: { message: ' ' }, + message: ' Rate limit exceeded ', + }, + 'generateText', + ); + + expect(aiService.webService.sendMessage).toHaveBeenCalledWith( + '#muzzlefeedback', + 'OpenAI 429 during generateText: Rate limit exceeded', + ); + }); + + it('uses a default alert message when OpenAI omits all error text', async () => { + await aiService.alertOnOpenAiRateLimit( + { + status: 429, + error: { message: ' ' }, + message: ' ', + }, + 'generateText', + ); + + expect(aiService.webService.sendMessage).toHaveBeenCalledWith( + '#muzzlefeedback', + 'OpenAI 429 during generateText: OpenAI returned a 429 response without an error message.', + ); + }); + + it('logs and suppresses Slack delivery failures', async () => { + (aiService.webService.sendMessage as Mock).mockRejectedValueOnce(new Error('slack failed')); + + await expect( + aiService.alertOnOpenAiRateLimit( + Object.assign(new Error('Rate limit exceeded'), { + status: 429, + error: { message: 'Please slow down.' }, + }), + 'generateText', + ), + ).resolves.toBeUndefined(); + + expect(aiService.aiServiceLogger.error).toHaveBeenCalledWith( + 'Failed to send OpenAI 429 alert to Slack', + expect.objectContaining({ + context: { + operation: 'generateText', + openAiErrorMessage: 'Please slow down.', + }, + error: expect.objectContaining({ + message: 'slack failed', + }), + }), + ); + }); + }); + describe('participate', () => { it('should be defined', () => { expect(aiService.participate).toBeDefined(); diff --git a/packages/backend/src/ai/ai.service.ts b/packages/backend/src/ai/ai.service.ts index 837a4ce1..064f40af 100644 --- a/packages/backend/src/ai/ai.service.ts +++ b/packages/backend/src/ai/ai.service.ts @@ -56,6 +56,32 @@ const isResponseOutputMessage = (block: ResponseOutputItem): block is ResponseOu const isResponseOutputText = (block: ResponseOutputText | ResponseOutputRefusal): block is ResponseOutputText => block.type === 'output_text'; +const getOpenAiStatusCode = (error: unknown): number | undefined => { + if (!isRecord(error)) { + return undefined; + } + + const status = Reflect.get(error, 'status'); + return typeof status === 'number' ? status : undefined; +}; + +const getOpenAiErrorMessage = (error: unknown): string | undefined => { + if (!isRecord(error)) { + return error instanceof Error ? error.message : undefined; + } + + const apiError = Reflect.get(error, 'error'); + if (isRecord(apiError)) { + const apiErrorMessage = Reflect.get(apiError, 'message'); + if (typeof apiErrorMessage === 'string' && apiErrorMessage.trim()) { + return apiErrorMessage.trim(); + } + } + + const message = Reflect.get(error, 'message'); + return typeof message === 'string' && message.trim() ? message.trim() : undefined; +}; + const extractAndParseOpenAiResponse = (response: OpenAI.Responses.Response): string | undefined => { const textBlock = response.output.find(isResponseOutputMessage); const outputText = textBlock?.content.find(isResponseOutputText)?.text; @@ -133,6 +159,28 @@ export class AIService { return this.slackPersistenceService.setCustomPrompt(userId, teamId, prompt); } + /** + * Sends a Slack alert to #muzzlefeedback when OpenAI returns a 429 rate-limit response. + * Expects OpenAI-style errors with a numeric `status` field and an optional nested `error.message`. + * Returns without side effects for non-429 errors. + */ + public async alertOnOpenAiRateLimit(error: unknown, operation: string): Promise { + if (getOpenAiStatusCode(error) !== 429) { + return; + } + + const errorMessage = getOpenAiErrorMessage(error) ?? 'OpenAI returned a 429 response without an error message.'; + + try { + await this.webService.sendMessage('#muzzlefeedback', `OpenAI 429 during ${operation}: ${errorMessage}`); + } catch (slackError) { + logError(this.aiServiceLogger, 'Failed to send OpenAI 429 alert to Slack', slackError, { + operation, + openAiErrorMessage: errorMessage, + }); + } + } + public clearCustomPrompt(userId: string, teamId: string): Promise { return this.slackPersistenceService.clearCustomPrompt(userId, teamId); } @@ -163,6 +211,7 @@ export class AIService { } }) .catch(async (e) => { + await this.alertOnOpenAiRateLimit(e, 'generateText'); logError(this.aiServiceLogger, 'Failed to generate AI text response', e, { userId, teamId, @@ -205,7 +254,11 @@ export class AIService { input: REDPLOY_MOONBEAM_TEXT_PROMPT, user: 'Moonbeam', }) - .then((x) => extractAndParseOpenAiResponse(x)); + .then((x) => extractAndParseOpenAiResponse(x)) + .catch(async (error) => { + await this.alertOnOpenAiRateLimit(error, 'redeployMoonbeam'); + throw error; + }); const aiImage = this.gemini.models .generateContent({ @@ -373,6 +426,7 @@ export class AIService { return extractAndParseOpenAiResponse(x); }) .catch(async (e) => { + await this.alertOnOpenAiRateLimit(e, 'generateCorpoSpeak'); logError(this.aiServiceLogger, 'Failed to generate corpo-speak response', e, { prompt: text, }); @@ -467,6 +521,7 @@ export class AIService { }); }) .catch(async (e) => { + await this.alertOnOpenAiRateLimit(e, 'promptWithHistory'); logError(this.aiServiceLogger, 'Failed to process prompt with history', e, { userId: request.user_id, teamId: request.team_id, @@ -556,6 +611,7 @@ export class AIService { } }) .catch(async (e) => { + await this.alertOnOpenAiRateLimit(e, 'participate'); logError(this.aiServiceLogger, 'Failed to generate AI participation response', e, { teamId, channelId, diff --git a/packages/backend/src/ai/memory/memory.job.spec.ts b/packages/backend/src/ai/memory/memory.job.spec.ts index 072eaefe..676adbb0 100644 --- a/packages/backend/src/ai/memory/memory.job.spec.ts +++ b/packages/backend/src/ai/memory/memory.job.spec.ts @@ -1,6 +1,29 @@ import { beforeEach, describe, expect, it, vi } from 'vitest'; import { MemoryJob } from './memory.job'; +type ExtractMemoriesInvoker = { + extractMemories: ( + teamId: string, + channelId: string, + conversationHistory: string, + participantSlackIds: string[], + ) => Promise; +}; + +const invokeExtractMemories = ( + job: MemoryJob, + teamId: string, + channelId: string, + conversationHistory: string, + participantSlackIds: string[], +): Promise => + (job as unknown as ExtractMemoriesInvoker).extractMemories( + teamId, + channelId, + conversationHistory, + participantSlackIds, + ); + describe('MemoryJob', () => { let job: MemoryJob; let memoryPersistenceService: { @@ -21,6 +44,7 @@ describe('MemoryJob', () => { warn: ReturnType; }; let aiService: { + alertOnOpenAiRateLimit: ReturnType; openAi: { responses: { create: ReturnType; @@ -48,6 +72,7 @@ describe('MemoryJob', () => { warn: vi.fn(), }; aiService = { + alertOnOpenAiRateLimit: vi.fn().mockResolvedValue(undefined), openAi: { responses: { create: vi.fn(), @@ -65,16 +90,7 @@ describe('MemoryJob', () => { it('returns early when extraction lock exists', async () => { redis.getValue.mockResolvedValue('1'); - await ( - job as never as { - extractMemories: ( - teamId: string, - channelId: string, - conversationHistory: string, - participantSlackIds: string[], - ) => Promise; - } - ).extractMemories('T1', 'C1', 'history', ['U1']); + await invokeExtractMemories(job, 'T1', 'C1', 'history', ['U1']); expect(jobLogger.info).toHaveBeenCalled(); }); @@ -84,16 +100,7 @@ describe('MemoryJob', () => { output: [{ type: 'message', content: [{ type: 'output_text', text: 'NONE' }] }], }); - await ( - job as never as { - extractMemories: ( - teamId: string, - channelId: string, - conversationHistory: string, - participantSlackIds: string[], - ) => Promise; - } - ).extractMemories('T1', 'C1', 'history', ['U1']); + await invokeExtractMemories(job, 'T1', 'C1', 'history', ['U1']); expect(memoryPersistenceService.saveMemories).not.toHaveBeenCalled(); }); @@ -117,16 +124,7 @@ describe('MemoryJob', () => { ], }); - await ( - job as never as { - extractMemories: ( - teamId: string, - channelId: string, - conversationHistory: string, - participantSlackIds: string[], - ) => Promise; - } - ).extractMemories('T1', 'C1', 'history', ['U123ABC']); + await invokeExtractMemories(job, 'T1', 'C1', 'history', ['U123ABC']); expect(memoryPersistenceService.saveMemories).toHaveBeenCalled(); expect(memoryPersistenceService.reinforceMemory).toHaveBeenCalledWith(10); @@ -152,17 +150,20 @@ describe('MemoryJob', () => { ], }); - await ( - job as never as { - extractMemories: ( - teamId: string, - channelId: string, - conversationHistory: string, - participantSlackIds: string[], - ) => Promise; - } - ).extractMemories('T1', 'C1', 'history', ['U123ABC']); + await invokeExtractMemories(job, 'T1', 'C1', 'history', ['U123ABC']); expect(jobLogger.warn).toHaveBeenCalled(); }); + + it('alerts on OpenAI 429 errors during extraction', async () => { + const rateLimitError = Object.assign(new Error('Rate limit exceeded'), { + status: 429, + error: { message: 'Too many requests.' }, + }); + aiService.openAi.responses.create.mockRejectedValue(rateLimitError); + + await invokeExtractMemories(job, 'T1', 'C1', 'history', ['U1']); + + expect(aiService.alertOnOpenAiRateLimit).toHaveBeenCalledWith(rateLimitError, 'memory extraction'); + }); }); diff --git a/packages/backend/src/ai/memory/memory.job.ts b/packages/backend/src/ai/memory/memory.job.ts index 1fbb739f..1ac41f8e 100644 --- a/packages/backend/src/ai/memory/memory.job.ts +++ b/packages/backend/src/ai/memory/memory.job.ts @@ -111,6 +111,10 @@ export class MemoryJob { instructions: prompt, input: conversationHistory, }) + .catch(async (error) => { + await this.aiService.alertOnOpenAiRateLimit(error, 'memory extraction'); + throw error; + }) .then((response) => extractAndParseOpenAiResponse(response)); if (!result) {