Skip to content

Commit e130725

Browse files
committed
fix(mothership): keep a message the server never admitted instead of dropping it
Two sends were lost without an error: - A send whose POST got no response (offline, Wi-Fi drop, waking a laptop) reconnected to the stream it would have opened. That stream does not exist, so the 404 read as "finished", the turn finalized as a success, and the refetched transcript no longer held the message. A queued follow-up was lost the same way, since it had already left the queue. - A send refused with 409 because another turn held the chat (started in another tab, or one this surface lost track of) reconnected to that turn under the new message's bubble, then vanished when it finished. Both now hand the message back under its id. An unreachable send is held in the queue, so it is not redispatched into the same failure, and goes out when the browser is back online (or when the user sends it). A send that found the chat busy waits in the queue behind that turn, which the chat shows as running, and goes out when it ends. Reusing the id keeps a retry deduplicated if the server did admit the first attempt.
1 parent 65b247f commit e130725

4 files changed

Lines changed: 235 additions & 30 deletions

File tree

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

Lines changed: 112 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1921,6 +1921,118 @@ describe('useChat remount send recovery', () => {
19211921
})
19221922
})
19231923

1924+
describe('a send the server never admitted', () => {
1925+
const idleHistory = (id: string): MothershipChatHistory => ({
1926+
id,
1927+
mode: 'agent',
1928+
title: 'Not admitted',
1929+
messages: [],
1930+
activeStreamId: null,
1931+
resources: [],
1932+
})
1933+
1934+
/**
1935+
* The POST fails at the network layer until the network is back, and the
1936+
* stream it would have opened does not exist.
1937+
*/
1938+
const network = { online: false }
1939+
function stubUnreachableSend() {
1940+
network.online = false
1941+
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
1942+
const url = String(input)
1943+
if (url === '/api/mothership/chat' && init?.method === 'POST') {
1944+
state.postBodies.push(JSON.parse(String(init.body)))
1945+
if (network.online) return emptySseResponse()
1946+
throw new TypeError('Failed to fetch')
1947+
}
1948+
if (url.includes('/api/mothership/chat/stream')) {
1949+
return Response.json({ error: 'Stream not found' }, { status: 404 })
1950+
}
1951+
return fetchStub(input, init)
1952+
})
1953+
}
1954+
1955+
it('holds a message sent while offline and sends it under the same id once back online', async () => {
1956+
const history = idleHistory('chat-offline-send')
1957+
mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history }))
1958+
stubUnreachableSend()
1959+
const { getResult } = renderUseChatInChat(history.id, history)
1960+
1961+
await act(async () => {
1962+
await getResult().sendMessage('Written while offline')
1963+
})
1964+
await waitFor(() => !getResult().isSending)
1965+
1966+
const queued = useMothershipQueueStore.getState().queues[history.id] ?? []
1967+
expect(queued.map((message) => message.content)).toEqual(['Written while offline'])
1968+
expect(queued[0].retryRequired).toBe(true)
1969+
expect(queued[0].resumeUserMessageId).toBe(state.postBodies[0].userMessageId)
1970+
expect(getResult().error).not.toBeNull()
1971+
expect(state.postBodies).toHaveLength(1)
1972+
1973+
network.online = true
1974+
await act(async () => {
1975+
window.dispatchEvent(new Event('online'))
1976+
})
1977+
await waitFor(() => state.postBodies.length === 2)
1978+
1979+
expect(state.postBodies[1].message).toBe('Written while offline')
1980+
expect(state.postBodies[1].userMessageId).toBe(state.postBodies[0].userMessageId)
1981+
})
1982+
1983+
it('keeps a queued follow-up whose dispatch could not reach the server', async () => {
1984+
const history = idleHistory('chat-offline-queue')
1985+
mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history }))
1986+
stubUnreachableSend()
1987+
useMothershipQueueStore
1988+
.getState()
1989+
.enqueue(history.id, { id: 'queued-follow-up', content: 'Queued before the drop' })
1990+
renderUseChatInChat(history.id, history)
1991+
1992+
await waitFor(() => state.postBodies.length === 1)
1993+
await waitFor(
1994+
() => useMothershipQueueStore.getState().queues[history.id]?.[0]?.retryRequired === true
1995+
)
1996+
1997+
const queued = useMothershipQueueStore.getState().queues[history.id] ?? []
1998+
expect(queued.map((message) => message.content)).toEqual(['Queued before the drop'])
1999+
expect(queued[0].resumeUserMessageId).toBe(state.postBodies[0].userMessageId)
2000+
expect(state.postBodies).toHaveLength(1)
2001+
})
2002+
2003+
/**
2004+
* 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.
2007+
*/
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+
)
2019+
}
2020+
return emptySseResponse()
2021+
}
2022+
return fetchStub(input, init)
2023+
})
2024+
const { getResult } = renderUseChatInChat(history.id, history)
2025+
2026+
await act(async () => {
2027+
await getResult().sendMessage('Sent from the second tab')
2028+
})
2029+
await waitFor(() => state.postBodies.length === 2)
2030+
2031+
expect(state.postBodies[1].message).toBe('Sent from the second tab')
2032+
expect(state.postBodies[1].userMessageId).toBe(state.postBodies[0].userMessageId)
2033+
})
2034+
})
2035+
19242036
/**
19252037
* A withdrawn send belongs to the chat it was sent to. The cross-surface
19262038
* lanes deliver to whatever chat is mounted next, so routing a chat-bound

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

Lines changed: 103 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -205,13 +205,24 @@ interface FinalizeOptions {
205205
streamTerminal?: boolean
206206
}
207207

208+
/**
209+
* A send handed back to the caller instead of rendered. `userMessageId` is what
210+
* a retry reuses so the server deduplicates the two attempts. An `unreachable`
211+
* send is held in the queue until the browser is back online or the user sends
212+
* it: dispatching it again at once would fail the same way.
213+
*/
214+
interface WithdrawnSendResult {
215+
userMessageId: string
216+
unreachable?: boolean
217+
}
218+
208219
/**
209220
* `true` when the send owns the transcript (rendered, or handed to reconnect),
210221
* `false` when the caller should restore the queue entry, and the object form
211-
* when an unmount cleanup withdrew it — `userMessageId` is what a retry reuses
212-
* so the server deduplicates the two attempts.
222+
* when the send was withdrawn: by an unmount cleanup, because it never reached
223+
* the server, or because another turn held the chat.
213224
*/
214-
type StartSendMessageResult = boolean | { userMessageId: string }
225+
type StartSendMessageResult = boolean | WithdrawnSendResult
215226

