Queue bounded stream prewarm requests
This commit is contained in:
@@ -6,7 +6,7 @@
|
|||||||
"scripts": {
|
"scripts": {
|
||||||
"start": "bun server.js",
|
"start": "bun server.js",
|
||||||
"dev": "bun --hot 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": {
|
"dependencies": {
|
||||||
"@hono/node-server": "^1.14.0",
|
"@hono/node-server": "^1.14.0",
|
||||||
|
|||||||
@@ -67,6 +67,7 @@ import { ingest as collectVideoMetadata, syncListening, linkListening, startThum
|
|||||||
import { registerCatalogRoutes } from './recommendations.js';
|
import { registerCatalogRoutes } from './recommendations.js';
|
||||||
import QRCode from 'qrcode';
|
import QRCode from 'qrcode';
|
||||||
import { createYtdlpPool } from './ytdlp-pool.js';
|
import { createYtdlpPool } from './ytdlp-pool.js';
|
||||||
|
import { createWarmQueue } from './warm-queue.js';
|
||||||
import { dirname, join as pathJoin } from 'node:path';
|
import { dirname, join as pathJoin } from 'node:path';
|
||||||
|
|
||||||
// A media proxy must not die because one client's stream hit an edge case
|
// 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
|
// GET /api/streams/warm?v=<id> — fire-and-forget: resolve streams into
|
||||||
// streamCache so the real /api/streams a moment later is instant. Bounded
|
// streamCache so the real /api/streams a moment later is instant. Bounded
|
||||||
// so a scrolling user can't queue dozens of yt-dlp processes.
|
// so a scrolling user can't queue dozens of yt-dlp processes.
|
||||||
const WARM_MAX = 2;
|
const streamWarmQueue = createWarmQueue({
|
||||||
let warmActive = 0;
|
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) => {
|
app.get('/api/streams/warm', async (c) => {
|
||||||
const id = (c.req.query('v') || '').trim();
|
const id = (c.req.query('v') || '').trim();
|
||||||
if (!/^[A-Za-z0-9_-]{11}$/.test(id)) return c.body(null, 204);
|
if (!/^[A-Za-z0-9_-]{11}$/.test(id)) return c.body(null, 204);
|
||||||
if (warmActive >= WARM_MAX) return c.body(null, 204);
|
// Best effort: bounded queue preserves a small burst (e.g. the first three
|
||||||
try { if (await media.getReady(id)) return c.body(null, 204); } catch { /* fall through */ }
|
// search cards) while limiting background extractor work.
|
||||||
warmActive++;
|
streamWarmQueue.enqueue(id);
|
||||||
resolveStreams(id).catch(() => {}).finally(() => { warmActive--; });
|
|
||||||
return c.body(null, 204);
|
return c.body(null, 204);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
30
server/warm-queue.js
Normal file
30
server/warm-queue.js
Normal 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
39
server/warm-queue.test.js
Normal 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);
|
||||||
|
});
|
||||||
|
});
|
||||||
Reference in New Issue
Block a user