aboutsummaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-10-07 11:50:48 +0200
committerChristophe Besson <cbesson@gmail.com>2026-10-07 11:50:48 +0200
commit00caeb3e10bfbf86892e7f31ad81f91ce8e50da7 (patch)
treecc9b9bf52377fb379a0eb41c14da50050b308ded
parent219093c823494bc8d760c99a7a55c4a23059abb5 (diff)
downloadmeshbay-00caeb3e10bfbf86892e7f31ad81f91ce8e50da7.tar.gz
fix(hub): let a download wait out a reconnect instead of failing
Chunks sent while the reconnect's connect() runs throw at once, and six retries 1.5 s apart ran out before the reconnect landed. Wait for it, up to two minutes, without spending retries. Follow-ups parked in §15.3. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
-rw-r--r--docs/MESHBAY_DESIGN.md4
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/file-utils.js45
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/transport.js2
-rw-r--r--packages/meshbay-hub/tests/test_downloads.py87
4 files changed, 131 insertions, 7 deletions
diff --git a/docs/MESHBAY_DESIGN.md b/docs/MESHBAY_DESIGN.md
index 8d6a18c..1730641 100644
--- a/docs/MESHBAY_DESIGN.md
+++ b/docs/MESHBAY_DESIGN.md
@@ -4051,6 +4051,10 @@ process runs it — `systemctl --user` on Linux, Task Scheduler on Windows.
| **A group key rotated from the Node page leaves the chat key where it was** | `ops.set_gek(rotate=True)` replaces the group key and opens no chat epoch; the MNP `gek_rotate` handler opened one itself (`_new_chat_epoch`), and it was the only door that did — but no client ever sent it, and it is gone since 6.0. The removals that matter (revoke, unpin, device revoke) open an epoch in `ops` for every door, so what is missing is the follow-through for an operator who rotates by hand: §4.5's "rotate after a removal" means it for chat too. The fix is the `_after_removal` shape — `open_chat_epoch` inside `ops.set_gek` when `rotated`. Found while removing the MNP message |
| **iPhone playback is untested, and the player ignores what ManagedMediaSource asks of it** | An iPhone has no `MediaSource`, only `ManagedMediaSource` (iOS 17.1+), which `video-player.js` now falls back to; below 17.1 nothing plays. Nobody has watched a film on one yet. The player listens for neither `startstreaming`/`endstreaming` nor `bufferedchange`, and the browser may evict buffered ranges on its own, ahead of the playhead included: the read-ahead (`pump()`, `currentRange()`) assumes a buffer only it shrinks. AirPlay is switched off on that path, because the source never opens otherwise |
| **The transcoded-seek test passes without transcoding** | `test_a_transcoded_video_keeps_accurate_seeking` (`test_stream_seek_audio_alignment.py`) forces the re-encode branch by swapping the module's `BROWSER_INCOMPATIBLE_VIDEO_CODECS`, then checks only that the result has no audio gap. The copy path also leaves no gap on that clip, so pointing the swap at a module the streaming code does not read still passes: the test cannot tell that the branch it is named after never ran. It should assert the re-encode happened (the `re-encoding` log line, or the encoder in the ffmpeg argv). Found by breaking the swap on purpose while moving the streaming code |
+| **A WebRTC session drops during many parallel downloads** | Seen once from the desktop client with two downloads running and seven queued: the channel closed, ICE went to `failed`, and the first reconnect attempts got `404 Node not connected` then `504 Node did not respond with WebRTC answer`, so the node had lost the hub too. Cause unknown: the node's log was not available |
+| **Any request sent during a reconnect's own connect() fails at once** | `_sendAndWait` skips `waitForReconnect` while `_inReconnectAttempt` is set, and that flag is transport-wide, not the handshake's: a chat send or an index request made in that window throws "DataChannel not open (state: connecting)". File chunks now wait the reconnect out in `_fetchChunkResilient`; nothing else does |
+| **A failed download leaves its in-flight chunks rejecting unhandled** | When `pipelinedDownload` throws, the other promises of its window are abandoned and each prints "Uncaught (in promise)" in the console. Noise, but in the log a dropped connection is read from |
+| **A queued transfer logs "no answer" every minute** | The lease watchdog re-asks after 60 s of silence whatever the last answer was, so a lease the node has answered `queued` and that has not moved prints "no answer for transfer ... asking again" once a minute for as long as it waits, which reads as a fault |
---
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/file-utils.js b/packages/meshbay-hub/src/meshbay_hub/static/file-utils.js
index 2ff2bd4..294455a 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/file-utils.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/file-utils.js
@@ -297,12 +297,23 @@ function _saveBlob(blob, filename) {
// problem, and the file already on disk (writable has real bytes in it by
// now) is worth more than an all-or-nothing download: retry the same chunk
// instead of letting one bad moment abort the whole transfer. Each retry
-// re-enters transport.fetchChunk, whose own _sendAndWait waits out an
-// in-flight reconnect before trying again, so this loop is mostly just
-// giving that reconnect the time and the attempts to land.
+// re-enters transport.fetchChunk. These attempts are for a request that
+// failed on a connection that is otherwise fine; a reconnect is waited out
+// separately (below) and costs none of them.
const CHUNK_RETRY_ATTEMPTS = 6;
const CHUNK_RETRY_DELAY_MS = 1500;
+// How long one chunk waits for a reconnect to land before its download fails.
+//
+// _sendAndWait does not wait while the reconnect's own connect() is running
+// (`_inReconnectAttempt`, which has to let the handshake through), so every
+// chunk asked for in that window throws "DataChannel not open (state:
+// connecting)" at once. Six attempts 1.5 s apart were gone in under ten
+// seconds, while one connect() can take longer than that on its own: found
+// live, a reconnect whose first attempt ended in a 504 failed the download
+// that was running, and the reconnect landed a moment later.
+const RECONNECT_WAIT_MS = 120000;
+
// How long one megabyte may take to reach the disk before we call it stuck.
//
// Every other await on this path is bounded and says so when it expires:
@@ -339,15 +350,34 @@ function _isRetryableTransportError(err) {
|| (err.message || '').startsWith('DataChannel not open');
}
-async function _fetchChunkResilient(transport, fileId, index, tr = '') {
+/** True once a reconnect that was in flight has landed. Gives up at
+ * `deadline`, or as soon as the transfer is cancelled. */
+async function _sitOutReconnect(transport, signal, deadline) {
+ while (transport.reconnecting && Date.now() < deadline
+ && !(signal && signal.aborted)) {
+ await transport.waitForReconnect(Math.min(250, deadline - Date.now()));
+ }
+ return !transport.reconnecting && transport.connected;
+}
+
+async function _fetchChunkResilient(transport, fileId, index, tr = '', signal = null) {
let lastErr;
- for (let attempt = 0; attempt < CHUNK_RETRY_ATTEMPTS; attempt++) {
+ const deadline = Date.now() + RECONNECT_WAIT_MS;
+ for (let attempt = 0; attempt < CHUNK_RETRY_ATTEMPTS;) {
try {
return await transport.fetchChunk(fileId, index, tr);
} catch (err) {
if (!_isRetryableTransportError(err)) throw err;
lastErr = err;
- if (attempt < CHUNK_RETRY_ATTEMPTS - 1) {
+ if (transport.reconnecting
+ && await _sitOutReconnect(transport, signal, deadline)) continue;
+ if (signal && signal.aborted) {
+ const abort = new Error('Cancelled');
+ abort.name = 'AbortError';
+ throw abort;
+ }
+ attempt++;
+ if (attempt < CHUNK_RETRY_ATTEMPTS) {
await new Promise((r) => setTimeout(r, CHUNK_RETRY_DELAY_MS));
}
}
@@ -364,7 +394,8 @@ async function pipelinedDownload(transport, gekKey, fileId, totalChunks, onChunk
const fire = () => {
while (nextSend < totalChunks && nextSend - nextRecv < PIPELINE_WINDOW) {
- inflight[nextSend] = _fetchChunkResilient(transport, fileId, nextSend, tr);
+ inflight[nextSend] = _fetchChunkResilient(transport, fileId, nextSend, tr,
+ signal);
nextSend++;
}
};
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport.js b/packages/meshbay-hub/src/meshbay_hub/static/transport.js
index 5ed6fa9..9dcea4f 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/transport.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/transport.js
@@ -621,6 +621,8 @@ class MeshBayTransport {
}
get connected() { return this._connected; }
+ /** True while _reconnectLoop is running, backoff included. */
+ get reconnecting() { return !!this._reconnectPromise; }
set onChat(fn) { this._onChat = fn; }
set onStreamInit(fn) { this._onStreamInit = fn; }
diff --git a/packages/meshbay-hub/tests/test_downloads.py b/packages/meshbay-hub/tests/test_downloads.py
index 7062d39..2a49a49 100644
--- a/packages/meshbay-hub/tests/test_downloads.py
+++ b/packages/meshbay-hub/tests/test_downloads.py
@@ -655,3 +655,90 @@ console.log(JSON.stringify({ worst }));
assert worst == 0, (
f"{worst} chunks already written were still held when a later one was — "
"the download keeps the whole file in memory until it ends")
+
+
+def _run_resilient(tmp_path, scenario):
+ """The real retry functions, lifted out and run against a transport stub
+ whose reconnect takes as long as the scenario says."""
+ src = (STATIC / "file-utils.js").read_text(encoding="utf-8")
+ fns = src[src.index("function _isRetryableTransportError"):
+ src.index("async function pipelinedDownload")]
+ script = tmp_path / "resilient.mjs"
+ script.write_text("""
+// The old budget, scaled down: six attempts, 10 ms apart, so a 300 ms
+// reconnect outlasts it by far, the way a real one outlasted 6 x 1.5 s.
+const CHUNK_RETRY_ATTEMPTS = 6;
+const CHUNK_RETRY_DELAY_MS = 10;
+const RECONNECT_WAIT_MS = 1000;
+const sleep = (ms) => new Promise((r) => setTimeout(r, ms));
+""" + fns + scenario, encoding="utf-8")
+ proc = subprocess.run(["node", str(script)], capture_output=True, text=True,
+ encoding="utf-8")
+ assert proc.returncode == 0, proc.stderr
+ return json.loads(proc.stdout)
+
+
+_RECONNECTING_TRANSPORT = """
+function transport(reconnectMs) {
+ const until = Date.now() + reconnectMs;
+ return {
+ calls: 0,
+ get reconnecting() { return Date.now() < until; },
+ get connected() { return Date.now() >= until; },
+ async waitForReconnect(ms) { await sleep(Math.min(ms, Math.max(0, until - Date.now()))); },
+ async fetchChunk() {
+ this.calls++;
+ // What _send throws while the reconnect's own connect() is running.
+ if (Date.now() < until) throw new Error('DataChannel not open (state: connecting)');
+ return { ct: 'chunk' };
+ },
+ };
+}
+"""
+
+
+def test_a_chunk_waits_out_a_reconnect_longer_than_its_retries(tmp_path):
+ """Found live: the reconnect's first attempt ended in a 504, the download
+ that was running failed with "DataChannel not open (state: connecting)",
+ and the connection came back a moment later."""
+ out = _run_resilient(tmp_path, _RECONNECTING_TRANSPORT + """
+const t = transport(300);
+const out = {};
+try { out.got = (await _fetchChunkResilient(t, 'f', 0)).ct; }
+catch (err) { out.err = err.message; }
+out.calls = t.calls;
+console.log(JSON.stringify(out));
+""")
+ assert out.get("got") == "chunk", out
+ # Waited, rather than spinning through attempts while the reconnect ran.
+ assert out["calls"] <= 3, out
+
+
+def test_a_reconnect_that_never_lands_still_fails_the_chunk(tmp_path):
+ out = _run_resilient(tmp_path, _RECONNECTING_TRANSPORT + """
+const t = transport(60000);
+const started = Date.now();
+const out = {};
+try { await _fetchChunkResilient(t, 'f', 0); out.err = 'none'; }
+catch (err) { out.err = err.message; }
+out.ms = Date.now() - started;
+console.log(JSON.stringify(out));
+""")
+ assert out["err"].startswith("DataChannel not open"), out
+ assert out["ms"] < 5000, f"gave up after {out['ms']} ms, past RECONNECT_WAIT_MS"
+
+
+def test_cancelling_during_a_reconnect_does_not_wait_for_it(tmp_path):
+ out = _run_resilient(tmp_path, _RECONNECTING_TRANSPORT + """
+const t = transport(60000);
+const signal = { aborted: false };
+setTimeout(() => { signal.aborted = true; }, 50);
+const started = Date.now();
+const out = {};
+try { await _fetchChunkResilient(t, 'f', 0, '', signal); out.err = 'none'; }
+catch (err) { out.err = err.name; }
+out.ms = Date.now() - started;
+console.log(JSON.stringify(out));
+""")
+ assert out["err"] == "AbortError", out
+ assert out["ms"] < 900, f"a cancel took {out['ms']} ms to be noticed"