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
@@ -1,6 +1,5 @@
'use client'

import { useEffect } from 'react'
import {
DropdownMenu,
DropdownMenuContent,
Expand Down Expand Up @@ -53,11 +52,6 @@ export function ModelSelector() {
else setNewChatEffort(choice)
}

useEffect(() => {
if (chatId) return
return () => useMothershipEffortStore.getState().setNewChatEffort(null)
}, [chatId])

const effortLabel = options.find((option) => option.value === effort)?.label ?? effort

return (
Expand Down
163 changes: 163 additions & 0 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 @@ -91,7 +91,9 @@ import { MothershipHandoffStorage } from '@/lib/core/utils/browser-storage'
import { MOTHERSHIP_STREAM_REPLAY_HEADER } from '@/lib/mothership/constants'
import type { MothershipStreamV1EventEnvelope } from '@/lib/mothership/generated/mothership-stream-v1'
import { getChatResourceSelectionId } from '@/lib/mothership/resources/types'
import { ChatSurfaceProvider } from '@/app/workspace/[workspaceId]/home/components/chat-surface-context'
import { collectCitedMessageSources } from '@/app/workspace/[workspaceId]/home/components/message-content/message-sources'
import { ModelSelector } from '@/app/workspace/[workspaceId]/home/components/user-input/components/model-selector'
import {
readQueuedSendHandoffState,
writeQueuedSendHandoffState,
Expand Down Expand Up @@ -472,6 +474,77 @@ function renderHomeLikeSurface(): {
}
}

/** Holds the chat POST of the first send until the test settles it. */
function holdFirstSend(): PromiseWithResolvers<Response> {
const post = Promise.withResolvers<Response>()
vi.stubGlobal('fetch', (input: RequestInfo | URL, init?: RequestInit) => {
if (String(input) !== '/api/mothership/chat' || init?.method !== 'POST') {
return fetchStub(input, init)
}
state.postBodies.push(JSON.parse(String(init.body)))
return post.promise
})
return post
}

/**
* Mounts a chatless surface shaped like `home.tsx`, with real composers: the empty-state one
* swaps for the chat view's once messages show, and the chat view's names the resolved chat.
*/
function renderComposerSwap(): {
container: HTMLElement
getResult: () => ReturnType<typeof useChat>
shownEffort: () => string | null | undefined
visit: (pathname: string) => void
} {
;(globalThis as { IS_REACT_ACT_ENVIRONMENT?: boolean }).IS_REACT_ACT_ENVIRONMENT = true
useMothershipEffortStore.getState().reset()
queryClient = new QueryClient({ defaultOptions: { queries: { retry: false } } })
const container = document.createElement('div')
const root = createRoot(container)
mountedRoots.push(root)
let result: ReturnType<typeof useChat> | undefined

function HomeLike() {
result = useChat('ws-1', undefined)
return result.messages.length > 0 ? (
<section key='chat'>
<ChatSurfaceProvider chatId={result.resolvedChatId}>
<ModelSelector />
</ChatSurfaceProvider>
</section>
) : (
<main key='empty'>
<ModelSelector />
</main>
)
}

const render = () =>
act(() => {
root.render(
<QueryClientProvider client={queryClient}>
<HomeLike />
</QueryClientProvider>
)
})
render()

return {
container,
getResult: () => {
if (result === undefined) throw new Error('Hook result is not ready')
return result
},
shownEffort: () =>
container.querySelector('[aria-label="Reasoning effort"]')?.getAttribute('aria-description'),
visit: (pathname) => {
mockUsePathname.mockReturnValue(pathname)
render()
},
}
}

/**
* Mounts the hook under StrictMode with a handoff already in storage, mirroring
* `home.tsx`'s consume-and-auto-send effect. This is the production-shaped
Expand Down Expand Up @@ -4701,6 +4774,96 @@ describe('useChat remount send recovery', () => {
])
})

it('keeps the new-chat effort across the composer swap of a first send that fails', async () => {
const post = holdFirstSend()
const surface = renderComposerSwap()
act(() => useMothershipEffortStore.getState().setNewChatEffort('low'))
expect(surface.shownEffort()).toBe('Low')

await act(async () => {
void surface.getResult().sendMessage('Plan the launch')
})
await waitFor(() => state.postBodies.length === 1)
expect(state.postBodies[0].effort).toBe('low')
expect(surface.container.querySelector('section')).not.toBeNull()
expect(surface.shownEffort()).toBe('Low')

await act(async () => {
post.reject(new TypeError('Failed to fetch'))
})
await waitFor(() => surface.container.querySelector('main') !== null)

expect(useMothershipEffortStore.getState().newChatEffort).toBe('low')
expect(surface.shownEffort()).toBe('Low')
})

it('keeps a new-chat effort picked while the first send is pending when that send fails', async () => {
const post = holdFirstSend()
const surface = renderComposerSwap()
act(() => useMothershipEffortStore.getState().setNewChatEffort('low'))
await act(async () => {
void surface.getResult().sendMessage('Plan the launch')
})
await waitFor(() => state.postBodies.length === 1)
act(() => useMothershipEffortStore.getState().setNewChatEffort('medium'))

await act(async () => {
post.reject(new TypeError('Failed to fetch'))
})
await waitFor(() => surface.container.querySelector('main') !== null)

expect(surface.shownEffort()).toBe('Medium')
})

it('starts the next new chat at the default after a first send stopped before admission', async () => {
const surface = renderComposerSwap()
act(() => useMothershipEffortStore.getState().setNewChatEffort('low'))
await act(async () => {
void surface.getResult().sendMessage('Plan the launch')
})
await waitFor(() => state.postBodies.length === 1)
await act(async () => {
await surface.getResult().stopGeneration()
})
await waitFor(() => surface.getResult().resolvedChatId === DEDUPED_CHAT_ID)

surface.visit(`/workspace/ws-1/chat/${DEDUPED_CHAT_ID}`)
surface.visit('/workspace/ws-1/home')
await waitFor(() => surface.container.querySelector('main') !== null)

expect(surface.shownEffort()).toBe('High')
})

it.each(['leaves the page', 'opens another chat'] as const)(
'drops an unsent new-chat effort when the surface %s',
(leave) => {
useMothershipEffortStore.getState().reset()
;(globalThis as { IS_REACT_ACT_ENVIRONMENT?: boolean }).IS_REACT_ACT_ENVIRONMENT = true
queryClient = new QueryClient({ defaultOptions: { queries: { retry: false } } })
const root = createRoot(document.createElement('div'))
mountedRoots.push(root)
function Surface({ chatId }: { chatId?: string }) {
useChat('ws-1', chatId)
return null
}
const render = (chatId?: string) =>
act(() =>
root.render(
<QueryClientProvider client={queryClient}>
<Surface chatId={chatId} />
</QueryClientProvider>
)
)
render()
useMothershipEffortStore.getState().setNewChatEffort('low')

if (leave === 'leaves the page') act(() => root.unmount())
else render('chat-other')

expect(useMothershipEffortStore.getState().newChatEffort).toBeNull()
}
)

it('loads the saved transcript once when its own stream completes', async () => {
const chatId = 'chat-own-completion'
const history: MothershipChatHistory = {
Expand Down
24 changes: 13 additions & 11 deletions apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1001,14 +1001,22 @@ export function useChat(
new Set())
const streamReaderRef = useRef<ReadableStreamDefaultReader<Uint8Array> | null>(null)
const chatIdRef = useRef<string | undefined>(initialChatId)
/** Cleared on unmount, so a late rollback cannot hand a pick to a surface the user left. */
/** Cleared on unmount, so late async work cannot act on a surface the user left. */
const surfaceMountedRef = useRef(true)
useEffect(() => {
surfaceMountedRef.current = true
return () => {
surfaceMountedRef.current = false
}
}, [])
/* The new-chat effort pick belongs to this surface, not to one composer: it outlives the swap
from the empty-state composer to the chat view during a first send, so a withdrawn send
leaves it in place. It drops when the surface unmounts or switches chats, and when it adopts
a chat (`adoptResolvedChatId`). */
useEffect(() => {
if (initialChatId) return
return () => useMothershipEffortStore.getState().setNewChatEffort(null)
}, [initialChatId])
const tableViewContextsRef = useRef({
scopeId: desktopScopeId,
views: new Map<string, MothershipTableViewContext>(),
Expand Down Expand Up @@ -1278,6 +1286,10 @@ export function useChat(
const resolvedDesktopScopeId = desktopChatScopeId(scopeKey, chatId)
if (wasPending) {
useChatPanelStore.getState().migrate(pendingDesktopScopeId, resolvedDesktopScopeId)
// Leaving the new chat. An admitted send has already moved the pick onto its chat; any
// other way out (a Stop before admission, a recovered handoff) must not carry it into
// the next new chat.
useMothershipEffortStore.getState().setNewChatEffort(null)
}
const activeActivityTracker = resourceActivityTrackerRef.current
if (activeActivityTracker?.generation === streamGenRef.current) {
Expand Down Expand Up @@ -3736,16 +3748,6 @@ export function useChat(
}

const rollbackOptimisticSend = () => {
// A withdrawn first send hands its pick back to the new-chat composer for the retry,
// only while that surface is still open on the new chat.
if (
!requestChatId &&
effortChoice &&
surfaceMountedRef.current &&
!chatIdRef.current &&
!selectedChatIdRef.current
)
useMothershipEffortStore.getState().setNewChatEffort(effortChoice)
if (requestChatId) {
upsertChatHistory(requestChatId, (current) => ({
...current,
Expand Down
23 changes: 22 additions & 1 deletion apps/sim/hooks/queries/mothership-chats.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ const { suspendBrowserScope, suspendTerminalScope, clearChat } = vi.hoisted(() =
}))

vi.mock('@/stores/mothership-queue/store', () => ({
useMothershipQueueStore: { getState: () => ({ clearChat }) },
useMothershipQueueStore: { getState: () => ({ clearChat, cleared: {} }) },
}))

vi.mock('@tanstack/react-query', () => reactQueryMock)
Expand All @@ -26,6 +26,7 @@ vi.mock('@/lib/terminal/transport', () => ({

import type { MothershipEffort } from '@/lib/mothership/model-options'
import {
fetchMothershipChatHistory,
useDeleteMothershipChats,
useSetMothershipChatEffort,
} from '@/hooks/queries/mothership-chats'
Expand Down Expand Up @@ -100,6 +101,26 @@ describe('tasks query boundary parsing', () => {
})
})

it('loads a chat from a server that predates the effort field', async () => {
vi.mocked(fetch).mockResolvedValueOnce(
jsonResponse({
success: true,
chat: {
id: 'chat-1',
title: null,
mode: 'agent',
messages: [],
activeStreamId: null,
resources: [],
},
})
)

const history = await fetchMothershipChatHistory('chat-1')

expect(history.effort).toBeNull()
})

it('keeps the latest effort pick when an earlier queued save of the same value fails', async () => {
const tanstack =
await vi.importActual<typeof import('@tanstack/react-query')>('@tanstack/react-query')
Expand Down
3 changes: 2 additions & 1 deletion apps/sim/lib/api/contracts/mothership-chats.ts
Original file line number Diff line number Diff line change
Expand Up @@ -420,7 +420,8 @@ export const getMothershipChatResponseSchema = z.object({
messages: z.array(z.unknown()),
activeStreamId: z.string().nullable(),
resources: z.array(z.unknown()),
effort: mothershipChatEffortChoiceSchema,
/** Optional so a client still loads chats from a server that predates the field. */
effort: mothershipChatEffortChoiceSchema.optional(),
createdAt: z.union([z.string(), z.date()]).nullable().optional(),
updatedAt: z.union([z.string(), z.date()]).nullable().optional(),
streamSnapshot: mothershipChatStreamSnapshotSchema.optional(),
Expand Down
16 changes: 16 additions & 0 deletions apps/sim/lib/events/sse-endpoint.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,22 @@ describe('createWorkspaceSSE', () => {
expect(unsubscribe).toHaveBeenCalledTimes(1)
})

it('never subscribes when the request aborted before the stream started', async () => {
const controller = new AbortController()
controller.abort()
const subscribe = vi.fn(() => () => {})
const { body } = await openConnection(controller.signal, [{ subscribe }])
let closed = false
void drain(body).then(() => {
closed = true
})

await vi.advanceTimersByTimeAsync(0)

expect(closed).toBe(true)
expect(subscribe).not.toHaveBeenCalled()
})

it('runs every teardown when one unsubscribe throws', async () => {
const first = vi.fn(() => {
throw new Error('unsubscribe failed')
Expand Down
7 changes: 7 additions & 0 deletions apps/sim/lib/events/sse-endpoint.ts
Original file line number Diff line number Diff line change
Expand Up @@ -213,6 +213,13 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig):
)
}

// An abort listener never fires for a signal that is already aborted, so a client that left
// while the route was authorizing would otherwise hold its subscriptions until rotation.
if (request.signal.aborted) {
close('aborted')
return
}

try {
// The runtime sends the status and headers with the first body chunk, so the stream writes
// one as soon as it opens. A client reads its state once the stream opens, so it opens once
Expand Down
4 changes: 2 additions & 2 deletions apps/sim/stores/mothership-effort/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,8 +15,8 @@ interface MothershipEffortState {
setModel: (model: ModelSelection['model']) => void
setFastMode: (fastMode: boolean) => void
/**
* The effort picked in a composer whose chat does not exist yet. Its first send records
* it on the new chat; leaving that composer unsent drops it.
* The effort picked on a chat surface whose chat does not exist yet. Its first send records
* it on the new chat; leaving that surface unsent drops it.
*/
newChatEffort: MothershipEffort | null
setNewChatEffort: (effort: MothershipEffort | null) => void
Expand Down
Loading