Skip to content

Commit 75ae554

Browse files
committed
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).
1 parent 79615cb commit 75ae554

4 files changed

Lines changed: 70 additions & 16 deletions

File tree

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

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -136,6 +136,7 @@ import { useChatPanelStore } from '@/stores/chat-panel/store'
136136
import { useMothershipEffortStore } from '@/stores/mothership-effort/store'
137137
import {
138138
liveQueueKey,
139+
liveQueuePosition,
139140
reusedRequestId,
140141
useMothershipQueueStore,
141142
} from '@/stores/mothership-queue/store'
@@ -4364,7 +4365,7 @@ export function useChat(
43644365
chatless surface, whose key dies with the mount, goes to the
43654366
cross-surface lanes. */
43664367
/** The new-chat queue may have moved to its chat while the POST was out. */
4367-
const requeueKey = liveQueueKey(activeChatKey)
4368+
const { chatKey: requeueKey, index: requeueIndex } = liveQueuePosition(activeChatKey, 0)
43684369
const chatless = requeueKey.startsWith(PENDING_CHAT_KEY_PREFIX)
43694370
if (result.reason === 'withdrawn' && chatless) {
43704371
handOffWithdrawnSend({ ...payload, userMessageId: result.userMessageId })
@@ -4374,7 +4375,7 @@ export function useChat(
43744375
it, so anything queued while its POST was out was written after it. The one
43754376
exception is a held send adopted from a dead mount of this surface in that
43764377
window, which can be older; it lands behind this one. */
4377-
useMothershipQueueStore.getState().insertAt(requeueKey, 0, {
4378+
useMothershipQueueStore.getState().insertAt(requeueKey, requeueIndex, {
43784379
...createQueuedMessage(payload, result.userMessageId),
43794380
...requeuedFields(result.reason, 0, chatless ? heldSendSurface : undefined),
43804381
admissionUnknown: result.admissionUnknown,
@@ -4933,7 +4934,10 @@ export function useChat(
49334934
const withdrawnUserMessageId = withdrawn?.userMessageId
49344935
/* The send may have waited on a Stop that saw the new chat's first message
49354936
admitted, which moved this queue to that chat. */
4936-
const restoreKey = liveQueueKey(dispatchChatKey)
4937+
const { chatKey: restoreKey, index: restoreIndex } = liveQueuePosition(
4938+
dispatchChatKey,
4939+
originalIndex
4940+
)
49374941
const chatless = restoreKey.startsWith(PENDING_CHAT_KEY_PREFIX)
49384942
const savedHandoff = readQueuedSendHandoffState()
49394943
const retainedHandoff =
@@ -4981,7 +4985,7 @@ export function useChat(
49814985
}
49824986
/** Once restored, the queue owns recovery; a second handoff reader must not resend it. */
49834987
clearQueuedSendHandoffState(msg.id)
4984-
useMothershipQueueStore.getState().insertAt(restoreKey, originalIndex, {
4988+
useMothershipQueueStore.getState().insertAt(restoreKey, restoreIndex, {
49854989
/* Only this outcome's policy applies: what an earlier one set (a hold, a
49864990
retry delay, a surface) must not outlive it. */
49874991
...withoutRequeueFields(dispatched),

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

Lines changed: 23 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 { liveQueueKey, 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 => ({
@@ -190,6 +194,24 @@ describe('useMothershipQueueStore', () => {
190194
expect(liveQueueKey('chat-X')).toBe('chat-X')
191195
})
192196

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', 0)
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+
193215
it('merges into an existing destination bucket instead of overwriting', () => {
194216
useMothershipQueueStore.getState().enqueue('chat-X', message('existing-1'))
195217
useMothershipQueueStore.getState().enqueue('chat-X', message('existing-2'))

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

Lines changed: 29 additions & 9 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,7 +53,7 @@ const initialState = {
4953
queues: {} as Record<string, QueuedMothershipMessage[]>,
5054
editing: {} as Record<string, string>,
5155
cleared: {} as Record<string, number>,
52-
migratedTo: {} as Record<string, string>,
56+
migratedTo: {} as Record<string, QueueMigration>,
5357
}
5458

5559
/**
@@ -108,14 +112,27 @@ const setQueueForChat = (
108112
next.length === 0 ? omitKey(queues, chatKey) : { ...queues, [chatKey]: next }
109113

110114
/**
111-
* The queue key a write captured before an `await` should use now: the key a
112-
* new-chat queue migrated to once its chat became known, if it did.
115+
* Where a position a write captured before an `await` lies now: in the queue a
116+
* new-chat queue migrated to once its chat became known, if it did, behind the
117+
* messages that queue already held.
113118
*/
114-
export function liveQueueKey(chatKey: string): string {
119+
export function liveQueuePosition(
120+
chatKey: string,
121+
index: number
122+
): { chatKey: string; index: number } {
115123
const { migratedTo } = useMothershipQueueStore.getState()
116-
let key = chatKey
117-
for (let hops = 0; hops < 8 && migratedTo[key] !== undefined; hops++) key = migratedTo[key]
118-
return key
124+
let position = { chatKey, index }
125+
for (let hops = 0; hops < 8; hops++) {
126+
const migration = migratedTo[position.chatKey]
127+
if (!migration) break
128+
position = { chatKey: migration.key, index: position.index + migration.behind }
129+
}
130+
return position
131+
}
132+
133+
/** The queue key a write captured before an `await` should use now (see `liveQueuePosition`). */
134+
export function liveQueueKey(chatKey: string): string {
135+
return liveQueuePosition(chatKey, 0).chatKey
119136
}
120137

121138
export const useMothershipQueueStore = create<MothershipQueueState>()(
@@ -203,7 +220,10 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
203220
migrate: (fromKey, toKey) =>
204221
set((state) => {
205222
if (fromKey === toKey) return state
206-
const migratedTo = { ...state.migratedTo, [fromKey]: toKey }
223+
const migratedTo = {
224+
...state.migratedTo,
225+
[fromKey]: { key: toKey, behind: state.queues[toKey]?.length ?? 0 },
226+
}
207227
const fromQueue = state.queues[fromKey]
208228
const fromEditing = state.editing[fromKey]
209229
if (!fromQueue && fromEditing === undefined) return { migratedTo }

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

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,13 @@ export type QueuedMessageEditPatch = Pick<
5656
| 'assistantSearchLevel'
5757
>
5858

59+
/** A new-chat queue's move to its chat's key. */
60+
export interface QueueMigration {
61+
key: string
62+
/** Messages the chat's queue already held, which stay ahead of the moved ones. */
63+
behind: number
64+
}
65+
5966
export interface MothershipQueueState {
6067
queues: Record<string, QueuedMothershipMessage[]>
6168
editing: Record<string, string>
@@ -68,9 +75,10 @@ export interface MothershipQueueState {
6875
/**
6976
* Where each new-chat key's queue moved when its chat became known
7077
* (`migrate`). A write that captured the old key before an `await` follows
71-
* this (`liveQueueKey`), so it lands in the chat's queue, not a dead key.
78+
* this (`liveQueueKey`, `liveQueuePosition`), so it lands in the chat's queue,
79+
* not a dead key.
7280
*/
73-
migratedTo: Record<string, string>
81+
migratedTo: Record<string, QueueMigration>
7482

7583
enqueue: (chatKey: string, message: QueuedMothershipMessage) => void
7684
insertAt: (chatKey: string, index: number, message: QueuedMothershipMessage) => void

0 commit comments

Comments
 (0)