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
723 changes: 723 additions & 0 deletions apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx

Large diffs are not rendered by default.

187 changes: 147 additions & 40 deletions apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts

Large diffs are not rendered by default.

119 changes: 119 additions & 0 deletions apps/sim/hooks/queries/mothership-chat-history-restore.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,119 @@
import {
apiClientRequestMock,
apiClientRequestMockFns,
} from '@sim/testing/mocks/api-client-request.mock'
import { reactQueryMock } from '@sim/testing/mocks/react-query.mock'
import { beforeEach, describe, expect, it, vi } from 'vitest'

vi.mock('@/lib/api/client/request', () => apiClientRequestMock)
vi.mock('@tanstack/react-query', () => reactQueryMock)

import {
fetchMothershipChatHistory,
type MothershipChatHistory,
useRestoreMothershipChat,
} from '@/hooks/queries/mothership-chats'
import { useMothershipQueueStore } from '@/stores/mothership-queue/store'

const mockRequestJson = apiClientRequestMockFns.mockRequestJson

const history: MothershipChatHistory = {
id: 'chat-1',
mode: 'agent',
title: 'Restored',
messages: [],
activeStreamId: null,
resources: [],
}

/** Whether the queue store takes a send for the chat, i.e. whether its delete still holds. */
function takesSends(chatId: string): boolean {
useMothershipQueueStore.getState().enqueue(chatId, { id: 'probe', content: 'probe' })
return useMothershipQueueStore.getState().queues[chatId] !== undefined
}

/** A server answer the test releases when it chooses. */
function deferredAnswer(value: unknown): () => void {
let answer!: () => void
mockRequestJson.mockReturnValue(
new Promise((resolve) => {
answer = () => resolve(value)
})
)
return answer
}

/** The options `useRestoreMothershipChat` hands to `useMutation` (the mock returns them). */
interface RestoreMutation {
mutationFn: (chatId: string) => Promise<void>
onMutate?: (chatId: string) => { deleteSeen?: number }
onSuccess: (data: undefined, chatId: string, context?: { deleteSeen?: number }) => void
}

async function restore(chatId: string, whileInFlight: () => void = () => {}) {
const mutation = useRestoreMothershipChat() as unknown as RestoreMutation
const context = mutation.onMutate?.(chatId)
const answer = deferredAnswer({ success: true })
const done = mutation.mutationFn(chatId)
whileInFlight()
answer()
await done
mutation.onSuccess(undefined, chatId, context)
}

describe('lifting a chat delete', () => {
beforeEach(() => {
useMothershipQueueStore.getState().reset()
mockRequestJson.mockReset()
})

it('reopens a chat this tab saw deleted once the server returns it again', async () => {
useMothershipQueueStore.getState().clearChat(history.id)
mockRequestJson.mockResolvedValue({ chat: history })

await fetchMothershipChatHistory(history.id)

expect(takesSends(history.id)).toBe(true)
})

it('keeps the delete when the read returning the chat began before it', async () => {
const answer = deferredAnswer({ chat: history })
const read = fetchMothershipChatHistory(history.id)
useMothershipQueueStore.getState().clearChat(history.id)
answer()
await read

expect(takesSends(history.id)).toBe(false)
})

it('keeps a newer delete that lands while a read after an earlier one is in flight', async () => {
useMothershipQueueStore.getState().clearChat(history.id)
const answer = deferredAnswer({ chat: history })
const read = fetchMothershipChatHistory(history.id)
useMothershipQueueStore.getState().reopenChat(history.id)
useMothershipQueueStore.getState().clearChat(history.id)
answer()
await read

expect(takesSends(history.id)).toBe(false)
})

it('reopens a chat restored from Recently Deleted', async () => {
useMothershipQueueStore.getState().clearChat(history.id)

await restore(history.id)

expect(takesSends(history.id)).toBe(true)
})

it('keeps a delete that lands while the restore is in flight', async () => {
useMothershipQueueStore.getState().clearChat(history.id)

await restore(history.id, () => {
useMothershipQueueStore.getState().reopenChat(history.id)
useMothershipQueueStore.getState().clearChat(history.id)
})

expect(takesSends(history.id)).toBe(false)
})
})
24 changes: 23 additions & 1 deletion apps/sim/hooks/queries/mothership-chats.ts
Original file line number Diff line number Diff line change
Expand Up @@ -281,7 +281,7 @@ export function useOrganizationMothershipChats(
})
}

