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 44b38bdfb7..c2cb7d5751 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 @@ -3,6 +3,8 @@ import { type UIMessageChunk } from 'ai'; import { type WorkspaceEntity } from 'src/engine/core-modules/workspace/workspace.entity'; import { StreamAgentChatJob } from 'src/engine/metadata-modules/ai/ai-chat/jobs/stream-agent-chat.job'; import { type StreamAgentChatJobData } from 'src/engine/metadata-modules/ai/ai-chat/jobs/stream-agent-chat-job.types'; +import { AiExceptionCode } from 'src/engine/metadata-modules/ai/ai.exception'; + type PublishedEvent = { type: string } & Record; @@ -294,6 +296,36 @@ describe('StreamAgentChatJob', () => { ); }); + it('persists the error and unblocks the thread when the workspace is missing', async () => { + const { job, publishedEvents, threadRepository } = buildJob({ + workspaceFound: false, + }); + + await expect(job.handle(jobData)).rejects.toMatchObject({ + code: AiExceptionCode.WORKSPACE_NOT_FOUND, + }); + + expect(publishedEvents[publishedEvents.length - 1]).toMatchObject({ + type: 'stream-error', + code: AiExceptionCode.WORKSPACE_NOT_FOUND, + }); + expect(threadRepository.update).toHaveBeenCalledWith( + 'workspace-id', + { id: 'thread-id' }, + { + lastStreamError: expect.objectContaining({ + code: AiExceptionCode.WORKSPACE_NOT_FOUND, + }), + }, + ); + expect(threadRepository.update).toHaveBeenCalledWith( + 'workspace-id', + { id: 'thread-id', activeStreamId: 'stream-id' }, + { activeStreamId: null }, + ); + }); + + 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 3045a32ecc..2bf046325b 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 @@ -21,6 +21,10 @@ import { AgentMessageRole } from 'src/engine/metadata-modules/ai/ai-agent-execut import { computeCostBreakdown } from 'src/engine/metadata-modules/ai/ai-billing/utils/compute-cost-breakdown.util'; import { convertDollarsToBillingCredits } from 'src/engine/metadata-modules/ai/ai-billing/utils/convert-dollars-to-billing-credits.util'; import { extractCacheCreationTokens } from 'src/engine/metadata-modules/ai/ai-billing/utils/extract-cache-creation-tokens.util'; +import { + AiException, + AiExceptionCode, +} from 'src/engine/metadata-modules/ai/ai.exception'; import { AgentChatThreadEntity } from 'src/engine/metadata-modules/ai/ai-chat/entities/agent-chat-thread.entity'; import { AgentChatCancelSubscriberService } from 'src/engine/metadata-modules/ai/ai-chat/services/agent-chat-cancel-subscriber.service'; import { AgentChatEventPublisherService } from 'src/engine/metadata-modules/ai/ai-chat/services/agent-chat-event-publisher.service'; @@ -64,25 +68,6 @@ export class StreamAgentChatJob { async handle(data: StreamAgentChatJobData): Promise { await this.eventPublisherService.resetStreamState(data.threadId); - const workspace = await this.workspaceRepository.findOne({ - where: { id: data.workspaceId }, - }); - - if (!workspace) { - this.logger.error(`Workspace ${data.workspaceId} not found`); - await this.eventPublisherService.publish({ - threadId: data.threadId, - workspaceId: data.workspaceId, - event: { - type: 'stream-error', - code: 'WORKSPACE_NOT_FOUND', - message: `Workspace ${data.workspaceId} not found`, - }, - }); - - return; - } - const abortController = new AbortController(); const cancelChannel = getCancelChannel(data.threadId); @@ -91,6 +76,17 @@ export class StreamAgentChatJob { }); try { + const workspace = await this.workspaceRepository.findOne({ + where: { id: data.workspaceId }, + }); + + if (!workspace) { + throw new AiException( + `Workspace ${data.workspaceId} not found`, + AiExceptionCode.WORKSPACE_NOT_FOUND, + ); + } + await this.executeStream(data, workspace, abortController.signal); } catch (error) { this.logger.error( 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 cd4431d4b9..919d8e7650 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 @@ -11,6 +11,7 @@ export enum AiExceptionCode { AGENT_EXECUTION_FAILED = 'AGENT_EXECUTION_FAILED', INVALID_AGENT_INPUT = 'INVALID_AGENT_INPUT', THREAD_NOT_FOUND = 'THREAD_NOT_FOUND', + WORKSPACE_NOT_FOUND = 'WORKSPACE_NOT_FOUND', INVALID_CHAT_THREAD_TITLE = 'INVALID_CHAT_THREAD_TITLE', MESSAGE_NOT_FOUND = 'MESSAGE_NOT_FOUND', QUESTION_NOT_PENDING = 'QUESTION_NOT_PENDING', @@ -36,6 +37,8 @@ const getAiExceptionUserFriendlyMessage = (code: AiExceptionCode) => { return msg`Invalid agent input.`; case AiExceptionCode.THREAD_NOT_FOUND: return msg`Chat thread not found.`; + case AiExceptionCode.WORKSPACE_NOT_FOUND: + return msg`Workspace not found.`; case AiExceptionCode.INVALID_CHAT_THREAD_TITLE: return msg`Chat thread title cannot be empty.`; case AiExceptionCode.MESSAGE_NOT_FOUND: 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 1286dc3809..5b394de6c1 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 @@ -23,6 +23,7 @@ export const aiGraphqlApiExceptionHandler = (error: Error) => { switch (error.code) { case AiExceptionCode.AGENT_NOT_FOUND: case AiExceptionCode.THREAD_NOT_FOUND: + case AiExceptionCode.WORKSPACE_NOT_FOUND: case AiExceptionCode.MESSAGE_NOT_FOUND: case AiExceptionCode.ROLE_NOT_FOUND: throw new NotFoundError(error);