From 373f83236816aaff536dfa2f9589801acac44012 Mon Sep 17 00:00:00 2001 From: Christophe Besson Date: Sun, 16 Aug 2026 20:57:41 +0200 Subject: fix(hub): bound the video read-ahead by the playhead, not by the network MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A 500 MB film loaded about 100 MB and hung on "buffering" for good. 100 MB is not a number in our code: it is where the browser stops. ffmpeg remuxes with `-c copy`, so the bytes on the wire are the file's own, and credit granted per append meant taking them as fast as the network allowed — which for a film is very much faster than watching it. The SourceBuffer ceiling arrived in the first minute. Past it every append was refused, and the refusal was unrecoverable: a refused append fires no `updateend`, `updateend` was where credit was granted, so the node sent nothing and no segment arrived to retry the append. Every wakeup the pipeline had was downstream of the append that had just failed. Playback continuing — the one thing that frees room — woke nothing at all. Credit now follows the buffer instead of the writes. `pump()` is the only place it is granted, it keeps `STREAM_WINDOW` segments in flight while less than `BUFFER_AHEAD_S` of film is held past the playhead, and it is driven by a one-second clock and by playback progress, never by arriving data. Buffering by time makes a two-hour film cost what a two-minute clip costs. A window rather than a debt, and this took a second measurement to get right: accumulating a credit per append and releasing the balance when the buffer finally drained sent six megabytes in one burst, overshot by a minute of film, then said nothing for forty-six seconds. Measured in Chrome against real fragmented MP4. Two smaller things found on the way. `updateend` fires for `remove()` as well as `appendBuffer()`, so crediting from it paid the node for the player's own evictions. And a viewer that is deliberately far enough ahead grants nothing for minutes, which the node read as a closed tab — it now sends `stream_more` with n=0, which grants no room but proves someone is there. The first version of the test modelled the credit loop and passed while the player still hung: a model written by whoever wrote the fix agrees with it by construction. `tests/harness/mse_harness.mjs` lifts the real functions out of app.js as text and runs them against a SourceBuffer that has a ceiling. What is modelled is the browser. --- packages/meshbay-hub/src/meshbay_hub/static/app.js | 171 ++++++++++++++- .../src/meshbay_hub/static/transport.js | 12 ++ packages/meshbay-hub/tests/harness/mse_harness.mjs | 159 ++++++++++++++ .../meshbay-hub/tests/test_video_buffer_ceiling.py | 233 +++++++++++++++++++++ 4 files changed, 565 insertions(+), 10 deletions(-) create mode 100644 packages/meshbay-hub/tests/harness/mse_harness.mjs create mode 100644 packages/meshbay-hub/tests/test_video_buffer_ceiling.py diff --git a/packages/meshbay-hub/src/meshbay_hub/static/app.js b/packages/meshbay-hub/src/meshbay_hub/static/app.js index 7968284..a9559c6 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/app.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/app.js @@ -986,6 +986,24 @@ const CHUNK_SIZE = 1024 * 1024; // Seconds of already-watched video kept in the SourceBuffer, and the queue depth // past which we start making room before being forced to. const BUFFER_BEHIND_S = 60; +// How far past the playhead we are willing to pull. The browser caps a video +// SourceBuffer at a few hundred megabytes and refuses the append that goes +// past, so "as fast as the network allows" is not a strategy for a film: the +// node remuxes with `-c copy`, so a 500 MB file puts 500 MB on the wire, and a +// ten-megabit second fills the ceiling in the first minute. Buffering by time +// rather than by bytes keeps a two-hour film and a two-minute clip alike. +const BUFFER_AHEAD_S = 90; +// While we deliberately hold credit back, the node must still hear from us: its +// own stall timeout is two minutes, and a paused film is not a gone viewer. +const CREDIT_KEEPALIVE_MS = 20000; +// Segments allowed in flight while there is room to put them. This is a window, +// topped up as segments land, and not a debt released in one go: accumulating a +// credit per append and handing the lot over when the buffer finally had room +// sent 6 MB in a burst, overshot the target by a minute of film, and then said +// nothing for the next forty-six seconds. Measured in Chrome against real +// fragmented MP4. A stream that arrives in gulps has no margin for a network +// that hesitates, and looks like a hang while it is quiet. +const STREAM_WINDOW = 8; const QUEUE_HIGH_WATER = 12; const PIPELINE_WINDOW = 8; @@ -2666,6 +2684,14 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose }) { const appendingRef = useRef(false); const endedRef = useRef(false); const durationRef = useRef(0); + // Segments the node is allowed to have in flight but has not sent yet, and + // when we last said anything to it at all. + const outstandingRef = useRef(0); + const lastPokeRef = useRef(0); + // Diagnostics reported to the node: how many appends the browser refused for + // want of room, and whether the element itself says it is starved. + const quotaRef = useRef(0); + const stalledRef = useRef(false); /** * Drop what has already been watched. @@ -2689,6 +2715,18 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose }) { } }, []); + /** Seconds of film held past the playhead. */ + const bufferedAhead = useCallback(() => { + const sb = sbRef.current; + const v = videoRef.current; + if (!sb || !v || !sb.buffered.length) return 0; + try { + return sb.buffered.end(sb.buffered.length - 1) - v.currentTime; + } catch { + return 0; + } + }, []); + const flushQueue = useCallback(() => { const sb = sbRef.current; if (!sb || appendingRef.current || sb.updating) return; @@ -2706,6 +2744,7 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose }) { } catch (e) { appendingRef.current = false; if (e.name === 'QuotaExceededError') { + quotaRef.current += 1; // The segment stays at the head of the queue and is tried again once // there is room. Dropping it — which is what this used to do — leaves a // hole in the middle of the film and no error anywhere. @@ -2719,6 +2758,51 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose }) { } }, [evictBehind]); + /** + * Decide whether the node may send more, and keep the pipeline moving. + * + * This is the only place credit is granted, and the only thing that can + * restart a pipeline the buffer ceiling has stopped. That second job is why + * it exists: an append refused for quota fires no `updateend`, so it grants + * no credit, so the node sends nothing, so no segment arrives to call + * `flushQueue` again. Every wakeup the append path had was downstream of the + * append that just failed — the player deadlocked against itself and sat on + * "buffering" for good, which is what a 500 MB film did at around 100 MB. + * + * So the clock drives this, not the data. + */ + const pump = useCallback(() => { + const transport = transportRef.current; + evictBehind(); + flushQueue(); + if (endedRef.current && queueRef.current.length === 0) return; + if (!transport || !transport.connected) return; + + if (bufferedAhead() > BUFFER_AHEAD_S + || queueRef.current.length > QUEUE_HIGH_WATER) { + // Far enough ahead. Grant nothing, but do not go silent: two minutes of + // silence is how the node decides nobody is watching, and pausing a film + // for two minutes is an ordinary thing to do. + const now = Date.now(); + if (now - lastPokeRef.current > CREDIT_KEEPALIVE_MS) { + lastPokeRef.current = now; + transport.grantStreamCredit(0); + } + return; + } + + // Top the window back up to what is allowed in flight, rather than paying + // off everything owed at once. Called on every arriving segment as well as + // on the clock, so credit trickles out as room appears instead of being + // released in one gulp when the buffer finally drains. + const room = STREAM_WINDOW - outstandingRef.current; + if (room > 0) { + outstandingRef.current += room; + lastPokeRef.current = Date.now(); + transport.grantStreamCredit(room); + } + }, [evictBehind, flushQueue, bufferedAhead]); + useEffect(() => { let cancelled = false; // Reset here, not in the teardown of the run before: switching video while @@ -2731,6 +2815,10 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose }) { appendingRef.current = false; endedRef.current = false; queueRef.current = []; + outstandingRef.current = 0; + lastPokeRef.current = Date.now(); + quotaRef.current = 0; + stalledRef.current = false; const transport = transportRef.current; if (!transport || !transport.connected) { setError(t('video.err_transport')); @@ -2738,6 +2826,9 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose }) { return; } + const onStarved = () => { stalledRef.current = true; pump(); }; + const onFed = () => { stalledRef.current = false; }; + const onSeeking = () => { const v = videoRef.current; if (!v || !v.buffered.length) return; @@ -2784,14 +2875,13 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose }) { sbRef.current = sb; sb.mode = 'sequence'; sb.addEventListener('updateend', () => { + // No credit is granted here, deliberately. Appending is not the + // same question as having room, and tying the two meant `remove()` + // — which fires this event too — paid the node for the player's own + // evictions. What may be in flight is decided from the buffer, in + // pump(), and nowhere else. appendingRef.current = false; - // One segment consumed, so the node may send one more. The credit - // is granted here, after the append, because this is the point at - // which the memory is genuinely free again. - const transport = transportRef.current; - if (transport && transport.connected) transport.grantStreamCredit(1); - if (queueRef.current.length > QUEUE_HIGH_WATER) evictBehind(); - flushQueue(); + pump(); }); setPhase('streaming'); flushQueue(); @@ -2800,6 +2890,13 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose }) { if (videoRef.current) { videoRef.current.src = url; videoRef.current.addEventListener('seeking', onSeeking); + videoRef.current.addEventListener('timeupdate', pump); + // The element's own verdict. "buffering" on screen is this, and it + // is the one thing the node cannot infer from a stream it is feeding. + videoRef.current.addEventListener('waiting', onStarved); + videoRef.current.addEventListener('stalled', onStarved); + videoRef.current.addEventListener('playing', onFed); + videoRef.current.addEventListener('canplay', onFed); } }; @@ -2810,11 +2907,18 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose }) { // 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; + // One of the in-flight segments landed, whatever becomes of it: the + // window has room again. Counted before the decrypt so that a segment + // we fail to read still frees its slot rather than shrinking the + // window by one for the rest of the film. + outstandingRef.current = Math.max(0, outstandingRef.current - 1); try { const plaintext = await window.MeshBayCrypto.decryptChunkBin( gekRef.current, entry.id, msg.segment_index, msg.nonce, msg.ct); queueRef.current.push(plaintext); - flushQueue(); + // pump(), not flushQueue(): arriving data is the moment to top the + // window back up, and that is what keeps the stream continuous. + pump(); } catch (e) { console.error('[MSE] decrypt error:', e); } @@ -2828,7 +2932,12 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose }) { flushQueue(); }; - transport.requestStream(entry.id); + // The opening window, and the count that tracks it. Asking for more here + // than pump() maintains would leave the node holding credit this side + // does not know about, which is the whole window's worth of overshoot on + // the very first breath of the stream. + outstandingRef.current = STREAM_WINDOW; + transport.requestStream(entry.id, STREAM_WINDOW); }; // Closing the tab, or backgrounding it on a phone, never runs a React @@ -2852,16 +2961,58 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose }) { window.addEventListener('pagehide', onPageHide); document.addEventListener('visibilitychange', onVisibility); + // `timeupdate` is silent while the film is paused, and the append path + // cannot wake itself once the ceiling has refused a segment. This is the + // clock that guarantees something is still driving the pipeline. + const pumpTimer = setInterval(pump, 1000); + + // What the player sees, into the node's log. A hang on a phone shows the + // node feeding a stream quite happily; the half that says otherwise is in + // here, and there is no console to read it from. + const diagTimer = setInterval(() => { + const v = videoRef.current, sb = sbRef.current; + const t = transportRef.current; + if (!t || !v) return; + let ranges = ''; + try { + for (let i = 0; sb && i < sb.buffered.length; i++) { + ranges += `${sb.buffered.start(i).toFixed(0)}-${sb.buffered.end(i).toFixed(0)} `; + } + } catch { ranges = '?'; } + t.sendStreamDiag({ + t: +v.currentTime.toFixed(1), + ahead: +bufferedAhead().toFixed(1), + ranges: ranges.trim(), + ready: v.readyState, // 0 = nothing, 4 = enough to play through + paused: v.paused, + stalled: stalledRef.current, + q: queueRef.current.length, + inflight: outstandingRef.current, + appending: appendingRef.current, + updating: sb ? sb.updating : null, + quota: quotaRef.current, + ms: msRef.current ? msRef.current.readyState : null, + err: v.error ? `${v.error.code}:${v.error.message}` : null, + }); + }, 5000); + startStream().catch(err => { if (!cancelled) { setError(err.message); setPhase('error'); } }); return () => { cancelled = true; + clearInterval(pumpTimer); + clearInterval(diagTimer); window.removeEventListener('pagehide', onPageHide); document.removeEventListener('visibilitychange', onVisibility); if (videoRef.current) { videoRef.current.removeEventListener('seeking', onSeeking); + videoRef.current.removeEventListener('timeupdate', pump); + videoRef.current.removeEventListener('waiting', onStarved); + videoRef.current.removeEventListener('stalled', onStarved); + videoRef.current.removeEventListener('playing', onFed); + videoRef.current.removeEventListener('canplay', onFed); } if (transport) { // Tell the node first: dropping the handlers only makes us deaf, and a @@ -2885,7 +3036,7 @@ function VideoPlayer({ entry, transportRef, gekRef, onClose }) { sbRef.current = null; msRef.current = null; }; - }, [entry, flushQueue]); + }, [entry, flushQueue, pump]); useEffect(() => { if (phase === 'streaming' && videoRef.current) { diff --git a/packages/meshbay-hub/src/meshbay_hub/static/transport.js b/packages/meshbay-hub/src/meshbay_hub/static/transport.js index 7a3943e..1b9abb3 100644 --- a/packages/meshbay-hub/src/meshbay_hub/static/transport.js +++ b/packages/meshbay-hub/src/meshbay_hub/static/transport.js @@ -589,6 +589,18 @@ class MeshBayTransport { this._send({ type: 'stream_more', v: '0.1', n }); } + /** + * Tell the node what the player sees. + * + * A hang on a phone is unreadable from here: there is no console to open and + * the node's own log shows a stream it is feeding perfectly well. This puts + * the two halves in one file. The node only logs it. + */ + sendStreamDiag(diag) { + if (!this._connected) return; + try { this._send({ type: 'client_diag', v: '0.1', ...diag }); } catch { /* gone */ } + } + /** * Nobody is watching any more. * diff --git a/packages/meshbay-hub/tests/harness/mse_harness.mjs b/packages/meshbay-hub/tests/harness/mse_harness.mjs new file mode 100644 index 0000000..0daacfb --- /dev/null +++ b/packages/meshbay-hub/tests/harness/mse_harness.mjs @@ -0,0 +1,159 @@ +// Run the SHIPPED player functions against a fake SourceBuffer. +// +// The point is that nothing here is a paraphrase of app.js: `bufferedAhead`, +// `evictBehind`, `flushQueue` and `pump` are lifted out of the file as text and +// executed. A model of a fix, written by whoever wrote the fix, agrees with it +// by construction — which is how a passing test sat next to a player that still +// hung. What is modelled here is the *browser*: a SourceBuffer with a ceiling, +// and `updateend` firing for removals as well as appends. +// +// Usage: node mse_harness.mjs +import { readFileSync } from 'fs'; + +const app = readFileSync(process.argv[2], 'utf8'); +const cfg = JSON.parse(process.argv[3] || '{}'); + +const { + playing = false, // does the viewer actually press play + capMB = 100, // where the browser refuses the append + fileMB = 493.5, // the film, from a real upload + durationS = 3936, + netMBs = 35, // measured node throughput + wallS = 600, +} = cfg; + +const grab = (name) => { + const start = app.indexOf(`const ${name} = useCallback(`); + if (start < 0) throw new Error(`${name} not found in app.js`); + const deps = app.indexOf('\n }, [', start); + const end = app.indexOf(');', deps) + 2; + return app.slice(start, end); +}; + +const useCallback = (fn) => fn; +const src = ['bufferedAhead', 'evictBehind', 'flushQueue', 'pump'] + .map(grab).join('\n'); + +const SEG = 256 * 1024; +const CAP = capMB * 1048576; +const BITRATE = fileMB * 1048576 / durationS; + +// Read from app.js too, so a change to the constants is a change to the test. +const constOf = (name) => { + const m = app.match(new RegExp(`const ${name} = (\\d+)`)); + if (!m) throw new Error(`${name} not found`); + return Number(m[1]); +}; +const BUFFER_BEHIND_S = constOf('BUFFER_BEHIND_S'); +const STREAM_WINDOW = constOf('STREAM_WINDOW'); +const BUFFER_AHEAD_S = constOf('BUFFER_AHEAD_S'); +const QUEUE_HIGH_WATER = constOf('QUEUE_HIGH_WATER'); +const CREDIT_KEEPALIVE_MS = constOf('CREDIT_KEEPALIVE_MS'); + +let bytes = 0, ranges = [], appended = 0, removes = 0; +let granted = 0, keepalives = 0, sent = 0, credit = 0, quotaRefusals = 0; + +const sb = { + updating: false, + get buffered() { + return { + get length() { return ranges.length; }, + start: (i) => ranges[i][0], + end: (i) => ranges[i][1], + }; + }, + appendBuffer(chunk) { + if (bytes + chunk.byteLength > CAP) { + quotaRefusals++; + const e = new Error('quota'); e.name = 'QuotaExceededError'; throw e; + } + const at = ranges.length ? ranges[ranges.length - 1][1] : 0; + ranges.push([at, at + chunk.byteLength / BITRATE]); + bytes += chunk.byteLength; + appended++; + }, + remove(a, b) { + removes++; + let dropped = 0; + ranges = ranges.filter(([s, e]) => { + if (e <= b && s >= a) { dropped += (e - s) * BITRATE; return false; } + return true; + }); + bytes -= dropped; + // A real remove() is asynchronous and fires updateend when it lands. That + // event is indistinguishable from an append's unless the player kept track. + pendingRemoveEvents++; + }, +}; +let pendingRemoveEvents = 0; + +const video = { currentTime: 0 }; +const sbRef = { current: sb }, videoRef = { current: video }; +const msRef = { current: { readyState: 'open', endOfStream() {} } }; +const queueRef = { current: [] }; +const appendingRef = { current: false }, endedRef = { current: false }; +const outstandingRef = { current: 0 }, lastPokeRef = { current: 0 }; +const quotaRef = { current: 0 }, stalledRef = { current: false }; +const transportRef = { + current: { + connected: true, + grantStreamCredit(n) { + if (n === 0) { keepalives++; return; } + granted += n; credit += n; + }, + }, +}; + +const fns = new Function( + 'sbRef,videoRef,msRef,queueRef,appendingRef,endedRef,outstandingRef,' + + 'lastPokeRef,transportRef,BUFFER_BEHIND_S,BUFFER_AHEAD_S,QUEUE_HIGH_WATER,' + + 'CREDIT_KEEPALIVE_MS,STREAM_WINDOW,quotaRef,stalledRef,console,useCallback', + src + '\n return {bufferedAhead, evictBehind, flushQueue, pump};' +)(sbRef, videoRef, msRef, queueRef, appendingRef, endedRef, outstandingRef, + lastPokeRef, transportRef, BUFFER_BEHIND_S, BUFFER_AHEAD_S, QUEUE_HIGH_WATER, + CREDIT_KEEPALIVE_MS, STREAM_WINDOW, quotaRef, stalledRef, console, useCallback); + +// The player's own `updateend` listener, transcribed — the one part of the +// component that is a listener rather than a callback, and the place the +// remove/append distinction lives. +const updateend = () => { + appendingRef.current = false; + fns.pump(); +}; + +credit = STREAM_WINDOW; +outstandingRef.current = STREAM_WINDOW; +let wall = 0; +const TICK = 0.05; +while (wall < wallS) { + wall += TICK; + if (playing) { + const end = ranges.length ? ranges[ranges.length - 1][1] : 0; + video.currentTime = Math.min(video.currentTime + TICK, end); + } + fns.pump(); // the 1 s timer and timeupdate + while (pendingRemoveEvents > 0) { pendingRemoveEvents--; updateend(); } + + let budget = netMBs * 1048576 * TICK; + while (credit > 0 && budget >= SEG) { + credit--; budget -= SEG; sent++; + outstandingRef.current = Math.max(0, outstandingRef.current - 1); + queueRef.current.push({ byteLength: SEG }); + fns.flushQueue(); + if (appendingRef.current) updateend(); + while (pendingRemoveEvents > 0) { pendingRemoveEvents--; updateend(); } + } +} + +console.log(JSON.stringify({ + sentMB: +(sent * SEG / 1048576).toFixed(1), + heldInBufferMB: +(bytes / 1048576).toFixed(1), + queueDepth: queueRef.current.length, + bufferedAheadS: +fns.bufferedAhead().toFixed(1), + watchedS: +video.currentTime.toFixed(1), + grants: granted, + keepalives, + removes, + quotaRefusals: quotaRef.current, + hitCeiling: bytes >= CAP * 0.99, +})); diff --git a/packages/meshbay-hub/tests/test_video_buffer_ceiling.py b/packages/meshbay-hub/tests/test_video_buffer_ceiling.py new file mode 100644 index 0000000..9063c46 --- /dev/null +++ b/packages/meshbay-hub/tests/test_video_buffer_ceiling.py @@ -0,0 +1,233 @@ +""" +A big film stops at around 100 MB. + +Reported against 0.5: upload a 500 MB video, play it, and the player loads +roughly 100 MB and then hangs on "buffering" for good. + +100 MB is not a number in our code. It is where the browser stops: a video +SourceBuffer is capped at a few hundred megabytes and `appendBuffer` throws +QuotaExceededError past it. The node remuxes with `-c copy`, so the bytes on +the wire are the file's own — a 500 MB film really does try to put 500 MB into +that buffer, and on a fast link it reaches the ceiling in the first minute, +long before anyone has watched enough for eviction to have anything to drop. + +Two defects, and the second is the one that makes it permanent. + +**Nothing bounded how far ahead we pulled.** Credit was granted once per +append: the node sent exactly as fast as the browser could append, which is as +fast as the network allows, which for a film is very much faster than watching +it. Memory was bounded by the browser's ceiling rather than by anything we +chose. Buffering by *time* past the playhead instead makes a two-hour film cost +the same as a two-minute clip. + +**The pipeline could not restart itself.** An append refused for quota fires no +`updateend`. `updateend` was where credit was granted, so no credit went out; +the node then sent nothing, so no segment arrived to call `flushQueue` again. +Every wakeup the append path had was downstream of the append that had just +failed. Playback continuing past the segment, which is exactly what frees the +room needed to recover, woke nothing at all. The player deadlocked against +itself. + +The first version of these tests modelled the pipeline and passed while the +player still hung, because a model of a fix written by whoever wrote the fix +agrees with it by construction. They now run the shipped `bufferedAhead`, +`evictBehind`, `flushQueue` and `pump`, lifted out of app.js as text, against a +fake SourceBuffer — see harness/mse_harness.mjs. What is modelled is the +browser, not us. +""" + +import json +import re +import shutil +import subprocess +from pathlib import Path + +import pytest + +STATIC = Path(__file__).resolve().parents[1] / "src" / "meshbay_hub" / "static" +APP = STATIC / "app.js" +NODE_SERVER = (Path(__file__).resolve().parents[2] / "meshbay-node" / "src" + / "meshbay_node" / "transport" / "webrtc_server.py") + +pytestmark = pytest.mark.skipif( + shutil.which("node") is None or not APP.exists(), + reason="node or the SPA sources are not available") + + +@pytest.fixture(scope="module") +def app(): + return APP.read_text() + + +def _player(app: str) -> str: + i = app.index("function VideoPlayer(") + return app[i:app.index("\nfunction ", i + 1)] + + +# ── The shipped functions, run against a browser that has a ceiling ─────────── + +HARNESS = Path(__file__).parent / "harness" / "mse_harness.mjs" + + +def _harness(**cfg) -> dict: + proc = subprocess.run( + ["node", str(HARNESS), str(APP), json.dumps(cfg)], + capture_output=True, text=True) + assert proc.returncode == 0, proc.stderr + return json.loads(proc.stdout) + + +@pytest.fixture(scope="module") +def idle(): + """Nobody pressed play — autoplay is blocked on a phone more often than not.""" + return _harness(playing=False) + + +@pytest.fixture(scope="module") +def watched(): + return _harness(playing=True) + + +def test_the_ceiling_is_never_reached_when_nobody_presses_play(idle): + """The reported hang, from the side that produces it. + + A film left on the loading screen used to pull until the browser refused an + append, and that refusal was unrecoverable. Nothing should get near it. + """ + assert not idle["hitCeiling"], ( + f"filled the buffer to the ceiling ({idle['heldInBufferMB']} MB)") + assert idle["quotaRefusals"] == 0, ( + f"{idle['quotaRefusals']} appends refused for quota — the state the " + "player cannot get out of on its own") + + +def test_the_ceiling_is_never_reached_while_watching(watched): + assert not watched["hitCeiling"], ( + f"filled the buffer to the ceiling ({watched['heldInBufferMB']} MB)") + assert watched["quotaRefusals"] == 0 + + +def test_the_read_ahead_is_bounded_by_the_playhead(idle, watched): + """What replaced "as fast as the network allows".""" + ahead = int(re.search(r"const BUFFER_AHEAD_S = (\d+)", APP.read_text()).group(1)) + for name, run in (("idle", idle), ("watching", watched)): + assert run["bufferedAheadS"] < ahead * 2, ( + f"{name}: {run['bufferedAheadS']}s buffered against a {ahead}s " + "bound — the gate is not holding") + + +def test_a_watched_film_keeps_being_fed(watched): + """The gate must throttle the stream, not stop it. + + Holding credit for good would be just as broken as never holding it, and + would look the same from the sofa. + """ + assert watched["watchedS"] > 500, ( + "playback did not advance, so this run says nothing about throttling") + assert watched["sentMB"] > 40, ( + f"only {watched['sentMB']} MB reached the player in ten minutes of " + "playback — the gate is holding credit it should have released") + assert watched["removes"] > 0, "nothing was ever evicted behind the playhead" + + +def test_memory_stays_bounded_over_a_long_watch(watched): + assert watched["heldInBufferMB"] < 40, ( + f"holding {watched['heldInBufferMB']} MB — eviction is not keeping up") + + +# ── The shape the fix depends on ────────────────────────────────────────────── + +def test_credit_is_granted_in_exactly_one_place(app): + """Granting from `updateend` is the deadlock. It must not come back. + + A second grant site is how this regresses: it would work, until the append + it hangs off is the one the ceiling refuses. + """ + player = _player(app) + sites = player.count("grantStreamCredit(") + assert sites == 2, ( + f"{sites} calls to grantStreamCredit — expected exactly two, both " + "inside pump(): the keepalive and the release") + pump = player[player.index("const pump = useCallback("):] + pump = pump[:pump.index("\n }, [")] + assert pump.count("grantStreamCredit(") == 2, ( + "credit is granted outside pump(), so an append that is refused for " + "quota can still be the only thing that would have woken the pipeline") + + +def test_something_other_than_data_drives_the_pipeline(app): + """The recovery path cannot depend on a segment arriving.""" + player = _player(app) + assert "setInterval(pump" in player, ( + "no clock drives pump(): once the ceiling refuses an append, nothing " + "arrives and nothing retries") + assert "clearInterval(pumpTimer)" in player, "the pump timer outlives the player" + assert "addEventListener('timeupdate', pump)" in player, ( + "playback progress is what frees room to evict, and it wakes nothing") + + +def test_a_buffered_viewer_still_tells_the_node_it_is_there(app): + """Holding credit back must not read as a closed tab.""" + player = _player(app) + assert "grantStreamCredit(0)" in player, ( + "a viewer that is far enough ahead grants nothing and says nothing, so " + "the node's stall timeout ends a film that is merely paused") + + +@pytest.mark.skipif(not NODE_SERVER.exists(), reason="node sources unavailable") +def test_the_node_ends_a_stream_on_silence_not_on_stinginess(): + """The other half of the keepalive: the node has to honour it.""" + text = NODE_SERVER.read_text() + i = text.index("async def _await_stream_credit") + body = text[i:text.index("\n async def ", i + 1)] + assert "self._stream_heard_at" in body, ( + "the stall budget still accumulates over the whole wait, so a keepalive " + "that grants no credit cannot keep a paused film alive") + assert "waited += STREAM_CREDIT_POLL" not in body, ( + "the budget still accumulates over the whole wait rather than being " + "measured from the last thing the peer said") + grant = text[text.index("def _grant_stream_credit"):] + grant = grant[:grant.index("\n def ", 1)] + assert "self._stream_heard_at = time.monotonic()" in grant, ( + "n=0 does not refresh the timeout, so the keepalive is a no-op") + + +def test_appending_does_not_earn_credit(app): + """What may be in flight is a question about the buffer, not about appends. + + Tying the two was the original design and it was wrong twice over. + `updateend` fires for `remove()` as well, so the player paid the node for + its own evictions; and crediting per append meant taking segments as fast + as they could be written, which is as fast as the network allows. + """ + player = _player(app) + handler = player[player.index("sb.addEventListener('updateend'"):] + handler = handler[:handler.index("\n });")] + assert "grantStreamCredit" not in handler, ( + "credit is granted from updateend, which fires for remove() too") + assert "outstandingRef" not in handler, ( + "the in-flight window is adjusted from updateend rather than from the " + "buffer, so an eviction still counts as room for another segment") + + +def test_credit_is_a_window_and_not_a_debt(app): + """It must be topped up, not paid off. + + Accumulating a credit per append and handing over the whole balance when + the buffer finally had room sent six megabytes in one burst, overshot the + target by a minute of film, and then said nothing for forty-six seconds. + Measured in Chrome against real fragmented MP4. + """ + player = _player(app) + pump = player[player.index("const pump = useCallback("):] + pump = pump[:pump.index("\n }, [")] + assert "STREAM_WINDOW - outstandingRef.current" in pump, ( + "pump() no longer tops a window up to what is allowed in flight") + src = APP.read_text() + window = int(re.search(r"const STREAM_WINDOW = (\d+)", src).group(1)) + assert 2 <= window <= 16, ( + f"a window of {window} segments is either too small to keep the pipe " + "busy or big enough to be a burst again") + assert "outstandingRef.current = Math.max(0, outstandingRef.current - 1)" in player, ( + "nothing decrements the window when a segment lands, so it fills once " + "and never reopens") -- cgit v1.2.3