Skip to content

Commit 9d3a096

Browse files
committed
fix(mothership): keep held and busy-refused sends exactly once across remounts
- A first message held offline on the new-chat page sat under that mount's queue key, which dies with the mount, so a reload or remount before the network returned stranded it. Held sends on a chatless surface now carry the surface they belong to, and the next chatless mount of that surface adopts them. - Held sends are released for every chat when the browser comes back online, and on mount when it already is, so a send held in a chat the user is not viewing (or one whose `online` event fired with no surface mounted) still goes out. - A send refused because the chat is busy is handed back only after the chat's running turn has been read, so the queue cannot redispatch it before that turn ends. A busy refusal that does not name the running turn no longer reads as a deduplicated send, which reconnected to a stream that never existed and lost the message.
1 parent e7e28ba commit 9d3a096

4 files changed

Lines changed: 231 additions & 57 deletions

File tree

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

Lines changed: 125 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -1935,14 +1935,18 @@ describe('useChat remount send recovery', () => {
19351935
* The POST fails at the network layer until the network is back, and the
19361936
* stream it would have opened does not exist.
19371937
*/
1938-
const network = { online: false }
1938+
const network = { online: false, acceptedPosts: 0 }
19391939
function stubUnreachableSend() {
19401940
network.online = false
1941+
network.acceptedPosts = 0
19411942
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
19421943
const url = String(input)
19431944
if (url === '/api/mothership/chat' && init?.method === 'POST') {
19441945
state.postBodies.push(JSON.parse(String(init.body)))
1945-
if (network.online) return emptySseResponse()
1946+
if (network.online) {
1947+
network.acceptedPosts++
1948+
return emptySseResponse()
1949+
}
19461950
throw new TypeError('Failed to fetch')
19471951
}
19481952
if (url.includes('/api/mothership/chat/stream')) {
@@ -2002,33 +2006,133 @@ describe('useChat remount send recovery', () => {
20022006

20032007
/**
20042008
* Another tab's turn holds the chat (this one missed its start). The server
2005-
* refuses the send naming that turn; the message must wait for it rather than
2006-
* render that turn's answer under itself and then vanish.
2009+
* refuses the send, naming that turn or, when its stream id is unreadable,
2010+
* nothing. The message must go out exactly once, under its id, after that turn
2011+
* ends: never rendered under the other turn's answer, never lost, never resent
2012+
* while the turn still runs.
20072013
*/
2008-
it('sends a message again after the turn that held the chat ends', async () => {
2009-
const history = idleHistory('chat-busy-elsewhere')
2010-
mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history }))
2011-
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
2012-
if (String(input) === '/api/mothership/chat' && init?.method === 'POST') {
2013-
state.postBodies.push(JSON.parse(String(init.body)))
2014-
if (state.postBodies.length === 1) {
2015-
return Response.json(
2016-
{ error: 'A response is already running', activeStreamId: 'turn-from-another-tab' },
2017-
{ status: 409 }
2018-
)
2014+
it.each([
2015+
['names', 'turn-from-another-tab'],
2016+
['does not name', undefined],
2017+
] as const)(
2018+
'sends a message once, after the turn that held the chat ends, when the refusal %s it',
2019+
async (_names, refusalStreamId) => {
2020+
const history = idleHistory(`chat-busy-${refusalStreamId ?? 'unnamed'}`)
2021+
let otherTurnRunning = true
2022+
mockRequestJson.mockImplementation(() =>
2023+
Promise.resolve({
2024+
chat: {
2025+
...history,
2026+
activeStreamId: otherTurnRunning ? 'turn-from-another-tab' : null,
2027+
},
2028+
})
2029+
)
2030+
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
2031+
const url = String(input)
2032+
if (url === '/api/mothership/chat' && init?.method === 'POST') {
2033+
state.postBodies.push(JSON.parse(String(init.body)))
2034+
if (otherTurnRunning) {
2035+
return Response.json(
2036+
{
2037+
error: 'A response is already in progress for this chat.',
2038+
...(refusalStreamId ? { activeStreamId: refusalStreamId } : {}),
2039+
},
2040+
{ status: 409 }
2041+
)
2042+
}
2043+
return emptySseResponse()
20192044
}
2020-
return emptySseResponse()
2021-
}
2022-
return fetchStub(input, init)
2045+
if (url.includes('/api/mothership/chat/stream')) {
2046+
if (url.includes('batch=true')) {
2047+
return Response.json({
2048+
success: true,
2049+
events: [],
2050+
status: otherTurnRunning ? 'streaming' : 'complete',
2051+
})
2052+
}
2053+
return emptySseResponse()
2054+
}
2055+
return fetchStub(input, init)
2056+
})
2057+
const { getResult } = renderUseChatInChat(history.id, history)
2058+
2059+
await act(async () => {
2060+
await getResult().sendMessage('Sent from the second tab')
2061+
})
2062+
await act(async () => {
2063+
await sleep(1500)
2064+
})
2065+
expect(state.postBodies).toHaveLength(1)
2066+
expect(useMothershipQueueStore.getState().queues[history.id]?.[0]?.content).toBe(
2067+
'Sent from the second tab'
2068+
)
2069+
2070+
otherTurnRunning = false
2071+
await waitFor(() => state.postBodies.length === 2, 5000)
2072+
await act(async () => {
2073+
await sleep(500)
2074+
})
2075+
2076+
expect(state.postBodies).toHaveLength(2)
2077+
expect(state.postBodies[1].message).toBe('Sent from the second tab')
2078+
expect(state.postBodies[1].userMessageId).toBe(state.postBodies[0].userMessageId)
2079+
}
2080+
)
2081+
2082+
/**
2083+
* A chatless surface keys its queue by mount, so a held first message would
2084+
* be stranded by a reload or remount before the network returns. The next
2085+
* chatless mount of the same surface adopts it.
2086+
*/
2087+
it('carries a first message held offline over to the next new-chat surface', async () => {
2088+
stubUnreachableSend()
2089+
const first = renderUseChat()
2090+
await act(async () => {
2091+
await first.getResult().sendMessage('First message, sent offline')
20232092
})
2024-
const { getResult } = renderUseChatInChat(history.id, history)
2093+
await waitFor(() => allQueuedMessages().some((message) => message.retryRequired === true))
2094+
first.unmount()
2095+
2096+
const second = renderUseChat()
2097+
await waitFor(() =>
2098+
second
2099+
.getResult()
2100+
.messageQueue.some((message) => message.content === 'First message, sent offline')
2101+
)
20252102

2103+
network.online = true
20262104
await act(async () => {
2027-
await getResult().sendMessage('Sent from the second tab')
2105+
window.dispatchEvent(new Event('online'))
20282106
})
2107+
await waitFor(() => network.acceptedPosts === 1)
2108+
await act(async () => {
2109+
await sleep(300)
2110+
})
2111+
2112+
expect(network.acceptedPosts).toBe(1)
2113+
expect(new Set(state.postBodies.map((body) => body.userMessageId)).size).toBe(1)
2114+
expect(state.postBodies.at(-1)?.message).toBe('First message, sent offline')
2115+
})
2116+
2117+
/** The `online` event can fire while no surface for the chat is mounted. */
2118+
it('sends a held message when its chat mounts after the network came back', async () => {
2119+
const history = idleHistory('chat-held-while-away')
2120+
mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history }))
2121+
stubUnreachableSend()
2122+
const first = renderUseChatInChat(history.id, history)
2123+
await act(async () => {
2124+
await first.getResult().sendMessage('Held while I was elsewhere')
2125+
})
2126+
await waitFor(
2127+
() => useMothershipQueueStore.getState().queues[history.id]?.[0]?.retryRequired === true
2128+
)
2129+
first.unmount()
2130+
2131+
network.online = true
2132+
renderUseChatInChat(history.id, history)
20292133
await waitFor(() => state.postBodies.length === 2)
20302134

2031-
expect(state.postBodies[1].message).toBe('Sent from the second tab')
2135+
expect(state.postBodies[1].message).toBe('Held while I was elsewhere')
20322136
expect(state.postBodies[1].userMessageId).toBe(state.postBodies[0].userMessageId)
20332137
})
20342138
})

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

