Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -25,9 +25,13 @@ describe('requeuedFields', () => {
})

it.each(['stop-failed', 'failed'] as const)(
'leaves a %s send for the user, on any surface',
'leaves a %s send for the user, adoptable by its chatless surface',
(reason) => {
expect(requeuedFields(reason, 4, 'ws-1:home')).toEqual({ retryRequired: true })
expect(requeuedFields(reason, 4, 'ws-1:home')).toEqual({
retryRequired: true,
heldSurface: 'ws-1:home',
})
expect(requeuedFields(reason, 4, undefined)).toEqual({ retryRequired: true })
}
)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,9 +42,10 @@ export function sendRetry(attempt: number): ScheduledRetry {
* - `stop-failed` and `failed` wait for the user;
* - `withdrawn` goes out again as soon as the queue drains.
*
* A message held by a chatless surface carries that surface (`chatlessSurface`),
* whose queue key dies with its mount, so the next mount of it adopts the
* message. Only sends that wait on the network or the server are held that way.
* A message re-queued on a chatless surface carries that surface
* (`chatlessSurface`), whose queue key dies with its mount, so the next mount of
* it adopts the message. That holds for every reason: a re-queue can land after
* the surface unmounted, when nothing else marks the dead key's queue.
*/
export function requeuedFields(
reason: RequeueReason,
Expand All @@ -60,7 +61,7 @@ export function requeuedFields(
return { ...sendRetry(previousAttempts + 1), ...surface }
case 'stop-failed':
case 'failed':
return { retryRequired: true }
return { retryRequired: true, ...surface }
case 'withdrawn':
return {}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2829,6 +2829,57 @@ describe('useChat remount send recovery', () => {
expect(allQueuedMessages()).toHaveLength(0)
})

/**
* Send-now on the new-chat surface stops the first message, which the Stop sees
* admitted into a chat: the surface moves to that chat and its queue moves with
* it. A busy refusal of the follow-up arriving after that must go back to the
* chat's queue, where the surface shows and retries it, not to the dead
* new-chat key it was dispatched from.
*/
it('re-queues a Send-now refused as busy in the chat the new-chat surface moved to', async () => {
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
const url = String(input)
if (
url === '/api/mothership/chat' &&
init?.method === 'POST' &&
state.postBodies.length > 0
) {
state.postBodies.push(JSON.parse(String(init.body)))
return Response.json(
{ error: 'A response is already in progress for this chat.' },
{ status: 409 }
)
}
return fetchStub(input, init)
})
const { getResult } = renderUseChat()
await act(async () => {
void getResult().sendMessage('inspect the workspace')
})
await waitFor(() => state.postBodies.length === 1)
await act(async () => {
void getResult().sendMessage('Follow-up')
})
await waitFor(() => allQueuedMessages().length === 1)

await act(async () => {
void getResult()
.sendNow()
.catch(() => {})
})
await waitFor(() => state.postBodies.length === 2)
await act(async () => {
await sleep(200)
})

const queues = useMothershipQueueStore.getState().queues
expect(queues[DEDUPED_CHAT_ID]?.map((message) => message.content)).toEqual(['Follow-up'])
expect(
Object.entries(queues).filter(([key, queue]) => key.startsWith('pending::') && queue.length)
).toEqual([])
expect(getResult().messageQueue.map((message) => message.content)).toEqual(['Follow-up'])
})

it('stopping a chat preserves an unrelated manual workflow execution', async () => {
const executionStore = useExecutionStore.getState()
executionStore.setIsExecuting('manual-workflow', true)
Expand Down
39 changes: 27 additions & 12 deletions apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -134,7 +134,12 @@ import { workflowKeys } from '@/hooks/queries/workflows'
import { snapAllSmoothText } from '@/hooks/use-smooth-text'
import { useChatPanelStore } from '@/stores/chat-panel/store'
import { useMothershipEffortStore } from '@/stores/mothership-effort/store'
import { reusedRequestId, useMothershipQueueStore } from '@/stores/mothership-queue/store'
import {
liveQueueKey,
liveQueuePosition,
reusedRequestId,
useMothershipQueueStore,
} from '@/stores/mothership-queue/store'
import type {
QueuedMothershipMessage,
QueuedSendHandoffSeed,
Expand Down Expand Up @@ -336,7 +341,7 @@ const PERSISTED_TURN_REFETCH_BASE_MS = 250
const PERSISTED_TURN_REFETCH_MAX_DELAY_MS = 5_000
/** How long a finished turn's save is waited for; a slow save still lands well inside it. */
const PERSISTED_TURN_WAIT_MS = 120_000
/** Pacing for re-sending a message refused because the chat was busy, or that could not reach Sim. */
/** How long a Stop's abort request may take before the Stop counts as failed. */
const STOP_REQUEST_TIMEOUT_MS = 15_000
const DETACHED_CHAT_RETRY_BASE_MS = 1000
const DETACHED_CHAT_RETRY_MAX_MS = 30_000
Expand Down Expand Up @@ -4359,7 +4364,9 @@ export function useChat(
whichever one they opened next. Only a send an unmount withdrew from a
chatless surface, whose key dies with the mount, goes to the
cross-surface lanes. */
const chatless = activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)
/** The new-chat queue may have moved to its chat while the POST was out. */
const { chatKey: requeueKey, index: requeueIndex } = liveQueuePosition(activeChatKey, [])
const chatless = requeueKey.startsWith(PENDING_CHAT_KEY_PREFIX)
if (result.reason === 'withdrawn' && chatless) {
handOffWithdrawnSend({ ...payload, userMessageId: result.userMessageId })
return
Expand All @@ -4368,7 +4375,7 @@ export function useChat(
it, so anything queued while its POST was out was written after it. The one
exception is a held send adopted from a dead mount of this surface in that
window, which can be older; it lands behind this one. */
useMothershipQueueStore.getState().insertAt(activeChatKey, 0, {
useMothershipQueueStore.getState().insertAt(requeueKey, requeueIndex, {
...createQueuedMessage(payload, result.userMessageId),
...requeuedFields(result.reason, 0, chatless ? heldSendSurface : undefined),
admissionUnknown: result.admissionUnknown,
Expand Down Expand Up @@ -4899,8 +4906,10 @@ export function useChat(
const dispatchChatKey = chatKeyRef.current
const queueAtStart =
useMothershipQueueStore.getState().queues[dispatchChatKey] ?? EMPTY_MESSAGE_QUEUE
let originalIndex = queueAtStart.findIndex((queued) => queued.id === msg.id)
if (originalIndex === -1) {
const startIndex = queueAtStart.findIndex((queued) => queued.id === msg.id)
/** What was queued ahead of it, which it goes back behind if it is restored. */
let aheadIds = queueAtStart.slice(0, Math.max(0, startIndex)).map((queued) => queued.id)
if (startIndex === -1) {
queuedMessageDispatchIds.delete(msg.id)
return
}
Expand All @@ -4913,7 +4922,7 @@ export function useChat(
return
}
removedFromQueue = true
useMothershipQueueStore.getState().remove(dispatchChatKey, msg.id)
useMothershipQueueStore.getState().remove(liveQueueKey(dispatchChatKey), msg.id)
}

/* What actually went out. `msg` is the snapshot from when the dispatch was
Expand All @@ -4925,7 +4934,13 @@ export function useChat(
withdrawn?: WithdrawnSendResult
) => {
const withdrawnUserMessageId = withdrawn?.userMessageId
const chatless = dispatchChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)
/* The send may have waited on a Stop that saw the new chat's first message
admitted, which moved this queue to that chat. */
const { chatKey: restoreKey, index: restoreIndex } = liveQueuePosition(
dispatchChatKey,
aheadIds
)
const chatless = restoreKey.startsWith(PENDING_CHAT_KEY_PREFIX)
const savedHandoff = readQueuedSendHandoffState()
const retainedHandoff =
savedHandoff?.id === msg.id
Expand Down Expand Up @@ -4972,7 +4987,7 @@ export function useChat(
}
/** Once restored, the queue owns recovery; a second handoff reader must not resend it. */
clearQueuedSendHandoffState(msg.id)
useMothershipQueueStore.getState().insertAt(dispatchChatKey, originalIndex, {
useMothershipQueueStore.getState().insertAt(restoreKey, restoreIndex, {
/* Only this outcome's policy applies: what an earlier one set (a hold, a
retry delay, a surface) must not outlive it. */
...withoutRequeueFields(dispatched),
Expand All @@ -4996,7 +5011,7 @@ export function useChat(
if (currentIndex === -1) {
return
}
originalIndex = currentIndex
aheadIds = queueAtSend.slice(0, currentIndex).map((queued) => queued.id)

// Re-read live: the user may have applied an in-place edit (`replaceAt`)
// between dispatch scheduling and this send.
Expand Down Expand Up @@ -5057,12 +5072,12 @@ export function useChat(
if (!history) {
useMothershipQueueStore
.getState()
.deferRetry(chatKey, msg.id, sendRetry((msg.sendRetries ?? 0) + 1))
.deferRetry(liveQueueKey(chatKey), msg.id, sendRetry((msg.sendRetries ?? 0) + 1))
return true
}
clearQueuedSendHandoffState(msg.id)
clearQueuedSendHandoffClaim(msg.id)
useMothershipQueueStore.getState().remove(chatKey, msg.id)
useMothershipQueueStore.getState().remove(liveQueueKey(chatKey), msg.id)
return true
},
[queryClient]
Expand Down
3 changes: 2 additions & 1 deletion apps/sim/lib/mothership/events.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { createLogger } from '@sim/logger'
import type { WorkspaceSearchFilters } from '@/lib/api/contracts/knowledge/search'
import type { AssistantSearchLevel } from '@/lib/mothership/generated/assistant'
import { sendPayload } from '@/app/workspace/[workspaceId]/home/hooks/send-queue-policy'
import type {
ChatRequestMode,
FileAttachmentForApi,
Expand Down Expand Up @@ -56,7 +57,7 @@ export interface MothershipSendMessageDetail {
* this to decide whether to persist a handoff instead.
*/
export function sendMothershipMessage(payload: SendPayload, resumeUserMessageId?: string): boolean {
const { content, ...payloadFields } = payload
const { content, ...payloadFields } = sendPayload(payload)
const trimmed = content.trim()
if (!trimmed && !payloadFields.fileAttachments?.length) {
logger.warn('sendMothershipMessage called with empty message')
Expand Down
50 changes: 49 additions & 1 deletion apps/sim/stores/mothership-queue/store.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,9 @@
import { beforeEach, describe, expect, it } from 'vitest'
import { useMothershipQueueStore } from '@/stores/mothership-queue/store'
import {
liveQueueKey,
liveQueuePosition,
useMothershipQueueStore,
} from '@/stores/mothership-queue/store'
import type { QueuedMothershipMessage } from '@/stores/mothership-queue/types'

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

describe('migrate', () => {
it('points a late write at the chat a new-chat queue moved to, even an empty one', () => {
useMothershipQueueStore.getState().migrate('pending::empty', 'chat-X')
useMothershipQueueStore.getState().migrate('chat-X', 'chat-X')

expect(liveQueueKey('pending::empty')).toBe('chat-X')
expect(liveQueueKey('pending::never-moved')).toBe('pending::never-moved')
expect(liveQueueKey('chat-X')).toBe('chat-X')
})

it('keeps a late write behind the messages the chat queue already held', () => {
useMothershipQueueStore.getState().enqueue('chat-Y', message('older-1'))
useMothershipQueueStore.getState().enqueue('chat-Y', message('older-2'))
useMothershipQueueStore.getState().enqueue('pending::moved', message('moved'))
useMothershipQueueStore.getState().migrate('pending::moved', 'chat-Y')

const position = liveQueuePosition('pending::moved', [])
useMothershipQueueStore.getState().insertAt(position.chatKey, position.index, message('late'))

expect(position).toEqual({ chatKey: 'chat-Y', index: 2 })
expect(useMothershipQueueStore.getState().queues['chat-Y']?.map((m) => m.id)).toEqual([
'older-1',
'older-2',
'late',
'moved',
])
})

it('keeps a late write in order when messages ahead of it were removed meanwhile', () => {
useMothershipQueueStore.getState().enqueue('chat-Z', message('older'))
useMothershipQueueStore.getState().enqueue('pending::sent', message('later'))
useMothershipQueueStore.getState().migrate('pending::sent', 'chat-Z')
useMothershipQueueStore.getState().remove('chat-Z', 'older')

const position = liveQueuePosition('pending::sent', [])
useMothershipQueueStore
.getState()
.insertAt(position.chatKey, position.index, message('follow-up'))

expect(useMothershipQueueStore.getState().queues['chat-Z']?.map((m) => m.id)).toEqual([
'follow-up',
'later',
])
})

it('merges into an existing destination bucket instead of overwriting', () => {
useMothershipQueueStore.getState().enqueue('chat-X', message('existing-1'))
useMothershipQueueStore.getState().enqueue('chat-X', message('existing-2'))
Expand Down
57 changes: 54 additions & 3 deletions apps/sim/stores/mothership-queue/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,11 @@ import { toError } from '@sim/utils/errors'
import { toRecord, toRecordOrNull } from '@sim/utils/object'
import { create } from 'zustand'
import { createJSONStorage, devtools, persist } from 'zustand/middleware'
import type { MothershipQueueState, QueuedMothershipMessage } from '@/stores/mothership-queue/types'
import type {
MothershipQueueState,
QueuedMothershipMessage,
QueueMigration,
} from '@/stores/mothership-queue/types'

const logger = createLogger('MothershipQueueStore')

Expand Down Expand Up @@ -49,6 +53,7 @@ const initialState = {
queues: {} as Record<string, QueuedMothershipMessage[]>,
editing: {} as Record<string, string>,
cleared: {} as Record<string, number>,
migratedTo: {} as Record<string, QueueMigration>,
}

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

/**
* The queue key a write captured before an `await` should use now: the key a
* new-chat queue migrated to once its chat became known, if it did.
*/
export function liveQueueKey(chatKey: string): string {
const { migratedTo } = useMothershipQueueStore.getState()
let key = chatKey
for (let hops = 0; hops < 8 && migratedTo[key] !== undefined; hops++) key = migratedTo[key].key
return key
}

/**
* Where a message goes back into its queue after a write captured before an
* `await`: in the queue's live key, right after the last message still there
* that was ahead of it (`aheadIds`, plus whatever a chat's queue already held
* when a new-chat queue moved into it), else at the head. Anchoring on ids, not
* an index, keeps it in order however the queue changed meanwhile.
*/
export function liveQueuePosition(
chatKey: string,
aheadIds: readonly string[]
): { chatKey: string; index: number } {
const { migratedTo, queues } = useMothershipQueueStore.getState()
const ahead = new Set(aheadIds)
let key = chatKey
for (let hops = 0; hops < 8; hops++) {
const migration = migratedTo[key]
if (!migration) break
for (const id of migration.ahead) ahead.add(id)
key = migration.key
}
const queue = queues[key] ?? []
let index = 0
queue.forEach((message, position) => {
if (ahead.has(message.id)) index = position + 1
})
return { chatKey: key, index }
}

export const useMothershipQueueStore = create<MothershipQueueState>()(
devtools(
persist(
Expand Down Expand Up @@ -191,9 +235,16 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
migrate: (fromKey, toKey) =>
set((state) => {
if (fromKey === toKey) return state
const migratedTo = {
...state.migratedTo,
[fromKey]: {
key: toKey,
ahead: (state.queues[toKey] ?? []).map((message) => message.id),
},
}
const fromQueue = state.queues[fromKey]
const fromEditing = state.editing[fromKey]
if (!fromQueue && fromEditing === undefined) return state
if (!fromQueue && fromEditing === undefined) return { migratedTo }

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

releaseHeldUntilOnline: () =>
Expand Down
Loading
Loading