export async function fetchMothershipChatHistory(
async function readMothershipChatHistory(
chatId: string,
signal?: AbortSignal
): Promise<MothershipChatHistory> {
Expand Down Expand Up @@ -309,6 +309,22 @@ export async function fetchMothershipChatHistory(
return parseChatHistory(await copilotRes.json())
}

/**
* Reads a chat from the server. A chat this tab saw deleted that the server
* returns again was restored, so it takes queued sends again. Only a read that
* began after the delete counts: one already in flight can return the chat from
* before it.
*/
export async function fetchMothershipChatHistory(
chatId: string,
signal?: AbortSignal
): Promise<MothershipChatHistory> {
const deleteSeen = useMothershipQueueStore.getState().cleared[chatId]
const history = await readMothershipChatHistory(chatId, signal)
if (deleteSeen !== undefined) useMothershipQueueStore.getState().reopenChat(chatId, deleteSeen)
return history
}

export function mothershipChatHistoryQueryOptions(chatId: string | undefined) {
return queryOptions({
queryKey: mothershipChatKeys.detail(chatId),
Expand Down Expand Up @@ -365,6 +381,12 @@ export function useRestoreMothershipChat(owner?: MothershipChatOwner) {
const queryClient = useQueryClient()
return useMutation({
mutationFn: restoreChat,
/** The delete this restore undoes; one that lands while it is in flight stays. */
onMutate: (chatId) => ({ deleteSeen: useMothershipQueueStore.getState().cleared[chatId] }),
onSuccess: (_data, chatId, context) => {
if (context?.deleteSeen === undefined) return
useMothershipQueueStore.getState().reopenChat(chatId, context.deleteSeen)
},
onSettled: () => {
queryClient.invalidateQueries({ queryKey: mothershipChatKeys.ownerLists(owner) })
},
Expand Down
22 changes: 22 additions & 0 deletions apps/sim/hooks/use-mothership-chat-events.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import {
handleMothershipChatStatusEvent,
resyncMothershipChatCaches,
} from '@/hooks/use-mothership-chat-events'
import { useMothershipQueueStore } from '@/stores/mothership-queue/store'

describe('handleMothershipChatStatusEvent', () => {
const queryClient = {
Expand Down Expand Up @@ -136,6 +137,27 @@ describe('handleMothershipChatStatusEvent', () => {
expect(suspendTerminalScope).toHaveBeenCalledWith('chat-1')
})

it('drops the queue of a chat deleted elsewhere and takes sends again once it is restored', () => {
useMothershipQueueStore.getState().reset()
const queued = { id: 'm1', content: 'follow-up' }
const publish = (type: 'deleted' | 'created') =>
handleMothershipChatStatusEvent(
queryClient,
'ws-1',
JSON.stringify({ chatId: 'chat-1', type, timestamp: Date.now() })
)

useMothershipQueueStore.getState().enqueue('chat-1', queued)
publish('deleted')
expect(useMothershipQueueStore.getState().queues['chat-1']).toBeUndefined()
useMothershipQueueStore.getState().enqueue('chat-1', queued)
expect(useMothershipQueueStore.getState().queues['chat-1']).toBeUndefined()

publish('created')
useMothershipQueueStore.getState().enqueue('chat-1', queued)
expect(useMothershipQueueStore.getState().queues['chat-1']?.map((m) => m.id)).toEqual(['m1'])
})

it('keeps started task detail when a stale started stream is older than the active stream', () => {
queryClient.getQueryData.mockReturnValue({
id: 'chat-1',
Expand Down
5 changes: 5 additions & 0 deletions apps/sim/hooks/use-mothership-chat-events.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import {
type MothershipChatOwner,
mothershipChatKeys,
} from '@/hooks/queries/mothership-chats'
import { useMothershipQueueStore } from '@/stores/mothership-queue/store'

const logger = createLogger('MothershipChatEvents')

Expand Down Expand Up @@ -119,8 +120,12 @@ export function handleMothershipChatStatusEvent(
// mutation would leave pages and PTYs running indefinitely.
void suspendDesktopChatScopes(payload.chatId)
queryClient.removeQueries({ queryKey: mothershipChatKeys.detail(payload.chatId) })
/** This tab's queue for the chat goes too, and no later send may bring it back. */
useMothershipQueueStore.getState().clearChat(payload.chatId)
return
}
/** A restore is published as `created`; the chat takes queued sends again. */
if (payload.type === 'created') useMothershipQueueStore.getState().reopenChat(payload.chatId)
if (payload.type === 'renamed') {
/**
* The lists invalidated above carry the title every surface renders; the
Expand Down
27 changes: 27 additions & 0 deletions apps/sim/stores/mothership-queue/store.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -88,5 +88,32 @@ describe('useMothershipQueueStore', () => {
])
expect(useMothershipQueueStore.getState().queues['pending::abc']).toBeUndefined()
})

it('lifts only the delete a restore saw, never a later one', () => {
useMothershipQueueStore.getState().clearChat('chat-X')
const seen = useMothershipQueueStore.getState().cleared['chat-X']
useMothershipQueueStore.getState().reopenChat('chat-X')
useMothershipQueueStore.getState().clearChat('chat-X')

useMothershipQueueStore.getState().reopenChat('chat-X', seen)
useMothershipQueueStore.getState().enqueue('chat-X', message('after-stale-restore'))
expect(useMothershipQueueStore.getState().queues['chat-X']).toBeUndefined()

const latest = useMothershipQueueStore.getState().cleared['chat-X']
useMothershipQueueStore.getState().reopenChat('chat-X', latest)
useMothershipQueueStore.getState().enqueue('chat-X', message('after-restore'))
expect(useMothershipQueueStore.getState().queues['chat-X']?.map((m) => m.id)).toEqual([
'after-restore',
])
})

it('does not move a new chat surface queue into a chat deleted meanwhile', () => {
useMothershipQueueStore.getState().enqueue('pending::abc', message('pending-1'))
useMothershipQueueStore.getState().clearChat('chat-X')
useMothershipQueueStore.getState().migrate('pending::abc', 'chat-X')
const state = useMothershipQueueStore.getState()
expect(state.queues['chat-X']).toBeUndefined()
expect(state.queues['pending::abc']).toBeUndefined()
})
})
})
36 changes: 26 additions & 10 deletions apps/sim/stores/mothership-queue/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,10 +41,13 @@ const sessionStorageAdapter = {
},
}

/** Numbers each delete, so a restore or read can tell the delete it saw from a later one. */
let deleteCount = 0

const initialState = {
queues: {} as Record<string, QueuedMothershipMessage[]>,
editing: {} as Record<string, string>,
cleared: {} as Record<string, true>,
cleared: {} as Record<string, number>,
}

const omitKey = <V>(record: Record<string, V>, key: string): Record<string, V> => {
Expand All @@ -67,13 +70,15 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
...initialState,

enqueue: (chatKey, message) =>
set((state) => ({
cleared: omitKey(state.cleared, chatKey),
queues: setQueueForChat(state.queues, chatKey, [
...(state.queues[chatKey] ?? []),
message,
]),
})),
set((state) => {
if (state.cleared[chatKey]) return state
Comment thread
waleedlatif1 marked this conversation as resolved.
Comment thread
waleedlatif1 marked this conversation as resolved.
return {
queues: setQueueForChat(state.queues, chatKey, [
...(state.queues[chatKey] ?? []),
message,
]),
}
}),

insertAt: (chatKey, index, message) =>
set((state) => {
Expand All @@ -99,6 +104,8 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
retryRequired: _retry,
heldUntilOnline: _held,
heldSurface: _surface,
busyRetries: _busyRetries,
notBefore: _notBefore,
...rest
} = next[index]
next[index] = {
Expand Down Expand Up @@ -144,7 +151,8 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
if (!fromQueue && fromEditing === undefined) return state

const queues = omitKey(state.queues, fromKey)
if (fromQueue && fromQueue.length > 0) {
/** A chat deleted meanwhile takes nothing: its queue is gone with it. */
if (fromQueue && fromQueue.length > 0 && !state.cleared[toKey]) {
// Merge defensively in case a stale bucket survived in
// sessionStorage. FIFO: existing first, then the resolved stream.
const existing = state.queues[toKey] ?? []
Expand Down Expand Up @@ -201,9 +209,17 @@ export const useMothershipQueueStore = create<MothershipQueueState>()(
set((state) => ({
queues: omitKey(state.queues, chatKey),
editing: omitKey(state.editing, chatKey),
cleared: { ...state.cleared, [chatKey]: true },
cleared: { ...state.cleared, [chatKey]: ++deleteCount },
})),

reopenChat: (chatKey, deleteToken) =>
set((state) => {
const current = state.cleared[chatKey]
if (current === undefined) return state
if (deleteToken !== undefined && deleteToken !== current) return state
return { cleared: omitKey(state.cleared, chatKey) }
}),

reset: () => set(initialState),
}),
{
Expand Down
16 changes: 13 additions & 3 deletions apps/sim/stores/mothership-queue/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,10 @@ export type QueuedMothershipMessage = QueuedMessage & {
* mount: the next chatless surface for the same owner and workflow adopts it.
*/
heldSurface?: string
/** Busy refusals so far; paces the next retry. */
busyRetries?: number
/** Epoch ms before which a busy-refused message is not sent again. */
notBefore?: number
/**
* Message id of a prior attempt at this send that an unmount cleanup
* withdrew. Reused when the entry is dispatched so the server deduplicates
Expand All @@ -49,10 +53,11 @@ export interface MothershipQueueState {
queues: Record<string, QueuedMothershipMessage[]>
editing: Record<string, string>
/**
* Chats cleared this session (deleted). A late restore of a send dispatched
* before the clear does not recreate their queue; a new enqueue lifts it.
* Chats cleared this session (deleted), each with the token of its latest
* delete. No write recreates their queue (a late restore, or a failed send
* handed back); restoring the chat lifts it.
*/
cleared: Record<string, true>
cleared: Record<string, number>

enqueue: (chatKey: string, message: QueuedMothershipMessage) => void
insertAt: (chatKey: string, index: number, message: QueuedMothershipMessage) => void
Expand All @@ -65,5 +70,10 @@ export interface MothershipQueueState {
/** Moves the sends a dead chatless mount of `surface` held onto `toKey`. */
adoptHeldSends: (toKey: string, surface: string) => void
clearChat: (chatKey: string) => void
/**
* Lifts `cleared` for a restored chat. Given the delete token an operation
* saw when it began, lifts only that delete, never one that landed after it.
*/
reopenChat: (chatKey: string, deleteToken?: number) => void
reset: () => void
}
Loading