Skip to content

Commit f0e5700

Browse files
authored
fix(mothership): send a late queue write to the chat a new-chat queue moved to (#8741)
* fix(mothership): send a late queue write to the chat a new-chat queue moved to On the new-chat surface, a Send-now whose Stop saw the first message admitted moved the queue to that chat, but the follow-up's busy refusal was re-queued under the dead new-chat key it was dispatched from: gone from the chat, never retried there, and liable to be adopted into a new chat later. migrate now records where a key moved, and every write that captured a key before an await (the dispatch's removal and restore, the direct send's re-queue, the history check's defer and drop) resolves it at write time (liveQueueKey). Every re-queue on a chatless surface also carries the surface, so one that lands after the surface unmounted can still be adopted. * fix(mothership): keep a late write behind what the chat's queue already held When the new-chat queue moved into a chat queue that already had messages, migrate put those first, but a late write still used its index in the new-chat queue and could land ahead of them. The move now records how many messages it went behind, and late writes resolve their position, not just their key (liveQueuePosition). * fix(mothership): anchor a late re-queue on the messages ahead of it, not an index A restored message went back at an index captured at dispatch, offset by how many messages a chat's queue held when the new-chat queue moved into it. Removing any of those while the POST was out shifted it behind a newer message. It now goes right after the last message still queued that was ahead of it (at dispatch, or in the chat's queue before the move), else at the head.
1 parent af5501b commit f0e5700

8 files changed

Lines changed: 208 additions & 23 deletions

File tree

‎apps/sim/app/workspace/[workspaceId]/home/hooks/send-queue-policy.test.ts‎

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,9 +25,13 @@ describe('requeuedFields', () => {
2525
})
2626

2727
it.each(['stop-failed', 'failed'] as const)(
28-
'leaves a %s send for the user, on any surface',
28+
'leaves a %s send for the user, adoptable by its chatless surface',
2929
(reason) => {
30-
expect(requeuedFields(reason, 4, 'ws-1:home')).toEqual({ retryRequired: true })
30+
expect(requeuedFields(reason, 4, 'ws-1:home')).toEqual({
31+
retryRequired: true,
32+
heldSurface: 'ws-1:home',
33+
})
34+
expect(requeuedFields(reason, 4, undefined)).toEqual({ retryRequired: true })
3135
}
3236
)
3337

‎apps/sim/app/workspace/[workspaceId]/home/hooks/send-queue-policy.ts‎

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -42,9 +42,10 @@ export function sendRetry(attempt: number): ScheduledRetry {
4242
* - `stop-failed` and `failed` wait for the user;
4343
* - `withdrawn` goes out again as soon as the queue drains.
4444
*
45-
* A message held by a chatless surface carries that surface (`chatlessSurface`),
46-
* whose queue key dies with its mount, so the next mount of it adopts the
47-
* message. Only sends that wait on the network or the server are held that way.
45+
* A message re-queued on a chatless surface carries that surface
46+
* (`chatlessSurface`), whose queue key dies with its mount, so the next mount of
47+
* it adopts the message. That holds for every reason: a re-queue can land after
48+
* the surface unmounted, when nothing else marks the dead key's queue.
4849
*/
4950
export function requeuedFields(
5051
reason: RequeueReason,
@@ -60,7 +61,7 @@ export function requeuedFields(
6061
return { ...sendRetry(previousAttempts + 1), ...surface }
6162
case 'stop-failed':
6263
case 'failed':
63-
return { retryRequired: true }
64+
return { retryRequired: true, ...surface }
6465
case 'withdrawn':
6566
return {}
6667
}

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

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2829,6 +2829,57 @@ describe('useChat remount send recovery', () => {
28292829
expect(allQueuedMessages()).toHaveLength(0)
28302830
})
28312831

2832+
/**
2833+
* Send-now on the new-chat surface stops the first message, which the Stop sees
2834+
* admitted into a chat: the surface moves to that chat and its queue moves with
2835+
* it. A busy refusal of the follow-up arriving after that must go back to the
2836+
* chat's queue, where the surface shows and retries it, not to the dead
2837+
* new-chat key it was dispatched from.
2838+
*/
2839+
it('re-queues a Send-now refused as busy in the chat the new-chat surface moved to', async () => {
2840+
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
2841+
const url = String(input)
2842+
if (
2843+
url === '/api/mothership/chat' &&
2844+
init?.method === 'POST' &&
2845+
state.postBodies.length > 0
2846+
) {
2847+
state.postBodies.push(JSON.parse(String(init.body)))
2848+
return Response.json(
2849+
{ error: 'A response is already in progress for this chat.' },
2850+
{ status: 409 }
2851+
)
2852+
}
2853+
return fetchStub(input, init)
2854+
})
2855+
const { getResult } = renderUseChat()
2856+
await act(async () => {
2857+
void getResult().sendMessage('inspect the workspace')
2858+
})
2859+
await waitFor(() => state.postBodies.length === 1)
2860+
await act(async () => {
2861+
void getResult().sendMessage('Follow-up')
2862+
})
2863+
await waitFor(() => allQueuedMessages().length === 1)
2864+
2865+
await act(async () => {
2866+
void getResult()
2867+
.sendNow()
2868+
.catch(() => {})
2869+
})
2870+
await waitFor(() => state.postBodies.length === 2)
2871+
await act(async () => {
2872+
await sleep(200)
2873+
})
2874+
2875+
const queues = useMothershipQueueStore.getState().queues
2876+
expect(queues[DEDUPED_CHAT_ID]?.map((message) => message.content)).toEqual(['Follow-up'])
2877+
expect(
2878+
Object.entries(queues).filter(([key, queue]) => key.startsWith('pending::') && queue.length)
2879+
).toEqual([])
2880+
expect(getResult().messageQueue.map((message) => message.content)).toEqual(['Follow-up'])
2881+
})
2882+
28322883
it('stopping a chat preserves an unrelated manual workflow execution', async () => {
28332884
const executionStore = useExecutionStore.getState()
28342885
executionStore.setIsExecuting('manual-workflow', true)

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

Lines changed: 27 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -134,7 +134,12 @@ import { workflowKeys } from '@/hooks/queries/workflows'
134134
import { snapAllSmoothText } from '@/hooks/use-smooth-text'
135135
import { useChatPanelStore } from '@/stores/chat-panel/store'
136136
import { useMothershipEffortStore } from '@/stores/mothership-effort/store'
137-
import { reusedRequestId, useMothershipQueueStore } from '@/stores/mothership-queue/store'
137+
import {
138+
liveQueueKey,
139+
liveQueuePosition,
140+
reusedRequestId,
141+
useMothershipQueueStore,
142+
} from '@/stores/mothership-queue/store'
138143
import type {
139144
QueuedMothershipMessage,
140145
QueuedSendHandoffSeed,
@@ -336,7 +341,7 @@ const PERSISTED_TURN_REFETCH_BASE_MS = 250
336341
const PERSISTED_TURN_REFETCH_MAX_DELAY_MS = 5_000
337342
/** How long a finished turn's save is waited for; a slow save still lands well inside it. */
338343
const PERSISTED_TURN_WAIT_MS = 120_000
339-
/** Pacing for re-sending a message refused because the chat was busy, or that could not reach Sim. */
344+
/** How long a Stop's abort request may take before the Stop counts as failed. */
340345
const STOP_REQUEST_TIMEOUT_MS = 15_000
341346
const DETACHED_CHAT_RETRY_BASE_MS = 1000
342347
const DETACHED_CHAT_RETRY_MAX_MS = 30_000
@@ -4359,7 +4364,9 @@ export function useChat(
43594364
whichever one they opened next. Only a send an unmount withdrew from a
43604365
chatless surface, whose key dies with the mount, goes to the
43614366
cross-surface lanes. */
4362-
const chatless = activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)
4367+
/** The new-chat queue may have moved to its chat while the POST was out. */
4368+
const { chatKey: requeueKey, index: requeueIndex } = liveQueuePosition(activeChatKey, [])
4369+
const chatless = requeueKey.startsWith(PENDING_CHAT_KEY_PREFIX)
43634370
if (result.reason === 'withdrawn' && chatless) {
43644371
handOffWithdrawnSend({ ...payload, userMessageId: result.userMessageId })
43654372
return
@@ -4368,7 +4375,7 @@ export function useChat(
43684375
it, so anything queued while its POST was out was written after it. The one
43694376
exception is a held send adopted from a dead mount of this surface in that
43704377
window, which can be older; it lands behind this one. */
4371-
useMothershipQueueStore.getState().insertAt(activeChatKey, 0, {
4378+
useMothershipQueueStore.getState().insertAt(requeueKey, requeueIndex, {
43724379
...createQueuedMessage(payload, result.userMessageId),
43734380
...requeuedFields(result.reason, 0, chatless ? heldSendSurface : undefined),
43744381
admissionUnknown: result.admissionUnknown,
@@ -4899,8 +4906,10 @@ export function useChat(
48994906
const dispatchChatKey = chatKeyRef.current
49004907
const queueAtStart =
49014908
useMothershipQueueStore.getState().queues[dispatchChatKey] ?? EMPTY_MESSAGE_QUEUE
4902-
let originalIndex = queueAtStart.findIndex((queued) => queued.id === msg.id)
4903-
if (originalIndex === -1) {
4909+
const startIndex = queueAtStart.findIndex((queued) => queued.id === msg.id)
4910+
/** What was queued ahead of it, which it goes back behind if it is restored. */
4911+
let aheadIds = queueAtStart.slice(0, Math.max(0, startIndex)).map((queued) => queued.id)
4912+
if (startIndex === -1) {
49044913
queuedMessageDispatchIds.delete(msg.id)
49054914
return
49064915
}
@@ -4913,7 +4922,7 @@ export function useChat(
49134922
return
49144923
}
49154924
removedFromQueue = true
4916-
useMothershipQueueStore.getState().remove(dispatchChatKey, msg.id)
4925+
useMothershipQueueStore.getState().remove(liveQueueKey(dispatchChatKey), msg.id)
49174926
}
49184927

49194928
/* What actually went out. `msg` is the snapshot from when the dispatch was
@@ -4925,7 +4934,13 @@ export function useChat(
49254934
withdrawn?: WithdrawnSendResult
49264935
) => {
49274936
const withdrawnUserMessageId = withdrawn?.userMessageId
4928-
const chatless = dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)
4937+
/* The send may have waited on a Stop that saw the new chat's first message
4938+
admitted, which moved this queue to that chat. */
4939+
const { chatKey: restoreKey, index: restoreIndex } = liveQueuePosition(
4940+
dispatchChatKey,
4941+
aheadIds
4942+
)
4943+
const chatless = restoreKey.startsWith(PENDING_CHAT_KEY_PREFIX)
49294944
const savedHandoff = readQueuedSendHandoffState()
49304945
const retainedHandoff =
49314946
savedHandoff?.id === msg.id
@@ -4972,7 +4987,7 @@ export function useChat(
49724987
}
49734988
/** Once restored, the queue owns recovery; a second handoff reader must not resend it. */
49744989
clearQueuedSendHandoffState(msg.id)
4975-
useMothershipQueueStore.getState().insertAt(dispatchChatKey, originalIndex, {
4990+
useMothershipQueueStore.getState().insertAt(restoreKey, restoreIndex, {
49764991
/* Only this outcome's policy applies: what an earlier one set (a hold, a
49774992
retry delay, a surface) must not outlive it. */
49784993
...withoutRequeueFields(dispatched),
@@ -4996,7 +5011,7 @@ export function useChat(
49965011
if (currentIndex === -1) {
49975012
return
49985013
}
4999-
originalIndex = currentIndex
5014+
aheadIds = queueAtSend.slice(0, currentIndex).map((queued) => queued.id)
50005015

50015016
// Re-read live: the user may have applied an in-place edit (`replaceAt`)
50025017
// between dispatch scheduling and this send.
@@ -5057,12 +5072,12 @@ export function useChat(
50575072
if (!history) {
50585073
useMothershipQueueStore
50595074
.getState()
5060-
.deferRetry(chatKey, msg.id, sendRetry((msg.sendRetries ?? 0) + 1))
5075+
.deferRetry(liveQueueKey(chatKey), msg.id, sendRetry((msg.sendRetries ?? 0) + 1))
50615076
return true
50625077
}
50635078
clearQueuedSendHandoffState(msg.id)
50645079
clearQueuedSendHandoffClaim(msg.id)
5065-
useMothershipQueueStore.getState().remove(chatKey, msg.id)
5080+
useMothershipQueueStore.getState().remove(liveQueueKey(chatKey), msg.id)
50665081
return true
50675082
},
50685083
[queryClient]

‎apps/sim/lib/mothership/events.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { createLogger } from '@sim/logger'
22
import type { WorkspaceSearchFilters } from '@/lib/api/contracts/knowledge/search'
33
import type { AssistantSearchLevel } from '@/lib/mothership/generated/assistant'
4+
import { sendPayload } from '@/app/workspace/[workspaceId]/home/hooks/send-queue-policy'
45
import type {
56
ChatRequestMode,
67
FileAttachmentForApi,
@@ -56,7 +57,7 @@ export interface MothershipSendMessageDetail {
5657
* this to decide whether to persist a handoff instead.
5758
*/
5859
export function sendMothershipMessage(payload: SendPayload, resumeUserMessageId?: string): boolean {
59-
const { content, ...payloadFields } = payload
60+
const { content, ...payloadFields } = sendPayload(payload)
6061
const trimmed = content.trim()
6162
if (!trimmed && !payloadFields.fileAttachments?.length) {
6263
logger.warn('sendMothershipMessage called with empty message')

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

Lines changed: 49 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,9 @@
11
import { beforeEach, describe, expect, it } from 'vitest'
2-
import { useMothershipQueueStore } from '@/stores/mothership-queue/store'
2+
import {
3+
liveQueueKey,
4+
liveQueuePosition,
5+
useMothershipQueueStore,
6+
} from '@/stores/mothership-queue/store'
37
import type { QueuedMothershipMessage } from '@/stores/mothership-queue/types'
48

59
const message = (id: string, content = `content-${id}`): QueuedMothershipMessage => ({
@@ -181,6 +185,50 @@ describe('useMothershipQueueStore', () => {
181185
})
182186

183187
describe('migrate', () => {
188+
it('points a late write at the chat a new-chat queue moved to, even an empty one', () => {
189+
useMothershipQueueStore.getState().migrate('pending::empty', 'chat-X')
190+
useMothershipQueueStore.getState().migrate('chat-X', 'chat-X')
191+
192+
expect(liveQueueKey('pending::empty')).toBe('chat-X')
193+
expect(liveQueueKey('pending::never-moved')).toBe('pending::never-moved')
194+
expect(liveQueueKey('chat-X')).toBe('chat-X')
195+
})
196+
197+
it('keeps a late write behind the messages the chat queue already held', () => {
198+
useMothershipQueueStore.getState().enqueue('chat-Y', message('older-1'))
199+
useMothershipQueueStore.getState().enqueue('chat-Y', message('older-2'))
200+
useMothershipQueueStore.getState().enqueue('pending::moved', message('moved'))
201+
useMothershipQueueStore.getState().migrate('pending::moved', 'chat-Y')
202+
203+
const position = liveQueuePosition('pending::moved', [])
204+
useMothershipQueueStore.getState().insertAt(position.chatKey, position.index, message('late'))
205+
206+
expect(position).toEqual({ chatKey: 'chat-Y', index: 2 })
207+
expect(useMothershipQueueStore.getState().queues['chat-Y']?.map((m) => m.id)).toEqual([
208+
'older-1',
209+
'older-2',
210+
'late',
211+
'moved',
212+
])
213+
})
214+
215+
it('keeps a late write in order when messages ahead of it were removed meanwhile', () => {
216+
useMothershipQueueStore.getState().enqueue('chat-Z', message('older'))
217+
useMothershipQueueStore.getState().enqueue('pending::sent', message('later'))
218+
useMothershipQueueStore.getState().migrate('pending::sent', 'chat-Z')
219+
useMothershipQueueStore.getState().remove('chat-Z', 'older')
220+
221+
const position = liveQueuePosition('pending::sent', [])
222+
useMothershipQueueStore
223+
.getState()
224+
.insertAt(position.chatKey, position.index, message('follow-up'))
225+
226+
expect(useMothershipQueueStore.getState().queues['chat-Z']?.map((m) => m.id)).toEqual([
227+
'follow-up',
228+
'later',
229+
])
230+
})
231+
184232
it('merges into an existing destination bucket instead of overwriting', () => {
185233
useMothershipQueueStore.getState().enqueue('chat-X', message('existing-1'))
186234
useMothershipQueueStore.getState().enqueue('chat-X', message('existing-2'))

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

Lines changed: 54 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,11 @@ import { toError } from '@sim/utils/errors'
33
import { toRecord, toRecordOrNull } from '@sim/utils/object'
44
import { create } from 'zustand'
55
import { createJSONStorage, devtools, persist } from 'zustand/middleware'
6-
import type { MothershipQueueState, QueuedMothershipMessage } from '@/stores/mothership-queue/types'
6+
import type {
7+
MothershipQueueState,
8+
QueuedMothershipMessage,
9+
QueueMigration,
10+
} from '@/stores/mothership-queue/types'
711

812
const logger = createLogger('MothershipQueueStore')
913

@@ -49,6 +53,7 @@ const initialState = {
4953
queues: {} as Record<string, QueuedMothershipMessage[]>,
5054
editing: {} as Record<string, string>,
5155
cleared: {} as Record<string, number>,
56+
migratedTo: {} as Record<string, QueueMigration>,
5257
}
5358

5459
/**
@@ -106,6 +111,45 @@ const setQueueForChat = (
106111
): Record<string, QueuedMothershipMessage[]> =>
107112
next.length === 0 ? omitKey(queues, chatKey) : { ...queues, [chatKey]: next }
108113

114+
/**
115+
* The queue key a write captured before an `await` should use now: the key a
116+
* new-chat queue migrated to once its chat became known, if it did.
117+
*/
118+
export function liveQueueKey(chatKey: string): string {
119+
const { migratedTo } = useMothershipQueueStore.getState()
120+
let key = chatKey
121+
for (let hops = 0; hops < 8 && migratedTo[key] !== undefined; hops++) key = migratedTo[key].key
122+
return key
123+
}
124+
125+
/**
126+
* Where a message goes back into its queue after a write captured before an
127+
* `await`: in the queue's live key, right after the last message still there
128+
* that was ahead of it (`aheadIds`, plus whatever a chat's queue already held
129+
* when a new-chat queue moved into it), else at the head. Anchoring on ids, not
130+
* an index, keeps it in order however the queue changed meanwhile.
131+
*/
132+
export function liveQueuePosition(
133+
chatKey: string,
134+
aheadIds: readonly string[]
135+
): { chatKey: string; index: number } {
136+
const { migratedTo, queues } = useMothershipQueueStore.getState()
137+
const ahead = new Set(aheadIds)
138+
let key = chatKey
139+
for (let hops = 0; hops < 8; hops++) {
140+
const migration = migratedTo[key]
141+
if (!migration) break
142+
for (const id of migration.ahead) ahead.add(id)
143+
key = migration.key
144+
}
145+
const queue = queues[key] ?? []
146+
let index = 0
147+
queue.forEach((message, position) => {
148+
if (ahead.has(message.id)) index = position + 1
149+
})
150+
return { chatKey: key, index }
151+
}
152+
109153
export const useMothershipQueueStore = create<MothershipQueueState>()(
110154
devtools(
111155
persist(
@@ -191,9 +235,16 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
191235
migrate: (fromKey, toKey) =>
192236
set((state) => {
193237
if (fromKey === toKey) return state
238+
const migratedTo = {
239+
...state.migratedTo,
240+
[fromKey]: {
241+
key: toKey,
242+
ahead: (state.queues[toKey] ?? []).map((message) => message.id),
243+
},
244+
}
194245
const fromQueue = state.queues[fromKey]
195246
const fromEditing = state.editing[fromKey]
196-
if (!fromQueue && fromEditing === undefined) return state
247+
if (!fromQueue && fromEditing === undefined) return { migratedTo }
197248

198249
const queues = omitKey(state.queues, fromKey)
199250
/** A chat deleted meanwhile takes nothing: its queue is gone with it. */
@@ -211,7 +262,7 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
211262
if (fromEditing !== undefined) {
212263
editing[toKey] = fromEditing
213264
}
214-
return { queues, editing }
265+
return { queues, editing, migratedTo }
215266
}),
216267

217268
releaseHeldUntilOnline: () =>

0 commit comments

Comments
 (0)