From f9700d5aa5f82e7d6683518a623fa79207cc5ac2 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Mon, 22 Jun 2026 22:29:16 +0200 Subject: [PATCH 01/19] Fix idle range completion scan --- lib/replicator.js | 5 ++++- test/replicate.js | 30 ++++++++++++++++++++++++++++++ 2 files changed, 34 insertions(+), 1 deletion(-) diff --git a/lib/replicator.js b/lib/replicator.js index 53d619c1..4236b024 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -2604,7 +2604,10 @@ module.exports = class Replicator { async _updateNonPrimary(updateAll) { // Check if running, if so skip it and the running one will issue another update for us (debounce) while (++this._updatesPending === 1) { - let len = Math.min(MAX_RANGES, this._ranges.length) + let len = + this._inflight.idle || updateAll + ? this._ranges.length + : Math.min(MAX_RANGES, this._ranges.length) for (let i = 0; i < len; i++) { const r = this._ranges[i] diff --git a/test/replicate.js b/test/replicate.js index bf9f874e..bc8df175 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -2891,6 +2891,36 @@ test('delayed updateAll timer doesnt keep event loop alive', async function (t) t.is(r._updateAllBump, null, 'timer reset to null') }) +test('idle range completion scans past max range window', async function (t) { + const core = await create(t) + + const pending = new Set() + const complete = [] + + for (let i = 0; i < 650; i++) { + pending.add(i) + + const done = core.download({ start: i, end: i + 1 }).done() + done.then(() => pending.delete(i), noop) + + if (i >= 300) complete.push(done) + } + + t.is(core.core.replicator._ranges.length, 650, 'all ranges are pending') + + core.core._setBitfieldRanges(300, 650, true) + await core.core.replicator._updateNonPrimary(false) + + const resolved = await Promise.race([ + Promise.all(complete).then(() => true), + new Promise((resolve) => setTimeout(resolve, 250, false)) + ]) + + t.ok(resolved, 'all locally complete ranges resolved') + t.is(pending.size, 300, 'only unavailable ranges remain pending') + t.is(core.core.replicator._ranges.length, 300, 'resolved ranges were removed') +}) + async function createAndDownload(t, core) { const b = await create(t, core.key) replicate(core, b, t, { teardown: false }) From 1414e1633bd2e84ab4125cd78d78fdd71f09db96 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Tue, 23 Jun 2026 13:57:28 +0200 Subject: [PATCH 02/19] Fix stuck range completion drain --- lib/replicator.js | 67 ++++++++++++++++++++++++++++++++++++++--------- test/replicate.js | 43 ++++++++++++++++++++++-------- 2 files changed, 86 insertions(+), 24 deletions(-) diff --git a/lib/replicator.js b/lib/replicator.js index 4236b024..c37d4ffd 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -2600,26 +2600,63 @@ module.exports = class Replicator { } } + _updateRanges(index, limit) { + if (this._ranges.length === 0) { + return { checked: 0, resolved: 0, hitMaxAfterResolve: false, index: 0 } + } + + index = Math.min(index, this._ranges.length - 1) + + let checked = 0 + let resolved = 0 + let hitMaxAfterResolve = false + let remaining = Math.min(limit, this._ranges.length) + + while (remaining-- > 0 && this._ranges.length > 0) { + if (index >= this._ranges.length) index = 0 + + const r = this._ranges[index] + + clampRange(this.core, r) + checked++ + + if (r.end !== -1 && r.start >= r.end) { + this._resolveRangeRequest(r) + resolved++ + if (this._ranges.length === MAX_RANGES) hitMaxAfterResolve = true + } else { + index++ + } + } + + index = this._ranges.length === 0 ? 0 : index % this._ranges.length + + return { checked, resolved, hitMaxAfterResolve, index } + } + // "slow" updates here - async but not allowed to ever throw async _updateNonPrimary(updateAll) { // Check if running, if so skip it and the running one will issue another update for us (debounce) while (++this._updatesPending === 1) { - let len = - this._inflight.idle || updateAll - ? this._ranges.length - : Math.min(MAX_RANGES, this._ranges.length) + let checkedSinceResolve = 0 + let rangeIndex = 0 - for (let i = 0; i < len; i++) { - const r = this._ranges[i] + while (this._ranges.length > 0) { + const drain = this._inflight.idle || updateAll + const limit = Math.min(MAX_RANGES, this._ranges.length - checkedSinceResolve) + const result = this._updateRanges(rangeIndex, limit) + const { checked, resolved, hitMaxAfterResolve } = result + rangeIndex = result.index - clampRange(this.core, r) + if (hitMaxAfterResolve) updateAll = true + if (resolved > 0) checkedSinceResolve = 0 + else checkedSinceResolve += checked - if (r.end !== -1 && r.start >= r.end) { - this._resolveRangeRequest(r) - i-- - if (len > this._ranges.length) len-- - if (this._ranges.length === MAX_RANGES) updateAll = true - } + if (checked === 0 || checkedSinceResolve >= this._ranges.length) break + if (!drain) break + + await yieldToLoop() + if (this.destroyed || this._updatesPending > 1) break } for (let i = 0; i < this._seeks.length; i++) { @@ -3387,6 +3424,10 @@ function incrementRx(stats1, stats2) { function noop() {} +function yieldToLoop() { + return new Promise((resolve) => setTimeout(resolve, 0)) +} + function backoff(times) { const sleep = times < 2 ? 200 : times < 5 ? 500 : times < 40 ? 1000 : 5000 return new Promise((resolve) => setTimeout(resolve, sleep)) diff --git a/test/replicate.js b/test/replicate.js index bc8df175..901c9328 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -2891,34 +2891,55 @@ test('delayed updateAll timer doesnt keep event loop alive', async function (t) t.is(r._updateAllBump, null, 'timer reset to null') }) -test('idle range completion scans past max range window', async function (t) { +test('idle range completion drains past max range window', async function (t) { const core = await create(t) + const clone = await create(t, core.key) const pending = new Set() const complete = [] + const totalLength = 150 + const availableStart = 100 + + for (let i = 0; i < totalLength; i++) { + await core.append('' + i) + } + + await core.clear(0, availableStart) + + const random = Math.random + Math.random = () => 0.999 + t.teardown(() => { + Math.random = random + }) + + replicate(core, clone, t) - for (let i = 0; i < 650; i++) { + for (let i = 0; i < totalLength; i++) { pending.add(i) - const done = core.download({ start: i, end: i + 1 }).done() + const done = clone.download({ start: i, end: i + 1 }).done() done.then(() => pending.delete(i), noop) - if (i >= 300) complete.push(done) + if (i >= availableStart) complete.push(done) } - t.is(core.core.replicator._ranges.length, 650, 'all ranges are pending') - - core.core._setBitfieldRanges(300, 650, true) - await core.core.replicator._updateNonPrimary(false) + t.is(clone.core.replicator._ranges.length, totalLength, 'all ranges are pending') + t.ok(await core.has(availableStart, totalLength), 'source has all available blocks') const resolved = await Promise.race([ Promise.all(complete).then(() => true), - new Promise((resolve) => setTimeout(resolve, 250, false)) + new Promise((resolve) => setTimeout(resolve, 1000, false)) ]) t.ok(resolved, 'all locally complete ranges resolved') - t.is(pending.size, 300, 'only unavailable ranges remain pending') - t.is(core.core.replicator._ranges.length, 300, 'resolved ranges were removed') + t.is(pending.size, availableStart, 'only unavailable ranges remain pending') + t.is(clone.core.replicator._ranges.length, availableStart, 'resolved ranges were removed') + + const outliers = [] + for (const range of clone.core.replicator._ranges) { + if (range.userStart >= availableStart) outliers.push(range.userStart) + } + t.alike(outliers, [], 'no available ranges remain pending') }) async function createAndDownload(t, core) { From 76376ac22eda082543407c38e1c18ae24b5e2bbc Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Tue, 23 Jun 2026 14:52:16 +0200 Subject: [PATCH 03/19] Remove unnecessary range update boundary state --- lib/replicator.js | 9 +++------ 1 file changed, 3 insertions(+), 6 deletions(-) diff --git a/lib/replicator.js b/lib/replicator.js index c37d4ffd..aa860e46 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -2602,14 +2602,13 @@ module.exports = class Replicator { _updateRanges(index, limit) { if (this._ranges.length === 0) { - return { checked: 0, resolved: 0, hitMaxAfterResolve: false, index: 0 } + return { checked: 0, resolved: 0, index: 0 } } index = Math.min(index, this._ranges.length - 1) let checked = 0 let resolved = 0 - let hitMaxAfterResolve = false let remaining = Math.min(limit, this._ranges.length) while (remaining-- > 0 && this._ranges.length > 0) { @@ -2623,7 +2622,6 @@ module.exports = class Replicator { if (r.end !== -1 && r.start >= r.end) { this._resolveRangeRequest(r) resolved++ - if (this._ranges.length === MAX_RANGES) hitMaxAfterResolve = true } else { index++ } @@ -2631,7 +2629,7 @@ module.exports = class Replicator { index = this._ranges.length === 0 ? 0 : index % this._ranges.length - return { checked, resolved, hitMaxAfterResolve, index } + return { checked, resolved, index } } // "slow" updates here - async but not allowed to ever throw @@ -2645,10 +2643,9 @@ module.exports = class Replicator { const drain = this._inflight.idle || updateAll const limit = Math.min(MAX_RANGES, this._ranges.length - checkedSinceResolve) const result = this._updateRanges(rangeIndex, limit) - const { checked, resolved, hitMaxAfterResolve } = result + const { checked, resolved } = result rangeIndex = result.index - if (hitMaxAfterResolve) updateAll = true if (resolved > 0) checkedSinceResolve = 0 else checkedSinceResolve += checked From 60823b1284cb22e03334a255d03d3ff8aff75af8 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Tue, 23 Jun 2026 15:16:52 +0200 Subject: [PATCH 04/19] Fix queued idle range drain race --- lib/replicator.js | 4 +-- test/replicate.js | 62 +++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 64 insertions(+), 2 deletions(-) diff --git a/lib/replicator.js b/lib/replicator.js index aa860e46..907a1009 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -2638,9 +2638,9 @@ module.exports = class Replicator { while (++this._updatesPending === 1) { let checkedSinceResolve = 0 let rangeIndex = 0 + const drain = this._inflight.idle || updateAll while (this._ranges.length > 0) { - const drain = this._inflight.idle || updateAll const limit = Math.min(MAX_RANGES, this._ranges.length - checkedSinceResolve) const result = this._updateRanges(rangeIndex, limit) const { checked, resolved } = result @@ -2653,7 +2653,7 @@ module.exports = class Replicator { if (!drain) break await yieldToLoop() - if (this.destroyed || this._updatesPending > 1) break + if (this.destroyed) break } for (let i = 0; i < this._seeks.length; i++) { diff --git a/test/replicate.js b/test/replicate.js index 901c9328..4baefc48 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -2942,6 +2942,68 @@ test('idle range completion drains past max range window', async function (t) { t.alike(outliers, [], 'no available ranges remain pending') }) +test('idle range completion keeps draining if update queues during yield', async function (t) { + const core = await create(t) + const pending = new Set() + const complete = [] + + const totalLength = 150 + const availableStart = 100 + const replicator = core.core.replicator + + for (let i = 0; i < totalLength; i++) { + pending.add(i) + + const done = core.download({ start: i, end: i + 1 }).done() + done.then(() => pending.delete(i), noop) + + if (i >= availableStart) complete.push(done) + } + + t.is(replicator._ranges.length, totalLength, 'all ranges are pending') + + core.core._setBitfieldRanges(availableStart, totalLength, true) + + const updateRanges = replicator._updateRanges + let updates = 0 + + replicator._updateRanges = function (index, limit) { + const result = updateRanges.call(this, index, limit) + + if (++updates === 1) { + setTimeout(() => { + this._inflight._active++ + this._updateNonPrimary(false).catch(noop) + }, 0) + } + + return result + } + + t.teardown(() => { + replicator._updateRanges = updateRanges + replicator._inflight._active = 0 + }) + + await replicator._updateNonPrimary(false) + replicator._inflight._active = 0 + + const resolved = await Promise.race([ + Promise.all(complete).then(() => true), + new Promise((resolve) => setTimeout(resolve, 1000, false)) + ]) + + t.ok(resolved, 'all locally complete ranges resolved') + t.is(pending.size, availableStart, 'only unavailable ranges remain pending') + t.is(replicator._ranges.length, availableStart, 'resolved ranges were removed') + + const outliers = [] + for (const range of replicator._ranges) { + if (range.userStart >= availableStart) outliers.push(range.userStart) + } + t.alike(outliers, [], 'no available ranges remain pending') +}) + async function createAndDownload(t, core) { const b = await create(t, core.key) replicate(core, b, t, { teardown: false }) From 5403eb5e7480440655351362b5fbc33e102d04f8 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Mon, 13 Jul 2026 22:21:22 +0200 Subject: [PATCH 05/19] Restore peer updates at max range boundary --- lib/replicator.js | 9 +++++-- test/replicate.js | 68 +++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 75 insertions(+), 2 deletions(-) diff --git a/lib/replicator.js b/lib/replicator.js index 907a1009..ee03a3b4 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -2602,13 +2602,14 @@ module.exports = class Replicator { _updateRanges(index, limit) { if (this._ranges.length === 0) { - return { checked: 0, resolved: 0, index: 0 } + return { checked: 0, resolved: 0, index: 0, updateAll: false } } index = Math.min(index, this._ranges.length - 1) let checked = 0 let resolved = 0 + let updateAll = false let remaining = Math.min(limit, this._ranges.length) while (remaining-- > 0 && this._ranges.length > 0) { @@ -2622,6 +2623,8 @@ module.exports = class Replicator { if (r.end !== -1 && r.start >= r.end) { this._resolveRangeRequest(r) resolved++ + // Crossing into the capped window exposes another range to peer scheduling. + if (this._ranges.length === MAX_RANGES) updateAll = true } else { index++ } @@ -2629,7 +2632,7 @@ module.exports = class Replicator { index = this._ranges.length === 0 ? 0 : index % this._ranges.length - return { checked, resolved, index } + return { checked, resolved, index, updateAll } } // "slow" updates here - async but not allowed to ever throw @@ -2646,6 +2649,8 @@ module.exports = class Replicator { const { checked, resolved } = result rangeIndex = result.index + if (result.updateAll) updateAll = true + if (resolved > 0) checkedSinceResolve = 0 else checkedSinceResolve += checked diff --git a/test/replicate.js b/test/replicate.js index 4baefc48..d2848201 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -2891,6 +2891,74 @@ test('delayed updateAll timer doesnt keep event loop alive', async function (t) t.is(r._updateAllBump, null, 'timer reset to null') }) +test('completed range exposes capped range to existing peer', async function (t) { + const writer = await create(t) + const totalLength = 65 + + await writer.append(new Array(totalLength).fill('a')) + + const firstPeer = await create(t, writer.key) + let streams = replicate(writer, firstPeer, t, { teardown: false }) + await firstPeer.get(0) + await unreplicate(streams) + + const cappedPeer = await create(t, writer.key) + streams = replicate(writer, cappedPeer, t, { teardown: false }) + await cappedPeer.get(totalLength - 1) + await unreplicate(streams) + + const downloader = await create(t, writer.key) + + t.absent(await firstPeer.has(totalLength - 1), 'first peer does not have capped block') + + const random = Math.random + Math.random = () => 0.999 + t.teardown(() => { + Math.random = random + }) + + const cappedDone = downloader + .download({ start: totalLength - 1, end: totalLength }) + .done() + .then(() => true, noop) + const firstDone = downloader.download({ start: 0, end: 1 }).done() + + for (let i = 1; i < totalLength - 1; i++) { + downloader.download({ start: i, end: i + 1 }) + } + + const replicator = downloader.core.replicator + t.is(replicator._ranges.length, totalLength, 'one range is outside the scan window') + + // Model unrelated inflight work so the idle fallback cannot trigger the peer update. + replicator._inflight._active++ + t.teardown(() => { + replicator._inflight._active-- + }) + + const cappedPeerAdded = new Promise((resolve) => downloader.once('peer-add', resolve)) + replicate(cappedPeer, downloader, t) + + const remoteCappedPeer = await cappedPeerAdded + while (!remoteCappedPeer.remoteBitfield.get(totalLength - 1)) await eventFlush() + + t.absent(await downloader.has(totalLength - 1), 'capped range was skipped initially') + + replicate(firstPeer, downloader, t) + await firstDone + + let timer = null + const downloaded = await Promise.race([ + cappedDone, + new Promise((resolve) => { + timer = setTimeout(resolve, 1000, false) + }) + ]) + clearTimeout(timer) + + t.ok(downloaded, 'existing peer is revisited for the newly exposed range') +}) + test('idle range completion drains past max range window', async function (t) { const core = await create(t) const clone = await create(t, core.key) From 37224a99e69f458658a8dbd1164e5cc9f677f204 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Tue, 14 Jul 2026 14:21:13 +0200 Subject: [PATCH 06/19] Yield range drains to pending seeks --- lib/replicator.js | 92 ++++++++++++++++++++++++++++------------------- test/replicate.js | 60 +++++++++++++++++++++++++++++++ 2 files changed, 116 insertions(+), 36 deletions(-) diff --git a/lib/replicator.js b/lib/replicator.js index ee03a3b4..40424f35 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -2015,6 +2015,8 @@ module.exports = class Replicator { this._active = 0 this._ifAvailable = 0 this._updatesPending = 0 + this._updatesRunning = false + this._updatesQueued = false this._applyingReorg = null this._manifestPeer = null this._notDownloadingLinger = notDownloadingLinger @@ -2637,60 +2639,78 @@ module.exports = class Replicator { // "slow" updates here - async but not allowed to ever throw async _updateNonPrimary(updateAll) { - // Check if running, if so skip it and the running one will issue another update for us (debounce) - while (++this._updatesPending === 1) { + this._updatesQueued = true + if (this._updatesRunning) { + if (this._inflight.idle || updateAll) this.queueUpdateAll() + return + } + + this._updatesRunning = true + + while (this._updatesQueued) { + this._updatesQueued = false + let checkedSinceResolve = 0 let rangeIndex = 0 const drain = this._inflight.idle || updateAll - while (this._ranges.length > 0) { - const limit = Math.min(MAX_RANGES, this._ranges.length - checkedSinceResolve) - const result = this._updateRanges(rangeIndex, limit) - const { checked, resolved } = result - rangeIndex = result.index + while (true) { + this._updatesPending = 1 - if (result.updateAll) updateAll = true + let continueRanges = false + if (this._ranges.length > 0) { + const limit = Math.min(MAX_RANGES, this._ranges.length - checkedSinceResolve) + const result = this._updateRanges(rangeIndex, limit) + const { checked, resolved } = result + rangeIndex = result.index - if (resolved > 0) checkedSinceResolve = 0 - else checkedSinceResolve += checked + if (result.updateAll) updateAll = true - if (checked === 0 || checkedSinceResolve >= this._ranges.length) break - if (!drain) break + if (resolved > 0) checkedSinceResolve = 0 + else checkedSinceResolve += checked - await yieldToLoop() - if (this.destroyed) break - } + continueRanges = checked > 0 && checkedSinceResolve < this._ranges.length && drain + } - for (let i = 0; i < this._seeks.length; i++) { - const s = this._seeks[i] + for (let i = 0; i < this._seeks.length; i++) { + const s = this._seeks[i] - let err = null - let res = null + let err = null + let res = null - try { - res = await s.seeker.update() - } catch (error) { - err = error - } + try { + res = await s.seeker.update() + } catch (error) { + err = error + } + + if (!res && !err) continue - if (!res && !err) continue + if (i < this._seeks.length - 1) this._seeks[i] = this._seeks.pop() + else this._seeks.pop() - if (i < this._seeks.length - 1) this._seeks[i] = this._seeks.pop() - else this._seeks.pop() + i-- - i-- + if (err) s.reject(err) + else s.resolve(res) + } - if (err) s.reject(err) - else s.resolve(res) + this._updatesPending = 0 + + if ((!continueRanges || this._seeks.length > 0) && (this._inflight.idle || updateAll)) { + this.queueUpdateAll() + } + if (!continueRanges) break + + await yieldToLoop() + if (this.destroyed) break } - // No additional updates scheduled - break - if (--this._updatesPending === 0) break - // Debounce the additional updates - continue - this._updatesPending = 0 + if (this.destroyed) break } - if (this._inflight.idle || updateAll) this.queueUpdateAll() + this._updatesPending = 0 + this._updatesRunning = false } _clearRequest(peer, req) { @@ -3427,7 +3447,7 @@ function incrementRx(stats1, stats2) { function noop() {} function yieldToLoop() { - return new Promise((resolve) => setTimeout(resolve, 0)) + return new Promise((resolve) => setImmediate(resolve)) } function backoff(times) { diff --git a/test/replicate.js b/test/replicate.js index d2848201..66203ace 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -2959,6 +2959,66 @@ test('completed range exposes capped range to existing peer', async function (t) t.ok(downloaded, 'existing peer is revisited for the newly exposed range') }) +test('processing ranges does not block seeks', async function (t) { + const core = await create(t) + const totalLength = 129 + + await core.append(new Array(totalLength).fill('a')) + + const clone = await create(t, core.key, { eagerUpgrade: false }) + const peerAdded = new Promise((resolve) => clone.once('peer-add', resolve)) + replicate(core, clone, t) + + const peer = await peerAdded + peer.paused = true + while (!peer.remoteSynced) await eventFlush() + + for (let i = 0; i < totalLength; i++) { + clone + .download({ start: i, end: i + 1 }) + .done() + .catch(noop) + } + + const replicator = clone.core.replicator + t.is(replicator._ranges.length, totalLength, 'range backlog spans multiple chunks') + + const updateRanges = replicator._updateRanges + const requestSeek = peer._requestSeek + let updates = 0 + let requestedAt = 0 + let seek = null + + replicator._updateRanges = function (index, limit) { + const result = updateRanges.call(this, index, limit) + + if (++updates === 1) { + setImmediate(() => { + peer.paused = false + seek = clone.seek(Math.floor(core.byteLength / 2)) + }) + } + + return result + } + + peer._requestSeek = function (s) { + const sent = requestSeek.call(this, s) + if (sent && requestedAt === 0) requestedAt = updates + return sent + } + + t.teardown(() => { + replicator._updateRanges = updateRanges + peer._requestSeek = requestSeek + }) + + await replicator._updateNonPrimary(true) + await seek + + t.is(requestedAt, 1, 'seek was requested during the first range yield') +}) + test('idle range completion drains past max range window', async function (t) { const core = await create(t) const clone = await create(t, core.key) From ffd7033a69c63ad7c6bd994393d3c2ef0063e914 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Tue, 14 Jul 2026 15:08:38 +0200 Subject: [PATCH 07/19] Assert seeks precede range scan completion --- test/replicate.js | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/test/replicate.js b/test/replicate.js index 66203ace..cbca2edf 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -2986,11 +2986,13 @@ test('processing ranges does not block seeks', async function (t) { const updateRanges = replicator._updateRanges const requestSeek = peer._requestSeek let updates = 0 - let requestedAt = 0 + let checked = 0 + let checkedAtRequest = totalLength let seek = null replicator._updateRanges = function (index, limit) { const result = updateRanges.call(this, index, limit) + checked += result.checked if (++updates === 1) { setImmediate(() => { @@ -3004,7 +3006,7 @@ test('processing ranges does not block seeks', async function (t) { peer._requestSeek = function (s) { const sent = requestSeek.call(this, s) - if (sent && requestedAt === 0) requestedAt = updates + if (sent && checkedAtRequest === totalLength) checkedAtRequest = checked return sent } @@ -3016,7 +3018,7 @@ test('processing ranges does not block seeks', async function (t) { await replicator._updateNonPrimary(true) await seek - t.is(requestedAt, 1, 'seek was requested during the first range yield') + t.ok(checkedAtRequest < totalLength, 'seek was requested before the range scan completed') }) test('idle range completion drains past max range window', async function (t) { From fe789a158a86334f89ea5c1eec244e17c480280e Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Tue, 14 Jul 2026 15:39:37 +0200 Subject: [PATCH 08/19] Restart range drain after cancellation --- lib/replicator.js | 4 ++++ test/replicate.js | 50 +++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 54 insertions(+) diff --git a/lib/replicator.js b/lib/replicator.js index 40424f35..c6ac253b 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -285,6 +285,10 @@ class RangeRequest extends Attachable { this.ranges[rangeIndex] = h } + if (!this.resolved && this.replicator._updatesRunning) { + this.replicator._updatesQueued = true + } + if (this.end === -1) { this.replicator._alwaysLatestBlock-- } diff --git a/test/replicate.js b/test/replicate.js index cbca2edf..154d1928 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -3134,6 +3134,56 @@ test('idle range completion keeps draining if update queues during yield', async t.alike(outliers, [], 'no available ranges remain pending') }) +test('idle range completion restarts if ranges cancel during yield', async function (t) { + const core = await create(t) + const downloads = [] + + const totalLength = 100 + const availableStart = 80 + const cancelLength = totalLength - availableStart + let completed = 0 + + for (let i = 0; i < totalLength; i++) { + const download = core.download({ start: i, end: i + 1 }) + const done = download.done() + + if (i >= availableStart) done.then(() => completed++, noop) + else done.catch(noop) + + if (i < cancelLength) downloads.push(download) + } + + const replicator = core.core.replicator + t.is(replicator._ranges.length, totalLength, 'range backlog spans multiple chunks') + + core.core._setBitfieldRanges(availableStart, totalLength, true) + + const updateRanges = replicator._updateRanges + let updates = 0 + + replicator._updateRanges = function (index, limit) { + const result = updateRanges.call(this, index, limit) + + if (++updates === 1) { + setImmediate(() => { + // Swap complete ranges behind the saved cursor while the drain is suspended. + for (const download of downloads) download.destroy() + }) + } + + return result + } + + t.teardown(() => { + replicator._updateRanges = updateRanges + }) + + await replicator._updateNonPrimary(false) + await eventFlush() + + t.is(completed, cancelLength, 'all complete ranges resolved') +}) + async function createAndDownload(t, core) { const b = await create(t, core.key) replicate(core, b, t, { teardown: false }) From 62096edfc509e1da2c75e012d9175eb059108949 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Wed, 15 Jul 2026 16:53:23 +0200 Subject: [PATCH 09/19] Simplify range drain state and tests --- lib/replicator.js | 14 ++++++------- test/replicate.js | 50 +++++++++-------------------------------------- 2 files changed, 16 insertions(+), 48 deletions(-) diff --git a/lib/replicator.js b/lib/replicator.js index c6ac253b..db3cd816 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -1545,8 +1545,8 @@ class Peer { } _requestSeek(s) { - // if replicator is updating the seeks etc, bail and wait for it to drain - if (this.replicator._updatesPending > 0) return false + // if replicator is updating the seeks, bail and wait for it to drain + if (this.replicator._updatingSeeks) return false if (this.replicator.pushOnly) return false const { length, fork } = this.core.state @@ -2018,7 +2018,7 @@ module.exports = class Replicator { this._hadPeers = false this._active = 0 this._ifAvailable = 0 - this._updatesPending = 0 + this._updatingSeeks = false this._updatesRunning = false this._updatesQueued = false this._applyingReorg = null @@ -2659,8 +2659,6 @@ module.exports = class Replicator { const drain = this._inflight.idle || updateAll while (true) { - this._updatesPending = 1 - let continueRanges = false if (this._ranges.length > 0) { const limit = Math.min(MAX_RANGES, this._ranges.length - checkedSinceResolve) @@ -2676,6 +2674,8 @@ module.exports = class Replicator { continueRanges = checked > 0 && checkedSinceResolve < this._ranges.length && drain } + this._updatingSeeks = true + for (let i = 0; i < this._seeks.length; i++) { const s = this._seeks[i] @@ -2699,7 +2699,7 @@ module.exports = class Replicator { else s.resolve(res) } - this._updatesPending = 0 + this._updatingSeeks = false if ((!continueRanges || this._seeks.length > 0) && (this._inflight.idle || updateAll)) { this.queueUpdateAll() @@ -2713,7 +2713,7 @@ module.exports = class Replicator { if (this.destroyed) break } - this._updatesPending = 0 + this._updatingSeeks = false this._updatesRunning = false } diff --git a/test/replicate.js b/test/replicate.js index 154d1928..251e8534 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -3025,15 +3025,11 @@ test('idle range completion drains past max range window', async function (t) { const core = await create(t) const clone = await create(t, core.key) - const pending = new Set() const complete = [] const totalLength = 150 const availableStart = 100 - for (let i = 0; i < totalLength; i++) { - await core.append('' + i) - } - + await core.append(new Array(totalLength).fill('a')) await core.clear(0, availableStart) const random = Math.random @@ -3045,16 +3041,12 @@ test('idle range completion drains past max range window', async function (t) { replicate(core, clone, t) for (let i = 0; i < totalLength; i++) { - pending.add(i) - const done = clone.download({ start: i, end: i + 1 }).done() - done.then(() => pending.delete(i), noop) - if (i >= availableStart) complete.push(done) + else done.catch(noop) } t.is(clone.core.replicator._ranges.length, totalLength, 'all ranges are pending') - t.ok(await core.has(availableStart, totalLength), 'source has all available blocks') const resolved = await Promise.race([ Promise.all(complete).then(() => true), @@ -3062,32 +3054,20 @@ test('idle range completion drains past max range window', async function (t) { ]) t.ok(resolved, 'all locally complete ranges resolved') - t.is(pending.size, availableStart, 'only unavailable ranges remain pending') - t.is(clone.core.replicator._ranges.length, availableStart, 'resolved ranges were removed') - - const outliers = [] - for (const range of clone.core.replicator._ranges) { - if (range.userStart >= availableStart) outliers.push(range.userStart) - } - t.alike(outliers, [], 'no available ranges remain pending') }) test('idle range completion keeps draining if update queues during yield', async function (t) { const core = await create(t) - const pending = new Set() - const complete = [] const totalLength = 150 const availableStart = 100 const replicator = core.core.replicator + let completed = 0 for (let i = 0; i < totalLength; i++) { - pending.add(i) - const done = core.download({ start: i, end: i + 1 }).done() - done.then(() => pending.delete(i), noop) - - if (i >= availableStart) complete.push(done) + if (i >= availableStart) done.then(() => completed++, noop) + else done.catch(noop) } t.is(replicator._ranges.length, totalLength, 'all ranges are pending') @@ -3101,10 +3081,10 @@ test('idle range completion keeps draining if update queues during yield', async const result = updateRanges.call(this, index, limit) if (++updates === 1) { - setTimeout(() => { + setImmediate(() => { this._inflight._active++ this._updateNonPrimary(false).catch(noop) - }, 0) + }) } return result @@ -3117,21 +3097,9 @@ test('idle range completion keeps draining if update queues during yield', async await replicator._updateNonPrimary(false) replicator._inflight._active = 0 + await eventFlush() - const resolved = await Promise.race([ - Promise.all(complete).then(() => true), - new Promise((resolve) => setTimeout(resolve, 1000, false)) - ]) - - t.ok(resolved, 'all locally complete ranges resolved') - t.is(pending.size, availableStart, 'only unavailable ranges remain pending') - t.is(replicator._ranges.length, availableStart, 'resolved ranges were removed') - - const outliers = [] - for (const range of replicator._ranges) { - if (range.userStart >= availableStart) outliers.push(range.userStart) - } - t.alike(outliers, [], 'no available ranges remain pending') + t.is(completed, totalLength - availableStart, 'all locally complete ranges resolved') }) test('idle range completion restarts if ranges cancel during yield', async function (t) { From aded7ec47a1fe00ff1402730f21a6d51575306df Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Wed, 15 Jul 2026 17:57:12 +0200 Subject: [PATCH 10/19] Remove redundant range empty checks --- lib/replicator.js | 23 ++++++++--------------- 1 file changed, 8 insertions(+), 15 deletions(-) diff --git a/lib/replicator.js b/lib/replicator.js index db3cd816..081afda2 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -2607,10 +2607,6 @@ module.exports = class Replicator { } _updateRanges(index, limit) { - if (this._ranges.length === 0) { - return { checked: 0, resolved: 0, index: 0, updateAll: false } - } - index = Math.min(index, this._ranges.length - 1) let checked = 0 @@ -2659,20 +2655,17 @@ module.exports = class Replicator { const drain = this._inflight.idle || updateAll while (true) { - let continueRanges = false - if (this._ranges.length > 0) { - const limit = Math.min(MAX_RANGES, this._ranges.length - checkedSinceResolve) - const result = this._updateRanges(rangeIndex, limit) - const { checked, resolved } = result - rangeIndex = result.index + const limit = Math.min(MAX_RANGES, this._ranges.length - checkedSinceResolve) + const result = this._updateRanges(rangeIndex, limit) + const { checked, resolved } = result + rangeIndex = result.index - if (result.updateAll) updateAll = true + if (result.updateAll) updateAll = true - if (resolved > 0) checkedSinceResolve = 0 - else checkedSinceResolve += checked + if (resolved > 0) checkedSinceResolve = 0 + else checkedSinceResolve += checked - continueRanges = checked > 0 && checkedSinceResolve < this._ranges.length && drain - } + const continueRanges = checked > 0 && checkedSinceResolve < this._ranges.length && drain this._updatingSeeks = true From 0e24e5aa4fae5d29c1397407f49578707c2d9e21 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Thu, 16 Jul 2026 11:31:04 +0200 Subject: [PATCH 11/19] Fix seek processing during range drains --- lib/replicator.js | 23 +++++++++----- test/replicate.js | 77 +++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 92 insertions(+), 8 deletions(-) diff --git a/lib/replicator.js b/lib/replicator.js index 511cb5f0..3972fb35 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -2667,6 +2667,7 @@ module.exports = class Replicator { let checkedSinceResolve = 0 let rangeIndex = 0 const drain = this._inflight.idle || updateAll + const checkedSeeks = new Set() while (true) { const limit = Math.min(MAX_RANGES, this._ranges.length - checkedSinceResolve) @@ -2683,8 +2684,11 @@ module.exports = class Replicator { this._updatingSeeks = true - for (let i = 0; i < this._seeks.length; i++) { - const s = this._seeks[i] + let checkedSeek = false + for (const s of this._seeks.slice()) { + if (checkedSeeks.has(s) || !this._seeks.includes(s)) continue + checkedSeeks.add(s) + checkedSeek = true let err = null let res = null @@ -2695,12 +2699,12 @@ module.exports = class Replicator { err = error } - if (!res && !err) continue + const seekIndex = this._seeks.indexOf(s) + if (seekIndex === -1 || (!res && !err)) continue - if (i < this._seeks.length - 1) this._seeks[i] = this._seeks.pop() - else this._seeks.pop() - - i-- + if (seekIndex < this._seeks.length - 1) { + this._seeks[seekIndex] = this._seeks.pop() + } else this._seeks.pop() if (err) s.reject(err) else s.resolve(res) @@ -2708,7 +2712,10 @@ module.exports = class Replicator { this._updatingSeeks = false - if ((!continueRanges || this._seeks.length > 0) && (this._inflight.idle || updateAll)) { + if ( + (!continueRanges || (checkedSeek && this._seeks.length > 0)) && + (this._inflight.idle || updateAll) + ) { this.queueUpdateAll() } if (!continueRanges) break diff --git a/test/replicate.js b/test/replicate.js index 5d770eea..5d3ac00f 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -3067,6 +3067,83 @@ test('processing ranges does not block seeks', async function (t) { t.ok(checkedAtRequest < totalLength, 'seek was requested before the range scan completed') }) +test('range drain checks unresolved seeks once', async function (t) { + const core = await create(t) + const totalRanges = 129 + + for (let i = 0; i < totalRanges; i++) { + core.download({ start: i, end: i + 1 }) + } + + const replicator = core.core.replicator + let seekUpdates = 0 + const seek = replicator.addSeek([], { + update() { + seekUpdates++ + return null + } + }) + seek.promise.catch(noop) + + const queueUpdateAll = replicator.queueUpdateAll + let peerUpdates = 0 + replicator.queueUpdateAll = function () { + peerUpdates++ + } + + t.teardown(() => { + replicator.queueUpdateAll = queueUpdateAll + if (seek.context) replicator.cancel(seek) + }) + + await replicator._updateNonPrimary(true) + + t.is(seekUpdates, 1, 'seek was checked once across all range chunks') + t.is(peerUpdates, 2, 'peers were updated before and after the range drain') +}) + +test('cancelling seek during local update preserves other seeks', async function (t) { + const core = await create(t) + const replicator = core.core.replicator + + let release = null + let startedResolve = null + const started = new Promise((resolve) => { + startedResolve = resolve + }) + + const first = replicator.addSeek([], { + update() { + startedResolve() + return new Promise((resolve) => { + release = resolve + }) + } + }) + first.promise.catch(noop) + + const second = replicator.addSeek([], { + update() { + return [1, 0] + } + }) + + t.teardown(() => { + if (first.context) replicator.cancel(first) + if (second.context) replicator.cancel(second) + }) + + const updating = replicator._updateNonPrimary(true) + await started + + replicator.cancel(first) + release([0, 0]) + await updating + + const result = await Promise.race([second.promise, eventFlush()]) + t.alike(result, [1, 0], 'remaining seek resolved') +}) + test('idle range completion drains past max range window', async function (t) { const core = await create(t) const clone = await create(t, core.key) From 51738b820c09a06def55e93dd7a755061accbee3 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Thu, 16 Jul 2026 11:56:00 +0200 Subject: [PATCH 12/19] Clear range test timeout --- test/replicate.js | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/test/replicate.js b/test/replicate.js index 5d3ac00f..2e7a5f40 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -3171,10 +3171,14 @@ test('idle range completion drains past max range window', async function (t) { t.is(clone.core.replicator._ranges.length, totalLength, 'all ranges are pending') + let timer = null const resolved = await Promise.race([ Promise.all(complete).then(() => true), - new Promise((resolve) => setTimeout(resolve, 1000, false)) + new Promise((resolve) => { + timer = setTimeout(resolve, 1000, false) + }) ]) + clearTimeout(timer) t.ok(resolved, 'all locally complete ranges resolved') }) From 176da453c02f39ad093213c8da9d68c8752e6951 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Wed, 22 Jul 2026 10:30:11 +0200 Subject: [PATCH 13/19] Simplify range and seek processing --- lib/replicator.js | 17 ++++++++++------- test/replicate.js | 4 ++-- 2 files changed, 12 insertions(+), 9 deletions(-) diff --git a/lib/replicator.js b/lib/replicator.js index 3972fb35..af3ed1b4 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -2671,11 +2671,15 @@ module.exports = class Replicator { while (true) { const limit = Math.min(MAX_RANGES, this._ranges.length - checkedSinceResolve) - const result = this._updateRanges(rangeIndex, limit) - const { checked, resolved } = result - rangeIndex = result.index + const { + checked, + resolved, + index, + updateAll: resultUpdateAll + } = this._updateRanges(rangeIndex, limit) + rangeIndex = index - if (result.updateAll) updateAll = true + if (resultUpdateAll) updateAll = true if (resolved > 0) checkedSinceResolve = 0 else checkedSinceResolve += checked @@ -2702,9 +2706,8 @@ module.exports = class Replicator { const seekIndex = this._seeks.indexOf(s) if (seekIndex === -1 || (!res && !err)) continue - if (seekIndex < this._seeks.length - 1) { - this._seeks[seekIndex] = this._seeks.pop() - } else this._seeks.pop() + const h = this._seeks.pop() + if (h !== s) this._seeks[seekIndex] = h if (err) s.reject(err) else s.resolve(res) diff --git a/test/replicate.js b/test/replicate.js index 2e7a5f40..0370bd03 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -2944,12 +2944,12 @@ test('completed range exposes capped range to existing peer', async function (t) await writer.append(new Array(totalLength).fill('a')) const firstPeer = await create(t, writer.key) - let streams = replicate(writer, firstPeer, t, { teardown: false }) + let streams = replicate(writer, firstPeer, t) await firstPeer.get(0) await unreplicate(streams) const cappedPeer = await create(t, writer.key) - streams = replicate(writer, cappedPeer, t, { teardown: false }) + streams = replicate(writer, cappedPeer, t) await cappedPeer.get(totalLength - 1) await unreplicate(streams) From 8d7720439fb40e693a93c21c3e366e7c279a27ff Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Wed, 22 Jul 2026 17:00:30 +0200 Subject: [PATCH 14/19] Clarify range and seek scheduling tests --- test/replicate.js | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/test/replicate.js b/test/replicate.js index 0370bd03..3b57d201 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -2937,7 +2937,7 @@ test('delayed updateAll timer doesnt keep event loop alive', async function (t) t.is(r._updateAllBump, null, 'timer reset to null') }) -test('completed range exposes capped range to existing peer', async function (t) { +test('crossing max range boundary reschedules existing peers', async function (t) { const writer = await create(t) const totalLength = 65 @@ -2990,6 +2990,7 @@ test('completed range exposes capped range to existing peer', async function (t) t.absent(await downloader.has(totalLength - 1), 'capped range was skipped initially') + // Completing the block-0 range crosses the 65 -> 64 boundary and must rescan existing peers. replicate(firstPeer, downloader, t) await firstDone @@ -3041,6 +3042,7 @@ test('processing ranges does not block seeks', async function (t) { checked += result.checked if (++updates === 1) { + // Add the seek while the drain is suspended between its first and second range batches. setImmediate(() => { peer.paused = false seek = clone.seek(Math.floor(core.byteLength / 2)) From 64bc2a510422fe55e029de2043857dde253c5051 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Wed, 22 Jul 2026 17:33:00 +0200 Subject: [PATCH 15/19] Simplify range drain seek test --- test/replicate.js | 14 +++----------- 1 file changed, 3 insertions(+), 11 deletions(-) diff --git a/test/replicate.js b/test/replicate.js index 3b57d201..f6f98ee4 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -3032,23 +3032,12 @@ test('processing ranges does not block seeks', async function (t) { const updateRanges = replicator._updateRanges const requestSeek = peer._requestSeek - let updates = 0 let checked = 0 let checkedAtRequest = totalLength - let seek = null replicator._updateRanges = function (index, limit) { const result = updateRanges.call(this, index, limit) checked += result.checked - - if (++updates === 1) { - // Add the seek while the drain is suspended between its first and second range batches. - setImmediate(() => { - peer.paused = false - seek = clone.seek(Math.floor(core.byteLength / 2)) - }) - } - return result } @@ -3063,6 +3052,9 @@ test('processing ranges does not block seeks', async function (t) { peer._requestSeek = requestSeek }) + peer.paused = false + const seek = clone.seek(Math.floor(core.byteLength / 2)) + await replicator._updateNonPrimary(true) await seek From 2407590adf963a661de44c22d9d34dcdbd0ef1e6 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Thu, 23 Jul 2026 15:19:36 +0200 Subject: [PATCH 16/19] Clarify range drain state handling --- lib/replicator.js | 1 - test/replicate.js | 2 ++ 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/lib/replicator.js b/lib/replicator.js index af3ed1b4..908fe148 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -2730,7 +2730,6 @@ module.exports = class Replicator { if (this.destroyed) break } - this._updatingSeeks = false this._updatesRunning = false } diff --git a/test/replicate.js b/test/replicate.js index f6f98ee4..a5eb42a8 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -3203,6 +3203,8 @@ test('idle range completion keeps draining if update queues during yield', async if (++updates === 1) { setImmediate(() => { + // Simulate an update arriving with new inflight work. The nested call only sets + // _updatesQueued; the original idle drain must still reach the complete tail ranges. this._inflight._active++ this._updateNonPrimary(false).catch(noop) }) From f62c8b9a668afcdda81a4173265b1a282aba957d Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Fri, 24 Jul 2026 13:01:32 +0200 Subject: [PATCH 17/19] Remove unnecessary range test catch --- test/replicate.js | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/test/replicate.js b/test/replicate.js index a5eb42a8..6e78846d 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -3206,7 +3206,7 @@ test('idle range completion keeps draining if update queues during yield', async // Simulate an update arriving with new inflight work. The nested call only sets // _updatesQueued; the original idle drain must still reach the complete tail ranges. this._inflight._active++ - this._updateNonPrimary(false).catch(noop) + this._updateNonPrimary(false) }) } From 7e81e3a4b8d20b40f2d99eb0e4b8ec323944202a Mon Sep 17 00:00:00 2001 From: Sean Zellmer Date: Tue, 28 Jul 2026 16:19:53 -0500 Subject: [PATCH 18/19] Assert re-queuing after drain --- test/replicate.js | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/test/replicate.js b/test/replicate.js index 6e78846d..f44d770c 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -3223,6 +3223,11 @@ test('idle range completion keeps draining if update queues during yield', async await eventFlush() t.is(completed, totalLength - availableStart, 'all locally complete ranges resolved') + + const singlePass = Math.ceil(totalLength / 64) + const repassBecauseOfResolving = Math.ceil(availableStart / 64) + const firstPassBatches = singlePass + repassBecauseOfResolving + t.is(updates, firstPassBatches + 1, 'does a second pass because _updateNonPrimary called mid update') }) test('idle range completion restarts if ranges cancel during yield', async function (t) { From 59b21e111ab551759c6212c8910cde76c7c18db7 Mon Sep 17 00:00:00 2001 From: Sean Zellmer Date: Tue, 28 Jul 2026 16:23:37 -0500 Subject: [PATCH 19/19] Lint `test/replicate.js` --- test/replicate.js | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/test/replicate.js b/test/replicate.js index f44d770c..ce4a99a2 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -3227,7 +3227,11 @@ test('idle range completion keeps draining if update queues during yield', async const singlePass = Math.ceil(totalLength / 64) const repassBecauseOfResolving = Math.ceil(availableStart / 64) const firstPassBatches = singlePass + repassBecauseOfResolving - t.is(updates, firstPassBatches + 1, 'does a second pass because _updateNonPrimary called mid update') + t.is( + updates, + firstPassBatches + 1, + 'does a second pass because _updateNonPrimary called mid update' + ) }) test('idle range completion restarts if ranges cancel during yield', async function (t) {