Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions apps/sim/app/api/mothership/chats/[chatId]/fork/route.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
7 changes: 6 additions & 1 deletion apps/sim/app/api/mothership/chats/[chatId]/fork/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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')
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand All @@ -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,
Expand Down Expand Up @@ -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'
)
}
}

Expand Down
1 change: 1 addition & 0 deletions apps/sim/lib/core/errors/retryable-infrastructure.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ const RETRYABLE_NETWORK_ERROR_CODES = new Set([
'ECONNREFUSED',
'EPIPE',
'ERR_STREAM_PREMATURE_CLOSE',
'EHOSTUNREACH',
'ENETDOWN',
'ENETRESET',
'ENETUNREACH',
Expand Down
124 changes: 109 additions & 15 deletions apps/sim/lib/mothership/chat/fork-worker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand All @@ -27,37 +26,132 @@ const request: ForkChatRequest = {
fileKeys: {},
}

type Answer = () => Promise<Response>

/** 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)
})
66 changes: 52 additions & 14 deletions apps/sim/lib/mothership/chat/fork-worker.ts
Original file line number Diff line number Diff line change
@@ -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<void> {
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
Comment thread
waleedlatif1 marked this conversation as resolved.
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
}
}
Loading