diff --git a/apps/sim/app/api/v2/chat/route.test.ts b/apps/sim/app/api/v2/chat/route.test.ts index 70b87a051f6..66d3084f8f8 100644 --- a/apps/sim/app/api/v2/chat/route.test.ts +++ b/apps/sim/app/api/v2/chat/route.test.ts @@ -267,7 +267,6 @@ describe('POST /api/v2/chat', () => { mockResolveOrCreateChat.mockResolvedValue({ chatId: SERVER_ISSUED_CHAT_ID, chat: chatRow(SERVER_ISSUED_CHAT_ID), - conversationHistory: [], isNew: true, }) }) @@ -444,7 +443,6 @@ describe('POST /api/v2/chat', () => { mockResolveOrCreateChat.mockResolvedValue({ chatId: OWNED_CONVERSATION_ID, chat: chatRow(OWNED_CONVERSATION_ID), - conversationHistory: [], isNew: false, }) @@ -469,38 +467,10 @@ describe('POST /api/v2/chat', () => { }) }) - it('posts only the current turn on a resumed conversation, never the stored transcript', async () => { - mockResolveOrCreateChat.mockResolvedValue({ - chatId: OWNED_CONVERSATION_ID, - chat: chatRow(OWNED_CONVERSATION_ID), - conversationHistory: [ - { role: 'user', content: 'first' }, - { role: 'assistant', content: 'first reply' }, - ], - isNew: false, - }) - - const response = await callChat({ - workspaceId: 'workspace-1', - message: 'and then?', - conversationId: OWNED_CONVERSATION_ID, - }) - - expect(response.status).toBe(200) - // Continuity is keyed by chatId downstream, exactly as the web send path - // and the Sim Chat block do. Replaying the transcript here would duplicate - // every prior turn. - expect(mockRunHeadlessCopilotLifecycle.mock.calls[0][0]).toMatchObject({ - message: 'and then?', - chatId: OWNED_CONVERSATION_ID, - }) - }) - it('answers 404 and runs nothing when the resolver refuses the named conversation', async () => { mockResolveOrCreateChat.mockResolvedValue({ chatId: OWNED_CONVERSATION_ID, chat: null, - conversationHistory: [], isNew: false, }) diff --git a/apps/sim/app/api/v2/chat/route.ts b/apps/sim/app/api/v2/chat/route.ts index df60c0d0f52..46b9beaa742 100644 --- a/apps/sim/app/api/v2/chat/route.ts +++ b/apps/sim/app/api/v2/chat/route.ts @@ -264,12 +264,11 @@ export const POST = withRouteHandler( // surface uses, and refuse every id that does not resolve with the same // response so the refusal carries no information about the id. Omitting // the id mints a server-issued conversation instead of trusting one. - // The resolved transcript is deliberately not forwarded: continuity is - // keyed by `chatId` downstream, exactly as the web send path and the Sim - // Chat block do, both of which post a single message with a chat id. + // Continuity is keyed by `chatId` downstream, exactly as the web send + // path and the Sim Chat block do, both of which post a single message + // with a chat id. const resolvedChat = await resolveOrCreateChat({ ...(conversationId ? { chatId: conversationId } : {}), - includeTranscript: false, userId, workspaceId, model: MOTHERSHIP_CHAT_DEFAULT_MODEL, diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx index 3b53040dbc2..95fe0c04e6c 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx @@ -82,6 +82,7 @@ import { } from '@/app/workspace/[workspaceId]/home/hooks/send-handoff' import { useChat } from '@/app/workspace/[workspaceId]/home/hooks/use-chat' import { type MothershipChatHistory, mothershipChatKeys } from '@/hooks/queries/mothership-chats' +import { handleMothershipChatStatusEvent } from '@/hooks/use-mothership-chat-events' import { useExecutionStore } from '@/stores/execution/store' import { useMothershipQueueStore } from '@/stores/mothership-queue/store' @@ -1821,4 +1822,81 @@ describe('useChat remount send recovery', () => { expect(getResult().messageQueue.map((entry) => entry.id)).toEqual(['unsent-entry']) expect(state.postBodies).toHaveLength(0) }) + + it('loads the saved transcript once when its own stream completes', async () => { + const chatId = 'chat-own-completion' + const history: MothershipChatHistory = { + id: chatId, + mode: 'agent', + title: 'Own stream', + messages: [], + activeStreamId: null, + resources: [], + } + const saved = [ + { id: 'saved-user', role: 'user', content: 'Summarize the run' }, + { id: 'saved-assistant', role: 'assistant', content: 'Done.' }, + ] + const detailRequests: string[] = [] + mockRequestJson.mockImplementation((contract: AnyApiRouteContract) => { + if (contract.path !== '/api/mothership/chats/[chatId]') { + return Promise.resolve({ chats: [] }) + } + detailRequests.push(chatId) + return Promise.resolve({ chat: { ...history, messages: saved } }) + }) + let stream: ReadableStreamDefaultController | undefined + let streamId: string | undefined + const emit = (event: Omit) => + stream?.enqueue( + new TextEncoder().encode( + `data: ${JSON.stringify({ v: 1, ts: '', stream: { streamId }, ...event })}\n\n` + ) + ) + vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => { + if (String(input) !== '/api/mothership/chat' || init?.method !== 'POST') { + return fetchStub(input, init) + } + streamId = JSON.parse(String(init.body)).userMessageId + return new Response( + new ReadableStream({ + start(controller) { + stream = controller + }, + }), + { headers: { 'Content-Type': 'text/event-stream', 'x-mothership-chat-id': chatId } } + ) + }) + const { getResult } = renderUseChatInChat(chatId, history) + + await act(async () => { + void getResult().sendMessage('Summarize the run') + }) + await waitFor(() => stream !== undefined) + emit({ seq: 1, type: 'text', payload: { channel: 'assistant', text: 'Done.' } }) + await waitFor( + () => + queryClient.getQueryData(mothershipChatKeys.detail(chatId)) + ?.activeStreamId === streamId + ) + /** The server publishes `completed` after persisting and before closing the stream. */ + handleMothershipChatStatusEvent(queryClient, 'ws-1', { + chatId, + type: 'completed', + streamId, + }) + emit({ seq: 2, type: 'complete', payload: { status: 'complete' } }) + stream?.close() + + await waitFor(() => !getResult().isSending && detailRequests.length > 0) + await act(async () => { + await sleep(50) + }) + expect(detailRequests).toHaveLength(1) + expect( + queryClient + .getQueryData(mothershipChatKeys.detail(chatId)) + ?.messages.map((message) => message.id) + ).toEqual(['saved-user', 'saved-assistant']) + }) }) diff --git a/apps/sim/hooks/use-mothership-chat-events.test.ts b/apps/sim/hooks/use-mothership-chat-events.test.ts index 28d2f44311f..940e04fe113 100644 --- a/apps/sim/hooks/use-mothership-chat-events.test.ts +++ b/apps/sim/hooks/use-mothership-chat-events.test.ts @@ -1,4 +1,5 @@ -import type { QueryClient } from '@tanstack/react-query' +import { sleep } from '@sim/utils/helpers' +import { QueryClient, QueryObserver } from '@tanstack/react-query' import { beforeEach, describe, expect, it, vi } from 'vitest' const { suspendBrowserScope, suspendTerminalScope } = vi.hoisted(() => ({ @@ -9,7 +10,7 @@ const { suspendBrowserScope, suspendTerminalScope } = vi.hoisted(() => ({ vi.mock('@/lib/browser-agent/transport', () => ({ suspendBrowserScope })) vi.mock('@/lib/terminal/transport', () => ({ suspendTerminalScope })) -import { mothershipChatKeys } from '@/hooks/queries/mothership-chats' +import { type MothershipChatHistory, mothershipChatKeys } from '@/hooks/queries/mothership-chats' import { handleMothershipChatStatusEvent, resyncMothershipChatCaches, @@ -184,6 +185,80 @@ describe('handleMothershipChatStatusEvent', () => { ) }) +describe('chat detail refetches driven by status events', () => { + function mountDetail(cached: MothershipChatHistory) { + const queryClient = new QueryClient() + const fetchTranscript = vi.fn(async () => cached) + queryClient.setQueryData(mothershipChatKeys.detail('chat-1'), cached) + const unsubscribe = new QueryObserver(queryClient, { + queryKey: mothershipChatKeys.detail('chat-1'), + queryFn: fetchTranscript, + staleTime: Number.POSITIVE_INFINITY, + }).subscribe(() => {}) + return { queryClient, fetchTranscript, unsubscribe } + } + + const liveStream: MothershipChatHistory = { + id: 'chat-1', + title: null, + messages: [ + { id: 'stream-1' }, + { id: 'live-assistant:stream-1' }, + ] as MothershipChatHistory['messages'], + activeStreamId: 'stream-1', + resources: [], + } + + it('does not reload the transcript when the viewer finishes its own live stream', async () => { + const { queryClient, fetchTranscript, unsubscribe } = mountDetail(liveStream) + + handleMothershipChatStatusEvent(queryClient, 'ws-1', { + chatId: 'chat-1', + type: 'completed', + streamId: 'stream-1', + }) + await sleep(0) + + expect(fetchTranscript).not.toHaveBeenCalled() + unsubscribe() + }) + + it('reloads the saved transcript when a cached mid-stream detail is opened after completion', async () => { + const queryClient = new QueryClient() + const fetchTranscript = vi.fn(async () => ({ ...liveStream, activeStreamId: null })) + queryClient.setQueryData(mothershipChatKeys.detail('chat-1'), liveStream) + + handleMothershipChatStatusEvent(queryClient, 'ws-1', { + chatId: 'chat-1', + type: 'completed', + streamId: 'stream-1', + }) + const unsubscribe = new QueryObserver(queryClient, { + queryKey: mothershipChatKeys.detail('chat-1'), + queryFn: fetchTranscript, + staleTime: Number.POSITIVE_INFINITY, + }).subscribe(() => {}) + + await vi.waitFor(() => expect(fetchTranscript).toHaveBeenCalledTimes(1)) + unsubscribe() + }) + + it('marks the detail stale on rename without reloading the transcript', async () => { + const { queryClient, fetchTranscript, unsubscribe } = mountDetail({ + ...liveStream, + messages: [], + activeStreamId: null, + }) + + handleMothershipChatStatusEvent(queryClient, 'ws-1', { chatId: 'chat-1', type: 'renamed' }) + await sleep(0) + + expect(fetchTranscript).not.toHaveBeenCalled() + expect(queryClient.getQueryState(mothershipChatKeys.detail('chat-1'))?.isInvalidated).toBe(true) + unsubscribe() + }) +}) + describe('resyncMothershipChatCaches', () => { const queryClient = { invalidateQueries: vi.fn().mockResolvedValue(undefined), diff --git a/apps/sim/hooks/use-mothership-chat-events.ts b/apps/sim/hooks/use-mothership-chat-events.ts index 505fa8a0f02..eca7b80ec88 100644 --- a/apps/sim/hooks/use-mothership-chat-events.ts +++ b/apps/sim/hooks/use-mothership-chat-events.ts @@ -33,12 +33,6 @@ interface ChatStatusEventPayload { streamId?: string } -const DETAIL_INVALIDATING_CHAT_STATUS_TYPES = new Set([ - 'started', - 'completed', - 'renamed', -]) - function isChatStatusEventType(value: unknown): value is ChatStatusEventType { return typeof value === 'string' && CHAT_STATUS_TYPE_SET.has(value) } @@ -69,7 +63,6 @@ function shouldSkipDetailInvalidationForStreamEvent( current: MothershipChatHistory | undefined, payload: ChatStatusEventPayload ) { - if (payload.type !== 'started' && payload.type !== 'completed') return false if (!current?.activeStreamId) return false if (!payload.streamId) return isLocalOptimisticActiveStream(current) if (payload.type === 'started' && current.activeStreamId === payload.streamId) return true @@ -128,17 +121,37 @@ export function handleMothershipChatStatusEvent( queryClient.removeQueries({ queryKey: mothershipChatKeys.detail(payload.chatId) }) return } - if (payload.type === 'started' || payload.type === 'completed') { - const current = queryClient.getQueryData( - mothershipChatKeys.detail(payload.chatId) - ) - if (shouldSkipDetailInvalidationForStreamEvent(current, payload)) { - return - } - } - if (payload.type && DETAIL_INVALIDATING_CHAT_STATUS_TYPES.has(payload.type)) { - queryClient.invalidateQueries({ queryKey: mothershipChatKeys.detail(payload.chatId) }) + if (payload.type === 'renamed') { + /** + * The lists invalidated above carry the title every surface renders; the + * detail only needs marking stale, not a full transcript reload. + */ + queryClient.invalidateQueries({ + queryKey: mothershipChatKeys.detail(payload.chatId), + refetchType: 'none', + }) + return } + if (payload.type !== 'started' && payload.type !== 'completed') return + const current = queryClient.getQueryData( + mothershipChatKeys.detail(payload.chatId) + ) + if (shouldSkipDetailInvalidationForStreamEvent(current, payload)) return + /** + * A completion of the cached live stream only marks the detail stale. The + * server persists the turn before closing the stream, so a surface rendering + * it refetches through its own finalization; the live message alone cannot + * tell this tab's stream from a server-loaded mid-stream snapshot, so any + * other cached copy reloads the saved transcript on its next mount. + */ + const completesCachedLiveStream = + payload.type === 'completed' && + current?.activeStreamId === payload.streamId && + isLocalOptimisticActiveStream(current) + queryClient.invalidateQueries({ + queryKey: mothershipChatKeys.detail(payload.chatId), + ...(completesCachedLiveStream ? { refetchType: 'none' as const } : {}), + }) } /** diff --git a/apps/sim/lib/mothership/chat/chat-mcp-servers.integration.ts b/apps/sim/lib/mothership/chat/chat-mcp-servers.integration.ts new file mode 100644 index 00000000000..7845d3ecaab --- /dev/null +++ b/apps/sim/lib/mothership/chat/chat-mcp-servers.integration.ts @@ -0,0 +1,166 @@ +/** Exercises MCP server inheritance against real PostgreSQL transcripts written by the append path. */ +import { db } from '@sim/db' +import { copilotChats, copilotMessages, user } from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { and, eq } from 'drizzle-orm' +import { afterAll, beforeAll, describe, expect, it } from 'vitest' +import { loadChatMcpServerIds, loadCopilotChatMessages } from '@/lib/mothership/chat/lifecycle' +import { appendCopilotChatMessages } from '@/lib/mothership/chat/messages-store' +import { + buildPersistedUserMessage, + type PersistedMessage, +} from '@/lib/mothership/chat/persisted-message' + +/** The transcript-based collection this query replaced, kept as the equivalence oracle. */ +function collectFromTranscript(messages: PersistedMessage[]): string[] { + const serverIds = new Set() + for (const message of messages) { + if (!Array.isArray(message.contexts)) continue + for (const ctx of message.contexts) { + if (ctx.kind === 'mcp' && typeof ctx.serverId === 'string' && ctx.serverId) { + serverIds.add(ctx.serverId) + } + } + } + return Array.from(serverIds) +} + +function userTurn(contexts: PersistedMessage['contexts'], content = 'turn'): PersistedMessage { + return buildPersistedUserMessage({ id: generateId(), content, contexts }) +} + +function assistantTurn(): PersistedMessage { + return { + id: generateId(), + role: 'assistant', + content: 'reply', + timestamp: new Date().toISOString(), + } +} + +describe('chat MCP server inheritance in PostgreSQL', () => { + const userId = generateId() + + async function createChat(): Promise { + const [chat] = await db + .insert(copilotChats) + .values({ userId, type: 'mothership' }) + .returning({ id: copilotChats.id }) + return chat.id + } + + beforeAll(async () => { + const now = new Date() + await db.insert(user).values({ + id: userId, + name: 'MCP inheritance fixture', + email: `${userId}@fixture.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + }) + + afterAll(async () => { + await db.delete(copilotChats).where(eq(copilotChats.userId, userId)) + await db.delete(user).where(eq(user.id, userId)) + await db.$client.end() + }) + + it('keeps servers tagged on earlier turns, in first-tagged order, matching the transcript', async () => { + const chatId = await createChat() + await appendCopilotChatMessages(chatId, [ + userTurn([ + { kind: 'mcp', serverId: 'server-b', label: 'B' }, + { kind: 'workflow', workflowId: 'wf-1', label: 'Workflow' }, + { kind: 'mcp', serverId: 'server-c', label: 'C' }, + ]), + assistantTurn(), + userTurn(undefined), + userTurn([ + { kind: 'mcp', serverId: 'server-a', label: 'A' }, + { kind: 'mcp', serverId: 'server-b', label: 'B again' }, + ]), + assistantTurn(), + ]) + + const serverIds = await loadChatMcpServerIds(chatId) + + expect(serverIds).toEqual(['server-b', 'server-c', 'server-a']) + expect(serverIds).toEqual(collectFromTranscript(await loadCopilotChatMessages(chatId))) + }) + + it('orders by transcript position, not by row insertion order', async () => { + const chatId = await createChat() + await db.insert(copilotMessages).values( + [ + { seq: 1, serverId: 'server-second' }, + { seq: 0, serverId: 'server-first' }, + ].map(({ seq, serverId }) => ({ + chatId, + messageId: generateId(), + role: 'user', + seq, + content: { role: 'user', content: 'x', contexts: [{ kind: 'mcp', serverId }] }, + })) + ) + + expect(await loadChatMcpServerIds(chatId)).toEqual(['server-first', 'server-second']) + }) + + it('ignores malformed contexts instead of failing the turn', async () => { + const chatId = await createChat() + await db.insert(copilotMessages).values([ + { + chatId, + messageId: generateId(), + role: 'user', + seq: 0, + content: { role: 'user', content: 'x', contexts: { kind: 'mcp', serverId: 'object' } }, + }, + { + chatId, + messageId: generateId(), + role: 'user', + seq: 1, + content: { + role: 'user', + content: 'y', + contexts: [ + 'mcp', + { kind: 'mcp', serverId: '' }, + { kind: 'mcp', serverId: 7 }, + { kind: 'mcp', serverId: 'server-ok' }, + ], + }, + }, + ]) + + expect(await loadChatMcpServerIds(chatId)).toEqual(['server-ok']) + }) + + it('does not inherit servers from a deleted message', async () => { + const chatId = await createChat() + const deleted = userTurn([{ kind: 'mcp', serverId: 'server-deleted', label: 'Gone' }]) + await appendCopilotChatMessages(chatId, [ + deleted, + userTurn([{ kind: 'mcp', serverId: 'server-kept', label: 'Kept' }]), + ]) + await db + .update(copilotMessages) + .set({ deletedAt: new Date() }) + .where(and(eq(copilotMessages.chatId, chatId), eq(copilotMessages.messageId, deleted.id))) + + expect(await loadChatMcpServerIds(chatId)).toEqual(['server-kept']) + }) + + it('never reads servers tagged in another chat', async () => { + const chatId = await createChat() + const otherChatId = await createChat() + await appendCopilotChatMessages(otherChatId, [ + userTurn([{ kind: 'mcp', serverId: 'server-other', label: 'Other' }]), + ]) + + expect(await loadChatMcpServerIds(chatId)).toEqual([]) + }) +}) diff --git a/apps/sim/lib/mothership/chat/lifecycle.test.ts b/apps/sim/lib/mothership/chat/lifecycle.test.ts index 5972113df5c..795d18b6d50 100644 --- a/apps/sim/lib/mothership/chat/lifecycle.test.ts +++ b/apps/sim/lib/mothership/chat/lifecycle.test.ts @@ -189,19 +189,8 @@ describe('lifecycle copilot chat reads (cutover to copilot_messages)', () => { }) }) - it('resolveOrCreateChat returns conversationHistory from the table for an existing chat', async () => { - dbChainMockFns.limit.mockResolvedValueOnce([chatRow]) - dbChainMockFns.orderBy.mockResolvedValueOnce([{ content: userMsg }, { content: asstMsg }]) - - const result = await resolveOrCreateChat({ chatId: CHAT_ID, userId: USER_ID, model: 'm' }) - - expect(result.isNew).toBe(false) - expect(result.conversationHistory).toEqual([userMsg, asstMsg]) - }) - it('resolveOrCreateChat refuses a resumed chat whose type is not the asserted one', async () => { dbChainMockFns.limit.mockResolvedValueOnce([{ ...chatRow, type: 'mothership' }]) - dbChainMockFns.orderBy.mockResolvedValueOnce([]) const result = await resolveOrCreateChat({ chatId: CHAT_ID, @@ -212,13 +201,11 @@ describe('lifecycle copilot chat reads (cutover to copilot_messages)', () => { // Same shape an unknown id resolves to: the refusal carries no reason. expect(result.chat).toBeNull() - expect(result.conversationHistory).toEqual([]) expect(result.isNew).toBe(false) }) it('resolveOrCreateChat resumes a chat whose type matches the asserted one', async () => { dbChainMockFns.limit.mockResolvedValueOnce([{ ...chatRow, type: 'mothership' }]) - dbChainMockFns.orderBy.mockResolvedValueOnce([{ content: userMsg }]) const result = await resolveOrCreateChat({ chatId: CHAT_ID, @@ -228,7 +215,6 @@ describe('lifecycle copilot chat reads (cutover to copilot_messages)', () => { }) expect(result.chat).not.toBeNull() - expect(result.conversationHistory).toEqual([userMsg]) }) }) @@ -278,7 +264,7 @@ describe('organization chat isolation', () => { }) it.each(['agent', 'assistant'] as const)( - 'retains the same transcript and resources when requesting %s for the next turn', + 'retains the same chat and resources when requesting %s for the next turn', async (mode) => { const resources = [{ type: 'table', id: 'table', title: 'Evidence' }] dbChainMockFns.limit.mockResolvedValueOnce([ @@ -290,7 +276,6 @@ describe('organization chat isolation', () => { resources, }, ]) - dbChainMockFns.orderBy.mockResolvedValueOnce([{ content: userMsg }, { content: asstMsg }]) const result = await resolveOrCreateChat({ chatId: CHAT_ID, userId: USER_ID, @@ -303,7 +288,6 @@ describe('organization chat isolation', () => { expect(result.chatId).toBe(CHAT_ID) expect(result.isNew).toBe(false) expect(result.chat?.resources).toEqual(resources) - expect(result.conversationHistory).toEqual([userMsg, asstMsg]) expect(mockAuthorizeOrganization).toHaveBeenCalledWith({ principal: orgPrincipal, input: { organizationId: 'org-1', mode }, @@ -316,7 +300,6 @@ describe('organization chat isolation', () => { dbChainMockFns.limit.mockResolvedValueOnce([ { ...chatRow, organizationId: 'org-2', type: 'mothership' }, ]) - dbChainMockFns.orderBy.mockResolvedValueOnce([]) const result = await resolveOrCreateChat({ chatId: CHAT_ID, userId: USER_ID, @@ -326,7 +309,6 @@ describe('organization chat isolation', () => { type: 'mothership', }) expect(result.chat).toBeNull() - expect(result.conversationHistory).toEqual([]) }) it('rejects mixed organization and workspace ownership before creating a chat', async () => { diff --git a/apps/sim/lib/mothership/chat/lifecycle.ts b/apps/sim/lib/mothership/chat/lifecycle.ts index 1ac0054c1d1..d58546c5b10 100644 --- a/apps/sim/lib/mothership/chat/lifecycle.ts +++ b/apps/sim/lib/mothership/chat/lifecycle.ts @@ -26,8 +26,7 @@ const logger = createLogger('CopilotChatLifecycle') export interface ChatLoadResult { chatId: string - chat: CopilotChatDetailRow | null - conversationHistory: unknown[] + chat: CopilotChatDetail | null isNew: boolean } @@ -98,6 +97,42 @@ export async function loadCopilotChatMessages(chatId: string): Promise stripToolResultOutput(row.content as PersistedMessage)) } +/** + * MCP server ids tagged (`/name`) on a chat's live messages, in first-tagged + * transcript order. Reads only the `mcp` entries of each message's `contexts` + * instead of materializing the whole transcript, and skips soft-deleted + * messages so a removed turn no longer enables its servers. + */ +export async function loadChatMcpServerIds(chatId: string): Promise { + const rows = await db + .select({ serverId: sql`mcp_context.value ->> 'serverId'` }) + .from(copilotMessages) + .crossJoinLateral( + sql`jsonb_array_elements( + case when jsonb_typeof(${copilotMessages.content} -> 'contexts') = 'array' + then ${copilotMessages.content} -> 'contexts' + else '[]'::jsonb + end + ) with ordinality as mcp_context(value, ordinal)` + ) + .where( + and( + eq(copilotMessages.chatId, chatId), + isNull(copilotMessages.deletedAt), + sql`mcp_context.value ->> 'kind' = 'mcp'`, + sql`jsonb_typeof(mcp_context.value -> 'serverId') = 'string'`, + sql`mcp_context.value ->> 'serverId' <> ''` + ) + ) + .orderBy( + sql`${copilotMessages.seq} asc nulls last`, + asc(copilotMessages.createdAt), + asc(copilotMessages.id), + sql`mcp_context.ordinal` + ) + return [...new Set(rows.map((row) => row.serverId))] +} + /** * Ownership + liveness predicate shared by the accessible-chat loaders: * the chat must belong to the user and not be soft-deleted. @@ -115,7 +150,7 @@ type CopilotChatAuthRow = Pick< 'id' | 'userId' | 'workflowId' | 'workspaceId' | 'organizationId' | 'type' > & { mode: ConversationMode } -export type CopilotChatDetailRow = Pick< +export type CopilotChatDetail = Pick< typeof copilotChats.$inferSelect, | 'id' | 'userId' @@ -128,8 +163,9 @@ export type CopilotChatDetailRow = Pick< | 'resources' | 'createdAt' | 'updatedAt' -> & { - mode: ConversationMode +> & { mode: ConversationMode } + +export type CopilotChatDetailRow = CopilotChatDetail & { /** Transcript assembled from `copilot_messages` (no longer a chat-row column). */ messages: unknown[] } @@ -271,35 +307,42 @@ export async function getAccessibleCopilotChat( } /** - * Load a copilot chat with the conversation transcript and resources after - * authorization, omitting copilot-only TOAST-able fields (`previewYaml`, - * `config`) and unused metadata (`model`, `pinned`, `lastSeenAt`). Use this for the mothership chat detail endpoint and the - * shared `resolveOrCreateChat` path — every column read here is consumed - * downstream, and dropping the others avoids per-request detoast overhead. + * Load a copilot chat's detail columns after authorization, without its + * transcript. The transcript is unbounded — no per-chat message cap on write + * and no pruning — so callers that only need the chat's scope and metadata + * (such as `resolveOrCreateChat`) must not pay to materialize it. */ -export async function getAccessibleCopilotChatWithMessages( +async function getAccessibleCopilotChatDetail( chatId: string, userId: string, - options?: { includeTranscript?: boolean; principal?: Principal } -): Promise { + principal?: Principal +): Promise { const [chat] = await db .select(copilotChatDetailColumns) .from(copilotChats) .where(ownedLiveChatWhere(chatId, userId)) .limit(1) - const authorized = await authorizeCopilotChatRow(chat, chatId, userId, options?.principal) - if (!authorized) return null + return authorizeCopilotChatRow(chat, chatId, userId, principal) +} - /** - * The transcript is unbounded — no per-chat message cap on write and no - * pruning — so a caller that only needs the chat's scope should not pay to - * materialize it. Every check `resolveOrCreateChat` runs reads detail - * columns only, so an empty list stays a truthful "not loaded" rather than - * "no messages" for the callers that opt out. - */ - const messages = options?.includeTranscript === false ? [] : await loadCopilotChatMessages(chatId) - return { ...authorized, messages } +/** + * Load a copilot chat with the conversation transcript and resources after + * authorization, omitting copilot-only TOAST-able fields (`previewYaml`, + * `config`) and unused metadata (`model`, `pinned`, `lastSeenAt`). Use this for + * the mothership chat detail endpoint — every column read here is consumed + * downstream, and dropping the others avoids per-request detoast overhead. + */ +export async function getAccessibleCopilotChatWithMessages( + chatId: string, + userId: string, + options?: { principal?: Principal } +): Promise { + const chat = await getAccessibleCopilotChatDetail(chatId, userId, options?.principal) + if (!chat) return null + + const messages = await loadCopilotChatMessages(chatId) + return { ...chat, messages } } /** @@ -323,11 +366,6 @@ export async function resolveOrCreateChat(params: { model: string type?: 'mothership' | 'copilot' title?: string - /** - * Skips loading the transcript on the resume path. For a caller that keys - * continuity by `chatId` alone and never reads `conversationHistory`. - */ - includeTranscript?: boolean }): Promise { const { chatId, @@ -340,7 +378,6 @@ export async function resolveOrCreateChat(params: { mode, type, title, - includeTranscript, } = params if (organizationId) { @@ -358,14 +395,11 @@ export async function resolveOrCreateChat(params: { } if (chatId) { - const chat = await getAccessibleCopilotChatWithMessages(chatId, userId, { - includeTranscript, - principal, - }) + const chat = await getAccessibleCopilotChatDetail(chatId, userId, principal) if (chat) { if ((organizationId ?? null) !== (chat.organizationId ?? null)) { - return { chatId, chat: null, conversationHistory: [], isNew: false } + return { chatId, chat: null, isNew: false } } if (workflowId && chat.workflowId !== workflowId) { logger.warn('Copilot chat workflow mismatch', { @@ -374,7 +408,7 @@ export async function resolveOrCreateChat(params: { requestWorkflowId: workflowId, chatWorkflowId: chat.workflowId, }) - return { chatId, chat: null, conversationHistory: [], isNew: false } + return { chatId, chat: null, isNew: false } } if (workspaceId && chat.workspaceId !== workspaceId) { @@ -384,7 +418,7 @@ export async function resolveOrCreateChat(params: { requestWorkspaceId: workspaceId, chatWorkspaceId: chat.workspaceId, }) - return { chatId, chat: null, conversationHistory: [], isNew: false } + return { chatId, chat: null, isNew: false } } if (type && chat.type !== type) { @@ -394,7 +428,7 @@ export async function resolveOrCreateChat(params: { requestType: type, chatType: chat.type, }) - return { chatId, chat: null, conversationHistory: [], isNew: false } + return { chatId, chat: null, isNew: false } } if (chat.workflowId) { @@ -405,17 +439,12 @@ export async function resolveOrCreateChat(params: { userId, workflowId: chat.workflowId, }) - return { chatId, chat: null, conversationHistory: [], isNew: false } + return { chatId, chat: null, isNew: false } } } } - return { - chatId, - chat: chat ?? null, - conversationHistory: chat && Array.isArray(chat.messages) ? chat.messages : [], - isNew: false, - } + return { chatId, chat, isNew: false } } const now = new Date() @@ -436,18 +465,8 @@ export async function resolveOrCreateChat(params: { if (!newChat) { logger.warn('Failed to create new copilot chat row', { userId, workflowId, workspaceId }) - return { - chatId: '', - chat: null, - conversationHistory: [], - isNew: true, - } + return { chatId: '', chat: null, isNew: true } } - return { - chatId: newChat.id, - chat: { ...newChat, messages: [] }, - conversationHistory: [], - isNew: true, - } + return { chatId: newChat.id, chat: newChat, isNew: true } } diff --git a/apps/sim/lib/mothership/chat/post.test.ts b/apps/sim/lib/mothership/chat/post.test.ts index 14d479c2282..9d5519a274f 100644 --- a/apps/sim/lib/mothership/chat/post.test.ts +++ b/apps/sim/lib/mothership/chat/post.test.ts @@ -220,7 +220,10 @@ import { DEFAULT_PERMISSION_GROUP_CONFIG } from '@/lib/permission-groups/fields' import { handleUnifiedChatPost } from './post' const { mockBuildCopilotRequestPayload: buildCopilotRequestPayload } = mothershipChatPayloadMockFns -const { mockResolveOrCreateChat: resolveOrCreateChat } = mothershipChatLifecycleMockFns +const { + mockResolveOrCreateChat: resolveOrCreateChat, + mockLoadChatMcpServerIds: loadChatMcpServerIds, +} = mothershipChatLifecycleMockFns const { mockAuthorizeOrganizationChat: authorizeOrganizationChat } = mothershipOrganizationChatsMockFns const { mockAppendCopilotChatMessages: appendCopilotChatMessages } = mothershipChatMessagesMockFns @@ -334,9 +337,9 @@ describe('handleUnifiedChatPost', () => { resolveOrCreateChat.mockResolvedValue({ chatId: 'chat-1', chat: { id: 'chat-1' }, - conversationHistory: [], isNew: true, }) + loadChatMcpServerIds.mockResolvedValue([]) finalizeAssistantTurn.mockResolvedValue({ found: true, updated: true, @@ -770,7 +773,6 @@ describe('handleUnifiedChatPost', () => { chatId: 'chat-1', chat: { id: 'chat-1' }, isNew: false, - conversationHistory: [{ role: 'user', content: 'Previous turn', requestMode: previousMode }], }) const response = await handleUnifiedChatPost( new NextRequest('http://localhost/api/mothership/chat', { @@ -1245,17 +1247,9 @@ describe('handleUnifiedChatPost', () => { resolveOrCreateChat.mockResolvedValue({ chatId: 'chat-1', chat: { id: 'chat-1' }, - conversationHistory: [ - { - id: 'msg-1', - role: 'user', - content: '/Docs search auth', - contexts: [{ kind: 'mcp', serverId: 'mcp-server-1', label: 'Docs' }], - }, - { id: 'msg-2', role: 'assistant', content: 'here you go' }, - ], isNew: false, }) + loadChatMcpServerIds.mockResolvedValue(['mcp-server-1']) const response = await handleUnifiedChatPost( new NextRequest('http://localhost/api/copilot/chat', { @@ -1283,16 +1277,9 @@ describe('handleUnifiedChatPost', () => { resolveOrCreateChat.mockResolvedValue({ chatId: 'chat-1', chat: { id: 'chat-1' }, - conversationHistory: [ - { - id: 'msg-1', - role: 'user', - content: '/Docs search auth', - contexts: [{ kind: 'mcp', serverId: 'mcp-server-1', label: 'Docs' }], - }, - ], isNew: false, }) + loadChatMcpServerIds.mockResolvedValue(['mcp-server-1']) const response = await handleUnifiedChatPost( new NextRequest('http://localhost/api/copilot/chat', { @@ -1771,7 +1758,6 @@ describe('handleUnifiedChatPost copilot.use capability gate', () => { resolveOrCreateChat.mockResolvedValue({ chatId: 'chat-1', chat: { id: 'chat-1' }, - conversationHistory: [], isNew: true, }) }) diff --git a/apps/sim/lib/mothership/chat/post.ts b/apps/sim/lib/mothership/chat/post.ts index 9822e89d962..f7d22257fdb 100644 --- a/apps/sim/lib/mothership/chat/post.ts +++ b/apps/sim/lib/mothership/chat/post.ts @@ -41,7 +41,11 @@ import { DESKTOP_TERMINAL_HINT_ID_MAX_LENGTH, DESKTOP_TERMINAL_HINT_TEXT_MAX_LENGTH, } from '@/lib/mothership/chat/desktop-capabilities' -import { type ChatLoadResult, resolveOrCreateChat } from '@/lib/mothership/chat/lifecycle' +import { + type ChatLoadResult, + loadChatMcpServerIds, + resolveOrCreateChat, +} from '@/lib/mothership/chat/lifecycle' import { authorizeOrganizationChat } from '@/lib/mothership/chat/organization-chats' import { buildCopilotRequestPayload } from '@/lib/mothership/chat/payload' import { @@ -489,27 +493,13 @@ function normalizeContexts(contexts: UnifiedChatRequest['contexts']) { * on a sent message showing only what the user actually typed that turn. */ function collectChatMcpServerIds( - conversationHistory: unknown[], + chatMcpServerIds: string[], currentContexts: UnifiedChatRequest['contexts'] ): string[] { - const serverIds = new Set() - - const collect = (contexts: unknown) => { - if (!Array.isArray(contexts)) return - for (const ctx of contexts) { - if (!ctx || typeof ctx !== 'object') continue - const { kind, serverId } = ctx as { kind?: unknown; serverId?: unknown } - if (kind === 'mcp' && typeof serverId === 'string' && serverId) { - serverIds.add(serverId) - } - } + const serverIds = new Set(chatMcpServerIds) + for (const ctx of currentContexts ?? []) { + if (ctx.kind === 'mcp' && ctx.serverId) serverIds.add(ctx.serverId) } - - for (const message of conversationHistory) { - collect((message as { contexts?: unknown } | null)?.contexts) - } - collect(currentContexts) - return Array.from(serverIds) } @@ -1153,7 +1143,6 @@ export async function handleUnifiedChatPost(req: NextRequest) { activeOtelRoot.setInputMessages({ userMessage: body.message }) let currentChat: ChatLoadResult['chat'] = null - let conversationHistory: unknown[] = [] let chatIsNew = false actualChatId = body.chatId @@ -1190,9 +1179,6 @@ export async function handleUnifiedChatPost(req: NextRequest) { currentChat = chatResult.chat actualChatId = chatResult.chatId || body.chatId chatIsNew = chatResult.isNew - conversationHistory = Array.isArray(chatResult.conversationHistory) - ? chatResult.conversationHistory - : [] if (body.chatId && !currentChat) { activeOtelRoot.span.setAttribute(TraceAttr.HttpStatusCode, 404) @@ -1317,13 +1303,21 @@ export async function handleUnifiedChatPost(req: NextRequest) { activeOtelRoot.context ) }) - const [agentContexts, userPermission, executionContext, personalCredentials] = - await Promise.all([ - agentContextsPromise, - userPermissionPromise, - executionContextPromise, - personalCredentialsPromise, - ]) + const chatMcpServerIdsPromise = + currentChat && !chatIsNew ? loadChatMcpServerIds(currentChat.id) : Promise.resolve([]) + const [ + agentContexts, + userPermission, + executionContext, + personalCredentials, + chatMcpServerIds, + ] = await Promise.all([ + agentContextsPromise, + userPermissionPromise, + executionContextPromise, + personalCredentialsPromise, + chatMcpServerIdsPromise, + ]) let workspaceContext: string | undefined if (personalCredentials) { workspaceContext = JSON.stringify({ @@ -1366,7 +1360,7 @@ export async function handleUnifiedChatPost(req: NextRequest) { [TraceAttr.CopilotContextsCount]: normalizedContexts.length, }, () => { - const mcpServerIds = collectChatMcpServerIds(conversationHistory, normalizedContexts) + const mcpServerIds = collectChatMcpServerIds(chatMcpServerIds, normalizedContexts) return branch.kind === 'workflow' ? branch.buildPayload({ message: body.message, diff --git a/apps/sim/lib/mothership/inbox/executor.test.ts b/apps/sim/lib/mothership/inbox/executor.test.ts index 0811194a75e..34484aca06e 100644 --- a/apps/sim/lib/mothership/inbox/executor.test.ts +++ b/apps/sim/lib/mothership/inbox/executor.test.ts @@ -154,7 +154,6 @@ describe('Inbox execution actor', () => { mockResolveOrCreateChat.mockResolvedValue({ chatId: 'chat-1', chat: { id: 'chat-1' }, - conversationHistory: [], isNew: true, }) dbChainMockFns.returning diff --git a/packages/testing/src/mocks/mothership-chat-lifecycle.mock.ts b/packages/testing/src/mocks/mothership-chat-lifecycle.mock.ts index 0572081eb65..e1bde38aa1d 100644 --- a/packages/testing/src/mocks/mothership-chat-lifecycle.mock.ts +++ b/packages/testing/src/mocks/mothership-chat-lifecycle.mock.ts @@ -13,6 +13,7 @@ import { vi } from 'vitest' */ export const mothershipChatLifecycleMockFns = { mockLoadCopilotChatMessages: vi.fn(), + mockLoadChatMcpServerIds: vi.fn(), mockGetAccessibleCopilotChatAuth: vi.fn(), mockGetAccessibleCopilotChatForCancellation: vi.fn(), mockGetAccessibleCopilotChat: vi.fn(), @@ -30,6 +31,7 @@ export const mothershipChatLifecycleMockFns = { */ export const mothershipChatLifecycleMock = { loadCopilotChatMessages: mothershipChatLifecycleMockFns.mockLoadCopilotChatMessages, + loadChatMcpServerIds: mothershipChatLifecycleMockFns.mockLoadChatMcpServerIds, getAccessibleCopilotChatAuth: mothershipChatLifecycleMockFns.mockGetAccessibleCopilotChatAuth, getAccessibleCopilotChatForCancellation: mothershipChatLifecycleMockFns.mockGetAccessibleCopilotChatForCancellation,