Skip to content

Commit 5ab33dd

Browse files
authored
fix(chat): reconnect silently stalled live streams (#8272)
1 parent 293e883 commit 5ab33dd

4 files changed

Lines changed: 208 additions & 2 deletions

File tree

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

Lines changed: 133 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1696,6 +1696,139 @@ describe('useChat remount send recovery', () => {
16961696
expect(state.postBodies[0]).toHaveProperty('effort')
16971697
})
16981698

1699+
it.each(['initial', 'tail'] as const)(
1700+
'recovers a silent %s connection after a tool group without refresh or resending',
1701+
async (connection) => {
1702+
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
1703+
let unmount: (() => void) | undefined
1704+
try {
1705+
const history: MothershipChatHistory = {
1706+
id: 'chat-silent-stream',
1707+
title: 'Silent stream',
1708+
messages: [],
1709+
activeStreamId: null,
1710+
resources: [],
1711+
}
1712+
const cancelled = vi.fn()
1713+
const cursors: string[] = []
1714+
let recovered = false
1715+
let tailReads = 0
1716+
let streamId = ''
1717+
const textEvent = (): MothershipStreamV1EventEnvelope => ({
1718+
v: 1,
1719+
seq: 3,
1720+
ts: new Date().toISOString(),
1721+
type: 'text',
1722+
stream: { streamId },
1723+
payload: { channel: 'assistant', text: 'The work continued.' },
1724+
})
1725+
vi.stubGlobal('fetch', async (input: RequestInfo | URL, init?: RequestInit) => {
1726+
const url = String(input)
1727+
if (url === '/api/mothership/chat' && init?.method === 'POST') {
1728+
const sent = JSON.parse(String(init.body))
1729+
state.postBodies.push(sent)
1730+
streamId = sent.userMessageId
1731+
const events: MothershipStreamV1EventEnvelope[] = [
1732+
{
1733+
v: 1,
1734+
seq: 1,
1735+
ts: new Date().toISOString(),
1736+
type: 'tool',
1737+
stream: { streamId },
1738+
payload: {
1739+
phase: 'call',
1740+
executor: 'go',
1741+
mode: 'sync',
1742+
toolName: 'run_code',
1743+
toolCallId: 'finished-tool',
1744+
arguments: { code: 'return 1' },
1745+
},
1746+
},
1747+
{
1748+
v: 1,
1749+
seq: 2,
1750+
ts: new Date().toISOString(),
1751+
type: 'tool',
1752+
stream: { streamId },
1753+
payload: {
1754+
phase: 'result',
1755+
toolName: 'run_code',
1756+
toolCallId: 'finished-tool',
1757+
success: true,
1758+
output: { value: 1 },
1759+
},
1760+
},
1761+
]
1762+
return new Response(
1763+
new ReadableStream<Uint8Array>({
1764+
start(controller) {
1765+
for (const event of events)
1766+
controller.enqueue(
1767+
new TextEncoder().encode(`data: ${JSON.stringify(event)}\n\n`)
1768+
)
1769+
if (connection === 'tail') controller.close()
1770+
},
1771+
cancel: cancelled,
1772+
}),
1773+
{
1774+
headers: {
1775+
'Content-Type': 'text/event-stream',
1776+
'x-mothership-chat-id': history.id,
1777+
},
1778+
}
1779+
)
1780+
}
1781+
if (url.includes('/api/mothership/chat/stream')) {
1782+
const params = new URL(url, 'https://sim.test').searchParams
1783+
if (params.get('batch') === 'true') {
1784+
cursors.push(params.get('after') ?? '')
1785+
recovered = connection === 'initial' || tailReads > 0
1786+
return Response.json({
1787+
success: true,
1788+
status: 'streaming',
1789+
events: recovered ? [{ eventId: 3, streamId, event: textEvent() }] : [],
1790+
})
1791+
}
1792+
tailReads++
1793+
return new Response(new ReadableStream<Uint8Array>({ cancel: cancelled }), {
1794+
headers: { 'Content-Type': 'text/event-stream' },
1795+
})
1796+
}
1797+
return fetchStub(input, init)
1798+
})
1799+
const mounted = renderUseChatInChat(history.id, history)
1800+
unmount = mounted.unmount
1801+
const { getResult } = mounted
1802+
await act(async () => {
1803+
void getResult().sendMessage('Keep working')
1804+
})
1805+
await act(async () => vi.advanceTimersByTimeAsync(0))
1806+
expect(
1807+
getResult()
1808+
.messages.flatMap((message) => message.contentBlocks ?? [])
1809+
.find((block) => block.toolCall?.id === 'finished-tool')?.toolCall?.status
1810+
).toBe('success')
1811+
expect(recovered).toBe(false)
1812+
await act(async () => vi.advanceTimersByTimeAsync(45_000))
1813+
expect(recovered).toBe(true)
1814+
expect(cursors.every((cursor) => cursor === '2')).toBe(true)
1815+
expect(cancelled).toHaveBeenCalledTimes(1)
1816+
expect(
1817+
getResult()
1818+
.messages.filter((message) => message.role === 'assistant')
1819+
.map((message) => message.content)
1820+
).toEqual(['The work continued.'])
1821+
expect(getResult().isSending).toBe(true)
1822+
expect(getResult().error).toBeNull()
1823+
expect(state.postBodies).toHaveLength(1)
1824+
expect(state.abortBodies).toHaveLength(0)
1825+
} finally {
1826+
unmount?.()
1827+
vi.useRealTimers()
1828+
}
1829+
}
1830+
)
1831+
16991832
it('recovers a running turn after reconnect exhaustion without reloading or resending', async () => {
17001833
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
17011834
try {

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -287,6 +287,8 @@ const RECONNECT_BASE_DELAY_MS = 1000
287287
const RECONNECT_MAX_DELAY_MS = 30_000
288288
const RECONNECT_EXHAUSTED_RECHECK_MS = 30_000
289289
const STREAM_BATCH_FETCH_TIMEOUT_MS = 10_000
290+
/** Both live transports heartbeat every 15s; three missed heartbeats trigger cursor recovery. */
291+
const STREAM_IDLE_TIMEOUT_MS = 45_000
290292
const STREAM_CHAT_ID_RESOLVE_TIMEOUT_MS = 10_000
291293
const CHAT_HISTORY_RECOVERY_TIMEOUT_MS = 10_000
292294
const STOP_REQUEST_TIMEOUT_MS = 15_000
@@ -2219,6 +2221,7 @@ export function useChat(
22192221

22202222
try {
22212223
await readSSELines(reader, {
2224+
idleTimeoutMs: STREAM_IDLE_TIMEOUT_MS,
22222225
onData: (raw) => {
22232226
if (state.sawCompleteEvent) return true
22242227
if (ops.isStale()) return

‎apps/sim/lib/core/utils/sse.test.ts‎

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -551,6 +551,55 @@ describe('readSSEEvents', () => {
551551
})
552552

553553
describe('readSSELines', () => {
554+
it('rejects a silent open connection and cancels its reader without waiting for cancellation', async () => {
555+
vi.useFakeTimers()
556+
const cancel = vi.fn(() => new Promise<void>(() => {}))
557+
const reader = new ReadableStream<Uint8Array>({ cancel }).getReader()
558+
const rejected = vi.fn()
559+
const reading = readSSELines(reader, { onData: vi.fn(), idleTimeoutMs: 45_000 }).catch(rejected)
560+
try {
561+
await vi.advanceTimersByTimeAsync(45_000)
562+
expect(rejected).toHaveBeenCalledWith(
563+
expect.objectContaining({ name: 'SSEIdleTimeoutError' })
564+
)
565+
expect(cancel).toHaveBeenCalledTimes(1)
566+
expect(vi.getTimerCount()).toBe(0)
567+
} finally {
568+
void reader.cancel()
569+
await reading
570+
vi.useRealTimers()
571+
}
572+
})
573+
574+
it('counts keepalive comments as activity while a long tool is running', async () => {
575+
vi.useFakeTimers()
576+
let controller!: ReadableStreamDefaultController<Uint8Array>
577+
const stream = new ReadableStream<Uint8Array>({
578+
start: (value) => {
579+
controller = value
580+
},
581+
})
582+
const onData = vi.fn()
583+
const rejected = vi.fn()
584+
const reading = readSSELines(stream, { onData, idleTimeoutMs: 45_000 }).catch(rejected)
585+
try {
586+
for (let heartbeat = 0; heartbeat < 6; heartbeat++) {
587+
await vi.advanceTimersByTimeAsync(15_000)
588+
controller.enqueue(new TextEncoder().encode(': keepalive\n\n'))
589+
await vi.advanceTimersByTimeAsync(0)
590+
}
591+
expect(rejected).not.toHaveBeenCalled()
592+
expect(onData).not.toHaveBeenCalled()
593+
controller.enqueue(new TextEncoder().encode('data: tool finished\n\n'))
594+
controller.close()
595+
await reading
596+
expect(onData).toHaveBeenCalledWith('tool finished')
597+
expect(vi.getTimerCount()).toBe(0)
598+
} finally {
599+
vi.useRealTimers()
600+
}
601+
})
602+
554603
it('delivers raw (un-parsed) data payloads', async () => {
555604
const stream = streamFromStringChunks(['data: raw-one\n\n', 'data: {"keep":"asString"}\n\n'])
556605
const lines: string[] = []

‎apps/sim/lib/core/utils/sse.ts‎

Lines changed: 23 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,8 @@ export interface ReadSSELinesOptions {
6565
onData: (rawData: string) => SSEStopSignal
6666
/** Aborts the read; checked before each chunk and between events. */
6767
signal?: AbortSignal
68+
/** Reconnect-capable consumers can bound a silent read; comments also keep it alive. */
69+
idleTimeoutMs?: number
6870
}
6971

7072
/**
@@ -133,16 +135,35 @@ function stripCarriageReturn(line: string): string {
133135
* @param options - The `onData` callback plus an optional `signal`.
134136
*/
135137
export async function readSSELines(source: SSESource, options: ReadSSELinesOptions): Promise<void> {
136-
const { onData, signal } = options
138+
const { onData, signal, idleTimeoutMs } = options
137139
const { reader, ownsLock } = toReader(source)
138140
const decoder = new TextDecoder()
139141
let buffer = ''
140142

143+
const readChunk = async () => {
144+
if (idleTimeoutMs === undefined) return reader.read()
145+
let timer: ReturnType<typeof setTimeout> | undefined
146+
try {
147+
return await new Promise<ReadableStreamReadResult<Uint8Array>>((resolve, reject) => {
148+
timer = setTimeout(() => {
149+
const error = new Error(`SSE connection was silent for ${idleTimeoutMs}ms`)
150+
error.name = 'SSEIdleTimeoutError'
151+
reject(error)
152+
/** Release the transport without waiting for an unresponsive underlying source. */
153+
void reader.cancel(error).catch(() => {})
154+
}, idleTimeoutMs)
155+
reader.read().then(resolve, reject)
156+
})
157+
} finally {
158+
clearTimeout(timer)
159+
}
160+
}
161+
141162
try {
142163
while (true) {
143164
if (signal?.aborted) break
144165

145-
const { done, value } = await reader.read()
166+
const { done, value } = await readChunk()
146167

147168
buffer += done ? decoder.decode() : decoder.decode(value, { stream: true })
148169
const lines = buffer.split('\n')

0 commit comments

Comments
 (0)