Skip to content

Commit d504f2e

Browse files
committed
fix(mothership): discard the worker copy of a fork that is not published
When anything after the worker copy failed, most often the final transaction that publishes the chat, the worker kept a conversation Sim had no chat for. The fork now asks the worker to clean that chat up, through the same cleanup endpoint chat deletion uses. It is best effort and a no-op when the worker never committed.
1 parent d7c42c0 commit d504f2e

3 files changed

Lines changed: 67 additions & 2 deletions

File tree

‎apps/sim/app/api/mothership/chats/[chatId]/fork/route.test.ts‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -342,6 +342,19 @@ describe('POST /api/mothership/chats/[chatId]/fork', () => {
342342
expect(dbChainMockFns.delete).not.toHaveBeenCalled()
343343
})
344344

345+
it('discards the worker copy when the fork cannot be published', async () => {
346+
dbChainMockFns.transaction.mockRejectedValueOnce(new Error('connection lost'))
347+
const res = await POST(createRequest('chat-1'), createRouteContext({ chatId: 'chat-1' }))
348+
expect(res.status).toBe(500)
349+
const [fork, discard] = mockFetchGo.mock.calls
350+
expect(fork[0]).toBe('http://mothership.test/api/chats/fork')
351+
expect(discard[0]).toBe('http://mothership.test/api/tasks/cleanup')
352+
expect(JSON.parse(discard[1].body)).toEqual({
353+
chatIds: [JSON.parse(fork[1].body).newChatId],
354+
})
355+
expect(mockPublishStatusChanged).not.toHaveBeenCalled()
356+
})
357+
345358
it('copies pre-cut uploads and drops only post-cut ghosts', async () => {
346359
// The source chat owns two more uploads (apple pre-cut, banana post-cut)
347360
// beside the kept one, plus one shared workspace-file resource. The fork

‎apps/sim/lib/mothership/chat/application/fork.ts‎

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,10 @@ import {
1919
planChatFileCopies,
2020
} from '@/lib/mothership/chat/fork-chat-files'
2121
import { planForkInlineImages } from '@/lib/mothership/chat/fork-inline-images'
22-
import { copyWorkerConversation } from '@/lib/mothership/chat/fork-worker'
22+
import {
23+
copyWorkerConversation,
24+
discardWorkerConversation,
25+
} from '@/lib/mothership/chat/fork-worker'
2326
import { loadCopilotChatMessages } from '@/lib/mothership/chat/lifecycle'
2427
import { appendCopilotChatMessages } from '@/lib/mothership/chat/messages-store'
2528
import {
@@ -88,6 +91,8 @@ export const forkChat = defineAuthorizedChatUseCase({
8891
const { chatId, upToMessageId } = input
8992
let preparedBlobs: ChatBlobCopyTask[] = []
9093
let published = false
94+
const newId = generateId()
95+
let workerCopyRequested = false
9196
try {
9297
const messages = await loadCopilotChatMessages(chatId)
9398
const forkIdx = messages.findIndex((m) => m.id === upToMessageId)
@@ -111,7 +116,6 @@ export const forkChat = defineAuthorizedChatUseCase({
111116
/** The source chat's chat-owned file ids (no cut) — the "is this resource a ghost?" test set for the rewrite. */
112117
const chatOwnedFileIds = new Set(chatOwnedFiles.map((row) => row.id))
113118

114-
const newId = generateId()
115119
/** Strip a leading "Fork | " so titles don't stack prefixes when forking a forked chat. */
116120
const baseTitle = (parent.title ?? 'New chat').replace(/^Fork \| /, '')
117121
const title = `Fork | ${baseTitle}`
@@ -129,6 +133,7 @@ export const forkChat = defineAuthorizedChatUseCase({
129133
).filter((resource) => resource.type !== 'file' || !failedIds.has(resource.id))
130134
const cutUser = [...forkedMessages].reverse().find((message) => message.role === 'user')
131135
if (!cutUser) throw new Error('The fork has no user message')
136+
workerCopyRequested = true
132137
await copyWorkerConversation({
133138
sourceChatId: chatId,
134139
newChatId: newId,
@@ -198,6 +203,8 @@ export const forkChat = defineAuthorizedChatUseCase({
198203
...(failed > 0 ? { failedFileCopies: failed } : {}),
199204
}
200205
} catch (error) {
206+
if (!published && workerCopyRequested)
207+
await discardWorkerConversation({ newChatId: newId, userId })
201208
if (!published) {
202209
await mapWithConcurrency(preparedBlobs, 4, async (task) => {
203210
try {

‎apps/sim/lib/mothership/chat/fork-worker.ts‎

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,14 @@
1+
import { createLogger } from '@sim/logger'
2+
import { getErrorMessage } from '@sim/utils/errors'
13
import { isRetryableNetworkError } from '@/lib/core/errors/retryable-infrastructure'
24
import { OrchestrationError } from '@/lib/core/orchestration/types'
35
import { type ForkChatRequest, ForkChatResponse } from '@/lib/mothership/generated/protocol'
46
import { fetchGo } from '@/lib/mothership/request/go/fetch'
57
import { mothershipRequestHeaders } from '@/lib/mothership/request/headers'
68
import { getMothershipBaseURL } from '@/lib/mothership/server/agent-url'
79

10+
const logger = createLogger('ForkWorker')
11+
812
/**
913
* How long one attempt waits for the copy. The worker refuses a cut above 50,000 events
1014
* with 413, and measured locally it copies 15,000 events in about 2 s and 50,000 in about
@@ -16,6 +20,9 @@ import { getMothershipBaseURL } from '@/lib/mothership/server/agent-url'
1620
*/
1721
const ATTEMPT_TIMEOUT_MS = 120_000
1822

23+
/** Cleanup is a small archive-and-purge of one chat, never a copy. */
24+
const DISCARD_TIMEOUT_MS = 10_000
25+
1926
/** Gateway failures: the request may never have reached a worker, so one more attempt is safe. */
2027
const RETRYABLE_STATUSES = new Set([502, 503, 504])
2128

@@ -76,3 +83,41 @@ export async function copyWorkerConversation(request: ForkChatRequest): Promise<
7683
return
7784
}
7885
}
86+
87+
/**
88+
* Removes the worker's copy of a fork that will not be published, so the worker keeps no
89+
* conversation Sim has no chat for. Best effort: a failure is logged, never thrown, since
90+
* the caller is already reporting the fork's own failure.
91+
*
92+
* The worker rolls back a copy whose caller disconnected, so after a timed-out attempt
93+
* this is normally a no-op. A worker that commits in the instant between its last
94+
* disconnect check and its commit, or one deployed before that rollback existed, can
95+
* commit after this call runs; that conversation then stays on the worker, unreachable
96+
* from Sim.
97+
*/
98+
export async function discardWorkerConversation(
99+
request: Pick<ForkChatRequest, 'newChatId' | 'userId'>
100+
): Promise<void> {
101+
try {
102+
const baseURL = await getMothershipBaseURL({ userId: request.userId })
103+
const response = await fetchGo(`${baseURL}/api/tasks/cleanup`, {
104+
method: 'POST',
105+
headers: mothershipRequestHeaders(),
106+
body: JSON.stringify({ chatIds: [request.newChatId] }),
107+
signal: AbortSignal.timeout(DISCARD_TIMEOUT_MS),
108+
spanName: 'sim → worker /api/tasks/cleanup',
109+
operation: 'discard_fork',
110+
})
111+
await response.body?.cancel().catch(() => undefined)
112+
if (!response.ok)
113+
logger.warn('Worker refused to discard an unpublished fork', {
114+
chatId: request.newChatId,
115+
status: response.status,
116+
})
117+
} catch (error) {
118+
logger.warn('Failed to discard an unpublished fork on the worker', {
119+
chatId: request.newChatId,
120+
error: getErrorMessage(error),
121+
})
122+
}
123+
}

0 commit comments

Comments
 (0)