Skip to content
Closed
Show file tree
Hide file tree
Changes from 4 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
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)
})
63 changes: 49 additions & 14 deletions apps/sim/lib/mothership/chat/fork-worker.ts
Original file line number Diff line number Diff line change
@@ -1,35 +1,70 @@
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. */
/** 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 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