From d3344eb3cd1c3decf41413c9da0cec3b903f2f6e Mon Sep 17 00:00:00 2001 From: Guilhem Lemouel Date: Tue, 15 Sep 2026 18:56:35 +0200 Subject: [PATCH] fix(chat): read a conversation forward from the start, and clear the polling deadline A poll with no row to resume from sent no cursor at all, and the messages endpoint answers that with the newest page rather than the oldest, so the rows before it were never read and the sweep dropped the temp rows covering them. Sequences start at 1, so reading from 0 puts every request on the forward branch. A read that stops at the request cap is now kept: each batch is a run of rows from the cursor, so the next tick carries on from where it stopped. The last poll of a turn has no tick after it, and there its prefix is dropped instead, since it would sit beside the temp rows showing the turn's start twice. `startPolling` armed a two-minute timeout it never held onto, and `stopPolling` cleared only the interval. A turn ending inside two minutes left the timeout armed to fire into the next turn on the same conversation and stop its polling part-way, which a non-streaming turn feels as intermediate rows drying up. The handle now lives on the turn's runtime and is cleared with the interval. Co-Authored-By: Claude Opus 5 (1M context) --- .../lib/components/flows/agentFormFields.ts | 8 +-- .../conversations/FlowChatManager.svelte.ts | 71 ++++++++++--------- .../conversations/FlowChatManager.test.ts | 31 +++++++- .../conversations/flowChatViewHost.svelte.ts | 8 +-- 4 files changed, 76 insertions(+), 42 deletions(-) diff --git a/frontend/src/lib/components/flows/agentFormFields.ts b/frontend/src/lib/components/flows/agentFormFields.ts index 4923e8661c..04076cd7e6 100644 --- a/frontend/src/lib/components/flows/agentFormFields.ts +++ b/frontend/src/lib/components/flows/agentFormFields.ts @@ -197,10 +197,6 @@ export function agentFieldIsSet( return true } -/** - * Whether the current schema carries this field at all. A linked step's schema is reduced to the - * flow-local inputs, which is what collapses its form to the Messages group on its own. - */ /** * Whether a run of this step would stream its answer, mirroring the worker's * `has_stream = user_wants_streaming && is_text_output`. Absence means on @@ -226,6 +222,10 @@ export function agentStreamingEnabled(value: Record | undefined): b return streaming?.value !== false } +/** + * Whether the current schema carries this field at all. A linked step's schema is reduced to the + * flow-local inputs, which is what collapses its form to the Messages group on its own. + */ export function agentFieldAppliesTo( spec: AgentFieldSpec, schemaProperties: Record | undefined diff --git a/frontend/src/lib/components/flows/conversations/FlowChatManager.svelte.ts b/frontend/src/lib/components/flows/conversations/FlowChatManager.svelte.ts index 2788a90f46..8962c77e82 100644 --- a/frontend/src/lib/components/flows/conversations/FlowChatManager.svelte.ts +++ b/frontend/src/lib/components/flows/conversations/FlowChatManager.svelte.ts @@ -75,12 +75,14 @@ const FOLLOW_RETRY_DELAY_MS = 500 /** How often a turn with no stream asks its job whether the run is over. */ const SETTLE_POLL_MS = 2000 +/** How long a conversation is polled before the reader is assumed to have left it running. */ +const POLL_MAX_MS = 2 * 60 * 1000 /** Rows per request when a poll reads a conversation. A batch shorter than this is how the * endpoint says there are no more. */ const POLL_PAGE_SIZE = 50 -/** Requests one poll will make. Bounds a conversation with more rows past the cursor than a - * poll should read in one go; a cursor that stops moving is caught separately. */ -const POLL_MAX_PAGES = 20 +/** Most requests one poll will make. Bounds how much of a conversation one tick reads; what it + * stops short of is left to the next tick, which resumes from the cursor it reached. */ +const POLL_MAX_REQUESTS = 20 /** * The agent events the worker streams, in the shape the transcript applies. The SDK names @@ -122,6 +124,9 @@ type TurnRuntime = { /** Stops this turn's `followJob`. */ follow?: AbortController pollingInterval?: ReturnType + /** Stops the interval above at `POLL_MAX_MS`, and is cleared with it: one left armed by a + * turn that ended early would fire into the next turn on the same conversation. */ + pollingDeadline?: ReturnType // What the turn has written so far — which row is open, and the text in it. Held here // rather than in the stream handler's locals: the typewriter reveals on animation // frames, long after the chunk that delivered the text was applied. @@ -162,11 +167,11 @@ export class FlowChatManager { */ allowsParallelTurns = $state(false) - /** Each conversation's rows, live ones included, so a turn keeps writing while the - * reader is in another chat. Doubles as the load cache: rows here are never re-fetched. */ /** Bumped when the chat is re-pointed at another flow. Work started before that must not * write what it fetched into the chat that replaced it. */ #generation = 0 + /** Each conversation's rows, live ones included, so a turn keeps writing while the + * reader is in another chat. Doubles as the load cache: rows here are never re-fetched. */ #rowsById = $state>({}) /** How far back each conversation has been paged. Held per conversation for the same * reason the rows are: a cached chat keeps its scrollback when the reader returns to it, @@ -471,10 +476,7 @@ export class FlowChatManager { runtime.turn = emptyTurnState(conversationId) runtime.follow?.abort() runtime.follow = undefined - if (runtime.pollingInterval) { - clearInterval(runtime.pollingInterval) - runtime.pollingInterval = undefined - } + this.stopPolling(conversationId) } const status = this.#liveStatus(conversationId) status.currentReasoning = '' @@ -935,11 +937,12 @@ export class FlowChatManager { // hand back its earliest and leave its answer behind, and the sweep below drops the // temp rows that were standing in for it. const response: ChatMessage[] = [] - let afterSeq = this.getLastPersistedMessageSeq(conversationId) + // Sequences start at 1, so a conversation with no row to resume from reads from 0 + // rather than with no cursor at all — without one the endpoint answers with the + // newest page instead of the oldest, which is not a prefix of anything. + let afterSeq = this.getLastPersistedMessageSeq(conversationId) ?? 0 let readWhole = false - let stalled = false - let requests = 0 - for (let page = 0; page < POLL_MAX_PAGES; page++) { + for (let request = 0; request < POLL_MAX_REQUESTS; request++) { const batch = await FlowConversationsService.listConversationMessages({ workspace: this.#workspace()!, conversationId: conversationId, @@ -950,31 +953,34 @@ export class FlowChatManager { // An interval tick already dispatched outlives `clearInterval`, and a turn's final // poll outlives its abort — either would put a forgotten flow's rows back. if (startedIn !== this.#generation) return - requests++ response.push(...batch) if (batch.length < POLL_PAGE_SIZE) { readWhole = true break } const furthest = Math.max(...batch.map((m) => m.created_seq)) - // A page that moved nothing would ask for the same rows forever, and is no more a - // finished read than the cap above. - if (afterSeq !== undefined && furthest <= afterSeq) { - stalled = true + if (furthest <= afterSeq) { + // A full page that leaves the cursor where it was would be asked for again on + // every tick and answered the same way, so nothing here or later can advance. + console.warn( + `Stopped reading conversation ${conversationId} at seq ${afterSeq}: a full ` + + `page of ${POLL_PAGE_SIZE} rows did not move the cursor` + ) + this.stopPolling(conversationId) break } afterSeq = furthest } - if (!readWhole) { - // Neither appended nor swept: half a conversation beside the temp rows standing - // in for it reads worse than those rows alone, and sweeping would drop the only - // copy of what was never read. Held rows are never re-fetched while the chat is - // open (see `loadMessages`), so a reload is what recovers this. + if (!readWhole && options?.removeTempMessages) { + // The last poll of a turn, and no tick comes after it to carry on from where this + // one stopped. Its temp rows have to stay, being the only copy of what the read + // did not reach, and the prefix it did read would sit beside them showing the + // start of the turn twice — so that prefix is dropped rather than the rows. A + // reload is what recovers it: `loadMessages` leaves rows already held alone. console.warn( - `Stopped reading conversation ${conversationId} after ${requests} request(s) ` + - `(${stalled ? 'the cursor stopped advancing' : `cap of ${POLL_MAX_PAGES}`}, ` + - `${response.length} rows, up to seq ${afterSeq}); leaving the transcript as it is` + `Read ${response.length} rows of conversation ${conversationId} at the end of a ` + + `turn without reaching its last one; leaving the transcript as it is` ) return } @@ -1011,12 +1017,9 @@ export class FlowChatManager { runtime.pollingInterval = setInterval(() => { this.pollConversationMessages(conversationId, { isNewConversation }) }, 500) // Poll every 0.5 seconds - setTimeout( - () => { - this.stopPolling(conversationId) - }, - 2 * 60 * 1000 - ) // Stop polling after 2 minutes + runtime.pollingDeadline = setTimeout(() => { + this.stopPolling(conversationId) + }, POLL_MAX_MS) } private stopPolling(conversationId: string) { @@ -1025,6 +1028,10 @@ export class FlowChatManager { clearInterval(runtime.pollingInterval) runtime.pollingInterval = undefined } + if (runtime?.pollingDeadline) { + clearTimeout(runtime.pollingDeadline) + runtime.pollingDeadline = undefined + } } // Message sending diff --git a/frontend/src/lib/components/flows/conversations/FlowChatManager.test.ts b/frontend/src/lib/components/flows/conversations/FlowChatManager.test.ts index e624ca7bf3..1ece20471d 100644 --- a/frontend/src/lib/components/flows/conversations/FlowChatManager.test.ts +++ b/frontend/src/lib/components/flows/conversations/FlowChatManager.test.ts @@ -204,6 +204,9 @@ describe('reading a turn longer than one page', () => { // re-fetch the same page and still look right from the rows alone. const calls = vi.mocked(FlowConversationsService.listConversationMessages).mock.calls expect((calls[1][0] as any).afterSeq).toBe(50) + // Nothing to resume from still reads forward from the start: with no cursor at all the + // endpoint answers with the newest page, which would skip everything before it. + expect((calls[0][0] as any).afterSeq).toBe(0) }) /** @@ -211,7 +214,7 @@ describe('reading a turn longer than one page', () => { * it, and applying it would both duplicate what the temp rows already show and sweep * away the only record of what was never read. */ - it('leaves the transcript alone when it could not read to the end', async () => { + it('keeps the temp rows a capped read did not reach', async () => { let seq = 0 vi.mocked(FlowConversationsService.listConversationMessages) .mockReset() @@ -234,9 +237,33 @@ describe('reading a turn longer than one page', () => { await (manager as any).pollConversationMessages('a', { removeTempMessages: true }) - // Untouched: neither the rows it managed to read nor the sweep were applied. + // Nothing after the last poll of a turn would finish the read, so the rows it did get + // are dropped rather than left showing the turn's start twice beside the temp row. expect(manager.messages).toEqual([streamed]) }) + + it('keeps what a capped read got when a later tick can finish it', async () => { + let seq = 0 + vi.mocked(FlowConversationsService.listConversationMessages) + .mockReset() + .mockImplementation((async () => { + const batch = assistantRows(seq + 1, 50) + seq += 50 + return batch + }) as any) + const manager = managerWithRows() + manager.selectedConversationId = 'a' + + await (manager as any).pollConversationMessages('a', {}) + + // 20 requests of 50 rows, and the rows are kept, so the next tick resumes from seq 1000 + // rather than reading the conversation from the start again. + expect(manager.messages).toHaveLength(1000) + expect(manager.messages.at(-1)?.content).toBe('row 1000') + await (manager as any).pollConversationMessages('a', {}) + const calls = vi.mocked(FlowConversationsService.listConversationMessages).mock.calls + expect((calls[20][0] as any).afterSeq).toBe(1000) + }) }) /** diff --git a/frontend/src/lib/components/flows/conversations/flowChatViewHost.svelte.ts b/frontend/src/lib/components/flows/conversations/flowChatViewHost.svelte.ts index 83900596fc..dd97b74ab9 100644 --- a/frontend/src/lib/components/flows/conversations/flowChatViewHost.svelte.ts +++ b/frontend/src/lib/components/flows/conversations/flowChatViewHost.svelte.ts @@ -469,10 +469,10 @@ export class FlowChatViewHost implements ChatViewHost { // nothing. What the reader wrote comes back, to the chat it was written in rather // than the one open by now. // - // The uploaded objects are deliberately left in place. Deleting them is itself a - // request that can fail, on a path that is already failing, and a resend uploads - // its own under a fresh prefix — so a lost send costs one prefix, not a growing - // number. The workspace's own storage retention is what collects them. + // The uploaded objects are deliberately left in place: deleting them is itself a + // request that can fail, on a path that is already failing. A resend uploads its + // own under a fresh prefix, so each failed send leaves one behind, and the + // workspace's own storage retention is what collects them. this.#restoreToComposer({ ...options, conversationId }) return false }