Skip to content

Commit e30a941

Browse files
committed
refactor(mothership): one hold field and one retry field on a queued send
A queued message's wait is now hold ('user' or 'online') instead of retryRequired plus heldUntilOnline, and its automatic retry is retry: { attempt, notBefore } instead of sendRetries plus notBefore, which folds ScheduledRetry away. Queues saved in the older shape are mapped when the session restores them.
1 parent f0e5700 commit e30a941

8 files changed

Lines changed: 116 additions & 81 deletions

File tree

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

Lines changed: 8 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -8,8 +8,7 @@ import {
88
describe('requeuedFields', () => {
99
it('holds an offline send for the network, on its chatless surface', () => {
1010
expect(requeuedFields('offline', 0, 'ws-1:home')).toEqual({
11-
retryRequired: true,
12-
heldUntilOnline: true,
11+
hold: 'online',
1312
heldSurface: 'ws-1:home',
1413
})
1514
})
@@ -18,20 +17,20 @@ describe('requeuedFields', () => {
1817
const before = Date.now()
1918
const fields = requeuedFields(reason, 2, undefined)
2019

21-
expect(fields.sendRetries).toBe(3)
22-
expect(fields.notBefore).toBeGreaterThan(before)
23-
expect(fields.retryRequired).toBeUndefined()
20+
expect(fields.retry?.attempt).toBe(3)
21+
expect(fields.retry?.notBefore).toBeGreaterThan(before)
22+
expect(fields.hold).toBeUndefined()
2423
expect(fields.heldSurface).toBeUndefined()
2524
})
2625

2726
it.each(['stop-failed', 'failed'] as const)(
2827
'leaves a %s send for the user, adoptable by its chatless surface',
2928
(reason) => {
3029
expect(requeuedFields(reason, 4, 'ws-1:home')).toEqual({
31-
retryRequired: true,
30+
hold: 'user',
3231
heldSurface: 'ws-1:home',
3332
})
34-
expect(requeuedFields(reason, 4, undefined)).toEqual({ retryRequired: true })
33+
expect(requeuedFields(reason, 4, undefined)).toEqual({ hold: 'user' })
3534
}
3635
)
3736

@@ -48,10 +47,8 @@ describe('withoutRequeueFields', () => {
4847
content: 'hello',
4948
resumeUserMessageId: 'attempt-1',
5049
admissionUnknown: true,
51-
retryRequired: true,
52-
heldUntilOnline: true,
53-
sendRetries: 2,
54-
notBefore: 123,
50+
hold: 'online',
51+
retry: { attempt: 2, notBefore: 123 },
5552
heldSurface: 'ws-1:home',
5653
})
5754
).toEqual({

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

Lines changed: 10 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import { backoffWithJitter } from '@sim/utils/retry'
22
import type { SendPayload } from '@/app/workspace/[workspaceId]/home/types'
3-
import type { QueuedMothershipMessage, ScheduledRetry } from '@/stores/mothership-queue/types'
3+
import type { QueuedMothershipMessage, SendRetry } from '@/stores/mothership-queue/types'
44

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

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

2522
const SEND_RETRY_BASE_MS = 1_000
2623
const SEND_RETRY_MAX_MS = 30_000
2724

28-
/** Queue fields for the `attempt`th automatic retry of a message: when it may be sent again. */
29-
export function sendRetry(attempt: number): ScheduledRetry {
25+
/** The `attempt`th automatic retry of a message: when it may be sent again. */
26+
export function sendRetry(attempt: number): SendRetry {
3027
return {
31-
sendRetries: attempt,
28+
attempt,
3229
notBefore:
3330
Date.now() +
3431
backoffWithJitter(attempt, null, { baseMs: SEND_RETRY_BASE_MS, maxMs: SEND_RETRY_MAX_MS }),
@@ -55,32 +52,25 @@ export function requeuedFields(
5552
const surface = chatlessSurface ? { heldSurface: chatlessSurface } : {}
5653
switch (reason) {
5754
case 'offline':
58-
return { retryRequired: true, heldUntilOnline: true, ...surface }
55+
return { hold: 'online', ...surface }
5956
case 'unreachable':
6057
case 'busy':
61-
return { ...sendRetry(previousAttempts + 1), ...surface }
58+
return { retry: sendRetry(previousAttempts + 1), ...surface }
6259
case 'stop-failed':
6360
case 'failed':
64-
return { retryRequired: true, ...surface }
61+
return { hold: 'user', ...surface }
6562
case 'withdrawn':
6663
return {}
6764
}
6865
}
6966

7067
/**
7168
* A queue entry without the fields an earlier outcome set, so a re-queue applies
72-
* only the policy for the outcome it is handling. A stale `heldUntilOnline`, for
69+
* only the policy for the outcome it is handling. A stale `online` hold, for
7370
* one, would let the browser coming online send a message waiting for the user.
7471
*/
7572
export function withoutRequeueFields(entry: QueuedMothershipMessage): QueuedMothershipMessage {
76-
const {
77-
retryRequired: _retryRequired,
78-
heldUntilOnline: _heldUntilOnline,
79-
sendRetries: _sendRetries,
80-
notBefore: _notBefore,
81-
heldSurface: _heldSurface,
82-
...rest
83-
} = entry
73+
const { hold: _hold, retry: _retry, heldSurface: _heldSurface, ...rest } = entry
8474
return rest
8575
}
8676

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

Lines changed: 14 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -1769,15 +1769,15 @@ describe('useChat remount send recovery', () => {
17691769
await act(async () => {
17701770
await sending
17711771
})
1772-
await waitFor(() => allQueuedMessages().some((message) => message.retryRequired === true))
1772+
await waitFor(() => allQueuedMessages().some((message) => message.hold === 'user'))
17731773
expect(state.postBodies).toHaveLength(1)
17741774
expect(allQueuedMessages()).toEqual([
17751775
expect.objectContaining({ id: queued.id, content: queued.content }),
17761776
])
17771777
expect(getResult().error).toBe('Previous response is still shutting down.')
17781778
const failed = allQueuedMessages()[0]
17791779
expect(failed).toMatchObject({
1780-
retryRequired: true,
1780+
hold: 'user',
17811781
queuedSendHandoff: {
17821782
stopRequired: true,
17831783
supersededStreamId: state.postBodies[0].userMessageId,
@@ -1966,7 +1966,7 @@ describe('useChat remount send recovery', () => {
19661966
expect(allQueuedMessages()).toEqual([
19671967
expect.objectContaining({
19681968
id: 'queued-correction',
1969-
retryRequired: true,
1969+
hold: 'user',
19701970
queuedSendHandoff: expect.objectContaining({
19711971
userMessageId: 'prepared-correction-request',
19721972
supersededStreamId: 'previous-response',
@@ -1998,7 +1998,7 @@ describe('useChat remount send recovery', () => {
19981998
useMothershipQueueStore.getState().enqueue('chat-a', {
19991999
id: 'earlier-correction',
20002000
content: 'inspect the second invoice instead',
2001-
retryRequired: true,
2001+
hold: 'user',
20022002
queuedSendHandoff: {
20032003
id: 'earlier-correction',
20042004
chatId: 'chat-a',
@@ -2038,7 +2038,7 @@ describe('useChat remount send recovery', () => {
20382038
expect(state.abortBodies[0]?.streamId).toBe(newerStreamId)
20392039
expect(state.postBodies).toHaveLength(1)
20402040
expect(allQueuedMessages()[0]).toMatchObject({
2041-
retryRequired: true,
2041+
hold: 'user',
20422042
queuedSendHandoff: {
20432043
supersededStreamId: newerStreamId,
20442044
userMessageId: 'prepared-correction',
@@ -2358,8 +2358,7 @@ describe('useChat remount send recovery', () => {
23582358
content: 'written while offline',
23592359
resumeUserMessageId: 'offline-attempt',
23602360
admissionUnknown: true,
2361-
retryRequired: true,
2362-
heldUntilOnline: true,
2361+
hold: 'online',
23632362
})
23642363

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

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

23822380
it('keeps a resumed message uneditable when its Send-now Stop does not settle', async () => {
@@ -3192,7 +3190,7 @@ describe('useChat remount send recovery', () => {
31923190

31933191
const queued = useMothershipQueueStore.getState().queues[history.id] ?? []
31943192
expect(queued.map((message) => message.content)).toEqual(['Written while offline'])
3195-
expect(queued[0].retryRequired).toBe(true)
3193+
expect(queued[0].hold).toBe('online')
31963194
expect(queued[0].resumeUserMessageId).toBe(state.postBodies[0].userMessageId)
31973195
expect(getResult().error).not.toBeNull()
31983196
expect(state.postBodies).toHaveLength(1)
@@ -3343,8 +3341,7 @@ describe('useChat remount send recovery', () => {
33433341
expect(state.postBodies).toHaveLength(0)
33443342
expect(useMothershipQueueStore.getState().queues[history.id]?.[0]).toMatchObject({
33453343
content: 'Never prepared',
3346-
retryRequired: true,
3347-
heldUntilOnline: true,
3344+
hold: 'online',
33483345
})
33493346
}
33503347
)
@@ -3435,7 +3432,7 @@ describe('useChat remount send recovery', () => {
34353432

34363433
await waitFor(() => state.postBodies.length === 1)
34373434
await waitFor(
3438-
() => useMothershipQueueStore.getState().queues[history.id]?.[0]?.retryRequired === true
3435+
() => useMothershipQueueStore.getState().queues[history.id]?.[0]?.hold !== undefined
34393436
)
34403437

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

35363533
const second = renderUseChat()
@@ -4307,7 +4304,7 @@ describe('useChat remount send recovery', () => {
43074304
await first.getResult().sendMessage('Held while I was elsewhere')
43084305
})
43094306
await waitFor(
4310-
() => useMothershipQueueStore.getState().queues[history.id]?.[0]?.retryRequired === true
4307+
() => useMothershipQueueStore.getState().queues[history.id]?.[0]?.hold !== undefined
43114308
)
43124309
first.unmount()
43134310

@@ -4354,7 +4351,7 @@ describe('useChat remount send recovery', () => {
43544351
id: 'withdrawn-entry',
43554352
content: 'check the trace for this req',
43564353
resumeUserMessageId: 'accepted-request',
4357-
retryRequired: true,
4354+
hold: 'user',
43584355
})
43594356
const { getResult } = renderUseChatInChat('chat-a', {
43604357
id: 'chat-a',
@@ -4391,7 +4388,7 @@ describe('useChat remount send recovery', () => {
43914388
useMothershipQueueStore.getState().enqueue('chat-a', {
43924389
id: 'unsent-entry',
43934390
content: 'check the trace for this req',
4394-
retryRequired: true,
4391+
hold: 'user',
43954392
resumeUserMessageId: 'unsent-request',
43964393
})
43974394
const { getResult } = renderUseChatInChat('chat-a', {

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

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -4994,7 +4994,7 @@ export function useChat(
49944994
...(retainedHandoff ? { queuedSendHandoff: retainedHandoff } : {}),
49954995
...requeuedFields(
49964996
withdrawn?.reason ?? 'failed',
4997-
dispatched.sendRetries ?? 0,
4997+
dispatched.retry?.attempt ?? 0,
49984998
chatless ? heldSendSurface : undefined
49994999
),
50005000
...(withdrawnUserMessageId ? { resumeUserMessageId: withdrawnUserMessageId } : {}),
@@ -5072,7 +5072,7 @@ export function useChat(
50725072
if (!history) {
50735073
useMothershipQueueStore
50745074
.getState()
5075-
.deferRetry(liveQueueKey(chatKey), msg.id, sendRetry((msg.sendRetries ?? 0) + 1))
5075+
.deferRetry(chatKey, msg.id, sendRetry((msg.retry?.attempt ?? 0) + 1))
50765076
return true
50775077
}
50785078
clearQueuedSendHandoffState(msg.id)
@@ -5101,9 +5101,9 @@ export function useChat(
51015101
const queueState = useMothershipQueueStore.getState()
51025102
const activeChatKey = chatKeyRef.current
51035103
const msg = queueState.queues[activeChatKey]?.[0]
5104-
if (!msg || msg.retryRequired) continue
5104+
if (!msg || msg.hold) continue
51055105
/** An automatic retry waits out its delay; the drain effect wakes it. */
5106-
if (msg.notBefore !== undefined && msg.notBefore > Date.now()) continue
5106+
if (msg.retry && msg.retry.notBefore > Date.now()) continue
51075107
// Pause draining if the head is bound to the composer; dispatching now
51085108
// would race the eventual submit. The next kick on edit-resolve resumes us.
51095109
if (queueState.editing[activeChatKey] === msg.id) continue
@@ -5276,8 +5276,8 @@ export function useChat(
52765276
// `notifyTurnEnded`. Idempotent — the dispatch loop dedupes.
52775277
const chatHistoryReady = chatHistory !== undefined
52785278
const remoteActiveStreamId = chatHistory?.activeStreamId ?? null
5279-
const queueHeadHeld = messageQueue[0]?.retryRequired === true
5280-
const queueHeadNotBefore = messageQueue[0]?.notBefore
5279+
const queueHeadHeld = messageQueue[0]?.hold !== undefined
5280+
const queueHeadNotBefore = messageQueue[0]?.retry?.notBefore
52815281
const [sendRetryWakeup, setSendRetryWakeup] = useState(0)
52825282
useEffect(() => {
52835283
if (!scopeKey) return

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

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,34 @@ describe('useMothershipQueueStore rehydration', () => {
1010
sessionStorage.clear()
1111
})
1212

13+
it('restores the hold and retry fields of a queue saved in their older shape', async () => {
14+
sessionStorage.setItem(
15+
'mothership-queue',
16+
JSON.stringify({
17+
state: {
18+
queues: {
19+
'chat-A': [
20+
{ id: 'for-user', content: 'a', retryRequired: true },
21+
{ id: 'for-network', content: 'b', retryRequired: true, heldUntilOnline: true },
22+
{ id: 'retrying', content: 'c', sendRetries: 2, notBefore: 1_000 },
23+
{ id: 'plain', content: 'd' },
24+
],
25+
},
26+
},
27+
version: 0,
28+
})
29+
)
30+
31+
await useMothershipQueueStore.persist.rehydrate()
32+
33+
expect(useMothershipQueueStore.getState().queues['chat-A']).toEqual([
34+
{ id: 'for-user', content: 'a', hold: 'user' },
35+
{ id: 'for-network', content: 'b', hold: 'online' },
36+
{ id: 'retrying', content: 'c', retry: { attempt: 2, notBefore: 1_000 } },
37+
{ id: 'plain', content: 'd' },
38+
])
39+
})
40+
1341
it('treats a resumed message saved before the edit guard as possibly sent', async () => {
1442
sessionStorage.setItem(
1543
'mothership-queue',

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -126,7 +126,7 @@ describe('useMothershipQueueStore', () => {
126126
useMothershipQueueStore.getState().enqueue('chat-A', {
127127
id: 'm1',
128128
content: 'original',
129-
retryRequired: true,
129+
hold: 'user',
130130
resumeUserMessageId: 'prior-request',
131131
/** Its Stop never settled, so it was never sent: the server cannot hold it. */
132132
admissionUnknown: false,
@@ -152,7 +152,7 @@ describe('useMothershipQueueStore', () => {
152152
})
153153
expect(edited?.queuedSendHandoff?.userMessageId).toBeUndefined()
154154
expect(edited?.resumeUserMessageId).toBeUndefined()
155-
expect(edited?.retryRequired).toBeUndefined()
155+
expect(edited?.hold).toBeUndefined()
156156
})
157157

158158
it('strips queuedSendHandoff on edit so a fresh handoff is minted at send time', () => {

0 commit comments

Comments
 (0)