diff options
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)); |