From 7f42068bfa939803127072f5b9bb13041791a720 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Fri, 2 Oct 2026 01:21:25 -0700 Subject: [PATCH 1/5] fix(mothership): pass typed fork refusals through and stop overlapping retries --- .../chats/[chatId]/fork/route.test.ts | 10 +++ .../mothership/chats/[chatId]/fork/route.ts | 7 ++- .../lib/mothership/chat/fork-worker.test.ts | 34 +++++++++- apps/sim/lib/mothership/chat/fork-worker.ts | 62 ++++++++++++++----- 4 files changed, 97 insertions(+), 16 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..352634cae26 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,16 @@ describe('POST /api/mothership/chats/[chatId]/fork', () => { expect(mockPublishStatusChanged).not.toHaveBeenCalled() }) + it.each([404, 409, 413])( + 'passes a worker %i refusal through without publishing a chat', + 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) + expect(dbChainMockFns.transaction).not.toHaveBeenCalled() + } + ) + 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..2dcff081e9d 100644 --- a/apps/sim/lib/mothership/chat/fork-worker.test.ts +++ b/apps/sim/lib/mothership/chat/fork-worker.test.ts @@ -8,6 +8,7 @@ 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' @@ -58,6 +59,37 @@ it.each(['missing-receipt', 'wrong-chat', 'unavailable'])( ) }) 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(fetchWorker).toHaveBeenCalledTimes(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 (status, code) => { + fetchWorker.mockResolvedValue(Response.json({ error: 'refused' }, { status })) + const failure = await copyWorkerConversation(request).catch((error: unknown) => error) + expect(asOrchestrationError(failure)?.code).toBe(code) + expect(fetchWorker).toHaveBeenCalledTimes(1) +}) + +it.each([500, 400])('does not repeat a fork the worker failed with %i', async (status) => { + fetchWorker.mockResolvedValue(new Response('', { status })) + const failure = await copyWorkerConversation(request).catch((error: unknown) => error) + expect(asOrchestrationError(failure)).toBeNull() + expect(fetchWorker).toHaveBeenCalledTimes(1) +}) + +it.each([502, 504])('retries a %i gateway failure once', async (status) => { + fetchWorker.mockResolvedValueOnce(new Response('', { status })) + await copyWorkerConversation(request) + expect(fetchWorker).toHaveBeenCalledTimes(2) +}) + +it('does not start a second copy while a timed-out one may still be running', async () => { + fetchWorker.mockRejectedValue(new DOMException('The operation timed out.', 'TimeoutError')) + await expect(copyWorkerConversation(request)).rejects.toThrow() + expect(fetchWorker).toHaveBeenCalledTimes(1) +}) diff --git a/apps/sim/lib/mothership/chat/fork-worker.ts b/apps/sim/lib/mothership/chat/fork-worker.ts index bda6145b4dd..c37e305840b 100644 --- a/apps/sim/lib/mothership/chat/fork-worker.ts +++ b/apps/sim/lib/mothership/chat/fork-worker.ts @@ -1,35 +1,69 @@ +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. */ +/** Above the worker's 30 s fork budget, so a slow copy finishes instead of racing its own retry. */ +const ATTEMPT_TIMEOUT_MS = 45_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 failed connection 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 { + await response.body?.cancel() + outcome = { kind: 'status', status: response.status } } - 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 + // fetch reports a refused or dropped connection, including mid-body, as a TypeError. + if (canRetry && error instanceof TypeError) 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 3f142fa40c62a0b9b3338f2b9f8e3e6ca6d66e26 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Fri, 2 Oct 2026 01:31:10 -0700 Subject: [PATCH 2/5] fix(mothership): classify a fork refusal even when its body fails to cancel --- .../chats/[chatId]/fork/route.test.ts | 14 ++- .../lib/mothership/chat/fork-worker.test.ts | 90 +++++++++++++------ apps/sim/lib/mothership/chat/fork-worker.ts | 3 +- 3 files changed, 71 insertions(+), 36 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 352634cae26..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,15 +316,11 @@ describe('POST /api/mothership/chats/[chatId]/fork', () => { expect(mockPublishStatusChanged).not.toHaveBeenCalled() }) - it.each([404, 409, 413])( - 'passes a worker %i refusal through without publishing a chat', - 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) - expect(dbChainMockFns.transaction).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({ diff --git a/apps/sim/lib/mothership/chat/fork-worker.test.ts b/apps/sim/lib/mothership/chat/fork-worker.test.ts index 2dcff081e9d..997d2ccadbe 100644 --- a/apps/sim/lib/mothership/chat/fork-worker.test.ts +++ b/apps/sim/lib/mothership/chat/fork-worker.test.ts @@ -15,8 +15,6 @@ 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(), @@ -28,39 +26,64 @@ const request: ForkChatRequest = { fileKeys: {}, } +type Answer = () => Promise + +/** 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(['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 new TypeError('Connection ended') + } + : 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 () => { + const reply: Answer = async () => { if (failure === 'unavailable') throw new TypeError('Worker unavailable') return Response.json( failure === 'wrong-chat' ? { chatId: generateId(), sourceThroughSeq: 7 } : { ok: true } ) - }) + } + answer(reply, reply) await expect(copyWorkerConversation(request)).rejects.toThrow() // Only the unreachable worker may never have seen the request; an answer is final. - expect(fetchWorker).toHaveBeenCalledTimes(failure === 'unavailable' ? 2 : 1) + expect(received).toHaveLength(failure === 'unavailable' ? 2 : 1) } ) @@ -68,28 +91,43 @@ it.each([ [404, 'not_found'], [409, 'conflict'], [413, 'payload_too_large'], -] as const)('classifies a worker %i refusal as %s without retrying it', async (status, code) => { - fetchWorker.mockResolvedValue(Response.json({ error: 'refused' }, { status })) +] 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(code) - expect(fetchWorker).toHaveBeenCalledTimes(1) + 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 (status) => { - fetchWorker.mockResolvedValue(new Response('', { status })) +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(fetchWorker).toHaveBeenCalledTimes(1) + expect(received).toHaveLength(1) }) -it.each([502, 504])('retries a %i gateway failure once', async (status) => { - fetchWorker.mockResolvedValueOnce(new Response('', { status })) +it.each([502, 504])('retries a %i gateway failure once', async (code) => { + answer(status(code)) await copyWorkerConversation(request) - expect(fetchWorker).toHaveBeenCalledTimes(2) + expect(received).toEqual([request, request]) }) it('does not start a second copy while a timed-out one may still be running', async () => { - fetchWorker.mockRejectedValue(new DOMException('The operation timed out.', 'TimeoutError')) + answer(async () => { + throw new DOMException('The operation timed out.', 'TimeoutError') + }, receipt) await expect(copyWorkerConversation(request)).rejects.toThrow() - expect(fetchWorker).toHaveBeenCalledTimes(1) + 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 c37e305840b..08cb37aaea2 100644 --- a/apps/sim/lib/mothership/chat/fork-worker.ts +++ b/apps/sim/lib/mothership/chat/fork-worker.ts @@ -50,8 +50,9 @@ export async function copyWorkerConversation(request: ForkChatRequest): Promise< }) if (response.ok) outcome = { kind: 'receipt', body: await response.json() } else { - await response.body?.cancel() 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) } } catch (error) { // fetch reports a refused or dropped connection, including mid-body, as a TypeError. From 0f0449de1acefc201a42972553a38781f6ded6aa Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Fri, 2 Oct 2026 01:37:05 -0700 Subject: [PATCH 3/5] fix(mothership): retry a fork only after a recognized socket failure --- .../sim/lib/mothership/chat/fork-worker.test.ts | 17 +++++++++++++++-- apps/sim/lib/mothership/chat/fork-worker.ts | 8 ++++---- 2 files changed, 19 insertions(+), 6 deletions(-) diff --git a/apps/sim/lib/mothership/chat/fork-worker.test.ts b/apps/sim/lib/mothership/chat/fork-worker.test.ts index 997d2ccadbe..635adfa409b 100644 --- a/apps/sim/lib/mothership/chat/fork-worker.test.ts +++ b/apps/sim/lib/mothership/chat/fork-worker.test.ts @@ -28,6 +28,11 @@ const request: ForkChatRequest = { 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[] = [] @@ -62,7 +67,7 @@ it.each(['lost-response', 'temporary-error'])( answer( failure === 'lost-response' ? async () => { - throw new TypeError('Connection ended') + throw socketFailure('ECONNRESET') } : status(503) ) @@ -75,7 +80,7 @@ it.each(['missing-receipt', 'wrong-chat', 'unavailable'])( 'refuses an unconfirmed copy: %s', async (failure) => { const reply: Answer = async () => { - if (failure === 'unavailable') throw new TypeError('Worker unavailable') + if (failure === 'unavailable') throw socketFailure('ECONNREFUSED') return Response.json( failure === 'wrong-chat' ? { chatId: generateId(), sourceThroughSeq: 7 } : { ok: true } ) @@ -131,3 +136,11 @@ it('does not start a second copy while a timed-out one may still be running', as 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 08cb37aaea2..259467c340d 100644 --- a/apps/sim/lib/mothership/chat/fork-worker.ts +++ b/apps/sim/lib/mothership/chat/fork-worker.ts @@ -1,3 +1,4 @@ +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' @@ -30,8 +31,8 @@ type Attempt = { kind: 'receipt'; body: unknown } | { kind: 'status'; status: nu /** * A lost copy acknowledgement retries the same destination and immutable request; the - * worker answers a repeat with the first attempt's receipt. Only a failed connection or a - * gateway status is retried: a timed-out attempt may still be copying. + * 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 }) @@ -55,8 +56,7 @@ export async function copyWorkerConversation(request: ForkChatRequest): Promise< await response.body?.cancel().catch(() => undefined) } } catch (error) { - // fetch reports a refused or dropped connection, including mid-body, as a TypeError. - if (canRetry && error instanceof TypeError) continue + if (canRetry && isRetryableNetworkError(error)) continue throw error } if (outcome.kind === 'status') { From 9d295c792fcdba6807eef018e2641603b8beb678 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Fri, 2 Oct 2026 01:43:13 -0700 Subject: [PATCH 4/5] fix(core): treat EHOSTUNREACH as a retryable network failure --- apps/sim/lib/core/errors/retryable-infrastructure.ts | 1 + apps/sim/lib/mothership/chat/fork-worker.test.ts | 11 +++++++++++ 2 files changed, 12 insertions(+) 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/fork-worker.test.ts b/apps/sim/lib/mothership/chat/fork-worker.test.ts index 635adfa409b..63bf3de8351 100644 --- a/apps/sim/lib/mothership/chat/fork-worker.test.ts +++ b/apps/sim/lib/mothership/chat/fork-worker.test.ts @@ -61,6 +61,17 @@ beforeEach(() => { ) }) +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) => { From b0ebe2c3bbd1084bf149d8114fe4c55fed40714a Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Fri, 2 Oct 2026 10:59:55 -0700 Subject: [PATCH 5/5] fix(mothership): show a fork refusal's own message The route passes a worker 409 or 413 through with a message the person can act on, but the fork action replaced every failure with "Failed to fork chat". It now shows the refusal's message for those two statuses. The per-attempt timeout's comment no longer claims it sits above a 30 s worker budget: that budget bounds each statement, not the copy. A timed-out attempt leaves nothing behind because the worker rolls back a copy whose caller has gone (mothership#575). --- .../components/message-actions/message-actions.tsx | 12 ++++++++++-- apps/sim/lib/mothership/chat/fork-worker.ts | 5 ++++- 2 files changed, 14 insertions(+), 3 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' + ) } } diff --git a/apps/sim/lib/mothership/chat/fork-worker.ts b/apps/sim/lib/mothership/chat/fork-worker.ts index 259467c340d..ece4e3e1965 100644 --- a/apps/sim/lib/mothership/chat/fork-worker.ts +++ b/apps/sim/lib/mothership/chat/fork-worker.ts @@ -5,7 +5,10 @@ import { fetchGo } from '@/lib/mothership/request/go/fetch' import { mothershipRequestHeaders } from '@/lib/mothership/request/headers' import { getMothershipBaseURL } from '@/lib/mothership/server/agent-url' -/** Above the worker's 30 s fork budget, so a slow copy finishes instead of racing its own retry. */ +/** + * How long one attempt waits for the copy. The worker rolls a copy back once its caller's + * connection closes, so an attempt that times out leaves no conversation behind. + */ const ATTEMPT_TIMEOUT_MS = 45_000 /** Gateway failures: the request may never have reached a worker, so one more attempt is safe. */