diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/drivers/bullmq.driver.ts b/packages/twenty-server/src/engine/core-modules/message-queue/drivers/bullmq.driver.ts index 8770297de4..dad792fc2e 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/drivers/bullmq.driver.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/drivers/bullmq.driver.ts @@ -30,6 +30,7 @@ import { getJobKey } from 'src/engine/core-modules/message-queue/utils/get-job-k import { type MetricsService } from 'src/engine/core-modules/metrics/metrics.service'; import { MetricsKeys } from 'src/engine/core-modules/metrics/types/metrics-keys.type'; import { applyWorkspaceSentryContextFromJobData } from 'src/engine/core-modules/sentry/utils/apply-workspace-sentry-context-from-job-data.util'; +import { type TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service'; export type BullMQDriverOptions = QueueOptions; @@ -47,10 +48,14 @@ export class BullMQDriver MessageQueue, Worker >; + private workerOptionsMap: Partial< + Record + > = {}; constructor( private options: BullMQDriverOptions, private metricsService: MetricsService, + private twentyConfigService: TwentyConfigService, ) {} onModuleInit() { @@ -89,7 +94,7 @@ export class BullMQDriver } async onModuleDestroy() { - const workers = Object.entries(this.workerMap); + const workers = Object.entries(this.workerMap) as [MessageQueue, Worker][]; const queues = Object.values(this.queueMap); if (workers.length > 0) { @@ -101,7 +106,11 @@ export class BullMQDriver let workerCloseError: unknown; try { - await Promise.all(workers.map(([, worker]) => worker.close())); + await Promise.all( + workers.map(([queueName, worker]) => + this.closeWorker(queueName, worker), + ), + ); } catch (error) { workerCloseError = error; } @@ -125,6 +134,34 @@ export class BullMQDriver this.logger.log('Message queue shutdown complete'); } + private async closeWorker( + queueName: MessageQueue, + worker: Worker, + ): Promise { + if (!this.workerOptionsMap[queueName]?.boundedShutdownDrain) { + await worker.close(); + + return; + } + + const shutdownTimeoutMs = this.twentyConfigService.get( + 'AI_STREAM_SHUTDOWN_DRAIN_MS', + ); + + const abortTimer = setTimeout(() => { + this.logger.warn( + `Queue ${queueName} still has active jobs after draining for ${shutdownTimeoutMs}ms, aborting them`, + ); + worker.cancelAllJobs('worker shutdown'); + }, shutdownTimeoutMs); + + try { + await worker.close(); + } finally { + clearTimeout(abortTimer); + } + } + work( queueName: MessageQueue, handler: (job: MessageQueueJob) => Promise, @@ -138,15 +175,20 @@ export class BullMQDriver ...(isDefined(options?.lockDuration) ? { lockDuration: options.lockDuration } : {}), + ...(isDefined(options?.maxStalledCount) + ? { maxStalledCount: options.maxStalledCount } + : {}), metrics: { maxDataPoints: MetricsTime.ONE_WEEK, collectInterval: 60000, }, }; + this.workerOptionsMap[queueName] = options; + this.workerMap[queueName] = new Worker( queueName, - async (job) => + async (job, _token, abortSignal) => Sentry.withIsolationScope(async () => { applyWorkspaceSentryContextFromJobData(job.data); @@ -169,7 +211,12 @@ export class BullMQDriver this.logger.log( `Processing job ${job.id} with name ${job.name} on queue ${queueName}${workspaceSuffix}`, ); - await handler({ data: job.data, id: job.id ?? '', name: job.name }); + await handler({ + data: job.data, + id: job.id ?? '', + name: job.name, + abortSignal, + }); const timeEnd = performance.now(); const executionTime = timeEnd - timeStart; diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/interfaces/message-queue-job.interface.ts b/packages/twenty-server/src/engine/core-modules/message-queue/interfaces/message-queue-job.interface.ts index 3bcb413d61..6c40404039 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/interfaces/message-queue-job.interface.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/interfaces/message-queue-job.interface.ts @@ -3,6 +3,11 @@ export interface MessageQueueJob { id: string; name: string; data: T; + abortSignal?: AbortSignal; +} + +export interface MessageQueueJobContext { + abortSignal?: AbortSignal; } export interface MessageQueueCronJobData< diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/interfaces/message-queue-module-options.interface.ts b/packages/twenty-server/src/engine/core-modules/message-queue/interfaces/message-queue-module-options.interface.ts index f9f8652053..6888049e5d 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/interfaces/message-queue-module-options.interface.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/interfaces/message-queue-module-options.interface.ts @@ -1,5 +1,6 @@ import { type BullMQDriverOptions } from 'src/engine/core-modules/message-queue/drivers/bullmq.driver'; import { type MetricsService } from 'src/engine/core-modules/metrics/metrics.service'; +import { type TwentyConfigService } from 'src/engine/core-modules/twenty-config/twenty-config.service'; export enum MessageQueueDriverType { BullMQ = 'bull-mq', @@ -10,6 +11,7 @@ export interface BullMQDriverFactoryOptions { type: MessageQueueDriverType.BullMQ; options: BullMQDriverOptions; metricsService: MetricsService; + twentyConfigService: TwentyConfigService; } export interface SyncDriverFactoryOptions { diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/interfaces/message-queue-worker-options.interface.ts b/packages/twenty-server/src/engine/core-modules/message-queue/interfaces/message-queue-worker-options.interface.ts index 777d122050..4fb070c2db 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/interfaces/message-queue-worker-options.interface.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/interfaces/message-queue-worker-options.interface.ts @@ -1,4 +1,6 @@ export interface MessageQueueWorkerOptions { concurrency?: number; lockDuration?: number; + maxStalledCount?: number; + boundedShutdownDrain?: boolean; } diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-core.module.ts b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-core.module.ts index 94e65694f9..b6160b1e65 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-core.module.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-core.module.ts @@ -96,7 +96,11 @@ export class MessageQueueCoreModule extends ConfigurableModuleClass { static async createDriver(config: typeof OPTIONS_TYPE) { switch (config.type) { case MessageQueueDriverType.BullMQ: { - return new BullMQDriver(config.options, config.metricsService); + return new BullMQDriver( + config.options, + config.metricsService, + config.twentyConfigService, + ); } case MessageQueueDriverType.Sync: { return new SyncDriver(); diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-worker-options.constant.ts b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-worker-options.constant.ts index 24b8e9b5be..8d5ee0defa 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-worker-options.constant.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue-worker-options.constant.ts @@ -8,6 +8,8 @@ export const QUEUE_WORKER_OPTIONS: Partial< [MessageQueue.aiStreamQueue]: { concurrency: 20, lockDuration: AI_STREAM_LOCK_DURATION_MS, + maxStalledCount: 0, + boundedShutdownDrain: true, }, [MessageQueue.logicFunctionQueue]: { concurrency: 10 }, }; diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.explorer.ts b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.explorer.ts index fe90a46de3..d6bf9eb222 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.explorer.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.explorer.ts @@ -208,7 +208,9 @@ export class MessageQueueExplorer implements OnModuleInit { for (const processMethodName of processMethodNames) { try { // @ts-expect-error legacy noImplicitAny - await instance[processMethodName].call(instance, job.data); + await instance[processMethodName].call(instance, job.data, { + abortSignal: job.abortSignal, + }); } catch (err) { if (shouldCaptureException(err)) { this.exceptionHandlerService.captureExceptions([err]); diff --git a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.module-factory.ts b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.module-factory.ts index 58bdf4ab02..2d30c80cc0 100644 --- a/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.module-factory.ts +++ b/packages/twenty-server/src/engine/core-modules/message-queue/message-queue.module-factory.ts @@ -15,7 +15,7 @@ import { type TwentyConfigService } from 'src/engine/core-modules/twenty-config/ * @param metricsService */ export const messageQueueModuleFactory = async ( - _twentyConfigService: TwentyConfigService, + twentyConfigService: TwentyConfigService, redisClientService: RedisClientService, metricsService: MetricsService, ): Promise => { @@ -29,6 +29,7 @@ export const messageQueueModuleFactory = async ( connection: redisClientService.getQueueClient(), }, metricsService, + twentyConfigService, } satisfies BullMQDriverFactoryOptions; } default: diff --git a/packages/twenty-server/src/engine/core-modules/twenty-config/config-variables.ts b/packages/twenty-server/src/engine/core-modules/twenty-config/config-variables.ts index 6d4f302595..ab6092b449 100644 --- a/packages/twenty-server/src/engine/core-modules/twenty-config/config-variables.ts +++ b/packages/twenty-server/src/engine/core-modules/twenty-config/config-variables.ts @@ -1280,6 +1280,17 @@ export class ConfigVariables { @IsOptional() SERVER_KEEP_ALIVE_TIMEOUT_MS = 65000; + @ConfigVariablesMetadata({ + group: ConfigVariablesGroup.SERVER_CONFIG, + description: + 'How long (ms) a worker shutdown waits for active AI chat stream jobs to finish before aborting them into a retryable interrupted state. Must be lower than the pod terminationGracePeriodSeconds so the abort and clean exit fit before SIGKILL (default: 300000)', + type: ConfigVariableType.NUMBER, + isEnvOnly: true, + }) + @CastToPositiveNumber() + @IsOptional() + AI_STREAM_SHUTDOWN_DRAIN_MS = 300_000; + @ConfigVariablesMetadata({ group: ConfigVariablesGroup.SERVER_CONFIG, description: 'Base URL for the server', diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/constants/stream-interrupted-code.constant.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/constants/stream-interrupted-code.constant.ts deleted file mode 100644 index 98be7902d1..0000000000 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/constants/stream-interrupted-code.constant.ts +++ /dev/null @@ -1 +0,0 @@ -export const STREAM_INTERRUPTED_CODE = 'STREAM_INTERRUPTED'; diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/__tests__/stream-agent-chat.job.spec.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/__tests__/stream-agent-chat.job.spec.ts index 238d8efafb..6a8d5e2ce1 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/__tests__/stream-agent-chat.job.spec.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/__tests__/stream-agent-chat.job.spec.ts @@ -99,9 +99,11 @@ describe('StreamAgentChatJob', () => { const publishedEvents: PublishedEvent[] = []; const threadRepository = { - findOne: jest - .fn() - .mockResolvedValue({ id: 'thread-id', deletedAt: null }), + findOne: jest.fn().mockResolvedValue({ + id: 'thread-id', + deletedAt: null, + activeStreamId: 'stream-id', + }), update: jest.fn().mockImplementation((_workspaceId, _criteria, values) => Promise.resolve({ affected: @@ -388,6 +390,119 @@ describe('StreamAgentChatJob', () => { ); }); + it('bails out without streaming when the thread no longer holds the claim for this stream', async () => { + const { + job, + publishedEvents, + threadRepository, + eventPublisherService, + agentChatStreamingService, + } = buildJob(); + + threadRepository.findOne.mockResolvedValueOnce({ + id: 'thread-id', + deletedAt: null, + activeStreamId: 'newer-stream-id', + }); + + await job.handle(jobData); + + expect(publishedEvents).toHaveLength(0); + expect(eventPublisherService.resetStreamState).not.toHaveBeenCalled(); + expect(threadRepository.update).not.toHaveBeenCalled(); + expect( + agentChatStreamingService.flushNextQueuedMessage, + ).not.toHaveBeenCalled(); + }); + + it('bails out without streaming when the thread was deleted', async () => { + const { job, publishedEvents, threadRepository, eventPublisherService } = + buildJob(); + + threadRepository.findOne.mockResolvedValueOnce(null); + + await job.handle(jobData); + + expect(publishedEvents).toHaveLength(0); + expect(eventPublisherService.resetStreamState).not.toHaveBeenCalled(); + expect(threadRepository.update).not.toHaveBeenCalled(); + }); + + it('persists the interrupted error and publishes the terminal sequence when aborted by a worker shutdown', async () => { + let triggerShutdown: (() => void) | undefined; + + const { + job, + publishedEvents, + threadRepository, + agentChatStreamingService, + } = buildJob({ + chatStream: createFakeChatStream({ + onFirstChunk: () => triggerShutdown?.(), + }), + }); + + const shutdownController = new AbortController(); + + triggerShutdown = () => shutdownController.abort(); + + await expect( + job.handle(jobData, { abortSignal: shutdownController.signal }), + ).rejects.toMatchObject({ code: AiExceptionCode.STREAM_INTERRUPTED }); + + const eventTypes = publishedEvents.map((event) => event.type); + const streamErrorIndex = eventTypes.indexOf('stream-error'); + const queueUpdatedIndex = eventTypes.indexOf('queue-updated'); + + expect(publishedEvents[streamErrorIndex]).toMatchObject({ + type: 'stream-error', + code: AiExceptionCode.STREAM_INTERRUPTED, + }); + expect(queueUpdatedIndex).toBeGreaterThan(streamErrorIndex); + expect(eventTypes).not.toContain('message-persisted'); + expect(threadRepository.update).toHaveBeenCalledWith( + 'workspace-id', + { id: 'thread-id' }, + { + lastStreamError: expect.objectContaining({ + code: AiExceptionCode.STREAM_INTERRUPTED, + }), + }, + ); + expect(threadRepository.update).toHaveBeenCalledWith( + 'workspace-id', + { id: 'thread-id', activeStreamId: 'stream-id' }, + { activeStreamId: null }, + ); + expect( + agentChatStreamingService.flushNextQueuedMessage, + ).not.toHaveBeenCalled(); + }); + + it('keeps user-cancel semantics when a shutdown signal is wired but never aborted', async () => { + let triggerUserCancel: (() => void) | undefined; + + const { job, publishedEvents, agentChatStreamingService, cancelCallbacks } = + buildJob({ + chatStream: createFakeChatStream({ + onFirstChunk: () => triggerUserCancel?.(), + }), + }); + + triggerUserCancel = () => cancelCallbacks.forEach((callback) => callback()); + + const shutdownController = new AbortController(); + + await job.handle(jobData, { abortSignal: shutdownController.signal }); + + expect(publishedEvents.map((event) => event.type)).not.toContain( + 'stream-error', + ); + expect( + agentChatStreamingService.flushNextQueuedMessage, + ).not.toHaveBeenCalled(); + }); + it('resolves without flushing the queue when the stream is cancelled', async () => { let triggerCancel: (() => void) | undefined; diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/stream-agent-chat.job.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/stream-agent-chat.job.ts index c43188dd72..54259f8c8a 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/stream-agent-chat.job.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/jobs/stream-agent-chat.job.ts @@ -12,6 +12,8 @@ import { isDefined } from 'twenty-shared/utils'; import { Repository } from 'typeorm'; import { v5 as uuidv5 } from 'uuid'; +import { type MessageQueueJobContext } from 'src/engine/core-modules/message-queue/interfaces/message-queue-job.interface'; + import { Process } from 'src/engine/core-modules/message-queue/decorators/process.decorator'; import { Processor } from 'src/engine/core-modules/message-queue/decorators/processor.decorator'; import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; @@ -68,7 +70,23 @@ export class StreamAgentChatJob { ) {} @Process(STREAM_AGENT_CHAT_JOB_NAME) - async handle(data: StreamAgentChatJobData): Promise { + async handle( + data: StreamAgentChatJobData, + context?: MessageQueueJobContext, + ): Promise { + const thread = await this.threadRepository.findOne(data.workspaceId, { + where: { id: data.threadId }, + select: ['id', 'activeStreamId'], + }); + + if (thread?.activeStreamId !== data.streamId) { + this.logger.warn( + `Skipping stream ${data.streamId} for thread ${data.threadId}: the thread no longer holds this claim`, + ); + + return; + } + await this.eventPublisherService.resetStreamState(data.threadId); const abortController = new AbortController(); @@ -82,6 +100,19 @@ export class StreamAgentChatJob { abortController.abort(); }); + context?.abortSignal?.addEventListener( + 'abort', + () => { + abortController.abort( + new AiException( + 'The response was interrupted before it could finish.', + AiExceptionCode.STREAM_INTERRUPTED, + ), + ); + }, + { once: true }, + ); + try { const workspace = await this.workspaceRepository.findOne({ where: { id: data.workspaceId }, @@ -268,7 +299,15 @@ export class StreamAgentChatJob { abortSignal.addEventListener( 'abort', () => { - void streamFinishedPromise.then(() => resolve()); + const reason = abortSignal.reason; + + void streamFinishedPromise.then(() => { + if (reason instanceof AiException) { + reject(reason); + } else { + resolve(); + } + }); }, { once: true }, ); diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/resolvers/agent-chat-subscription.resolver.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/resolvers/agent-chat-subscription.resolver.ts index 6a378d9204..4e9fa272dd 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/resolvers/agent-chat-subscription.resolver.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/resolvers/agent-chat-subscription.resolver.ts @@ -100,29 +100,23 @@ export class AgentChatSubscriptionResolver { }); } - // Reaping from the keep-alive loop means a user watching a stream whose - // worker died sees the interrupted state without having to interact; - // reapDeadStream publishes the stream-error event they are subscribed to. private async reapWatchedStreamIfDead( workspaceId: string, threadId: string, ): Promise { - try { - const thread = await this.threadRepository.findOne(workspaceId, { + const thread = await this.threadRepository + .findOne(workspaceId, { where: { id: threadId }, select: ['id', 'activeStreamId'], - }); + }) + .catch(() => null); - if (!isDefined(thread) || !isDefined(thread.activeStreamId)) { - return; - } - - await this.agentChatStreamingService.reapDeadStream({ - thread, - workspaceId, - }); - } catch { - // The keep-alive tick must never die with the check + if (!isDefined(thread) || !isDefined(thread.activeStreamId)) { + return; } + + await this.agentChatStreamingService + .reapDeadStream({ thread, workspaceId }) + .catch(() => {}); } } diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-streaming.service.claim.spec.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-streaming.service.claim.spec.ts index 00c9e01cd8..adfacf8489 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-streaming.service.claim.spec.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/__tests__/agent-chat-streaming.service.claim.spec.ts @@ -1,5 +1,5 @@ import { type WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; -import { STREAM_INTERRUPTED_CODE } from 'src/engine/metadata-modules/ai/ai-chat/constants/stream-interrupted-code.constant'; +import { AiExceptionCode } from 'src/engine/metadata-modules/ai/ai.exception'; import { AgentChatStreamingService } from 'src/engine/metadata-modules/ai/ai-chat/services/agent-chat-streaming.service'; describe('AgentChatStreamingService claim & reap', () => { @@ -175,7 +175,7 @@ describe('AgentChatStreamingService claim & reap', () => { }); expect(reaped).toEqual( - expect.objectContaining({ code: STREAM_INTERRUPTED_CODE }), + expect.objectContaining({ code: AiExceptionCode.STREAM_INTERRUPTED }), ); expect(threadRepository.update).toHaveBeenCalledWith( 'workspace-id', @@ -183,7 +183,7 @@ describe('AgentChatStreamingService claim & reap', () => { expect.objectContaining({ activeStreamId: null, lastStreamError: expect.objectContaining({ - code: STREAM_INTERRUPTED_CODE, + code: AiExceptionCode.STREAM_INTERRUPTED, }), }), ); @@ -193,7 +193,7 @@ describe('AgentChatStreamingService claim & reap', () => { expect(publishedEvents).toContainEqual( expect.objectContaining({ type: 'stream-error', - code: STREAM_INTERRUPTED_CODE, + code: AiExceptionCode.STREAM_INTERRUPTED, }), ); }); diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat-streaming.service.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat-streaming.service.ts index 0d00945a2e..dbe29f5e76 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat-streaming.service.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat-streaming.service.ts @@ -16,7 +16,6 @@ import { InjectMessageQueue } from 'src/engine/core-modules/message-queue/decora import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queue.constants'; import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service'; import { type WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; -import { STREAM_INTERRUPTED_CODE } from 'src/engine/metadata-modules/ai/ai-chat/constants/stream-interrupted-code.constant'; import { AgentMessageRole, AgentMessageStatus, @@ -81,7 +80,7 @@ export class AgentChatStreamingService { } const interruptedError: AgentChatThreadLastStreamError = { - code: STREAM_INTERRUPTED_CODE, + code: AiExceptionCode.STREAM_INTERRUPTED, message: 'The response was interrupted before it could finish.', failedAt: new Date().toISOString(), }; diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai.exception.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai.exception.ts index 36a7d4b1d2..1fa1ccb9dc 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai.exception.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai.exception.ts @@ -22,6 +22,7 @@ export enum AiExceptionCode { ROLE_NOT_FOUND = 'ROLE_NOT_FOUND', ROLE_CANNOT_BE_ASSIGNED_TO_AGENTS = 'ROLE_CANNOT_BE_ASSIGNED_TO_AGENTS', NO_FAILED_TURN_TO_RETRY = 'NO_FAILED_TURN_TO_RETRY', + STREAM_INTERRUPTED = 'STREAM_INTERRUPTED', } const getAiExceptionUserFriendlyMessage = (code: AiExceptionCode) => { @@ -60,6 +61,8 @@ const getAiExceptionUserFriendlyMessage = (code: AiExceptionCode) => { return msg`This role cannot be assigned to agents.`; case AiExceptionCode.NO_FAILED_TURN_TO_RETRY: return msg`There is no failed message to retry.`; + case AiExceptionCode.STREAM_INTERRUPTED: + return msg`The response was interrupted before it could finish.`; default: assertUnreachable(code); } diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/utils/ai-graphql-api-exception-handler.util.ts b/packages/twenty-server/src/engine/metadata-modules/ai/utils/ai-graphql-api-exception-handler.util.ts index b54dff9c84..30411ef6e0 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/utils/ai-graphql-api-exception-handler.util.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/utils/ai-graphql-api-exception-handler.util.ts @@ -42,6 +42,7 @@ export const aiGraphqlApiExceptionHandler = (error: Error) => { case AiExceptionCode.AGENT_EXECUTION_FAILED: case AiExceptionCode.API_KEY_NOT_CONFIGURED: case AiExceptionCode.USER_WORKSPACE_ID_NOT_FOUND: + case AiExceptionCode.STREAM_INTERRUPTED: throw new InternalServerError(error); default: { return assertUnreachable(error.code);