aboutsummaryrefslogtreecommitdiffstats
path: root/packages
diff options
context:
space:
mode:
Diffstat (limited to 'packages')
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/static/video-player.js37
-rw-r--r--packages/meshbay-hub/tests/harness/window_leak.mjs7
-rw-r--r--packages/meshbay-hub/tests/test_video_seek.py65
3 files changed, 106 insertions, 3 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/static/video-player.js b/packages/meshbay-hub/src/meshbay_hub/static/video-player.js
index 2a4d2ea..4fa652b 100644
--- a/packages/meshbay-hub/src/meshbay_hub/static/video-player.js
+++ b/packages/meshbay-hub/src/meshbay_hub/static/video-player.js
@@ -438,6 +438,10 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose, onDownload }) {
const castRestartPendingRef = useRef(false);
const castDeviceRef = useRef(null);
const castRestartGenRef = useRef(0);
+ // Segments of the *new* stream that arrive while its seek is still landing
+ // (`reinitAt`/`resumeAt` wait on the SourceBuffer). Null when not holding.
+ const heldRef = useRef(null);
+ const holdGenRef = useRef(0);
const landingPlayheadRef = useRef(false);
// The current Screen Wake Lock sentinel, if the browser granted one — see
// the effect below. Null on any platform/context that does not support it,
@@ -764,6 +768,7 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose, onDownload }) {
ranges: describeRanges(),
});
awaitingInitRef.current = true;
+ heldRef.current = null; holdGenRef.current++;
// A seek supersedes a resume that had been asked for and not landed:
// this one empties the buffer, and taking the resume branch on its
// `stream_init` would leave the film we navigated away from in place.
@@ -798,6 +803,7 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose, onDownload }) {
});
resumingRef.current = true;
awaitingInitRef.current = true;
+ heldRef.current = null; holdGenRef.current++;
seekTargetRef.current = null;
outstandingRef.current = STREAM_WINDOW;
console.log('[resume] request', +target.toFixed(1));
@@ -1064,7 +1070,23 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose, onDownload }) {
castRestartPendingRef.current = true;
}
const land = resuming ? resumeAt : reinitAt;
- land(msg.start || 0).catch(() => {
+ // What follows this message on the ordered channel is the new
+ // stream, its header first. The landing waits on the SourceBuffer,
+ // and those segments used to arrive in that gap and be dropped as
+ // the old film's: the player kept its header from the start of the
+ // film and never noticed, but a cast relay restarted here got a
+ // stream beginning at a moof — no ftyp, no moov — and the receiver
+ // gave up on it within a second (measured on a phone). So they are
+ // held, and taken in order once the landing is done.
+ const hold = ++holdGenRef.current;
+ heldRef.current = [];
+ land(msg.start || 0).then(async () => {
+ while (holdGenRef.current === hold && heldRef.current && heldRef.current.length) {
+ await consumeSegment(heldRef.current.shift());
+ }
+ if (holdGenRef.current === hold) heldRef.current = null;
+ }).catch(() => {
+ if (holdGenRef.current === hold) heldRef.current = null;
setError(t('video.err_transport'));
setPhase('error');
});
@@ -1178,6 +1200,12 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose, onDownload }) {
// for credit that cannot come. A race, which is why the same seek
// worked twice and hung on the third.
outstandingRef.current = Math.max(0, outstandingRef.current - 1);
+ // The new stream, arriving while its seek lands: kept, not dropped.
+ // Counted above already, so the replay must not count it again.
+ if (heldRef.current) {
+ if (!msg.file_id || msg.file_id === entry.id) heldRef.current.push(msg);
+ return;
+ }
if (awaitingInitRef.current) {
console.log('[seek] dropping segment (awaitingInit), outstanding:', outstandingRef.current);
}
@@ -1190,6 +1218,13 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose, onDownload }) {
// would otherwise be decrypted against the wrong file — which fails,
// loudly, in the console, for something that is simply not ours.
if (msg.file_id && msg.file_id !== entry.id) return;
+ await consumeSegment(msg);
+ };
+
+ // Everything a segment goes through once it is known to belong to the
+ // stream being played: by arrival, or replayed after a landing.
+ const consumeSegment = async (msg) => {
+ if (cancelled) return;
try {
const plaintext = await window.MeshBayCrypto.decryptChunkBin(
gekRef.current, entry.id, msg.segment_index, msg.nonce, msg.ct);
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"