From 37af9891221731bf2ee1853fc71b8c64de8eeac4 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Fri, 2 Oct 2026 12:33:38 -0700 Subject: [PATCH 1/9] fix(core): treat EHOSTUNREACH as a retryable network failure An unreachable host is a connection that never opened, so the request never reached its destination and one more attempt is as safe as after ECONNREFUSED or ENETUNREACH. --- apps/sim/lib/core/errors/retryable-infrastructure.ts | 1 + 1 file changed, 1 insertion(+) 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', From 329b20be39eeac012011424d0ba3d26456677443 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Fri, 2 Oct 2026 12:33:39 -0700 Subject: [PATCH 2/9] fix(mothership): pass typed fork refusals through and stop overlapping retries The worker refuses a fork it cannot make with 404 (the chat or message is gone), 409 (the response has not finished) or 413 (the cut is above its ceiling), but every non-2xx became a generic error and a 500. Those are now classified and passed through by the fork route with a message the person can act on. The retry loop's own catch swallowed every first-attempt failure, so a 400, a 500, a malformed receipt and a timed-out attempt were all sent again. A timed-out attempt may still be copying, so only a recognized socket failure or a 502/503/504 gets one more attempt now. The per-attempt timeout is sized from measured copy latency at the worker's ceiling with headroom, so a legitimate fork finishes inside one attempt and a timed-out one is abandoned (the worker rolls back a copy whose caller left). --- .../chats/[chatId]/fork/route.test.ts | 6 + .../mothership/chats/[chatId]/fork/route.ts | 7 +- .../lib/mothership/chat/fork-worker.test.ts | 124 +++++++++++++++--- apps/sim/lib/mothership/chat/fork-worker.ts | 71 ++++++++-- 4 files changed, 178 insertions(+), 30 deletions(-) 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..a36bb60743f 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 @@ -316,6 +316,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, 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/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..f5e4887ea73 100644 --- a/apps/sim/lib/mothership/chat/fork-worker.ts +++ b/apps/sim/lib/mothership/chat/fork-worker.ts @@ -1,35 +1,78 @@ +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. */ +/** + * 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 about 2 s and 50,000 in about + * 9 s; production databases are slower by an estimated factor of up to eight, which puts + * the ceiling near 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 + +/** 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 } } From af1108d8fa9cef31ff19104e3a46c4e1de245faa Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Fri, 2 Oct 2026 12:33:49 -0700 Subject: [PATCH 3/9] fix(mothership): show a fork refusal's own message A 409 or 413 from the fork route tells the person what to do (wait for the response to finish, or fork from an earlier message), so the toast shows that message instead of a generic failure. --- .../components/message-actions/message-actions.tsx | 12 ++++++++++-- 1 file changed, 10 insertions(+), 2 deletions(-) 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' + ) } } From 3fdbd6b59f14f94e2ebb6c42ae596c3e622211a3 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Fri, 2 Oct 2026 12:33:56 -0700 Subject: [PATCH 4/9] fix(mothership): discard the worker copy of a fork that is not published When anything after the worker copy failed, most often the final transaction that publishes the chat, the worker kept a conversation Sim had no chat for. The fork now asks the worker to clean that chat up, through the same cleanup endpoint chat deletion uses. It is best effort and a no-op when the worker never committed. --- .../chats/[chatId]/fork/route.test.ts | 13 ++++++ .../lib/mothership/chat/application/fork.ts | 11 ++++- apps/sim/lib/mothership/chat/fork-worker.ts | 45 +++++++++++++++++++ 3 files changed, 67 insertions(+), 2 deletions(-) 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 a36bb60743f..58bbe020139 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 @@ -342,6 +342,19 @@ 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('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/lib/mothership/chat/application/fork.ts b/apps/sim/lib/mothership/chat/application/fork.ts index d3c2ebcec3b..47527c92bc4 100644 --- a/apps/sim/lib/mothership/chat/application/fork.ts +++ b/apps/sim/lib/mothership/chat/application/fork.ts @@ -19,7 +19,10 @@ 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 { @@ -88,6 +91,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) @@ -111,7 +116,6 @@ export const forkChat = defineAuthorizedChatUseCase({ /** 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}` @@ -129,6 +133,7 @@ export const forkChat = defineAuthorizedChatUseCase({ ).filter((resource) => resource.type !== 'file' || !failedIds.has(resource.id)) 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, @@ -198,6 +203,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/fork-worker.ts b/apps/sim/lib/mothership/chat/fork-worker.ts index f5e4887ea73..496e2952750 100644 --- a/apps/sim/lib/mothership/chat/fork-worker.ts +++ b/apps/sim/lib/mothership/chat/fork-worker.ts @@ -1,3 +1,5 @@ +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' @@ -5,6 +7,8 @@ import { fetchGo } from '@/lib/mothership/request/go/fetch' import { mothershipRequestHeaders } from '@/lib/mothership/request/headers' import { getMothershipBaseURL } from '@/lib/mothership/server/agent-url' +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 about 2 s and 50,000 in about @@ -16,6 +20,9 @@ import { getMothershipBaseURL } from '@/lib/mothership/server/agent-url' */ 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]) @@ -76,3 +83,41 @@ export async function copyWorkerConversation(request: ForkChatRequest): Promise< 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), + }) + } +} From b22f5d8656c7692a7a298bda89d437058ce6a1db Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Fri, 2 Oct 2026 12:34:05 -0700 Subject: [PATCH 5/9] fix(mothership): point every reference a fork copies at what the fork holds - A file whose blob copy failed is never published, yet the fork's messages and the worker's history were rewritten to its id and key. References to it now stay on the source file, and its resource tab is dropped as before. - In-app /workspace//files/ links were left on the source file; the fork stays in the same workspace, so only the file id moves. - Tool-call arguments and display titles kept the source file's id or key. - A Sources tab for a response past the cut was copied, pointing at a message the fork does not have. --- .../chats/[chatId]/fork/route.test.ts | 97 +++++++++++++++++++ .../lib/mothership/chat/application/fork.ts | 33 ++++--- .../chat/rewrite-file-references.ts | 88 +++++++++++++++-- 3 files changed, 199 insertions(+), 19 deletions(-) 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 58bbe020139..e8d1365b44d 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 @@ -355,6 +355,103 @@ describe('POST /api/mothership/chats/[chatId]/fork', () => { 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/lib/mothership/chat/application/fork.ts b/apps/sim/lib/mothership/chat/application/fork.ts index 47527c92bc4..2d41725ec67 100644 --- a/apps/sim/lib/mothership/chat/application/fork.ts +++ b/apps/sim/lib/mothership/chat/application/fork.ts @@ -26,6 +26,7 @@ import { 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' @@ -100,18 +101,24 @@ 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)) @@ -125,12 +132,12 @@ 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 @@ -143,8 +150,8 @@ 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. */ 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) => ({ From 5298c1c3879adab4aada396a1091b3d9db300ae7 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Fri, 2 Oct 2026 12:34:06 -0700 Subject: [PATCH 6/9] fix(cleanup): purge a chat's inline images and only its own attachments - Chat images are stored under the chat's id but were never deleted with the chat. Their keys are rebuilt from the assistant messages that published them, the same way a fork finds them to copy. - Message attachments were deleted as copilot storage whatever their key. Attachments in Chat are workspace-bucket files owned by their workspace_files row, which a fork can share with its source (a failed copy, a deleted file), and the copilot bucket falls back to the workspace bucket on GCS and can be configured to it on S3. Only copilot keys are deleted from messages now. --- apps/sim/lib/cleanup/chat-cleanup.test.ts | 76 +++++++++++++++++++ apps/sim/lib/cleanup/chat-cleanup.ts | 55 +++++++++++++- .../chat/application/inline-images.test.ts | 2 +- .../chat/fork-inline-images.test.ts | 2 +- .../lib/mothership/chat/fork-inline-images.ts | 61 ++++++--------- .../lib/mothership/chat/inline-image-key.ts | 47 ++++++++++++ .../mothership/chat/inline-image-storage.ts | 15 +--- 7 files changed, 201 insertions(+), 57 deletions(-) create mode 100644 apps/sim/lib/cleanup/chat-cleanup.test.ts create mode 100644 apps/sim/lib/mothership/chat/inline-image-key.ts 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..48ef0d28962 --- /dev/null +++ b/apps/sim/lib/cleanup/chat-cleanup.test.ts @@ -0,0 +1,76 @@ +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`; returns the deleted keys by context. */ +async function purge(messages: Record[]) { + queueTableRows(workspaceFiles, []) + queueTableRows( + copilotMessages, + messages.map((content) => ({ chatId, content })) + ) + const cleanup = await prepareChatCleanup([chatId], 'test') + queueTableRows(copilotChats, []) + 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('deletes only copilot-storage attachment keys 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('assistant/org-1/user-1/u1/shared.png'), + attachment('copilot/1234/legacy.png'), + ], + }, + ]) + expect(deleted).toEqual({ copilot: ['copilot/1234/legacy.png'] }) + }) + + 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..f5fc0dac3de 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 { isRecordLike } from '@sim/utils/object' import { and, inArray, isNull } 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') @@ -33,9 +39,47 @@ interface FileRef { } /** - * Collect all file storage keys for the given chat IDs from two sources: + * 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 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, for copilot-storage keys only + * 3. the chat-scoped inline images each assistant message published + * + * An attachment whose key belongs to another storage context 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 +121,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,6 +136,7 @@ export async function collectChatFiles(chatIds: string[]): Promise { (attachment as Record).key ) { const key = (attachment as Record).key as string + if (tryInferContextFromKey(key) !== 'copilot') continue if (!seen.has(key)) { seen.add(key) files.push({ key, context: 'copilot', chatId: row.chatId }) 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/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` -} From a286f0a98798e2435c399f3f718fe363df08ecbd Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Fri, 2 Oct 2026 12:53:56 -0700 Subject: [PATCH 7/9] fix(cleanup): delete an organization attachment once no chat references it Organization Chat attachments (assistant/ keys) have no workspace_files row and a fork carries the same key, so the purge now deletes one under its own storage context only after confirming no remaining chat of that organization still references it. --- apps/sim/lib/cleanup/chat-cleanup.test.ts | 33 +++++++- apps/sim/lib/cleanup/chat-cleanup.ts | 93 ++++++++++++++++++++--- 2 files changed, 111 insertions(+), 15 deletions(-) diff --git a/apps/sim/lib/cleanup/chat-cleanup.test.ts b/apps/sim/lib/cleanup/chat-cleanup.test.ts index 48ef0d28962..ef054b4696b 100644 --- a/apps/sim/lib/cleanup/chat-cleanup.test.ts +++ b/apps/sim/lib/cleanup/chat-cleanup.test.ts @@ -15,8 +15,14 @@ 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`; returns the deleted keys by context. */ -async function purge(messages: Record[]) { +/** + * 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, @@ -24,6 +30,10 @@ async function purge(messages: Record[]) { ) 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]) @@ -38,7 +48,7 @@ describe('chat purge storage', () => { storageServiceMockFns.mockDeleteFiles.mockResolvedValue({ deleted: 1, failed: [] }) }) - it('deletes only copilot-storage attachment keys as copilot storage', async () => { + 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([ @@ -47,7 +57,6 @@ describe('chat purge storage', () => { content: 'look', fileAttachments: [ attachment('workspace/ws-1/1700-abc-shared.png'), - attachment('assistant/org-1/user-1/u1/shared.png'), attachment('copilot/1234/legacy.png'), ], }, @@ -55,6 +64,22 @@ describe('chat purge storage', () => { 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([ { diff --git a/apps/sim/lib/cleanup/chat-cleanup.ts b/apps/sim/lib/cleanup/chat-cleanup.ts index f5fc0dac3de..ab9745dd093 100644 --- a/apps/sim/lib/cleanup/chat-cleanup.ts +++ b/apps/sim/lib/cleanup/chat-cleanup.ts @@ -3,7 +3,7 @@ import { copilotChats, copilotMessages, workspaceFiles } from '@sim/db/schema' import { createLogger } from '@sim/logger' import { chunkArray } from '@sim/utils/helpers' import { isRecordLike } from '@sim/utils/object' -import { and, inArray, isNull } from 'drizzle-orm' +import { and, eq, inArray, isNull, or, sql } from 'drizzle-orm' import { env } from '@/lib/core/config/env' import { inlineChatImageKey, @@ -36,6 +36,19 @@ 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 } /** @@ -71,15 +84,15 @@ function inlineChatImageKeys(chatId: string, content: Record): /** * 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, for copilot-storage keys only + * 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 * - * An attachment whose key belongs to another storage context 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. + * 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[] = [] @@ -136,8 +149,12 @@ export async function collectChatFiles(chatIds: string[]): Promise { (attachment as Record).key ) { const key = (attachment as Record).key as string - if (tryInferContextFromKey(key) !== 'copilot') continue - 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 }) } @@ -149,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[], @@ -271,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) { From fc33334180c4aa2ac651b77e75c7f36dbebb2fdd Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Fri, 2 Oct 2026 13:01:50 -0700 Subject: [PATCH 8/9] fix(mothership): order fork publication against a purge of its source A 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 after it deletes the source row. The fork's publish transaction now holds the source row with FOR KEY SHARE: a purge that already removed it refuses the fork (404, worker copy discarded), and one that has not waits for the commit and then sees the fork's references. --- .../chats/[chatId]/fork/route.test.ts | 21 +++++++++++++++++++ .../lib/mothership/chat/application/fork.ts | 14 +++++++++++++ 2 files changed, 35 insertions(+) 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 e8d1365b44d..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) => { diff --git a/apps/sim/lib/mothership/chat/application/fork.ts b/apps/sim/lib/mothership/chat/application/fork.ts index 2d41725ec67..97b15c298ce 100644 --- a/apps/sim/lib/mothership/chat/application/fork.ts +++ b/apps/sim/lib/mothership/chat/application/fork.ts @@ -156,6 +156,20 @@ export const forkChat = defineAuthorizedChatUseCase({ /** 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({ From ef3e2721aa92bb9df670e717b7d2a8e34611d600 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Fri, 2 Oct 2026 13:01:50 -0700 Subject: [PATCH 9/9] chore(mothership): quote the final copy latency behind the fork attempt timeout --- apps/sim/lib/mothership/chat/fork-worker.ts | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/apps/sim/lib/mothership/chat/fork-worker.ts b/apps/sim/lib/mothership/chat/fork-worker.ts index 496e2952750..d5570db0c2b 100644 --- a/apps/sim/lib/mothership/chat/fork-worker.ts +++ b/apps/sim/lib/mothership/chat/fork-worker.ts @@ -11,9 +11,9 @@ 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 about 2 s and 50,000 in about - * 9 s; production databases are slower by an estimated factor of up to eight, which puts - * the ceiling near 70 s. 120 s leaves headroom above that, and the load balancers on both + * 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.