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
198 changes: 193 additions & 5 deletions apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -2151,11 +2151,6 @@ describe('useChat remount send recovery', () => {
expect(allQueuedMessages()).toHaveLength(0)
})

/**
* After a Stop of the first message, only the follow-up was the user's
* intent: the Stop's POST is left to the server, nothing withdraws it, and the
* next mount sends just the follow-up, once.
*/
/**
* A first message held at the queue head after a remount may already be a
* turn on the server. Editing it would send different text under a new id,
Expand Down Expand Up @@ -2204,6 +2199,199 @@ describe('useChat remount send recovery', () => {
})
})

/** A chat with a turn running, so anything sent to it waits in its queue. */
function renderBusyChat(id: string) {
const history: MothershipChatHistory = {
id,
mode: 'agent',
title: 'Busy',
messages: [],
activeStreamId: 'turn-still-running',
resources: [],
}
mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history }))
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
if (String(input).includes('/api/mothership/chat/stream')) {
if (String(input).includes('batch=true')) {
return Response.json({ success: true, events: [], status: 'streaming' })
}
return new Response(new ReadableStream<Uint8Array>(), {
headers: { 'Content-Type': 'text/event-stream' },
})
}
return fetchStub(input, init)
})
return { history, ...renderUseChatInChat(id, history) }
}

/**
* A send another surface withdrew arrives here under its original id, through
* the send event or the stored handoff. Queued behind a running turn, it may
* already be a turn on the server, so it can't be edited either.
*/
it('does not let a withdrawn send handed to a busy chat be edited', async () => {
const { history, getResult } = renderBusyChat('chat-busy-on-handoff')
await waitFor(() => getResult().isSending)
await act(async () => {
await getResult().sendMessage('handed over from another surface', undefined, undefined, {
resumeUserMessageId: 'withdrawn-attempt',
})
})
const queued = useMothershipQueueStore.getState().queues[history.id]?.[0]
expect(queued?.resumeUserMessageId).toBe('withdrawn-attempt')

let edited: ReturnType<ReturnType<typeof useChat>['editQueuedMessage']>
await act(async () => {
edited = getResult().editQueuedMessage(queued?.id ?? '')
})

expect(edited).toBeUndefined()
expect(getResult().editingQueuedId).toBeNull()
})

/** A follow-up whose dispatch got no answer may have reached the server too. */
it('does not let a queued follow-up be edited after its send got no answer', async () => {
const history: MothershipChatHistory = {
id: 'chat-follow-up-unanswered',
mode: 'agent',
title: 'Unanswered',
messages: [],
activeStreamId: null,
resources: [],
}
mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history }))
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
if (String(input) === '/api/mothership/chat' && init?.method === 'POST') {
state.postBodies.push(JSON.parse(String(init.body)))
throw new TypeError('Failed to fetch')
}
return fetchStub(input, init)
})
useMothershipQueueStore
.getState()
.enqueue(history.id, { id: 'follow-up', content: 'and the second invoice' })
const { getResult } = renderUseChatInChat(history.id, history)
await waitFor(
() =>
useMothershipQueueStore.getState().queues[history.id]?.[0]?.resumeUserMessageId !==
undefined
)

let edited: ReturnType<ReturnType<typeof useChat>['editQueuedMessage']>
await act(async () => {
edited = getResult().editQueuedMessage('follow-up')
})

expect(edited).toBeUndefined()
expect(useMothershipQueueStore.getState().queues[history.id]?.[0]).toMatchObject({
content: 'and the second invoice',
resumeUserMessageId: state.postBodies[0].userMessageId,
})
})

/**
* Send-now on a resumed message whose Stop of the running turn does not
* settle sends nothing. That says nothing about the earlier attempt the
* message resumes, so it must stay uneditable.
*/
it('keeps a resumed message uneditable when its Send-now Stop does not settle', async () => {
state.abortSettlements = [false, false, false, false]
const { getResult } = renderUseChatInChat('chat-a')
await act(async () => {
void getResult().sendMessage('Original request')
})
await waitFor(() => state.postBodies.length === 1 && getResult().isSending)
await act(async () => {
await getResult().sendMessage('handed over from another surface', undefined, undefined, {
resumeUserMessageId: 'withdrawn-attempt',
})
})
await waitFor(() => useMothershipQueueStore.getState().queues['chat-a']?.length === 1)

await act(async () => {
await getResult()
.sendNow()
.catch(() => {})
await sleep(200)
})
const queued = useMothershipQueueStore.getState().queues['chat-a']?.[0]
let edited: ReturnType<ReturnType<typeof useChat>['editQueuedMessage']>
await act(async () => {
edited = getResult().editQueuedMessage(queued?.id ?? '')
})

expect(state.postBodies).toHaveLength(1)
expect(queued).toMatchObject({
content: 'handed over from another surface',
resumeUserMessageId: 'withdrawn-attempt',
admissionUnknown: true,
})
expect(edited).toBeUndefined()
})

