From 1d442e7c798b43053a0a41223729a4d9dc9ac435 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Thu, 14 May 2026 10:34:13 +0200 Subject: [PATCH 01/13] Skip peer update scans on idle close --- lib/replicator.js | 20 +++++++++++++--- test/replicate.js | 61 +++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 78 insertions(+), 3 deletions(-) diff --git a/lib/replicator.js b/lib/replicator.js index 369a6f08..c324008c 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -754,7 +754,7 @@ class Peer { if (this.remoteOpened === false) { this.replicator._ifAvailable-- - this.replicator.updateAll() + this.replicator._maybeResolveAvailable() return } @@ -2449,16 +2449,25 @@ module.exports = class Replicator { this.peers.splice(this.peers.indexOf(peer), 1) peer.wants.destroy() - if (this._manifestPeer === peer) this._manifestPeer = null + let needsUpdate = false + + if (this._manifestPeer === peer) { + this._manifestPeer = null + needsUpdate = true + } for (const req of this._inflight) { if (req.peer !== peer) continue this._inflight.remove(req.id, true) this._clearRequest(peer, req) + needsUpdate = true } this._onpeerupdate(false, peer) - this.updateAll() + + // If the peer had no requests to clear, there is nothing to reschedule. + if (needsUpdate && this.peers.length > 0) this.updateAll() + else this._maybeResolveAvailable() } _queueBlock(b) { @@ -2668,6 +2677,11 @@ module.exports = class Replicator { } } + _maybeResolveAvailable() { + this._checkUpgradeIfAvailable() + this._maybeResolveIfAvailableRanges() + } + _clearRequest(peer, req) { if (req.block !== null) { this._clearInflightBlock(this._blocks, req) diff --git a/test/replicate.js b/test/replicate.js index 9cdc933e..3daa29cc 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -793,6 +793,67 @@ test('multiplexing multiple times over the same stream', async function (t) { s3.destroy() }) +test('closing idle multiplexed stream does not update all peers', async function (t) { + const n1 = new NoiseSecretStream(true) + const n2 = new NoiseSecretStream(false) + + n1.rawStream.pipe(n2.rawStream).pipe(n1.rawStream) + + const cores = [] + const clones = [] + + for (let i = 0; i < 4; i++) { + const core = await create(t) + await core.append('block-' + i) + + const clone = await create(t, core.key) + + core.replicate(n1, { keepAlive: true }) + clone.replicate(n2, { keepAlive: true }) + + cores.push(core) + clones.push(clone) + } + + for (const clone of clones) await clone.get(0) + + await eventFlush() + + for (let i = 0; i < cores.length; i++) { + t.is(cores[i].peers.length, 1) + t.is(clones[i].peers.length, 1) + } + + let updates = 0 + const restore = [] + for (const core of cores.concat(clones)) { + const replicator = core.core.replicator + const updateAll = replicator.updateAll + restore.push(() => { + replicator.updateAll = updateAll + }) + + replicator.updateAll = function () { + updates++ + return updateAll.apply(this, arguments) + } + } + + t.teardown(() => { + for (const fn of restore) fn() + }) + + n1.destroy() + n2.destroy() + + await Promise.all([ + new Promise((resolve) => n1.once('close', resolve)), + new Promise((resolve) => n2.once('close', resolve)) + ]) + + t.is(updates, 0) +}) + test('destroying a stream and re-replicating works', async function (t) { const core = await create(t) From 38bc8644de3f140000798aeed750c486aa38cf00 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Fri, 15 May 2026 15:52:57 +0200 Subject: [PATCH 02/13] Fix close-path peer update scheduling --- lib/replicator.js | 12 +++- test/replicate.js | 157 ++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 167 insertions(+), 2 deletions(-) diff --git a/lib/replicator.js b/lib/replicator.js index c324008c..18d408ed 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -2465,8 +2465,16 @@ module.exports = class Replicator { this._onpeerupdate(false, peer) - // If the peer had no requests to clear, there is nothing to reschedule. - if (needsUpdate && this.peers.length > 0) this.updateAll() + // A close can also wake pending work that was not assigned to the closing peer. + const hasPendingUpdate = + needsUpdate || + this._queued.length > 0 || + this._seeks.length > 0 || + this._ranges.length > 0 || + this._reorgs.length > 0 || + this._maybeUpdate() + + if (hasPendingUpdate && this.peers.length > 0) this.updateAll() else this._maybeResolveAvailable() } diff --git a/test/replicate.js b/test/replicate.js index 3daa29cc..191f69b9 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -793,6 +793,52 @@ test('multiplexing multiple times over the same stream', async function (t) { s3.destroy() }) +test('closing peer with inflight block reschedules on remaining peer', async function (t) { + const writer = await create(t) + await writer.append(['block-0', 'block-1', 'block-2']) + + const mirror = await create(t, writer.key) + replicate(writer, mirror, t) + await mirror.download({ start: 0, end: writer.length }).done() + + const clone = await create(t, writer.key) + + const [slowWriterStream, slowCloneStream] = makeStreamPair(t, { latency: [200, 200] }) + writer.replicate(slowWriterStream, { keepAlive: true }) + clone.replicate(slowCloneStream, { keepAlive: true }) + + const mirrorStreams = replicate(mirror, clone, t, { keepAlive: true, teardown: false }) + t.teardown(() => destroyStreams(mirrorStreams)) + + await clone.update({ wait: true }) + await eventFlush() + + const slowPeer = await waitForPeerForStream(clone, slowWriterStream) + const mirrorPeer = await waitForPeerForStream(clone, mirrorStreams[0]) + + let allowMirror = false + const requestBlock = mirrorPeer._requestBlock + mirrorPeer._requestBlock = function () { + if (!allowMirror) return false + return requestBlock.apply(this, arguments) + } + + t.teardown(() => { + mirrorPeer._requestBlock = requestBlock + }) + + const block = clone.get(0) + await waitForInflightBlockOnPeer(clone, slowPeer) + + allowMirror = true + await destroyStreams([slowWriterStream, slowCloneStream]) + + t.alike( + await withTimeout(block, 1000, 'block did not resume after peer close'), + b4a.from('block-0') + ) +}) + test('closing idle multiplexed stream does not update all peers', async function (t) { const n1 = new NoiseSecretStream(true) const n2 = new NoiseSecretStream(false) @@ -854,6 +900,50 @@ test('closing idle multiplexed stream does not update all peers', async function t.is(updates, 0) }) +test('closing idle peer schedules pending range on remaining peer', async function (t) { + const writer = await create(t) + const batch = [] + for (let i = 0; i < 32; i++) batch.push('block-' + i) + await writer.append(batch) + + const sparse = await create(t, writer.key) + const clone = await create(t, writer.key) + + const writerStreams = replicate(writer, clone, t, { teardown: false }) + const sparseStreams = replicate(sparse, clone, t, { teardown: false }) + + t.teardown(() => destroyStreams(writerStreams)) + t.teardown(() => destroyStreams(sparseStreams)) + + await clone.update({ wait: true }) + await eventFlush() + + const writerPeer = findPeerForStream(clone, writerStreams[0]) + + let allowRange = false + const requestRange = writerPeer._requestRange + writerPeer._requestRange = function () { + if (!allowRange) return false + return requestRange.apply(this, arguments) + } + + t.teardown(() => { + writerPeer._requestRange = requestRange + }) + + const download = clone.download({ start: 0, end: writer.length }) + + await new Promise((resolve) => setTimeout(resolve, 100)) + t.is(clone.core.replicator._inflight.idle, true, 'range has no inflight requests yet') + + allowRange = true + await destroyStreams(sparseStreams) + + await withTimeout(download.done(), 500, 'download did not resume after peer close') + + t.is(clone.contiguousLength, writer.length) +}) + test('destroying a stream and re-replicating works', async function (t) { const core = await create(t) @@ -2955,6 +3045,73 @@ async function createAndDownload(t, core) { return b } +function findPeerForStream(core, remoteStream) { + for (const peer of core.core.replicator.peers) { + if (b4a.equals(peer.stream.remotePublicKey, remoteStream.publicKey)) return peer + } + + throw new Error('peer not found') +} + +async function waitForPeerForStream(core, remoteStream) { + const deadline = Date.now() + 5000 + + while (Date.now() < deadline) { + for (const peer of core.core.replicator.peers) { + if (b4a.equals(peer.stream.remotePublicKey, remoteStream.publicKey)) return peer + } + + await new Promise((resolve) => setImmediate(resolve)) + } + + throw new Error('peer not found') +} + +async function waitForInflightBlockOnPeer(core, peer) { + const deadline = Date.now() + 1000 + + while (Date.now() < deadline) { + const req = core.core.replicator._inflight._requests.find((req) => { + return req && req.peer === peer && req.block + }) + + if (req) return req + + await new Promise((resolve) => setImmediate(resolve)) + } + + throw new Error('block request was not inflight on expected peer') +} + +async function withTimeout(promise, ms, message) { + let timeout = null + + try { + return await Promise.race([ + promise, + new Promise((resolve, reject) => { + timeout = setTimeout(() => reject(new Error(message)), ms) + }) + ]) + } finally { + clearTimeout(timeout) + } +} + +function destroyStreams(streams) { + return Promise.all( + streams.map((s) => { + if (s.destroyed) return null + + return new Promise((resolve) => { + s.on('error', () => {}) + s.on('close', resolve) + s.destroy() + }) + }) + ) +} + async function waitForRequestBlock(core) { while (true) { const reqBlock = core.core.replicator._inflight._requests.find((req) => req && req.block) From b8083d17009e6078e8ff2d3479ff207a717e36f5 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Mon, 18 May 2026 15:50:46 +0200 Subject: [PATCH 03/13] Tighten idle close update scan test --- test/replicate.js | 72 ++++++++++++++++++++--------------------------- 1 file changed, 30 insertions(+), 42 deletions(-) diff --git a/test/replicate.js b/test/replicate.js index 191f69b9..39ed5f18 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -840,64 +840,52 @@ test('closing peer with inflight block reschedules on remaining peer', async fun }) test('closing idle multiplexed stream does not update all peers', async function (t) { - const n1 = new NoiseSecretStream(true) - const n2 = new NoiseSecretStream(false) - - n1.rawStream.pipe(n2.rawStream).pipe(n1.rawStream) - - const cores = [] - const clones = [] - - for (let i = 0; i < 4; i++) { - const core = await create(t) - await core.append('block-' + i) - - const clone = await create(t, core.key) + const writer = await create(t) + await writer.append('block') - core.replicate(n1, { keepAlive: true }) - clone.replicate(n2, { keepAlive: true }) + const clone = await create(t, writer.key) + const streams = [] - cores.push(core) - clones.push(clone) + for (let i = 0; i < 3; i++) { + const peer = new Promise((resolve) => clone.once('peer-add', resolve)) + streams.push(replicate(writer, clone, t, { keepAlive: true })) + await peer } - for (const clone of clones) await clone.get(0) - + await clone.get(0) await eventFlush() - for (let i = 0; i < cores.length; i++) { - t.is(cores[i].peers.length, 1) - t.is(clones[i].peers.length, 1) - } + t.is(clone.peers.length, 3, 'clone has multiple peers') + + const replicator = clone.core.replicator + const updateAll = replicator.updateAll + const updatePeer = replicator._updatePeer let updates = 0 - const restore = [] - for (const core of cores.concat(clones)) { - const replicator = core.core.replicator - const updateAll = replicator.updateAll - restore.push(() => { - replicator.updateAll = updateAll - }) + let peerScans = 0 - replicator.updateAll = function () { - updates++ - return updateAll.apply(this, arguments) - } + replicator.updateAll = function () { + updates++ + return updateAll.apply(this, arguments) + } + + replicator._updatePeer = function () { + peerScans++ + return updatePeer.apply(this, arguments) } t.teardown(() => { - for (const fn of restore) fn() + replicator.updateAll = updateAll + replicator._updatePeer = updatePeer }) - n1.destroy() - n2.destroy() - - await Promise.all([ - new Promise((resolve) => n1.once('close', resolve)), - new Promise((resolve) => n2.once('close', resolve)) - ]) + const peerRemoved = new Promise((resolve) => clone.once('peer-remove', resolve)) + await unreplicate(streams[0]) + await peerRemoved + t.is(clone.peers.length, 2, 'clone still has remaining peers') t.is(updates, 0) + t.is(peerScans, 0) }) test('closing idle peer schedules pending range on remaining peer', async function (t) { From 6a5faed794d47061ed63db73901d8f40735d8aa3 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Mon, 18 May 2026 16:48:17 +0200 Subject: [PATCH 04/13] Use peer-add events in close reschedule tests --- test/replicate.js | 34 +++++++--------------------------- 1 file changed, 7 insertions(+), 27 deletions(-) diff --git a/test/replicate.js b/test/replicate.js index 39ed5f18..ca526d3e 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -804,18 +804,19 @@ test('closing peer with inflight block reschedules on remaining peer', async fun const clone = await create(t, writer.key) const [slowWriterStream, slowCloneStream] = makeStreamPair(t, { latency: [200, 200] }) + const slowPeerWait = new Promise((resolve) => clone.once('peer-add', resolve)) writer.replicate(slowWriterStream, { keepAlive: true }) clone.replicate(slowCloneStream, { keepAlive: true }) + const slowPeer = await slowPeerWait + const mirrorPeerWait = new Promise((resolve) => clone.once('peer-add', resolve)) const mirrorStreams = replicate(mirror, clone, t, { keepAlive: true, teardown: false }) t.teardown(() => destroyStreams(mirrorStreams)) + const mirrorPeer = await mirrorPeerWait await clone.update({ wait: true }) await eventFlush() - const slowPeer = await waitForPeerForStream(clone, slowWriterStream) - const mirrorPeer = await waitForPeerForStream(clone, mirrorStreams[0]) - let allowMirror = false const requestBlock = mirrorPeer._requestBlock mirrorPeer._requestBlock = function () { @@ -897,7 +898,10 @@ test('closing idle peer schedules pending range on remaining peer', async functi const sparse = await create(t, writer.key) const clone = await create(t, writer.key) + const writerPeerWait = new Promise((resolve) => clone.once('peer-add', resolve)) const writerStreams = replicate(writer, clone, t, { teardown: false }) + const writerPeer = await writerPeerWait + const sparseStreams = replicate(sparse, clone, t, { teardown: false }) t.teardown(() => destroyStreams(writerStreams)) @@ -906,8 +910,6 @@ test('closing idle peer schedules pending range on remaining peer', async functi await clone.update({ wait: true }) await eventFlush() - const writerPeer = findPeerForStream(clone, writerStreams[0]) - let allowRange = false const requestRange = writerPeer._requestRange writerPeer._requestRange = function () { @@ -3033,28 +3035,6 @@ async function createAndDownload(t, core) { return b } -function findPeerForStream(core, remoteStream) { - for (const peer of core.core.replicator.peers) { - if (b4a.equals(peer.stream.remotePublicKey, remoteStream.publicKey)) return peer - } - - throw new Error('peer not found') -} - -async function waitForPeerForStream(core, remoteStream) { - const deadline = Date.now() + 5000 - - while (Date.now() < deadline) { - for (const peer of core.core.replicator.peers) { - if (b4a.equals(peer.stream.remotePublicKey, remoteStream.publicKey)) return peer - } - - await new Promise((resolve) => setImmediate(resolve)) - } - - throw new Error('peer not found') -} - async function waitForInflightBlockOnPeer(core, peer) { const deadline = Date.now() + 1000 From e61c1a009e5dab7d162bbe0d9facacdc2e409d7e Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Mon, 18 May 2026 17:05:18 +0200 Subject: [PATCH 05/13] Remove redundant close test teardown --- test/replicate.js | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/test/replicate.js b/test/replicate.js index ca526d3e..886ce8dc 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -899,13 +899,10 @@ test('closing idle peer schedules pending range on remaining peer', async functi const clone = await create(t, writer.key) const writerPeerWait = new Promise((resolve) => clone.once('peer-add', resolve)) - const writerStreams = replicate(writer, clone, t, { teardown: false }) + const writerStreams = replicate(writer, clone, t) const writerPeer = await writerPeerWait - const sparseStreams = replicate(sparse, clone, t, { teardown: false }) - - t.teardown(() => destroyStreams(writerStreams)) - t.teardown(() => destroyStreams(sparseStreams)) + const sparseStreams = replicate(sparse, clone, t) await clone.update({ wait: true }) await eventFlush() From d2d0e2d289eaa36c87ea9d7077bc42dfeef366f6 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Mon, 18 May 2026 17:17:42 +0200 Subject: [PATCH 06/13] Use unreplicate in close reschedule tests --- test/replicate.js | 22 ++++------------------ 1 file changed, 4 insertions(+), 18 deletions(-) diff --git a/test/replicate.js b/test/replicate.js index 886ce8dc..aad38c5f 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -811,7 +811,7 @@ test('closing peer with inflight block reschedules on remaining peer', async fun const mirrorPeerWait = new Promise((resolve) => clone.once('peer-add', resolve)) const mirrorStreams = replicate(mirror, clone, t, { keepAlive: true, teardown: false }) - t.teardown(() => destroyStreams(mirrorStreams)) + t.teardown(() => unreplicate(mirrorStreams)) const mirrorPeer = await mirrorPeerWait await clone.update({ wait: true }) @@ -832,7 +832,7 @@ test('closing peer with inflight block reschedules on remaining peer', async fun await waitForInflightBlockOnPeer(clone, slowPeer) allowMirror = true - await destroyStreams([slowWriterStream, slowCloneStream]) + await unreplicate([slowWriterStream, slowCloneStream]) t.alike( await withTimeout(block, 1000, 'block did not resume after peer close'), @@ -899,7 +899,7 @@ test('closing idle peer schedules pending range on remaining peer', async functi const clone = await create(t, writer.key) const writerPeerWait = new Promise((resolve) => clone.once('peer-add', resolve)) - const writerStreams = replicate(writer, clone, t) + replicate(writer, clone, t) const writerPeer = await writerPeerWait const sparseStreams = replicate(sparse, clone, t) @@ -924,7 +924,7 @@ test('closing idle peer schedules pending range on remaining peer', async functi t.is(clone.core.replicator._inflight.idle, true, 'range has no inflight requests yet') allowRange = true - await destroyStreams(sparseStreams) + await unreplicate(sparseStreams) await withTimeout(download.done(), 500, 'download did not resume after peer close') @@ -3063,20 +3063,6 @@ async function withTimeout(promise, ms, message) { } } -function destroyStreams(streams) { - return Promise.all( - streams.map((s) => { - if (s.destroyed) return null - - return new Promise((resolve) => { - s.on('error', () => {}) - s.on('close', resolve) - s.destroy() - }) - }) - ) -} - async function waitForRequestBlock(core) { while (true) { const reqBlock = core.core.replicator._inflight._requests.find((req) => req && req.block) From 68cbe0c2b46bd0e5ef42bed029c5b6373247d73e Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Mon, 18 May 2026 17:49:11 +0200 Subject: [PATCH 07/13] Use peer pause in close reschedule tests --- test/replicate.js | 22 ++++++---------------- 1 file changed, 6 insertions(+), 16 deletions(-) diff --git a/test/replicate.js b/test/replicate.js index aad38c5f..80f91785 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -817,21 +817,16 @@ test('closing peer with inflight block reschedules on remaining peer', async fun await clone.update({ wait: true }) await eventFlush() - let allowMirror = false - const requestBlock = mirrorPeer._requestBlock - mirrorPeer._requestBlock = function () { - if (!allowMirror) return false - return requestBlock.apply(this, arguments) - } + mirrorPeer.paused = true t.teardown(() => { - mirrorPeer._requestBlock = requestBlock + mirrorPeer.paused = false }) const block = clone.get(0) await waitForInflightBlockOnPeer(clone, slowPeer) - allowMirror = true + mirrorPeer.paused = false await unreplicate([slowWriterStream, slowCloneStream]) t.alike( @@ -907,15 +902,10 @@ test('closing idle peer schedules pending range on remaining peer', async functi await clone.update({ wait: true }) await eventFlush() - let allowRange = false - const requestRange = writerPeer._requestRange - writerPeer._requestRange = function () { - if (!allowRange) return false - return requestRange.apply(this, arguments) - } + writerPeer.paused = true t.teardown(() => { - writerPeer._requestRange = requestRange + writerPeer.paused = false }) const download = clone.download({ start: 0, end: writer.length }) @@ -923,7 +913,7 @@ test('closing idle peer schedules pending range on remaining peer', async functi await new Promise((resolve) => setTimeout(resolve, 100)) t.is(clone.core.replicator._inflight.idle, true, 'range has no inflight requests yet') - allowRange = true + writerPeer.paused = false await unreplicate(sparseStreams) await withTimeout(download.done(), 500, 'download did not resume after peer close') From bd62b177a5d795e3900e797e2ba892c3a29011d7 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Mon, 18 May 2026 17:55:49 +0200 Subject: [PATCH 08/13] Assert sparse peer removal in close range test --- test/replicate.js | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/test/replicate.js b/test/replicate.js index 80f91785..2c1e2721 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -897,7 +897,9 @@ test('closing idle peer schedules pending range on remaining peer', async functi replicate(writer, clone, t) const writerPeer = await writerPeerWait + const sparsePeerWait = new Promise((resolve) => clone.once('peer-add', resolve)) const sparseStreams = replicate(sparse, clone, t) + const sparsePeer = await sparsePeerWait await clone.update({ wait: true }) await eventFlush() @@ -913,9 +915,13 @@ test('closing idle peer schedules pending range on remaining peer', async functi await new Promise((resolve) => setTimeout(resolve, 100)) t.is(clone.core.replicator._inflight.idle, true, 'range has no inflight requests yet') + const peerRemoved = new Promise((resolve) => clone.once('peer-remove', resolve)) + writerPeer.paused = false await unreplicate(sparseStreams) + t.is(await peerRemoved, sparsePeer) + await withTimeout(download.done(), 500, 'download did not resume after peer close') t.is(clone.contiguousLength, writer.length) From b17110b5bb42f7d9f93a85b1a2d4c204d281233f Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Mon, 18 May 2026 18:19:00 +0200 Subject: [PATCH 09/13] Use upload event in close reschedule test --- test/replicate.js | 45 ++++++++++++++++++++------------------------- 1 file changed, 20 insertions(+), 25 deletions(-) diff --git a/test/replicate.js b/test/replicate.js index 2c1e2721..0541c41d 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -803,11 +803,9 @@ test('closing peer with inflight block reschedules on remaining peer', async fun const clone = await create(t, writer.key) - const [slowWriterStream, slowCloneStream] = makeStreamPair(t, { latency: [200, 200] }) - const slowPeerWait = new Promise((resolve) => clone.once('peer-add', resolve)) - writer.replicate(slowWriterStream, { keepAlive: true }) - clone.replicate(slowCloneStream, { keepAlive: true }) - const slowPeer = await slowPeerWait + const writerPeerWait = new Promise((resolve) => clone.once('peer-add', resolve)) + const writerStreams = replicate(writer, clone, t, { keepAlive: true }) + const writerPeer = await writerPeerWait const mirrorPeerWait = new Promise((resolve) => clone.once('peer-add', resolve)) const mirrorStreams = replicate(mirror, clone, t, { keepAlive: true, teardown: false }) @@ -823,11 +821,24 @@ test('closing peer with inflight block reschedules on remaining peer', async fun mirrorPeer.paused = false }) - const block = clone.get(0) - await waitForInflightBlockOnPeer(clone, slowPeer) + const peerRemoved = new Promise((resolve) => clone.once('peer-remove', resolve)) + const uploaded = new Promise((resolve, reject) => { + writer.once('upload', async function (index) { + try { + t.is(index, 0) + + mirrorPeer.paused = false + await unreplicate(writerStreams) + t.is(await peerRemoved, writerPeer) + resolve() + } catch (err) { + reject(err) + } + }) + }) - mirrorPeer.paused = false - await unreplicate([slowWriterStream, slowCloneStream]) + const block = clone.get(0) + await uploaded t.alike( await withTimeout(block, 1000, 'block did not resume after peer close'), @@ -3028,22 +3039,6 @@ async function createAndDownload(t, core) { return b } -async function waitForInflightBlockOnPeer(core, peer) { - const deadline = Date.now() + 1000 - - while (Date.now() < deadline) { - const req = core.core.replicator._inflight._requests.find((req) => { - return req && req.peer === peer && req.block - }) - - if (req) return req - - await new Promise((resolve) => setImmediate(resolve)) - } - - throw new Error('block request was not inflight on expected peer') -} - async function withTimeout(promise, ms, message) { let timeout = null From ad0a00f1a6f763b2269c8407df1704599116e6d1 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Mon, 18 May 2026 21:59:25 +0200 Subject: [PATCH 10/13] Wait for append in close reschedule tests --- test/replicate.js | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/test/replicate.js b/test/replicate.js index 0541c41d..871bafc5 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -802,6 +802,7 @@ test('closing peer with inflight block reschedules on remaining peer', async fun await mirror.download({ start: 0, end: writer.length }).done() const clone = await create(t, writer.key) + const cloneAppend = new Promise((resolve) => clone.once('append', resolve)) const writerPeerWait = new Promise((resolve) => clone.once('peer-add', resolve)) const writerStreams = replicate(writer, clone, t, { keepAlive: true }) @@ -812,7 +813,7 @@ test('closing peer with inflight block reschedules on remaining peer', async fun t.teardown(() => unreplicate(mirrorStreams)) const mirrorPeer = await mirrorPeerWait - await clone.update({ wait: true }) + await cloneAppend await eventFlush() mirrorPeer.paused = true @@ -903,6 +904,7 @@ test('closing idle peer schedules pending range on remaining peer', async functi const sparse = await create(t, writer.key) const clone = await create(t, writer.key) + const cloneAppend = new Promise((resolve) => clone.once('append', resolve)) const writerPeerWait = new Promise((resolve) => clone.once('peer-add', resolve)) replicate(writer, clone, t) @@ -912,7 +914,7 @@ test('closing idle peer schedules pending range on remaining peer', async functi const sparseStreams = replicate(sparse, clone, t) const sparsePeer = await sparsePeerWait - await clone.update({ wait: true }) + await cloneAppend await eventFlush() writerPeer.paused = true From af1c4d247304ae0ebc6454bed94f36c4d94d4883 Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Mon, 18 May 2026 22:13:38 +0200 Subject: [PATCH 11/13] Remove close test timeout guards --- test/replicate.js | 22 ++-------------------- 1 file changed, 2 insertions(+), 20 deletions(-) diff --git a/test/replicate.js b/test/replicate.js index 871bafc5..d354c333 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -841,10 +841,7 @@ test('closing peer with inflight block reschedules on remaining peer', async fun const block = clone.get(0) await uploaded - t.alike( - await withTimeout(block, 1000, 'block did not resume after peer close'), - b4a.from('block-0') - ) + t.alike(await block, b4a.from('block-0')) }) test('closing idle multiplexed stream does not update all peers', async function (t) { @@ -935,7 +932,7 @@ test('closing idle peer schedules pending range on remaining peer', async functi t.is(await peerRemoved, sparsePeer) - await withTimeout(download.done(), 500, 'download did not resume after peer close') + await download.done() t.is(clone.contiguousLength, writer.length) }) @@ -3041,21 +3038,6 @@ async function createAndDownload(t, core) { return b } -async function withTimeout(promise, ms, message) { - let timeout = null - - try { - return await Promise.race([ - promise, - new Promise((resolve, reject) => { - timeout = setTimeout(() => reject(new Error(message)), ms) - }) - ]) - } finally { - clearTimeout(timeout) - } -} - async function waitForRequestBlock(core) { while (true) { const reqBlock = core.core.replicator._inflight._requests.find((req) => req && req.block) From c53dc6ecc9c82d8f31a61a968ac8c6c07808d67e Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Mon, 18 May 2026 22:36:20 +0200 Subject: [PATCH 12/13] Move idle close optimization to follow-up --- lib/replicator.js | 28 +++------------------------ test/replicate.js | 49 ----------------------------------------------- 2 files changed, 3 insertions(+), 74 deletions(-) diff --git a/lib/replicator.js b/lib/replicator.js index 18d408ed..369a6f08 100644 --- a/lib/replicator.js +++ b/lib/replicator.js @@ -754,7 +754,7 @@ class Peer { if (this.remoteOpened === false) { this.replicator._ifAvailable-- - this.replicator._maybeResolveAvailable() + this.replicator.updateAll() return } @@ -2449,33 +2449,16 @@ module.exports = class Replicator { this.peers.splice(this.peers.indexOf(peer), 1) peer.wants.destroy() - let needsUpdate = false - - if (this._manifestPeer === peer) { - this._manifestPeer = null - needsUpdate = true - } + if (this._manifestPeer === peer) this._manifestPeer = null for (const req of this._inflight) { if (req.peer !== peer) continue this._inflight.remove(req.id, true) this._clearRequest(peer, req) - needsUpdate = true } this._onpeerupdate(false, peer) - - // A close can also wake pending work that was not assigned to the closing peer. - const hasPendingUpdate = - needsUpdate || - this._queued.length > 0 || - this._seeks.length > 0 || - this._ranges.length > 0 || - this._reorgs.length > 0 || - this._maybeUpdate() - - if (hasPendingUpdate && this.peers.length > 0) this.updateAll() - else this._maybeResolveAvailable() + this.updateAll() } _queueBlock(b) { @@ -2685,11 +2668,6 @@ module.exports = class Replicator { } } - _maybeResolveAvailable() { - this._checkUpgradeIfAvailable() - this._maybeResolveIfAvailableRanges() - } - _clearRequest(peer, req) { if (req.block !== null) { this._clearInflightBlock(this._blocks, req) diff --git a/test/replicate.js b/test/replicate.js index d354c333..20b1b02f 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -844,55 +844,6 @@ test('closing peer with inflight block reschedules on remaining peer', async fun t.alike(await block, b4a.from('block-0')) }) -test('closing idle multiplexed stream does not update all peers', async function (t) { - const writer = await create(t) - await writer.append('block') - - const clone = await create(t, writer.key) - const streams = [] - - for (let i = 0; i < 3; i++) { - const peer = new Promise((resolve) => clone.once('peer-add', resolve)) - streams.push(replicate(writer, clone, t, { keepAlive: true })) - await peer - } - - await clone.get(0) - await eventFlush() - - t.is(clone.peers.length, 3, 'clone has multiple peers') - - const replicator = clone.core.replicator - const updateAll = replicator.updateAll - const updatePeer = replicator._updatePeer - - let updates = 0 - let peerScans = 0 - - replicator.updateAll = function () { - updates++ - return updateAll.apply(this, arguments) - } - - replicator._updatePeer = function () { - peerScans++ - return updatePeer.apply(this, arguments) - } - - t.teardown(() => { - replicator.updateAll = updateAll - replicator._updatePeer = updatePeer - }) - - const peerRemoved = new Promise((resolve) => clone.once('peer-remove', resolve)) - await unreplicate(streams[0]) - await peerRemoved - - t.is(clone.peers.length, 2, 'clone still has remaining peers') - t.is(updates, 0) - t.is(peerScans, 0) -}) - test('closing idle peer schedules pending range on remaining peer', async function (t) { const writer = await create(t) const batch = [] From b9865a885c01dfcaf8254d5f461ed975b1355b4a Mon Sep 17 00:00:00 2001 From: marcus-pousette-hp Date: Mon, 18 May 2026 22:55:56 +0200 Subject: [PATCH 13/13] Clarify close reschedule test setup --- test/replicate.js | 1 + 1 file changed, 1 insertion(+) diff --git a/test/replicate.js b/test/replicate.js index 20b1b02f..6ddd3063 100644 --- a/test/replicate.js +++ b/test/replicate.js @@ -816,6 +816,7 @@ test('closing peer with inflight block reschedules on remaining peer', async fun await cloneAppend await eventFlush() + // Force the initial block request onto the writer so closing it exercises rescheduling. mirrorPeer.paused = true t.teardown(() => {