Queue bounded stream prewarm requests
This commit is contained in:
@@ -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",
|
||||
|
||||
@@ -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
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