From 29e5eb638d5134ffb984f7e814a3c4eb4a1c72d1 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Sat, 10 Oct 2026 12:52:11 -0700 Subject: [PATCH 1/2] fix(inbox): scope reply threading to the receiving workspace --- apps/sim/app/api/webhooks/agentmail/route.ts | 16 +- .../agentmail/thread-scope.integration.ts | 288 ++++++++++++++++++ .../sim/lib/mothership/chat/messages-store.ts | 14 +- .../sim/lib/mothership/inbox/executor.test.ts | 34 +-- apps/sim/lib/mothership/inbox/executor.ts | 55 ++-- 5 files changed, 368 insertions(+), 39 deletions(-) create mode 100644 apps/sim/app/api/webhooks/agentmail/thread-scope.integration.ts diff --git a/apps/sim/app/api/webhooks/agentmail/route.ts b/apps/sim/app/api/webhooks/agentmail/route.ts index 7e287f5386b..fc5767d8bbb 100644 --- a/apps/sim/app/api/webhooks/agentmail/route.ts +++ b/apps/sim/app/api/webhooks/agentmail/route.ts @@ -185,13 +185,20 @@ export const POST = withRouteHandler(async (req: Request) => { const emailMessageId = message.message_id const inReplyTo = message.in_reply_to || null + // Message ids are sender-controlled and shared with every recipient of a mail, + // so a task is only ever matched inside the workspace that owns this inbox. const [existingResult, isAllowed, recentCount, parentTaskResult, isEntitled] = await Promise.all([ emailMessageId ? db .select({ id: mothershipInboxTask.id }) .from(mothershipInboxTask) - .where(eq(mothershipInboxTask.emailMessageId, emailMessageId)) + .where( + and( + eq(mothershipInboxTask.workspaceId, result.id), + eq(mothershipInboxTask.emailMessageId, emailMessageId) + ) + ) .limit(1) : Promise.resolve([]), isSenderAllowed(fromEmail, result.id), @@ -200,7 +207,12 @@ export const POST = withRouteHandler(async (req: Request) => { ? db .select({ chatId: mothershipInboxTask.chatId }) .from(mothershipInboxTask) - .where(eq(mothershipInboxTask.responseMessageId, inReplyTo)) + .where( + and( + eq(mothershipInboxTask.workspaceId, result.id), + eq(mothershipInboxTask.responseMessageId, inReplyTo) + ) + ) .limit(1) : Promise.resolve([]), hasWorkspaceInboxAccess(result.id), diff --git a/apps/sim/app/api/webhooks/agentmail/thread-scope.integration.ts b/apps/sim/app/api/webhooks/agentmail/thread-scope.integration.ts new file mode 100644 index 00000000000..8541dba80c7 --- /dev/null +++ b/apps/sim/app/api/webhooks/agentmail/thread-scope.integration.ts @@ -0,0 +1,288 @@ +/** + * An inbox reply continues its parent task's chat, and the parent is found by the + * `In-Reply-To` header — a value every recipient of the agent's reply holds. Against + * real PostgreSQL, a reply delivered to one workspace's inbox must never reach a chat, + * task, or transcript owned by another workspace, while a reply inside the workspace + * still continues its thread. The webhook is signed with the real Svix scheme. + */ +import { authBanMock } from '@sim/testing/mocks/auth-ban.mock' +import { + billingAttributionMock, + billingAttributionMockFns, +} from '@sim/testing/mocks/billing-attribution.mock' +import { + billingSubscriptionMock, + billingSubscriptionMockFns, +} from '@sim/testing/mocks/billing-subscription.mock' +import { + mothershipHeadlessLifecycleMock, + mothershipHeadlessLifecycleMockFns, +} from '@sim/testing/mocks/mothership-headless-lifecycle.mock' +import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest' + +const { mockExecuteInboxTask, mockSendInboxResponse } = vi.hoisted(() => ({ + mockExecuteInboxTask: vi.fn(), + mockSendInboxResponse: vi.fn(), +})) + +vi.mock('@/lib/auth/ban', () => authBanMock) +vi.mock('@/lib/billing/core/billing-attribution', () => billingAttributionMock) +vi.mock('@/lib/billing/core/subscription', () => billingSubscriptionMock) +vi.mock('@/lib/mothership/request/lifecycle/headless', () => mothershipHeadlessLifecycleMock) +vi.mock('@/lib/mothership/request/lifecycle/start', () => ({ + requestChatTitle: async () => null, +})) +vi.mock('@/lib/mothership/inbox/response', () => ({ + sendInboxResponse: mockSendInboxResponse, +})) +vi.mock('@/lib/mothership/inbox/executor', () => ({ + executeInboxTask: mockExecuteInboxTask, +})) + +import { db } from '@sim/db' +import { + copilotChats, + copilotMessages, + mothershipInboxTask, + mothershipInboxWebhook, + permissions, + user, + workspace, +} from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { and, eq, inArray } from 'drizzle-orm' +import { NextRequest } from 'next/server' +import { Webhook } from 'svix' +import { POST } from '@/app/api/webhooks/agentmail/route' + +const { executeInboxTask } = await vi.importActual< + typeof import('@/lib/mothership/inbox/executor') +>('@/lib/mothership/inbox/executor') +const { mockRunHeadlessCopilotLifecycle } = mothershipHeadlessLifecycleMockFns + +interface InboxWorkspace { + id: string + ownerId: string + ownerEmail: string + inboxId: string + secret: string +} + +const victim = inboxWorkspace() +const attacker = inboxWorkspace() +const victimChatId = generateId() +const victimReplyMessageId = `<${generateId()}@agentmail.to>` +const victimChatUpdatedAt = new Date('2026-01-01T00:00:00.000Z') + +function inboxWorkspace(): InboxWorkspace { + const id = generateId() + return { + id, + ownerId: generateId(), + ownerEmail: `${id}@inbox-thread-scope.test`, + inboxId: `${id}@agentmail.to`, + secret: `whsec_${Buffer.from(generateId()).toString('base64')}`, + } +} + +async function seedWorkspace(ws: InboxWorkspace): Promise { + const now = new Date() + await db.insert(user).values({ + id: ws.ownerId, + name: 'Inbox thread scope fixture', + email: ws.ownerEmail, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + await db.insert(workspace).values({ + id: ws.id, + name: 'Inbox thread scope fixture', + ownerId: ws.ownerId, + billedAccountUserId: ws.ownerId, + inboxEnabled: true, + inboxAddress: ws.inboxId, + inboxProviderId: ws.inboxId, + }) + await db.insert(permissions).values({ + id: generateId(), + userId: ws.ownerId, + entityType: 'workspace', + entityId: ws.id, + permissionType: 'admin', + }) + await db.insert(mothershipInboxWebhook).values({ + id: generateId(), + workspaceId: ws.id, + webhookId: generateId(), + secret: ws.secret, + }) +} + +/** Delivers a `message.received` webhook to `ws`'s inbox from its owner, signed with its secret. */ +async function deliver( + ws: InboxWorkspace, + message: { messageId: string; inReplyTo?: string } +): Promise { + const body = JSON.stringify({ + event_type: 'message.received', + message: { + message_id: message.messageId, + thread_id: generateId(), + inbox_id: ws.inboxId, + from: `Owner <${ws.ownerEmail}>`, + to: [ws.inboxId], + subject: 'Re: quarterly numbers', + text: 'Please continue this thread.', + created_at: new Date().toISOString(), + ...(message.inReplyTo ? { in_reply_to: message.inReplyTo } : {}), + }, + }) + const svixId = `msg_${generateId()}` + const timestamp = new Date() + return POST( + new NextRequest('http://localhost:3000/api/webhooks/agentmail', { + method: 'POST', + headers: { + 'content-type': 'application/json', + 'svix-id': svixId, + 'svix-timestamp': String(Math.floor(timestamp.getTime() / 1000)), + 'svix-signature': new Webhook(ws.secret).sign(svixId, timestamp, body), + }, + body, + }) + ) +} + +async function taskFor(workspaceId: string, emailMessageId: string) { + const [task] = await db + .select() + .from(mothershipInboxTask) + .where( + and( + eq(mothershipInboxTask.workspaceId, workspaceId), + eq(mothershipInboxTask.emailMessageId, emailMessageId) + ) + ) + return task +} + +async function victimChatState() { + const [chat] = await db + .select({ updatedAt: copilotChats.updatedAt }) + .from(copilotChats) + .where(eq(copilotChats.id, victimChatId)) + const messages = await db + .select({ id: copilotMessages.id }) + .from(copilotMessages) + .where(eq(copilotMessages.chatId, victimChatId)) + return { updatedAt: chat?.updatedAt, messageCount: messages.length } +} + +describe('inbox reply threading stays inside the receiving workspace', () => { + beforeAll(async () => { + await seedWorkspace(victim) + await seedWorkspace(attacker) + await db.insert(copilotChats).values({ + id: victimChatId, + userId: victim.ownerId, + workspaceId: victim.id, + type: 'mothership', + updatedAt: victimChatUpdatedAt, + }) + await db.insert(mothershipInboxTask).values({ + id: generateId(), + workspaceId: victim.id, + fromEmail: victim.ownerEmail, + subject: 'quarterly numbers', + status: 'completed', + chatId: victimChatId, + responseMessageId: victimReplyMessageId, + }) + }) + + beforeEach(() => { + mockExecuteInboxTask.mockResolvedValue(undefined) + mockSendInboxResponse.mockResolvedValue(null) + billingSubscriptionMockFns.mockHasWorkspaceInboxAccess.mockResolvedValue(true) + billingAttributionMockFns.mockResolveBillingAttribution.mockResolvedValue({}) + mockRunHeadlessCopilotLifecycle.mockResolvedValue({ + success: true, + content: 'Here are the numbers.', + contentBlocks: [], + toolCalls: [], + }) + }) + + afterAll(async () => { + const ownerIds = [victim.ownerId, attacker.ownerId] + await db.delete(permissions).where(inArray(permissions.userId, ownerIds)) + await db.delete(workspace).where(inArray(workspace.id, [victim.id, attacker.id])) + await db.delete(user).where(inArray(user.id, ownerIds)) + }) + + it("does not adopt another workspace's chat for a reply to that workspace's message", async () => { + const messageId = `<${generateId()}@attacker.test>` + + const response = await deliver(attacker, { messageId, inReplyTo: victimReplyMessageId }) + + expect(response.status).toBe(200) + const task = await taskFor(attacker.id, messageId) + expect(task?.inReplyTo).toBe(victimReplyMessageId) + expect(task?.chatId).toBeNull() + }) + + it('continues the parent chat for a reply inside the same workspace', async () => { + const messageId = `<${generateId()}@victim.test>` + + await deliver(victim, { messageId, inReplyTo: victimReplyMessageId }) + + expect((await taskFor(victim.id, messageId))?.chatId).toBe(victimChatId) + }) + + it("accepts a message whose id another workspace's inbox already received", async () => { + const sharedMessageId = `<${generateId()}@cc-both.test>` + + await deliver(attacker, { messageId: sharedMessageId }) + await deliver(victim, { messageId: sharedMessageId }) + + expect(await taskFor(attacker.id, sharedMessageId)).toMatchObject({ status: 'received' }) + expect(await taskFor(victim.id, sharedMessageId)).toMatchObject({ status: 'received' }) + }) + + it("runs a task carrying another workspace's chat in a fresh chat of its own workspace", async () => { + const taskId = generateId() + await db.insert(mothershipInboxTask).values({ + id: taskId, + workspaceId: attacker.id, + fromEmail: attacker.ownerEmail, + subject: 'Re: quarterly numbers', + bodyText: 'Append me to the victim transcript.', + status: 'received', + chatId: victimChatId, + inReplyTo: victimReplyMessageId, + }) + + await executeInboxTask(taskId) + + expect(mockRunHeadlessCopilotLifecycle).toHaveBeenCalledOnce() + const [, options] = mockRunHeadlessCopilotLifecycle.mock.calls[0] + expect(options).toMatchObject({ workspaceId: attacker.id }) + expect(options.chatId).not.toBe(victimChatId) + const [runChat] = await db + .select({ workspaceId: copilotChats.workspaceId }) + .from(copilotChats) + .where(eq(copilotChats.id, options.chatId)) + expect(runChat?.workspaceId).toBe(attacker.id) + + const [task] = await db + .select({ status: mothershipInboxTask.status, chatId: mothershipInboxTask.chatId }) + .from(mothershipInboxTask) + .where(eq(mothershipInboxTask.id, taskId)) + expect(task).toEqual({ status: 'completed', chatId: options.chatId }) + expect(await victimChatState()).toEqual({ + updatedAt: victimChatUpdatedAt, + messageCount: 0, + }) + }) +}) diff --git a/apps/sim/lib/mothership/chat/messages-store.ts b/apps/sim/lib/mothership/chat/messages-store.ts index c303c502472..c4b8374f67d 100644 --- a/apps/sim/lib/mothership/chat/messages-store.ts +++ b/apps/sim/lib/mothership/chat/messages-store.ts @@ -99,16 +99,26 @@ export async function appendCopilotChatMessages( * because the user deleted the conversation after asking — resurrecting it with * a reply would undo that deletion, and the caller still has its reply in the * response. Throws on a write failure. + * + * `scope.workspaceId` confines the write to a chat in that workspace; a chat + * elsewhere receives nothing, exactly like a deleted one. */ export async function persistCopilotChatTurn( chatId: string, - messages: PersistedMessage[] + messages: PersistedMessage[], + scope?: { workspaceId: string } ): Promise { await db.transaction(async (tx) => { const [updated] = await tx .update(copilotChats) .set({ updatedAt: new Date() }) - .where(and(eq(copilotChats.id, chatId), isNull(copilotChats.deletedAt))) + .where( + and( + eq(copilotChats.id, chatId), + isNull(copilotChats.deletedAt), + ...(scope ? [eq(copilotChats.workspaceId, scope.workspaceId)] : []) + ) + ) .returning({ model: copilotChats.model }) if (!updated) return await appendCopilotChatMessages(chatId, messages, { chatModel: updated.model ?? null }, tx) diff --git a/apps/sim/lib/mothership/inbox/executor.test.ts b/apps/sim/lib/mothership/inbox/executor.test.ts index b825db1c975..bc42535ec1e 100644 --- a/apps/sim/lib/mothership/inbox/executor.test.ts +++ b/apps/sim/lib/mothership/inbox/executor.test.ts @@ -133,6 +133,12 @@ const WORKSPACE = { inboxMountedSecrets: ['INBOX_KEY'], } +/** Queues a task and, when it continues a chat, that chat as found in the task's workspace. */ +function queueInboxTask(task: { chatId: string | null }) { + queueTableRows(schemaMock.mothershipInboxTask, [task]) + if (task.chatId) queueTableRows(schemaMock.copilotChats, [{ id: task.chatId }]) +} + describe('Inbox execution actor', () => { beforeEach(() => { resetDbChainMock() @@ -153,15 +159,13 @@ describe('Inbox execution actor', () => { chat: { id: 'chat-1' }, isNew: true, }) - dbChainMockFns.returning - .mockResolvedValueOnce([{ id: 'task-1' }]) - .mockResolvedValueOnce([{ model: 'claude-opus-4-8' }]) + dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'task-1' }]) }) it('sends a valid inbox turn without eagerly loading integration schemas', async () => { const chatId = '44444444-4444-4444-8444-444444444444' const workspaceId = '55555555-5555-4555-8555-555555555555' - queueTableRows(schemaMock.mothershipInboxTask, [{ ...INBOX_TASK, chatId, workspaceId }]) + queueInboxTask({ ...INBOX_TASK, chatId, workspaceId }) queueTableRows(schemaMock.workspace, [{ ...WORKSPACE, id: workspaceId }]) queueTableRows(schemaMock.user, [{ id: 'member-1' }]) mockGetUserEntityPermissions.mockResolvedValue('write') @@ -186,7 +190,7 @@ describe('Inbox execution actor', () => { }) it('gives a workspace member their own raw-secret authority', async () => { - queueTableRows(schemaMock.mothershipInboxTask, [INBOX_TASK]) + queueInboxTask(INBOX_TASK) queueTableRows(schemaMock.workspace, [WORKSPACE]) queueTableRows(schemaMock.user, [{ id: 'member-1' }]) mockGetUserEntityPermissions.mockResolvedValue('write') @@ -209,7 +213,7 @@ describe('Inbox execution actor', () => { }) it('does not lend a read-only member write authority', async () => { - queueTableRows(schemaMock.mothershipInboxTask, [INBOX_TASK]) + queueInboxTask(INBOX_TASK) queueTableRows(schemaMock.workspace, [WORKSPACE]) queueTableRows(schemaMock.user, [{ id: 'member-1' }]) mockCheckWorkspaceAccess.mockResolvedValue({ permission: 'read' }) @@ -231,7 +235,7 @@ describe('Inbox execution actor', () => { * secret actor already refuses for a direct mount. */ it('caps an external sender at read even when the owner is an admin', async () => { - queueTableRows(schemaMock.mothershipInboxTask, [INBOX_TASK]) + queueInboxTask(INBOX_TASK) queueTableRows(schemaMock.workspace, [WORKSPACE]) queueTableRows(schemaMock.user, []) mockCheckWorkspaceAccess.mockResolvedValue({ permission: 'admin' }) @@ -254,7 +258,7 @@ describe('Inbox execution actor', () => { }) it('stamps the shared mothership model on the chat it creates for a task', async () => { - queueTableRows(schemaMock.mothershipInboxTask, [{ ...INBOX_TASK, chatId: null }]) + queueInboxTask({ ...INBOX_TASK, chatId: null }) queueTableRows(schemaMock.workspace, [WORKSPACE]) queueTableRows(schemaMock.user, [{ id: 'member-1' }]) mockGetUserEntityPermissions.mockResolvedValue('write') @@ -271,7 +275,7 @@ describe('Inbox execution actor', () => { }) it('leaves an external sender with no permission at none rather than promoting to read', async () => { - queueTableRows(schemaMock.mothershipInboxTask, [INBOX_TASK]) + queueInboxTask(INBOX_TASK) queueTableRows(schemaMock.workspace, [WORKSPACE]) queueTableRows(schemaMock.user, []) mockCheckWorkspaceAccess.mockResolvedValue({ permission: null }) @@ -285,9 +289,7 @@ describe('Inbox execution actor', () => { it.each(['member', 'external'])( 'makes inbox attachments readable without increasing %s tool authority', async (actor) => { - queueTableRows(schemaMock.mothershipInboxTask, [ - { ...INBOX_TASK, hasAttachments: true, agentmailMessageId: 'mail-1' }, - ]) + queueInboxTask({ ...INBOX_TASK, hasAttachments: true, agentmailMessageId: 'mail-1' }) queueTableRows(schemaMock.workspace, [WORKSPACE]) queueTableRows(schemaMock.user, actor === 'member' ? [{ id: 'member-1' }] : []) mockGetUserEntityPermissions.mockResolvedValue('write') @@ -343,9 +345,7 @@ describe('Inbox execution actor', () => { it.each(['download', 'binding', 'declared-size', 'actual-size'])( 'keeps a valid sibling readable when an attachment fails during %s', async (failure) => { - queueTableRows(schemaMock.mothershipInboxTask, [ - { ...INBOX_TASK, hasAttachments: true, agentmailMessageId: 'mail-1' }, - ]) + queueInboxTask({ ...INBOX_TASK, hasAttachments: true, agentmailMessageId: 'mail-1' }) queueTableRows(schemaMock.workspace, [WORKSPACE]) queueTableRows(schemaMock.user, [{ id: 'member-1' }]) mockGetUserEntityPermissions.mockResolvedValue('write') @@ -409,9 +409,7 @@ describe('Inbox execution actor', () => { ) it('does not promise readable attachments when metadata cannot be loaded', async () => { - queueTableRows(schemaMock.mothershipInboxTask, [ - { ...INBOX_TASK, hasAttachments: true, agentmailMessageId: 'mail-1' }, - ]) + queueInboxTask({ ...INBOX_TASK, hasAttachments: true, agentmailMessageId: 'mail-1' }) queueTableRows(schemaMock.workspace, [WORKSPACE]) queueTableRows(schemaMock.user, [{ id: 'member-1' }]) mockGetUserEntityPermissions.mockResolvedValue('write') diff --git a/apps/sim/lib/mothership/inbox/executor.ts b/apps/sim/lib/mothership/inbox/executor.ts index b79eef17960..c738667b8b5 100644 --- a/apps/sim/lib/mothership/inbox/executor.ts +++ b/apps/sim/lib/mothership/inbox/executor.ts @@ -6,7 +6,7 @@ import { and, eq, isNull, sql } from 'drizzle-orm' import { getActivelyBannedUserIds, isEmailBlocked } from '@/lib/auth/ban' import { resolveBillingAttribution } from '@/lib/billing/core/billing-attribution' import { resolveOrCreateChat } from '@/lib/mothership/chat/lifecycle' -import { appendCopilotChatMessages } from '@/lib/mothership/chat/messages-store' +import { persistCopilotChatTurn } from '@/lib/mothership/chat/messages-store' import { buildPersistedAssistantMessage, buildPersistedUserMessage, @@ -19,7 +19,7 @@ import * as agentmail from '@/lib/mothership/inbox/agentmail-client' import { prepareInboxAttachments } from '@/lib/mothership/inbox/attachments' import { formatEmailAsMessage } from '@/lib/mothership/inbox/format' import { sendInboxResponse } from '@/lib/mothership/inbox/response' -import type { AgentMailAttachment } from '@/lib/mothership/inbox/types' +import type { AgentMailAttachment, InboxTask } from '@/lib/mothership/inbox/types' import { runHeadlessCopilotLifecycle } from '@/lib/mothership/request/lifecycle/headless' import { requestChatTitle } from '@/lib/mothership/request/lifecycle/start' import type { OrchestratorResult } from '@/lib/mothership/request/types' @@ -83,7 +83,7 @@ export async function executeInboxTask(taskId: string): Promise { return } - let chatId = inboxTask.chatId + let chatId = await resolveWorkspaceChatId(inboxTask) let responseSent = false try { @@ -279,6 +279,7 @@ export async function executeInboxTask(taskId: string): Promise { if (chatId) { await persistChatMessages( chatId, + ws.id, userMessageId, messageContent, { @@ -347,6 +348,38 @@ export async function executeInboxTask(taskId: string): Promise { } } +/** + * The chat a task continues, only when it is a live chat in the task's + * workspace. A reply inherits its parent task's chat, so a chat from any other + * workspace (or one since deleted) is dropped and the run starts a fresh chat + * instead of writing across tenants. + */ +async function resolveWorkspaceChatId( + inboxTask: Pick +): Promise { + if (!inboxTask.chatId) return null + + const [chat] = await db + .select({ id: copilotChats.id }) + .from(copilotChats) + .where( + and( + eq(copilotChats.id, inboxTask.chatId), + eq(copilotChats.workspaceId, inboxTask.workspaceId), + isNull(copilotChats.deletedAt) + ) + ) + .limit(1) + if (chat) return chat.id + + logger.warn('Dropping inbox task chat that is not live in the task workspace', { + taskId: inboxTask.id, + chatId: inboxTask.chatId, + workspaceId: inboxTask.workspaceId, + }) + return null +} + /** * Resolve the execution and raw-secret actors independently. Workspace members * execute and mount secrets as themselves. External senders retain the existing @@ -423,6 +456,7 @@ async function resolveInboxExecutionActor( */ async function persistChatMessages( chatId: string, + workspaceId: string, userMessageId: string, userContent: string, result: OrchestratorResult, @@ -439,20 +473,7 @@ async function persistChatMessages( // Best-effort: the email response is the primary deliverable, so a failure // here is logged (in the catch below) rather than failing the task. - await db.transaction(async (tx) => { - const [updated] = await tx - .update(copilotChats) - .set({ updatedAt: new Date() }) - .where(eq(copilotChats.id, chatId)) - .returning({ model: copilotChats.model }) - if (!updated) return - await appendCopilotChatMessages( - chatId, - [userMessage, assistantMessage], - { chatModel: updated.model ?? null }, - tx - ) - }) + await persistCopilotChatTurn(chatId, [userMessage, assistantMessage], { workspaceId }) } catch (err) { logger.error('Failed to persist chat messages', { chatId, From 56feab9d23cd6277d10ff3953a84cfe8c24b34d9 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Sat, 10 Oct 2026 12:58:11 -0700 Subject: [PATCH 2/2] fix(inbox): resolve the inherited chat inside the guarded run and cover the scoped persist --- .../agentmail/thread-scope.integration.ts | 33 +++++++++++++++---- apps/sim/lib/mothership/inbox/executor.ts | 6 ++-- 2 files changed, 31 insertions(+), 8 deletions(-) diff --git a/apps/sim/app/api/webhooks/agentmail/thread-scope.integration.ts b/apps/sim/app/api/webhooks/agentmail/thread-scope.integration.ts index 8541dba80c7..c165e118950 100644 --- a/apps/sim/app/api/webhooks/agentmail/thread-scope.integration.ts +++ b/apps/sim/app/api/webhooks/agentmail/thread-scope.integration.ts @@ -53,6 +53,8 @@ import { generateId } from '@sim/utils/id' import { and, eq, inArray } from 'drizzle-orm' import { NextRequest } from 'next/server' import { Webhook } from 'svix' +import { persistCopilotChatTurn } from '@/lib/mothership/chat/messages-store' +import { buildPersistedUserMessage } from '@/lib/mothership/chat/persisted-message' import { POST } from '@/app/api/webhooks/agentmail/route' const { executeInboxTask } = await vi.importActual< @@ -150,7 +152,8 @@ async function deliver( 'svix-signature': new Webhook(ws.secret).sign(svixId, timestamp, body), }, body, - }) + }), + { params: Promise.resolve({}) } ) } @@ -167,16 +170,20 @@ async function taskFor(workspaceId: string, emailMessageId: string) { return task } +async function messageCount(chatId: string): Promise { + const messages = await db + .select({ id: copilotMessages.id }) + .from(copilotMessages) + .where(eq(copilotMessages.chatId, chatId)) + return messages.length +} + async function victimChatState() { const [chat] = await db .select({ updatedAt: copilotChats.updatedAt }) .from(copilotChats) .where(eq(copilotChats.id, victimChatId)) - const messages = await db - .select({ id: copilotMessages.id }) - .from(copilotMessages) - .where(eq(copilotMessages.chatId, victimChatId)) - return { updatedAt: chat?.updatedAt, messageCount: messages.length } + return { updatedAt: chat?.updatedAt, messageCount: await messageCount(victimChatId) } } describe('inbox reply threading stays inside the receiving workspace', () => { @@ -280,6 +287,20 @@ describe('inbox reply threading stays inside the receiving workspace', () => { .from(mothershipInboxTask) .where(eq(mothershipInboxTask.id, taskId)) expect(task).toEqual({ status: 'completed', chatId: options.chatId }) + expect(await messageCount(options.chatId)).toBe(2) + expect(await victimChatState()).toEqual({ + updatedAt: victimChatUpdatedAt, + messageCount: 0, + }) + }) + + it("writes nothing when a turn is persisted into another workspace's chat", async () => { + await persistCopilotChatTurn( + victimChatId, + [buildPersistedUserMessage({ id: generateId(), content: 'Injected turn' })], + { workspaceId: attacker.id } + ) + expect(await victimChatState()).toEqual({ updatedAt: victimChatUpdatedAt, messageCount: 0, diff --git a/apps/sim/lib/mothership/inbox/executor.ts b/apps/sim/lib/mothership/inbox/executor.ts index c738667b8b5..776a8254aad 100644 --- a/apps/sim/lib/mothership/inbox/executor.ts +++ b/apps/sim/lib/mothership/inbox/executor.ts @@ -83,19 +83,21 @@ export async function executeInboxTask(taskId: string): Promise { return } - let chatId = await resolveWorkspaceChatId(inboxTask) + let chatId: string | null = null let responseSent = false try { - const [[claimed], actor] = await Promise.all([ + const [[claimed], actor, workspaceChatId] = await Promise.all([ db .update(mothershipInboxTask) .set({ status: 'processing', processingStartedAt: new Date() }) .where(and(eq(mothershipInboxTask.id, taskId), eq(mothershipInboxTask.status, 'received'))) .returning({ id: mothershipInboxTask.id }), resolveInboxExecutionActor(inboxTask.fromEmail, ws), + resolveWorkspaceChatId(inboxTask), ]) const userId = actor.executionUserId + chatId = workspaceChatId if (!claimed) { logger.info('Task already claimed by another execution, skipping', { taskId })