Skip to content

Commit 69d9ac6

Browse files
authored
fix(mothership): close a pre-aborted SSE stream, keep the new-chat effort across a failed first send (#8754)
* 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. * 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.
1 parent 9a1aca0 commit 69d9ac6

8 files changed

Lines changed: 225 additions & 21 deletions

File tree

‎apps/sim/app/workspace/[workspaceId]/home/components/user-input/components/model-selector.tsx‎

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,5 @@
11
'use client'
22

3-
import { useEffect } from 'react'
43
import {
54
DropdownMenu,
65
DropdownMenuContent,
@@ -53,11 +52,6 @@ export function ModelSelector() {
5352
else setNewChatEffort(choice)
5453
}
5554

56-
useEffect(() => {
57-
if (chatId) return
58-
return () => useMothershipEffortStore.getState().setNewChatEffort(null)
59-
}, [chatId])
60-
6155
const effortLabel = options.find((option) => option.value === effort)?.label ?? effort
6256

6357
return (

‎apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx‎

Lines changed: 163 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -91,7 +91,9 @@ import { MothershipHandoffStorage } from '@/lib/core/utils/browser-storage'
9191
import { MOTHERSHIP_STREAM_REPLAY_HEADER } from '@/lib/mothership/constants'
9292
import type { MothershipStreamV1EventEnvelope } from '@/lib/mothership/generated/mothership-stream-v1'
9393
import { getChatResourceSelectionId } from '@/lib/mothership/resources/types'
94+
import { ChatSurfaceProvider } from '@/app/workspace/[workspaceId]/home/components/chat-surface-context'
9495
import { collectCitedMessageSources } from '@/app/workspace/[workspaceId]/home/components/message-content/message-sources'
96+
import { ModelSelector } from '@/app/workspace/[workspaceId]/home/components/user-input/components/model-selector'
9597
import {
9698
readQueuedSendHandoffState,
9799
writeQueuedSendHandoffState,
@@ -472,6 +474,77 @@ function renderHomeLikeSurface(): {
472474
}
473475
}
474476

477+
/** Holds the chat POST of the first send until the test settles it. */
478+
function holdFirstSend(): PromiseWithResolvers<Response> {
479+
const post = Promise.withResolvers<Response>()
480+
vi.stubGlobal('fetch', (input: RequestInfo | URL, init?: RequestInit) => {
481+
if (String(input) !== '/api/mothership/chat' || init?.method !== 'POST') {
482+
return fetchStub(input, init)
483+
}
484+
state.postBodies.push(JSON.parse(String(init.body)))
485+
return post.promise
486+
})
487+
return post
488+
}
489+
490+
/**
491+
* Mounts a chatless surface shaped like `home.tsx`, with real composers: the empty-state one
492+
* swaps for the chat view's once messages show, and the chat view's names the resolved chat.
493+
*/
494+
function renderComposerSwap(): {
495+
container: HTMLElement
496+
getResult: () => ReturnType<typeof useChat>
497+
shownEffort: () => string | null | undefined
498+
visit: (pathname: string) => void
499+
} {
500+
;(globalThis as { IS_REACT_ACT_ENVIRONMENT?: boolean }).IS_REACT_ACT_ENVIRONMENT = true
501+
useMothershipEffortStore.getState().reset()
502+
queryClient = new QueryClient({ defaultOptions: { queries: { retry: false } } })
503+
const container = document.createElement('div')
504+
const root = createRoot(container)
505+
mountedRoots.push(root)
506+
let result: ReturnType<typeof useChat> | undefined
507+
508+
function HomeLike() {
509+
result = useChat('ws-1', undefined)
510+
return result.messages.length > 0 ? (
511+
<section key='chat'>
512+
<ChatSurfaceProvider chatId={result.resolvedChatId}>
513+
<ModelSelector />
514+
</ChatSurfaceProvider>
515+
</section>
516+
) : (
517+
<main key='empty'>
518+
<ModelSelector />
519+
</main>
520+
)
521+
}
522+
523+
const render = () =>
524+
act(() => {
525+
root.render(
526+
<QueryClientProvider client={queryClient}>
527+
<HomeLike />
528+
</QueryClientProvider>
529+
)
530+
})
531+
render()
532+
533+
return {
534+
container,
535+
getResult: () => {
536+
if (result === undefined) throw new Error('Hook result is not ready')
537+
return result
538+
},
539+
shownEffort: () =>
540+
container.querySelector('[aria-label="Reasoning effort"]')?.getAttribute('aria-description'),
541+
visit: (pathname) => {
542+
mockUsePathname.mockReturnValue(pathname)
543+
render()
544+
},
545+
}
546+
}
547+
475548
/**
476549
* Mounts the hook under StrictMode with a handoff already in storage, mirroring
477550
* `home.tsx`'s consume-and-auto-send effect. This is the production-shaped
@@ -4701,6 +4774,96 @@ describe('useChat remount send recovery', () => {
47014774
])
47024775
})
47034776

4777+
it('keeps the new-chat effort across the composer swap of a first send that fails', async () => {
4778+
const post = holdFirstSend()
4779+
const surface = renderComposerSwap()
4780+
act(() => useMothershipEffortStore.getState().setNewChatEffort('low'))
4781+
expect(surface.shownEffort()).toBe('Low')
4782+
4783+
await act(async () => {
4784+
void surface.getResult().sendMessage('Plan the launch')
4785+
})
4786+
await waitFor(() => state.postBodies.length === 1)
4787+
expect(state.postBodies[0].effort).toBe('low')
4788+
expect(surface.container.querySelector('section')).not.toBeNull()
4789+
expect(surface.shownEffort()).toBe('Low')
4790+
4791+
await act(async () => {
4792+
post.reject(new TypeError('Failed to fetch'))
4793+
})
4794+
await waitFor(() => surface.container.querySelector('main') !== null)
4795+
4796+
expect(useMothershipEffortStore.getState().newChatEffort).toBe('low')
4797+
expect(surface.shownEffort()).toBe('Low')
4798+
})
4799+
4800+
it('keeps a new-chat effort picked while the first send is pending when that send fails', async () => {
4801+
const post = holdFirstSend()
4802+
const surface = renderComposerSwap()
4803+
act(() => useMothershipEffortStore.getState().setNewChatEffort('low'))
4804+
await act(async () => {
4805+
void surface.getResult().sendMessage('Plan the launch')
4806+
})
4807+
await waitFor(() => state.postBodies.length === 1)
4808+
act(() => useMothershipEffortStore.getState().setNewChatEffort('medium'))
4809+
4810+
await act(async () => {
4811+
post.reject(new TypeError('Failed to fetch'))
4812+
})
4813+
await waitFor(() => surface.container.querySelector('main') !== null)
4814+
4815+
expect(surface.shownEffort()).toBe('Medium')
4816+
})
4817+
4818+
it('starts the next new chat at the default after a first send stopped before admission', async () => {
4819+
const surface = renderComposerSwap()
4820+
act(() => useMothershipEffortStore.getState().setNewChatEffort('low'))
4821+
await act(async () => {
4822+
void surface.getResult().sendMessage('Plan the launch')
4823+
})
4824+
await waitFor(() => state.postBodies.length === 1)
4825+
await act(async () => {
4826+
await surface.getResult().stopGeneration()
4827+
})
4828+
await waitFor(() => surface.getResult().resolvedChatId === DEDUPED_CHAT_ID)
4829+
4830+
surface.visit(`/workspace/ws-1/chat/${DEDUPED_CHAT_ID}`)
4831+
surface.visit('/workspace/ws-1/home')
4832+
await waitFor(() => surface.container.querySelector('main') !== null)
4833+
4834+
expect(surface.shownEffort()).toBe('High')
4835+
})
4836+
4837+
it.each(['leaves the page', 'opens another chat'] as const)(
4838+
'drops an unsent new-chat effort when the surface %s',
4839+
(leave) => {
4840+
useMothershipEffortStore.getState().reset()
4841+
;(globalThis as { IS_REACT_ACT_ENVIRONMENT?: boolean }).IS_REACT_ACT_ENVIRONMENT = true
4842+
queryClient = new QueryClient({ defaultOptions: { queries: { retry: false } } })
4843+
const root = createRoot(document.createElement('div'))
4844+
mountedRoots.push(root)
4845+
function Surface({ chatId }: { chatId?: string }) {
4846+
useChat('ws-1', chatId)
4847+
return null
4848+
}
4849+
const render = (chatId?: string) =>
4850+
act(() =>
4851+
root.render(
4852+
<QueryClientProvider client={queryClient}>
4853+
<Surface chatId={chatId} />
4854+
</QueryClientProvider>
4855+
)
4856+
)
4857+
render()
4858+
useMothershipEffortStore.getState().setNewChatEffort('low')
4859+
4860+
if (leave === 'leaves the page') act(() => root.unmount())
4861+
else render('chat-other')
4862+
4863+
expect(useMothershipEffortStore.getState().newChatEffort).toBeNull()
4864+
}
4865+
)
4866+
47044867
it('loads the saved transcript once when its own stream completes', async () => {
47054868
const chatId = 'chat-own-completion'
47064869
const history: MothershipChatHistory = {

‎apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts‎

Lines changed: 13 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1001,14 +1001,22 @@ export function useChat(
10011001
new Set())
10021002
const streamReaderRef = useRef<ReadableStreamDefaultReader<Uint8Array> | null>(null)
10031003
const chatIdRef = useRef<string | undefined>(initialChatId)
1004-
/** Cleared on unmount, so a late rollback cannot hand a pick to a surface the user left. */
1004+
/** Cleared on unmount, so late async work cannot act on a surface the user left. */
10051005
const surfaceMountedRef = useRef(true)
10061006
useEffect(() => {
10071007
surfaceMountedRef.current = true
10081008
return () => {
10091009
surfaceMountedRef.current = false
10101010
}
10111011
}, [])
1012+
/* The new-chat effort pick belongs to this surface, not to one composer: it outlives the swap
1013+
from the empty-state composer to the chat view during a first send, so a withdrawn send
1014+
leaves it in place. It drops when the surface unmounts or switches chats, and when it adopts
1015+
a chat (`adoptResolvedChatId`). */
1016+
useEffect(() => {
1017+
if (initialChatId) return
1018+
return () => useMothershipEffortStore.getState().setNewChatEffort(null)
1019+
}, [initialChatId])
10121020
const tableViewContextsRef = useRef({
10131021
scopeId: desktopScopeId,
10141022
views: new Map<string, MothershipTableViewContext>(),
@@ -1278,6 +1286,10 @@ export function useChat(
12781286
const resolvedDesktopScopeId = desktopChatScopeId(scopeKey, chatId)
12791287
if (wasPending) {
12801288
useChatPanelStore.getState().migrate(pendingDesktopScopeId, resolvedDesktopScopeId)
1289+
// Leaving the new chat. An admitted send has already moved the pick onto its chat; any
1290+
// other way out (a Stop before admission, a recovered handoff) must not carry it into
1291+
// the next new chat.
1292+
useMothershipEffortStore.getState().setNewChatEffort(null)
12811293
}
12821294
const activeActivityTracker = resourceActivityTrackerRef.current
12831295
if (activeActivityTracker?.generation === streamGenRef.current) {
@@ -3736,16 +3748,6 @@ export function useChat(
37363748
}
37373749

37383750
const rollbackOptimisticSend = () => {
3739-
// A withdrawn first send hands its pick back to the new-chat composer for the retry,
3740-
// only while that surface is still open on the new chat.
3741-
if (
3742-
!requestChatId &&
3743-
effortChoice &&
3744-
surfaceMountedRef.current &&
3745-
!chatIdRef.current &&
3746-
!selectedChatIdRef.current
3747-
)
3748-
useMothershipEffortStore.getState().setNewChatEffort(effortChoice)
37493751
if (requestChatId) {
37503752
upsertChatHistory(requestChatId, (current) => ({
37513753
...current,

‎apps/sim/hooks/queries/mothership-chats.test.ts‎

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@ const { suspendBrowserScope, suspendTerminalScope, clearChat } = vi.hoisted(() =
1111
}))
1212

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

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

2727
import type { MothershipEffort } from '@/lib/mothership/model-options'
2828
import {
29+
fetchMothershipChatHistory,
2930
useDeleteMothershipChats,
3031
useSetMothershipChatEffort,
3132
} from '@/hooks/queries/mothership-chats'
@@ -100,6 +101,26 @@ describe('tasks query boundary parsing', () => {
100101
})
101102
})
102103

104+
it('loads a chat from a server that predates the effort field', async () => {
105+
vi.mocked(fetch).mockResolvedValueOnce(
106+
jsonResponse({
107+
success: true,
108+
chat: {
109+
id: 'chat-1',
110+
title: null,
111+
mode: 'agent',
112+
messages: [],
113+
activeStreamId: null,
114+
resources: [],
115+
},
116+
})
117+
)
118+
119+
const history = await fetchMothershipChatHistory('chat-1')
120+
121+
expect(history.effort).toBeNull()
122+
})
123+
103124
it('keeps the latest effort pick when an earlier queued save of the same value fails', async () => {
104125
const tanstack =
105126
await vi.importActual<typeof import('@tanstack/react-query')>('@tanstack/react-query')

‎apps/sim/lib/api/contracts/mothership-chats.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -420,7 +420,8 @@ export const getMothershipChatResponseSchema = z.object({
420420
messages: z.array(z.unknown()),
421421
activeStreamId: z.string().nullable(),
422422
resources: z.array(z.unknown()),
423-
effort: mothershipChatEffortChoiceSchema,
423+
/** Optional so a client still loads chats from a server that predates the field. */
424+
effort: mothershipChatEffortChoiceSchema.optional(),
424425
createdAt: z.union([z.string(), z.date()]).nullable().optional(),
425426
updatedAt: z.union([z.string(), z.date()]).nullable().optional(),
426427
streamSnapshot: mothershipChatStreamSnapshotSchema.optional(),

‎apps/sim/lib/events/sse-endpoint.test.ts‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -133,6 +133,22 @@ describe('createWorkspaceSSE', () => {
133133
expect(unsubscribe).toHaveBeenCalledTimes(1)
134134
})
135135

136+
it('never subscribes when the request aborted before the stream started', async () => {
137+
const controller = new AbortController()
138+
controller.abort()
139+
const subscribe = vi.fn(() => () => {})
140+
const { body } = await openConnection(controller.signal, [{ subscribe }])
141+
let closed = false
142+
void drain(body).then(() => {
143+
closed = true
144+
})
145+
146+
await vi.advanceTimersByTimeAsync(0)
147+
148+
expect(closed).toBe(true)
149+
expect(subscribe).not.toHaveBeenCalled()
150+
})
151+
136152
it('runs every teardown when one unsubscribe throws', async () => {
137153
const first = vi.fn(() => {
138154
throw new Error('unsubscribe failed')

‎apps/sim/lib/events/sse-endpoint.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -213,6 +213,13 @@ export function createSSEStream(request: NextRequest, config: SSEStreamConfig):
213213
)
214214
}
215215

216+
// An abort listener never fires for a signal that is already aborted, so a client that left
217+
// while the route was authorizing would otherwise hold its subscriptions until rotation.
218+
if (request.signal.aborted) {
219+
close('aborted')
220+
return
221+
}
222+
216223
try {
217224
// The runtime sends the status and headers with the first body chunk, so the stream writes
218225
// one as soon as it opens. A client reads its state once the stream opens, so it opens once

‎apps/sim/stores/mothership-effort/store.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,8 +15,8 @@ interface MothershipEffortState {
1515
setModel: (model: ModelSelection['model']) => void
1616
setFastMode: (fastMode: boolean) => void
1717
/**
18-
* The effort picked in a composer whose chat does not exist yet. Its first send records
19-
* it on the new chat; leaving that composer unsent drops it.
18+
* The effort picked on a chat surface whose chat does not exist yet. Its first send records
19+
* it on the new chat; leaving that surface unsent drops it.
2020
*/
2121
newChatEffort: MothershipEffort | null
2222
setNewChatEffort: (effort: MothershipEffort | null) => void

0 commit comments

Comments
 (0)