Skip to content

Commit d05ae7a

Browse files
authored
fix(mothership): make chat fork complete, consistent and safe to retry (#8583)
* 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. * 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). * 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. * 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. * 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/<id>/files/<fileId> 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. * 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. * 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. * 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. * chore(mothership): quote the final copy latency behind the fork attempt timeout
1 parent 44016cd commit d05ae7a

15 files changed

Lines changed: 790 additions & 113 deletions

File tree

‎apps/sim/app/api/mothership/chats/[chatId]/fork/route.test.ts‎

Lines changed: 137 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -185,6 +185,7 @@ describe('POST /api/mothership/chats/[chatId]/fork', () => {
185185
queueTableRows(copilotChats, [chat])
186186
queueTableRows(copilotChats, [chat])
187187
queueTableRows(member, [{ role: 'member' }])
188+
queueTableRows(copilotChats, [{ id: chat.id }])
188189
const attachments = [
189190
{
190191
id: 'upload-1',
@@ -215,6 +216,26 @@ describe('POST /api/mothership/chats/[chatId]/fork', () => {
215216
expect(mockAssertActiveWorkspaceAccess).not.toHaveBeenCalled()
216217
})
217218

219+
it('refuses to publish a fork whose source chat was purged while it copied', async () => {
220+
// The fork shares its attachment keys with the source, and cleanup deletes an unreferenced
221+
// key only after deleting the source row: publishing without the source would leave the
222+
// fork pointing at bytes that cleanup is about to delete.
223+
dbChainMockFns.limit.mockReset()
224+
const chat = { ...parentRow, workspaceId: null, organizationId: 'org-1', resources: [] }
225+
queueTableRows(copilotChats, [chat])
226+
queueTableRows(copilotChats, [chat])
227+
queueTableRows(member, [{ role: 'member' }])
228+
queueTableRows(copilotChats, [])
229+
const res = await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' }))
230+
expect(res.status).toBe(404)
231+
expect(mockAppendCopilotChatMessages).not.toHaveBeenCalled()
232+
expect(mockPublishStatusChanged).not.toHaveBeenCalled()
233+
expect(mockFetchGo.mock.calls.map(([url]) => url)).toEqual([
234+
'http://mothership.test/api/chats/fork',
235+
'http://mothership.test/api/tasks/cleanup',
236+
])
237+
})
238+
218239
it.each(['membership', 'capability'] as const)(
219240
'denies organization forks after %s revocation before copying',
220241
async (revocation) => {
@@ -316,6 +337,12 @@ describe('POST /api/mothership/chats/[chatId]/fork', () => {
316337
expect(mockPublishStatusChanged).not.toHaveBeenCalled()
317338
})
318339

340+
it.each([404, 409, 413])('passes a worker %i refusal through', async (status) => {
341+
mockFetchGo.mockResolvedValue(Response.json({ error: 'refused' }, { status }))
342+
const res = await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' }))
343+
expect(res.status).toBe(status)
344+
})
345+
319346
it('surfaces failed blob copies and excludes their metadata from publication', async () => {
320347
mockExecuteChatFileBlobCopies.mockResolvedValue({
321348
copied: 1,
@@ -336,6 +363,116 @@ describe('POST /api/mothership/chats/[chatId]/fork', () => {
336363
expect(dbChainMockFns.delete).not.toHaveBeenCalled()
337364
})
338365

366+
it('discards the worker copy when the fork cannot be published', async () => {
367+
dbChainMockFns.transaction.mockRejectedValueOnce(new Error('connection lost'))
368+
const res = await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' }))
369+
expect(res.status).toBe(500)
370+
const [fork, discard] = mockFetchGo.mock.calls
371+
expect(fork[0]).toBe('http://mothership.test/api/chats/fork')
372+
expect(discard[0]).toBe('http://mothership.test/api/tasks/cleanup')
373+
expect(JSON.parse(discard[1].body)).toEqual({
374+
chatIds: [JSON.parse(fork[1].body).newChatId],
375+
})
376+
expect(mockPublishStatusChanged).not.toHaveBeenCalled()
377+
})
378+
379+
it('keeps references on the source file when its copy fails', async () => {
380+
const oldKey = 'workspace/ws-1/old-cat.png'
381+
mockListForkableChatFiles.mockResolvedValue([
382+
{ id: OLD_FILE_ID, key: oldKey, messageId: 'msg-1', workspaceId: 'ws-1' },
383+
])
384+
mockPlanChatFileCopies.mockReturnValue({
385+
idMap: new Map([[OLD_FILE_ID, NEW_FILE_ID]]),
386+
keyMap: new Map([[oldKey, 'workspace/ws-1/new-cat.png']]),
387+
blobTasks: [
388+
{
389+
copyId: NEW_FILE_ID,
390+
sourceKey: oldKey,
391+
targetKey: 'workspace/ws-1/new-cat.png',
392+
context: 'mothership',
393+
fileName: 'cat.png',
394+
contentType: 'image/png',
395+
},
396+
],
397+
})
398+
mockExecuteChatFileBlobCopies.mockResolvedValue({
399+
copied: 0,
400+
failed: 1,
401+
failedCopyIds: [NEW_FILE_ID],
402+
})
403+
const res = await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' }))
404+
expect(res.status).toBe(200)
405+
// The copy is never published, so neither Sim's transcript nor the worker's may name it.
406+
expect(mockAppendCopilotChatMessages.mock.calls[0][1][0].content).toBe(
407+
`See ![cat](/api/files/view/${OLD_FILE_ID})`
408+
)
409+
const forkRequest = JSON.parse(mockFetchGo.mock.calls[0][1].body)
410+
expect(forkRequest.fileIds).toEqual({})
411+
expect(forkRequest.fileKeys).toEqual({})
412+
})
413+
414+
it('re-points in-app file links and tool-call arguments at the copied file', async () => {
415+
mockLoadCopilotChatMessages.mockResolvedValue([
416+
{ ...threeMessages[0], content: `Open /workspace/ws-1/files/${OLD_FILE_ID}` },
417+
{
418+
...threeMessages[1],
419+
contentBlocks: [
420+
{
421+
type: 'tool',
422+
toolCall: {
423+
id: 'call-1',
424+
name: 'sim_cli',
425+
state: 'success',
426+
params: { fileId: OLD_FILE_ID, args: ['files', 'read', OLD_FILE_ID], n: 2 },
427+
display: { title: `Read /api/files/view/${OLD_FILE_ID}` },
428+
},
429+
},
430+
],
431+
},
432+
])
433+
mockPlanChatFileCopies.mockReturnValue({
434+
idMap: new Map([[OLD_FILE_ID, NEW_FILE_ID]]),
435+
keyMap: new Map(),
436+
blobTasks: [],
437+
})
438+
const res = await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' }))
439+
expect(res.status).toBe(200)
440+
const [user, assistant] = mockAppendCopilotChatMessages.mock.calls[0][1]
441+
expect(user.content).toBe(`Open /workspace/ws-1/files/${NEW_FILE_ID}`)
442+
expect(assistant.contentBlocks[0].toolCall).toMatchObject({
443+
params: { fileId: NEW_FILE_ID, args: ['files', 'read', NEW_FILE_ID], n: 2 },
444+
display: { title: `Read /api/files/view/${NEW_FILE_ID}` },
445+
})
446+
})
447+
448+
it('drops a Sources tab whose response is past the cut', async () => {
449+
const kept = { type: 'sources', id: 'cited-sources', title: 'Sources' }
450+
dbChainMockFns.limit.mockResolvedValue([
451+
{ ...parentRow, resources: [{ ...kept, sources: { messageId: 'msg-3' } }] },
452+
])
453+
await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' }))
454+
expect(dbChainMockFns.values).toHaveBeenCalledWith(expect.objectContaining({ resources: [] }))
455+
456+
dbChainMockFns.values.mockClear()
457+
dbChainMockFns.limit.mockResolvedValue([
458+
{
459+
...parentRow,
460+
resources: [{ ...kept, sources: { messageId: 'live-id', requestId: 'req-2' } }],
461+
},
462+
])
463+
mockLoadCopilotChatMessages.mockResolvedValue([
464+
threeMessages[0],
465+
{ ...threeMessages[1], requestId: 'req-2' },
466+
threeMessages[2],
467+
])
468+
await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' }))
469+
expect(dbChainMockFns.values).toHaveBeenCalledWith(
470+
expect.objectContaining({
471+
resources: [{ ...kept, sources: { messageId: 'live-id', requestId: 'req-2' } }],
472+
})
473+
)
474+
})
475+
339476
it('copies pre-cut uploads and drops only post-cut ghosts', async () => {
340477
// The source chat owns two more uploads (apple pre-cut, banana post-cut)
341478
// beside the kept one, plus one shared workspace-file resource. The fork

‎apps/sim/app/api/mothership/chats/[chatId]/fork/route.ts‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ import { createLogger } from '@sim/logger'
22
import { type NextRequest, NextResponse } from 'next/server'
33
import { forkMothershipChatContract } from '@/lib/api/contracts/mothership-chats'
44
import { parseRequest } from '@/lib/api/server'
5-
import { asOrchestrationError } from '@/lib/core/orchestration/types'
5+
import { asOrchestrationError, statusForOrchestrationError } from '@/lib/core/orchestration/types'
66
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
77
import { forkChat } from '@/lib/mothership/chat/application/fork'
88
import {
@@ -38,6 +38,11 @@ export const POST = withRouteHandler(
3838
if (classified?.code === 'not_found' || classified?.code === 'forbidden')
3939
return NextResponse.json({ error: 'Chat not found' }, { status: 404 })
4040
if (classified?.code === 'validation') return createBadRequestResponse(classified.message)
41+
if (classified?.code === 'conflict' || classified?.code === 'payload_too_large')
42+
return NextResponse.json(
43+
{ error: classified.message },
44+
{ status: statusForOrchestrationError(classified.code) }
45+
)
4146
logger.error('Error forking chat:', error)
4247
return createInternalServerErrorResponse('Failed to fork chat')
4348
}

‎apps/sim/app/workspace/[workspaceId]/components/message-actions/message-actions.tsx‎

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import {
1919
useCopyToClipboard,
2020
} from '@sim/emcn'
2121
import { useParams, useRouter } from 'next/navigation'
22+
import { isApiClientError } from '@/lib/api/client/errors'
2223
import { isLiveAssistantMessageId } from '@/lib/mothership/chat/live-message-id'
2324
import { organizationRoutes } from '@/lib/navigation/paths'
2425
import { useChatSurface } from '@/app/workspace/[workspaceId]/home/components/chat-surface-context'
@@ -40,6 +41,9 @@ interface MessageActionsProps {
4041
messageId?: string
4142
}
4243

44+
/** Fork refusals whose message tells the person what to do: the response is unfinished, or the chat is too long. */
45+
const FORK_REFUSAL_STATUSES = new Set([409, 413])
46+
4347
export const MessageActions = memo(function MessageActions({
4448
content,
4549
getCopyContent,
@@ -143,8 +147,12 @@ export const MessageActions = memo(function MessageActions({
143147
useFolderStore.getState().clearChatSelection()
144148
router.push(`/workspace/${params.workspaceId}/chat/${result.id}`)
145149
}
146-
} catch {
147-
toast.error('Failed to fork chat')
150+
} catch (error) {
151+
toast.error(
152+
isApiClientError(error) && FORK_REFUSAL_STATUSES.has(error.status)
153+
? error.message
154+
: 'Failed to fork chat'
155+
)
148156
}
149157
}
150158

Lines changed: 101 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,101 @@
1+
import { copilotChats, copilotMessages, workspaceFiles } from '@sim/db/schema'
2+
import { queueTableRows, resetDbChainMock } from '@sim/testing/mocks/database.mock'
3+
import { storageServiceMockFns } from '@sim/testing/mocks/storage-service.mock'
4+
import { uploadsMock, uploadsMockFns } from '@sim/testing/mocks/uploads.mock'
5+
import { beforeEach, describe, expect, it, vi } from 'vitest'
6+
7+
vi.mock('@/lib/uploads', () => uploadsMock)
8+
9+
import { prepareChatCleanup } from '@/lib/cleanup/chat-cleanup'
10+
import { inlineChatImageKey } from '@/lib/mothership/chat/inline-image-key'
11+
12+
const chatId = '3f0c2a52-8a43-4d4b-9b5f-0b3c7a1e2d10'
13+
14+
function attachment(key: string) {
15+
return { id: 'wf_x', key, filename: 'x.png', media_type: 'image/png', size: 1 }
16+
}
17+
18+
/**
19+
* Purges one deleted chat whose message rows are `messages`, while `remaining` are the message
20+
* rows other chats of the organization still hold; returns the deleted keys by context.
21+
*/
22+
async function purge(
23+
messages: Record<string, unknown>[],
24+
remaining: Record<string, unknown>[] = []
25+
) {
26+
queueTableRows(workspaceFiles, [])
27+
queueTableRows(
28+
copilotMessages,
29+
messages.map((content) => ({ chatId, content }))
30+
)
31+
const cleanup = await prepareChatCleanup([chatId], 'test')
32+
queueTableRows(copilotChats, [])
33+
queueTableRows(
34+
copilotMessages,
35+
remaining.map((content) => ({ content }))
36+
)
37+
await cleanup.execute()
38+
return Object.fromEntries(
39+
storageServiceMockFns.mockDeleteFiles.mock.calls.map(([keys, context]) => [context, keys])
40+
)
41+
}
42+
43+
describe('chat purge storage', () => {
44+
beforeEach(() => {
45+
resetDbChainMock()
46+
uploadsMockFns.mockIsUsingCloudStorage.mockReturnValue(true)
47+
storageServiceMockFns.mockDeleteFiles.mockReset()
48+
storageServiceMockFns.mockDeleteFiles.mockResolvedValue({ deleted: 1, failed: [] })
49+
})
50+
51+
it('never deletes a workspace attachment as copilot storage', async () => {
52+
// A fork carries its source's key when the file's copy failed or the file was deleted, and
53+
// the copilot bucket can be the workspace bucket, so a workspace key here is another chat's file.
54+
const deleted = await purge([
55+
{
56+
role: 'user',
57+
content: 'look',
58+
fileAttachments: [
59+
attachment('workspace/ws-1/1700-abc-shared.png'),
60+
attachment('copilot/1234/legacy.png'),
61+
],
62+
},
63+
])
64+
expect(deleted).toEqual({ copilot: ['copilot/1234/legacy.png'] })
65+
})
66+
67+
it('deletes an organization attachment only once no remaining chat references it', async () => {
68+
const shared = 'assistant/org-1/user-1/u1/shared.png'
69+
const own = 'assistant/org-1/user-1/u2/own.png'
70+
const message = {
71+
role: 'user',
72+
content: 'look',
73+
fileAttachments: [attachment(shared), attachment(own)],
74+
}
75+
// A fork of this chat still holds the shared upload.
76+
const deleted = await purge(
77+
[message],
78+
[{ role: 'user', content: 'fork', fileAttachments: [attachment(shared)] }]
79+
)
80+
expect(deleted).toEqual({ mothership: [own] })
81+
})
82+
83+
it('deletes the inline images an assistant message published', async () => {
84+
const deleted = await purge([
85+
{
86+
role: 'assistant',
87+
requestId: 'req-1',
88+
content: 'Here: ![chart](files/chart.png) and `![code](files/no.png)`',
89+
contentBlocks: [{ type: 'text', content: '![other](/tmp/out.png)' }],
90+
},
91+
{ role: 'user', requestId: 'req-2', content: '![u](files/u.png)' },
92+
{ role: 'assistant', content: '![x](files/x.png)' },
93+
])
94+
expect(deleted).toEqual({
95+
mothership: [
96+
inlineChatImageKey(chatId, 'req-1', 'files/chart.png'),
97+
inlineChatImageKey(chatId, 'req-1', '/tmp/out.png'),
98+
],
99+
})
100+
})
101+
})

0 commit comments

Comments
 (0)