Queue bounded stream prewarm requests

This commit is contained in:
Jonathan Sykes
2026-10-07 03:13:30 +08:00
parent b4c776dcfe
commit b77938ab4a
4 changed files with 82 additions and 7 deletions

View File

@@ -6,7 +6,7 @@
"scripts": {
"start": "bun server.js",
"dev": "bun --hot server.js",
"test": "bun test ./recommendations.test.js && bun test --timeout 60000 ./media-cache.test.js && bun test ./notes.test.js && bun test ./transcriptions.test.js && bun test ./admin-analytics.test.js && bun test ./remote.test.js && bun test ./party.test.js && bun test ./uploads.test.js && bun test ./innertube.test.js && bun test ./related.test.js && bun test ./ytdlp-pool.test.js && bun test ./p2p-db.test.js && bun test ./p2p-admit.test.js && bun test ./p2p-retention.test.js && bun test ./p2p-routes.test.js && bun test ./p2p-hub.test.js && bun test --timeout 60000 ./p2p-intake.test.js && bun test ./flags.test.js"
"test": "bun test ./recommendations.test.js && bun test --timeout 60000 ./media-cache.test.js && bun test ./notes.test.js && bun test ./transcriptions.test.js && bun test ./admin-analytics.test.js && bun test ./remote.test.js && bun test ./party.test.js && bun test ./uploads.test.js && bun test ./innertube.test.js && bun test ./related.test.js && bun test ./ytdlp-pool.test.js && bun test ./warm-queue.test.js && bun test ./p2p-db.test.js && bun test ./p2p-admit.test.js && bun test ./p2p-retention.test.js && bun test ./p2p-routes.test.js && bun test ./p2p-hub.test.js && bun test --timeout 60000 ./p2p-intake.test.js && bun test ./flags.test.js"
},
"dependencies": {
"@hono/node-server": "^1.14.0",

View File

@@ -67,6 +67,7 @@ import { ingest as collectVideoMetadata, syncListening, linkListening, startThum
import { registerCatalogRoutes } from './recommendations.js';
import QRCode from 'qrcode';
import { createYtdlpPool } from './ytdlp-pool.js';
import { createWarmQueue } from './warm-queue.js';
import { dirname, join as pathJoin } from 'node:path';
// A media proxy must not die because one client's stream hit an edge case
@@ -728,15 +729,20 @@ function pickFormat(formats, { formatId, wantAudio, wantHeight }) {
// GET /api/streams/warm?v=<id> — fire-and-forget: resolve streams into
// streamCache so the real /api/streams a moment later is instant. Bounded
// so a scrolling user can't queue dozens of yt-dlp processes.
const WARM_MAX = 2;
let warmActive = 0;
const streamWarmQueue = createWarmQueue({
concurrency: 2,
maxQueued: 8,
run: async (id) => {
try { if (await media.getReady(id)) return; } catch { /* fall through */ }
await resolveStreams(id);
},
});
app.get('/api/streams/warm', async (c) => {
const id = (c.req.query('v') || '').trim();
if (!/^[A-Za-z0-9_-]{11}$/.test(id)) return c.body(null, 204);
if (warmActive >= WARM_MAX) return c.body(null, 204);
try { if (await media.getReady(id)) return c.body(null, 204); } catch { /* fall through */ }
warmActive++;
resolveStreams(id).catch(() => {}).finally(() => { warmActive--; });
// Best effort: bounded queue preserves a small burst (e.g. the first three
// search cards) while limiting background extractor work.
streamWarmQueue.enqueue(id);
return c.body(null, 204);
});

30
server/warm-queue.js Normal file
View File

@@ -0,0 +1,30 @@
// Bounded scheduler for best-effort stream prewarming. Requests beyond the
// active limit wait in a short FIFO instead of being silently discarded.
export function createWarmQueue({ concurrency = 2, maxQueued = 8, run }) {
const queued = [];
const known = new Set();
let active = 0;
function pump() {
while (active < concurrency && queued.length) {
const id = queued.shift();
active++;
Promise.resolve().then(() => run(id)).catch(() => {}).finally(() => {
active--;
known.delete(id);
pump();
});
}
}
function enqueue(id) {
if (known.has(id)) return true;
if (active >= concurrency && queued.length >= maxQueued) return false;
known.add(id);
queued.push(id);
pump();
return true;
}
return { enqueue, get pending() { return queued.length + active; } };
}

39
server/warm-queue.test.js Normal file
View File

@@ -0,0 +1,39 @@
import { describe, expect, test } from 'bun:test';
import { createWarmQueue } from './warm-queue.js';
describe('stream warm queue', () => {
test('runs at the concurrency limit and drains queued requests FIFO', async () => {
const started = [];
const releases = new Map();
const q = createWarmQueue({ concurrency: 2, maxQueued: 2, run: (id) => {
started.push(id);
return new Promise((resolve) => releases.set(id, resolve));
} });
q.enqueue('a'); q.enqueue('b'); q.enqueue('c');
expect(q.pending).toBe(3);
await Promise.resolve();
expect(started).toEqual(['a', 'b']);
releases.get('a')();
await new Promise((resolve) => setTimeout(resolve, 0));
expect(started).toEqual(['a', 'b', 'c']);
releases.get('b')(); releases.get('c')();
await new Promise((resolve) => setTimeout(resolve, 0));
expect(q.pending).toBe(0);
});
test('deduplicates in-flight ids and rejects work past the queue bound', async () => {
const releases = [];
const q = createWarmQueue({ concurrency: 1, maxQueued: 1, run: () => new Promise((resolve) => releases.push(resolve)) });
expect(q.enqueue('a')).toBe(true);
await Promise.resolve();
expect(q.enqueue('a')).toBe(true);
expect(q.enqueue('b')).toBe(true);
expect(q.enqueue('c')).toBe(false);
releases[0]();
await new Promise((resolve) => setTimeout(resolve, 0));
releases[1]();
await new Promise((resolve) => setTimeout(resolve, 0));
expect(q.pending).toBe(0);
});
});