Skip to content

Commit c8aa1d2

Browse files
committed
fix(mothership): send released held messages through the queue's own drain rules
Releasing held sends on mount kicked the queue dispatcher directly, which skips the drain's guards. After a reload the chat history is not loaded yet, so a follow-up queued behind a still-running turn went out at once and was refused as busy. The release now only clears the hold; the drain effect, which waits for the history and for the running turn to end, sends a released head, and now also re-runs when the head's hold clears.
1 parent 9d3a096 commit c8aa1d2

2 files changed

Lines changed: 52 additions & 15 deletions

File tree

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

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2114,6 +2114,48 @@ describe('useChat remount send recovery', () => {
21142114
expect(state.postBodies.at(-1)?.message).toBe('First message, sent offline')
21152115
})
21162116

2117+
/**
2118+
* Releasing held sends on mount must not bypass the queue's own rules: a
2119+
* follow-up queued behind a turn that is still running (here, restored after a
2120+
* reload) waits for that turn instead of being sent into a busy chat.
2121+
*/
2122+
it('keeps a queued follow-up waiting on mount while the chat is still running', async () => {
2123+
const history: MothershipChatHistory = {
2124+
...idleHistory('chat-still-running'),
2125+
activeStreamId: 'turn-still-running',
2126+
}
2127+
mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history }))
2128+
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
2129+
const url = String(input)
2130+
if (url === '/api/mothership/chat' && init?.method === 'POST') {
2131+
state.postBodies.push(JSON.parse(String(init.body)))
2132+
return emptySseResponse()
2133+
}
2134+
if (url.includes('/api/mothership/chat/stream')) {
2135+
if (url.includes('batch=true')) {
2136+
return Response.json({ success: true, events: [], status: 'streaming' })
2137+
}
2138+
return new Response(new ReadableStream<Uint8Array>(), {
2139+
headers: { 'Content-Type': 'text/event-stream' },
2140+
})
2141+
}
2142+
return fetchStub(input, init)
2143+
})
2144+
useMothershipQueueStore
2145+
.getState()
2146+
.enqueue(history.id, { id: 'queued-before-reload', content: 'Queued before the reload' })
2147+
renderUseChatInChat(history.id)
2148+
2149+
await act(async () => {
2150+
await sleep(1000)
2151+
})
2152+
2153+
expect(state.postBodies).toHaveLength(0)
2154+
expect(useMothershipQueueStore.getState().queues[history.id]?.[0]?.content).toBe(
2155+
'Queued before the reload'
2156+
)
2157+
})
2158+
21172159
/** The `online` event can fire while no surface for the chat is mounted. */
21182160
it('sends a held message when its chat mounts after the network came back', async () => {
21192161
const history = idleHistory('chat-held-while-away')

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

Lines changed: 10 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -5043,23 +5043,16 @@ export function useChat(
50435043
}, [])
50445044

50455045
/**
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.
5046+
* Sends held because the server could not be reached are released once the
5047+
* browser is online: on the `online` event, and on mount in case it fired while
5048+
* no chat surface was listening. The queue drain below sends a released head
5049+
* under its usual rules (history loaded, no running turn). A chatless surface
5050+
* first adopts what a dead mount of the same surface held, since that mount's
5051+
* queue key died with it.
50505052
*/
50515053
useEffect(() => {
50525054
if (typeof window === 'undefined') return
5053-
const releaseHeldSends = () => {
5054-
useMothershipQueueStore.getState().releaseHeldUntilOnline()
5055-
if (
5056-
useMothershipQueueStore.getState().queues[chatKeyRef.current]?.length &&
5057-
!sendingRef.current &&
5058-
!pendingStopPromiseRef.current
5059-
) {
5060-
void enqueueQueueDispatchRef.current({ type: 'send_head' })
5061-
}
5062-
}
5055+
const releaseHeldSends = () => useMothershipQueueStore.getState().releaseHeldUntilOnline()
50635056
if (chatKey.startsWith(PENDING_CHAT_KEY_PREFIX)) {
50645057
useMothershipQueueStore.getState().adoptHeldSends(chatKey, heldSendSurface)
50655058
}
@@ -5091,9 +5084,10 @@ export function useChat(
50915084
// `notifyTurnEnded`. Idempotent — the dispatch loop dedupes.
50925085
const chatHistoryReady = chatHistory !== undefined
50935086
const remoteActiveStreamId = chatHistory?.activeStreamId ?? null
5087+
const queueHeadHeld = messageQueue[0]?.retryRequired === true
50945088
useEffect(() => {
50955089
if (!scopeKey) return
5096-
if (messageQueue.length === 0) return
5090+
if (messageQueue.length === 0 || queueHeadHeld) return
50975091
if (sendingRef.current || pendingStopPromiseRef.current) return
50985092
if (queueDispatchTaskRef.current) return
50995093
if (resolvedChatId && !chatHistoryReady) return
@@ -5104,6 +5098,7 @@ export function useChat(
51045098
organizationId,
51055099
scopeKey,
51065100
messageQueue.length,
5101+
queueHeadHeld,
51075102
resolvedChatId,
51085103
chatHistoryReady,
51095104
remoteActiveStreamId,

0 commit comments

Comments
 (0)