/**
* A held message the server then refuses as busy is known not to be a turn
* there: the server answers a retry of an admitted id as a duplicate, never
* as busy. The user can edit it again.
*/
it('lets a held message be edited again once the server refuses it as busy', async () => {
const history: MothershipChatHistory = {
id: 'chat-held-then-refused',
mode: 'agent',
title: 'Refused',
messages: [],
activeStreamId: null,
resources: [],
}
mockRequestJson.mockImplementation(() => Promise.resolve({ chat: history }))
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
const url = String(input)
if (url === '/api/mothership/chat' && init?.method === 'POST') {
state.postBodies.push(JSON.parse(String(init.body)))
return Response.json(
{
error: 'A response is already in progress for this chat.',
activeStreamId: 'turn-from-another-tab',
},
{ status: 409 }
)
}
if (url.includes('/api/mothership/chat/stream')) {
if (url.includes('batch=true')) {
return Response.json({ success: true, events: [], status: 'streaming' })
}
return new Response(new ReadableStream<Uint8Array>(), {
headers: { 'Content-Type': 'text/event-stream' },
})
}
return fetchStub(input, init)
})
useMothershipQueueStore.getState().enqueue(history.id, {
id: 'held-first',
content: 'inspect the workspace',
resumeUserMessageId: 'first-attempt',
admissionUnknown: true,
})
const { getResult } = renderUseChatInChat(history.id, history)
await waitFor(() => state.postBodies.length === 1)
await waitFor(
() => useMothershipQueueStore.getState().queues[history.id]?.[0]?.id === 'held-first'
)

let edited: ReturnType<ReturnType<typeof useChat>['editQueuedMessage']>
await act(async () => {
edited = getResult().editQueuedMessage('held-first')
})

expect(edited?.content).toBe('inspect the workspace')
expect(getResult().editingQueuedId).toBe('held-first')
})

/**
* After a Stop of the first message, only the follow-up was the user's
* intent: the Stop's POST is left to the server, nothing withdraws it, and the
* next mount sends just the follow-up, once.
*/
it('sends only the follow-up after a failed Stop when the new-chat surface remounts', async () => {
stubFirstPostPendingThenAdmitted()
const first = renderHomeLikeSurface()
Expand Down
34 changes: 31 additions & 3 deletions apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -234,6 +234,18 @@ interface WithdrawnSendResult {
busy?: boolean
/** Not sent at all (its Stop handoff failed); kept queued for the user to send. */
held?: boolean
/**
* The server refused this id outright (busy, or a predecessor still shutting
* down). It answers a retry of an admitted id as a duplicate instead, so the
* server is known not to have it, and its queue entry can be edited.
*/
notAdmitted?: boolean
/**
* This attempt never reached the server (its Stop did not settle). That says
* nothing about an earlier attempt the message resumes, whose uncertainty it
* keeps.
*/
neverSent?: boolean
}

/**
Expand Down Expand Up @@ -3887,7 +3899,7 @@ export function useChat(
setError(getErrorMessage(err, 'Failed to stop the previous response'))
/* Nothing was sent. Hand the message back so it stays in its chat's queue
even if the user has switched chats since the Stop began. */
return { userMessageId, held: true }
return { userMessageId, held: true, neverSent: true }
}
}

