// One catalog for every discovery source. Media-file caching is independent. import { db } from './db.js'; import { mergeMissing } from './media-meta-core.js'; const VIDEO_ID = /^[\w-]{11}$/; const SOURCES = new Set(['search', 'search-cache', 'client-search', 'channel', 'streams', 'playlist', 'profile', 'sync', 'related', 'backfill', 'collector', 'device-play']); 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`; export const channelKey = c => String(c.channelId || c.channel || '').toLowerCase().replace(/\s*-\s*topic$|vevo$|\s+official$/g, '').trim(); function safeThumbnail(url) { try { const u = new URL(url); return u.protocol === 'https:' && /^(?:i|i\d)\.ytimg\.com$/.test(u.hostname) && !u.port && !u.username && !u.password; } catch { return false; } } export function normalizeCard(v) { if (!v || typeof v.id !== 'string' || !VIDEO_ID.test(v.id) || v.custom || v.upload) return null; const title = text(v.title); if (!title || title === '(untitled)') return null; let thumbnail = canonicalThumb(v.id); try { const u = new URL(v.thumbnail); // Never fetch arbitrary client URLs (including redirects or private hosts). if (safeThumbnail(u.href)) thumbnail = u.href; } catch { /* canonical YouTube art */ } const card = { id: v.id, title, channel: text(v.channel), channelId: text(v.channelId, 64), duration: positive(v.duration), thumbnail }; try { const u = new URL(v.channelUrl); if (u.protocol === 'https:' && /^(www\.)?youtube\.com$/.test(u.hostname)) card.channelUrl = u.href; } catch { /* optional */ } for (const key of ['tags', 'categories']) if (Array.isArray(v[key])) card[key] = v[key].slice(0, 30).map(x => text(x, 80)).filter(Boolean); if (v.description) card.description = text(v.description, 1200); if (positive(v.viewCount ?? v.view_count)) card.viewCount = positive(v.viewCount ?? v.view_count); return card; } export function extractCards(value, limit = 5000) { const cards = new Map(); let visited = 0; function visit(v, depth) { if (!v || typeof v !== 'object' || depth > 10 || ++visited > 40000 || cards.size >= limit) return; if (v.id && v.title) { const c = normalizeCard(v); if (c) { cards.set(c.id, c); return; } } for (const child of Object.values(v)) visit(child, depth + 1); } visit(value, 0); return [...cards.values()]; } let ingestion = Promise.resolve(); let writes = 0; export function ingest(cards, source = 'search', { fillMissing = false } = {}) { const task = ingestion.then(() => ingestBatch(cards, source, fillMissing)); ingestion = task.catch(() => {}); return task; } async function ingestBatch(cards, source, fillMissing) { source = SOURCES.has(source) ? source : 'playlist'; const unique = extractCards(cards); const now = Date.now(); // Chunking avoids enormous SQL batches on imported playlists/profiles. for (let i = 0; i < unique.length; i += 100) { const batch = unique.slice(i, i + 100); const existing = await db.execute({ sql: `SELECT id, card FROM video_meta WHERE id IN (${batch.map(() => '?').join(',')})`, args: batch.map(c => c.id) }); const old = new Map(existing.rows.map(r => [r.id, JSON.parse(r.card)])); await db.batch(batch.flatMap(c => { // Sparse playlist cards must not erase richer extractor metadata. const merged = fillMissing ? mergeMissing(old.get(c.id) || {}, c) : { ...(old.get(c.id) || {}), ...Object.fromEntries(Object.entries(c).filter(([, v]) => v !== '' && v !== 0 && (!Array.isArray(v) || v.length))) }; return [{ sql: `INSERT INTO video_meta (id,card,hay,seen,updated_at) VALUES (?,?,?,1,?) ON CONFLICT(id) DO UPDATE SET card=excluded.card,hay=excluded.hay,seen=seen+1,updated_at=excluded.updated_at`, args: [c.id, JSON.stringify(merged), `${merged.title} ${merged.channel || ''} ${(merged.tags || []).join(' ')}`.toLowerCase(), now] }, { sql: `INSERT INTO video_meta_sources (video_id,source,last_seen) VALUES (?,?,?) ON CONFLICT(video_id,source) DO UPDATE SET discoveries=discoveries+1,last_seen=excluded.last_seen`, args: [c.id, source, now] }, { sql: `INSERT INTO video_thumbnails (video_id,url) VALUES (?,?) ON CONFLICT(video_id) DO UPDATE SET retry_at=0 WHERE data IS NULL AND retry_at=9007199254740991`, args: [c.id, merged.thumbnail] }, { sql: 'INSERT INTO video_channels (video_id,channel,updated_at) VALUES (?,?,?) ON CONFLICT(video_id) DO UPDATE SET channel=excluded.channel,updated_at=excluded.updated_at', args: [c.id, channelKey(merged), now] }]; }), 'write'); } if (unique.length) kickThumbnails(); if (unique.length && ++writes % 100 === 0) await trimCatalog(); return unique.length; } export async function trimCatalog(max = Number(process.env.VIDEO_META_MAX) || 500000) { const count = Number((await db.execute('SELECT COUNT(*) AS n FROM video_meta')).rows[0].n); if (count <= max) return 0; await db.execute({ sql: `DELETE FROM video_meta WHERE id IN (SELECT id FROM video_meta ORDER BY EXISTS(SELECT 1 FROM listening_daily l WHERE l.video_id=video_meta.id),updated_at LIMIT ?)`, args: [count - max] }); await db.batch([ '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; } let running = null; let timer = null; let backfillRunning = false; const THUMB_BUDGET = Number(process.env.VIDEO_THUMB_MAX_BYTES) || 512 * 1024 ** 2; export async function drainThumbnails({ fetchImage = fetch, batchSize = 12, budget = THUMB_BUDGET } = {}) { const rows = (await db.execute({ sql: 'SELECT video_id,url,attempts FROM video_thumbnails WHERE data IS NULL AND retry_at<=? ORDER BY retry_at LIMIT ?', args: [Date.now(), batchSize] })).rows; // Two workers, bounded response size, timeout and no redirects. for (let i = 0; i < rows.length; i += 2) await Promise.all(rows.slice(i, i + 2).map(async r => { try { if (!safeThumbnail(r.url)) throw new Error('invalid thumbnail host'); const res = await fetchImage(r.url, { redirect: 'error', signal: AbortSignal.timeout(8000) }); const mime = (res.headers.get('content-type') || '').split(';')[0]; if (!res.ok || !['image/jpeg', 'image/png', 'image/webp'].includes(mime) || Number(res.headers.get('content-length')) > 1024 ** 2) throw new Error('invalid thumbnail'); if (!res.body) throw new Error('empty thumbnail'); const reader = res.body.getReader(); const chunks = []; let size = 0; try { for (;;) { const { done, value } = await reader.read(); if (done) break; size += value.byteLength; if (size > 1024 ** 2) throw new Error('thumbnail too large'); chunks.push(value); } } finally { await reader.cancel().catch(() => {}); } if (!size) throw new Error('empty thumbnail'); const bytes = new Uint8Array(size); let offset = 0; for (const chunk of chunks) { bytes.set(chunk, offset); offset += chunk.byteLength; } await db.execute({ sql: 'UPDATE video_thumbnails SET data=?,mime=?,size=?,fetched_at=?,attempts=0 WHERE video_id=?', args: [bytes, mime, size, Date.now(), r.video_id] }); } catch { const attempts = Number(r.attempts) + 1; await db.execute({ sql: 'UPDATE video_thumbnails SET attempts=?,retry_at=? WHERE video_id=?', args: [attempts, Date.now() + Math.min(86400000, 60000 * 2 ** Math.min(attempts, 10)), r.video_id] }); } })); let total = Number((await db.execute('SELECT COALESCE(SUM(size),0) AS n FROM video_thumbnails')).rows[0].n); if (total > budget) { const oldest = (await db.execute('SELECT video_id,size FROM video_thumbnails WHERE data IS NOT NULL ORDER BY fetched_at')).rows; for (const r of oldest) { if (total <= budget) break; // Re-ingestion can request it again; background retries never churn art // that the storage budget explicitly evicted. await db.execute({ sql: 'UPDATE video_thumbnails SET data=NULL,size=0,retry_at=9007199254740991 WHERE video_id=?', args: [r.video_id] }); total -= Number(r.size); } } return rows.length; } function kickThumbnails() { if (!timer || running) return; running = drainThumbnails().catch(e => console.warn('[catalog] thumbnails:', e.message)).finally(() => { running = null; }); } export function startThumbnails() { if (timer) return; const work = () => { kickThumbnails(); if (!backfillRunning) { backfillRunning = true; backfillStep().catch(e => console.warn('[catalog] backfill:', e.message)).finally(() => { backfillRunning = false; }); } }; timer = setInterval(work, 10000); timer.unref?.(); work(); } // A durable cursor per legacy table; one small page per tick, resumable after // a deploy. Capture the ceiling once so ongoing discoveries cannot prolong it. export async function backfillStep(pageSize = 30) { const tables = [ ['video_meta', 'card'], ['search_cache', 'results'], ['media_cache', 'meta'], ['video_history', null], ['playlists', 'data'], ['shared_playlists', 'data'], ['profiles', 'data'], ]; for (const [table, field] of tables) { await db.execute(`INSERT OR IGNORE INTO catalog_backfill (source,ceiling) SELECT '${table}',COALESCE(MAX(rowid),0) FROM ${table}`); const state = (await db.execute({ sql: 'SELECT cursor,ceiling FROM catalog_backfill WHERE source=?', args: [table] })).rows[0]; if (Number(state.cursor) >= Number(state.ceiling)) continue; const rows = (await db.execute({ sql: `SELECT rowid AS rid,${field || 'video_id AS id,title,channel,thumbnail,duration'}${table === 'media_cache' ? ',video_id' : ''} FROM ${table} WHERE rowid>? AND rowid<=? ORDER BY rowid LIMIT ?`, args: [state.cursor, state.ceiling, pageSize] })).rows; for (const row of rows) { let value = row; if (field) { try { value = JSON.parse(row[field]); } catch { continue; } } if (table === 'media_cache') value = { ...value, id: row.video_id }; await ingest(extractCards(value), 'backfill'); } await db.execute({ sql: 'UPDATE catalog_backfill SET cursor=? WHERE source=?', args: [rows.length ? rows.at(-1).rid : state.ceiling, table] }); return rows.length; } return 0; } export async function withLocalThumbnails(cards) { if (!cards.length) return cards; const rows = await db.execute({ sql: `SELECT video_id FROM video_thumbnails WHERE data IS NOT NULL AND video_id IN (${cards.map(() => '?').join(',')})`, args: cards.map(c => c.id) }); const stored = new Set(rows.rows.map(r => r.video_id)); return cards.map(c => stored.has(c.id) ? { ...c, thumbnail: `/api/catalog/${c.id}/thumbnail` } : c); } // Monotonic daily snapshots: repeated saves/reloads never multiply play counts. // A play is counted by StatsCore only after 30 seconds actually listened. export async function syncListening(fingerprint, stats) { if (!stats || !stats.days || typeof stats.days !== 'object') return; const cutoff = new Date(Date.now() - 400 * 86400000).toISOString().slice(0, 10); const tomorrow = new Date(Date.now() + 86400000).toISOString().slice(0, 10); const entries = []; for (const [day, data] of Object.entries(stats.days).sort().reverse()) { if (entries.length >= 10000) break; if (!/^\d{4}-\d{2}-\d{2}$/.test(day) || day < cutoff || day > tomorrow || !data?.songs) continue; const parsed = new Date(day + 'T12:00:00Z'); if (!Number.isFinite(parsed.getTime()) || parsed.toISOString().slice(0, 10) !== day) continue; for (const [id, count] of Object.entries(data.songs)) { if (!VIDEO_ID.test(id) || !Number.isInteger(count) || count <= 0) continue; if (entries.length >= 10000) break; entries.push({ sql: `INSERT INTO listening_daily (fingerprint,day,video_id,plays) VALUES (?,?,?,?) ON CONFLICT(fingerprint,day,video_id) DO UPDATE SET plays=excluded.plays WHERE excluded.plays>plays`, args: [fingerprint, day, id, Math.min(count, 10000)] }); } } for (let i = 0; i < entries.length; i += 200) await db.batch(entries.slice(i, i + 200), 'write'); await db.execute({ sql: 'DELETE FROM listening_daily WHERE day < ?', args: [cutoff] }); const meta = Object.entries(stats.meta && typeof stats.meta === 'object' ? stats.meta : {}).slice(0, 600).map(([id, m]) => ({ id, title: m?.t, channel: m?.c })); await ingest(meta, 'sync'); } export async function linkListening(device, profile) { await db.batch([ { sql: `INSERT INTO listening_daily (fingerprint,day,video_id,plays) SELECT ?,day,video_id,plays FROM listening_daily WHERE fingerprint=? ON CONFLICT(fingerprint,day,video_id) DO UPDATE SET plays=MAX(plays,excluded.plays)`, args: [profile, device] }, { sql: 'DELETE FROM listening_daily WHERE fingerprint=?', args: [device] }, ], 'write'); }