From ec1c01d1bd33388ab59d3f72c8cef73ef99c86c0 Mon Sep 17 00:00:00 2001 From: Guilhem Lemouel Date: Tue, 15 Sep 2026 13:05:18 +0200 Subject: [PATCH] fix(chat): re-open a failed stream rather than freeing the composer MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `followJob` absorbs the server ending a stream on its own clock, but not the request to open one failing — a restarting server, a 502. Ending the turn there unlocked the composer while the flow kept running, and the next turn would write the same agent memory. Those attempts are retried from the offset already read, so the answer resumes rather than replaying, and the turn stays busy across the gap. Bounded, so a genuinely gone endpoint still reports itself instead of retrying forever under a chat that looks live. Co-Authored-By: Claude Opus 5 (1M context) --- .../conversations/FlowChatManager.svelte.ts | 120 +++++++++++------- .../conversations/FlowChatManager.test.ts | 31 ++++- 2 files changed, 103 insertions(+), 48 deletions(-) diff --git a/frontend/src/lib/components/flows/conversations/FlowChatManager.svelte.ts b/frontend/src/lib/components/flows/conversations/FlowChatManager.svelte.ts index 833f96c946..3840e024cb 100644 --- a/frontend/src/lib/components/flows/conversations/FlowChatManager.svelte.ts +++ b/frontend/src/lib/components/flows/conversations/FlowChatManager.svelte.ts @@ -65,6 +65,14 @@ type TurnStatus = { jobId?: string } +/** + * How many times a turn re-opens its stream after the request itself failed, and the first + * delay between attempts — doubled each time. Enough to outlast a server restart; bounded so + * a genuinely gone endpoint reports itself rather than retrying under a chat that looks live. + */ +const FOLLOW_RETRIES = 4 +const FOLLOW_RETRY_DELAY_MS = 500 + /** * The agent events the worker streams, in the shape the transcript applies. The SDK names * the same six events after the wire protocol; this is the rest of the app's vocabulary. @@ -1116,6 +1124,10 @@ export class FlowChatManager { * rather than starting a second run — which is what this used to do, leaving two runs * writing one conversation. It reconnects on an unclean close too, and buffers a chunk * that ends mid-line until the rest arrives. + * + * What it does not absorb is the stream request itself failing — a restarting server, a + * 502. Giving up there would free the composer while the flow runs on, and the next turn + * would write the same agent memory, so those are retried from the offset already read. */ async #followJob(currentConversationId: string, jobId: string) { const runtime = this.#liveRuntime(currentConversationId) @@ -1132,57 +1144,75 @@ export class FlowChatManager { pollDelayMs: get(enterpriseLicense) ? 50 : undefined }) - try { - for await (const update of followJob(api, jobId, { signal: controller.signal })) { - if (update.type === 'stream') { - // Stop polling since we are receiving last step streaming - this.stopPolling(currentConversationId) - // One chunk can carry several events, so each is applied in turn: a - // chunk holding a call and its result must produce both. - for (const streamed of update.events) { - const event = toStreamEvent(streamed) - if (event.kind === 'reasoning') { - status.isReasoningActive = true - runtime.reasoningReveal.push(event.content) - } else if (event.kind === 'token') { - runtime.replyReveal.push(event.content) - } else { - // Whatever the pacing still holds belongs to the row before the tool — - // thinking that led straight to the call included — so it is revealed - // before the event that closes that row. - if (event.kind === 'tool_call' || event.kind === 'tool_execution') { - this.#flushReveals(currentConversationId) - status.currentReasoning = '' - status.isReasoningActive = false + // Kept across attempts so a reconnect resumes after what is already on screen. + let streamOffset: number | undefined + for (let attempt = 0; ; attempt++) { + try { + for await (const update of followJob(api, jobId, { + signal: controller.signal, + streamOffset, + onOffset: (offset) => (streamOffset = offset) + })) { + if (update.type === 'stream') { + // Stop polling since we are receiving last step streaming + this.stopPolling(currentConversationId) + // One chunk can carry several events, so each is applied in turn: a + // chunk holding a call and its result must produce both. + for (const streamed of update.events) { + const event = toStreamEvent(streamed) + if (event.kind === 'reasoning') { + status.isReasoningActive = true + runtime.reasoningReveal.push(event.content) + } else if (event.kind === 'token') { + runtime.replyReveal.push(event.content) + } else { + // Whatever the pacing still holds belongs to the row before the tool — + // thinking that led straight to the call included — so it is revealed + // before the event that closes that row. + if (event.kind === 'tool_call' || event.kind === 'tool_execution') { + this.#flushReveals(currentConversationId) + status.currentReasoning = '' + status.isReasoningActive = false + } + const step = applyStreamEvent( + { rows: this.#rowsOf(currentConversationId), state: runtime.turn }, + event, + this.#newRowId + ) + this.#rowsById[currentConversationId] = step.rows + runtime.turn = step.state } - const step = applyStreamEvent( - { rows: this.#rowsOf(currentConversationId), state: runtime.turn }, - event, - this.#newRowId - ) - this.#rowsById[currentConversationId] = step.rows - runtime.turn = step.state } + continue } + // Anything still buffered would be dropped by the temp-row sweep below. + this.#flushReveals(currentConversationId) + // Do a final poll to get all messages from database + await this.pollConversationMessages(currentConversationId, { + removeTempMessages: true + }) + this.endTurn(currentConversationId, { settled: true }) + } + return + } catch (error) { + // A Stop, a conversation's turn ending, or the chat going away — the turn was + // already settled by whoever aborted it. + if (controller.signal.aborted) return + console.error('Error following the flow job:', error) + if (attempt < FOLLOW_RETRIES) { + await new Promise((resolve) => setTimeout(resolve, FOLLOW_RETRY_DELAY_MS * 2 ** attempt)) + if (controller.signal.aborted) return continue } - // Anything still buffered would be dropped by the temp-row sweep below. - this.#flushReveals(currentConversationId) - // Do a final poll to get all messages from database - await this.pollConversationMessages(currentConversationId, { - removeTempMessages: true - }) - this.endTurn(currentConversationId, { settled: true }) + // What the server said, not a constant: a turn dying on an expired session or a + // job the server cannot find is only actionable if the reader is told which. + sendUserToast( + `Stream error: ${error instanceof Error ? error.message : String(error)}`, + true + ) + this.endTurn(currentConversationId) + return } - } catch (error) { - // A Stop, a conversation's turn ending, or the chat going away — the turn was - // already settled by whoever aborted it. - if (controller.signal.aborted) return - console.error('Error following the flow job:', error) - // What the server said, not a constant: a turn dying on an expired session or a - // job the server cannot find is only actionable if the reader is told which. - sendUserToast(`Stream error: ${error instanceof Error ? error.message : String(error)}`, true) - this.endTurn(currentConversationId) } } diff --git a/frontend/src/lib/components/flows/conversations/FlowChatManager.test.ts b/frontend/src/lib/components/flows/conversations/FlowChatManager.test.ts index a0871c4786..c7f0bf0db8 100644 --- a/frontend/src/lib/components/flows/conversations/FlowChatManager.test.ts +++ b/frontend/src/lib/components/flows/conversations/FlowChatManager.test.ts @@ -20,7 +20,8 @@ vi.mock('$lib/stores', () => ({ /** Each `streamJob` the turn opens, and the updates the next one answers with. */ const { streamCalls, streamScript } = vi.hoisted(() => ({ streamCalls: [] as { jobId: string; streamOffset: number | undefined }[], - streamScript: [] as unknown[][] + /** Per opened stream: the updates it answers with, or 'throw' to fail the request. */ + streamScript: [] as (unknown[] | 'throw')[] })) // Only the transport is faked. `followJob` — which owns the re-attach, the offset and the @@ -30,7 +31,9 @@ vi.mock('windmill-chat', async (importOriginal) => { class FakeApi { async *streamJob(jobId: string, options: { streamOffset?: number } = {}) { streamCalls.push({ jobId, streamOffset: options.streamOffset }) - for (const update of streamScript.shift() ?? []) yield update + const next = streamScript.shift() + if (next === 'throw') throw new Error('stream request failed') + for (const update of next ?? []) yield update } } return { ...actual, WindmillChatApi: FakeApi } @@ -173,7 +176,7 @@ describe('an SSE timeout re-attaches instead of re-running', () => { delete (globalThis as any).location }) - function turnWith(script: unknown[][]) { + function turnWith(script: (unknown[] | 'throw')[]) { streamScript.push(...script) const manager = (live = managerWithRows()) const onRunFlow = vi.fn(async () => 'job-1') @@ -214,6 +217,28 @@ describe('an SSE timeout re-attaches instead of re-running', () => { expect(manager.currentJobId).toBe('job-1') }) + /** + * `followJob` absorbs the server ending a stream, but not the request to open one + * failing — a restarting server, a 502. Ending the turn there would free the composer + * while the flow runs on, and the next turn would write the same agent memory. + */ + it('re-opens a failed stream from its offset rather than ending the turn', async () => { + const { manager } = turnWith([ + [{ type: 'update', stream_offset: 7 }], + 'throw', + // Reachable again, with nothing more to say yet. + [] + ]) + + manager.inputMessage = 'ask something' + await manager.sendMessage(undefined, undefined, 'a') + + await vi.waitFor(() => expect(streamCalls.length).toBeGreaterThanOrEqual(3), { timeout: 5000 }) + + expect(streamCalls[2]).toEqual({ jobId: 'job-1', streamOffset: 7 }) + expect(manager.isConversationBusy('a')).toBe(true) + }) + /** * A chunk is not guaranteed to end on a line boundary. Split mid-JSON, the two halves * were parsed separately and both discarded, losing that token from the answer.