Expand Down Expand Up @@ -4004,7 +4016,7 @@ export function useChat(
}
if (viewOnSend)
setError('Previous response is still shutting down; queued message was restored.')
return { userMessageId, held: true }
return { userMessageId, held: true, notAdmitted: true }
}
/** Withdraws this refused send so the queue retries it, under the same id, later. */
const releaseRefusedSend = () => {
Expand Down Expand Up @@ -4040,7 +4052,7 @@ export function useChat(
exact: true,
refetchType: 'none',
})
return { userMessageId, busy: true }
return { userMessageId, busy: true, notAdmitted: true }
}
/* "Already sent" with no stream for it means the earlier attempt is still
in flight on the server (or died before starting a turn), not that a turn
Expand Down Expand Up @@ -4414,6 +4426,11 @@ export function useChat(
: {}),
...(result.held ? { retryRequired: true } : {}),
...(result.busy ? busyRetry(1) : {}),
admissionUnknown: result.notAdmitted
? false
: result.neverSent
? options?.resumeUserMessageId !== undefined
: true,
...((result.unreachable || result.busy) && activeChatKey.startsWith(PENDING_CHAT_KEY_PREFIX)
? { heldSurface: heldSendSurface }
: {}),
Expand Down Expand Up @@ -5043,6 +5060,17 @@ export function useChat(
? { heldSurface: heldSendSurface }
: {}),
...(withdrawnUserMessageId ? { resumeUserMessageId: withdrawnUserMessageId } : {}),
/* A refusal of this id settles it; an attempt that never left keeps the
earlier uncertainty; any other withdrawal may have reached the server. */
...(withdrawn
? {
admissionUnknown: withdrawn.notAdmitted
? false
: withdrawn.neverSent
? dispatched.admissionUnknown === true
: true,
}
: {}),
})
}

Expand Down
7 changes: 4 additions & 3 deletions apps/sim/app/workspace/[workspaceId]/home/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,9 +35,10 @@ export interface QueuedMessage {
assistantSearch?: WorkspaceSearchFilters
assistantSearchLevel?: AssistantSearchLevel
/**
* A first message withdrawn before the server answered. The server may
* already hold it as sent, so it goes out exactly as written, under its
* original id, and cannot be edited into a different message.
* An earlier attempt at this message got no answer, so the server may
* already hold it as sent. It goes out exactly as written, under that
* attempt's id, and cannot be edited into a different message. False once
* the server is known not to have it (it refused it, or it was never sent).
*/
admissionUnknown?: boolean
}
Expand Down
42 changes: 42 additions & 0 deletions apps/sim/stores/mothership-queue/store.dom.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
/**
* @vitest-environment jsdom
*/
import { beforeEach, describe, expect, it } from 'vitest'
import { useMothershipQueueStore } from '@/stores/mothership-queue/store'

describe('useMothershipQueueStore rehydration', () => {
beforeEach(() => {
useMothershipQueueStore.getState().reset()
sessionStorage.clear()
})

it('treats a resumed message saved before the edit guard as possibly sent', async () => {
sessionStorage.setItem(
'mothership-queue',
JSON.stringify({
state: {
queues: {
'chat-A': [
{ id: 'saved-before', content: 'original', resumeUserMessageId: 'attempt-1' },
{
id: 'refused',
content: 'original',
resumeUserMessageId: 'attempt-2',
admissionUnknown: false,
},
{ id: 'plain', content: 'never sent' },
],
},
},
version: 0,
})
)

await useMothershipQueueStore.persist.rehydrate()

const [savedBefore, refused, plain] = useMothershipQueueStore.getState().queues['chat-A'] ?? []
expect(savedBefore?.admissionUnknown).toBe(true)
expect(refused?.admissionUnknown).toBe(false)
expect(plain?.admissionUnknown).toBeUndefined()
})
})
20 changes: 20 additions & 0 deletions apps/sim/stores/mothership-queue/store.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,24 @@ describe('useMothershipQueueStore', () => {
})

describe('replaceAt', () => {
it('treats any message resuming an earlier attempt as possibly sent, unless told otherwise', () => {
useMothershipQueueStore
.getState()
.enqueue('chat-A', { id: 'resumed', content: 'original', resumeUserMessageId: 'attempt-1' })
useMothershipQueueStore.getState().insertAt('chat-A', 0, {
id: 'refused',
content: 'original',
resumeUserMessageId: 'attempt-2',
admissionUnknown: false,
})
useMothershipQueueStore.getState().replaceAt('chat-A', 'resumed', { content: 'edited' })
useMothershipQueueStore.getState().replaceAt('chat-A', 'refused', { content: 'edited' })

const [refused, resumed] = useMothershipQueueStore.getState().queues['chat-A'] ?? []
expect(resumed).toMatchObject({ content: 'original', resumeUserMessageId: 'attempt-1' })
expect(refused?.content).toBe('edited')
})

it('leaves a first message the server may already hold unchanged, with its id', () => {
useMothershipQueueStore.getState().enqueue('chat-A', {
id: 'm1',
Expand All @@ -50,6 +68,8 @@ describe('useMothershipQueueStore', () => {
content: 'original',
retryRequired: true,
resumeUserMessageId: 'prior-request',
/** Its Stop never settled, so it was never sent: the server cannot hold it. */
admissionUnknown: false,
queuedSendHandoff: {
id: 'm1',
chatId: 'chat-A',
Expand Down
Loading
Loading