From 3b76ec528f1fd0ec476282e6c8abea91d32434c9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?F=C3=A9lix=20Malfait?= Date: Thu, 2 Jul 2026 21:11:02 +0200 Subject: [PATCH] fix(ai): route missing-workspace stream failures through the standard error path (#22480) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Rationale When `StreamAgentChatJob` can't find the workspace, it publishes a transient `stream-error` event and **returns before the try/finally exists** (`stream-agent-chat.job.ts`). Consequences on main: - `activeStreamId` is never cleared → every subsequent send in that thread queues behind a dead claim, forever; - no `lastStreamError` is persisted → nothing renders after a reload, and Retry has nothing to retry; - nothing throws → **zero telemetry**. Sentry confirms: the "Workspace not found" issues that exist are all auth/Stripe paths — this path fails in complete silence. ## Why this is the root cause, not a symptom patch The job's catch/finally already implement the correct failure contract for *every other* error: persist a typed `lastStreamError`, publish the typed event, release the claim guarded on the observed streamId. The bug is that one code path bypasses that contract via an early return. The fix removes the bypass — the lookup moves inside the `try` and throws a typed `AiException(WORKSPACE_NOT_FOUND)` — rather than duplicating cleanup in the early-return branch (which would be the symptom patch, and would drift the next time the contract changes). The alternative "prevent the job from existing when the workspace is gone" isn't achievable: workspace deletion between enqueue and pickup is an inherent race, so the job must handle it regardless. ## User impact A workspace deleted/deactivated mid-flight currently bricks the thread silently (the user just sees sends vanish into a queue). With this, the failure is visible (typed error message), recoverable (standard failed-turn state), and observable (real exception in monitoring). ## Stack Based on #22479 (spec harness) — it extends the same spec file with the regression test. `WORKSPACE_NOT_FOUND` is a TypeScript enum member, not a GraphQL schema change: no client-sdk regeneration needed. ## Test plan - [x] Regression test: missing workspace → typed rejection, `lastStreamError` persisted, terminal event published, claim released - [ ] CI green https://claude.ai/code/session_01Lyi6zTema2FMVVh8MD6c38 --- _Generated by [Claude Code](https://claude.ai/code/session_01Lyi6zTema2FMVVh8MD6c38)_ Review in cubic --- .../__tests__/stream-agent-chat.job.spec.ts | 32 +++++++++++++++++ .../ai/ai-chat/jobs/stream-agent-chat.job.ts | 34 ++++++++----------- .../metadata-modules/ai/ai.exception.ts | 3 ++ .../ai-graphql-api-exception-handler.util.ts | 1 + 4 files changed, 51 insertions(+), 19 deletions(-) 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);