Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions .changeset/fix-prefetch-queue-stale-abort.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
---
'@electric-sql/client': patch
---

Fix stream stopping after tab visibility changes due to stale aborted requests in PrefetchQueue.

**Root cause:** When a page is hidden, the stream pauses and aborts in-flight prefetch requests. The aborted promises remained in the PrefetchQueue's internal Map. When the page became visible and the stream resumed, `consume()` returned the stale aborted promise, causing an AbortError to propagate to ShapeStream and stop syncing.

**The fix:**
- `PrefetchQueue.consume()` now checks if the request's abort signal is already aborted before returning it
- `PrefetchQueue.abort()` now clears the internal map after aborting controllers
- The fetch wrapper clears `prefetchQueue` after calling `abort()` to ensure fresh requests

Fixes #3460
8 changes: 7 additions & 1 deletion packages/typescript-client/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -724,7 +724,6 @@ export class ShapeStream<T extends Row<unknown> = Row>
async #requestShape(): Promise<void> {
if (this.#state === `pause-requested`) {
this.#state = `paused`

return
}

Expand Down Expand Up @@ -1252,6 +1251,13 @@ export class ShapeStream<T extends Row<unknown> = Row>
this.#started &&
(this.#state === `paused` || this.#state === `pause-requested`)
) {
// Don't resume if the user's signal is already aborted
// This can happen if the signal was aborted while we were paused
// (e.g., TanStack DB collection was GC'd)
if (this.options.signal?.aborted) {
return
}

// If we're resuming from pause-requested state, we need to set state back to active
// to prevent the pause from completing
if (this.#state === `pause-requested`) {
Expand Down
17 changes: 14 additions & 3 deletions packages/typescript-client/src/fetch.ts
Original file line number Diff line number Diff line change
Expand Up @@ -221,7 +221,7 @@ export function createFetchWithChunkBuffer(
): typeof fetch {
const { maxChunksToPrefetch } = prefetchOptions

let prefetchQueue: PrefetchQueue
let prefetchQueue: PrefetchQueue | undefined

const prefetchClient = async (...args: Parameters<typeof fetchClient>) => {
const url = args[0].toString()
Expand All @@ -233,7 +233,10 @@ export function createFetchWithChunkBuffer(
return prefetchedRequest
}

// Clear the prefetch queue after aborting to prevent returning
// stale/aborted requests on future calls with the same URL
prefetchQueue?.abort()
prefetchQueue = undefined

// perform request and fire off prefetch queue if request is eligible
const response = await fetchClient(...args)
Expand Down Expand Up @@ -340,16 +343,24 @@ class PrefetchQueue {

abort(): void {
this.#prefetchQueue.forEach(([_, aborter]) => aborter.abort())
this.#prefetchQueue.clear()
}

consume(...args: Parameters<typeof fetch>): Promise<Response> | void {
const url = args[0].toString()

const request = this.#prefetchQueue.get(url)?.[0]
const entry = this.#prefetchQueue.get(url)
// only consume if request is in queue and is the queue "head"
// if request is in the queue but not the head, the queue is being
// consumed out of order and should be restarted
if (!request || url !== this.#queueHeadUrl) return
if (!entry || url !== this.#queueHeadUrl) return

const [request, aborter] = entry
// Don't return aborted requests - they will reject with AbortError
if (aborter.signal.aborted) {
this.#prefetchQueue.delete(url)
return
}
this.#prefetchQueue.delete(url)

// fire off new prefetch since request has been consumed
Expand Down
72 changes: 72 additions & 0 deletions packages/typescript-client/test/fetch.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -525,6 +525,78 @@ describe(`createFetchWithChunkBuffer`, () => {
// no new prefetches since main request was aborted
expect(mockFetch).toHaveBeenCalledTimes(2)
})

it(`should not return aborted prefetch requests after queue is cleared`, async () => {
// This test verifies the fix for issue #3460:
// When a page is hidden, the stream pauses and aborts in-flight prefetches.
// Previously, the aborted promises stayed in the queue and would be returned
// on resume, causing the stream to incorrectly interpret the abort as fatal.
const fetchWrapper = createFetchWithChunkBuffer(mockFetch, {
maxChunksToPrefetch: 1,
})

const nextUrl = sortUrlParams(`${baseUrl}&handle=123&offset=0`)

// Initial response that triggers prefetch of nextUrl
const initialResponse = new Response(`initial chunk`, {
status: 200,
headers: responseHeaders({
[SHAPE_HANDLE_HEADER]: `123`,
[CHUNK_LAST_OFFSET_HEADER]: `0`,
}),
})

// Fresh response for when we request nextUrl after clearing queue
const freshResponse = new Response(`fresh chunk`, {
status: 200,
headers: responseHeaders({
[SHAPE_HANDLE_HEADER]: `123`,
[CHUNK_LAST_OFFSET_HEADER]: `1`,
}),
})

// Create an outer abort controller that we'll use to abort the stream
// This simulates ShapeStream's #requestAbortController
const outerAbortController = new AbortController()

mockFetch
.mockResolvedValueOnce(initialResponse) // Initial request
.mockImplementationOnce((_url, init) => {
// Return a promise that hangs until the signal is aborted
const signal = init?.signal as AbortSignal | undefined
return new Promise<Response>((_, reject) => {
if (signal) {
signal.addEventListener(`abort`, () => {
reject(new DOMException(`Aborted`, `AbortError`))
})
}
})
})
.mockResolvedValueOnce(freshResponse) // Fresh request after queue cleared

// Make initial request with abort signal - this triggers prefetch of nextUrl
// The prefetch inherits the signal through chainAborter
await fetchWrapper(baseUrl, { signal: outerAbortController.signal })
await sleep() // Let prefetch start

expect(mockFetch).toHaveBeenCalledTimes(2)
expect(mockFetch).toHaveBeenNthCalledWith(2, nextUrl, expect.anything())

// Simulate page hidden: abort the stream's request controller
// This propagates through chainAborter to abort the prefetch's internal aborter
outerAbortController.abort()
await sleep()

// Simulate page visible: request the same URL that was being prefetched
// Before fix: this would return the aborted promise and throw
// After fix: consume() checks aborter.signal.aborted and makes a fresh request
const result = await fetchWrapper(nextUrl)

expect(result).toBe(freshResponse)
// 4 calls: initial, prefetch (aborted), fresh request, new prefetch triggered by fresh response
expect(mockFetch).toHaveBeenCalledTimes(4)
expect(mockFetch).toHaveBeenNthCalledWith(3, nextUrl)
})
})

describe(`createFetchWithConsumedMessages`, () => {
Expand Down
Loading