Add admin analytics, metadata collection, and grouped lyric cues

Queue reviewable Whisper drafts from the song list and lyrics editor. Preserve line breaks within one timed cue across editing, saving, reporting, and service views.

Add storage and listening analytics with a durable metadata collector, related-search depth, video limits, thumbnail storage, and a browsable metadata library.
This commit is contained in:
Jonathan Sykes
2026-10-03 07:54:11 +08:00
parent 716b61ffee
commit 2b38717c05
23 changed files with 1181 additions and 48 deletions

149
server/admin-analytics.js Normal file
View File

@@ -0,0 +1,149 @@
import { randomUUID } from 'node:crypto';
import { statfs, stat } from 'node:fs/promises';
import { db } from './db.js';
import { ingest } from './video-catalog.js';
const LEASE = 120000;
const n = v => Number(v) || 0;
const clean = value => String(value || '').trim().slice(0, 200);
const parse = value => { try { return JSON.parse(value) || {}; } catch { return {}; } };
function publicJob(row) {
const state = parse(row.state);
return { id: row.id, query: row.query, maxVideos: n(row.max_videos), depth: n(row.depth), status: row.status,
collected: state.ids?.length || 0, enriched: state.enriched || 0, failed: state.failed || 0,
searches: state.searches || 0, current: state.current || '', error: row.error, errors: state.errors || [],
createdAt: n(row.created_at), updatedAt: n(row.updated_at) };
}
export function descriptiveMetadata(raw) {
const out = {};
for (const key of ['id','title','fulltitle','description','channel','channel_id','channel_url','uploader','uploader_id','uploader_url','upload_date','release_date','timestamp','duration','view_count','like_count','comment_count','tags','categories','language','live_status','availability','age_limit','license','chapters','thumbnails','thumbnail','webpage_url','artist','artists','album','track','release_year']) {
if (raw[key] != null) out[key] = raw[key];
}
// Retain format specs, excluding expiring media URLs and request headers.
if (Array.isArray(raw.formats)) out.formats = raw.formats.map(f => Object.fromEntries(['format_id','format_note','ext','width','height','fps','vcodec','acodec','abr','tbr','filesize','filesize_approx','audio_channels','asr'].filter(k => f[k] != null).map(k => [k, f[k]])));
if (JSON.stringify(out).length > 500000) throw Error('Extracted metadata exceeds 500 KB.');
return out;
}
function card(raw) {
return { ...raw, channel: raw.channel || raw.uploader || '', channelId: raw.channel_id,
channelUrl: raw.channel_url, thumbnail: raw.thumbnail || raw.thumbnails?.at(-1)?.url };
}
export async function storageAnalytics(paths = []) {
const queries = [
"SELECT status,COUNT(*) AS count,COALESCE(SUM(size),0) AS bytes,COALESCE(SUM(duration),0) AS seconds FROM media_cache GROUP BY status",
"SELECT kind,COUNT(*) AS count,COALESCE(SUM(size),0) AS bytes FROM uploads GROUP BY kind",
];
// Keep each aggregate on its own table, so thumbnail bytes and plays never multiply.
const [media, uploads, catalog, thumbs, details, sources, plays, top] = await Promise.all([
db.execute(queries[0]), db.execute(queries[1]),
db.execute('SELECT COUNT(*) AS videos,COALESCE(SUM(length(CAST(card AS BLOB))),0) AS bytes,MAX(updated_at) AS lastSeen FROM video_meta'),
db.execute('SELECT COUNT(*) AS total,SUM(CASE WHEN data IS NOT NULL THEN 1 ELSE 0 END) AS saved,COALESCE(SUM(size),0) AS bytes,SUM(CASE WHEN data IS NULL AND retry_at<9007199254740991 THEN 1 ELSE 0 END) AS pending FROM video_thumbnails'),
db.execute('SELECT COUNT(*) AS videos,COALESCE(SUM(length(CAST(metadata AS BLOB))),0) AS bytes FROM video_details'),
db.execute('SELECT source,COUNT(*) AS videos,SUM(discoveries) AS discoveries,MAX(last_seen) AS lastSeen FROM video_meta_sources GROUP BY source ORDER BY videos DESC'),
db.execute('SELECT COALESCE(SUM(plays),0) AS plays,COUNT(DISTINCT video_id) AS videos FROM listening_daily'),
db.execute('SELECT l.video_id,SUM(l.plays) AS plays,m.card FROM listening_daily l LEFT JOIN video_meta m ON m.id=l.video_id GROUP BY l.video_id ORDER BY plays DESC LIMIT 10'),
]);
const volumes = await Promise.all(paths.map(async ({ label, path }) => {
try { const [fs, file] = await Promise.all([statfs(path), stat(path)]); return { label, device: String(file.dev), total: n(fs.blocks) * n(fs.bsize), free: n(fs.bavail) * n(fs.bsize) }; }
catch { return { label, unavailable: true }; }
}));
return { media: media.rows, uploads: uploads.rows, catalog: catalog.rows[0], thumbnails: thumbs.rows[0], details: details.rows[0], sources: sources.rows, listening: plays.rows[0], topPlayed: top.rows.map(r => ({ id: r.video_id, plays: n(r.plays), title: parse(r.card).title || r.video_id })), volumes };
}
export function registerAnalyticsRoutes(app, { adminAuth, runYtdlp, paths = [], autoStart = true }) {
let running = null;
const owner = randomUUID();
async function runNext() {
if (running) return running;
running = execute().finally(() => { running = null; });
return running;
}
async function execute() {
const time = Date.now();
const job = (await db.execute({ sql: `UPDATE metadata_collections SET status='running',owner=?,lease_until=?,updated_at=?
WHERE id=(SELECT id FROM metadata_collections WHERE status='queued' OR (status='running' AND lease_until<?) ORDER BY created_at LIMIT 1) RETURNING *`, args: [owner, time + LEASE, time, time] })).rows[0];
if (!job) return false;
const state = { ids: [], queries: [{ q: job.query, level: 0 }], pending: [], enriched: 0, failed: 0, searches: 0, errors: [], ...parse(job.state) };
const seen = new Set(state.ids);
async function checkpoint(status = 'running', error = null) {
const result = await db.execute({ sql: "UPDATE metadata_collections SET state=?,status=?,error=?,lease_until=?,updated_at=? WHERE id=? AND owner=? AND status='running'", args: [JSON.stringify(state), status, error, Date.now() + LEASE, Date.now(), job.id, owner] });
if (!result.rowsAffected) { const error = Error('Collection cancelled or reassigned.'); error.name = 'CollectionStopped'; throw error; }
}
try {
while (state.pending.length || (state.queries.length && seen.size < job.max_videos && state.searches < 24)) {
if (!state.pending.length) {
const { q, level } = state.queries[0]; state.current = q; await checkpoint();
const remaining = job.max_videos - seen.size;
const count = Math.max(1, Math.ceil(remaining / (job.depth - level + 1)));
const output = await runYtdlp([`ytsearch${count}:${q}`, '--dump-json', '--flat-playlist', '--skip-download', '--no-warnings', '--ignore-errors'], { signal: AbortSignal.timeout(60000) });
await checkpoint(); // Recheck cancellation/ownership before storing discoveries.
const unique = new Map();
for (const raw of output.split('\n').filter(Boolean).map(parse)) {
if (/^[\w-]{11}$/.test(raw.id) && !seen.has(raw.id)) unique.set(raw.id, raw);
}
const results = [...unique.values()].slice(0, remaining);
await ingest(results.map(card), 'collector');
for (const raw of results) { if (seen.has(raw.id)) continue; seen.add(raw.id); state.ids.push(raw.id); state.pending.push({ id: raw.id, level }); }
state.queries.shift(); state.searches++; await checkpoint();
continue;
}
const item = state.pending[0]; state.current = item.id; await checkpoint();
try {
const raw = JSON.parse(await runYtdlp(['-J', '--skip-download', '--no-playlist', '--no-warnings', `https://www.youtube.com/watch?v=${item.id}`], { signal: AbortSignal.timeout(60000) }));
if (raw.id !== item.id) throw Error('Extractor returned a different video.');
await checkpoint();
const metadata = descriptiveMetadata(raw);
await ingest([card(raw)], 'collector');
await db.execute({ sql: 'INSERT INTO video_details (video_id,metadata,updated_at) VALUES (?,?,?) ON CONFLICT(video_id) DO UPDATE SET metadata=excluded.metadata,updated_at=excluded.updated_at', args: [item.id, JSON.stringify(metadata), Date.now()] });
state.enriched++;
if (item.level < job.depth && state.queries.length < 24) {
for (const q of [raw.channel || raw.uploader, ...(raw.tags || []).slice(0, 2)]) {
if (state.queries.length >= 24) break;
const term = clean(q); if (!term) continue;
state.visited ??= [job.query.toLowerCase()];
if (state.visited.includes(term.toLowerCase())) continue;
state.visited.push(term.toLowerCase()); state.queries.push({ q: term, level: item.level + 1 });
}
}
} catch (error) { if (error.name === 'CollectionStopped') throw error; state.failed++; if (state.errors.length < 20) state.errors.push({ id: item.id, error: String(error.message).slice(0, 200) }); }
state.pending.shift(); await checkpoint();
}
state.current = ''; await checkpoint('complete');
} catch (error) {
if (error.name === 'CollectionStopped') return true;
await db.execute({ sql: "UPDATE metadata_collections SET status='failed',state=?,error=?,updated_at=? WHERE id=? AND owner=? AND status='running'", args: [JSON.stringify(state), String(error.message).slice(0, 300), Date.now(), job.id, owner] });
}
return true;
}
const kick = () => runNext().catch(error => console.error('[metadata collector]', error.message));
if (autoStart) { const timer = setInterval(kick, 10000); timer.unref?.(); kick(); }
app.get('/api/admin/analytics', adminAuth, async c => c.json({ ok: true, ...await storageAnalytics(paths) }, 200, { 'Cache-Control': 'no-store' }));
app.get('/api/admin/metadata', adminAuth, async c => {
const query = clean(c.req.query('q')), offset = Math.max(0, Math.min(500000, parseInt(c.req.query('offset')) || 0));
const where = query ? 'WHERE m.hay LIKE ?' : '', args = query ? ['%' + query.replace(/[\\%_]/g, '\\$&') + '%'] : [];
const filter = where ? where + " ESCAPE '\\'" : '';
const rows = (await db.execute({ sql: `SELECT m.id,m.card,m.updated_at,d.updated_at AS enriched,t.size AS thumbnailBytes FROM video_meta m LEFT JOIN video_details d ON d.video_id=m.id LEFT JOIN video_thumbnails t ON t.video_id=m.id ${filter} ORDER BY m.updated_at DESC LIMIT 50 OFFSET ?`, args: [...args, offset] })).rows;
const total = n((await db.execute({ sql: `SELECT COUNT(*) AS n FROM video_meta m ${filter}`, args })).rows[0].n);
return c.json({ ok: true, total, offset, videos: rows.map(r => ({ ...parse(r.card), updatedAt: n(r.updated_at), enriched: !!r.enriched, thumbnailBytes: n(r.thumbnailBytes) })) });
});
app.get('/api/admin/metadata/:id', adminAuth, async c => {
const row = (await db.execute({ sql: 'SELECT m.card,d.metadata FROM video_meta m LEFT JOIN video_details d ON d.video_id=m.id WHERE m.id=?', args: [c.req.param('id')] })).rows[0];
return row ? c.json({ ok: true, metadata: row.metadata ? parse(row.metadata) : parse(row.card) }) : c.json({ ok: false, error: 'Not found.' }, 404);
});
app.get('/api/admin/collections', adminAuth, async c => c.json({ ok: true, jobs: (await db.execute('SELECT * FROM metadata_collections ORDER BY created_at DESC LIMIT 30')).rows.map(publicJob) }));
app.post('/api/admin/collections', adminAuth, async c => {
let body; try { body = await c.req.json(); } catch { return c.json({ ok: false, error: 'Invalid JSON.' }, 400); }
const query = clean(body?.query), limit = Number(body?.maxVideos), depth = Number(body?.depth);
if (!query || !Number.isInteger(limit) || limit < 1 || limit > 500 || !Number.isInteger(depth) || depth < 0 || depth > 3) return c.json({ ok: false, error: 'Enter a search, 1–500 videos and depth 0–3.' }, 400);
const time = Date.now(), id = randomUUID();
const result = await db.execute({ sql: "INSERT INTO metadata_collections (id,query,max_videos,depth,created_at,updated_at) SELECT ?,?,?,?,?,? WHERE (SELECT COUNT(*) FROM metadata_collections WHERE status IN ('queued','running'))<3", args: [id, query, limit, depth, time, time] });
if (!result.rowsAffected) return c.json({ ok: false, error: 'Three collections are already active.' }, 429);
if (autoStart) kick();
return c.json({ ok: true, id }, 202);
});
app.post('/api/admin/collections/:id/cancel', adminAuth, async c => {
const result = await db.execute({ sql: "UPDATE metadata_collections SET status='cancelled',updated_at=? WHERE id=? AND status IN ('queued','running')", args: [Date.now(), c.req.param('id')] });
return c.json({ ok: !!result.rowsAffected });
});
return { runNext };
}

