Skip to content
Open
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
12 changes: 6 additions & 6 deletions lib/replicator.js
Original file line number Diff line number Diff line change
Expand Up @@ -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()
}
Expand Down Expand Up @@ -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
}
}

Expand Down
160 changes: 160 additions & 0 deletions test/replicate.js
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 })
Expand Down
Loading