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
30 changes: 0 additions & 30 deletions apps/sim/app/api/v2/chat/route.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -267,7 +267,6 @@ describe('POST /api/v2/chat', () => {
mockResolveOrCreateChat.mockResolvedValue({
chatId: SERVER_ISSUED_CHAT_ID,
chat: chatRow(SERVER_ISSUED_CHAT_ID),
conversationHistory: [],
isNew: true,
})
})
Expand Down Expand Up @@ -444,7 +443,6 @@ describe('POST /api/v2/chat', () => {
mockResolveOrCreateChat.mockResolvedValue({
chatId: OWNED_CONVERSATION_ID,
chat: chatRow(OWNED_CONVERSATION_ID),
conversationHistory: [],
isNew: false,
})

Expand All @@ -469,38 +467,10 @@ describe('POST /api/v2/chat', () => {
})
})

it('posts only the current turn on a resumed conversation, never the stored transcript', async () => {
mockResolveOrCreateChat.mockResolvedValue({
chatId: OWNED_CONVERSATION_ID,
chat: chatRow(OWNED_CONVERSATION_ID),
conversationHistory: [
{ role: 'user', content: 'first' },
{ role: 'assistant', content: 'first reply' },
],
isNew: false,
})

const response = await callChat({
workspaceId: 'workspace-1',
message: 'and then?',
conversationId: OWNED_CONVERSATION_ID,
})

expect(response.status).toBe(200)
// Continuity is keyed by chatId downstream, exactly as the web send path
// and the Sim Chat block do. Replaying the transcript here would duplicate
// every prior turn.
expect(mockRunHeadlessCopilotLifecycle.mock.calls[0][0]).toMatchObject({
message: 'and then?',
chatId: OWNED_CONVERSATION_ID,
})
})

