Skip to content
Closed
Show file tree
Hide file tree
Changes from 1 commit
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
10 changes: 10 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,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,
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
34 changes: 33 additions & 1 deletion apps/sim/lib/mothership/chat/fork-worker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'

Expand Down Expand Up @@ -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)
Comment thread
waleedlatif1 marked this conversation as resolved.
Outdated
})

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)
})
62 changes: 48 additions & 14 deletions apps/sim/lib/mothership/chat/fork-worker.ts
Original file line number Diff line number Diff line change
@@ -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<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 {
await response.body?.cancel()
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
Outdated
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
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
Outdated
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