diff --git a/server/package.json b/server/package.json index 114ca94..6a6375e 100644 --- a/server/package.json +++ b/server/package.json @@ -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", diff --git a/server/server.js b/server/server.js index b904559..7596424 100644 --- a/server/server.js +++ b/server/server.js @@ -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= — 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); }); diff --git a/server/warm-queue.js b/server/warm-queue.js new file mode 100644 index 0000000..b11998d --- /dev/null +++ b/server/warm-queue.js @@ -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; } }; +} diff --git a/server/warm-queue.test.js b/server/warm-queue.test.js new file mode 100644 index 0000000..bdd589d --- /dev/null +++ b/server/warm-queue.test.js @@ -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); + }); +});