216227
interface StartSendMessageOptions {
217228
/** Awaited before dispatch. Defaults to the hook's in-flight stop, if any. */
@@ -3834,6 +3845,30 @@ export function useChat(
38343845
setError('Previous response is still shutting down; queued message was restored.')
38353846
return false
38363847
}
3848+
if (conflictStreamId !== userMessageId) {
3849+
/* 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+
rollbackOptimisticSend()
3853+
if (streamGenRef.current === gen) {
3854+
streamGenRef.current++
3855+
abortController.abort('send_conflict:chat_busy')
3856+
abortControllerRef.current = null
3857+
clearActiveTurn()
3858+
setTransportIdle()
3859+
}
3860+
if (requestChatId) {
3861+
const busyChatId = requestChatId
3862+
upsertChatHistory(busyChatId, (current) => ({
3863+
...current,
3864+
activeStreamId: conflictStreamId,
3865+
}))
3866+
void queryClient.invalidateQueries({
3867+
queryKey: mothershipChatKeys.detail(busyChatId),
3868+
})
3869+
}
3870+
return { userMessageId }
3871+
}
38373872
/* A send deduplicated against an earlier attempt comes back naming
38383873
the chat that attempt opened. Adopting it here spares a chatless
38393874
surface the stream-to-chat lookup and puts the user in the right
@@ -3953,6 +3988,27 @@ export function useChat(
39533988
return consumedByTranscript
39543989
}
39553990

3991+
if (!sendReachedServer) {
3992+
/* The POST got no answer, so nothing shows the server admitted it, and a
3993+
stream that was never opened reads as finished. Hand the message back
3994+
under the same id: if the server did admit it, the retry deduplicates
3995+
and the chat's own recovery shows the running turn. */
3996+
rollbackOptimisticSend()
3997+
if (gen !== undefined && streamGenRef.current === gen) {
3998+
streamGenRef.current++
3999+
admission?.controller.abort('send_unreachable:no_response')
4000+
abortControllerRef.current = null
4001+
clearActiveTurn()
4002+
setTransportIdle()
4003+
}
4004+
setError(
4005+
err instanceof TypeError
4006+
? 'Message not sent: Sim could not be reached. It will send when you are back online.'
4007+
: getErrorMessage(err, 'Failed to send message')
4008+
)
4009+
return { userMessageId, unreachable: true }
4010+
}
4011+
39564012
const activeStreamId = streamIdRef.current
39574013
if (activeStreamId && gen !== undefined && streamGenRef.current === gen) {
39584014
const succeeded = await retryReconnect({
@@ -4110,11 +4166,12 @@ export function useChat(
41104166
const result = await startSendMessage(message, fileAttachments, contexts, options)
41114167
if (typeof result !== 'object') return
41124168

4113-
/* An unmount cleanup withdrew the send. A chat-bound key is the stable
4114-
chat id, so re-queueing under the key this was sent to is the durable
4115-
retry — and keeps the message in that chat rather than following the
4116-
user into whichever one they opened next. Only a chatless surface,
4117-
whose key dies with the mount, goes to the cross-surface lanes. */
4169+
/* The send was withdrawn. A chat-bound key is the stable chat id, so
4170+
re-queueing under the key this was sent to is the durable retry, and
4171+
keeps the message in that chat rather than following the user into
4172+
whichever one they opened next. Only a send an unmount withdrew from a
4173+
chatless surface, whose key dies with the mount, goes to the
4174+
cross-surface lanes. */
41184175
const withdrawn = {
41194176
content: message,
41204177
fileAttachments,
@@ -4126,24 +4183,22 @@ export function useChat(
41264183
? { assistantSearchLevel: options?.assistantSearchLevel }
41274184
: {}),
41284185
}
4129-
if (activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) {
4186+
if (!result.unreachable && activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) {
41304187
handOffWithdrawnSend(withdrawn)
41314188
return
41324189
}
4133-
useMothershipQueueStore
4134-
.getState()
4135-
.enqueue(
4136-
activeChatKey,
4137-
createQueuedMessage(
4138-
message,
4139-
fileAttachments,
4140-
contexts,
4141-
result.userMessageId,
4142-
options?.requestMode,
4143-
options?.assistantSearch,
4144-
options?.assistantSearchLevel
4145-
)
4146-
)
4190+
useMothershipQueueStore.getState().enqueue(activeChatKey, {
4191+
...createQueuedMessage(
4192+
message,
4193+
fileAttachments,
4194+
contexts,
4195+
result.userMessageId,
4196+
options?.requestMode,
4197+
options?.assistantSearch,
4198+
options?.assistantSearchLevel
4199+
),
4200+
...(result.unreachable ? { retryRequired: true, heldUntilOnline: true } : {}),
4201+
})
41474202
},
41484203
[
41494204
workspaceId,
@@ -4696,9 +4751,10 @@ export function useChat(
46964751
let dispatched = msg
46974752
const restoreQueuedMessage = (
46984753
handoff?: QueuedSendHandoffSeed,
4699-
withdrawnUserMessageId?: string
4754+
withdrawn?: WithdrawnSendResult
47004755
) => {
4701-
const withdrawnByCleanup = withdrawnUserMessageId !== undefined
4756+
const withdrawnUserMessageId = withdrawn?.userMessageId
4757+
const retriesOnItsOwn = withdrawn !== undefined && !withdrawn.unreachable
47024758
const savedHandoff = readQueuedSendHandoffState()
47034759
const retainedHandoff =
47044760
savedHandoff?.id === msg.id
@@ -4714,7 +4770,7 @@ export function useChat(
47144770
if (!removedFromQueue) {
47154771
return
47164772
}
4717-
if (options.epoch !== queueDispatchEpochRef.current && !withdrawnByCleanup) {
4773+
if (options.epoch !== queueDispatchEpochRef.current && !retriesOnItsOwn) {
47184774
return
47194775
}
47204776
// If the user explicitly removed this message during dispatch, honor
@@ -4727,7 +4783,7 @@ export function useChat(
47274783
restore would strand this under the dead instance's key — hand it to
47284784
the next surface instead. A chat-bound key is the stable chat id, so
47294785
the queue itself is the durable retry. */
4730-
if (withdrawnByCleanup && dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) {
4786+
if (withdrawn && retriesOnItsOwn && dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) {
47314787
clearQueuedSendHandoffState(msg.id)
47324788
handOffWithdrawnSend({
47334789
content: dispatched.content,
@@ -4738,7 +4794,7 @@ export function useChat(
47384794
...(dispatched.assistantSearchLevel !== undefined
47394795
? { assistantSearchLevel: dispatched.assistantSearchLevel }
47404796
: {}),
4741-
userMessageId: withdrawnUserMessageId,
4797+
userMessageId: withdrawn.userMessageId,
47424798
})
47434799
return
47444800
}
@@ -4747,7 +4803,8 @@ export function useChat(
47474803
useMothershipQueueStore.getState().insertAt(dispatchChatKey, originalIndex, {
47484804
...dispatched,
47494805
...(retainedHandoff ? { queuedSendHandoff: retainedHandoff } : {}),
4750-
retryRequired: !withdrawnByCleanup,
4806+
retryRequired: !retriesOnItsOwn,
4807+
...(withdrawn?.unreachable ? { heldUntilOnline: true } : {}),
47514808
...(withdrawnUserMessageId ? { resumeUserMessageId: withdrawnUserMessageId } : {}),
47524809
})
47534810
}
@@ -4791,7 +4848,7 @@ export function useChat(
47914848
if (sendResult !== true) {
47924849
restoreQueuedMessage(
47934850
activeQueuedSendHandoff,
4794-
typeof sendResult === 'object' ? sendResult.userMessageId : undefined
4851+
typeof sendResult === 'object' ? sendResult : undefined
47954852
)
47964853
}
47974854
} catch {
@@ -4956,6 +5013,22 @@ export function useChat(
49565013
}
49575014
}, [])
49585015

5016+
/** A send held because the server could not be reached goes out once the browser is back online. */
5017+
useEffect(() => {
5018+
if (typeof window === 'undefined') return
5019+
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) {
5025+
void enqueueQueueDispatchRef.current({ type: 'send_head' })
5026+
}
5027+
}
5028+
window.addEventListener('online', releaseHeldSends)
5029+
return () => window.removeEventListener('online', releaseHeldSends)
5030+
}, [])
5031+
49595032
/** A recovered send already in history belongs to its accepted turn, even after Stop. */
49605033
useEffect(() => {
49615034
if (!chatHistory || chatHistory.id !== chatKeyRef.current) return

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

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,7 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
9393
queuedSendHandoff,
9494
resumeUserMessageId: _staleResume,
9595
retryRequired: _retry,
96+
heldUntilOnline: _held,
9697
...rest
9798
} = next[index]
9899
next[index] = {
@@ -151,6 +152,18 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
151152
return { queues, editing }
152153
}),
153154

155+
releaseHeldUntilOnline: (chatKey) =>
156+
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) }
165+
}),
166+
154167
clearChat: (chatKey) =>
155168
set((state) => ({
156169
queues: omitKey(state.queues, chatKey),

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

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,11 @@ export type QueuedMothershipMessage = QueuedMessage & {
1313
queuedSendHandoff?: QueuedSendHandoffSeed
1414
/** A failed dispatch remains queued until the user retries or edits it. */
1515
retryRequired?: boolean
16+
/**
17+
* The failed dispatch never reached the server, so the browser coming back
18+
* online releases it for dispatch too.
19+
*/
20+
heldUntilOnline?: boolean
1621
/**
1722
* Message id of a prior attempt at this send that an unmount cleanup
1823
* withdrew. Reused when the entry is dispatched so the server deduplicates
@@ -44,6 +49,8 @@ export interface MothershipQueueState {
4449
remove: (chatKey: string, id: string) => void
4550
setEditing: (chatKey: string, id: string | null) => void
4651
migrate: (fromKey: string, toKey: string) => void
52+
/** Releases the chat's sends held for the network for dispatch. */
53+
releaseHeldUntilOnline: (chatKey: string) => void
4754
clearChat: (chatKey: string) => void
4855
reset: () => void
4956
}

0 commit comments

Comments
 (0)