diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-09-24 00:10:41 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-09-24 00:10:41 +0200 |
| commit | 059eb0318daf27d98bd1c9532705ff406c0a10f8 (patch) | |
| tree | 662f03d0df9306d2b8f861ce394835d3748ae3f8 /packages/meshbay-hub/tests/harness | |
| parent | 2d657f40ebd697e4332c95d7a57bbb292ff46012 (diff) | |
| download | meshbay-0.15.tar.gz | |
fix(hub): Search reaches each group with one offer, and a busy hub is not a dead node0.15
The other half of the 4G failure. The page negotiated every group twice: the
sweep opened a connection, read the index and closed it, then the warm-up
opened the same group again. The sweep, the warm-up and the tiles each had a
concurrency ceiling of their own, and together they went past what the hub
admits per account. Whatever the hub refused was then reported as "node
unreachable" and remembered as down, which put that group last next time.
- Every connection goes through ConnectionPool, which holds the page's one
ceiling (six at once, sized for twenty groups on a phone) and keeps what
the sweep opened for the tiles. A visit costs one offer per group. A
refresh costs none for a connection that answers a four-second ping, and
a connection that died while the phone slept is replaced, not waited on.
- A connection whose index is being read is held against eviction. With
more groups than the pool keeps, it was otherwise the least recently used
one.
- Negotiations still under way when the page closes close what they get,
and a sweep cut short that way remembers nobody as down.
- transport.js sends an offer again on 429, 502 or 503, honouring
Retry-After, with jittered waits of about twenty seconds at worst. Search
counts each retry as progress. A 404, 403 or 504 still fails at once, so
a dead node costs no time.
The fan-out tests assumed a ceiling of three and were re-measured: four dead
groups of twelve now hold nothing back, even on a first visit. The pool and
the retry run as shipped code, lifted as text, against a fake clock. Each
guard was checked by removing it and seeing its test fail.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Diffstat (limited to 'packages/meshbay-hub/tests/harness')
3 files changed, 307 insertions, 2 deletions
diff --git a/packages/meshbay-hub/tests/harness/offer_retry_harness.mjs b/packages/meshbay-hub/tests/harness/offer_retry_harness.mjs new file mode 100644 index 0000000..3a4b568 --- /dev/null +++ b/packages/meshbay-hub/tests/harness/offer_retry_harness.mjs @@ -0,0 +1,84 @@ +// Run the SHIPPED `postOffer` against a fake hub and a fake clock. +// +// `postOffer` and the constants that govern it are lifted out of transport.js +// as text. What is modelled is the hub — a list of the statuses it answers in +// turn, each with an optional Retry-After — and the clock. +// +// Usage: node offer_retry_harness.mjs <path to transport.js> <json config> +import { readFileSync } from 'fs'; + +const src = readFileSync(process.argv[2], 'utf8'); +const cfg = JSON.parse(process.argv[3] || '{}'); + +const { + // What the hub answers, one entry per POST: a status, or [status, retryAfter]. + answers = [200], + // After this many ms the caller closes its transport. null: never. + closeAt = null, + runUntil = 600000, +} = cfg; + +let now = 0; +let nextId = 1; +let timers = []; +globalThis.setTimeout = (fn, ms) => { + const t = { at: now + (ms || 0), fn, id: nextId++ }; + timers.push(t); + return t.id; +}; +globalThis.clearTimeout = (id) => { timers = timers.filter((t) => t.id !== id); }; +const flush = async () => { for (let i = 0; i < 50; i++) await Promise.resolve(); }; +async function run(until) { + await flush(); + while (timers.length) { + const due = timers.reduce((a, b) => (b.at < a.at ? b : a)); + if (due.at > until) break; + timers = timers.filter((t) => t !== due); + now = due.at; + due.fn(); + await flush(); + } + now = until; +} +// Jitter at its midpoint, so the delays asserted are the nominal ones. +Math.random = () => 0.5; + +const posts = []; +const call = async () => { + const a = answers[Math.min(posts.length, answers.length - 1)]; + const [status, retryAfter] = Array.isArray(a) ? a : [a, null]; + posts.push(now); + return { + ok: status >= 200 && status < 300, + status, + headers: new Headers(retryAfter === null ? {} : { 'Retry-After': String(retryAfter) }), + json: async () => ({ detail: `status ${status}` }), + }; +}; + +const lift = (signature, end) => { + const start = src.indexOf(signature); + if (start < 0) throw new Error(`${signature} is gone from transport.js`); + return src.slice(start, src.indexOf(end, start) + end.length); +}; +const make = new Function( + `${lift('const OFFER_RETRY_STATUSES', ';\n')} + ${lift('const OFFER_RETRY_DELAYS_MS', ';\n')} + ${lift('const OFFER_RETRY_AFTER_MAX_MS', ';\n')} + ${lift('async function postOffer(', '\n}\n')} + return postOffer;`, +); +const postOffer = make(); + +let closed = false; +if (closeAt !== null) setTimeout(() => { closed = true; }, closeAt); +const retries = []; +let outcome = null; +postOffer(call, 'https://hub.example/v1/nodes/n/webrtc/offer', {}, { + isClosed: () => closed, + onRetry: (status, delay) => retries.push({ status, delay }), +}).then(() => { outcome = { result: 'answered', at: now }; }) + .catch((e) => { outcome = { result: 'failed', at: now, status: e.status ?? null }; }); + +await run(runUntil); +console.log(JSON.stringify({ ...(outcome || { result: 'pending' }), posts, retries })); diff --git a/packages/meshbay-hub/tests/harness/search_fanout_harness.mjs b/packages/meshbay-hub/tests/harness/search_fanout_harness.mjs index c05e23b..5f95ede 100644 --- a/packages/meshbay-hub/tests/harness/search_fanout_harness.mjs +++ b/packages/meshbay-hub/tests/harness/search_fanout_harness.mjs @@ -80,7 +80,7 @@ const session = { bundleKey: 'k' }; const _loadBundleKey = async () => 'k'; // One group's index, answered on the clock rather than over a network. -const fetchGroupIndex = (groupId) => new Promise((resolve, reject) => { +const fetchGroupIndex = (_pool, groupId) => new Promise((resolve, reject) => { const n = Number(String(groupId).split('-')[1]); const isDead = dead.includes(n); setTimeout( @@ -123,7 +123,7 @@ const list = Array.from({ length: groups }, (_, i) => ({ id: `g-${i}`, name: `G$ let finishedAt = null; fetchAllIndexes( - list, 'tok', 'alice', 'u1', + null, list, 'tok', 'alice', 'u1', () => {}, // progress counters only (results) => { shown.push({ at: now, groups: results.size }); }, ).then(() => { finishedAt = now; }); diff --git a/packages/meshbay-hub/tests/harness/search_pool_harness.mjs b/packages/meshbay-hub/tests/harness/search_pool_harness.mjs new file mode 100644 index 0000000..521a97e --- /dev/null +++ b/packages/meshbay-hub/tests/harness/search_pool_harness.mjs @@ -0,0 +1,221 @@ +// Run the SHIPPED `ConnectionPool`, `fetchGroupIndex` and `fetchAllIndexes` +// against a fake clock and fake nodes, and count offers. +// +// The question this answers is how many connections Search negotiates, how many +// at once, and whether one it is reading from can be closed under it — the +// things that decided whether a phone ran into the hub's per-account ceilings. +// The pool, the index fetch and the sweep are lifted out of search-page.js as +// text; what is modelled is the clock and the nodes (`connectToGroup`, which +// is where an offer is made). +// +// Usage: node search_pool_harness.mjs <path to search-page.js> <json config> +import { readFileSync } from 'fs'; + +const src = readFileSync(process.argv[2], 'utf8'); +const cfg = JSON.parse(process.argv[3] || '{}'); + +const { + scenario = 'sweep', + groups = 5, + // Group indexes whose node never answers; `connectToGroup` fails after deadMs. + dead = [], + connectMs = 2000, // a phone on 4G + deadMs = 10000, + indexMs = 300, + runUntil = 600000, +} = cfg; + +// ── The clock ──────────────────────────────────────────────────────────────── + +let now = 0; +let nextId = 1; +let timers = []; +globalThis.setTimeout = (fn, ms) => { + const t = { at: now + (ms || 0), fn, id: nextId++ }; + timers.push(t); + return t.id; +}; +globalThis.clearTimeout = (id) => { timers = timers.filter((t) => t.id !== id); }; +// The pool orders its connections by `Date.now()`, so that is the fake clock too. +Date.now = () => now; +const flush = async () => { for (let i = 0; i < 50; i++) await Promise.resolve(); }; +async function run(until) { + await flush(); + while (timers.length) { + const due = timers.reduce((a, b) => (b.at < a.at ? b : a)); + if (due.at > until) break; + timers = timers.filter((t) => t !== due); + now = due.at; + due.fn(); + await flush(); + } + now = until; +} +const sleep = (ms) => new Promise((r) => setTimeout(r, ms)); + +// ── The nodes ──────────────────────────────────────────────────────────────── + +let offers = 0; +let negotiating = 0; +let peakNegotiating = 0; +let readsOnClosed = 0; +let pings = 0; +const transports = []; +// Groups whose held connection has died quietly, as a phone's does in its sleep. +const silentlyDead = new Set(); + +class FakeTransport { + constructor(groupId) { + this.groupId = groupId; + this.closedFlag = false; + transports.push(this); + } + get connected() { return !this.closedFlag; } + close() { this.closedFlag = true; } + async ping(timeoutMs) { + pings += 1; + if (silentlyDead.has(this.groupId)) { + await sleep(timeoutMs); + throw new Error('Response timeout'); + } + await sleep(100); + } + async fetchIndex() { + // Indexes differ in size, so reads end at different times and overlap + // other groups' negotiations — which uniform timings never do. + const n = Number(String(this.groupId).split('-')[1]); + await sleep(indexMs * (1 + (n % 4))); + if (this.closedFlag) { readsOnClosed += 1; throw new Error('Transport closed'); } + return { entries: [{ id: `e-${this.groupId}` }] }; + } +} + +async function connectToGroup(_hub, groupId) { + offers += 1; + negotiating += 1; + peakNegotiating = Math.max(peakNegotiating, negotiating); + try { + const n = Number(String(groupId).split('-')[1]); + if (dead.includes(n)) { + await sleep(deadMs); + throw new Error('Connection timeout'); + } + await sleep(connectMs); + return { transport: new FakeTransport(groupId), ack: {} }; + } finally { + negotiating -= 1; + } +} + +globalThis.window = {}; +let remembered = null; +globalThis.localStorage = { getItem: () => null, setItem: (_k, v) => { remembered = v; } }; +const session = { bundleKey: 'k' }; +const _loadBundleKey = async () => 'k'; +const cacheGroupIndex = () => {}; + +// ── The code under test, as text ───────────────────────────────────────────── + +const constant = (name) => { + const m = new RegExp(`^const ${name} = (\\d+);`, 'm').exec(src); + if (!m) throw new Error(`${name} is gone from search-page.js`); + return `const ${name} = ${m[1]};`; +}; +const lift = (signature) => { + const start = src.indexOf(signature); + if (start < 0) throw new Error(`${signature} is gone from search-page.js`); + return src.slice(start, src.indexOf('\n}\n', start) + 3); +}; + +const make = new Function( + 'connectToGroup', 'window', 'session', '_loadBundleKey', 'cacheGroupIndex', + 'localStorage', + `${constant('MAX_IN_FLIGHT')} + ${constant('MAX_POOL_SIZE')} + ${constant('REUSE_PING_MS')} + const DOWN_KEY = 'harness'; + ${lift('class ConnectionPool {')} + ${lift('async function fetchGroupIndex(')} + ${lift('function lastKnownDown(')} + ${lift('function rememberDown(')} + ${lift('async function inFlight(')} + ${lift('async function fetchAllIndexes(')} + return { ConnectionPool, fetchAllIndexes, MAX_IN_FLIGHT, MAX_POOL_SIZE };`, +); +const code = make(connectToGroup, globalThis.window, session, _loadBundleKey, + cacheGroupIndex, globalThis.localStorage); + +// ── The scenarios ──────────────────────────────────────────────────────────── + +const list = Array.from({ length: groups }, (_, i) => ({ id: `g-${i}`, name: `G${i}` })); +let evicted = 0; +const pool = new code.ConnectionPool('https://hub.example', () => { evicted += 1; }); +const out = { maxInFlight: code.MAX_IN_FLIGHT, maxPool: code.MAX_POOL_SIZE }; + +// Groups whose index has landed — the only ones with tiles on screen. +const shown = new Set(); +const sweep = async () => { + const r = await code.fetchAllIndexes(pool, list, 'tok', 'alice', 'u1', () => {}, + (results) => { for (const id of results.keys()) shown.add(id); }); + return { found: r.results.size, unreachable: r.unreachable.length, at: now }; +}; + +if (scenario === 'sweep') { + sweep().then((r) => { Object.assign(out, r); }); + await run(runUntil); +} else if (scenario === 'refresh') { + // A sweep, then the Search page's own refresh button: the second must reuse. + sweep().then(async () => { + out.offersFirst = offers; + const again = await sweep(); + out.offersSecond = offers - out.offersFirst; + Object.assign(out, again); + }); + await run(runUntil); +} else if (scenario === 'slept') { + // A sweep, then the phone sleeps and two connections die without saying so. + sweep().then(async () => { + out.offersFirst = offers; + silentlyDead.add('g-0'); + silentlyDead.add('g-1'); + const again = await sweep(); + out.offersSecond = offers - out.offersFirst; + Object.assign(out, again); + }); + await run(runUntil); +} else if (scenario === 'tiles') { + // A tile of every group asks for its connection while the sweep is on. + sweep().then((r) => { Object.assign(out, r); }); + for (const g of list) pool.connect(g.id, 'tok', 'k', 'alice', 'u1').catch(() => {}); + await run(runUntil); +} else if (scenario === 'scroll') { + // Tiles of groups already on screen keep using their connections while the + // sweep is still reading other groups' indexes — which makes a connection an + // index is being read from the least recently used one in the pool. + let sweeping = true; + sweep().then((r) => { sweeping = false; Object.assign(out, r); }); + const touch = async () => { + while (sweeping) { + await sleep(250); + for (const id of shown) { + if (pool.has(id)) pool.connect(id, 'tok', 'k', 'alice', 'u1').catch(() => {}); + } + } + }; + touch(); + await run(runUntil); +} else if (scenario === 'unmount') { + // The page goes away while connections are still being negotiated. + sweep().catch(() => {}); + await run(connectMs / 2); + pool.closeAll(); + await run(runUntil); + out.openAfter = transports.filter((t) => !t.closedFlag).length; +} + +Object.assign(out, { + offers, peakNegotiating, readsOnClosed, pings, evicted, poolSize: pool.size, + remembered: remembered === null ? null : JSON.parse(remembered), + openTransports: transports.filter((t) => !t.closedFlag).length, +}); +console.log(JSON.stringify(out)); |