Skip to content

Commit 351b8b9

Browse files
authored
fix(mothership): keep a chatless surface's queue across remounts and contain hung chat hook tests (#8707)
* improvement(mothership): contain hung chat hook tests and name the delete-lifting actions - bound async act in the use-chat DOM tests, so a send that never settles fails its own test instead of every later one - note at the deduped stream lookup that a timeout abort is retried on purpose - split reopenChat into liftDelete (token required) and reopenRestoredChat * fix(mothership): keep a chatless surface's queue when it remounts A follow-up queued on the new-chat surface waits until the first message's chat is known. If the surface remounted first (for example after that message's Stop failed), the queue stayed under the dead mount's key, neither sent nor shown. The unmount now marks that queue held for its surface, and the next mount of the surface adopts it like other held sends. * fix(mothership): keep a withdrawn first message ahead of its follow-ups on remount Holding a dying chatless queue for the next mount let that mount send a follow-up before the first message's cross-surface handoff arrived, so the two went out in reverse order. When follow-ups are queued, the unmount now puts the withdrawn first message at their head, under its own id, and skips its separate handoff; both reach the same surface in the order written. A first message after a Stop is not withdrawn, so only the follow-up goes out.
1 parent da993c5 commit 351b8b9

8 files changed

Lines changed: 266 additions & 22 deletions

File tree

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

Lines changed: 144 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@
1919
* request it aborted was accepted.
2020
*/
2121

22-
import { act, type ReactNode, StrictMode, useEffect, useState } from 'react'
22+
import { type ReactNode, act as reactAct, StrictMode, useEffect, useState } from 'react'
2323
import { authClientMock, authClientMockFns } from '@sim/testing/mocks/auth-client.mock'
2424
import { libDesktopMock, libDesktopMockFns } from '@sim/testing/mocks/lib-desktop.mock'
2525
import { nextNavigationMock, nextNavigationMockFns } from '@sim/testing/mocks/next-navigation.mock'
@@ -103,6 +103,36 @@ import { useExecutionStore } from '@/stores/execution/store'
103103
import { useMothershipEffortStore } from '@/stores/mothership-effort/store'
104104
import { useMothershipQueueStore } from '@/stores/mothership-queue/store'
105105

106+
/** Captured before any test fakes timers, so the act budget below runs in real time. */
107+
const realSetTimeout = globalThis.setTimeout
108+
const realClearTimeout = globalThis.clearTimeout
109+
/** Well under the 10s test timeout: the scope must end before the runner abandons the test. */
110+
const ACT_BUDGET_MS = 6_000
111+
112+
/**
113+
* React's `act`, with async callbacks bounded. A callback that never settles (a
114+
* regressed send stuck reconnecting) used to run into the test timeout while
115+
* still inside React's act scope, and the renders of every later test queued
116+
* behind it ("Hook result is not ready"). Failing the act after a budget ends
117+
* the scope, so one regression is one red test. Sync callbacks stay synchronous.
118+
*/
119+
function act(callback: () => unknown): Promise<void> {
120+
return reactAct((): undefined | Promise<void> => {
121+
const result = callback()
122+
if (!(result instanceof Promise)) return undefined
123+
let budget: ReturnType<typeof setTimeout> | undefined
124+
const budgetSpent = new Promise<never>((_, reject) => {
125+
budget = realSetTimeout(
126+
() => reject(new Error(`act callback still pending after ${ACT_BUDGET_MS}ms`)),
127+
ACT_BUDGET_MS
128+
)
129+
})
130+
return Promise.race([result.then(() => undefined), budgetSpent]).finally(() =>
131+
realClearTimeout(budget)
132+
)
133+
})
134+
}
135+
106136
authClientMockFns.mockUseSession.mockImplementation(() => ({
107137
data: { user: { id: 'test-viewer' } },
108138
}))
@@ -2022,6 +2052,119 @@ describe('useChat remount send recovery', () => {
20222052
expect(state.abortBodies[0]).not.toHaveProperty('chatId')
20232053
})
20242054

2055+
/**
2056+
* The first POST on the new-chat surface never answers; later POSTs open a
2057+
* turn in the chat the first message created. The abort endpoint and the
2058+
* stream lookup fail, as they would for a Stop that cannot reach the server.
2059+
*/
2060+
function stubFirstPostPendingThenAdmitted() {
2061+
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
2062+
const url = String(input)
2063+
if (url === '/api/mothership/chat' && init?.method === 'POST') {
2064+
state.postBodies.push(JSON.parse(String(init.body)))
2065+
if (state.postBodies.length === 1) {
2066+
return new Promise<Response>((_, reject) => {
2067+
init.signal?.addEventListener('abort', () => reject(init.signal?.reason), {
2068+
once: true,
2069+
})
2070+
})
2071+
}
2072+
return new Response(
2073+
new ReadableStream<Uint8Array>({
2074+
start(controller) {
2075+
controller.close()
2076+
},
2077+
}),
2078+
{
2079+
status: 200,
2080+
headers: {
2081+
'Content-Type': 'text/event-stream',
2082+
'x-mothership-chat-id': DEDUPED_CHAT_ID,
2083+
},
2084+
}
2085+
)
2086+
}
2087+
if (url.includes('/api/copilot/chat/abort')) {
2088+
state.abortBodies.push(JSON.parse(String(init?.body)))
2089+
return Response.json({ error: 'Internal error' }, { status: 500 })
2090+
}
2091+
if (url.includes('/api/mothership/chat/stream') && state.postBodies.length === 1) {
2092+
return Response.json({ error: 'Internal error' }, { status: 500 })
2093+
}
2094+
return fetchStub(input, init)
2095+
})
2096+
}
2097+
2098+
/**
2099+
* A follow-up typed on the new-chat surface while the first message waits for
2100+
* the server is queued behind it. If the surface remounts then, the first
2101+
* message is withdrawn and both must reach the next mount in the order they
2102+
* were written: the first message (under its own id), then the follow-up.
2103+
*/
2104+
it('sends a withdrawn first message before its follow-up when the new-chat surface remounts', async () => {
2105+
stubFirstPostPendingThenAdmitted()
2106+
const first = renderHomeLikeSurface()
2107+
await act(async () => {
2108+
void first.getResult().sendMessage('inspect the workspace')
2109+
})
2110+
await waitFor(() => state.postBodies.length === 1)
2111+
await act(async () => {
2112+
void first.getResult().sendMessage('follow-up while admission pending')
2113+
})
2114+
await waitFor(() => allQueuedMessages().length === 1)
2115+
first.unmount()
2116+
2117+
const second = renderHomeLikeSurface()
2118+
await waitFor(() => state.postBodies.length >= 3, 4_000)
2119+
await act(async () => {
2120+
await sleep(300)
2121+
})
2122+
2123+
const afterRemount = state.postBodies.slice(1)
2124+
expect(afterRemount.map((body) => body.message)).toEqual([
2125+
'inspect the workspace',
2126+
'follow-up while admission pending',
2127+
])
2128+
expect(afterRemount[0].userMessageId).toBe(state.postBodies[0].userMessageId)
2129+
expect(afterRemount[1].chatId).toBe(DEDUPED_CHAT_ID)
2130+
expect(second.claimedByOwnListener()).toBe(0)
2131+
expect(allQueuedMessages()).toHaveLength(0)
2132+
})
2133+
2134+
/**
2135+
* After a Stop of the first message, only the follow-up was the user's
2136+
* intent: the Stop's POST is left to the server, nothing withdraws it, and the
2137+
* next mount sends just the follow-up, once.
2138+
*/
2139+
it('sends only the follow-up after a failed Stop when the new-chat surface remounts', async () => {
2140+
stubFirstPostPendingThenAdmitted()
2141+
const first = renderHomeLikeSurface()
2142+
await act(async () => {
2143+
void first.getResult().sendMessage('inspect the workspace')
2144+
})
2145+
await waitFor(() => state.postBodies.length === 1)
2146+
await act(async () => {
2147+
void first
2148+
.getResult()
2149+
.stopGeneration()
2150+
.catch(() => {})
2151+
void first.getResult().sendMessage('Sent while the Stop was failing')
2152+
await sleep(1_000)
2153+
})
2154+
first.unmount()
2155+
2156+
renderHomeLikeSurface()
2157+
await waitFor(() => state.postBodies.length >= 2, 4_000)
2158+
await act(async () => {
2159+
await sleep(500)
2160+
})
2161+
2162+
expect(state.postBodies.slice(1).map((body) => body.message)).toEqual([
2163+
'Sent while the Stop was failing',
2164+
])
2165+
expect(allQueuedMessages()).toHaveLength(0)
2166+
})
2167+
20252168
it('stopping a chat preserves an unrelated manual workflow execution', async () => {
20262169
const executionStore = useExecutionStore.getState()
20272170
executionStore.setIsExecuting('manual-workflow', true)

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

Lines changed: 57 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -268,6 +268,8 @@ interface PendingChatAdmission {
268268
chatKey: string
269269
controller: AbortController
270270
settled: Promise<string | undefined>
271+
/** What an unmount would withdraw, if it ran before the server answered. */
272+
send: WithdrawnSend
271273
}
272274

273275
/** A send an unmount cleanup withdrew, as handed to the next chat surface. */
@@ -771,6 +773,10 @@ export function useChat(
771773
const onlineEventsRef = useRef(0)
772774
/** Identifies this chatless surface across mounts, for the sends it holds. */
773775
const heldSendSurface = `${scopeKey}:${options?.workflowId ?? 'home'}`
776+
const heldSendSurfaceRef = useRef(heldSendSurface)
777+
heldSendSurfaceRef.current = heldSendSurface
778+
/** Withdrawn first messages the unmount queued itself, so their handoff is skipped. */
779+
const withdrawnHeldAtUnmountRef = useRef<Set<string> | null>(null)
774780
const onToolResultRef = useRef(options?.onToolResult)
775781
onToolResultRef.current = options?.onToolResult
776782
const onTitleUpdateRef = useRef(options?.onTitleUpdate)
@@ -3788,6 +3794,17 @@ export function useChat(
37883794
settled: new Promise((resolve) => {
37893795
resolveAdmission = resolve
37903796
}),
3797+
send: {
3798+
content: message,
3799+
userMessageId,
3800+
...(fileAttachments ? { fileAttachments } : {}),
3801+
...(contexts ? { contexts } : {}),
3802+
...(options?.requestMode ? { requestMode: options.requestMode } : {}),
3803+
...(options?.assistantSearch ? { assistantSearch: options.assistantSearch } : {}),
3804+
...(options?.assistantSearchLevel !== undefined
3805+
? { assistantSearchLevel: options.assistantSearchLevel }
3806+
: {}),
3807+
},
37913808
}
37923809
pendingChatAdmissionRef.current = admission
37933810
}
@@ -4031,9 +4048,10 @@ export function useChat(
40314048
message. Retry it later like a busy refusal; the server's claim settles.
40324049
This is checked before adopting the chat the answer names, so a retried
40334050
message stays under the key it was sent from. A lookup that fails for
4034-
another reason proves nothing either way, so it is retried too: the server
4035-
deduplicates the retry by id. Only a lookup this send aborted (Stop, or the
4036-
user moving on) is not retried. */
4051+
another reason proves nothing either way, so it is retried on purpose; that
4052+
includes the lookup's own timeout abort, which leaves this send's signal
4053+
untouched. The server deduplicates the retry by id. Only an abort of this
4054+
send itself (Stop, or the user moving on) is not retried. */
40374055
const dedupedStreamExists = await fetchStreamBatch(
40384056
conflictStreamId,
40394057
'0',
@@ -4249,6 +4267,8 @@ export function useChat(
42494267
*/
42504268
const handOffWithdrawnSend = useCallback(
42514269
(send: WithdrawnSend) => {
4270+
/** The unmount already queued it ahead of its follow-ups; see the unmount cleanup. */
4271+
if (withdrawnHeldAtUnmountRef.current?.delete(send.userMessageId)) return
42524272
if (
42534273
sendMothershipMessage(
42544274
send.content,
@@ -5311,6 +5331,40 @@ export function useChat(
53115331

53125332
useEffect(() => {
53135333
return () => {
5334+
/* A chatless mount's queue key dies with it, so messages still queued there
5335+
go to the next mount of this surface, as held sends do. A first message
5336+
this unmount withdraws (its POST not yet answered, and not stopped) goes
5337+
at their head: the follow-ups were written after it, and the next mount
5338+
would otherwise send them before its handoff arrives. Alone, it keeps
5339+
the usual cross-surface handoff. */
5340+
const deadKey = chatKeyRef.current
5341+
if (deadKey.startsWith(PENDING_CHAT_KEY_PREFIX)) {
5342+
const queueStore = useMothershipQueueStore.getState()
5343+
const withdrawing = pendingChatAdmissionRef.current
5344+
if (
5345+
withdrawing &&
5346+
withdrawing.chatKey === deadKey &&
5347+
abortControllerRef.current === withdrawing.controller &&
5348+
(queueStore.queues[deadKey]?.length ?? 0) > 0
5349+
) {
5350+
const { send } = withdrawing
5351+
queueStore.insertAt(deadKey, 0, {
5352+
id: generateId(),
5353+
content: send.content,
5354+
resumeUserMessageId: send.userMessageId,
5355+
...(send.fileAttachments ? { fileAttachments: send.fileAttachments } : {}),
5356+
...(send.contexts ? { contexts: send.contexts } : {}),
5357+
...(send.requestMode ? { requestMode: send.requestMode } : {}),
5358+
...(send.assistantSearch ? { assistantSearch: send.assistantSearch } : {}),
5359+
...(send.assistantSearchLevel !== undefined
5360+
? { assistantSearchLevel: send.assistantSearchLevel }
5361+
: {}),
5362+
})
5363+
withdrawnHeldAtUnmountRef.current ??= new Set()
5364+
withdrawnHeldAtUnmountRef.current.add(send.userMessageId)
5365+
}
5366+
queueStore.holdForSurface(deadKey, heldSendSurfaceRef.current)
5367+
}
53145368
cancelActiveStreamRecovery()
53155369
clearQueueDispatchState()
53165370
streamGenRef.current++

‎apps/sim/hooks/queries/mothership-chat-history-restore.test.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -90,7 +90,7 @@ describe('lifting a chat delete', () => {
9090
useMothershipQueueStore.getState().clearChat(history.id)
9191
const answer = deferredAnswer({ chat: history })
9292
const read = fetchMothershipChatHistory(history.id)
93-
useMothershipQueueStore.getState().reopenChat(history.id)
93+
useMothershipQueueStore.getState().reopenRestoredChat(history.id)
9494
useMothershipQueueStore.getState().clearChat(history.id)
9595
answer()
9696
await read
@@ -110,7 +110,7 @@ describe('lifting a chat delete', () => {
110110
useMothershipQueueStore.getState().clearChat(history.id)
111111

112112
await restore(history.id, () => {
113-
useMothershipQueueStore.getState().reopenChat(history.id)
113+
useMothershipQueueStore.getState().reopenRestoredChat(history.id)
114114
useMothershipQueueStore.getState().clearChat(history.id)
115115
})
116116

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -322,7 +322,7 @@ export async function fetchMothershipChatHistory(
322322
): Promise<MothershipChatHistory> {
323323
const deleteSeen = useMothershipQueueStore.getState().cleared[chatId]
324324
const history = await readMothershipChatHistory(chatId, signal)
325-
if (deleteSeen !== undefined) useMothershipQueueStore.getState().reopenChat(chatId, deleteSeen)
325+
if (deleteSeen !== undefined) useMothershipQueueStore.getState().liftDelete(chatId, deleteSeen)
326326
return history
327327
}
328328

@@ -386,7 +386,7 @@ export function useRestoreMothershipChat(owner?: MothershipChatOwner) {
386386
onMutate: (chatId) => ({ deleteSeen: useMothershipQueueStore.getState().cleared[chatId] }),
387387
onSuccess: (_data, chatId, context) => {
388388
if (context?.deleteSeen === undefined) return
389-
useMothershipQueueStore.getState().reopenChat(chatId, context.deleteSeen)
389+
useMothershipQueueStore.getState().liftDelete(chatId, context.deleteSeen)
390390
},
391391
onSettled: () => {
392392
queryClient.invalidateQueries({ queryKey: mothershipChatKeys.ownerLists(owner) })

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -130,7 +130,8 @@ export function handleMothershipChatStatusEvent(
130130
return
131131
}
132132
/** A restore is published as `created`; the chat takes queued sends again. */
133-
if (payload.type === 'created') useMothershipQueueStore.getState().reopenChat(payload.chatId)
133+
if (payload.type === 'created')
134+
useMothershipQueueStore.getState().reopenRestoredChat(payload.chatId)
134135
if (payload.type === 'renamed') {
135136
/**
136137
* The lists invalidated above carry the title every surface renders; the

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

Lines changed: 18 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,21 @@ describe('useMothershipQueueStore', () => {
7575
})
7676
})
7777

78+
describe('holdForSurface', () => {
79+
it("hands a dead chatless mount's queue to the next mount of its surface only", () => {
80+
useMothershipQueueStore.getState().enqueue('pending::dead', message('m1'))
81+
useMothershipQueueStore.getState().holdForSurface('pending::dead', 'ws-1:home')
82+
83+
useMothershipQueueStore.getState().adoptHeldSends('pending::other', 'ws-1:workflow-1')
84+
expect(useMothershipQueueStore.getState().queues['pending::other']).toBeUndefined()
85+
86+
useMothershipQueueStore.getState().adoptHeldSends('pending::next', 'ws-1:home')
87+
const state = useMothershipQueueStore.getState()
88+
expect(state.queues['pending::next']?.map((m) => m.id)).toEqual(['m1'])
89+
expect(state.queues['pending::dead']).toBeUndefined()
90+
})
91+
})
92+
7893
describe('migrate', () => {
7994
it('merges into an existing destination bucket instead of overwriting', () => {
8095
useMothershipQueueStore.getState().enqueue('chat-X', message('existing-1'))
@@ -92,15 +107,15 @@ describe('useMothershipQueueStore', () => {
92107
it('lifts only the delete a restore saw, never a later one', () => {
93108
useMothershipQueueStore.getState().clearChat('chat-X')
94109
const seen = useMothershipQueueStore.getState().cleared['chat-X']
95-
useMothershipQueueStore.getState().reopenChat('chat-X')
110+
useMothershipQueueStore.getState().reopenRestoredChat('chat-X')
96111
useMothershipQueueStore.getState().clearChat('chat-X')
97112

98-
useMothershipQueueStore.getState().reopenChat('chat-X', seen)
113+
useMothershipQueueStore.getState().liftDelete('chat-X', seen)
99114
useMothershipQueueStore.getState().enqueue('chat-X', message('after-stale-restore'))
100115
expect(useMothershipQueueStore.getState().queues['chat-X']).toBeUndefined()
101116

102117
const latest = useMothershipQueueStore.getState().cleared['chat-X']
103-
useMothershipQueueStore.getState().reopenChat('chat-X', latest)
118+
useMothershipQueueStore.getState().liftDelete('chat-X', latest)
104119
useMothershipQueueStore.getState().enqueue('chat-X', message('after-restore'))
105120
expect(useMothershipQueueStore.getState().queues['chat-X']?.map((m) => m.id)).toEqual([
106121
'after-restore',

0 commit comments

Comments
 (0)