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/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/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 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..ece4e3e1965 100644 --- a/apps/sim/lib/mothership/chat/fork-worker.ts +++ b/apps/sim/lib/mothership/chat/fork-worker.ts @@ -1,35 +1,73 @@ +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 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. */ +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 } }