From 3c266bf5b278ee85f3d958eabceb7a62c3c8f1ac Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 11:10:59 -0700 Subject: [PATCH 1/2] fix(mothership): close a pre-aborted SSE stream, keep the new-chat effort across a failed first send - createSSEStream closes at once when the request aborted before start(), instead of subscribing until rotation. - The new-chat effort pick is dropped by useChat when its chatless surface is left, not by each composer's unmount, so a failed first send keeps it and the pending chat view shows it. - The chat response's effort is optional, so a new client loads chats from a server that predates it. --- .../user-input/components/model-selector.tsx | 6 -- .../home/hooks/use-chat.dom.test.tsx | 89 +++++++++++++++++++ .../[workspaceId]/home/hooks/use-chat.ts | 7 ++ .../hooks/queries/mothership-chats.test.ts | 23 ++++- .../sim/lib/api/contracts/mothership-chats.ts | 3 +- apps/sim/lib/events/sse-endpoint.test.ts | 16 ++++ apps/sim/lib/events/sse-endpoint.ts | 7 ++ apps/sim/stores/mothership-effort/store.ts | 4 +- 8 files changed, 145 insertions(+), 10 deletions(-) diff --git a/apps/sim/app/workspace/[workspaceId]/home/components/user-input/components/model-selector.tsx b/apps/sim/app/workspace/[workspaceId]/home/components/user-input/components/model-selector.tsx index 83a7f75c1f1..6be1012843e 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/components/user-input/components/model-selector.tsx +++ b/apps/sim/app/workspace/[workspaceId]/home/components/user-input/components/model-selector.tsx @@ -1,6 +1,5 @@ 'use client' -import { useEffect } from 'react' import { DropdownMenu, DropdownMenuContent, @@ -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 ( diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx index 2035ec6d383..e999ddb8799 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx @@ -92,6 +92,7 @@ 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 { 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, @@ -4701,6 +4702,94 @@ describe('useChat remount send recovery', () => { ]) }) + it('keeps the new-chat effort across the composer swap of a first send that fails', async () => { + useMothershipEffortStore.getState().reset() + const post = Promise.withResolvers() + 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 + }) + ;(globalThis as { IS_REACT_ACT_ENVIRONMENT?: boolean }).IS_REACT_ACT_ENVIRONMENT = true + queryClient = new QueryClient({ defaultOptions: { queries: { retry: false } } }) + const container = document.createElement('div') + const root = createRoot(container) + mountedRoots.push(root) + let result: ReturnType | undefined + /** Like home.tsx: the empty-state composer swaps for the chat view's once messages show. */ + function HomeLike() { + result = useChat('ws-1', undefined) + return result.messages.length > 0 ? ( +
+ +
+ ) : ( +
+ +
+ ) + } + act(() => { + root.render( + + + + ) + }) + const shownEffort = () => + container.querySelector('[aria-label="Reasoning effort"]')?.getAttribute('aria-description') + act(() => useMothershipEffortStore.getState().setNewChatEffort('low')) + expect(shownEffort()).toBe('Low') + + await act(async () => { + void result?.sendMessage('Plan the launch') + }) + await waitFor(() => state.postBodies.length === 1) + expect(state.postBodies[0].effort).toBe('low') + expect(container.querySelector('section')).not.toBeNull() + expect(shownEffort()).toBe('Low') + + await act(async () => { + post.reject(new TypeError('Failed to fetch')) + }) + await waitFor(() => container.querySelector('main') !== null) + + expect(useMothershipEffortStore.getState().newChatEffort).toBe('low') + expect(shownEffort()).toBe('Low') + }) + + 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( + + + + ) + ) + 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 = { diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts index 2e12636dccf..53bdc041bdb 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -1009,6 +1009,13 @@ export function useChat( 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, and drops only when the + surface leaves the new chat. */ + useEffect(() => { + if (initialChatId) return + return () => useMothershipEffortStore.getState().setNewChatEffort(null) + }, [initialChatId]) const tableViewContextsRef = useRef({ scopeId: desktopScopeId, views: new Map(), diff --git a/apps/sim/hooks/queries/mothership-chats.test.ts b/apps/sim/hooks/queries/mothership-chats.test.ts index a366ade5303..ee7461b2cbb 100644 --- a/apps/sim/hooks/queries/mothership-chats.test.ts +++ b/apps/sim/hooks/queries/mothership-chats.test.ts @@ -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) @@ -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' @@ -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('@tanstack/react-query') diff --git a/apps/sim/lib/api/contracts/mothership-chats.ts b/apps/sim/lib/api/contracts/mothership-chats.ts index 11b7670a2cd..5b4574e5acd 100644 --- a/apps/sim/lib/api/contracts/mothership-chats.ts +++ b/apps/sim/lib/api/contracts/mothership-chats.ts @@ -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(), diff --git a/apps/sim/lib/events/sse-endpoint.test.ts b/apps/sim/lib/events/sse-endpoint.test.ts index e812b845e1f..bcb0e3c70be 100644 --- a/apps/sim/lib/events/sse-endpoint.test.ts +++ b/apps/sim/lib/events/sse-endpoint.test.ts @@ -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') diff --git a/apps/sim/lib/events/sse-endpoint.ts b/apps/sim/lib/events/sse-endpoint.ts index a509a884f9e..0b5afb6a260 100644 --- a/apps/sim/lib/events/sse-endpoint.ts +++ b/apps/sim/lib/events/sse-endpoint.ts @@ -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 diff --git a/apps/sim/stores/mothership-effort/store.ts b/apps/sim/stores/mothership-effort/store.ts index bf60d53618b..20c3f1cd664 100644 --- a/apps/sim/stores/mothership-effort/store.ts +++ b/apps/sim/stores/mothership-effort/store.ts @@ -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 From 9eb580bd6b77fc5fb22ca246a8f3b3f3ff42f6b5 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Wed, 7 Oct 2026 11:41:04 -0700 Subject: [PATCH 2/2] fix(mothership): drop the new-chat effort when the surface adopts a chat A first send stopped before admission adopted its chat without moving the pick, so the next new chat on the same Home mount showed and sent it. adoptResolvedChatId now drops the pick when the surface leaves the new chat. The rollback restore is gone: nothing clears the pick while a send is pending, and it overwrote a pick made during the send. --- .../home/hooks/use-chat.dom.test.tsx | 160 +++++++++++++----- .../[workspaceId]/home/hooks/use-chat.ts | 21 +-- 2 files changed, 125 insertions(+), 56 deletions(-) diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx index e999ddb8799..e61b0c6c965 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx @@ -91,6 +91,7 @@ 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 { @@ -473,6 +474,77 @@ function renderHomeLikeSurface(): { } } +/** Holds the chat POST of the first send until the test settles it. */ +function holdFirstSend(): PromiseWithResolvers { + const post = Promise.withResolvers() + 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 + 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 | undefined + + function HomeLike() { + result = useChat('ws-1', undefined) + return result.messages.length > 0 ? ( +
+ + + +
+ ) : ( +
+ +
+ ) + } + + const render = () => + act(() => { + root.render( + + + + ) + }) + 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 @@ -4703,61 +4775,63 @@ describe('useChat remount send recovery', () => { }) it('keeps the new-chat effort across the composer swap of a first send that fails', async () => { - useMothershipEffortStore.getState().reset() - const post = Promise.withResolvers() - 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 - }) - ;(globalThis as { IS_REACT_ACT_ENVIRONMENT?: boolean }).IS_REACT_ACT_ENVIRONMENT = true - queryClient = new QueryClient({ defaultOptions: { queries: { retry: false } } }) - const container = document.createElement('div') - const root = createRoot(container) - mountedRoots.push(root) - let result: ReturnType | undefined - /** Like home.tsx: the empty-state composer swaps for the chat view's once messages show. */ - function HomeLike() { - result = useChat('ws-1', undefined) - return result.messages.length > 0 ? ( -
- -
- ) : ( -
- -
- ) - } - act(() => { - root.render( - - - - ) - }) - const shownEffort = () => - container.querySelector('[aria-label="Reasoning effort"]')?.getAttribute('aria-description') + const post = holdFirstSend() + const surface = renderComposerSwap() act(() => useMothershipEffortStore.getState().setNewChatEffort('low')) - expect(shownEffort()).toBe('Low') + expect(surface.shownEffort()).toBe('Low') await act(async () => { - void result?.sendMessage('Plan the launch') + void surface.getResult().sendMessage('Plan the launch') }) await waitFor(() => state.postBodies.length === 1) expect(state.postBodies[0].effort).toBe('low') - expect(container.querySelector('section')).not.toBeNull() - expect(shownEffort()).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(() => container.querySelector('main') !== null) + await waitFor(() => surface.container.querySelector('main') !== null) expect(useMothershipEffortStore.getState().newChatEffort).toBe('low') - expect(shownEffort()).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)( diff --git a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts index 53bdc041bdb..627272e31ae 100644 --- a/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts +++ b/apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts @@ -1001,7 +1001,7 @@ export function useChat( new Set()) const streamReaderRef = useRef | null>(null) const chatIdRef = useRef(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 @@ -1010,8 +1010,9 @@ export function useChat( } }, []) /* 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, and drops only when the - surface leaves the new chat. */ + 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) @@ -1285,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) { @@ -3743,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,