Lines changed: 55 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -729,6 +729,8 @@ export function useChat(
729729
const pendingStopModeRef = useRef<StopGenerationMode | null>(null)
730730
const workflowIdRef = useRef(options?.workflowId)
731731
workflowIdRef.current = options?.workflowId
732+
/** Identifies this chatless surface across mounts, for the sends it holds. */
733+
const heldSendSurface = `${scopeKey}:${options?.workflowId ?? 'home'}`
732734
const onToolResultRef = useRef(options?.onToolResult)
733735
onToolResultRef.current = options?.onToolResult
734736
const onTitleUpdateRef = useRef(options?.onTitleUpdate)
@@ -3834,10 +3836,10 @@ export function useChat(
38343836
if (!response.ok) {
38353837
const errorData = await response.json().catch(() => ({}))
38363838
if (response.status === 409) {
3839+
/* A deduplicated send always names itself; a chat busy with another
3840+
turn names that turn, or nothing when its stream id is unreadable. */
38373841
const conflictStreamId =
3838-
typeof errorData.activeStreamId === 'string'
3839-
? errorData.activeStreamId
3840-
: userMessageId
3842+
typeof errorData.activeStreamId === 'string' ? errorData.activeStreamId : undefined
38413843
const supersededStreamId = queuedSendHandoff?.supersededStreamId ?? pendingStopStreamId
38423844
if (supersededStreamId && conflictStreamId === supersededStreamId) {
38433845
rollbackOptimisticSend()
@@ -3853,8 +3855,11 @@ export function useChat(
38533855
}
38543856
if (conflictStreamId !== userMessageId) {
38553857
/* Another turn holds the chat: one started in another tab, or one this
3856-
surface lost track of. This message was not admitted, so it waits in
3857-
the queue behind that turn, which the chat now shows as running. */
3858+
surface lost track of. This message was not admitted (the server
3859+
released its id), so it goes back to the queue under the same id. The
3860+
queue drains only while the chat is idle, so the chat's running turn is
3861+
read before the message is handed back: the chat then attaches to that
3862+
turn and the message goes out once, after it ends. */
38583863
rollbackOptimisticSend()
38593864
if (streamGenRef.current === gen) {
38603865
streamGenRef.current++
@@ -3865,13 +3870,15 @@ export function useChat(
38653870
}
38663871
if (requestChatId) {
38673872
const busyChatId = requestChatId
3868-
upsertChatHistory(busyChatId, (current) => ({
3869-
...current,
3870-
activeStreamId: conflictStreamId,
3871-
}))
3872-
void queryClient.invalidateQueries({
3873-
queryKey: mothershipChatKeys.detail(busyChatId),
3874-
})
3873+
if (conflictStreamId) {
3874+
upsertChatHistory(busyChatId, (current) => ({
3875+
...current,
3876+
activeStreamId: conflictStreamId,
3877+
}))
3878+
}
3879+
await queryClient
3880+
.refetchQueries({ queryKey: mothershipChatKeys.detail(busyChatId), exact: true })
3881+
.catch(() => {})
38753882
}
38763883
return { userMessageId }
38773884
}
@@ -4203,14 +4210,23 @@ export function useChat(
42034210
options?.assistantSearch,
42044211
options?.assistantSearchLevel
42054212
),
4206-
...(result.unreachable ? { retryRequired: true, heldUntilOnline: true } : {}),
4213+
...(result.unreachable
4214+
? {
4215+
retryRequired: true,
4216+
heldUntilOnline: true,
4217+
...(activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)
4218+
? { heldSurface: heldSendSurface }
4219+
: {}),
4220+
}
4221+
: {}),
42074222
})
42084223
},
42094224
[
42104225
workspaceId,
42114226
createQueuedMessage,
42124227
startSendMessage,
42134228
handOffWithdrawnSend,
4229+
heldSendSurface,
42144230
hasPendingChatAdmission,
42154231
]
42164232
)
@@ -4810,7 +4826,14 @@ export function useChat(
48104826
...dispatched,
48114827
...(retainedHandoff ? { queuedSendHandoff: retainedHandoff } : {}),
48124828
retryRequired: !retriesOnItsOwn,
4813-
...(withdrawn?.unreachable ? { heldUntilOnline: true } : {}),
4829+
...(withdrawn?.unreachable
4830+
? {
4831+
heldUntilOnline: true,
4832+
...(dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)
4833+
? { heldSurface: heldSendSurface }
4834+
: {}),
4835+
}
4836+
: {}),
48144837
...(withdrawnUserMessageId ? { resumeUserMessageId: withdrawnUserMessageId } : {}),
48154838
})
48164839
}
@@ -4865,7 +4888,7 @@ export function useChat(
48654888
userRemovedDuringDispatch.delete(msg.id)
48664889
}
48674890
},
4868-
[startSendMessage, handOffWithdrawnSend]
4891+
[startSendMessage, handOffWithdrawnSend, heldSendSurface]
48694892
)
48704893