it('answers 404 and runs nothing when the resolver refuses the named conversation', async () => {
mockResolveOrCreateChat.mockResolvedValue({
chatId: OWNED_CONVERSATION_ID,
chat: null,
conversationHistory: [],
isNew: false,
})

Expand Down
7 changes: 3 additions & 4 deletions apps/sim/app/api/v2/chat/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -264,12 +264,11 @@ export const POST = withRouteHandler(
// surface uses, and refuse every id that does not resolve with the same
// response so the refusal carries no information about the id. Omitting
// the id mints a server-issued conversation instead of trusting one.
// The resolved transcript is deliberately not forwarded: continuity is
// keyed by `chatId` downstream, exactly as the web send path and the Sim
// Chat block do, both of which post a single message with a chat id.
// Continuity is keyed by `chatId` downstream, exactly as the web send
// path and the Sim Chat block do, both of which post a single message
// with a chat id.
const resolvedChat = await resolveOrCreateChat({
...(conversationId ? { chatId: conversationId } : {}),
includeTranscript: false,
userId,
workspaceId,
model: MOTHERSHIP_CHAT_DEFAULT_MODEL,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,7 @@ import {
} from '@/app/workspace/[workspaceId]/home/hooks/send-handoff'
import { useChat } from '@/app/workspace/[workspaceId]/home/hooks/use-chat'
import { type MothershipChatHistory, mothershipChatKeys } from '@/hooks/queries/mothership-chats'
import { handleMothershipChatStatusEvent } from '@/hooks/use-mothership-chat-events'
import { useExecutionStore } from '@/stores/execution/store'
import { useMothershipQueueStore } from '@/stores/mothership-queue/store'

Expand Down Expand Up @@ -1821,4 +1822,81 @@ describe('useChat remount send recovery', () => {
expect(getResult().messageQueue.map((entry) => entry.id)).toEqual(['unsent-entry'])
expect(state.postBodies).toHaveLength(0)
})

it('loads the saved transcript once when its own stream completes', async () => {
const chatId = 'chat-own-completion'
const history: MothershipChatHistory = {
id: chatId,
mode: 'agent',
title: 'Own stream',
messages: [],
activeStreamId: null,
resources: [],
}
const saved = [
{ id: 'saved-user', role: 'user', content: 'Summarize the run' },
{ id: 'saved-assistant', role: 'assistant', content: 'Done.' },
]
const detailRequests: string[] = []
mockRequestJson.mockImplementation((contract: AnyApiRouteContract) => {
if (contract.path !== '/api/mothership/chats/[chatId]') {
return Promise.resolve({ chats: [] })
}
detailRequests.push(chatId)
return Promise.resolve({ chat: { ...history, messages: saved } })
})
let stream: ReadableStreamDefaultController<Uint8Array> | undefined
let streamId: string | undefined
const emit = (event: Omit<MothershipStreamV1EventEnvelope, 'v' | 'ts' | 'stream'>) =>
stream?.enqueue(
new TextEncoder().encode(
`data: ${JSON.stringify({ v: 1, ts: '', stream: { streamId }, ...event })}\n\n`
)
)
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
if (String(input) !== '/api/mothership/chat' || init?.method !== 'POST') {
return fetchStub(input, init)
}
streamId = JSON.parse(String(init.body)).userMessageId
return new Response(
new ReadableStream<Uint8Array>({
start(controller) {
stream = controller
},
}),
{ headers: { 'Content-Type': 'text/event-stream', 'x-mothership-chat-id': chatId } }
)
})
const { getResult } = renderUseChatInChat(chatId, history)

await act(async () => {
void getResult().sendMessage('Summarize the run')
})
await waitFor(() => stream !== undefined)
emit({ seq: 1, type: 'text', payload: { channel: 'assistant', text: 'Done.' } })
await waitFor(
() =>
queryClient.getQueryData<MothershipChatHistory>(mothershipChatKeys.detail(chatId))
?.activeStreamId === streamId
)
/** The server publishes `completed` after persisting and before closing the stream. */
handleMothershipChatStatusEvent(queryClient, 'ws-1', {
chatId,
type: 'completed',
streamId,
})
emit({ seq: 2, type: 'complete', payload: { status: 'complete' } })
stream?.close()

await waitFor(() => !getResult().isSending && detailRequests.length > 0)
await act(async () => {
await sleep(50)
})
expect(detailRequests).toHaveLength(1)
expect(
queryClient
.getQueryData<MothershipChatHistory>(mothershipChatKeys.detail(chatId))
?.messages.map((message) => message.id)
).toEqual(['saved-user', 'saved-assistant'])
})
})
79 changes: 77 additions & 2 deletions apps/sim/hooks/use-mothership-chat-events.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import type { QueryClient } from '@tanstack/react-query'
import { sleep } from '@sim/utils/helpers'
import { QueryClient, QueryObserver } from '@tanstack/react-query'
import { beforeEach, describe, expect, it, vi } from 'vitest'

const { suspendBrowserScope, suspendTerminalScope } = vi.hoisted(() => ({
Expand All @@ -9,7 +10,7 @@ const { suspendBrowserScope, suspendTerminalScope } = vi.hoisted(() => ({
vi.mock('@/lib/browser-agent/transport', () => ({ suspendBrowserScope }))
vi.mock('@/lib/terminal/transport', () => ({ suspendTerminalScope }))

import { mothershipChatKeys } from '@/hooks/queries/mothership-chats'
import { type MothershipChatHistory, mothershipChatKeys } from '@/hooks/queries/mothership-chats'
import {
handleMothershipChatStatusEvent,
resyncMothershipChatCaches,
Expand Down Expand Up @@ -184,6 +185,80 @@ describe('handleMothershipChatStatusEvent', () => {
)
})

describe('chat detail refetches driven by status events', () => {
function mountDetail(cached: MothershipChatHistory) {
const queryClient = new QueryClient()
const fetchTranscript = vi.fn(async () => cached)
queryClient.setQueryData(mothershipChatKeys.detail('chat-1'), cached)
const unsubscribe = new QueryObserver(queryClient, {
queryKey: mothershipChatKeys.detail('chat-1'),
queryFn: fetchTranscript,
staleTime: Number.POSITIVE_INFINITY,
}).subscribe(() => {})
return { queryClient, fetchTranscript, unsubscribe }
}

const liveStream: MothershipChatHistory = {
id: 'chat-1',
Comment thread
waleedlatif1 marked this conversation as resolved.
title: null,
messages: [
{ id: 'stream-1' },
{ id: 'live-assistant:stream-1' },
] as MothershipChatHistory['messages'],
activeStreamId: 'stream-1',
resources: [],
}

it('does not reload the transcript when the viewer finishes its own live stream', async () => {
const { queryClient, fetchTranscript, unsubscribe } = mountDetail(liveStream)

handleMothershipChatStatusEvent(queryClient, 'ws-1', {
chatId: 'chat-1',
type: 'completed',
streamId: 'stream-1',
})
await sleep(0)

expect(fetchTranscript).not.toHaveBeenCalled()
unsubscribe()
})

it('reloads the saved transcript when a cached mid-stream detail is opened after completion', async () => {
const queryClient = new QueryClient()
const fetchTranscript = vi.fn(async () => ({ ...liveStream, activeStreamId: null }))
queryClient.setQueryData(mothershipChatKeys.detail('chat-1'), liveStream)

handleMothershipChatStatusEvent(queryClient, 'ws-1', {
chatId: 'chat-1',
type: 'completed',
streamId: 'stream-1',
})
const unsubscribe = new QueryObserver(queryClient, {
queryKey: mothershipChatKeys.detail('chat-1'),
queryFn: fetchTranscript,
staleTime: Number.POSITIVE_INFINITY,
}).subscribe(() => {})

await vi.waitFor(() => expect(fetchTranscript).toHaveBeenCalledTimes(1))
unsubscribe()
})

it('marks the detail stale on rename without reloading the transcript', async () => {
const { queryClient, fetchTranscript, unsubscribe } = mountDetail({
...liveStream,
messages: [],
activeStreamId: null,
})

handleMothershipChatStatusEvent(queryClient, 'ws-1', { chatId: 'chat-1', type: 'renamed' })
await sleep(0)

expect(fetchTranscript).not.toHaveBeenCalled()
expect(queryClient.getQueryState(mothershipChatKeys.detail('chat-1'))?.isInvalidated).toBe(true)
unsubscribe()
})
})

describe('resyncMothershipChatCaches', () => {
const queryClient = {
invalidateQueries: vi.fn().mockResolvedValue(undefined),
Expand Down
47 changes: 30 additions & 17 deletions apps/sim/hooks/use-mothership-chat-events.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,12 +33,6 @@ interface ChatStatusEventPayload {
streamId?: string
}

const DETAIL_INVALIDATING_CHAT_STATUS_TYPES = new Set<ChatStatusEventType>([
'started',
'completed',
'renamed',
])

function isChatStatusEventType(value: unknown): value is ChatStatusEventType {
return typeof value === 'string' && CHAT_STATUS_TYPE_SET.has(value)
}
Expand Down Expand Up @@ -69,7 +63,6 @@ function shouldSkipDetailInvalidationForStreamEvent(
current: MothershipChatHistory | undefined,
payload: ChatStatusEventPayload
) {
if (payload.type !== 'started' && payload.type !== 'completed') return false
if (!current?.activeStreamId) return false
if (!payload.streamId) return isLocalOptimisticActiveStream(current)
if (payload.type === 'started' && current.activeStreamId === payload.streamId) return true
Expand Down Expand Up @@ -128,17 +121,37 @@ export function handleMothershipChatStatusEvent(
queryClient.removeQueries({ queryKey: mothershipChatKeys.detail(payload.chatId) })
return
}
if (payload.type === 'started' || payload.type === 'completed') {
const current = queryClient.getQueryData<MothershipChatHistory>(
mothershipChatKeys.detail(payload.chatId)
)
if (shouldSkipDetailInvalidationForStreamEvent(current, payload)) {
return
}
}
if (payload.type && DETAIL_INVALIDATING_CHAT_STATUS_TYPES.has(payload.type)) {
queryClient.invalidateQueries({ queryKey: mothershipChatKeys.detail(payload.chatId) })
if (payload.type === 'renamed') {
/**
* The lists invalidated above carry the title every surface renders; the
* detail only needs marking stale, not a full transcript reload.
*/
queryClient.invalidateQueries({
queryKey: mothershipChatKeys.detail(payload.chatId),
refetchType: 'none',
})
return
}
if (payload.type !== 'started' && payload.type !== 'completed') return
const current = queryClient.getQueryData<MothershipChatHistory>(
mothershipChatKeys.detail(payload.chatId)
)
if (shouldSkipDetailInvalidationForStreamEvent(current, payload)) return
/**
* A completion of the cached live stream only marks the detail stale. The
* server persists the turn before closing the stream, so a surface rendering
* it refetches through its own finalization; the live message alone cannot
* tell this tab's stream from a server-loaded mid-stream snapshot, so any
* other cached copy reloads the saved transcript on its next mount.
*/
const completesCachedLiveStream =
Comment thread
waleedlatif1 marked this conversation as resolved.
payload.type === 'completed' &&
current?.activeStreamId === payload.streamId &&
isLocalOptimisticActiveStream(current)
queryClient.invalidateQueries({
queryKey: mothershipChatKeys.detail(payload.chatId),
...(completesCachedLiveStream ? { refetchType: 'none' as const } : {}),
})
}

/**
Expand Down
Loading
Loading