diff --git a/server/LIVE_PACKET.md b/server/LIVE_PACKET.md index 9633630..9ae7f5d 100644 --- a/server/LIVE_PACKET.md +++ b/server/LIVE_PACKET.md @@ -4,6 +4,8 @@ Docker Compose sets `STREAMVAULT_LIVE_PACKET_AWARE=1` globally for current and f The `/api/live/:id/authorize`, `/index.m3u8`, and `/segment/:segmentId.ts` handlers retain their existing app authentication, signed native-playback tickets, reader admission, and private cache. Per-channel admission happens before async setup. In packet mode, the worker accepts only a loopback `/api/stream/:id` URL and supplies an Authorization header when configured, without putting it on argv. The existing `/api/stream/:id` route itself remains publicly accessible as before; the new signed HLS routes do not change that route's access policy. One PyAV demux/mux plus one FFmpeg HLS mux spans clean upstream Content-Length and chunked EOFs. The final packet of BOTH H.264 and AAC remains private at each source boundary: a socket cut can make PyAV flush a truncated ADTS audio frame even with `is_corrupt=False`. A replay must exactly match the committed A/V history and replace each withheld packet with a same-timestamp, same-flags payload extending its prefix before any successor output. This prevents publishing a partial AAC frame or permanently rejecting its later complete counterpart. Source reads, seam probes, staged files, published cache, readers, and channel slots have bounds. FFmpeg retains a 64-entry staging manifest to survive accelerated replay; Node removes each staged TS after its private served copy is published, so staging media do not accumulate with the manifest. A failed seam or worker exit after accepting source media invalidates the cache and latches that channel unavailable until process restart; it never reconnects raw media into the same presentation. An initial HTTP/transport failure before accepting any media leaves no playable segment and permits a fresh worker on a later authorization (not a hot retry loop). Once media is established, an unavailable HTTP reopen receives at most three attempts inside the same worker, with cancellable 250 ms gaps. It does not refresh byte/publication watchdogs or replace the mux/timeline; any eventual media still undergoes full A/V seam proof. Exhaustion and unsafe bodies invalidate the presentation as before. Native players receive an unavailable response; the client retries 503/429 within a 25-second authorization budget before considering legacy fallback, honoring `Retry-After` without retrying earlier than requested and cancelling obsolete player/backend generations. Authentication rejections never permit fallback. Chromium uses Hls.js/MSE when supported even if it advertises native HLS; iPhone/Safari retain native HLS, and native-only browsers keep their fallback. This avoids Chromium's short native buffer and source-replay underruns seen in the ESPN browser soak. Unsupported tracks (including subtitles) and incompatible codecs fail closed in this mode. +Stream-copy cuts HLS at source keyframes, so a two-second requested cadence does not impose a four-second GOP limit. Published segments must have a finite duration greater than zero and at most 12 seconds, with the existing per-segment byte bound unchanged. Every playlist advertises a fixed `TARGETDURATION:12`; it does not change as short and long GOPs arrive. Latched unsafe channels return immediate 501 without `Retry-After`, while real warmup/capacity still returns retryable 503. This status distinction never clears a latch, resumes an unproved seam, or restores stale ticket media. Production diagnostics record only channel ID and fixed retirement category. + The runtime image installs Debian `python3` and `python3-av`. Node invokes `/usr/bin/python3` so PyAV is visible; Python's `/usr/local/bin/python3` may not see apt-installed modules. FFmpeg and Python diagnostics are suppressed because URLs can contain credentials. Upstream responses are streamed with a 16 KiB read cap, a bounded socket stall timeout, finite seam history/probe, and a staged/published disk ceiling; cumulative response bytes are **not** capped, since a legitimate live HTTP body may stay open for hours. Source-read progress is marked while consuming a repeated prefix, so the HLS publication watchdog does not kill an actively catching-up stitch; an actually silent source still times out. A separate 45-second cap without any NEW published segment kills and latches the channel even if bytes (such as TS null packets) keep arriving; bytes alone are not playable progress. The Node channel owner, not an arbitrary Python wall-clock timer, retires idle viewers. Local integration tests exercise an authenticated loopback proxy, signed playlist **and** segment fetches, both clean EOF framings, an exact non-IDR continuation checked against a baseline packet inventory, and strict FFmpeg video/audio HLS decode. Negative cases include shifted AAC, a missing replay picture, incompatible video, false Content-Length, and an open silent source; previously cached segments are invalidated and the channel does not thrash-restart. Admission, reader cap, and shutdown cleanup are separately exercised. The original standalone matcher tests include ambiguous replay and bounded probing. A production-sized, natural-EOF fixture also exercises a fast two-session HLS manifest rollover on the Pi with the default poll interval and quotas. diff --git a/server/src/index.ts b/server/src/index.ts index 5b77399..2a17e1e 100644 --- a/server/src/index.ts +++ b/server/src/index.ts @@ -219,7 +219,10 @@ app.use(createArchiveRouter({ store: archiveStore, root: archiveRoot, secret: ar prioritize: (id, start, end) => archiveCapture.prioritizeWindow(id, start, end) })); const liveBuffer = createLiveBuffer({ root: process.env.STREAMVAULT_LIVE_BUFFER_DIR || path.join(process.env.TMPDIR || '/var/tmp', 'streamvault-live-buffer'), - packetAware: packetAwareLiveEnabled(process.env.STREAMVAULT_LIVE_PACKET_AWARE) }); + packetAware: packetAwareLiveEnabled(process.env.STREAMVAULT_LIVE_PACKET_AWARE), + // IDs are validated by packet ingestion; reasons are fixed server categories. + // Never log source URLs, tickets, or worker stderr. + onUnsafe: (reason, id) => console.warn('Live HLS unavailable', { channelId: id, reason }) }); app.use('/api/live', createLiveRouter(liveBuffer, id => { const channel = getChannelById(id); // The existing proxy validates upstream responses and hides credential-bearing diff --git a/server/src/live-buffer-policy.test.ts b/server/src/live-buffer-policy.test.ts new file mode 100644 index 0000000..7f7b9a2 --- /dev/null +++ b/server/src/live-buffer-policy.test.ts @@ -0,0 +1,114 @@ +import { afterEach, beforeEach, expect, it, vi } from 'vitest'; +import { EventEmitter } from 'node:events'; +import express from 'express'; +import { createServer, type Server } from 'node:http'; +import { mkdtemp, readdir, rm, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import path from 'node:path'; +const mocks = vi.hoisted(() => ({ spawn: vi.fn() })); +vi.mock('node:child_process', () => ({ spawn: mocks.spawn })); +import { createLiveBuffer, createLiveRouter, LIVE_HLS_MAX_SEGMENT_SECONDS } from './live-buffer.js'; + +let root: string; +let buffer: ReturnType; +let stage: string; +let server: Server; +let base: string; +const source = 'http://127.0.0.1:1/api/stream/channel-a'; +const auth = { headers: { authorization: 'Bearer synthetic-token' } }; +const unsafe = vi.fn(); +beforeEach(async () => { + process.env.STREAMVAULT_AUTH_TOKEN = 'synthetic-token'; + unsafe.mockClear(); mocks.spawn.mockReset(); + mocks.spawn.mockImplementation(() => { + const worker = Object.assign(new EventEmitter(), { exitCode: null as number | null, signalCode: null, + kill: vi.fn(() => { queueMicrotask(() => { worker.exitCode = 0; worker.emit('close', 0); }); return true; }) }); + return worker; + }); + root = await mkdtemp(path.join(tmpdir(), 'sv-live-policy-')); + buffer = createLiveBuffer({ root, packetAware: true, pollMs: 10, onUnsafe: unsafe }); + await buffer.playlist('channel-a', source); + await vi.waitFor(() => expect(mocks.spawn).toHaveBeenCalledTimes(1)); + const [channel] = await readdir(root); + const [ingest] = (await readdir(path.join(root, channel))).filter(name => name.startsWith('ingest-')); + stage = path.join(root, channel, ingest); + const app = express(); + app.use('/api/live', createLiveRouter(buffer, id => id === 'channel-a' ? source : null)); + server = createServer(app); + await new Promise(resolve => server.listen(0, '127.0.0.1', resolve)); + const address = server.address(); + if (!address || typeof address === 'string') throw Error('no listener'); + base = `http://127.0.0.1:${address.port}/api/live/channel-a`; +}); +afterEach(async () => { + await buffer.stop(); + await new Promise(resolve => { server.close(() => resolve()); server.closeAllConnections(); }); + await rm(root, { recursive: true, force: true }); + delete process.env.STREAMVAULT_AUTH_TOKEN; +}); +async function publish(index: number, duration: string) { + await writeFile(path.join(stage, `${index}.ts`), Buffer.alloc(188)); + await writeFile(path.join(stage, 'index.m3u8'), `#EXTM3U\n#EXTINF:${duration},\n${index}.ts\n`); +} +it('keeps a fixed twelve-second target as short and long GOPs arrive, including the upper bound', async () => { + expect(LIVE_HLS_MAX_SEGMENT_SECONDS).toBe(12); + for (const [index, duration] of ['2', '8.334', '8.333', '12'].entries()) { + await publish(index, duration); + await vi.waitFor(async () => { + const manifest = await buffer.playlist('channel-a', source); + expect(manifest).toContain(`#EXTINF:${Number(duration).toFixed(3)},`); + expect(manifest).toContain('#EXT-X-TARGETDURATION:12\n'); + }); + } + expect(unsafe).not.toHaveBeenCalled(); +}); +it.each([{ label: 'over twelve seconds', duration: '12.001' }, { label: 'zero', duration: '0' }, + { label: 'nonfinite', duration: '9'.repeat(400) }])('still latches an unpublishable duration: $label', async ({ duration }) => { + await publish(0, duration); + await vi.waitFor(() => expect(buffer.activeCount).toBe(0)); + expect(unsafe).toHaveBeenCalledExactlyOnceWith('unpublishable_stage', 'channel-a'); + expect(await buffer.playlist('channel-a', source)).toBeNull(); + expect(mocks.spawn).toHaveBeenCalledTimes(1); +}); +it('reports a latched unsafe channel immediately as 501 without retry or stale ticket media', async () => { + await publish(0, '2'); + await vi.waitFor(async () => expect(await buffer.playlist('channel-a', source)).toContain('segment/')); + const authorization = await fetch(`${base}/authorize`, auth); + expect(authorization.status).toBe(200); + const { playlistUrl } = await authorization.json(); + const playlist = new URL(playlistUrl, base); + const text = await (await fetch(playlist)).text(); + const segment = new URL(text.split('\n').find(line => line.startsWith('segment/'))!, playlist); + expect((await fetch(segment)).status).toBe(200); + await writeFile(path.join(stage, 'UNSAFE'), 'fixed-test-marker'); + await vi.waitFor(() => expect(buffer.activeCount).toBe(0)); + const started = performance.now(); + const rejected = await fetch(`${base}/authorize`, auth); + expect(rejected.status).toBe(501); + expect(rejected.headers.get('retry-after')).toBeNull(); + expect(performance.now() - started).toBeLessThan(1000); + expect(buffer.isUnsafe('channel-a')).toBe(true); + expect(buffer.isUnsafe('channel-b')).toBe(false); + for (let attempt = 0; attempt < 3; attempt++) { + const response = await fetch(playlist); + expect(response.status).toBe(501); + expect(response.headers.get('retry-after')).toBeNull(); + expect((await fetch(segment)).status).toBe(404); + expect(await buffer.playlist('channel-a', source)).toBeNull(); + } + expect((await fetch(`${base}/authorize`)).status).toBe(401); + expect((await fetch(`${base}/index.m3u8`)).status).toBe(401); + expect(unsafe).toHaveBeenCalledExactlyOnceWith('worker_unsafe', 'channel-a'); + expect(mocks.spawn).toHaveBeenCalledTimes(1); + expect(buffer.activeCount).toBe(0); +}, 10000); +it('still returns retryable 503 for genuine cold playlist and authorization warmup', async () => { + const playlist = await fetch(`${base}/index.m3u8`, auth); + expect(playlist.status).toBe(503); + expect(playlist.headers.get('retry-after')).toBe('2'); + const authorization = await fetch(`${base}/authorize`, auth); + expect(authorization.status).toBe(503); + expect(authorization.headers.get('retry-after')).toBe('2'); + expect(buffer.activeCount).toBe(1); + expect(unsafe).not.toHaveBeenCalled(); +}, 10000); diff --git a/server/src/live-buffer-startup.test.ts b/server/src/live-buffer-startup.test.ts index 1fcb188..13451ec 100644 --- a/server/src/live-buffer-startup.test.ts +++ b/server/src/live-buffer-startup.test.ts @@ -11,17 +11,18 @@ let started: number; let buffer: ReturnType; let worker: EventEmitter & { kill: ReturnType; exitCode: number | null; signalCode: string | null }; const reasons: string[] = []; +const unsafe = vi.fn((reason: string) => { reasons.push(reason); }); const inspected = vi.fn(async () => { throw Object.assign(new Error('absent'), { code: 'ENOENT' }); }); beforeEach(async () => { vi.useFakeTimers({ toFake: ['Date'] }); vi.setSystemTime(100_000); - reasons.length = 0; inspected.mockClear(); + reasons.length = 0; inspected.mockClear(); unsafe.mockClear(); root = await mkdtemp(path.join(tmpdir(), 'sv-live-startup-')); worker = Object.assign(new EventEmitter(), { exitCode: null as number | null, signalCode: null as string | null, kill: vi.fn(() => { queueMicrotask(() => { worker.exitCode = 0; worker.emit('close', 0); }); return true; }) }); mocks.spawn.mockReset(); mocks.spawn.mockReturnValue(worker); buffer = createLiveBuffer({ root, packetAware: true, pollMs: 10, stallMs: 15_000, - maxUnpublishedMs: 45_000, idleMs: 120_000, unsafeMarkerStat: inspected, onUnsafe: reason => reasons.push(reason) }); + maxUnpublishedMs: 45_000, idleMs: 120_000, unsafeMarkerStat: inspected, onUnsafe: unsafe }); await buffer.playlist('channel-a', 'http://127.0.0.1:1/api/stream/channel-a'); await vi.waitFor(() => expect(mocks.spawn).toHaveBeenCalledTimes(1)); started = Date.now(); @@ -41,6 +42,7 @@ it('does not poison a cold worker before its first source byte but still bounds vi.setSystemTime(started + 46_000); await vi.waitFor(() => expect(buffer.activeCount).toBe(0)); expect(reasons).toEqual(['no_new_segments']); + expect(unsafe).toHaveBeenCalledExactlyOnceWith('no_new_segments', 'channel-a'); expect(worker.kill).toHaveBeenCalledWith('SIGTERM'); expect(await buffer.playlist('channel-a', 'http://127.0.0.1:1/api/stream/channel-a')).toBeNull(); expect(mocks.spawn).toHaveBeenCalledTimes(1); diff --git a/server/src/live-buffer.test.ts b/server/src/live-buffer.test.ts index dd4f555..99dc1bc 100644 --- a/server/src/live-buffer.test.ts +++ b/server/src/live-buffer.test.ts @@ -102,7 +102,7 @@ describe('shared live rolling HLS HTTP', () => { expect(firstPlayableMs).toBeLessThan(9000); const manifest = await waitPlaylist(playlist, text => text.includes('#EXT-X-DISCONTINUITY') && (text.match(/segment\//g) || []).length >= 2); expect(manifest).not.toContain('#EXT-X-ENDLIST'); - expect(manifest).toContain('#EXT-X-TARGETDURATION:4'); + expect(manifest).toContain('#EXT-X-TARGETDURATION:12'); const paths = [...manifest.matchAll(/^(segment\/[^\s]+)$/gm)].map(match => match[1]); expect(paths.length).toBeGreaterThan(1); const segment = await fetch(new URL(paths[0], playlist)); diff --git a/server/src/live-buffer.ts b/server/src/live-buffer.ts index 4d890bb..1c1cb29 100644 --- a/server/src/live-buffer.ts +++ b/server/src/live-buffer.ts @@ -9,6 +9,9 @@ import { isAuthorizedRequest } from './security.js'; interface Segment { id: number; name: string; duration: number; bytes: number; discontinuity: boolean; } export const packetAwareLiveEnabled = (value: string | undefined): boolean => value === '1'; +// Stream-copy HLS cuts at source keyframes, not the requested two-second cadence. +// Keep the advertised upper bound fixed for the lifetime of every playlist. +export const LIVE_HLS_MAX_SEGMENT_SECONDS = 12; interface Channel { id: string; dir: string; epoch: string; segments: Segment[]; bytes: number; sequence: number; discontinuitySequence: number; lastAccess: number; lastPublish: number; worker?: ChildProcess; working?: string; @@ -20,7 +23,7 @@ export interface LiveBufferOptions { stallMs?: number; maxUnpublishedMs?: number; idleMs?: number; retryMs?: number; pollMs?: number; segmentSeconds?: number; minFreeBytes?: number; removeChannelDir?: (dir: string) => Promise; unsafeMarkerStat?: (file: string) => Promise; - onUnsafe?: (reason: string) => void; // fixed, non-URL diagnostic categories + onUnsafe?: (reason: string, id?: string) => void; // fixed, non-URL diagnostic categories packetAware?: boolean; } @@ -59,10 +62,11 @@ export function createLiveBuffer(options: LiveBufferOptions = {}) { const channels = new Map(); const unsafeChannels = new Set(); let unsafeCapacityExhausted = false; + const isUnsafe = (id: string): boolean => packetAware && (unsafeCapacityExhausted || unsafeChannels.has(id)); function markUnsafe(id: string, reason = 'worker_exit') { // Bound failed-ID accounting; if exhausted, reject ALL new packet channels // until restart rather than forgetting an unsafe one and retrying it. - options.onUnsafe?.(reason); + options.onUnsafe?.(reason, id); if (unsafeChannels.size >= maxChannels * 16) unsafeCapacityExhausted = true; else unsafeChannels.add(id); } @@ -148,7 +152,7 @@ export function createLiveBuffer(options: LiveBufferOptions = {}) { const key = `${ch.generation}-${name}`; const duration = Number(durationText); // Never conceal a unique segment by skipping an unpublishable input. - if (!Number.isFinite(duration) || duration <= 0 || duration > 4) { await retireUnsafe(ch); return; } + if (!Number.isFinite(duration) || duration <= 0 || duration > LIVE_HLS_MAX_SEGMENT_SECONDS) { await retireUnsafe(ch); return; } const stage = path.join(ch.working, name); let data: Buffer; try { @@ -293,7 +297,7 @@ export function createLiveBuffer(options: LiveBufferOptions = {}) { } async function get(id: string, url: string) { - if (packetAware && (unsafeCapacityExhausted || unsafeChannels.has(id))) return null; + if (isUnsafe(id)) return null; if (packetAware && (!/^https?:\/\/127\.0\.0\.1:\d+\/api\/stream\/[A-Za-z0-9_-]+(?:\?subs=1)?$/.test(url) || !/^[-A-Za-z0-9_]+$/.test(id))) return null; let ch = channels.get(id); @@ -318,12 +322,13 @@ export function createLiveBuffer(options: LiveBufferOptions = {}) { } return { + isUnsafe, get activeCount() { return channels.size; }, get activeReaders() { return activeReaders; }, async playlist(id: string, url: string, ticket?: string) { const ch = await get(id, url); if (!ch || !ch.segments.length) return null; - const rows = ['#EXTM3U', '#EXT-X-VERSION:3', '#EXT-X-TARGETDURATION:4', + const rows = ['#EXTM3U', '#EXT-X-VERSION:3', `#EXT-X-TARGETDURATION:${LIVE_HLS_MAX_SEGMENT_SECONDS}`, `#EXT-X-MEDIA-SEQUENCE:${ch.segments[0].id}`, `#EXT-X-DISCONTINUITY-SEQUENCE:${ch.discontinuitySequence}`]; for (const seg of ch.segments) { if (seg.discontinuity) rows.push('#EXT-X-DISCONTINUITY'); @@ -386,6 +391,7 @@ export function createLiveRouter(buffer: ReturnType, so const url = source(id); if (!url) { res.status(404).end(); return; } if (req.query.audio === '1') { res.status(422).end(); return; } + if (buffer.isUnsafe(id)) { res.status(501).end(); return; } // Native players may not retry an initial 503 playlist. Wait briefly for // a published segment; otherwise let the caller use its legacy TS path. if (authorizationWaiters >= 8) { res.set('Retry-After', '2').status(503).end(); return; } @@ -395,11 +401,13 @@ export function createLiveRouter(buffer: ReturnType, so const deadline = Date.now() + 5_000; while (Date.now() < deadline && !req.destroyed) { const manifest = await buffer.playlist(id, url); + if (buffer.isUnsafe(id)) break; if (manifest?.includes('\nsegment/')) { ready = true; break; } await new Promise(resolve => setTimeout(resolve, 100)); } } finally { authorizationWaiters--; } if (req.destroyed) return; + if (buffer.isUnsafe(id)) { res.status(501).end(); return; } if (!ready) { res.set('Retry-After', '2').status(503).end(); return; } const { ticket, expiresAt } = ticketFor(id); res.set('Cache-Control', 'private, no-store').json({ playlistUrl: `/api/live/${encodeURIComponent(id)}/index.m3u8?ticket=${ticket}`, expiresAt }); @@ -408,9 +416,11 @@ export function createLiveRouter(buffer: ReturnType, so const id = String(req.params.id); const url = prepare(req, res, id); if (!url) return; + if (buffer.isUnsafe(id)) { res.status(501).end(); return; } const ticket = typeof req.query.ticket === 'string' && allowed(req, id) ? req.query.ticket : undefined; try { const manifest = await buffer.playlist(id, url, ticket); + if (buffer.isUnsafe(id)) { res.status(501).end(); return; } if (!manifest) { res.set('Retry-After', '2').status(503).end(); return; } res.type('application/vnd.apple.mpegurl').send(manifest); } catch { res.status(503).end(); } diff --git a/server/src/live-packet.integration.test.ts b/server/src/live-packet.integration.test.ts index f3f57fb..e3d7dfd 100644 --- a/server/src/live-packet.integration.test.ts +++ b/server/src/live-packet.integration.test.ts @@ -54,6 +54,74 @@ async function waitFor(fn: () => Promise, ms = 12000): Promise { throw Error('packet live readiness deadline exceeded'); } +it('authorizes and serves authenticated eight-second GOP segments without an unsafe latch', async () => { + process.env.STREAMVAULT_AUTH_TOKEN = 'test-only-local-token'; + const root = await mkdtemp(path.join(tmpdir(), 'sv-packet-long-gop-')); + roots.push(root); + const sourceFile = path.join(root, 'source.ts'); + const encode = spawnSync('ffmpeg', ['-hide_banner', '-loglevel', 'error', + '-f', 'lavfi', '-i', 'testsrc2=size=160x90:rate=10:duration=24', + '-f', 'lavfi', '-i', 'sine=frequency=440:sample_rate=48000:duration=24', + '-c:v', 'libx264', '-preset', 'ultrafast', '-threads', '1', '-g', '80', '-keyint_min', '80', + '-sc_threshold', '0', '-bf', '0', '-c:a', 'aac', '-f', 'mpegts', sourceFile], { timeout: 15000 }); + expect(encode.status).toBe(0); + const body = await readFile(sourceFile); + const upstream = express(); + let requests = 0; + upstream.get('/api/stream/channel-a', (req, res) => { + if (req.header('authorization') !== 'Bearer test-only-local-token') { res.status(401).end(); return; } + requests++; + res.type('video/mp2t').write(body); + const packet = Buffer.alloc(188, 0xff); packet.set([0x47, 0x1f, 0xff, 0x10]); + const block = Buffer.concat(Array.from({ length: 100 }, () => packet)); + const pace = setInterval(() => { if (!res.destroyed) res.write(block); }, 100); + res.once('close', () => clearInterval(pace)); + }); + const proxy = await listen(upstream); + const failures: string[] = []; + const buffer = createLiveBuffer({ root: path.join(root, 'cache'), packetAware: true, pollMs: 50, + onUnsafe: reason => failures.push(reason) }); + buffers.push(buffer); + const app = express(); + app.use('/api/live', createLiveRouter(buffer, id => id === 'channel-a' ? `${proxy}/api/stream/${id}` : null)); + const base = await listen(app); + const endpoint = `${base}/api/live/channel-a`; + expect((await fetch(`${endpoint}/index.m3u8`)).status).toBe(401); + // Exercise the duration gate before authorization, independently of cold + // PyAV import/FFmpeg startup time on a busy recording host. + await waitFor(async () => { + const manifest = await buffer.playlist('channel-a', `${proxy}/api/stream/channel-a`); + return manifest || failures.length ? true : null; + }, 15000); + const authorization = await fetch(`${endpoint}/authorize`, { headers: { authorization: 'Bearer test-only-local-token' } }); + expect(authorization.status, `unsafe reasons=${failures.join(',')}`).toBe(200); + const { playlistUrl } = await authorization.json(); + const playlist = new URL(playlistUrl, base); + const manifest = await waitFor(async () => { + const res = await fetch(playlist); + if (res.status !== 200) return null; + const text = await res.text(); + return (text.match(/^segment\//gm) ?? []).length >= 2 ? text : null; + }); + expect(manifest.includes('#EXT-X-TARGETDURATION:12\n')).toBe(true); + const durations = [...manifest.matchAll(/#EXTINF:([\d.]+)/g)].map(match => Number(match[1])); + expect(durations.length).toBeGreaterThanOrEqual(2); + expect(durations.every(duration => duration >= 8 && duration <= 8.1)).toBe(true); + for (const segment of manifest.split('\n').filter(line => line.startsWith('segment/'))) { + expect(segment.includes('ticket=')).toBe(true); + expect((await fetch(new URL(segment.split('?')[0], playlist))).status).toBe(401); + const response = await fetch(new URL(segment, playlist)); + expect(response.status).toBe(200); + expect((await response.arrayBuffer()).byteLength).toBeGreaterThan(188); + } + const refreshed = await fetch(playlist); + expect(refreshed.status).toBe(200); + expect((await refreshed.text()).includes('#EXT-X-TARGETDURATION:12\n')).toBe(true); + expect(failures).toEqual([]); + expect(buffer.activeCount).toBe(1); + expect(requests).toBe(1); +}, 20000); + it('does not kill a progressing replay merely because HLS cannot publish duplicates yet', async () => { const [first, second] = await media(); const upstream = express(); @@ -584,7 +652,13 @@ with av.open(sys.argv[1]) as src, av.open(sys.argv[2], 'w', format='mpegts') as release(); await waitFor(async () => buffer.activeCount === 0 ? true : null, 10000); expect((await fetch(new URL(firstSegment, playlist))).status).toBe(404); - expect((await fetch(playlist)).status).toBe(503); + const rejected = await fetch(playlist); + expect(rejected.status).toBe(501); + expect(rejected.headers.get('retry-after')).toBeNull(); + const rejectedAuthorization = await fetch(`${base}/api/live/channel-a/authorize`, { headers: { authorization: 'Bearer test-only-local-token' } }); + expect(rejectedAuthorization.status).toBe(501); + expect(rejectedAuthorization.headers.get('retry-after')).toBeNull(); + expect(buffer.isUnsafe('channel-a')).toBe(true); await new Promise(resolve => setTimeout(resolve, 500)); expect(requests).toBe(2); }, 35000, diff --git a/src/hooks/usePlayer.liveHls.test.tsx b/src/hooks/usePlayer.liveHls.test.tsx index 5d9f42a..4706e9e 100644 --- a/src/hooks/usePlayer.liveHls.test.tsx +++ b/src/hooks/usePlayer.liveHls.test.tsx @@ -183,6 +183,49 @@ describe('signed live MSE startup lead', () => { }, ); + it.each([401, 403, 404, 410])('shows an actionable error for rejected legacy live HTTP %s rather than retrying forever', async code => { + mocks.authorize.mockResolvedValue(null as unknown as string); + await act(async () => { + hookRef.current?.play(); + await vi.waitFor(() => expect(mocks.mpegtsPlayer.load).toHaveBeenCalled()); + }); + const onError = mocks.mpegtsPlayer.on.mock.calls.find(([event]) => event === 'error')![1]; + await act(async () => onError('NetworkError', 'HttpStatusCodeInvalid', { code, msg: 'test rejection' })); + expect(usePlayerStore.getState().status).toBe('error'); + expect(usePlayerStore.getState().errorMessage).toMatch(/live stream.*retry/i); + expect(mocks.mpegtsPlayer.destroy).toHaveBeenCalled(); + await act(async () => { + onError('NetworkError', 'HttpStatusCodeInvalid', { code: 503 }); + video.dispatchEvent(new Event('error')); + video.dispatchEvent(new Event('waiting')); + await vi.advanceTimersByTimeAsync(65_000); + }); + expect(usePlayerStore.getState().status).toBe('error'); + expect(mocks.authorize).toHaveBeenCalledOnce(); + // A new explicit retry must still be able to attach and play normally. + mocks.authorize.mockResolvedValue('/api/live/live_test/index.m3u8?ticket=retry'); + ranges = [[40, 60]]; + await ready(); + expect(video.play).toHaveBeenCalledOnce(); + expect(usePlayerStore.getState().status).toBe('playing'); + }); + + it('does not shorten first-frame acquisition after an initial zero-time timeupdate', async () => { + await act(async () => hookRef.current?.play()); + video.currentTime = 0; + await act(async () => { + video.dispatchEvent(new Event('timeupdate')); + await vi.advanceTimersByTimeAsync(24_000); + }); + expect(mocks.authorize).toHaveBeenCalledOnce(); + video.currentTime = 1; + await act(async () => { + video.dispatchEvent(new Event('timeupdate')); + await vi.advanceTimersByTimeAsync(12_250); + }); + expect(mocks.authorize).toHaveBeenCalledTimes(2); + }); + it.each(['stop', 'channel switch', 'new authorization', 'new authorization with queued retry'])( 'cannot start from obsolete callbacks or recovery timers after %s', async cancellation => { await ready(); diff --git a/src/hooks/usePlayer.ts b/src/hooks/usePlayer.ts index b315c62..d8a3bc5 100644 --- a/src/hooks/usePlayer.ts +++ b/src/hooks/usePlayer.ts @@ -831,7 +831,8 @@ export function usePlayer(): { } } catch { /* ignore */ } - let lastMediaTime = -1; + // Zero is the unstarted media clock, not evidence of first-frame progress. + let lastMediaTime = 0; let startupReady = false; let startupStarted = false; let canPlay = false; @@ -1033,7 +1034,7 @@ export function usePlayer(): { }; video.onsuspend = () => log.debug('HTML5 event: suspend'); video.onerror = () => { - if (!isCurrentPlayback()) return; + if (!isCurrentPlayback() || usePlayerStore.getState().status === 'error') return; const err = video.error; const errMsg = err ? `code=${err.code} message="${err.message}"` : 'unknown'; log.error(`HTML5 event: error — ${errMsg}`); @@ -1147,6 +1148,15 @@ export function usePlayer(): { player.on(mpegts.Events.ERROR, (type: string, detail: string, info: unknown) => { if (!isCurrentPlayback() || activeMpegtsPlayer !== player) return; log.error(`mpegts ERROR: type=${type} detail=${detail}`, info); + const status = info && typeof info === 'object' && 'code' in info + ? Number(info.code) : 0; + if (isLiveTs && detail === 'HttpStatusCodeInvalid' && [401, 403, 404, 410].includes(status)) { + disableLiveStreamRecovery(); + activeMpegtsPlayer = null; + player.destroy(); + setError('This live stream is unavailable or was rejected upstream. Tap to retry.'); + return; + } if (isLiveTs) { setStatus('loading'); liveStreamRecovery.transportEnded('mpegts-error'); diff --git a/src/services/liveHls.test.ts b/src/services/liveHls.test.ts index d8df33c..fb3a00f 100644 --- a/src/services/liveHls.test.ts +++ b/src/services/liveHls.test.ts @@ -60,13 +60,14 @@ describe('shared live HLS playback', () => { expect(hls.destroy).toHaveBeenCalledOnce(); }); - it('targets 24 seconds at the served four-second target without enabling max-latency catch-up', async () => { + it('keeps 24-second headroom independent of long-GOP TARGETDURATION without catch-up seeks', async () => { const video = document.createElement('video'); vi.spyOn(video, 'canPlayType').mockReturnValue(''); const dispose = await attachLiveHls(video, '/api/live/test/index.m3u8', vi.fn()); - expect(hlsMock.instances[0].config.liveSyncDurationCount).toBe(6); - expect(Number(hlsMock.instances[0].config.liveSyncDurationCount) * 4).toBe(24); - expect(hlsMock.instances[0].config.liveMaxLatencyDurationCount).toBe(Infinity); + expect(hlsMock.instances[0].config.liveSyncDuration).toBe(24); + expect(hlsMock.instances[0].config.liveMaxLatencyDuration).toBe(Infinity); + expect(hlsMock.instances[0].config).not.toHaveProperty('liveSyncDurationCount'); + expect(hlsMock.instances[0].config).not.toHaveProperty('liveMaxLatencyDurationCount'); dispose(); }); diff --git a/src/services/liveHls.ts b/src/services/liveHls.ts index f0de690..4832f16 100644 --- a/src/services/liveHls.ts +++ b/src/services/liveHls.ts @@ -31,13 +31,12 @@ export async function attachLiveHls( const hls: Hls = new HlsPlayer({ enableWorker: true, backBufferLength: 30, - // Six advertised TARGETDURATIONs: 24 seconds on this server's target of - // four, clamped to the available window. Prioritize headroom over latency. - liveSyncDurationCount: 6, + // Keep the intentional 24-second, window-clamped headroom independent of + // TARGETDURATION: long-GOP feeds need a larger advertised segment bound. + liveSyncDuration: 24, // Preserve buffered playback rather than seeking nearer the edge after a - // replay burst. The default infinite catch-up threshold avoids a seek that - // discards the startup lead and emits a visible waiting event. - liveMaxLatencyDurationCount: Infinity, + // replay burst. Infinite maximum latency keeps catch-up seeks disabled. + liveMaxLatencyDuration: Infinity, }); let disposed = false; hls.on(HlsPlayer.Events.BUFFER_APPENDED, () => {