48714894
const runQueueDispatchLoop = useCallback(async () => {
@@ -5019,21 +5042,31 @@ export function useChat(
50195042
}
50205043
}, [])
50215044

5022-
/** A send held because the server could not be reached goes out once the browser is back online. */
5045+
/**
5046+
* Sends held because the server could not be reached go out once the browser is
5047+
* online: on the `online` event, and on mount in case it fired while no chat
5048+
* surface was listening. A chatless surface first adopts what a dead mount of the
5049+
* same surface held, since that mount's queue key died with it.
5050+
*/
50235051
useEffect(() => {
50245052
if (typeof window === 'undefined') return
50255053
const releaseHeldSends = () => {
5026-
const chatKey = chatKeyRef.current
5027-
const queue = useMothershipQueueStore.getState().queues[chatKey]
5028-
if (!queue?.some((message) => message.heldUntilOnline)) return
5029-
useMothershipQueueStore.getState().releaseHeldUntilOnline(chatKey)
5030-
if (!sendingRef.current && !pendingStopPromiseRef.current) {
5054+
useMothershipQueueStore.getState().releaseHeldUntilOnline()
5055+
if (
5056+
useMothershipQueueStore.getState().queues[chatKeyRef.current]?.length &&
5057+
!sendingRef.current &&
5058+
!pendingStopPromiseRef.current
5059+
) {
50315060
void enqueueQueueDispatchRef.current({ type: 'send_head' })
50325061
}
50335062
}
5063+
if (chatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) {
5064+
useMothershipQueueStore.getState().adoptHeldSends(chatKey, heldSendSurface)
5065+
}
5066+
if (navigator.onLine) releaseHeldSends()
50345067
window.addEventListener('online', releaseHeldSends)
50355068
return () => window.removeEventListener('online', releaseHeldSends)
5036-
}, [])
5069+
}, [chatKey, heldSendSurface])
50375070

