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 @@ -8,8 +8,7 @@ import {
describe('requeuedFields', () => {
it('holds an offline send for the network, on its chatless surface', () => {
expect(requeuedFields('offline', 0, 'ws-1:home')).toEqual({
retryRequired: true,
heldUntilOnline: true,
hold: 'online',
heldSurface: 'ws-1:home',
})
})
Expand All @@ -18,20 +17,20 @@ describe('requeuedFields', () => {
const before = Date.now()
const fields = requeuedFields(reason, 2, undefined)

expect(fields.sendRetries).toBe(3)
expect(fields.notBefore).toBeGreaterThan(before)
expect(fields.retryRequired).toBeUndefined()
expect(fields.retry?.attempt).toBe(3)
expect(fields.retry?.notBefore).toBeGreaterThan(before)
expect(fields.hold).toBeUndefined()
expect(fields.heldSurface).toBeUndefined()
})

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

Expand All @@ -48,10 +47,8 @@ describe('withoutRequeueFields', () => {
content: 'hello',
resumeUserMessageId: 'attempt-1',
admissionUnknown: true,
retryRequired: true,
heldUntilOnline: true,
sendRetries: 2,
notBefore: 123,
hold: 'online',
retry: { attempt: 2, notBefore: 123 },
heldSurface: 'ws-1:home',
})
).toEqual({
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { backoffWithJitter } from '@sim/utils/retry'
import type { SendPayload } from '@/app/workspace/[workspaceId]/home/types'
import type { QueuedMothershipMessage, ScheduledRetry } from '@/stores/mothership-queue/types'
import type { QueuedMothershipMessage, SendRetry } from '@/stores/mothership-queue/types'

/**
* Why a send came back to its caller instead of going out:
Expand All @@ -17,18 +17,15 @@ export type WithdrawalReason = 'withdrawn' | 'offline' | 'unreachable' | 'busy'
export type RequeueReason = WithdrawalReason | 'failed'

/** The queue fields that say when, and on which surface, a re-queued message goes out. */
type RequeueFields = Pick<
QueuedMothershipMessage,
'retryRequired' | 'heldUntilOnline' | 'sendRetries' | 'notBefore' | 'heldSurface'
>
type RequeueFields = Pick<QueuedMothershipMessage, 'hold' | 'retry' | 'heldSurface'>

const SEND_RETRY_BASE_MS = 1_000
const SEND_RETRY_MAX_MS = 30_000

/** Queue fields for the `attempt`th automatic retry of a message: when it may be sent again. */
export function sendRetry(attempt: number): ScheduledRetry {
/** The `attempt`th automatic retry of a message: when it may be sent again. */
export function sendRetry(attempt: number): SendRetry {
return {
sendRetries: attempt,
attempt,
notBefore:
Date.now() +
backoffWithJitter(attempt, null, { baseMs: SEND_RETRY_BASE_MS, maxMs: SEND_RETRY_MAX_MS }),
Expand All @@ -55,32 +52,25 @@ export function requeuedFields(
const surface = chatlessSurface ? { heldSurface: chatlessSurface } : {}
switch (reason) {
case 'offline':
return { retryRequired: true, heldUntilOnline: true, ...surface }
return { hold: 'online', ...surface }
case 'unreachable':
case 'busy':
return { ...sendRetry(previousAttempts + 1), ...surface }
return { retry: sendRetry(previousAttempts + 1), ...surface }
case 'stop-failed':
case 'failed':
return { retryRequired: true, ...surface }
return { hold: 'user', ...surface }
case 'withdrawn':
return {}
}
}

/**
* A queue entry without the fields an earlier outcome set, so a re-queue applies
* only the policy for the outcome it is handling. A stale `heldUntilOnline`, for
* only the policy for the outcome it is handling. A stale `online` hold, for
* one, would let the browser coming online send a message waiting for the user.
*/
export function withoutRequeueFields(entry: QueuedMothershipMessage): QueuedMothershipMessage {
const {
retryRequired: _retryRequired,
heldUntilOnline: _heldUntilOnline,
sendRetries: _sendRetries,
notBefore: _notBefore,
heldSurface: _heldSurface,
...rest
} = entry
const { hold: _hold, retry: _retry, heldSurface: _heldSurface, ...rest } = entry
return rest
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1769,15 +1769,15 @@ describe('useChat remount send recovery', () => {
await act(async () => {
await sending
})
await waitFor(() => allQueuedMessages().some((message) => message.retryRequired === true))
await waitFor(() => allQueuedMessages().some((message) => message.hold === 'user'))
expect(state.postBodies).toHaveLength(1)
expect(allQueuedMessages()).toEqual([
expect.objectContaining({ id: queued.id, content: queued.content }),
])
expect(getResult().error).toBe('Previous response is still shutting down.')
const failed = allQueuedMessages()[0]
expect(failed).toMatchObject({
retryRequired: true,
hold: 'user',
queuedSendHandoff: {
stopRequired: true,
supersededStreamId: state.postBodies[0].userMessageId,
Expand Down Expand Up @@ -1966,7 +1966,7 @@ describe('useChat remount send recovery', () => {
expect(allQueuedMessages()).toEqual([
expect.objectContaining({
id: 'queued-correction',
retryRequired: true,
hold: 'user',
queuedSendHandoff: expect.objectContaining({
userMessageId: 'prepared-correction-request',
supersededStreamId: 'previous-response',
Expand Down Expand Up @@ -1998,7 +1998,7 @@ describe('useChat remount send recovery', () => {
useMothershipQueueStore.getState().enqueue('chat-a', {
id: 'earlier-correction',
content: 'inspect the second invoice instead',
retryRequired: true,
hold: 'user',
queuedSendHandoff: {
id: 'earlier-correction',
chatId: 'chat-a',
Expand Down Expand Up @@ -2038,7 +2038,7 @@ describe('useChat remount send recovery', () => {
expect(state.abortBodies[0]?.streamId).toBe(newerStreamId)
expect(state.postBodies).toHaveLength(1)
expect(allQueuedMessages()[0]).toMatchObject({
retryRequired: true,
hold: 'user',
queuedSendHandoff: {
supersededStreamId: newerStreamId,
userMessageId: 'prepared-correction',
Expand Down Expand Up @@ -2358,8 +2358,7 @@ describe('useChat remount send recovery', () => {
content: 'written while offline',
resumeUserMessageId: 'offline-attempt',
admissionUnknown: true,
retryRequired: true,
heldUntilOnline: true,
hold: 'online',
})

await act(async () => {
Expand All @@ -2375,8 +2374,7 @@ describe('useChat remount send recovery', () => {

expect(state.postBodies).toHaveLength(1)
const queued = useMothershipQueueStore.getState().queues['chat-a']?.[0]
expect(queued).toMatchObject({ id: 'held-offline', retryRequired: true })
expect(queued?.heldUntilOnline).toBeUndefined()
expect(queued).toMatchObject({ id: 'held-offline', hold: 'user' })
})

it('keeps a resumed message uneditable when its Send-now Stop does not settle', async () => {
Expand Down Expand Up @@ -2868,9 +2866,12 @@ describe('useChat remount send recovery', () => {
.catch(() => {})
})
await waitFor(() => state.postBodies.length === 2)
await act(async () => {
await sleep(200)
})
/** The refusal goes back to a queue; the assertions below say which one. */
await waitFor(() =>
Object.values(useMothershipQueueStore.getState().queues).some((queue) =>
queue.some((message) => message.content === 'Follow-up')
)
)

const queues = useMothershipQueueStore.getState().queues
expect(queues[DEDUPED_CHAT_ID]?.map((message) => message.content)).toEqual(['Follow-up'])
Expand Down Expand Up @@ -3192,7 +3193,7 @@ describe('useChat remount send recovery', () => {

const queued = useMothershipQueueStore.getState().queues[history.id] ?? []
expect(queued.map((message) => message.content)).toEqual(['Written while offline'])
expect(queued[0].retryRequired).toBe(true)
expect(queued[0].hold).toBe('online')
expect(queued[0].resumeUserMessageId).toBe(state.postBodies[0].userMessageId)
expect(getResult().error).not.toBeNull()
expect(state.postBodies).toHaveLength(1)
Expand Down Expand Up @@ -3343,8 +3344,7 @@ describe('useChat remount send recovery', () => {
expect(state.postBodies).toHaveLength(0)
expect(useMothershipQueueStore.getState().queues[history.id]?.[0]).toMatchObject({
content: 'Never prepared',
retryRequired: true,
heldUntilOnline: true,
hold: 'online',
})
}
)
Expand Down Expand Up @@ -3435,7 +3435,7 @@ describe('useChat remount send recovery', () => {

await waitFor(() => state.postBodies.length === 1)
await waitFor(
() => useMothershipQueueStore.getState().queues[history.id]?.[0]?.retryRequired === true
() => useMothershipQueueStore.getState().queues[history.id]?.[0]?.hold !== undefined
)

const queued = useMothershipQueueStore.getState().queues[history.id] ?? []
Expand Down Expand Up @@ -3530,7 +3530,7 @@ describe('useChat remount send recovery', () => {
await act(async () => {
await first.getResult().sendMessage('First message, sent offline')
})
await waitFor(() => allQueuedMessages().some((message) => message.retryRequired === true))
await waitFor(() => allQueuedMessages().some((message) => message.hold !== undefined))
first.unmount()

const second = renderUseChat()
Expand Down Expand Up @@ -4307,7 +4307,7 @@ describe('useChat remount send recovery', () => {
await first.getResult().sendMessage('Held while I was elsewhere')
})
await waitFor(
() => useMothershipQueueStore.getState().queues[history.id]?.[0]?.retryRequired === true
() => useMothershipQueueStore.getState().queues[history.id]?.[0]?.hold !== undefined
)
first.unmount()

Expand Down Expand Up @@ -4354,7 +4354,7 @@ describe('useChat remount send recovery', () => {
id: 'withdrawn-entry',
content: 'check the trace for this req',
resumeUserMessageId: 'accepted-request',
retryRequired: true,
hold: 'user',
})
const { getResult } = renderUseChatInChat('chat-a', {
id: 'chat-a',
Expand Down Expand Up @@ -4391,7 +4391,7 @@ describe('useChat remount send recovery', () => {
useMothershipQueueStore.getState().enqueue('chat-a', {
id: 'unsent-entry',
content: 'check the trace for this req',
retryRequired: true,
hold: 'user',
resumeUserMessageId: 'unsent-request',
})
const { getResult } = renderUseChatInChat('chat-a', {
Expand Down
14 changes: 7 additions & 7 deletions apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4994,7 +4994,7 @@ export function useChat(
...(retainedHandoff ? { queuedSendHandoff: retainedHandoff } : {}),
...requeuedFields(
withdrawn?.reason ?? 'failed',
dispatched.sendRetries ?? 0,
dispatched.retry?.attempt ?? 0,
chatless ? heldSendSurface : undefined
),
...(withdrawnUserMessageId ? { resumeUserMessageId: withdrawnUserMessageId } : {}),
Expand Down Expand Up @@ -5072,12 +5072,12 @@ export function useChat(
if (!history) {
useMothershipQueueStore
.getState()
.deferRetry(liveQueueKey(chatKey), msg.id, sendRetry((msg.sendRetries ?? 0) + 1))
.deferRetry(chatKey, msg.id, sendRetry((msg.retry?.attempt ?? 0) + 1))
return true
}
clearQueuedSendHandoffState(msg.id)
clearQueuedSendHandoffClaim(msg.id)
useMothershipQueueStore.getState().remove(liveQueueKey(chatKey), msg.id)
useMothershipQueueStore.getState().remove(chatKey, msg.id)
return true
},
[queryClient]
Expand All @@ -5101,9 +5101,9 @@ export function useChat(
const queueState = useMothershipQueueStore.getState()
const activeChatKey = chatKeyRef.current
const msg = queueState.queues[activeChatKey]?.[0]
if (!msg || msg.retryRequired) continue
if (!msg || msg.hold) continue
/** An automatic retry waits out its delay; the drain effect wakes it. */
if (msg.notBefore !== undefined && msg.notBefore > Date.now()) continue
if (msg.retry && msg.retry.notBefore > Date.now()) continue
// Pause draining if the head is bound to the composer; dispatching now
// would race the eventual submit. The next kick on edit-resolve resumes us.
if (queueState.editing[activeChatKey] === msg.id) continue
Expand Down Expand Up @@ -5276,8 +5276,8 @@ export function useChat(
// `notifyTurnEnded`. Idempotent — the dispatch loop dedupes.
const chatHistoryReady = chatHistory !== undefined
const remoteActiveStreamId = chatHistory?.activeStreamId ?? null
const queueHeadHeld = messageQueue[0]?.retryRequired === true
const queueHeadNotBefore = messageQueue[0]?.notBefore
const queueHeadHeld = messageQueue[0]?.hold !== undefined
const queueHeadNotBefore = messageQueue[0]?.retry?.notBefore
const [sendRetryWakeup, setSendRetryWakeup] = useState(0)
useEffect(() => {
if (!scopeKey) return
Expand Down
28 changes: 28 additions & 0 deletions apps/sim/stores/mothership-queue/store.dom.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,34 @@ describe('useMothershipQueueStore rehydration', () => {
sessionStorage.clear()
})

it('restores the hold and retry fields of a queue saved in their older shape', async () => {
sessionStorage.setItem(
'mothership-queue',
JSON.stringify({
state: {
queues: {
'chat-A': [
{ id: 'for-user', content: 'a', retryRequired: true },
{ id: 'for-network', content: 'b', retryRequired: true, heldUntilOnline: true },
{ id: 'retrying', content: 'c', sendRetries: 2, notBefore: 1_000 },
{ id: 'plain', content: 'd' },
],
},
},
version: 0,
})
)

await useMothershipQueueStore.persist.rehydrate()

expect(useMothershipQueueStore.getState().queues['chat-A']).toEqual([
{ id: 'for-user', content: 'a', hold: 'user' },
{ id: 'for-network', content: 'b', hold: 'online' },
{ id: 'retrying', content: 'c', retry: { attempt: 2, notBefore: 1_000 } },
{ id: 'plain', content: 'd' },
])
})

it('treats a resumed message saved before the edit guard as possibly sent', async () => {
sessionStorage.setItem(
'mothership-queue',
Expand Down
22 changes: 20 additions & 2 deletions apps/sim/stores/mothership-queue/store.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -126,7 +126,7 @@ describe('useMothershipQueueStore', () => {
useMothershipQueueStore.getState().enqueue('chat-A', {
id: 'm1',
content: 'original',
retryRequired: true,
hold: 'user',
resumeUserMessageId: 'prior-request',
/** Its Stop never settled, so it was never sent: the server cannot hold it. */
admissionUnknown: false,
Expand All @@ -152,7 +152,7 @@ describe('useMothershipQueueStore', () => {
})
expect(edited?.queuedSendHandoff?.userMessageId).toBeUndefined()
expect(edited?.resumeUserMessageId).toBeUndefined()
expect(edited?.retryRequired).toBeUndefined()
expect(edited?.hold).toBeUndefined()
})

it('strips queuedSendHandoff on edit so a fresh handoff is minted at send time', () => {
Expand Down Expand Up @@ -212,6 +212,24 @@ describe('useMothershipQueueStore', () => {
])
})

it('keeps the first move when the same new-chat queue is migrated again', () => {
useMothershipQueueStore.getState().enqueue('chat-W', message('older'))
useMothershipQueueStore.getState().enqueue('pending::twice', message('moved'))
useMothershipQueueStore.getState().migrate('pending::twice', 'chat-W')
useMothershipQueueStore.getState().enqueue('chat-W', message('written-later'))
useMothershipQueueStore.getState().migrate('pending::twice', 'chat-W')

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

expect(useMothershipQueueStore.getState().queues['chat-W']?.map((m) => m.id)).toEqual([
'older',
'late',
'moved',
'written-later',
])
})

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'))
Expand Down
Loading
Loading