Skip to content

Commit 6adf4c6

Browse files
authored
fix(mothership): close the remaining ways a chat send is lost, resent hot, or polled after unmount (#8695)
* fix(mothership): close the remaining ways a chat send is lost, resent hot, or polled after unmount - A deleted chat's queue could come back through `enqueue`: a direct send to a chat deleted while its POST was failing (unreachable or busy) was re-queued, invisible, and would go out if the chat were restored. A cleared chat now refuses every queue write until it is restored from Recently Deleted. - A busy refusal re-queued its message for immediate redispatch. With Redis down every send is refused as busy without naming a turn, so the message was resent on every refusal (or stalled in the queue). Busy re-queues now wait out a jittered, growing delay (`backoffWithJitter`, 1s to 30s) before the drain sends them again under the same id. - A send waiting on the server's chat lock behind another tab's turn was aborted by a return-triggered recovery, and its message was lost. Recovery now leaves any POST that has not answered alone unless the chat already holds its message. When the busy refusal then arrives, the chat attaches to the turn holding it explicitly: the history read it relies on can repeat one the surface skipped while the POST was pending, so nothing else re-ran to attach, and the queued message waited behind a turn the tab never showed. - A retry told "already sent" while the server's earlier attempt with that id was still in flight (it had opened no stream yet) reattached to the missing stream, read the 404 as a finished turn, and dropped the message. Seen in the browser harness when Redis came back during a busy-refusal retry. Such a refusal is now retried later like a busy one. - A Send-now message whose Stop handoff failed (for example because the user switched chats while the Stop settled) returned a plain failure that the dispatch epoch check discarded. It is now handed back held, so it stays in its own chat's queue for the user to send; a surface that unmounted keeps resuming it from its stored handoff as before. - The saved-turn re-read kept refetching an unobserved chat for up to two minutes after the chat view unmounted. It now stops with the view. * fix(mothership): keep deleted chats empty across tabs and back off busy first messages - migrate no longer moves a new-chat queue into a chat deleted meanwhile - another tab's delete clears and tombstones the chat's queue; a restore (published as created) or the chat loading lifts it - a busy refusal on the new-chat surface waits out the backoff in its own queue instead of being handed to the surface's listener and resent at once - a deduplicated send whose stream check resolves after the user moved on leaves Stop on the new turn * fix(mothership): lift a chat's delete only when a later server read returns it The chat history effect also ran on history written locally, such as the rollback after a busy refusal, so a chat deleted in another tab could take queued sends again before any restore. The delete is now lifted only by a server read that began after it. * fix(mothership): apply the restore check to every server read of a chat Recovery on a return to the tab reads the chat directly, not through the history query, so a restore it found left the chat refusing follow-ups. The check now lives in fetchMothershipChatHistory itself. * test(mothership): use the central request mock in the restore read test * fix(mothership): lift only the chat delete a read or restore saw Each delete now gets a token. A history read or a restore lifts the delete it saw when it began, so a slow answer cannot reopen a chat deleted again while it was in flight. * fix(mothership): keep a send refused as busy after the user switched chats A chat switch detaches the view but leaves a send waiting on the chat lock running. Its busy refusal was dropped as a stale answer, so the message was lost. A conflict is now handled after a switch too, re-queuing the message in its own chat. The busy branch attaches through recovery's own history read instead of reading the chat twice. * fix(mothership): check a deduplicated send's stream before adopting its chat On the new-chat surface, an "already sent" answer names the chat the earlier attempt opened. Adopting it before finding no stream moved the surface to that chat while the retried message was queued under the new-chat key, where nothing sent it. The stream is now checked first. * fix(mothership): retry a deduplicated send whose stream lookup failed A lookup that failed with a network error or 5xx was read as proof that the earlier attempt's stream existed, so a send that never started was reattached and finalized instead of re-sent. Only a lookup this send aborted skips the retry now; the server deduplicates the retry by id.
1 parent f98df1b commit 6adf4c6

9 files changed

Lines changed: 1105 additions & 54 deletions

File tree

‎apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx‎

Lines changed: 723 additions & 0 deletions
Large diffs are not rendered by default.

‎apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts‎

Lines changed: 147 additions & 40 deletions
Large diffs are not rendered by default.
Lines changed: 119 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,119 @@
1+
import {
2+
apiClientRequestMock,
3+
apiClientRequestMockFns,
4+
} from '@sim/testing/mocks/api-client-request.mock'
5+
import { reactQueryMock } from '@sim/testing/mocks/react-query.mock'
6+
import { beforeEach, describe, expect, it, vi } from 'vitest'
7+
8+
vi.mock('@/lib/api/client/request', () => apiClientRequestMock)
9+
vi.mock('@tanstack/react-query', () => reactQueryMock)
10+
11+
import {
12+
fetchMothershipChatHistory,
13+
type MothershipChatHistory,
14+
useRestoreMothershipChat,
15+
} from '@/hooks/queries/mothership-chats'
16+
import { useMothershipQueueStore } from '@/stores/mothership-queue/store'
17+
18+
const mockRequestJson = apiClientRequestMockFns.mockRequestJson
19+
20+
const history: MothershipChatHistory = {
21+
id: 'chat-1',
22+
mode: 'agent',
23+
title: 'Restored',
24+
messages: [],
25+
activeStreamId: null,
26+
resources: [],
27+
}
28+
29+
/** Whether the queue store takes a send for the chat, i.e. whether its delete still holds. */
30+
function takesSends(chatId: string): boolean {
31+
useMothershipQueueStore.getState().enqueue(chatId, { id: 'probe', content: 'probe' })
32+
return useMothershipQueueStore.getState().queues[chatId] !== undefined
33+
}
34+
35+
/** A server answer the test releases when it chooses. */
36+
function deferredAnswer(value: unknown): () => void {
37+
let answer!: () => void
38+
mockRequestJson.mockReturnValue(
39+
new Promise((resolve) => {
40+
answer = () => resolve(value)
41+
})
42+
)
43+
return answer
44+
}
45+
46+
/** The options `useRestoreMothershipChat` hands to `useMutation` (the mock returns them). */
47+
interface RestoreMutation {
48+
mutationFn: (chatId: string) => Promise<void>
49+
onMutate?: (chatId: string) => { deleteSeen?: number }
50+
onSuccess: (data: undefined, chatId: string, context?: { deleteSeen?: number }) => void
51+
}
52+
53+
async function restore(chatId: string, whileInFlight: () => void = () => {}) {
54+
const mutation = useRestoreMothershipChat() as unknown as RestoreMutation
55+
const context = mutation.onMutate?.(chatId)
56+
const answer = deferredAnswer({ success: true })
57+
const done = mutation.mutationFn(chatId)
58+
whileInFlight()
59+
answer()
60+
await done
61+
mutation.onSuccess(undefined, chatId, context)
62+
}
63+
64+
describe('lifting a chat delete', () => {
65+
beforeEach(() => {
66+
useMothershipQueueStore.getState().reset()
67+
mockRequestJson.mockReset()
68+
})
69+
70+
it('reopens a chat this tab saw deleted once the server returns it again', async () => {
71+
useMothershipQueueStore.getState().clearChat(history.id)
72+
mockRequestJson.mockResolvedValue({ chat: history })
73+
74+
await fetchMothershipChatHistory(history.id)
75+
76+
expect(takesSends(history.id)).toBe(true)
77+
})
78+
79+
it('keeps the delete when the read returning the chat began before it', async () => {
80+
const answer = deferredAnswer({ chat: history })
81+
const read = fetchMothershipChatHistory(history.id)
82+
useMothershipQueueStore.getState().clearChat(history.id)
83+
answer()
84+
await read
85+
86+
expect(takesSends(history.id)).toBe(false)
87+
})
88+
89+
it('keeps a newer delete that lands while a read after an earlier one is in flight', async () => {
90+
useMothershipQueueStore.getState().clearChat(history.id)
91+
const answer = deferredAnswer({ chat: history })
92+
const read = fetchMothershipChatHistory(history.id)
93+
useMothershipQueueStore.getState().reopenChat(history.id)
94+
useMothershipQueueStore.getState().clearChat(history.id)
95+
answer()
96+
await read
97+
98+
expect(takesSends(history.id)).toBe(false)
99+
})
100+
101+
it('reopens a chat restored from Recently Deleted', async () => {
102+
useMothershipQueueStore.getState().clearChat(history.id)
103+
104+
await restore(history.id)
105+
106+
expect(takesSends(history.id)).toBe(true)
107+
})
108+
109+
it('keeps a delete that lands while the restore is in flight', async () => {
110+
useMothershipQueueStore.getState().clearChat(history.id)
111+
112+
await restore(history.id, () => {
113+
useMothershipQueueStore.getState().reopenChat(history.id)
114+
useMothershipQueueStore.getState().clearChat(history.id)
115+
})
116+
117+
expect(takesSends(history.id)).toBe(false)
118+
})
119+
})

‎apps/sim/hooks/queries/mothership-chats.ts‎

Lines changed: 23 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -281,7 +281,7 @@ export function useOrganizationMothershipChats(
281281
})
282282
}
283283

