summaryrefslogtreecommitdiffstats
path: root/packages
diff options
context:
space:
mode:
Diffstat (limited to 'packages')
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/search-page.js38
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/transport.js26
-rw-r--r--packages/meshbay-hub/tests/harness/search_connect_harness.mjs134
-rw-r--r--packages/meshbay-hub/tests/test_search_connect_deadline.py124
4 files changed, 317 insertions, 5 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/search-page.js b/packages/meshbay-hub/src/meshbay_hub/static/search-page.js
index 593760c..679b3f9 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/search-page.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/search-page.js
@@ -27,7 +27,22 @@ const BATCH_SIZE = 3;
// good. grenet (one shared group) never hit it; cbesson (many) always did.
const MAX_POOL_SIZE = 12;
const DEBOUNCE_MS = 200;
-const SEARCH_TIMEOUT = 10000;
+// How long a connection attempt may make **no progress**, not how long it may
+// take. Search used a flat 10 s for the whole of `connect()`, which is the hub
+// round trip, ICE gathering (capped at 4 s in transport.js), DTLS, the
+// DataChannel opening and the handshake's own round trips. Opening the same
+// group from the sidebar has no such deadline and gets the transport's own 30 s,
+// so a link slow enough to need twelve seconds — a phone on 4G — failed here and
+// worked from there, for the same work against the same node.
+//
+// A node that is not answering produces no progress event and still fails in
+// `SEARCH_STALL_MS`, which is what keeps the fan-out bounded: `fetchAllIndexes`
+// goes in batches of three and waits for the slowest of each, so the number that
+// matters for a page full of unreachable groups is this one, unchanged.
+// `SEARCH_MAX_MS` bounds the other case — a node that answers ICE and then stops
+// — because a deadline that only ever resets has none.
+const SEARCH_STALL_MS = 10000;
+const SEARCH_MAX_MS = 30000;
const SEARCH_VIDEO_ROOT = '__search__';
const SEARCH_AUDIO_ROOT = '__search__';
const SEARCH_PHOTO_ROOTS = ['__search_photos__'];
@@ -55,20 +70,33 @@ async function connectToGroup(hubBase, groupId, token, bundleKey, username, user
for (const n of nodesData.nodes) {
const transport = new window.MeshBayTransport(hubBase, live);
transport.onNeedToken = async () => (await ensureFreshToken()) || token;
- let timer;
+ let stallTimer, capTimer;
+ const stopTimers = () => {
+ clearTimeout(stallTimer);
+ clearTimeout(capTimer);
+ transport.onConnectProgress = null;
+ };
try {
const ack = await Promise.race([
transport.connect(
n.node_id, live, groupId, null, null, bundleKey,
username, userId, null),
new Promise((_, reject) => {
- timer = setTimeout(() => reject(new Error('Connection timeout')), SEARCH_TIMEOUT);
+ const giveUp = () => reject(new Error('Connection timeout'));
+ stallTimer = setTimeout(giveUp, SEARCH_STALL_MS);
+ capTimer = setTimeout(giveUp, SEARCH_MAX_MS);
+ // Each step the transport reports — the peer answering ICE, the
+ // channel opening — buys another window, never more than the cap.
+ transport.onConnectProgress = () => {
+ clearTimeout(stallTimer);
+ stallTimer = setTimeout(giveUp, SEARCH_STALL_MS);
+ };
}),
]);
- clearTimeout(timer);
+ stopTimers();
return { transport, ack };
} catch (e) {
- clearTimeout(timer);
+ stopTimers();
lastErr = e;
try { transport.close(); } catch {}
if (e.reason && e.reason !== 'not_hosted') throw e;
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport.js b/packages/meshbay-hub/src/meshbay_hub/static/transport.js
index 1a28fb0..72c831f 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/transport.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/transport.js
@@ -582,6 +582,23 @@ class MeshBayTransport {
get newNodeBundleRecovery() { return this._newNodeBundleRecovery || null; }
set newNodeBundleRecovery(v) { this._newNodeBundleRecovery = v; }
+ /**
+ * Tell a caller that connecting is getting somewhere.
+ *
+ * Only a caller that imposes its own deadline on `connect()` sets
+ * `onConnectProgress`, and only Search does: everywhere else a connection is
+ * opened one at a time and waits out the budgets in here. It exists so that
+ * deadline can measure *stalling* rather than elapsed time — a link slow
+ * enough to need twelve seconds is not the same thing as a node that is not
+ * answering, and a fixed number cannot tell them apart.
+ *
+ * Never lets a caller's callback break the connection it is reporting on.
+ */
+ _noteConnectProgress(phase) {
+ if (!this.onConnectProgress) return;
+ try { this.onConnectProgress(phase); } catch { /* the caller's problem */ }
+ }
+
async connect(nodeId, jwtToken, groupId, gekRaw, sessionKeys, bundleKey, username,
userId, joinCode, recoveryKey) {
// Remembered for _reconnectLoop, which calls connect() again with these
@@ -633,6 +650,7 @@ class MeshBayTransport {
clearTimeout(timeout);
this._connected = true;
trace('channel_open', {});
+ this._noteConnectProgress('channel_open');
resolve();
};
});
@@ -695,6 +713,14 @@ class MeshBayTransport {
pc.oniceconnectionstatechange = () => {
console.log('[MeshBay] ICE state:', pc.iceConnectionState);
trace('ice_state', { state: pc.iceConnectionState });
+ // `connected`/`completed` is the first moment this browser knows the peer
+ // is there at all: a candidate pair answered. `checking` is not — it is
+ // this side trying addresses that may all be dead.
+ if (pc === this._pc
+ && (pc.iceConnectionState === 'connected'
+ || pc.iceConnectionState === 'completed')) {
+ this._noteConnectProgress('ice_connected');
+ }
};
// Diagnostic-only: a periodic health ping and a resume-triggered one, so
diff --git a/packages/meshbay-hub/tests/harness/search_connect_harness.mjs b/packages/meshbay-hub/tests/harness/search_connect_harness.mjs
new file mode 100644
index 0000000..2479d46
--- /dev/null
+++ b/packages/meshbay-hub/tests/harness/search_connect_harness.mjs
@@ -0,0 +1,134 @@
+// Run the SHIPPED `connectToGroup` against a fake transport and a fake clock.
+//
+// Nothing here is a paraphrase of search-page.js: the function and the two
+// constants that govern its deadline are lifted out of the file as text and
+// executed. What is modelled is the *environment* — a clock the test advances,
+// a hub that lists nodes, and a transport that reports progress and resolves
+// whenever the scenario says it does.
+//
+// A real clock would make these scenarios a minute each and would hide the
+// thing worth asserting, which is the moment the deadline fires rather than
+// roughly how long it took.
+//
+// Usage: node search_connect_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 {
+ // When the transport reports progress, in ms from the start of the attempt.
+ progressAt = [],
+ // When connect() resolves. null: never.
+ resolveAt = null,
+ // How many nodes the hub lists for the group.
+ nodes = 1,
+ // How far to run the clock before giving up on the harness itself.
+ runUntil = 300000,
+} = 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);
+};
+
+// Promise jobs queued by the code under test have to run between one timer and
+// the next, or a `.then` that schedules a timer would be seen after the clock
+// has already passed it.
+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;
+}
+
+// ── The environment ──────────────────────────────────────────────────────────
+
+let closed = 0;
+
+class FakeTransport {
+ constructor() {
+ this.onConnectProgress = null;
+ this.onNeedToken = null;
+ }
+
+ connect() {
+ const started = now;
+ for (const at of progressAt) {
+ setTimeout(() => {
+ // Exactly what transport.js does: it never lets the caller's callback
+ // break the connection it is reporting on.
+ if (this.onConnectProgress) {
+ try { this.onConnectProgress('test'); } catch { /* caller's problem */ }
+ }
+ }, at);
+ }
+ return new Promise((resolve) => {
+ if (resolveAt !== null) {
+ setTimeout(() => resolve({ ok: true, at: now - started }), resolveAt);
+ }
+ });
+ }
+
+ close() { closed += 1; }
+}
+
+globalThis.window = { MeshBayTransport: FakeTransport };
+
+const hubFetch = async () => ({
+ nodes: Array.from({ length: nodes }, (_, i) => ({ node_id: `node-${i}` })),
+});
+const ensureFreshToken = async () => 'fresh-token';
+
+// ── 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 fnStart = src.indexOf('async function connectToGroup(');
+if (fnStart < 0) throw new Error('connectToGroup is gone from search-page.js');
+const fnEnd = src.indexOf('\n}\n', fnStart) + 3;
+const fnText = src.slice(fnStart, fnEnd);
+
+const make = new Function(
+ 'hubFetch', 'ensureFreshToken', 'window',
+ `${constant('SEARCH_STALL_MS')}
+ ${constant('SEARCH_MAX_MS')}
+ ${fnText}
+ return connectToGroup;`,
+);
+const connectToGroup = make(hubFetch, ensureFreshToken, globalThis.window);
+
+// ── The scenario ─────────────────────────────────────────────────────────────
+
+let outcome = null;
+connectToGroup('https://hub.example', 'g1', 'tok', null, 'alice', 'u1')
+ .then(() => { outcome = { result: 'connected', at: now }; })
+ .catch((e) => { outcome = { result: 'failed', at: now, error: String(e.message) }; });
+
+await run(runUntil);
+
+console.log(JSON.stringify({
+ ...(outcome || { result: 'still-waiting', at: now }),
+ closed,
+}));
diff --git a/packages/meshbay-hub/tests/test_search_connect_deadline.py b/packages/meshbay-hub/tests/test_search_connect_deadline.py
new file mode 100644
index 0000000..d41cbb5
--- /dev/null
+++ b/packages/meshbay-hub/tests/test_search_connect_deadline.py
@@ -0,0 +1,124 @@
+"""
+Search gave up on a connection that was still getting somewhere.
+
+Opening a group from the sidebar has no deadline of its own: it gets the
+transport's, which allows 30 s for the DataChannel to open and 30 s for each
+request after that. Search wrapped the same `connect()` in a flat 10 s — and
+that 10 s has to cover the hub round trip, ICE gathering (capped at 4 s in
+transport.js), DTLS, the channel opening and the handshake's own round trips.
+On a phone on 4G the budget is met by luck rather than by margin, and the same
+group then fails in Search and opens from the sidebar, against the same node.
+
+Raising the number is the wrong repair. `fetchAllIndexes` fans out in batches of
+three and waits for the slowest of each, so a page of unreachable groups costs
+(batches x the deadline) of spinner: a bigger number makes every dead group
+slower to report for the sake of the live ones.
+
+So the deadline measures **stalling** instead of elapsed time. A node that is
+not there produces no progress and still fails in `SEARCH_STALL_MS`, unchanged;
+a node that answers ICE and opens a channel buys another window each time, up to
+`SEARCH_MAX_MS`, which exists because a deadline that only ever resets has none.
+
+These run the shipped `connectToGroup`, lifted out of search-page.js as text,
+against a fake clock and a transport whose progress the scenario dictates — see
+harness/search_connect_harness.mjs. A real clock would make each of these a
+minute and would blur the one thing worth asserting, which is *when* the
+deadline fires.
+"""
+
+import json
+import shutil
+import subprocess
+from pathlib import Path
+
+import pytest
+
+STATIC = Path(__file__).resolve().parents[1] / "src" / "meshbay_hub" / "static"
+SEARCH_PAGE = STATIC / "search-page.js"
+HARNESS = Path(__file__).parent / "harness" / "search_connect_harness.mjs"
+
+pytestmark = pytest.mark.skipif(
+ shutil.which("node") is None or not SEARCH_PAGE.exists(),
+ reason="node or the SPA sources are not available")
+
+STALL_MS = 10000
+CAP_MS = 30000
+
+
+def _attempt(**cfg) -> dict:
+ proc = subprocess.run(
+ ["node", str(HARNESS), str(SEARCH_PAGE), json.dumps(cfg)],
+ capture_output=True, text=True)
+ assert proc.returncode == 0, proc.stderr
+ return json.loads(proc.stdout)
+
+
+def test_the_constants_are_what_these_scenarios_assume():
+ """
+ Read back, not restated: every number below is a multiple of these two, and
+ a change to either without a change here would leave the scenarios asserting
+ about a deadline the code no longer has.
+ """
+ code = SEARCH_PAGE.read_text(encoding="utf-8")
+ assert f"const SEARCH_STALL_MS = {STALL_MS};" in code
+ assert f"const SEARCH_MAX_MS = {CAP_MS};" in code
+
+
+def test_a_node_that_never_answers_still_fails_in_one_window():
+ """
+ The fan-out's cost, and the reason the repair is not a bigger number.
+
+ Nothing reports progress, so this is the whole budget a dead group spends —
+ and with batches of three it is what a page full of them costs.
+ """
+ out = _attempt(progressAt=[], resolveAt=None)
+ assert out["result"] == "failed"
+ assert out["at"] == STALL_MS
+ assert out["closed"] == 1, "the transport of a failed attempt must be closed"
+
+
+def test_a_slow_link_that_keeps_moving_is_not_cut_off():
+ """
+ The defect. ICE answers at 8 s, the channel opens at 16 s, the handshake
+ lands at 20 s — every step late, none of them stalled. The flat deadline
+ failed this at 10 s; the sidebar, with no deadline of its own, did not.
+ """
+ out = _attempt(progressAt=[8000, 16000], resolveAt=20000)
+ assert out["result"] == "connected"
+ assert out["at"] == 20000
+ assert out["closed"] == 0
+
+
+def test_a_node_that_answers_and_then_stops_fails_one_window_later():
+ """Progress buys a window, not an exemption."""
+ out = _attempt(progressAt=[8000], resolveAt=None)
+ assert out["result"] == "failed"
+ assert out["at"] == 8000 + STALL_MS
+
+
+def test_progress_that_never_finishes_is_bounded_by_the_cap():
+ """
+ Without the cap this is the connection that never ends: something reports
+ progress just often enough that the stall window never expires.
+ """
+ out = _attempt(progressAt=list(range(5000, 55000, 5000)), resolveAt=None)
+ assert out["result"] == "failed"
+ assert out["at"] == CAP_MS
+
+
+def test_the_walk_over_several_nodes_pays_one_window_each():
+ """
+ Two nodes, neither answering: the deadline is per attempt, so the group
+ costs two windows and both transports are closed. A deadline that leaked
+ between attempts would show up here as one window, or as an open transport.
+ """
+ out = _attempt(progressAt=[], resolveAt=None, nodes=2)
+ assert out["result"] == "failed"
+ assert out["at"] == 2 * STALL_MS
+ assert out["closed"] == 2
+
+
+def test_a_fast_connection_is_unaffected():
+ out = _attempt(progressAt=[], resolveAt=3000)
+ assert out["result"] == "connected"
+ assert out["at"] == 3000