Skip to content

Commit 98a8ddd

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 e130725 commit 98a8ddd

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)
@@ -3828,10 +3830,10 @@ export function useChat(
38283830
if (!response.ok) {
38293831
const errorData = await response.json().catch(() => ({}))
38303832
if (response.status === 409) {
3833+
/* A deduplicated send always names itself; a chat busy with another
3834+
turn names that turn, or nothing when its stream id is unreadable. */
38313835
const conflictStreamId =
3832-
typeof errorData.activeStreamId === 'string'
3833-
? errorData.activeStreamId
3834-
: userMessageId
3836+
typeof errorData.activeStreamId === 'string' ? errorData.activeStreamId : undefined
38353837
const supersededStreamId = queuedSendHandoff?.supersededStreamId ?? pendingStopStreamId
38363838
if (supersededStreamId && conflictStreamId === supersededStreamId) {
38373839
rollbackOptimisticSend()
@@ -3847,8 +3849,11 @@ export function useChat(
38473849
}
38483850
if (conflictStreamId !== userMessageId) {
38493851
/* Another turn holds the chat: one started in another tab, or one this
3850-
surface lost track of. This message was not admitted, so it waits in
3851-
the queue behind that turn, which the chat now shows as running. */
3852+
surface lost track of. This message was not admitted (the server
3853+
released its id), so it goes back to the queue under the same id. The
3854+
queue drains only while the chat is idle, so the chat's running turn is
3855+
read before the message is handed back: the chat then attaches to that
3856+
turn and the message goes out once, after it ends. */
38523857
rollbackOptimisticSend()
38533858
if (streamGenRef.current === gen) {
38543859
streamGenRef.current++
@@ -3859,13 +3864,15 @@ export function useChat(
38593864
}
38603865
if (requestChatId) {
38613866
const busyChatId = requestChatId
3862-
upsertChatHistory(busyChatId, (current) => ({
3863-
...current,
3864-
activeStreamId: conflictStreamId,
3865-
}))
3866-
void queryClient.invalidateQueries({
3867-
queryKey: mothershipChatKeys.detail(busyChatId),
3868-
})
3867+
if (conflictStreamId) {
3868+
upsertChatHistory(busyChatId, (current) => ({
3869+
...current,
3870+
activeStreamId: conflictStreamId,
3871+
}))
3872+
}
3873+
await queryClient
3874+
.refetchQueries({ queryKey: mothershipChatKeys.detail(busyChatId), exact: true })
3875+
.catch(() => {})
38693876
}
38703877
return { userMessageId }
38713878
}
@@ -4197,14 +4204,23 @@ export function useChat(
41974204
options?.assistantSearch,
41984205
options?.assistantSearchLevel
41994206
),
4200-
...(result.unreachable ? { retryRequired: true, heldUntilOnline: true } : {}),
4207+
...(result.unreachable
4208+
? {
4209+
retryRequired: true,
4210+
heldUntilOnline: true,
4211+
...(activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)
4212+
? { heldSurface: heldSendSurface }
4213+
: {}),
4214+
}
4215+
: {}),
42014216
})
42024217
},
42034218
[
42044219
workspaceId,
42054220
createQueuedMessage,
42064221
startSendMessage,
42074222
handOffWithdrawnSend,
4223+
heldSendSurface,
42084224
hasPendingChatAdmission,
42094225
]
42104226
)
@@ -4804,7 +4820,14 @@ export function useChat(
48044820
...dispatched,
48054821
...(retainedHandoff ? { queuedSendHandoff: retainedHandoff } : {}),
48064822
retryRequired: !retriesOnItsOwn,
4807-
...(withdrawn?.unreachable ? { heldUntilOnline: true } : {}),
4823+
...(withdrawn?.unreachable
4824+
? {
4825+
heldUntilOnline: true,
4826+
...(dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)
4827+
? { heldSurface: heldSendSurface }
4828+
: {}),
4829+
}
4830+
: {}),
48084831
...(withdrawnUserMessageId ? { resumeUserMessageId: withdrawnUserMessageId } : {}),
48094832
})
48104833
}
@@ -4859,7 +4882,7 @@ export function useChat(
48594882
userRemovedDuringDispatch.delete(msg.id)
48604883
}
48614884
},
4862-
[startSendMessage, handOffWithdrawnSend]
4885+
[startSendMessage, handOffWithdrawnSend, heldSendSurface]
48634886
)
48644887

48654888
const runQueueDispatchLoop = useCallback(async () => {
@@ -5013,21 +5036,31 @@ export function useChat(
50135036
}
50145037
}, [])
50155038

5016-
/** A send held because the server could not be reached goes out once the browser is back online. */
5039+
/**
5040+
* Sends held because the server could not be reached go out once the browser is
5041+
* online: on the `online` event, and on mount in case it fired while no chat
5042+
* surface was listening. A chatless surface first adopts what a dead mount of the
5043+
* same surface held, since that mount's queue key died with it.
5044+
*/
50175045
useEffect(() => {
50185046
if (typeof window === 'undefined') return
50195047
const releaseHeldSends = () => {
5020-
const chatKey = chatKeyRef.current
5021-
const queue = useMothershipQueueStore.getState().queues[chatKey]
5022-
if (!queue?.some((message) => message.heldUntilOnline)) return
5023-
useMothershipQueueStore.getState().releaseHeldUntilOnline(chatKey)
5024-
if (!sendingRef.current && !pendingStopPromiseRef.current) {
5048+
useMothershipQueueStore.getState().releaseHeldUntilOnline()
5049+
if (
5050+
useMothershipQueueStore.getState().queues[chatKeyRef.current]?.length &&
5051+
!sendingRef.current &&
5052+
!pendingStopPromiseRef.current
5053+
) {
50255054
void enqueueQueueDispatchRef.current({ type: 'send_head' })
50265055
}
50275056
}
5057+
if (chatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) {
5058+
useMothershipQueueStore.getState().adoptHeldSends(chatKey, heldSendSurface)
5059+
}
5060+
if (navigator.onLine) releaseHeldSends()
50285061
window.addEventListener('online', releaseHeldSends)
50295062
return () => window.removeEventListener('online', releaseHeldSends)
5030-
}, [])
5063+
}, [chatKey, heldSendSurface])
50315064

50325065
/** A recovered send already in history belongs to its accepted turn, even after Stop. */
50335066
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)