View File

@@ -0,0 +1,83 @@
import { test, expect, beforeAll, afterAll } from 'bun:test';
import { mkdtempSync, rmSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { Hono } from 'hono';
const root = mkdtempSync(join(tmpdir(), 'ytp-analytics-'));
process.env.DB_PATH = join(root, 'test.db');
const { db, initDb, upsertMedia } = await import('./db.js');
const { registerAnalyticsRoutes, storageAnalytics } = await import('./admin-analytics.js');
const catalog = await import('./video-catalog.js');
const originalFetch = globalThis.fetch;
globalThis.fetch = async () => new Response(new Uint8Array([1,2,3]), { headers: { 'Content-Type': 'image/jpeg' } });
let app, runner, calls = [], serial = 0;
const adminAuth = async (c, next) => c.req.header('x-test-admin') === 'yes' ? next() : c.json({ ok: false }, 401);
const get = path => app.request(path, { headers: { 'x-test-admin': 'yes' } });
const post = (path, body = {}) => app.request(path, { method: 'POST', headers: { 'x-test-admin': 'yes', 'Content-Type': 'application/json' }, body: JSON.stringify(body) });
const fakeExtractor = async args => {
calls.push(args);
if (args[0].startsWith('ytsearch')) return Array.from({ length: Number(args[0].match(/^ytsearch(\d+)/)[1]) }, () => ({ id: String(++serial).padStart(11, '0'), title: 'Worship ' + serial, channel: 'Channel ' + serial })).map(JSON.stringify).join('\n');
const id = args.at(-1).split('v=')[1];
return JSON.stringify({ id, title: 'Enriched ' + id, channel: 'Related ' + id, tags: ['Topic ' + id], duration: 60, description: 'Full description '.repeat(200), view_count: 1234, formats: [{ format_id: '140', acodec: 'aac', url: 'https://temporary.example', http_headers: { Cookie: 'omitted' } }] });
};
beforeAll(async () => { await initDb(); app = new Hono(); runner = registerAnalyticsRoutes(app, { adminAuth, runYtdlp: fakeExtractor, autoStart: false, paths: [{ label: 'DB', path: root }] }); });
afterAll(async () => { await catalog.drainThumbnails(); globalThis.fetch = originalFetch; db.close(); rmSync(root, { recursive: true, force: true }); });
test('admin routes reject unauthenticated requests and validate collector limits', async () => {
for (const path of ['/api/admin/analytics','/api/admin/metadata','/api/admin/metadata/00000000001','/api/admin/collections']) expect((await app.request(path)).status).toBe(401);
expect((await app.request('/api/admin/collections', { method: 'POST' })).status).toBe(401);
for (const body of [{ query: '', maxVideos: 5, depth: 0 }, { query: 'worship', maxVideos: 501, depth: 0 }, { query: 'worship', maxVideos: 5, depth: 4 }]) expect((await post('/api/admin/collections', body)).status).toBe(400);
for (let i = 0; i < 3; i++) expect((await post('/api/admin/collections', { query: 'queue ' + i, maxVideos: 1, depth: 0 })).status).toBe(202);
expect((await post('/api/admin/collections', { query: 'queue 4', maxVideos: 1, depth: 0 })).status).toBe(429);
for (const job of (await (await get('/api/admin/collections')).json()).jobs) await post('/api/admin/collections/' + job.id + '/cancel');
});
test('depth zero stores full descriptive metadata and actual thumbnail bytes', async () => {
calls = []; await post('/api/admin/collections', { query: 'root worship', maxVideos: 4, depth: 0 }); await runner.runNext(); await catalog.drainThumbnails();
const job = (await (await get('/api/admin/collections')).json()).jobs.find(j => j.query === 'root worship');
expect(job.status).toBe('complete'); expect(job.collected).toBe(4); expect(job.enriched).toBe(4); expect(calls.filter(a => a[0].startsWith('ytsearch'))).toHaveLength(1);
const videos = (await (await get('/api/admin/metadata')).json()).videos;
expect(videos).toHaveLength(4); expect(videos.every(v => v.enriched && v.thumbnailBytes === 3)).toBe(true);
const detail = (await (await get('/api/admin/metadata/' + videos[0].id)).json()).metadata;
expect(detail.description.length).toBeGreaterThan(1200); expect(detail.view_count).toBe(1234); expect(detail.formats[0].url).toBeUndefined(); expect(detail.formats[0].http_headers).toBeUndefined();
});
test('related depth follows channels/topics but caps unique videos across all searches', async () => {
calls = []; await post('/api/admin/collections', { query: 'deep worship', maxVideos: 10, depth: 2 }); await runner.runNext();
const job = (await (await get('/api/admin/collections')).json()).jobs.find(j => j.query === 'deep worship');
expect(job.status).toBe('complete'); expect(job.collected).toBe(10); expect(job.enriched).toBe(10);
const searches = calls.filter(a => a[0].startsWith('ytsearch')); expect(searches.length).toBeGreaterThan(1); expect(searches[1][0]).toContain('Related'); expect(searches.length).toBeLessThanOrEqual(24);
});
test('expired jobs resume pending videos without repeating their search', async () => {
await db.execute({ sql: "INSERT INTO metadata_collections (id,query,max_videos,depth,status,state,lease_until,created_at,updated_at) VALUES ('restart','resume',1,0,'running',?,0,0,0)", args: [JSON.stringify({ ids: ['99999999999'], queries: [], pending: [{ id: '99999999999', level: 0 }] })] });
calls = []; await runner.runNext(); const job = (await (await get('/api/admin/collections')).json()).jobs.find(j => j.id === 'restart');
expect(job.status).toBe('complete'); expect(job.collected).toBe(1); expect(job.enriched).toBe(1); expect(calls).toHaveLength(1);
});
test('storage totals count table aggregates once and use actual disk capacity', async () => {
await upsertMedia('aaaaaaaaaaa', { status: 'ready', size: 2048, duration: 120 });
await db.execute("INSERT INTO listening_daily (fingerprint,day,video_id,plays) VALUES ('listener','2026-10-03','aaaaaaaaaaa',3)");
const result = await storageAnalytics([{ label: 'DB', path: root }, { label: 'Absent', path: root + '/missing' }]);
expect(Number(result.media.find(r => r.status === 'ready').bytes)).toBe(2048); expect(Number(result.listening.plays)).toBe(3); expect(result.volumes[0].free).toBeGreaterThan(0); expect(result.volumes[1].unavailable).toBe(true);
expect(result.sources.find(r => r.source === 'collector')).toBeDefined(); expect((await get('/api/admin/metadata?q=Related&offset=0')).status).toBe(200); expect((await get('/api/admin/metadata?q=%25')).status).toBe(200);
});
test('extraction errors remain visible while discovered cards stay stored', async () => {
const failureApp = new Hono(); const failed = registerAnalyticsRoutes(failureApp, { adminAuth, autoStart: false, runYtdlp: async a => a[0].startsWith('ytsearch') ? JSON.stringify({ id: 'failure0001', title: 'Discovered before failure' }) : Promise.reject(Error('Unavailable video')) });
await post('/api/admin/collections', { query: 'failure', maxVideos: 1, depth: 0 }); await failed.runNext();
const job = (await (await get('/api/admin/collections')).json()).jobs.find(j => j.query === 'failure');
expect(job.status).toBe('complete'); expect(job.failed).toBe(1); expect(job.collected).toBe(1); expect(job.errors[0].error).toBe('Unavailable video');
});
test('duplicate search cards do not consume slots before unique results', async () => {
const a = { id: 'unique00001', title: 'Unique A' }, b = { id: 'unique00002', title: 'Unique B' };
const worker = registerAnalyticsRoutes(new Hono(), { adminAuth, autoStart: false, runYtdlp: async args => args[0].startsWith('ytsearch') ? [a,a,b].map(JSON.stringify).join('\n') : JSON.stringify(args.at(-1).endsWith(a.id) ? a : b) });
await post('/api/admin/collections', { query: 'duplicates', maxVideos: 2, depth: 0 }); await worker.runNext();
const job = (await (await get('/api/admin/collections')).json()).jobs.find(j => j.query === 'duplicates');
expect(job.collected).toBe(2); expect(job.enriched).toBe(2);
});
test('cancelling an in-flight search prevents storing its returned cards', async () => {
let started, release;
const began = new Promise(resolve => { started = resolve; }), pause = new Promise(resolve => { release = resolve; });
const worker = registerAnalyticsRoutes(new Hono(), { adminAuth, autoStart: false, runYtdlp: async () => { started(); await pause; return JSON.stringify({ id: 'cancel00001', title: 'Should not save' }); } });
const queued = await (await post('/api/admin/collections', { query: 'cancel while searching', maxVideos: 1, depth: 0 })).json();
const work = worker.runNext(); await began; await post('/api/admin/collections/' + queued.id + '/cancel'); release(); await work;
const job = (await (await get('/api/admin/collections')).json()).jobs.find(j => j.id === queued.id);
expect(job.status).toBe('cancelled'); expect(job.failed).toBe(0);
expect((await db.execute("SELECT * FROM video_meta WHERE id='cancel00001'")).rows).toHaveLength(0);
});

View File

@@ -157,6 +157,28 @@ export async function initDb() {
);
CREATE INDEX IF NOT EXISTS idx_video_channel ON video_channels (channel,updated_at DESC,video_id);
CREATE TABLE IF NOT EXISTS video_details (
video_id TEXT PRIMARY KEY, metadata TEXT NOT NULL, updated_at INTEGER NOT NULL
);
CREATE TABLE IF NOT EXISTS metadata_collections (
id TEXT PRIMARY KEY, query TEXT NOT NULL, max_videos INTEGER NOT NULL, depth INTEGER NOT NULL,
status TEXT NOT NULL DEFAULT 'queued', state TEXT NOT NULL DEFAULT '{}', error TEXT,
lease_until INTEGER NOT NULL DEFAULT 0, owner TEXT,
created_at INTEGER NOT NULL, updated_at INTEGER NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_collections_queue ON metadata_collections (status,created_at);
CREATE TABLE IF NOT EXISTS lyric_transcriptions (
id TEXT PRIMARY KEY, video_id TEXT NOT NULL, status TEXT NOT NULL DEFAULT 'queued',
base_rev INTEGER NOT NULL DEFAULT 0, result TEXT, error TEXT, stage TEXT NOT NULL DEFAULT 'queued',
attempts INTEGER NOT NULL DEFAULT 0, lease_token TEXT, lease_until INTEGER NOT NULL DEFAULT 0,
created_at INTEGER NOT NULL, updated_at INTEGER NOT NULL
);
CREATE UNIQUE INDEX IF NOT EXISTS idx_transcription_active ON lyric_transcriptions (video_id)
WHERE status IN ('queued','running');
CREATE INDEX IF NOT EXISTS idx_transcription_queue ON lyric_transcriptions (status,created_at);
CREATE TABLE IF NOT EXISTS lyrics_worker_state (id INTEGER PRIMARY KEY CHECK(id=1), last_seen INTEGER NOT NULL);
-- Shared per-video documents (kind = lyrics | chapters), visible to every
-- user. The live copy is here; every save also lands in video_note_revs
-- as a full snapshot, which is the server-side backup and undo history.

View File

@@ -42,6 +42,8 @@ import { createHash, createHmac, randomBytes, timingSafeEqual } from 'node:crypt
import { mkdirSync, readdirSync, unlinkSync, writeFileSync, readFileSync } from 'node:fs';
import { join } from 'node:path';
import { getCookie, setCookie, deleteCookie } from 'hono/cookie';
import { registerAnalyticsRoutes } from './admin-analytics.js';
import { registerTranscriptionRoutes } from './transcriptions.js';
export const NOTE_KINDS = new Set(['lyrics', 'chapters']);
// A YouTube id or one of the server's own uploads (see uploads.js).
@@ -59,6 +61,12 @@ function cleanText(v, max) {
return String(v == null ? '' : v).replace(/[\u0000-\u001f\u007f]+/g, ' ').trim().slice(0, max);
}
function cleanLyricText(value) {
return String(value ?? '').replace(/\r\n?/g, '\n')
.replace(/[\u0000-\u0009\u000b-\u001f\u007f]+/g, ' ')
.split('\n').map(part => part.trim()).filter(Boolean).join('\n').trim().slice(0, MAX_LINE_CHARS);
}
function cleanTime(v) {
if (v === null || v === undefined || v === '') return null;
const n = Number(v);
@@ -76,7 +84,7 @@ export function sanitizeLyrics(input) {
const lines = [];
for (const l of rawLines) {
if (!l || typeof l !== 'object') continue;
const text = cleanText(l.text, MAX_LINE_CHARS);
const text = cleanLyricText(l.text);
if (!text) continue;
const kind = l.kind === 'section' || l.kind === 'cue' ? l.kind : 'line';
lines.push({ t: cleanTime(l.t), text, kind });
@@ -426,7 +434,7 @@ export function registerNoteRoutes(app, deps) {
if (!who) return c.json({ ok: false, error: 'link an online profile to report a lyric line' }, 401);
if (who.invalid) return c.json({ ok: false, error: 'invalid API token' }, 401);
if (flagOverBudget(who.by)) return c.json({ ok: false, error: 'too many reports — wait a few minutes' }, 429);
const text = cleanText(body.text, MAX_LINE_CHARS);
const text = cleanLyricText(body.text);
if (!text) return c.json({ ok: false, error: 'which line? (text is missing)' }, 400);
try {
const lyr = (await db.getNotes(id)).lyrics;
@@ -662,6 +670,18 @@ export function registerNoteRoutes(app, deps) {
await next();
};
if (deps.analytics) registerAnalyticsRoutes(app, { adminAuth: requireAdminOrToken, runYtdlp, ...deps.analytics });
registerTranscriptionRoutes(app, {
db, adminAuth: requireAdminOrToken, sanitizeLyrics,
workerEnabled: !!workerToken && workerToken.length >= 24,
workerAuth: async (c, next) => {
const token = (c.req.header('authorization') || '').match(/^Bearer\s+(\S+)$/i)?.[1];
if (!workerToken || workerToken.length < 24 || !token || !safeEqual(token, workerToken)) return c.json({ ok: false, error: 'lyrics worker token required' }, 401);
await next();
},
});
// Saved (server-cached) videos with whether each already has lyrics — the
// work list for batch lyric injection.
app.get('/api/admin/media', requireAdminOrToken, async (c) => {

View File

@@ -84,6 +84,12 @@ describe('pure helpers', () => {
expect(d.offset).toBe(30);
});
test('lyrics retain tight line breaks while other metadata strips controls', () => {
const result = N.sanitizeLyrics({ lines: [{ t: 12, text: ' Because You are God\r\n You can do anything\u0007 ', kind: 'line' }], tags: ['Key\nG'] });
expect(result.lines).toEqual([{ t: 12, text: 'Because You are God\nYou can do anything', kind: 'line' }]);
expect(result.tags).toEqual(['Key G']);
});
test('sanitizeChapters requires time + title and sorts', () => {
const d = N.sanitizeChapters({ items: [{ t: 30, title: 'B' }, { t: 5, title: 'A', note: 'n' }, { t: null, title: 'no time' }] });
expect(d.items).toEqual([{ t: 5, title: 'A', note: 'n' }, { t: 30, title: 'B', note: '' }]);
@@ -245,3 +251,19 @@ describe('routes', () => {
expect(await res.text()).toContain('admin');
});
});
test('saving a grouped cue preserves one timestamp in current and revision documents', async () => {
const id = 'groupedCue1';
const lines = [{ t: 12, text: 'Because You are God\nYou can do anything', kind: 'line' }];
const response = await app.request(`/api/notes/${id}/lyrics`, json('PUT', { data: { lines }, baseRev: 0, profile: 'josh' }));
expect(response.status).toBe(200);
const current = await (await app.request(`/api/notes/${id}`)).json();
expect(current.lyrics.data.lines).toEqual(lines);
const revision = await (await app.request(`/api/notes/${id}/lyrics/revs/1`)).json();
expect(revision.data.lines).toEqual(lines);
const flag = await app.request(`/api/notes/${id}/flags`, json('POST', { text: lines[0].text, profile: 'josh', reason: 'words' }));
expect(flag.status).toBe(200);
const flags = await (await app.request(`/api/notes/${id}/flags`)).json();
expect(flags.flags[0].text).toBe(lines[0].text);
});

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 ./remote.test.js && bun test ./party.test.js && bun test ./uploads.test.js && bun test ./innertube.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 ./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"
},
"dependencies": {
"@hono/node-server": "^1.14.0",

View File

@@ -84,7 +84,7 @@ export function registerCatalogRoutes(app, { resolveListener = async (c, name, f
? c.req.raw.clone() : null;
await next();
if (!c.res.ok) return;
const discovery = c.req.method === 'GET' && !path.startsWith('/api/catalog/') && path !== '/api/recommendations';
const discovery = c.req.method === 'GET' && !/^\/api\/admin\/(analytics|metadata|collections)(\/|$)/.test(path) && !path.startsWith('/api/catalog/') && path !== '/api/recommendations';
try {
if (discovery && c.res.headers.get('content-type')?.includes('application/json')) {
const result = await c.res.clone().json();

View File

@@ -2328,6 +2328,11 @@ const notes = registerNoteRoutes(app, {
backupDir: pathJoin(dirname(process.env.DB_PATH || './data/ytplayer.db'), 'backups'),
adminHtmlPath: './public/admin.html',
workerToken: process.env.LYRICS_WORKER_TOKEN || '',
analytics: { paths: [
{ label: 'Video cache', path: MEDIA_DIR },
{ label: 'Uploads', path: process.env.UPLOAD_DIR || pathJoin(dirname(process.env.DB_PATH || './data/ytplayer.db'), 'uploads') },
{ label: 'Database', path: dirname(process.env.DB_PATH || './data/ytplayer.db') },
] },
});
// ============================================================================

79
server/transcriptions.js Normal file
View File

@@ -0,0 +1,79 @@
// Explicit Whisper requests produce reviewable drafts, never published notes.
import { randomUUID } from 'node:crypto';
import { db as sql } from './db.js';
const ID = /^([\w-]{11}|upl_[a-f0-9]{12})$/;
const LEASE_MS = 120000;
function publicJob(row) {
if (!row) return null;
return { id: row.id, videoId: row.video_id, status: row.status, stage: row.stage,
baseRev: Number(row.base_rev), result: row.result ? JSON.parse(row.result) : null,
error: row.error, createdAt: Number(row.created_at), updatedAt: Number(row.updated_at) };
}
export function registerTranscriptionRoutes(app, { db, adminAuth, workerAuth, workerEnabled, sanitizeLyrics, now = Date.now }) {
const seen = () => sql.execute({ sql: 'INSERT INTO lyrics_worker_state (id,last_seen) VALUES (1,?) ON CONFLICT(id) DO UPDATE SET last_seen=excluded.last_seen', args: [now()] });
async function source(id) {
if (id.startsWith('upl_')) return db.getUpload(id);
const row = await db.getMedia(id); return row?.status === 'ready' ? row : null;
}
app.post('/api/admin/transcriptions/:video', adminAuth, async c => {
const id = c.req.param('video');
if (!ID.test(id)) return c.json({ ok: false, error: 'invalid video id' }, 400);
if (!workerEnabled) return c.json({ ok: false, error: 'Whisper is unavailable: the lyrics worker is not configured.' }, 503);
const audio = await source(id);
if (!audio) return c.json({ ok: false, error: 'Save this video on the server before transcribing its audio.' }, 409);
if (Number(audio.duration) > 3600) return c.json({ ok: false, error: 'Whisper requests are limited to one hour of audio.' }, 422);
const notes = await db.getNotes(id);
// One active job per song and a bounded queue, including across restarts.
const existing = (await sql.execute({ sql: "SELECT * FROM lyric_transcriptions WHERE video_id=? AND status IN ('queued','running')", args: [id] })).rows[0];
if (existing) return c.json({ ok: true, job: publicJob(existing) });
const active = Number((await sql.execute("SELECT COUNT(*) AS n FROM lyric_transcriptions WHERE status IN ('queued','running')")).rows[0].n);
if (active >= 25) return c.json({ ok: false, error: 'The transcription queue is full. Try again after a job finishes.' }, 429);
const time = now();
await sql.execute({ sql: `INSERT OR IGNORE INTO lyric_transcriptions (id,video_id,base_rev,created_at,updated_at)
SELECT ?,?,?,?,? WHERE (SELECT COUNT(*) FROM lyric_transcriptions WHERE status IN ('queued','running')) < 25`,
args: [randomUUID(), id, notes.lyrics?.rev || 0, time, time] });
const row = (await sql.execute({ sql: "SELECT * FROM lyric_transcriptions WHERE video_id=? AND status IN ('queued','running')", args: [id] })).rows[0];
return row ? c.json({ ok: true, job: publicJob(row) }, 202) : c.json({ ok: false, error: 'The transcription queue is full.' }, 429);
});
app.get('/api/admin/transcriptions/:video', adminAuth, async c => {
const id = c.req.param('video');
if (!ID.test(id)) return c.json({ ok: false, error: 'invalid video id' }, 400);
const row = (await sql.execute({ sql: 'SELECT * FROM lyric_transcriptions WHERE video_id=? ORDER BY created_at DESC,rowid DESC LIMIT 1', args: [id] })).rows[0];
const worker = (await sql.execute('SELECT last_seen FROM lyrics_worker_state WHERE id=1')).rows[0];
return c.json({ ok: true, job: publicJob(row), workerOnline: !!worker && now() - Number(worker.last_seen) < 90000, enabled: !!workerEnabled }, 200, { 'Cache-Control': 'no-store' });
});
app.post('/api/lyrics-worker/claim', workerAuth, async c => {
await seen(); const time = now();
await sql.execute({ sql: "UPDATE lyric_transcriptions SET status='failed',stage='failed',error='The worker stopped repeatedly. Please retry.',updated_at=? WHERE status='running' AND lease_until<? AND attempts>=3", args: [time, time] });
// UPDATE RETURNING makes claiming atomic for multiple worker processes.
const row = (await sql.execute({ sql: `UPDATE lyric_transcriptions SET status='running',stage='loading-model',attempts=attempts+1,
lease_token=?,lease_until=?,updated_at=? WHERE id=(SELECT id FROM lyric_transcriptions
WHERE (status='queued' OR (status='running' AND lease_until<?)) AND attempts<3 ORDER BY created_at LIMIT 1) RETURNING *`,
args: [randomUUID(), time + LEASE_MS, time, time] })).rows[0];
return c.json({ ok: true, job: row ? { ...publicJob(row), lease: row.lease_token,
audioPath: row.video_id.startsWith('upl_') ? `/api/uploads/${row.video_id}` : `/api/media/${row.video_id}?a=1` } : null });
});
app.post('/api/lyrics-worker/jobs/:job', workerAuth, async c => {
let body;
try { const raw = await c.req.text(); if (raw.length > 400000) return c.json({ ok: false, error: 'transcript too large' }, 413); body = JSON.parse(raw); }
catch { return c.json({ ok: false, error: 'invalid JSON' }, 400); }
if (!body || typeof body.lease !== 'string') return c.json({ ok: false, error: 'missing lease' }, 400);
const time = now();
let result = null, status = 'running', error = null;
let stage = ['loading-model', 'downloading-audio', 'transcribing'].includes(body.stage) ? body.stage : 'transcribing';
if (body.status === 'complete') {
try { result = sanitizeLyrics(body.result); if (!result.lines.length) throw Error('No sung lyrics were found.'); }
catch (e) { return c.json({ ok: false, error: e.message }, 400); }
result.tags = ['auto-transcribed (whisper)']; status = 'complete'; stage = 'complete';
} else if (body.status === 'failed') { status = 'failed'; stage = 'failed'; error = String(body.error || 'Transcription failed.').slice(0, 300); }
const changed = await sql.execute({ sql: `UPDATE lyric_transcriptions SET status=?,stage=?,result=?,error=?,lease_until=?,updated_at=?
WHERE id=? AND status='running' AND lease_token=? AND lease_until>=?`,
args: [status, stage, result ? JSON.stringify(result) : null, error, time + LEASE_MS, time, c.req.param('job'), body.lease, time] });
if (!changed.rowsAffected) return c.json({ ok: false, error: 'This job lease has expired or finished.' }, 409);
await seen();
// Bound completed draft retention; active jobs are never deleted.
await sql.execute({ sql: "DELETE FROM lyric_transcriptions WHERE status IN ('complete','failed') AND updated_at<?", args: [time - 30 * 86400000] });
return c.json({ ok: true });
});
}

View File

@@ -0,0 +1,97 @@
import { test, expect, beforeAll, afterAll } from 'bun:test';
import { mkdtempSync, rmSync, writeFileSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { Hono } from 'hono';
const root = mkdtempSync(join(tmpdir(), 'ytp-transcriptions-'));
process.env.DB_PATH = join(root, 'test.db');
const db = await import('./db.js');
const { registerNoteRoutes } = await import('./notes.js');
const VID = '0gfX0dFLaBc';
const TOKEN = 'worker-test-token-'.repeat(3);
const worker = { Authorization: 'Bearer ' + TOKEN };
const json = (body = {}, headers = {}) => ({ method: 'POST', headers: { 'Content-Type': 'application/json', ...headers }, body: JSON.stringify(body) });
let app, cookie, job, lease;
beforeAll(async () => {
await db.initDb(); await db.upsertMedia(VID, { status: 'ready', duration: 240 });
await db.saveNote({ videoId: VID, kind: 'lyrics', baseRev: 0, source: 'user', updatedBy: 'fixture', data: { lines: [{ t: 1, text: 'Original saved lyrics', kind: 'line' }], tags: [], offset: 0 } });
writeFileSync(join(root, 'admin.html'), '<title>admin</title>');
app = new Hono();
registerNoteRoutes(app, { db, getProfile: async () => null, profileNameRe: /^[\w-]{3,40}$/, runYtdlp: async () => '{}', adminPassword: 'test-password', workerToken: TOKEN, backupDir: join(root, 'backups'), adminHtmlPath: join(root, 'admin.html') });
const login = await app.request('/api/admin/login', json({ password: 'test-password' }));
cookie = login.headers.get('set-cookie').split(';')[0];
});
afterAll(() => { db.db.close(); rmSync(root, { recursive: true, force: true }); });
const admin = () => ({ Cookie: cookie });
const draft = { lines: [{ t: 2, text: 'Because You are God', kind: 'line' }, { t: 6, text: 'You can do anything', kind: 'line' }], tags: ['wrong source'], offset: 0 };
test('admin and worker authentication are enforced separately', async () => {
expect((await app.request('/api/admin/transcriptions/' + VID, json())).status).toBe(401);
expect((await app.request('/api/admin/transcriptions/' + VID)).status).toBe(401);
expect((await app.request('/api/lyrics-worker/claim', json({}, admin()))).status).toBe(401);
});
test('two button requests deduplicate and preserve published lyrics', async () => {
const first = await app.request('/api/admin/transcriptions/' + VID, json({}, admin()));
expect(first.status).toBe(202); job = (await first.json()).job;
const second = await (await app.request('/api/admin/transcriptions/' + VID, json({}, admin()))).json();
expect(second.job.id).toBe(job.id); expect(job.baseRev).toBe(1);
expect((await db.getNotes(VID)).lyrics.data.lines[0].text).toBe('Original saved lyrics');
await db.initDb();
expect((await (await app.request('/api/admin/transcriptions/' + VID, { headers: admin() })).json()).job.id).toBe(job.id);
});
test('claiming is atomic and does not expose leases to the admin UI', async () => {
const responses = await Promise.all([app.request('/api/lyrics-worker/claim', json({}, worker)), app.request('/api/lyrics-worker/claim', json({}, worker))]);
const jobs = await Promise.all(responses.map(r => r.json()));
expect(jobs.filter(r => r.job)).toHaveLength(1); lease = jobs.find(r => r.job).job.lease;
expect(jobs.find(r => r.job).job.audioPath).toBe('/api/media/' + VID + '?a=1');
const view = await (await app.request('/api/admin/transcriptions/' + VID, { headers: admin() })).json();
expect(view.job.lease).toBeUndefined(); expect(view.workerOnline).toBe(true);
});
test('heartbeat advances stages and invalid leases cannot finish jobs', async () => {
expect((await app.request('/api/lyrics-worker/jobs/' + job.id, json({ lease: 'wrong', status: 'complete', result: draft }, worker))).status).toBe(409);
expect((await app.request('/api/lyrics-worker/jobs/' + job.id, json({ lease, stage: 'transcribing' }, worker))).status).toBe(200);
const view = await (await app.request('/api/admin/transcriptions/' + VID, { headers: admin() })).json();
expect(view.job.stage).toBe('transcribing');
});
test('completion returns a sanitized review draft without overwriting saved lyrics', async () => {
expect((await app.request('/api/lyrics-worker/jobs/' + job.id, json({ lease, status: 'complete', result: draft }, worker))).status).toBe(200);
const view = await (await app.request('/api/admin/transcriptions/' + VID, { headers: admin() })).json();
expect(view.job.status).toBe('complete'); expect(view.job.result.lines).toHaveLength(2);
expect(view.job.result.tags).toEqual(['auto-transcribed (whisper)']);
expect((await db.getNotes(VID)).lyrics.rev).toBe(1);
expect((await app.request('/api/lyrics-worker/jobs/' + job.id, json({ lease, status: 'failed' }, worker))).status).toBe(409);
});
test('expired jobs are reclaimed with a new lease; stale results are rejected', async () => {
const queued = await (await app.request('/api/admin/transcriptions/' + VID, json({}, admin()))).json();
expect(queued.job.id).not.toBe(job.id); job = queued.job;
const first = (await (await app.request('/api/lyrics-worker/claim', json({}, worker))).json()).job;
await db.db.execute({ sql: 'UPDATE lyric_transcriptions SET lease_until=0 WHERE id=?', args: [job.id] });
const reclaimed = (await (await app.request('/api/lyrics-worker/claim', json({}, worker))).json()).job;
expect(reclaimed.lease).not.toBe(first.lease);
expect((await app.request('/api/lyrics-worker/jobs/' + job.id, json({ lease: first.lease, status: 'complete', result: draft }, worker))).status).toBe(409);
lease = reclaimed.lease;
});
test('empty transcripts and oversized drafts are rejected, and failures retain saved lyrics', async () => {
expect((await app.request('/api/lyrics-worker/jobs/' + job.id, json({ lease, status: 'complete', result: { lines: [] } }, worker))).status).toBe(400);
expect((await app.request('/api/lyrics-worker/jobs/' + job.id, json({ lease, result: 'x'.repeat(400001) }, worker))).status).toBe(413);
expect((await app.request('/api/lyrics-worker/jobs/' + job.id, json({ lease, status: 'failed', error: 'No sung vocals were detected.' }, worker))).status).toBe(200);
const view = await (await app.request('/api/admin/transcriptions/' + VID, { headers: admin() })).json();
expect(view.job.error).toBe('No sung vocals were detected.'); expect((await db.getNotes(VID)).lyrics.rev).toBe(1);
});
test('uncached songs, invalid IDs, and overlong audio give usable errors', async () => {
expect((await app.request('/api/admin/transcriptions/nope', json({}, admin()))).status).toBe(400);
expect((await app.request('/api/admin/transcriptions/aaaaaaaaaaa', json({}, admin()))).status).toBe(409);
await db.upsertMedia('bbbbbbbbbbb', { status: 'ready', duration: 4000 });
expect((await app.request('/api/admin/transcriptions/bbbbbbbbbbb', json({}, admin()))).status).toBe(422);
});
test('uploads use the upload media endpoint and repeated worker crashes fail visibly', async () => {
const id = 'upl_0123456789ab';
await db.createUpload({ id, kind: 'audio', title: 'Fixture upload', duration: 30, ext: 'm4a', mime: 'audio/mp4', size: 100 });
await app.request('/api/admin/transcriptions/' + id, json({}, admin()));
const claim = (await (await app.request('/api/lyrics-worker/claim', json({}, worker))).json()).job;
expect(claim.audioPath).toBe('/api/uploads/' + id);
await db.db.execute({ sql: 'UPDATE lyric_transcriptions SET lease_until=0,attempts=3 WHERE id=?', args: [claim.id] });
await app.request('/api/lyrics-worker/claim', json({}, worker));
const view = await (await app.request('/api/admin/transcriptions/' + id, { headers: admin() })).json();
expect(view.job.status).toBe('failed');
});

View File

@@ -2,7 +2,7 @@
import { db } from './db.js';
const VIDEO_ID = /^[\w-]{11}$/;
const SOURCES = new Set(['search', 'search-cache', 'client-search', 'channel', 'streams', 'playlist', 'profile', 'sync', 'related', 'backfill']);
const SOURCES = new Set(['search', 'search-cache', 'client-search', 'channel', 'streams', 'playlist', 'profile', 'sync', 'related', 'backfill', 'collector']);
const text = (v, n = 300) => typeof v === 'string' ? v.trim().slice(0, n) : '';
const positive = v => Number.isFinite(Number(v)) && Number(v) > 0 ? Number(v) : 0;
const canonicalThumb = id => `https://i.ytimg.com/vi/${id}/hqdefault.jpg`;
@@ -90,6 +90,7 @@ export async function trimCatalog(max = Number(process.env.VIDEO_META_MAX) || 50
'DELETE FROM video_meta_sources WHERE video_id NOT IN (SELECT id FROM video_meta)',
'DELETE FROM video_thumbnails WHERE video_id NOT IN (SELECT id FROM video_meta)',
'DELETE FROM video_channels WHERE video_id NOT IN (SELECT id FROM video_meta)',
'DELETE FROM video_details WHERE video_id NOT IN (SELECT id FROM video_meta)',
], 'write');
return count - max;
}