50385071
/** A recovered send already in history belongs to its accepted turn, even after Stop. */
50395072
useEffect(() => {

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

Lines changed: 39 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -94,6 +94,7 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
9494
resumeUserMessageId: _staleResume,
9595
retryRequired: _retry,
9696
heldUntilOnline: _held,
97+
heldSurface: _surface,
9798
...rest
9899
} = next[index]
99100
next[index] = {
@@ -143,7 +144,11 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
143144
// Merge defensively in case a stale bucket survived in
144145
// sessionStorage. FIFO: existing first, then the resolved stream.
145146
const existing = state.queues[toKey] ?? []
146-
queues[toKey] = [...existing, ...fromQueue]
147+
/** A chat-bound key is stable, so its messages no longer need a surface to adopt them. */
148+
queues[toKey] = [
149+
...existing,
150+
...fromQueue.map(({ heldSurface: _surface, ...message }) => message),
151+
]
147152
}
148153
const editing = omitKey(state.editing, fromKey)
149154
if (fromEditing !== undefined) {
@@ -152,16 +157,40 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
152157
return { queues, editing }
153158
}),
154159

155-
releaseHeldUntilOnline: (chatKey) =>
160+
releaseHeldUntilOnline: () =>
156161
set((state) => {
157-
const current = state.queues[chatKey] ?? []
158-
if (!current.some((m) => m.heldUntilOnline)) return state
159-
const next = current.map((message) => {
160-
if (!message.heldUntilOnline) return message
161-
const { retryRequired: _retry, heldUntilOnline: _held, ...rest } = message
162-
return rest
163-
})
164-
return { queues: setQueueForChat(state.queues, chatKey, next) }
162+
let released = false
163+
const queues: Record<string, QueuedMothershipMessage[]> = {}
164+
for (const [chatKey, queue] of Object.entries(state.queues)) {
165+
queues[chatKey] = queue.map((message) => {
166+
if (!message.heldUntilOnline) return message
167+
released = true
168+
const { retryRequired: _retry, heldUntilOnline: _held, ...rest } = message
169+
return rest
170+
})
171+
}
172+
return released ? { queues } : state
173+
}),
174+
175+
adoptHeldSends: (toKey, surface) =>
176+
set((state) => {
177+
const adopted: QueuedMothershipMessage[] = []
178+
let queues = state.queues
179+
for (const [chatKey, queue] of Object.entries(state.queues)) {
180+
if (chatKey === toKey) continue
181+
const held = queue.filter((message) => message.heldSurface === surface)
182+
if (held.length === 0) continue
183+
adopted.push(...held)
184+
queues = setQueueForChat(
185+
queues,
186+
chatKey,
187+
queue.filter((message) => message.heldSurface !== surface)
188+
)
189+
}
190+
if (adopted.length === 0) return state
191+
return {
192+
queues: setQueueForChat(queues, toKey, [...(queues[toKey] ?? []), ...adopted]),
193+
}
165194
}),
166195

167196
clearChat: (chatKey) =>

0 commit comments

Comments
 (0)