diff --git a/lib/replicator.js b/lib/replicator.js index 53d619c1..98ffeebc 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -2435,6 +2435,7 @@ module.exports = class Replicator { _requestDone(id, roundtrip) { this._inflight.remove(id, roundtrip) + if (this._inflight.idle) this.queueUpdateAll() if (this.isDownloading() === true) return for (const peer of this.peers) peer.signalUpgrade() } @@ -2604,18 +2605,17 @@ 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 remaining = Math.min(MAX_RANGES, this._ranges.length) - for (let i = 0; i < len; i++) { - const r = this._ranges[i] + const ranges = new RandomIterator(this._ranges) + for (const r of ranges) { + if (remaining-- <= 0) break clampRange(this.core, r) 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 (this._ranges.length >= MAX_RANGES) updateAll = true } } diff --git a/test/replicate.js b/test/replicate.js index bf9f874e..ec13c119 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -11,6 +11,7 @@ const { } = require('./helpers') const { makeStreamPair } = require('./helpers/networking.js') const crypto = require('hypercore-crypto') +const { once } = require('events') const Hypercore = require('../') const DEBUG = false @@ -2891,6 +2892,165 @@ test('delayed updateAll timer doesnt keep event loop alive', async function (t) t.is(r._updateAllBump, null, 'timer reset to null') }) +test('processing ranges doesnt block seeks', async function (t) { + const core = await create(t) + const clone = await create(t, core.key) + + // Well above MAX_RANGES so draining the backlog takes many batches: 50k/64 + const totalLength = 50_000 + + await core.append(new Array(totalLength).fill('a')) + + // Make all range requests stay in `_ranges` forever so `_updateNonPrimary` has an oversized backlog to scan. + await core.clear(0, totalLength) + + replicate(core, clone, t) + + // Populate ranges + for (let i = 0; i < totalLength; i++) { + clone + .download({ start: i, end: i + 1 }) + .done() + .catch(noop) + } + + t.is(clone.core.replicator._ranges.length, totalLength, 'large range backlog is pending') + + const bytesBefore = core.byteLength + + // Appending triggers an upgrade on clone, kicking off `_updateNonPrimary` to attempt the backlog. + // Seeking to `bytesBefore` needs that same upgrade (so clone doesnt know about it). + await core.append('end') + + // Baseline (no range backlog) resolves this in ~1ms; 900ms+ with the backlog + const BUDGET = 100 + + const start = Date.now() + let delta = -1 + const seek = clone.seek(bytesBefore).then(() => { + delta = Date.now() - start + }) + t.comment('seek requested') + await once(clone, 'append') + + // Too allow merkle IO for seek + await new Promise((resolve) => setImmediate(resolve)) + + t.is(clone.core.replicator._seeks.length, 0, 'seek resolved after upgrade received') + await seek + + t.ok(delta <= BUDGET, `seek resolved within ${BUDGET}ms despite a ${totalLength}-range backlog`) + t.comment('seek time', delta) + t.is(clone.core.replicator._ranges.length, totalLength, 'range backlog is still pending') +}) + +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) + + 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) + } + + 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, 1000, false)) + ]) + + 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.skip('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 })