From f73476d71f9b7268650c257db0eed3ab49b94849 Mon Sep 17 00:00:00 2001 From: Christophe Besson Date: Sat, 3 Oct 2026 10:29:59 +0200 Subject: fix(client): a seek's first segments are held while it lands, not dropped They are the new stream, header first. Dropped, a cast relay restarted at the landing got no ftyp/moov and the receiver gave up; they are now replayed in order once reinitAt/resumeAt is done. Co-Authored-By: Claude Opus 5.5 --- packages/meshbay-hub/tests/harness/window_leak.mjs | 7 ++- packages/meshbay-hub/tests/test_video_seek.py | 65 ++++++++++++++++++++++ 2 files changed, 70 insertions(+), 2 deletions(-) (limited to 'packages/meshbay-hub/tests') diff --git a/packages/meshbay-hub/tests/harness/window_leak.mjs b/packages/meshbay-hub/tests/harness/window_leak.mjs index d5185de..f05d8c0 100644 --- a/packages/meshbay-hub/tests/harness/window_leak.mjs +++ b/packages/meshbay-hub/tests/harness/window_leak.mjs @@ -41,20 +41,23 @@ let cancelled = false; const pump = () => {}; const flushQueue = () => {}; const gekRef = { current: null }; +// Not holding: these are the abandoned stream's segments, before any stream_init. +const heldRef = { current: null }; +const consumeSegment = async () => {}; const window_ = { MeshBayCrypto: { decryptChunkBin: async () => new Uint8Array(4) } }; const silentConsole = { log() {}, warn() {}, error() {} }; const handler = new Function( 'msg', 'cancelled', 'awaitingInitRef', 'outstandingRef', 'queueRef', - 'entry', 'gekRef', 'pump', 'flushQueue', 'console', 'window', + 'entry', 'gekRef', 'pump', 'flushQueue', 'console', 'window', 'heldRef', 'consumeSegment', `return (async () => {${source}})();`); const deliver = (n) => Promise.all( Array.from({ length: n }, (_, i) => handler( { file_id: entry.id, segment_index: i, nonce: 'n', ct: 'c' }, cancelled, awaitingInitRef, outstandingRef, queueRef, entry, gekRef, - pump, flushQueue, silentConsole, window_))); + pump, flushQueue, silentConsole, window_, heldRef, consumeSegment))); const run = async () => { // The whole window arrives while reinitAt is still awaiting its updateends. diff --git a/packages/meshbay-hub/tests/test_video_seek.py b/packages/meshbay-hub/tests/test_video_seek.py index 3289eda..3bcdb71 100644 --- a/packages/meshbay-hub/tests/test_video_seek.py +++ b/packages/meshbay-hub/tests/test_video_seek.py @@ -320,3 +320,68 @@ def test_the_resume_strings_exist_everywhere(locale): text = (STATIC / "locales" / f"{locale}.js").read_text(encoding="utf-8") for key in ("video.resumed_at", "video.from_start"): assert key in text, f"{locale} is missing {key}" + + +def test_a_seeks_first_segments_are_held_while_it_lands_not_dropped(tmp_path): + """What follows a `stream_init` is the new stream, its header first. + + The landing (`reinitAt`/`resumeAt`) waits on the SourceBuffer, and those + segments used to arrive in that gap and be dropped as the old film's. The + player never noticed — its SourceBuffer kept the header from the start of + the film — but a cast relay restarted at that landing was handed a stream + beginning at a moof, and the receiver gave up within a second. Measured on + a phone: two segments dropped one millisecond after `stream_init`, a relay + header of 0 bytes; with them held, a 2 KB header and the film on the TV. + + The real handler and the real drain run here: held while landing, counted + once, another file's segment refused, and replayed in arrival order. + """ + import json + import subprocess + + if shutil.which("node") is None: + pytest.skip("node is not available") + app = APP.read_text(encoding="utf-8") + start = app.index("transport.onStreamData = async (msg) => {") + handler = app[app.index("{", start) + 1:app.index("\n };", start)] + drain_at = app.index("land(msg.start || 0).then(async () => {") + drain = app[app.index("{", drain_at) + 1:app.index("}).catch(", drain_at)] + + script = tmp_path / "hold.mjs" + script.write_text( + f"const handlerSrc = {json.dumps(handler)}, drainSrc = {json.dumps(drain)};\n" + """ +const consumed = []; +const env = { + cancelled: false, entry: { id: 'film' }, + outstandingRef: { current: 8 }, awaitingInitRef: { current: true }, + heldRef: { current: [] }, holdGenRef: { current: 1 }, hold: 1, + consumeSegment: async (m) => { consumed.push(m.segment_index); }, + console: { log() {} }, +}; +const names = Object.keys(env); +const handler = new Function(...names, 'msg', `return (async () => {${handlerSrc}})();`); +const drain = new Function(...names, `return (async () => {${drainSrc}})();`); +const call = (fn, msg) => fn(...names.map((k) => env[k]), msg); + +// Landing: the new stream arrives, plus one stray segment from another file. +for (const [i, file] of [[0, 'film'], [1, 'film'], [2, 'other'], [3, 'film']]) { + await call(handler, { file_id: file, segment_index: i }); +} +const held = env.heldRef.current.map((m) => m.segment_index); +const consumedWhileLanding = consumed.length; +const outstanding = env.outstandingRef.current; +// The landing completes. +env.awaitingInitRef.current = false; +await call(drain); +console.log(JSON.stringify({ held, consumedWhileLanding, outstanding, consumed, + heldAfter: env.heldRef.current })); +""", encoding="utf-8") + proc = subprocess.run(["node", str(script)], capture_output=True, text=True, + encoding="utf-8", timeout=60) + assert proc.returncode == 0, proc.stderr + out = json.loads(proc.stdout) + assert out["held"] == [0, 1, 3], "the new stream's segments were dropped or mixed" + assert out["consumedWhileLanding"] == 0, "a segment was taken before the landing finished" + assert out["outstanding"] == 4, "each arrival is counted once, held or not" + assert out["consumed"] == [0, 1, 3], "the replay must keep arrival order, header first" + assert out["heldAfter"] is None, "holding must stop once the landing has drained" -- cgit v1.2.3