Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions server/LIVE_PACKET.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
5 changes: 4 additions & 1 deletion server/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
114 changes: 114 additions & 0 deletions server/src/live-buffer-policy.test.ts
Original file line number Diff line number Diff line change
@@ -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<typeof createLiveBuffer>;
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<void>(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<void>(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);
6 changes: 4 additions & 2 deletions server/src/live-buffer-startup.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,17 +11,18 @@ let started: number;
let buffer: ReturnType<typeof createLiveBuffer>;
let worker: EventEmitter & { kill: ReturnType<typeof vi.fn>; 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();
Expand All @@ -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);
Expand Down
2 changes: 1 addition & 1 deletion server/src/live-buffer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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));
Expand Down
Loading
Loading