diff --git a/apps/sim/app/api/mothership/chats/[chatId]/fork/route.test.ts b/apps/sim/app/api/mothership/chats/[chatId]/fork/route.test.ts index c0568a6dae3..4e0a2ee3985 100644 --- a/apps/sim/app/api/mothership/chats/[chatId]/fork/route.test.ts +++ b/apps/sim/app/api/mothership/chats/[chatId]/fork/route.test.ts @@ -185,6 +185,7 @@ describe('POST /api/mothership/chats/[chatId]/fork', () => { queueTableRows(copilotChats, [chat]) queueTableRows(copilotChats, [chat]) queueTableRows(member, [{ role: 'member' }]) + queueTableRows(copilotChats, [{ id: chat.id }]) const attachments = [ { id: 'upload-1', @@ -215,6 +216,26 @@ describe('POST /api/mothership/chats/[chatId]/fork', () => { expect(mockAssertActiveWorkspaceAccess).not.toHaveBeenCalled() }) + it('refuses to publish a fork whose source chat was purged while it copied', async () => { + // The fork shares its attachment keys with the source, and cleanup deletes an unreferenced + // key only after deleting the source row: publishing without the source would leave the + // fork pointing at bytes that cleanup is about to delete. + dbChainMockFns.limit.mockReset() + const chat = { ...parentRow, workspaceId: null, organizationId: 'org-1', resources: [] } + queueTableRows(copilotChats, [chat]) + queueTableRows(copilotChats, [chat]) + queueTableRows(member, [{ role: 'member' }]) + queueTableRows(copilotChats, []) + const res = await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' })) + expect(res.status).toBe(404) + expect(mockAppendCopilotChatMessages).not.toHaveBeenCalled() + expect(mockPublishStatusChanged).not.toHaveBeenCalled() + expect(mockFetchGo.mock.calls.map(([url]) => url)).toEqual([ + 'http://mothership.test/api/chats/fork', + 'http://mothership.test/api/tasks/cleanup', + ]) + }) + it.each(['membership', 'capability'] as const)( 'denies organization forks after %s revocation before copying', async (revocation) => { @@ -316,6 +337,12 @@ describe('POST /api/mothership/chats/[chatId]/fork', () => { expect(mockPublishStatusChanged).not.toHaveBeenCalled() }) + it.each([404, 409, 413])('passes a worker %i refusal through', async (status) => { + mockFetchGo.mockResolvedValue(Response.json({ error: 'refused' }, { status })) + const res = await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' })) + expect(res.status).toBe(status) + }) + it('surfaces failed blob copies and excludes their metadata from publication', async () => { mockExecuteChatFileBlobCopies.mockResolvedValue({ copied: 1, @@ -336,6 +363,116 @@ describe('POST /api/mothership/chats/[chatId]/fork', () => { expect(dbChainMockFns.delete).not.toHaveBeenCalled() }) + it('discards the worker copy when the fork cannot be published', async () => { + dbChainMockFns.transaction.mockRejectedValueOnce(new Error('connection lost')) + const res = await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' })) + expect(res.status).toBe(500) + const [fork, discard] = mockFetchGo.mock.calls + expect(fork[0]).toBe('http://mothership.test/api/chats/fork') + expect(discard[0]).toBe('http://mothership.test/api/tasks/cleanup') + expect(JSON.parse(discard[1].body)).toEqual({ + chatIds: [JSON.parse(fork[1].body).newChatId], + }) + expect(mockPublishStatusChanged).not.toHaveBeenCalled() + }) + + it('keeps references on the source file when its copy fails', async () => { + const oldKey = 'workspace/ws-1/old-cat.png' + mockListForkableChatFiles.mockResolvedValue([ + { id: OLD_FILE_ID, key: oldKey, messageId: 'msg-1', workspaceId: 'ws-1' }, + ]) + mockPlanChatFileCopies.mockReturnValue({ + idMap: new Map([[OLD_FILE_ID, NEW_FILE_ID]]), + keyMap: new Map([[oldKey, 'workspace/ws-1/new-cat.png']]), + blobTasks: [ + { + copyId: NEW_FILE_ID, + sourceKey: oldKey, + targetKey: 'workspace/ws-1/new-cat.png', + context: 'mothership', + fileName: 'cat.png', + contentType: 'image/png', + }, + ], + }) + mockExecuteChatFileBlobCopies.mockResolvedValue({ + copied: 0, + failed: 1, + failedCopyIds: [NEW_FILE_ID], + }) + const res = await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' })) + expect(res.status).toBe(200) + // The copy is never published, so neither Sim's transcript nor the worker's may name it. + expect(mockAppendCopilotChatMessages.mock.calls[0][1][0].content).toBe( + `See ![cat](/api/files/view/${OLD_FILE_ID})` + ) + const forkRequest = JSON.parse(mockFetchGo.mock.calls[0][1].body) + expect(forkRequest.fileIds).toEqual({}) + expect(forkRequest.fileKeys).toEqual({}) + }) + + it('re-points in-app file links and tool-call arguments at the copied file', async () => { + mockLoadCopilotChatMessages.mockResolvedValue([ + { ...threeMessages[0], content: `Open /workspace/ws-1/files/${OLD_FILE_ID}` }, + { + ...threeMessages[1], + contentBlocks: [ + { + type: 'tool', + toolCall: { + id: 'call-1', + name: 'sim_cli', + state: 'success', + params: { fileId: OLD_FILE_ID, args: ['files', 'read', OLD_FILE_ID], n: 2 }, + display: { title: `Read /api/files/view/${OLD_FILE_ID}` }, + }, + }, + ], + }, + ]) + mockPlanChatFileCopies.mockReturnValue({ + idMap: new Map([[OLD_FILE_ID, NEW_FILE_ID]]), + keyMap: new Map(), + blobTasks: [], + }) + const res = await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' })) + expect(res.status).toBe(200) + const [user, assistant] = mockAppendCopilotChatMessages.mock.calls[0][1] + expect(user.content).toBe(`Open /workspace/ws-1/files/${NEW_FILE_ID}`) + expect(assistant.contentBlocks[0].toolCall).toMatchObject({ + params: { fileId: NEW_FILE_ID, args: ['files', 'read', NEW_FILE_ID], n: 2 }, + display: { title: `Read /api/files/view/${NEW_FILE_ID}` }, + }) + }) + + it('drops a Sources tab whose response is past the cut', async () => { + const kept = { type: 'sources', id: 'cited-sources', title: 'Sources' } + dbChainMockFns.limit.mockResolvedValue([ + { ...parentRow, resources: [{ ...kept, sources: { messageId: 'msg-3' } }] }, + ]) + await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' })) + expect(dbChainMockFns.values).toHaveBeenCalledWith(expect.objectContaining({ resources: [] })) + + dbChainMockFns.values.mockClear() + dbChainMockFns.limit.mockResolvedValue([ + { + ...parentRow, + resources: [{ ...kept, sources: { messageId: 'live-id', requestId: 'req-2' } }], + }, + ]) + mockLoadCopilotChatMessages.mockResolvedValue([ + threeMessages[0], + { ...threeMessages[1], requestId: 'req-2' }, + threeMessages[2], + ]) + await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' })) + expect(dbChainMockFns.values).toHaveBeenCalledWith( + expect.objectContaining({ + resources: [{ ...kept, sources: { messageId: 'live-id', requestId: 'req-2' } }], + }) + ) + }) + it('copies pre-cut uploads and drops only post-cut ghosts', async () => { // The source chat owns two more uploads (apple pre-cut, banana post-cut) // beside the kept one, plus one shared workspace-file resource. The fork diff --git a/apps/sim/app/api/mothership/chats/[chatId]/fork/route.ts b/apps/sim/app/api/mothership/chats/[chatId]/fork/route.ts index ffecce2ca36..312f27b26ce 100644 --- a/apps/sim/app/api/mothership/chats/[chatId]/fork/route.ts +++ b/apps/sim/app/api/mothership/chats/[chatId]/fork/route.ts @@ -2,7 +2,7 @@ import { createLogger } from '@sim/logger' import { type NextRequest, NextResponse } from 'next/server' import { forkMothershipChatContract } from '@/lib/api/contracts/mothership-chats' import { parseRequest } from '@/lib/api/server' -import { asOrchestrationError } from '@/lib/core/orchestration/types' +import { asOrchestrationError, statusForOrchestrationError } from '@/lib/core/orchestration/types' import { withRouteHandler } from '@/lib/core/utils/with-route-handler' import { forkChat } from '@/lib/mothership/chat/application/fork' import { @@ -38,6 +38,11 @@ export const POST = withRouteHandler( if (classified?.code === 'not_found' || classified?.code === 'forbidden') return NextResponse.json({ error: 'Chat not found' }, { status: 404 }) if (classified?.code === 'validation') return createBadRequestResponse(classified.message) + if (classified?.code === 'conflict' || classified?.code === 'payload_too_large') + return NextResponse.json( + { error: classified.message }, + { status: statusForOrchestrationError(classified.code) } + ) logger.error('Error forking chat:', error) return createInternalServerErrorResponse('Failed to fork chat') } diff --git a/apps/sim/app/workspace/[workspaceId]/components/message-actions/message-actions.tsx b/apps/sim/app/workspace/[workspaceId]/components/message-actions/message-actions.tsx index cf99a9d5a72..b82ab88f86d 100644 --- a/apps/sim/app/workspace/[workspaceId]/components/message-actions/message-actions.tsx +++ b/apps/sim/app/workspace/[workspaceId]/components/message-actions/message-actions.tsx @@ -19,6 +19,7 @@ import { useCopyToClipboard, } from '@sim/emcn' import { useParams, useRouter } from 'next/navigation' +import { isApiClientError } from '@/lib/api/client/errors' import { isLiveAssistantMessageId } from '@/lib/mothership/chat/live-message-id' import { organizationRoutes } from '@/lib/navigation/paths' import { useChatSurface } from '@/app/workspace/[workspaceId]/home/components/chat-surface-context' @@ -40,6 +41,9 @@ interface MessageActionsProps { messageId?: string } +/** Fork refusals whose message tells the person what to do: the response is unfinished, or the chat is too long. */ +const FORK_REFUSAL_STATUSES = new Set([409, 413]) + export const MessageActions = memo(function MessageActions({ content, getCopyContent, @@ -143,8 +147,12 @@ export const MessageActions = memo(function MessageActions({ useFolderStore.getState().clearChatSelection() router.push(`/workspace/${params.workspaceId}/chat/${result.id}`) } - } catch { - toast.error('Failed to fork chat') + } catch (error) { + toast.error( + isApiClientError(error) && FORK_REFUSAL_STATUSES.has(error.status) + ? error.message + : 'Failed to fork chat' + ) } } diff --git a/apps/sim/lib/cleanup/chat-cleanup.test.ts b/apps/sim/lib/cleanup/chat-cleanup.test.ts new file mode 100644 index 00000000000..ef054b4696b --- /dev/null +++ b/apps/sim/lib/cleanup/chat-cleanup.test.ts @@ -0,0 +1,101 @@ +import { copilotChats, copilotMessages, workspaceFiles } from '@sim/db/schema' +import { queueTableRows, resetDbChainMock } from '@sim/testing/mocks/database.mock' +import { storageServiceMockFns } from '@sim/testing/mocks/storage-service.mock' +import { uploadsMock, uploadsMockFns } from '@sim/testing/mocks/uploads.mock' +import { beforeEach, describe, expect, it, vi } from 'vitest' + +vi.mock('@/lib/uploads', () => uploadsMock) + +import { prepareChatCleanup } from '@/lib/cleanup/chat-cleanup' +import { inlineChatImageKey } from '@/lib/mothership/chat/inline-image-key' + +const chatId = '3f0c2a52-8a43-4d4b-9b5f-0b3c7a1e2d10' + +function attachment(key: string) { + return { id: 'wf_x', key, filename: 'x.png', media_type: 'image/png', size: 1 } +} + +/** + * Purges one deleted chat whose message rows are `messages`, while `remaining` are the message + * rows other chats of the organization still hold; returns the deleted keys by context. + */ +async function purge( + messages: Record[], + remaining: Record[] = [] +) { + queueTableRows(workspaceFiles, []) + queueTableRows( + copilotMessages, + messages.map((content) => ({ chatId, content })) + ) + const cleanup = await prepareChatCleanup([chatId], 'test') + queueTableRows(copilotChats, []) + queueTableRows( + copilotMessages, + remaining.map((content) => ({ content })) + ) + await cleanup.execute() + return Object.fromEntries( + storageServiceMockFns.mockDeleteFiles.mock.calls.map(([keys, context]) => [context, keys]) + ) +} + +describe('chat purge storage', () => { + beforeEach(() => { + resetDbChainMock() + uploadsMockFns.mockIsUsingCloudStorage.mockReturnValue(true) + storageServiceMockFns.mockDeleteFiles.mockReset() + storageServiceMockFns.mockDeleteFiles.mockResolvedValue({ deleted: 1, failed: [] }) + }) + + it('never deletes a workspace attachment as copilot storage', async () => { + // A fork carries its source's key when the file's copy failed or the file was deleted, and + // the copilot bucket can be the workspace bucket, so a workspace key here is another chat's file. + const deleted = await purge([ + { + role: 'user', + content: 'look', + fileAttachments: [ + attachment('workspace/ws-1/1700-abc-shared.png'), + attachment('copilot/1234/legacy.png'), + ], + }, + ]) + expect(deleted).toEqual({ copilot: ['copilot/1234/legacy.png'] }) + }) + + it('deletes an organization attachment only once no remaining chat references it', async () => { + const shared = 'assistant/org-1/user-1/u1/shared.png' + const own = 'assistant/org-1/user-1/u2/own.png' + const message = { + role: 'user', + content: 'look', + fileAttachments: [attachment(shared), attachment(own)], + } + // A fork of this chat still holds the shared upload. + const deleted = await purge( + [message], + [{ role: 'user', content: 'fork', fileAttachments: [attachment(shared)] }] + ) + expect(deleted).toEqual({ mothership: [own] }) + }) + + it('deletes the inline images an assistant message published', async () => { + const deleted = await purge([ + { + role: 'assistant', + requestId: 'req-1', + content: 'Here: ![chart](files/chart.png) and `![code](files/no.png)`', + contentBlocks: [{ type: 'text', content: '![other](/tmp/out.png)' }], + }, + { role: 'user', requestId: 'req-2', content: '![u](files/u.png)' }, + { role: 'assistant', content: '![x](files/x.png)' }, + ]) + expect(deleted).toEqual({ + mothership: [ + inlineChatImageKey(chatId, 'req-1', 'files/chart.png'), + inlineChatImageKey(chatId, 'req-1', '/tmp/out.png'), + ], + }) + }) +}) diff --git a/apps/sim/lib/cleanup/chat-cleanup.ts b/apps/sim/lib/cleanup/chat-cleanup.ts index b3910219bb3..ab9745dd093 100644 --- a/apps/sim/lib/cleanup/chat-cleanup.ts +++ b/apps/sim/lib/cleanup/chat-cleanup.ts @@ -2,11 +2,17 @@ import { dbFor } from '@sim/db' import { copilotChats, copilotMessages, workspaceFiles } from '@sim/db/schema' import { createLogger } from '@sim/logger' import { chunkArray } from '@sim/utils/helpers' -import { and, inArray, isNull } from 'drizzle-orm' +import { isRecordLike } from '@sim/utils/object' +import { and, eq, inArray, isNull, or, sql } from 'drizzle-orm' import { env } from '@/lib/core/config/env' +import { + inlineChatImageKey, + inlineChatImageReferences, +} from '@/lib/mothership/chat/inline-image-key' import { SIM_AGENT_API_URL } from '@/lib/mothership/constants' import type { StorageContext } from '@/lib/uploads' import { isUsingCloudStorage, StorageService } from '@/lib/uploads' +import { tryInferContextFromKey } from '@/lib/uploads/utils/file-utils' const logger = createLogger('ChatCleanup') @@ -30,12 +36,63 @@ interface FileRef { key: string context: ChatScopedContext chatId: string + /** + * An organization Chat attachment (`assistant//…`): no `workspace_files` row owns it, + * and a fork carries the same key, so it is deleted only once no remaining chat of that + * organization references it. + */ + organizationId?: string +} + +/** The organization an `assistant//…` attachment key was uploaded under. */ +function organizationAttachmentOwner(key: string): string | undefined { + if (tryInferContextFromKey(key) !== 'mothership') return undefined + const [, organizationId] = key.split('/') + return organizationId || undefined +} + +/** + * The chat images an assistant message published under its request id, keyed by chat id, + * so they are purged with the chat. A row this cannot read as such a message has none. + */ +function inlineChatImageKeys(chatId: string, content: Record): string[] { + if (content.role !== 'assistant' || typeof content.requestId !== 'string') return [] + const published = inlineChatImageReferences({ + role: 'assistant', + requestId: content.requestId, + content: typeof content.content === 'string' ? content.content : '', + contentBlocks: Array.isArray(content.contentBlocks) + ? content.contentBlocks.flatMap((block) => + isRecordLike(block) && block.type === 'text' && typeof block.content === 'string' + ? [{ type: 'text' as const, content: block.content }] + : [] + ) + : undefined, + }) + if (!published) return [] + const keys: string[] = [] + for (const reference of published.references) { + try { + keys.push(inlineChatImageKey(chatId, published.requestId, reference)) + } catch { + // A chat or request id outside the key grammar never had an image stored under it. + } + } + return keys } /** - * Collect all file storage keys for the given chat IDs from two sources: + * Collect all file storage keys for the given chat IDs from three sources: * 1. workspaceFiles rows with chatId FK (chat-scoped contexts only) - * 2. fileAttachments[].key inside each copilot_messages.content + * 2. fileAttachments[].key inside each copilot_messages.content: copilot keys, and + * organization attachments under their own context (see {@link FileRef.organizationId}) + * 3. the chat-scoped inline images each assistant message published + * + * A workspace attachment is owned by its `workspace_files` row (source 1, or the workspace + * file lifecycle), never by the message: a fork carries the same attachment when its copy + * failed or its file was deleted, and the copilot bucket falls back to the workspace bucket + * on GCS (and may be configured to it on S3), so deleting such a key as copilot storage + * could delete a file another chat or the workspace still uses. */ export async function collectChatFiles(chatIds: string[]): Promise { const files: FileRef[] = [] @@ -77,6 +134,12 @@ export async function collectChatFiles(chatIds: string[]): Promise { for (const row of messageRows) { const msg = row.content if (!msg || typeof msg !== 'object') continue + for (const key of inlineChatImageKeys(row.chatId, msg as Record)) { + if (!seen.has(key)) { + seen.add(key) + files.push({ key, context: 'mothership', chatId: row.chatId }) + } + } const attachments = (msg as Record).fileAttachments if (!Array.isArray(attachments)) continue for (const attachment of attachments) { @@ -86,7 +149,12 @@ export async function collectChatFiles(chatIds: string[]): Promise { (attachment as Record).key ) { const key = (attachment as Record).key as string - if (!seen.has(key)) { + if (seen.has(key)) continue + const organizationId = organizationAttachmentOwner(key) + if (organizationId) { + seen.add(key) + files.push({ key, context: 'mothership', chatId: row.chatId, organizationId }) + } else if (tryInferContextFromKey(key) === 'copilot') { seen.add(key) files.push({ key, context: 'copilot', chatId: row.chatId }) } @@ -98,6 +166,55 @@ export async function collectChatFiles(chatIds: string[]): Promise { return files } +/** + * Organization attachment keys that a remaining chat still references, so they outlive the + * chats being purged. Runs after the caller deleted those chats' rows, so any match is + * another chat: a fork, or a chat that is only soft-deleted and may be restored. + */ +async function organizationAttachmentsStillReferenced(files: FileRef[]): Promise> { + const keysByOrganization = new Map() + for (const file of files) { + if (!file.organizationId) continue + const keys = keysByOrganization.get(file.organizationId) + if (keys) keys.push(file.key) + else keysByOrganization.set(file.organizationId, [file.key]) + } + const referenced = new Set() + for (const [organizationId, keys] of keysByOrganization) { + for (const chunk of chunkArray(keys, CHAT_FILE_COLLECT_CHUNK_SIZE)) { + const wanted = new Set(chunk) + const rows = await cleanupDb + .select({ content: copilotMessages.content }) + .from(copilotMessages) + .innerJoin(copilotChats, eq(copilotChats.id, copilotMessages.chatId)) + .where( + and( + eq(copilotChats.organizationId, organizationId), + or( + ...chunk.map( + (key) => + sql`${copilotMessages.content} @> ${JSON.stringify({ fileAttachments: [{ key }] })}::jsonb` + ) + ) + ) + ) + for (const { content } of rows) { + const attachments = isRecordLike(content) ? content.fileAttachments : undefined + if (!Array.isArray(attachments)) continue + for (const attachment of attachments) { + if ( + isRecordLike(attachment) && + typeof attachment.key === 'string' && + wanted.has(attachment.key) + ) + referenced.add(attachment.key) + } + } + } + } + return referenced +} + /** Groups files by storage context so each context can use one batch DELETE call. */ export async function deleteStorageFiles( files: FileRef[], @@ -220,7 +337,12 @@ export async function prepareChatCleanup( ) } const confirmedChatIds = chatIds.filter((id) => !survivors.has(id)) - const confirmedFiles = files.filter((file) => !survivors.has(file.chatId)) + const sharedKeys = await organizationAttachmentsStillReferenced( + files.filter((file) => !survivors.has(file.chatId)) + ) + const confirmedFiles = files.filter( + (file) => !survivors.has(file.chatId) && !sharedKeys.has(file.key) + ) // Call copilot backend if (confirmedChatIds.length > 0) { diff --git a/apps/sim/lib/core/errors/retryable-infrastructure.ts b/apps/sim/lib/core/errors/retryable-infrastructure.ts index 73051cb11c0..3890d19eb8c 100644 --- a/apps/sim/lib/core/errors/retryable-infrastructure.ts +++ b/apps/sim/lib/core/errors/retryable-infrastructure.ts @@ -23,6 +23,7 @@ const RETRYABLE_NETWORK_ERROR_CODES = new Set([ 'ECONNREFUSED', 'EPIPE', 'ERR_STREAM_PREMATURE_CLOSE', + 'EHOSTUNREACH', 'ENETDOWN', 'ENETRESET', 'ENETUNREACH', diff --git a/apps/sim/lib/mothership/chat/application/fork.ts b/apps/sim/lib/mothership/chat/application/fork.ts index d3c2ebcec3b..97b15c298ce 100644 --- a/apps/sim/lib/mothership/chat/application/fork.ts +++ b/apps/sim/lib/mothership/chat/application/fork.ts @@ -19,10 +19,14 @@ import { planChatFileCopies, } from '@/lib/mothership/chat/fork-chat-files' import { planForkInlineImages } from '@/lib/mothership/chat/fork-inline-images' -import { copyWorkerConversation } from '@/lib/mothership/chat/fork-worker' +import { + copyWorkerConversation, + discardWorkerConversation, +} from '@/lib/mothership/chat/fork-worker' import { loadCopilotChatMessages } from '@/lib/mothership/chat/lifecycle' import { appendCopilotChatMessages } from '@/lib/mothership/chat/messages-store' import { + publishedFileRefMaps, rewriteMessageFileRefs, rewriteResourceFileRefs, } from '@/lib/mothership/chat/rewrite-file-references' @@ -88,6 +92,8 @@ export const forkChat = defineAuthorizedChatUseCase({ const { chatId, upToMessageId } = input let preparedBlobs: ChatBlobCopyTask[] = [] let published = false + const newId = generateId() + let workerCopyRequested = false try { const messages = await loadCopilotChatMessages(chatId) const forkIdx = messages.findIndex((m) => m.id === upToMessageId) @@ -95,23 +101,28 @@ export const forkChat = defineAuthorizedChatUseCase({ throw new OrchestrationError('validation', 'Message not found in chat') } const forkedMessages = messages.slice(0, forkIdx + 1) + const keptMessageIds = new Set(forkedMessages.map((m) => m.id)) + const keptRequestIds = new Set( + forkedMessages.flatMap((m) => (m.requestId ? [m.requestId] : [])) + ) + /** The Sources panel reads its response by message or request id; a response past the cut is not in the fork. */ + const addressesKeptMessage = ({ sources }: MothershipResource) => + !!sources && + (keptMessageIds.has(sources.messageId) || + (!!sources.requestId && keptRequestIds.has(sources.requestId))) /** Single workspace_files read per fork: every chat-owned upload. The copied set is timeline-cut to the kept message slice in memory (files born after the fork point stay behind). */ const chatOwnedFiles = context.workspaceId ? await listForkableChatFiles(db, chatId) : [] - const sourceFiles = filterForkableChatFiles( - chatOwnedFiles, - new Set(forkedMessages.map((m) => m.id)) - ) + const sourceFiles = filterForkableChatFiles(chatOwnedFiles, keptMessageIds) /** Resources are stored as a jsonb array on the chat row. They carry no timestamps, so they can't be timeline-cut like messages — instead, file resources whose chat-owned file is NOT copied (uploads born after the cut) are dropped in the rewrite below; everything else is copied. */ const parentResources = sanitizeChatResources( Array.isArray(parent.resources) ? (parent.resources as MothershipResource[]) : [] - ) + ).filter((resource) => resource.type !== 'sources' || addressesKeptMessage(resource)) /** The source chat's chat-owned file ids (no cut) — the "is this resource a ghost?" test set for the rewrite. */ const chatOwnedFileIds = new Set(chatOwnedFiles.map((row) => row.id)) - const newId = generateId() /** Strip a leading "Fork | " so titles don't stack prefixes when forking a forked chat. */ const baseTitle = (parent.title ?? 'New chat').replace(/^Fork \| /, '') const title = `Fork | ${baseTitle}` @@ -121,14 +132,15 @@ export const forkChat = defineAuthorizedChatUseCase({ preparedBlobs = [...plan.blobTasks, ...planForkInlineImages(forkedMessages, chatId, newId)] const { failed, failedCopyIds } = await executeChatFileBlobCopies(preparedBlobs) const failedIds = new Set(failedCopyIds) - const maps = { fileIds: plan.idMap, fileKeys: plan.keyMap } - const newChatResources = rewriteResourceFileRefs( - parentResources, - maps, - chatOwnedFileIds - ).filter((resource) => resource.type !== 'file' || !failedIds.has(resource.id)) + const maps = { + ...publishedFileRefMaps(plan, failedIds), + workspaceId: context.workspaceId, + } + /** A chat-owned file whose copy failed is a ghost here too: no published copy stands in for it. */ + const newChatResources = rewriteResourceFileRefs(parentResources, maps, chatOwnedFileIds) const cutUser = [...forkedMessages].reverse().find((message) => message.role === 'user') if (!cutUser) throw new Error('The fork has no user message') + workerCopyRequested = true await copyWorkerConversation({ sourceChatId: chatId, newChatId: newId, @@ -138,12 +150,26 @@ export const forkChat = defineAuthorizedChatUseCase({ userId, upToMessageId: cutUser.id, includeResponse: forkedMessages.at(-1)?.role === 'assistant', - fileIds: Object.fromEntries(plan.idMap), - fileKeys: Object.fromEntries(plan.keyMap), + fileIds: Object.fromEntries(maps.fileIds), + fileKeys: Object.fromEntries(maps.fileKeys), }) /** Publish only after both the file bytes and the worker conversation are prepared. */ await db.transaction(async (tx) => { + /** + * The fork can share keys with its source (organization attachments, files whose copy + * failed), and chat cleanup deletes a shared key once no remaining chat references it, + * checking only after it deletes the source row. Holding the source row until commit + * orders the two: a purge that already removed it refuses this fork, and one that has + * not waits for this commit and then sees the fork's references. + */ + const [source] = await tx + .select({ id: copilotChats.id }) + .from(copilotChats) + .where(eq(copilotChats.id, chatId)) + .for('key share') + .limit(1) + if (!source) throw new OrchestrationError('not_found', 'Chat not found') const [row] = await tx .insert(copilotChats) .values({ @@ -198,6 +224,8 @@ export const forkChat = defineAuthorizedChatUseCase({ ...(failed > 0 ? { failedFileCopies: failed } : {}), } } catch (error) { + if (!published && workerCopyRequested) + await discardWorkerConversation({ newChatId: newId, userId }) if (!published) { await mapWithConcurrency(preparedBlobs, 4, async (task) => { try { diff --git a/apps/sim/lib/mothership/chat/application/inline-images.test.ts b/apps/sim/lib/mothership/chat/application/inline-images.test.ts index d34295c72d4..b004fee8a2d 100644 --- a/apps/sim/lib/mothership/chat/application/inline-images.test.ts +++ b/apps/sim/lib/mothership/chat/application/inline-images.test.ts @@ -46,8 +46,8 @@ import { materializeStreamImage, readInlineChatImage, } from '@/lib/mothership/chat/application/inline-images' +import { inlineChatImageKey } from '@/lib/mothership/chat/inline-image-key' import { inlineChatImageUrl } from '@/lib/mothership/chat/inline-image-reference' -import { inlineChatImageKey } from '@/lib/mothership/chat/inline-image-storage' import { GET } from '@/app/api/mothership/chats/[chatId]/images/[requestId]/route' const mocks = { diff --git a/apps/sim/lib/mothership/chat/fork-inline-images.test.ts b/apps/sim/lib/mothership/chat/fork-inline-images.test.ts index 63ad5be033a..9131a08f7cb 100644 --- a/apps/sim/lib/mothership/chat/fork-inline-images.test.ts +++ b/apps/sim/lib/mothership/chat/fork-inline-images.test.ts @@ -4,7 +4,7 @@ import { describe, expect, it, vi } from 'vitest' vi.mock('@/lib/uploads/core/storage-service', () => storageServiceMock) import { planForkInlineImages } from '@/lib/mothership/chat/fork-inline-images' -import { inlineChatImageKey } from '@/lib/mothership/chat/inline-image-storage' +import { inlineChatImageKey } from '@/lib/mothership/chat/inline-image-key' import type { PersistedMessage } from '@/lib/mothership/chat/persisted-message' const message = (content: string): PersistedMessage => ({ diff --git a/apps/sim/lib/mothership/chat/fork-inline-images.ts b/apps/sim/lib/mothership/chat/fork-inline-images.ts index d0d495883bd..17dd744ba33 100644 --- a/apps/sim/lib/mothership/chat/fork-inline-images.ts +++ b/apps/sim/lib/mothership/chat/fork-inline-images.ts @@ -1,14 +1,10 @@ import { OrchestrationError } from '@/lib/core/orchestration/types' import type { ChatBlobCopyTask } from '@/lib/mothership/chat/fork-chat-files' import { - inlineImageRequestIdSchema, - inlineImageSourceSchema, -} from '@/lib/mothership/chat/inline-image-reference' -import { - INLINE_CHAT_IMAGE_MAX_BYTES, inlineChatImageKey, -} from '@/lib/mothership/chat/inline-image-storage' -import { collectMarkdownImageSources } from '@/lib/mothership/chat/markdown-images' + inlineChatImageReferences, +} from '@/lib/mothership/chat/inline-image-key' +import { INLINE_CHAT_IMAGE_MAX_BYTES } from '@/lib/mothership/chat/inline-image-storage' import type { PersistedMessage } from '@/lib/mothership/chat/persisted-message' /** Only retained assistant image references are copied; their Markdown and request IDs stay unchanged. */ @@ -19,38 +15,25 @@ export function planForkInlineImages( ): ChatBlobCopyTask[] { const tasks = new Map() for (const message of messages) { - if ( - message.role !== 'assistant' || - !message.requestId || - !inlineImageRequestIdSchema.safeParse(message.requestId).success - ) - continue - const contents = [ - message.content, - ...(message.contentBlocks ?? []) - .filter((block) => block.type === 'text') - .map((block) => block.content ?? ''), - ] - for (const content of contents) { - for (const reference of collectMarkdownImageSources(content)) { - if (!inlineImageSourceSchema.safeParse(reference).success) continue - const sourceKey = inlineChatImageKey(sourceChatId, message.requestId, reference) - tasks.set(sourceKey, { - copyId: sourceKey, - sourceKey, - targetKey: inlineChatImageKey(newChatId, message.requestId, reference), - context: 'mothership', - fileName: 'chat-image.webp', - contentType: 'image/webp', - maxBytes: INLINE_CHAT_IMAGE_MAX_BYTES, - persistMetadata: false, - }) - if (tasks.size > 200) - throw new OrchestrationError( - 'payload_too_large', - 'A chat fork can copy at most 200 inline images.' - ) - } + const published = inlineChatImageReferences(message) + if (!published) continue + for (const reference of published.references) { + const sourceKey = inlineChatImageKey(sourceChatId, published.requestId, reference) + tasks.set(sourceKey, { + copyId: sourceKey, + sourceKey, + targetKey: inlineChatImageKey(newChatId, published.requestId, reference), + context: 'mothership', + fileName: 'chat-image.webp', + contentType: 'image/webp', + maxBytes: INLINE_CHAT_IMAGE_MAX_BYTES, + persistMetadata: false, + }) + if (tasks.size > 200) + throw new OrchestrationError( + 'payload_too_large', + 'A chat fork can copy at most 200 inline images.' + ) } } return [...tasks.values()] diff --git a/apps/sim/lib/mothership/chat/fork-worker.test.ts b/apps/sim/lib/mothership/chat/fork-worker.test.ts index 9d159598215..63bf3de8351 100644 --- a/apps/sim/lib/mothership/chat/fork-worker.test.ts +++ b/apps/sim/lib/mothership/chat/fork-worker.test.ts @@ -8,14 +8,13 @@ import { } from '@sim/testing/mocks/mothership-go-fetch.mock' import { generateId } from '@sim/utils/id' import { beforeEach, expect, it, vi } from 'vitest' +import { asOrchestrationError } from '@/lib/core/orchestration/types' import { copyWorkerConversation } from '@/lib/mothership/chat/fork-worker' import type { ForkChatRequest } from '@/lib/mothership/generated/protocol' vi.mock('@/lib/mothership/request/go/fetch', () => mothershipGoFetchMock) vi.mock('@/lib/mothership/server/agent-url', () => mothershipAgentUrlMock) -const fetchWorker = mothershipGoFetchMockFns.mockFetchGo - const request: ForkChatRequest = { sourceChatId: generateId(), newChatId: generateId(), @@ -27,37 +26,132 @@ const request: ForkChatRequest = { fileKeys: {}, } +type Answer = () => Promise + +/** What undici throws for a refused or reset socket: a TypeError with the syscall code on `cause`. */ +function socketFailure(code: string): TypeError { + return new TypeError('fetch failed', { cause: Object.assign(new Error(code), { code }) }) +} + +/** A fake worker: it records every fork request it receives and answers from a script. */ +let received: unknown[] = [] +let script: Answer[] = [] +const receipt: Answer = async () => + Response.json({ chatId: request.newChatId, sourceThroughSeq: 7 }) + +function answer(...answers: Answer[]) { + script = answers +} + +const status = + (code: number, body: BodyInit | null = ''): Answer => + async () => + new Response(body, { status: code }) + beforeEach(() => { mothershipAgentUrlMockFns.mockGetMothershipBaseURL.mockResolvedValue('http://worker.test') - fetchWorker.mockReset() - fetchWorker.mockImplementation(async () => - Response.json({ chatId: request.newChatId, sourceThroughSeq: 7 }) + received = [] + script = [] + mothershipGoFetchMockFns.mockFetchGo.mockReset() + mothershipGoFetchMockFns.mockFetchGo.mockImplementation( + async (_url: string, init: { body: string }) => { + received.push(JSON.parse(init.body)) + return (script.shift() ?? receipt)() + } ) }) +it.each(['EHOSTUNREACH', 'ENETUNREACH'])( + 'retries a fork the worker never received (%s)', + async (code) => { + answer(async () => { + throw socketFailure(code) + }) + await copyWorkerConversation(request) + expect(received).toEqual([request, request]) + } +) + it.each(['lost-response', 'temporary-error'])( 'retries the same immutable fork after %s', async (failure) => { - if (failure === 'lost-response') - fetchWorker.mockRejectedValueOnce(new TypeError('Connection ended')) - else fetchWorker.mockResolvedValueOnce(new Response('', { status: 503 })) + answer( + failure === 'lost-response' + ? async () => { + throw socketFailure('ECONNRESET') + } + : status(503) + ) await copyWorkerConversation(request) - expect(fetchWorker).toHaveBeenCalledTimes(2) - expect(fetchWorker.mock.calls[0][1].body).toBe(fetchWorker.mock.calls[1][1].body) - expect(JSON.parse(fetchWorker.mock.calls[1][1].body)).toEqual(request) + expect(received).toEqual([request, request]) } ) it.each(['missing-receipt', 'wrong-chat', 'unavailable'])( 'refuses an unconfirmed copy: %s', async (failure) => { - fetchWorker.mockImplementation(async () => { - if (failure === 'unavailable') throw new TypeError('Worker unavailable') + const reply: Answer = async () => { + if (failure === 'unavailable') throw socketFailure('ECONNREFUSED') return Response.json( failure === 'wrong-chat' ? { chatId: generateId(), sourceThroughSeq: 7 } : { ok: true } ) - }) + } + answer(reply, reply) await expect(copyWorkerConversation(request)).rejects.toThrow() - expect(fetchWorker).toHaveBeenCalledTimes(2) + // Only the unreachable worker may never have seen the request; an answer is final. + expect(received).toHaveLength(failure === 'unavailable' ? 2 : 1) } ) + +it.each([ + [404, 'not_found'], + [409, 'conflict'], + [413, 'payload_too_large'], +] as const)('classifies a worker %i refusal as %s without retrying it', async (code, kind) => { + answer(status(code, JSON.stringify({ error: 'refused' })), receipt) + const failure = await copyWorkerConversation(request).catch((error: unknown) => error) + expect(asOrchestrationError(failure)?.code).toBe(kind) + expect(received).toHaveLength(1) +}) + +it('classifies a refusal whose body fails to cancel', async () => { + const body = new ReadableStream({ + cancel() { + throw new TypeError('Body already closed') + }, + }) + answer(status(404, body), receipt) + const failure = await copyWorkerConversation(request).catch((error: unknown) => error) + expect(asOrchestrationError(failure)?.code).toBe('not_found') + expect(received).toHaveLength(1) +}) + +it.each([500, 400])('does not repeat a fork the worker failed with %i', async (code) => { + answer(status(code), receipt) + const failure = await copyWorkerConversation(request).catch((error: unknown) => error) + expect(failure).toBeInstanceOf(Error) + expect(asOrchestrationError(failure)).toBeNull() + expect(received).toHaveLength(1) +}) + +it.each([502, 504])('retries a %i gateway failure once', async (code) => { + answer(status(code)) + await copyWorkerConversation(request) + expect(received).toEqual([request, request]) +}) + +it('does not start a second copy while a timed-out one may still be running', async () => { + answer(async () => { + throw new DOMException('The operation timed out.', 'TimeoutError') + }, receipt) + await expect(copyWorkerConversation(request)).rejects.toThrow() + expect(received).toHaveLength(1) +}) + +it('does not repeat a fork after a TypeError that is not a socket failure', async () => { + answer(async () => { + throw new TypeError('Body is unusable: Body has already been read') + }, receipt) + await expect(copyWorkerConversation(request)).rejects.toThrow() + expect(received).toHaveLength(1) +}) diff --git a/apps/sim/lib/mothership/chat/fork-worker.ts b/apps/sim/lib/mothership/chat/fork-worker.ts index bda6145b4dd..d5570db0c2b 100644 --- a/apps/sim/lib/mothership/chat/fork-worker.ts +++ b/apps/sim/lib/mothership/chat/fork-worker.ts @@ -1,35 +1,123 @@ +import { createLogger } from '@sim/logger' +import { getErrorMessage } from '@sim/utils/errors' +import { isRetryableNetworkError } from '@/lib/core/errors/retryable-infrastructure' +import { OrchestrationError } from '@/lib/core/orchestration/types' import { type ForkChatRequest, ForkChatResponse } from '@/lib/mothership/generated/protocol' import { fetchGo } from '@/lib/mothership/request/go/fetch' import { mothershipRequestHeaders } from '@/lib/mothership/request/headers' import { getMothershipBaseURL } from '@/lib/mothership/server/agent-url' -/** A lost copy acknowledgement retries the same destination and immutable request. */ +const logger = createLogger('ForkWorker') + +/** + * How long one attempt waits for the copy. The worker refuses a cut above 50,000 events + * with 413, and measured locally it copies 15,000 events in 1.5–2 s and 50,000 in 5–9 s; + * production databases are slower by an estimated factor of up to eight, which puts the + * ceiling at 40–70 s. 120 s leaves headroom above that, and the load balancers on both + * sides keep idle connections far longer. A legitimate fork therefore finishes inside one + * attempt, and one that does not is abandoned: the worker rolls a copy back once its + * caller's connection closes, so a timed-out attempt is never retried and never published. + */ +const ATTEMPT_TIMEOUT_MS = 120_000 + +/** Cleanup is a small archive-and-purge of one chat, never a copy. */ +const DISCARD_TIMEOUT_MS = 10_000 + +/** Gateway failures: the request may never have reached a worker, so one more attempt is safe. */ +const RETRYABLE_STATUSES = new Set([502, 503, 504]) + +/** The worker's typed refusals, each one the caller can act on. */ +function workerRefusal(status: number): Error { + if (status === 404) return new OrchestrationError('not_found', 'Chat not found') + if (status === 409) + return new OrchestrationError( + 'conflict', + 'The selected response has not finished. Retry the fork once it completes.' + ) + if (status === 413) + return new OrchestrationError( + 'payload_too_large', + 'This conversation is too long to fork. Fork from an earlier message.' + ) + return new Error('The conversation could not be copied. Retry the fork.') +} + +type Attempt = { kind: 'receipt'; body: unknown } | { kind: 'status'; status: number } + +/** + * A lost copy acknowledgement retries the same destination and immutable request; the + * worker answers a repeat with the first attempt's receipt. Only a refused, reset or dropped + * socket or a gateway status is retried: a timed-out attempt may still be copying. + */ export async function copyWorkerConversation(request: ForkChatRequest): Promise { const baseURL = await getMothershipBaseURL({ userId: request.userId }) const body = JSON.stringify(request) - for (let attempt = 0; attempt < 2; attempt++) { + for (let attempt = 0; ; attempt++) { + const canRetry = attempt === 0 + let outcome: Attempt try { const response = await fetchGo(`${baseURL}/api/chats/fork`, { method: 'POST', headers: mothershipRequestHeaders(), body, - signal: AbortSignal.timeout(15_000), + signal: AbortSignal.timeout(ATTEMPT_TIMEOUT_MS), spanName: 'sim → worker /api/chats/fork', operation: 'fork_chat', }) - if (!response.ok) { - if (response.status >= 500 && attempt === 0) { - await response.body?.cancel() - continue - } - throw new Error('The conversation could not be copied. Retry the fork.') + if (response.ok) outcome = { kind: 'receipt', body: await response.json() } + else { + outcome = { kind: 'status', status: response.status } + // Releasing the unread error body is best effort; the status is already the answer. + await response.body?.cancel().catch(() => undefined) } - const receipt = ForkChatResponse.parse(await response.json()) - if (receipt.chatId !== request.newChatId) - throw new Error('The fork returned a different chat') - return } catch (error) { - if (attempt === 1) throw error + if (canRetry && isRetryableNetworkError(error)) continue + throw error + } + if (outcome.kind === 'status') { + if (canRetry && RETRYABLE_STATUSES.has(outcome.status)) continue + throw workerRefusal(outcome.status) } + const receipt = ForkChatResponse.parse(outcome.body) + if (receipt.chatId !== request.newChatId) throw new Error('The fork returned a different chat') + return + } +} + +/** + * Removes the worker's copy of a fork that will not be published, so the worker keeps no + * conversation Sim has no chat for. Best effort: a failure is logged, never thrown, since + * the caller is already reporting the fork's own failure. + * + * The worker rolls back a copy whose caller disconnected, so after a timed-out attempt + * this is normally a no-op. A worker that commits in the instant between its last + * disconnect check and its commit, or one deployed before that rollback existed, can + * commit after this call runs; that conversation then stays on the worker, unreachable + * from Sim. + */ +export async function discardWorkerConversation( + request: Pick +): Promise { + try { + const baseURL = await getMothershipBaseURL({ userId: request.userId }) + const response = await fetchGo(`${baseURL}/api/tasks/cleanup`, { + method: 'POST', + headers: mothershipRequestHeaders(), + body: JSON.stringify({ chatIds: [request.newChatId] }), + signal: AbortSignal.timeout(DISCARD_TIMEOUT_MS), + spanName: 'sim → worker /api/tasks/cleanup', + operation: 'discard_fork', + }) + await response.body?.cancel().catch(() => undefined) + if (!response.ok) + logger.warn('Worker refused to discard an unpublished fork', { + chatId: request.newChatId, + status: response.status, + }) + } catch (error) { + logger.warn('Failed to discard an unpublished fork on the worker', { + chatId: request.newChatId, + error: getErrorMessage(error), + }) } } diff --git a/apps/sim/lib/mothership/chat/inline-image-key.ts b/apps/sim/lib/mothership/chat/inline-image-key.ts new file mode 100644 index 00000000000..aa35c449011 --- /dev/null +++ b/apps/sim/lib/mothership/chat/inline-image-key.ts @@ -0,0 +1,47 @@ +import { createHash } from 'node:crypto' +import { + INLINE_CHAT_IMAGE_PREFIX, + inlineImageRequestIdSchema, + inlineImageSourceSchema, + normalizeInlineFileReference, +} from '@/lib/mothership/chat/inline-image-reference' +import { collectMarkdownImageSources } from '@/lib/mothership/chat/markdown-images' +import type { PersistedMessage } from '@/lib/mothership/chat/persisted-message' + +/** Stable per-message source identity makes replay independent of sandbox lifetime. */ +export function inlineChatImageKey(chatId: string, requestId: string, reference: string): string { + inlineImageRequestIdSchema.parse(chatId) + inlineImageRequestIdSchema.parse(requestId) + const digest = createHash('sha256').update(normalizeInlineFileReference(reference)).digest('hex') + return `${INLINE_CHAT_IMAGE_PREFIX}${chatId}/${requestId}/${digest}.webp` +} + +/** + * The inline image references an assistant message published under its request id: every + * first-party Markdown image source in its text and text blocks. Anything else (user text, + * a message with no request receipt, remote URLs) never had a chat image stored. + */ +export function inlineChatImageReferences( + message: Pick +): { requestId: string; references: string[] } | null { + if ( + message.role !== 'assistant' || + !message.requestId || + !inlineImageRequestIdSchema.safeParse(message.requestId).success + ) + return null + const contents = [ + message.content, + ...(message.contentBlocks ?? []) + .filter((block) => block.type === 'text') + .map((block) => block.content ?? ''), + ] + const references = new Set() + for (const content of contents) { + // Every Markdown image, inline or reference-style, opens with `![`; most messages have none. + if (typeof content !== 'string' || !content.includes('![')) continue + for (const reference of collectMarkdownImageSources(content)) + if (inlineImageSourceSchema.safeParse(reference).success) references.add(reference) + } + return { requestId: message.requestId, references: [...references] } +} diff --git a/apps/sim/lib/mothership/chat/inline-image-storage.ts b/apps/sim/lib/mothership/chat/inline-image-storage.ts index 0747cd7c3df..40f253ceda5 100644 --- a/apps/sim/lib/mothership/chat/inline-image-storage.ts +++ b/apps/sim/lib/mothership/chat/inline-image-storage.ts @@ -1,11 +1,6 @@ -import { createHash } from 'node:crypto' import sharp from 'sharp' import { OrchestrationError } from '@/lib/core/orchestration/types' -import { - INLINE_CHAT_IMAGE_PREFIX, - inlineImageRequestIdSchema, - normalizeInlineFileReference, -} from '@/lib/mothership/chat/inline-image-reference' +import { inlineChatImageKey } from '@/lib/mothership/chat/inline-image-key' import { isObjectNotFoundError } from '@/lib/uploads/core/errors' import { downloadFile, uploadFile } from '@/lib/uploads/core/storage-service' @@ -112,11 +107,3 @@ export async function loadInlineChatImage( throw error } } - -/** Stable per-message source identity makes replay independent of sandbox lifetime. */ -export function inlineChatImageKey(chatId: string, requestId: string, reference: string): string { - inlineImageRequestIdSchema.parse(chatId) - inlineImageRequestIdSchema.parse(requestId) - const digest = createHash('sha256').update(normalizeInlineFileReference(reference)).digest('hex') - return `${INLINE_CHAT_IMAGE_PREFIX}${chatId}/${requestId}/${digest}.webp` -} diff --git a/apps/sim/lib/mothership/chat/rewrite-file-references.ts b/apps/sim/lib/mothership/chat/rewrite-file-references.ts index e7ebea3dc2f..1ecc8fa7d2e 100644 --- a/apps/sim/lib/mothership/chat/rewrite-file-references.ts +++ b/apps/sim/lib/mothership/chat/rewrite-file-references.ts @@ -1,3 +1,5 @@ +import { isPlainRecord } from '@sim/utils/object' +import type { PersistedContentBlock } from '@/lib/api/contracts/copilot-messages' import type { PersistedMessage } from '@/lib/mothership/chat/persisted-message' import type { MothershipResource } from '@/lib/mothership/resources/types' import { rewriteForkContentRefs } from '@/ee/workspace-forking/lib/remap/remap-content-refs' @@ -10,6 +12,33 @@ import { rewriteForkContentRefs } from '@/ee/workspace-forking/lib/remap/remap-c export interface ChatFileRefMaps { fileIds: ReadonlyMap fileKeys: ReadonlyMap + /** + * The workspace both chats live in. A fork stays in its source's workspace, so in-app + * `/workspace//files/` links keep their workspace and only the file id moves. + */ + workspaceId?: string | null +} + +/** + * The copy plan's id/key maps restricted to copies whose bytes were prepared. A failed copy + * is never published, so the fork's messages, resources and worker history keep naming the + * source file (alive while the source chat is) instead of an id or key that never exists. + */ +export function publishedFileRefMaps( + plan: { + idMap: ReadonlyMap + keyMap: ReadonlyMap + blobTasks: readonly { copyId: string; targetKey: string }[] + }, + failedCopyIds: ReadonlySet +): { fileIds: Map; fileKeys: Map } { + const failedKeys = new Set( + plan.blobTasks.filter((task) => failedCopyIds.has(task.copyId)).map((task) => task.targetKey) + ) + return { + fileIds: new Map([...plan.idMap].filter(([, copyId]) => !failedCopyIds.has(copyId))), + fileKeys: new Map([...plan.keyMap].filter(([, copyKey]) => !failedKeys.has(copyKey))), + } } function hasMappings(maps: ChatFileRefMaps): boolean { @@ -17,15 +46,64 @@ function hasMappings(maps: ChatFileRefMaps): boolean { } function rewriteText(text: string, maps: ChatFileRefMaps): string { - return rewriteForkContentRefs(text, { fileIds: maps.fileIds, fileKeys: maps.fileKeys }) + return rewriteForkContentRefs(text, { + fileIds: maps.fileIds, + fileKeys: maps.fileKeys, + ...(maps.workspaceId ? { workspaceId: { from: maps.workspaceId, to: maps.workspaceId } } : {}), + }) +} + +/** + * Tool arguments name a file by its bare id or storage key (`{ fileId }`, `["files", "read", id]`), + * so a string that IS a mapped id or key is replaced whole; any other string gets the URL grammar. + */ +function rewriteToolValue(value: unknown, maps: ChatFileRefMaps): unknown { + if (typeof value === 'string') + return maps.fileIds.get(value) ?? maps.fileKeys.get(value) ?? rewriteText(value, maps) + if (Array.isArray(value)) return value.map((entry) => rewriteToolValue(entry, maps)) + if (isPlainRecord(value)) + return Object.fromEntries( + Object.entries(value).map(([key, entry]) => [key, rewriteToolValue(entry, maps)]) + ) + return value +} + +function rewriteBlock(block: PersistedContentBlock, maps: ChatFileRefMaps): PersistedContentBlock { + const toolCall = block.toolCall + return { + ...block, + ...(block.content ? { content: rewriteText(block.content, maps) } : {}), + ...(toolCall + ? { + toolCall: { + ...toolCall, + ...(toolCall.params + ? { params: rewriteToolValue(toolCall.params, maps) as Record } + : {}), + ...(toolCall.activityDescription + ? { activityDescription: rewriteText(toolCall.activityDescription, maps) } + : {}), + ...(toolCall.display?.title + ? { + display: { + ...toolCall.display, + title: rewriteText(toolCall.display.title, maps), + }, + } + : {}), + }, + } + : {}), + } } /** * Re-point every file reference in a copied transcript at the copied files, so * the fork is self-contained (it survives the original chat's deletion). * Rewrites: free-text URLs in `content` and text content blocks (serve/view/ - * in-app/`sim:file` forms, via the shared fork grammar), attachment chip - * ids+keys, and `@`-mention context chip file ids. References to anything not + * in-app/`sim:file` forms, via the shared fork grammar), tool-call arguments + * and display text, attachment chip ids+keys, and `@`-mention context chip + * file ids. References to anything not * in the maps (shared workspace files, workflows, other chats) pass through * unchanged. Pure; returns the input array untouched when there is nothing to * rewrite. @@ -50,9 +128,7 @@ export function rewriteMessageFileRefs( content: rewriteText(message.content, maps), } if (message.contentBlocks?.length) { - rewritten.contentBlocks = message.contentBlocks.map((block) => - block.content ? { ...block, content: rewriteText(block.content, maps) } : block - ) + rewritten.contentBlocks = message.contentBlocks.map((block) => rewriteBlock(block, maps)) } if (message.fileAttachments?.length) { rewritten.fileAttachments = message.fileAttachments.map((att) => ({