284-
export async function fetchMothershipChatHistory(
284+
async function readMothershipChatHistory(
285285
chatId: string,
286286
signal?: AbortSignal
287287
): Promise<MothershipChatHistory> {
@@ -309,6 +309,22 @@ export async function fetchMothershipChatHistory(
309309
return parseChatHistory(await copilotRes.json())
310310
}
311311

312+
/**
313+
* Reads a chat from the server. A chat this tab saw deleted that the server
314+
* returns again was restored, so it takes queued sends again. Only a read that
315+
* began after the delete counts: one already in flight can return the chat from
316+
* before it.
317+
*/
318+
export async function fetchMothershipChatHistory(
319+
chatId: string,
320+
signal?: AbortSignal
321+
): Promise<MothershipChatHistory> {
322+
const deleteSeen = useMothershipQueueStore.getState().cleared[chatId]
323+
const history = await readMothershipChatHistory(chatId, signal)
324+
if (deleteSeen !== undefined) useMothershipQueueStore.getState().reopenChat(chatId, deleteSeen)
325+
return history
326+
}
327+
312328
export function mothershipChatHistoryQueryOptions(chatId: string | undefined) {
313329
return queryOptions({
314330
queryKey: mothershipChatKeys.detail(chatId),
@@ -365,6 +381,12 @@ export function useRestoreMothershipChat(owner?: MothershipChatOwner) {
365381
const queryClient = useQueryClient()
366382
return useMutation({
367383
mutationFn: restoreChat,
384+
/** The delete this restore undoes; one that lands while it is in flight stays. */
385+
onMutate: (chatId) => ({ deleteSeen: useMothershipQueueStore.getState().cleared[chatId] }),
386+
onSuccess: (_data, chatId, context) => {
387+
if (context?.deleteSeen === undefined) return
388+
useMothershipQueueStore.getState().reopenChat(chatId, context.deleteSeen)
389+
},
368390
onSettled: () => {
369391
queryClient.invalidateQueries({ queryKey: mothershipChatKeys.ownerLists(owner) })
370392
},

‎apps/sim/hooks/use-mothership-chat-events.test.ts‎

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import {
1515
handleMothershipChatStatusEvent,
1616
resyncMothershipChatCaches,
1717
} from '@/hooks/use-mothership-chat-events'
18+
import { useMothershipQueueStore } from '@/stores/mothership-queue/store'
1819

1920
describe('handleMothershipChatStatusEvent', () => {
2021
const queryClient = {
@@ -136,6 +137,27 @@ describe('handleMothershipChatStatusEvent', () => {
136137
expect(suspendTerminalScope).toHaveBeenCalledWith('chat-1')
137138
})
138139

140+
it('drops the queue of a chat deleted elsewhere and takes sends again once it is restored', () => {
141+
useMothershipQueueStore.getState().reset()
142+
const queued = { id: 'm1', content: 'follow-up' }
143+
const publish = (type: 'deleted' | 'created') =>
144+
handleMothershipChatStatusEvent(
145+
queryClient,
146+
'ws-1',
147+
JSON.stringify({ chatId: 'chat-1', type, timestamp: Date.now() })
148+
)
149+
150+
useMothershipQueueStore.getState().enqueue('chat-1', queued)
151+
publish('deleted')
152+
expect(useMothershipQueueStore.getState().queues['chat-1']).toBeUndefined()
153+
useMothershipQueueStore.getState().enqueue('chat-1', queued)
154+
expect(useMothershipQueueStore.getState().queues['chat-1']).toBeUndefined()
155+
156+
publish('created')
157+
useMothershipQueueStore.getState().enqueue('chat-1', queued)
158+
expect(useMothershipQueueStore.getState().queues['chat-1']?.map((m) => m.id)).toEqual(['m1'])
159+
})
160+
139161
it('keeps started task detail when a stale started stream is older than the active stream', () => {
140162
queryClient.getQueryData.mockReturnValue({
141163
id: 'chat-1',

‎apps/sim/hooks/use-mothership-chat-events.ts‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import {
1010
type MothershipChatOwner,
1111
mothershipChatKeys,
1212
} from '@/hooks/queries/mothership-chats'
13+
import { useMothershipQueueStore } from '@/stores/mothership-queue/store'
1314

1415
const logger = createLogger('MothershipChatEvents')
1516

@@ -119,8 +120,12 @@ export function handleMothershipChatStatusEvent(
119120
// mutation would leave pages and PTYs running indefinitely.
120121
void suspendDesktopChatScopes(payload.chatId)
121122
queryClient.removeQueries({ queryKey: mothershipChatKeys.detail(payload.chatId) })
123+
/** This tab's queue for the chat goes too, and no later send may bring it back. */
124+
useMothershipQueueStore.getState().clearChat(payload.chatId)
122125
return
123126
}
127+
/** A restore is published as `created`; the chat takes queued sends again. */
128+
if (payload.type === 'created') useMothershipQueueStore.getState().reopenChat(payload.chatId)
124129
if (payload.type === 'renamed') {
125130
/**
126131
* The lists invalidated above carry the title every surface renders; the

‎apps/sim/stores/mothership-queue/store.test.ts‎

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -88,5 +88,32 @@ describe('useMothershipQueueStore', () => {
8888
])
8989
expect(useMothershipQueueStore.getState().queues['pending::abc']).toBeUndefined()
9090
})
91+
92+
it('lifts only the delete a restore saw, never a later one', () => {
93+
useMothershipQueueStore.getState().clearChat('chat-X')
94+
const seen = useMothershipQueueStore.getState().cleared['chat-X']
95+
useMothershipQueueStore.getState().reopenChat('chat-X')
96+
useMothershipQueueStore.getState().clearChat('chat-X')
97+
98+
useMothershipQueueStore.getState().reopenChat('chat-X', seen)
99+
useMothershipQueueStore.getState().enqueue('chat-X', message('after-stale-restore'))
100+
expect(useMothershipQueueStore.getState().queues['chat-X']).toBeUndefined()
101+
102+
const latest = useMothershipQueueStore.getState().cleared['chat-X']
103+
useMothershipQueueStore.getState().reopenChat('chat-X', latest)
104+
useMothershipQueueStore.getState().enqueue('chat-X', message('after-restore'))
105+
expect(useMothershipQueueStore.getState().queues['chat-X']?.map((m) => m.id)).toEqual([
106+
'after-restore',
107+
])
108+
})
109+
110+
it('does not move a new chat surface queue into a chat deleted meanwhile', () => {
111+
useMothershipQueueStore.getState().enqueue('pending::abc', message('pending-1'))
112+
useMothershipQueueStore.getState().clearChat('chat-X')
113+
useMothershipQueueStore.getState().migrate('pending::abc', 'chat-X')
114+
const state = useMothershipQueueStore.getState()
115+
expect(state.queues['chat-X']).toBeUndefined()
116+
expect(state.queues['pending::abc']).toBeUndefined()
117+
})
91118
})
92119
})

‎apps/sim/stores/mothership-queue/store.ts‎

Lines changed: 26 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -41,10 +41,13 @@ const sessionStorageAdapter = {
4141
},
4242
}
4343

44+
/** Numbers each delete, so a restore or read can tell the delete it saw from a later one. */
45+
let deleteCount = 0
46+
4447
const initialState = {
4548
queues: {} as Record<string, QueuedMothershipMessage[]>,
4649
editing: {} as Record<string, string>,
47-
cleared: {} as Record<string, true>,
50+
cleared: {} as Record<string, number>,
4851
}
4952

5053
const omitKey = <V>(record: Record<string, V>, key: string): Record<string, V> => {
@@ -67,13 +70,15 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
6770
...initialState,
6871

6972
enqueue: (chatKey, message) =>
70-
set((state) => ({
71-
cleared: omitKey(state.cleared, chatKey),
72-
queues: setQueueForChat(state.queues, chatKey, [
73-
...(state.queues[chatKey] ?? []),
74-
message,
75-
]),
76-
})),
73+
set((state) => {
74+
if (state.cleared[chatKey]) return state
75+
return {
76+
queues: setQueueForChat(state.queues, chatKey, [
77+
...(state.queues[chatKey] ?? []),
78+
message,
79+
]),
80+
}
81+
}),
7782

7883
insertAt: (chatKey, index, message) =>
7984
set((state) => {
@@ -99,6 +104,8 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
99104
retryRequired: _retry,
100105
heldUntilOnline: _held,
101106
heldSurface: _surface,
107+
busyRetries: _busyRetries,
108+
notBefore: _notBefore,
102109
...rest
103110
} = next[index]
104111
next[index] = {
@@ -144,7 +151,8 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
144151
if (!fromQueue && fromEditing === undefined) return state
145152

146153
const queues = omitKey(state.queues, fromKey)
147-
if (fromQueue && fromQueue.length > 0) {
154+
/** A chat deleted meanwhile takes nothing: its queue is gone with it. */
155+
if (fromQueue && fromQueue.length > 0 && !state.cleared[toKey]) {
148156
// Merge defensively in case a stale bucket survived in
149157
// sessionStorage. FIFO: existing first, then the resolved stream.
150158
const existing = state.queues[toKey] ?? []
@@ -201,9 +209,17 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
201209
set((state) => ({
202210
queues: omitKey(state.queues, chatKey),
203211
editing: omitKey(state.editing, chatKey),
204-
cleared: { ...state.cleared, [chatKey]: true },
212+
cleared: { ...state.cleared, [chatKey]: ++deleteCount },
205213
})),
206214

215+
reopenChat: (chatKey, deleteToken) =>
216+
set((state) => {
217+
const current = state.cleared[chatKey]
218+
if (current === undefined) return state
219+
if (deleteToken !== undefined && deleteToken !== current) return state
220+
return { cleared: omitKey(state.cleared, chatKey) }
221+
}),
222+
207223
reset: () => set(initialState),
208224
}),
209225
{

‎apps/sim/stores/mothership-queue/types.ts‎

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,10 @@ export type QueuedMothershipMessage = QueuedMessage & {
2424
* mount: the next chatless surface for the same owner and workflow adopts it.
2525
*/
2626
heldSurface?: string
27+
/** Busy refusals so far; paces the next retry. */
28+
busyRetries?: number
29+
/** Epoch ms before which a busy-refused message is not sent again. */
30+
notBefore?: number
2731
/**
2832
* Message id of a prior attempt at this send that an unmount cleanup
2933
* withdrew. Reused when the entry is dispatched so the server deduplicates
@@ -49,10 +53,11 @@ export interface MothershipQueueState {
4953
queues: Record<string, QueuedMothershipMessage[]>
5054
editing: Record<string, string>
5155
/**
52-
* Chats cleared this session (deleted). A late restore of a send dispatched
53-
* before the clear does not recreate their queue; a new enqueue lifts it.
56+
* Chats cleared this session (deleted), each with the token of its latest
57+
* delete. No write recreates their queue (a late restore, or a failed send
58+
* handed back); restoring the chat lifts it.
5459
*/
55-
cleared: Record<string, true>
60+
cleared: Record<string, number>
5661

5762
enqueue: (chatKey: string, message: QueuedMothershipMessage) => void
5863
insertAt: (chatKey: string, index: number, message: QueuedMothershipMessage) => void
@@ -65,5 +70,10 @@ export interface MothershipQueueState {
6570
/** Moves the sends a dead chatless mount of `surface` held onto `toKey`. */
6671
adoptHeldSends: (toKey: string, surface: string) => void
6772
clearChat: (chatKey: string) => void
73+
/**
74+
* Lifts `cleared` for a restored chat. Given the delete token an operation
75+
* saw when it began, lifts only that delete, never one that landed after it.
76+
*/
77+
reopenChat: (chatKey: string, deleteToken?: number) => void
6878
reset: () => void
6979
}

0 commit comments

Comments
 (0)