Collect discovered video metadata and recommend videos from listening history
This commit is contained in:
23
server/db.js
23
server/db.js
@@ -134,6 +134,29 @@ export async function initDb() {
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_video_meta_upd ON video_meta (updated_at);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS video_meta_sources (
|
||||
video_id TEXT NOT NULL, source TEXT NOT NULL, discoveries INTEGER NOT NULL DEFAULT 1,
|
||||
last_seen INTEGER NOT NULL, PRIMARY KEY (video_id, source)
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS video_thumbnails (
|
||||
video_id TEXT PRIMARY KEY, url TEXT NOT NULL, data BLOB, mime TEXT,
|
||||
size INTEGER NOT NULL DEFAULT 0, fetched_at INTEGER NOT NULL DEFAULT 0,
|
||||
retry_at INTEGER NOT NULL DEFAULT 0, attempts INTEGER NOT NULL DEFAULT 0
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_thumb_pending ON video_thumbnails (retry_at) WHERE data IS NULL;
|
||||
CREATE TABLE IF NOT EXISTS listening_daily (
|
||||
fingerprint TEXT NOT NULL, day TEXT NOT NULL, video_id TEXT NOT NULL,
|
||||
plays INTEGER NOT NULL, PRIMARY KEY (fingerprint, day, video_id)
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_listening_video ON listening_daily (video_id, day);
|
||||
CREATE TABLE IF NOT EXISTS catalog_backfill (
|
||||
source TEXT PRIMARY KEY, cursor INTEGER NOT NULL DEFAULT 0, ceiling INTEGER NOT NULL
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS video_channels (
|
||||
video_id TEXT PRIMARY KEY, channel TEXT NOT NULL, updated_at INTEGER NOT NULL DEFAULT 0
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS idx_video_channel ON video_channels (channel,updated_at DESC,video_id);
|
||||
|
||||
-- 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.
|
||||
|
||||
@@ -6,7 +6,7 @@
|
||||
"scripts": {
|
||||
"start": "bun server.js",
|
||||
"dev": "bun --hot server.js",
|
||||
"test": "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 ./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",
|
||||
|
||||
129
server/recommendations.js
Normal file
129
server/recommendations.js
Normal file
@@ -0,0 +1,129 @@
|
||||
import { db } from './db.js';
|
||||
import { ingest, extractCards, withLocalThumbnails, channelKey, syncListening } from './video-catalog.js';
|
||||
|
||||
const tokens = s => new Set(String(s || '').toLowerCase().match(/[\p{L}\p{N}]{3,}/gu)?.filter(w => !['the', 'and', 'for', 'with', 'official', 'video', 'music', 'audio', 'lyrics', 'live'].includes(w)) || []);
|
||||
const channelOf = channelKey;
|
||||
|
||||
export function rankRecommendations(candidates, { limit = 12 } = {}) {
|
||||
const seeds = candidates.filter(c => c.personal > 0).sort((a, b) => b.personal - a.personal).slice(0, 20);
|
||||
const channels = new Map(), interests = new Map();
|
||||
for (const c of seeds) {
|
||||
const weight = Math.log2(1 + c.personal);
|
||||
const channel = channelOf(c.card);
|
||||
if (channel) channels.set(channel, (channels.get(channel) || 0) + weight);
|
||||
for (const w of tokens(`${c.card.title} ${(c.card.tags || []).join(' ')} ${(c.card.categories || []).join(' ')}`)) interests.set(w, (interests.get(w) || 0) + weight);
|
||||
}
|
||||
const ranked = candidates.map(c => {
|
||||
const words = tokens(`${c.card.title} ${(c.card.tags || []).join(' ')} ${(c.card.categories || []).join(' ')}`);
|
||||
let overlap = 0; for (const w of words) overlap += interests.get(w) || 0;
|
||||
overlap /= Math.max(1, words.size);
|
||||
const channel = channels.get(channelOf(c.card)) || 0;
|
||||
const popularity = Math.log2(1 + c.plays);
|
||||
const score = 3 * Math.log2(1 + c.personal) + 2 * Math.log2(1 + channel) + Math.log2(1 + overlap)
|
||||
+ popularity + 0.6 * Math.log2(1 + c.recent);
|
||||
const reason = c.personal > 0 ? 'One of your most played' : channel > 0 ? `More from ${c.card.channel || 'a channel you enjoy'}`
|
||||
: overlap > 0 ? 'Similar to your most played' : c.plays > 0 ? 'Popular with listeners' : 'Recently discovered';
|
||||
return { ...c.card, reason, score, personal: c.personal, plays: c.plays };
|
||||
}).sort((a, b) => b.score - a.score || a.id.localeCompare(b.id));
|
||||
const result = [], channelCounts = new Map();
|
||||
let familiar = 0;
|
||||
// Reserve room for discoveries when the catalog has them, and vary channels.
|
||||
const hasDiscovery = ranked.some(c => !c.personal);
|
||||
for (const card of ranked) {
|
||||
const channel = channelOf(card);
|
||||
if (channel && (channelCounts.get(channel) || 0) >= 3) continue;
|
||||
if (hasDiscovery && card.personal && familiar >= Math.ceil(limit / 2)) continue;
|
||||
result.push(card); if (card.personal) familiar++;
|
||||
if (channel) channelCounts.set(channel, (channelCounts.get(channel) || 0) + 1);
|
||||
if (result.length >= limit) break;
|
||||
}
|
||||
return result.map(({ score, personal, plays, ...card }) => card);
|
||||
}
|
||||
|
||||
export async function recommend(fingerprint = '', limit = 12) {
|
||||
const recent = new Date(Date.now() - 30 * 86400000).toISOString().slice(0, 10);
|
||||
// Popular/personal seeds plus recent discoveries; no media_cache join/filter.
|
||||
const rows = (await db.execute({ sql: `WITH popular AS (
|
||||
SELECT video_id FROM listening_daily GROUP BY video_id ORDER BY SUM(plays) DESC LIMIT 300
|
||||
), personal AS (
|
||||
SELECT video_id FROM listening_daily WHERE fingerprint=? GROUP BY video_id ORDER BY SUM(plays) DESC LIMIT 100
|
||||
), discovered AS (
|
||||
SELECT id AS video_id FROM video_meta ORDER BY updated_at DESC LIMIT 600
|
||||
), ids AS (
|
||||
SELECT video_id FROM popular UNION SELECT video_id FROM personal UNION SELECT video_id FROM discovered
|
||||
), plays AS (
|
||||
SELECT video_id,SUM(plays) AS plays,
|
||||
SUM(CASE WHEN fingerprint=? THEN plays ELSE 0 END) AS personal,
|
||||
SUM(CASE WHEN day>=? THEN plays ELSE 0 END) AS recent
|
||||
FROM listening_daily WHERE video_id IN (SELECT video_id FROM ids) GROUP BY video_id
|
||||
) SELECT m.card,m.hay,COALESCE(p.plays,0) AS plays,COALESCE(p.personal,0) AS personal,
|
||||
COALESCE(p.recent,0) AS recent FROM video_meta m LEFT JOIN plays p ON p.video_id=m.id
|
||||
WHERE m.id IN (SELECT video_id FROM ids)`, args: [fingerprint, fingerprint, recent] })).rows;
|
||||
const candidates = rows.map(r => ({ ...r, card: JSON.parse(r.card), plays: Number(r.plays), personal: Number(r.personal), recent: Number(r.recent) }));
|
||||
// Include older, unplayed videos from seed channels, even in a large catalog.
|
||||
const seedChannels = [...new Set(candidates.filter(c => c.personal > 0).sort((a, b) => b.personal - a.personal).slice(0, 10).map(c => channelOf(c.card)).filter(Boolean))];
|
||||
if (seedChannels.length) {
|
||||
const extra = (await db.execute({ sql: `WITH nearby AS (
|
||||
${seedChannels.map(() => 'SELECT video_id FROM (SELECT video_id FROM video_channels WHERE channel=? ORDER BY updated_at DESC LIMIT 30)').join(' UNION ')}
|
||||
) SELECT m.card,COALESCE(SUM(l.plays),0) AS plays,
|
||||
COALESCE(SUM(CASE WHEN l.fingerprint=? THEN l.plays ELSE 0 END),0) AS personal,
|
||||
COALESCE(SUM(CASE WHEN l.day>=? THEN l.plays ELSE 0 END),0) AS recent
|
||||
FROM nearby n JOIN video_meta m ON m.id=n.video_id LEFT JOIN listening_daily l ON l.video_id=n.video_id GROUP BY m.id`, args: [...seedChannels, fingerprint, recent] })).rows;
|
||||
const known = new Set(candidates.map(c => c.card.id));
|
||||
for (const r of extra) { const card = JSON.parse(r.card); if (!known.has(card.id)) { known.add(card.id); candidates.push({ card, plays: Number(r.plays), personal: Number(r.personal), recent: Number(r.recent) }); } }
|
||||
}
|
||||
return withLocalThumbnails(rankRecommendations(candidates, { limit }));
|
||||
}
|
||||
|
||||
// Capture cards from successful discovery responses, including future routes.
|
||||
// Incoming profile/sync payloads also carry cards absent from the response.
|
||||
export function registerCatalogRoutes(app, { resolveListener = async (c, name, fp) => name ? c.json({ ok: false, error: 'profile unavailable' }, 401) : fp } = {}) {
|
||||
app.use('/api/*', async (c, next) => {
|
||||
const path = c.req.path;
|
||||
const incoming = c.req.method === 'POST' && /^\/api\/(user\/sync|profile\/(create|save)|playlist\/share)$/.test(path)
|
||||
? 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';
|
||||
try {
|
||||
if (discovery && c.res.headers.get('content-type')?.includes('application/json')) {
|
||||
const result = await c.res.clone().json();
|
||||
if (typeof result.data === 'string') { try { result.data = JSON.parse(result.data); } catch { /* invalid profile */ } }
|
||||
await ingest(extractCards(result), path.includes('search') ? 'search-cache' : path.includes('channel') ? 'channel' : path.includes('streams') ? 'streams' : 'playlist');
|
||||
}
|
||||
if (incoming) {
|
||||
const body = await incoming.json();
|
||||
const payload = typeof body.data === 'string' ? JSON.parse(body.data) : body;
|
||||
await ingest(extractCards(payload), path.includes('profile') ? 'profile' : 'sync');
|
||||
if (path.startsWith('/api/profile/')) {
|
||||
const name = body.name || (await c.res.clone().json()).name;
|
||||
if (name) await syncListening('profile:' + name.toLowerCase(), body.data?.stats);
|
||||
}
|
||||
}
|
||||
} catch (e) { console.warn('[catalog] collection:', e.message); }
|
||||
});
|
||||
app.get('/api/recommendations', async c => {
|
||||
const fp = String(c.req.query('fp') || '').slice(0, 200);
|
||||
const limit = Math.max(1, Math.min(24, Math.floor(Number(c.req.query('limit')) || 12)));
|
||||
try {
|
||||
const listener = await resolveListener(c, c.req.query('profile') || '', fp);
|
||||
if (typeof listener !== 'string') return listener;
|
||||
return c.json({ ok: true, results: await recommend(listener, limit) }, 200, { 'Cache-Control': 'no-store' });
|
||||
}
|
||||
catch (e) { console.warn('[recommendations]', e.message); return c.json({ ok: false, error: 'Recommendations are unavailable right now.' }, 503); }
|
||||
});
|
||||
app.post('/api/catalog/collect', async c => {
|
||||
if (Number(c.req.header('content-length')) > 512000) return c.json({ ok: false, error: 'metadata too large' }, 413);
|
||||
let body;
|
||||
try { const raw = await c.req.text(); if (raw.length > 512000) return c.json({ ok: false, error: 'metadata too large' }, 413); body = JSON.parse(raw); }
|
||||
catch { return c.json({ ok: false, error: 'invalid JSON' }, 400); }
|
||||
if (!Array.isArray(body.cards) || body.cards.length > 500) return c.json({ ok: false, error: 'expected at most 500 cards' }, 400);
|
||||
try { return c.json({ ok: true, collected: await ingest(body.cards, 'client-search') }); }
|
||||
catch { return c.json({ ok: false, error: 'Could not save metadata.' }, 503); }
|
||||
});
|
||||
app.get('/api/catalog/:id/thumbnail', async c => {
|
||||
const r = (await db.execute({ sql: 'SELECT data,mime FROM video_thumbnails WHERE video_id=? AND data IS NOT NULL', args: [c.req.param('id')] })).rows[0];
|
||||
if (!r) return c.notFound();
|
||||
return new Response(r.data, { headers: { 'Content-Type': r.mime, 'Cache-Control': 'public, max-age=86400', 'X-Content-Type-Options': 'nosniff' } });
|
||||
});
|
||||
}
|
||||
191
server/recommendations.test.js
Normal file
191
server/recommendations.test.js
Normal file
@@ -0,0 +1,191 @@
|
||||
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-recommendations-'));
|
||||
process.env.DB_PATH = join(root, 'catalog.db');
|
||||
const { initDb, db } = await import('./db.js');
|
||||
const catalog = await import('./video-catalog.js');
|
||||
const { recommend, rankRecommendations, registerCatalogRoutes } = await import('./recommendations.js');
|
||||
const cache = await import('./search-cache.js');
|
||||
const a = { id: 'aaaaaaaaaaa', title: 'Quiet piano morning', channel: 'Piano studio', duration: 240, tags: ['piano', 'acoustic'] };
|
||||
const b = { id: 'bbbbbbbbbbb', title: 'Evening at the piano', channel: 'Piano studio', duration: 320 };
|
||||
const c = { id: 'ccccccccccc', title: 'Ocean documentary', channel: 'Nature films', duration: 900 };
|
||||
const day = new Date().toISOString().slice(0, 10);
|
||||
beforeAll(initDb);
|
||||
afterAll(() => { db.close(); rmSync(root, { recursive: true, force: true }); });
|
||||
const json = body => ({ method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify(body) });
|
||||
|
||||
test('normalizes metadata, excludes custom files, and refuses arbitrary thumbnail hosts', () => {
|
||||
expect(catalog.normalizeCard({ ...a, thumbnail: 'http://127.0.0.1/secrets' }).thumbnail).toBe('https://i.ytimg.com/vi/aaaaaaaaaaa/hqdefault.jpg');
|
||||
expect(catalog.normalizeCard({ ...a, thumbnail: 'https://i.ytimg.com:1234/private' }).thumbnail).toContain('/vi/aaaaaaaaaaa/');
|
||||
expect(catalog.normalizeCard({ ...a, id: '../anything' })).toBeNull();
|
||||
expect(catalog.normalizeCard({ ...a, custom: true })).toBeNull();
|
||||
expect(catalog.extractCards({ id: 'playlist1', title: 'My playlist', videos: [a, a, b] })).toHaveLength(2);
|
||||
});
|
||||
|
||||
test('all discovery sources preserve richer metadata on sparse updates', async () => {
|
||||
await catalog.ingest([a, b, c], 'search');
|
||||
await catalog.ingest([{ id: a.id, title: a.title }], 'playlist');
|
||||
const row = (await db.execute({ sql: 'SELECT card FROM video_meta WHERE id=?', args: [a.id] })).rows[0];
|
||||
expect(JSON.parse(row.card).tags).toEqual(a.tags);
|
||||
expect(JSON.parse(row.card).channel).toBe(a.channel);
|
||||
expect((await db.execute({ sql: 'SELECT source FROM video_meta_sources WHERE video_id=?', args: [a.id] })).rows.map(r => r.source).sort()).toEqual(['playlist', 'search']);
|
||||
});
|
||||
|
||||
test('daily listening snapshots are idempotent and monotonic, with malformed dates/counts ignored', async () => {
|
||||
const stats = { days: { [day]: { songs: { [a.id]: 8, [c.id]: 2, invalid: 999 } }, tomorrow: { songs: { [a.id]: 999 } } }, meta: {} };
|
||||
await catalog.syncListening('listener-a', stats);
|
||||
await catalog.syncListening('listener-a', stats);
|
||||
await catalog.syncListening('listener-a', { days: { [day]: { songs: { [a.id]: 3 } } } });
|
||||
const rows = (await db.execute('SELECT * FROM listening_daily')).rows;
|
||||
expect(rows).toHaveLength(2);
|
||||
expect(Number(rows.find(r => r.video_id === a.id).plays)).toBe(8);
|
||||
});
|
||||
|
||||
test('server ranks personal favorites and unplayed channel matches without a media cache', async () => {
|
||||
const cards = await recommend('listener-a');
|
||||
expect(cards[0].id).toBe(a.id);
|
||||
expect(cards.find(v => v.id === b.id).reason).toBe('More from Piano studio');
|
||||
expect((await db.execute('SELECT COUNT(*) AS n FROM media_cache')).rows[0].n).toBe(0);
|
||||
expect((await recommend())[0].id).toBe(a.id);
|
||||
});
|
||||
|
||||
test('ranking diversifies channels and mixes discoveries with most played', () => {
|
||||
const candidates = Array.from({ length: 20 }, (_, i) => ({ card: { ...a, id: String(i).padStart(11, '0'), channel: 'Channel ' + Math.floor(i / 5) }, personal: i < 10 ? 50 : 0, plays: 50 - i, recent: 5 }));
|
||||
const cards = rankRecommendations(candidates, { limit: 12 });
|
||||
expect(cards.filter(v => v.reason === 'One of your most played').length).toBeLessThanOrEqual(6);
|
||||
for (const channel of new Set(cards.map(v => v.channel))) expect(cards.filter(v => v.channel === channel).length).toBeLessThanOrEqual(3);
|
||||
expect(cards.some(v => v.reason !== 'One of your most played')).toBe(true);
|
||||
});
|
||||
|
||||
test('metadata and recommendations survive search-result cache eviction', async () => {
|
||||
await cache.put('piano', [a, b]);
|
||||
await cache.trim({ maxRows: 0 });
|
||||
expect(await cache.get('piano')).toBeNull();
|
||||
expect((await cache.searchVideos('piano')).length).toBe(2);
|
||||
expect((await recommend('listener-a')).some(v => v.id === b.id)).toBe(true);
|
||||
});
|
||||
|
||||
test('thumbnail bytes are persisted and served locally without an upstream redirect', async () => {
|
||||
const calls = [];
|
||||
await catalog.drainThumbnails({ fetchImage: async (url, options) => { calls.push({ url, options }); return new Response(new Uint8Array([255, 216, 255, 217]), { headers: { 'Content-Type': 'image/jpeg' } }); } });
|
||||
expect(calls).toHaveLength(3);
|
||||
expect(calls.every(c => c.options.redirect === 'error')).toBe(true);
|
||||
const app = new Hono(); registerCatalogRoutes(app);
|
||||
const res = await app.request('/api/catalog/' + a.id + '/thumbnail');
|
||||
expect(res.status).toBe(200);
|
||||
expect(res.headers.get('content-type')).toBe('image/jpeg');
|
||||
expect(new Uint8Array(await res.arrayBuffer())).toEqual(new Uint8Array([255, 216, 255, 217]));
|
||||
expect((await recommend('listener-a'))[0].thumbnail).toBe('/api/catalog/' + a.id + '/thumbnail');
|
||||
});
|
||||
|
||||
test('failed image jobs remain durable, retry later, and respect the thumbnail byte budget', async () => {
|
||||
const v = { ...a, id: 'ddddddddddd' }; await catalog.ingest([v], 'related');
|
||||
await catalog.drainThumbnails({ fetchImage: async () => { throw new Error('offline'); } });
|
||||
const row = (await db.execute({ sql: 'SELECT * FROM video_thumbnails WHERE video_id=?', args: [v.id] })).rows[0];
|
||||
expect(row.data).toBeNull(); expect(Number(row.attempts)).toBe(1); expect(Number(row.retry_at)).toBeGreaterThan(Date.now());
|
||||
let called = false;
|
||||
await catalog.drainThumbnails({ fetchImage: async () => { called = true; throw new Error('unexpected'); }, budget: 4 });
|
||||
expect(called).toBe(false);
|
||||
expect(Number((await db.execute('SELECT SUM(size) AS n FROM video_thumbnails')).rows[0].n)).toBeLessThanOrEqual(4);
|
||||
expect((await db.execute('SELECT COUNT(*) AS n FROM video_meta')).rows[0].n).toBe(4);
|
||||
});
|
||||
|
||||
test('successful cached searches, channels, streams, and playlists all feed the catalog', async () => {
|
||||
const app = new Hono(); registerCatalogRoutes(app);
|
||||
const paths = ['/api/search', '/api/channel', '/api/streams', '/api/playlist/shared', '/api/playlist/expand', '/api/profile/load', '/api/p2p/available'];
|
||||
for (let i = 0; i < paths.length; i++) {
|
||||
const card = { ...b, id: 'surface' + String(i).padStart(4, '0') };
|
||||
app.get(paths[i], c => c.json({ ok: true, data: { playlists: [{ videos: [card] }] } }));
|
||||
}
|
||||
for (let i = 0; i < paths.length; i++) {
|
||||
const card = { ...b, id: 'surface' + String(i).padStart(4, '0') };
|
||||
expect((await app.request(paths[i])).status).toBe(200);
|
||||
expect((await db.execute({ sql: 'SELECT id FROM video_meta WHERE id=?', args: [card.id] })).rows).toHaveLength(1);
|
||||
}
|
||||
});
|
||||
|
||||
test('profile and playlist POST bodies are collected only after successful writes', async () => {
|
||||
const app = new Hono(); registerCatalogRoutes(app);
|
||||
app.post('/api/profile/save', c => c.json({ ok: true }));
|
||||
app.post('/api/playlist/share', c => c.json({ ok: false }, 403));
|
||||
const card = { ...b, id: 'postprofile' };
|
||||
await app.request('/api/profile/save', json({ data: { playlists: [{ videos: [card] }] } }));
|
||||
expect((await db.execute({ sql: 'SELECT id FROM video_meta WHERE id=?', args: [card.id] })).rows).toHaveLength(1);
|
||||
const rejected = { ...b, id: 'deniedvideo' };
|
||||
await app.request('/api/playlist/share', json({ playlist: { videos: [rejected] } }));
|
||||
expect((await db.execute({ sql: 'SELECT id FROM video_meta WHERE id=?', args: [rejected.id] })).rows).toHaveLength(0);
|
||||
});
|
||||
|
||||
test('client cached search collection validates input and stores new metadata', async () => {
|
||||
const app = new Hono(); registerCatalogRoutes(app);
|
||||
const card = { ...b, id: 'cachedvideo' };
|
||||
const res = await app.request('/api/catalog/collect', json({ cards: [card] }));
|
||||
expect(await res.json()).toEqual({ ok: true, collected: 1 });
|
||||
expect((await app.request('/api/catalog/collect', json({ cards: Array(501).fill(card) }))).status).toBe(400);
|
||||
const picks = await (await app.request('/api/recommendations?fp=listener-a&limit=4')).json();
|
||||
expect(picks.results).toHaveLength(4);
|
||||
});
|
||||
|
||||
test('legacy backfill progresses with a durable cursor and finds unplayed search metadata', async () => {
|
||||
const card = { ...b, id: 'legacyvideo' };
|
||||
await cache.put('legacy', [card]);
|
||||
await db.execute({ sql: 'INSERT INTO media_cache (video_id,meta) VALUES (?,?)', args: ['legacymedia', JSON.stringify({ title: 'Older cached piano', channel: 'Piano studio', duration: 60 })] });
|
||||
for (let i = 0; i < 25; i++) await catalog.backfillStep(2);
|
||||
expect((await db.execute({ sql: 'SELECT id FROM video_meta WHERE id=?', args: [card.id] })).rows).toHaveLength(1);
|
||||
expect((await db.execute({ sql: 'SELECT id FROM video_meta WHERE id=?', args: ['legacymedia'] })).rows).toHaveLength(1);
|
||||
const state = (await db.execute("SELECT * FROM catalog_backfill WHERE source='search_cache'")).rows[0];
|
||||
expect(Number(state.cursor)).toBe(Number(state.ceiling));
|
||||
const before = (await db.execute('SELECT COUNT(*) AS n FROM video_meta_sources')).rows[0].n;
|
||||
await catalog.backfillStep(2);
|
||||
expect((await db.execute('SELECT COUNT(*) AS n FROM video_meta_sources')).rows[0].n).toBe(before);
|
||||
});
|
||||
|
||||
test('linking devices to one profile merges daily counts without multiplying synced history', async () => {
|
||||
const stats = { days: { [day]: { songs: { [a.id]: 20 } } } };
|
||||
await catalog.syncListening('device:one', stats); await catalog.syncListening('device:two', stats);
|
||||
await catalog.linkListening('device:one', 'profile:linked'); await catalog.linkListening('device:two', 'profile:linked');
|
||||
await catalog.syncListening('profile:linked', stats);
|
||||
const rows = (await db.execute("SELECT * FROM listening_daily WHERE fingerprint IN ('device:one','device:two','profile:linked')")).rows;
|
||||
expect(rows).toHaveLength(1); expect(Number(rows[0].plays)).toBe(20);
|
||||
});
|
||||
|
||||
test('profile recommendations use the profile gate and separate profile identities', async () => {
|
||||
const app = new Hono();
|
||||
registerCatalogRoutes(app, { resolveListener: async (c, name, fp) => {
|
||||
if (!name) return fp;
|
||||
if (c.req.header('x-profile-secret') !== 'test-secret') return c.json({ ok: false }, 401);
|
||||
return 'profile:' + name;
|
||||
} });
|
||||
expect((await app.request('/api/recommendations?profile=linked')).status).toBe(401);
|
||||
const linked = await (await app.request('/api/recommendations?profile=linked', { headers: { 'X-Profile-Secret': 'test-secret' } })).json();
|
||||
expect(linked.results.find(c => c.id === a.id).reason).toBe('One of your most played');
|
||||
const other = await (await app.request('/api/recommendations?profile=other', { headers: { 'X-Profile-Secret': 'test-secret' } })).json();
|
||||
expect(other.results.find(c => c.id === a.id).reason).not.toBe('One of your most played');
|
||||
});
|
||||
|
||||
test('concurrent sparse discovery does not erase metadata; corrupt thumbnail jobs cannot fetch private hosts', async () => {
|
||||
const card = { ...a, id: 'concurrent1' };
|
||||
await Promise.all([catalog.ingest([card]), catalog.ingest([{ id: card.id, title: card.title }], 'playlist')]);
|
||||
const row = (await db.execute({ sql: 'SELECT card FROM video_meta WHERE id=?', args: [card.id] })).rows[0];
|
||||
expect(JSON.parse(row.card).tags).toEqual(a.tags);
|
||||
await db.execute({ sql: 'UPDATE video_thumbnails SET url=?,retry_at=0,data=NULL WHERE video_id=?', args: ['http://127.0.0.1/private', card.id] });
|
||||
const urls = [];
|
||||
await catalog.drainThumbnails({ fetchImage: async url => { urls.push(url); return new Response(new Uint8Array([1]), { headers: { 'Content-Type': 'image/jpeg' } }); } });
|
||||
expect(urls.some(url => url.includes('127.0.0.1'))).toBe(false);
|
||||
});
|
||||
|
||||
test('indexed channel candidates include older unplayed matches beyond the recent discovery window', async () => {
|
||||
const seed = { id: 'oldchannel1', title: 'Favorite strings', channel: 'Rare ensemble' };
|
||||
const match = { id: 'oldchannel2', title: 'Archival strings', channel: 'Rare ensemble' };
|
||||
await catalog.ingest([seed, match]);
|
||||
await catalog.syncListening('rare-listener', { days: { [day]: { songs: { [seed.id]: 10 } } } });
|
||||
await catalog.ingest(Array.from({ length: 650 }, (_, i) => ({ id: 'archive' + String(i).padStart(4, '0'), title: 'Archive entry ' + i, channel: 'Other channel ' + i })), 'search');
|
||||
await db.execute({ sql: 'UPDATE video_meta SET updated_at=0 WHERE id=?', args: [match.id] });
|
||||
expect((await recommend('rare-listener')).some(c => c.id === match.id)).toBe(true);
|
||||
const plan = (await db.execute({ sql: 'EXPLAIN QUERY PLAN SELECT video_id FROM video_channels WHERE channel=? ORDER BY updated_at DESC LIMIT 30', args: ['piano studio'] })).rows;
|
||||
expect(plan.some(r => String(r.detail).includes('idx_video_channel'))).toBe(true);
|
||||
});
|
||||
@@ -2,6 +2,7 @@
|
||||
// repeat query — from any device, or after a restart — never waits on YouTube.
|
||||
// Capped by row count (default 250,000 queries) AND bytes (default 2 GiB), LRU.
|
||||
import { db } from './db.js';
|
||||
import { ingest, withLocalThumbnails, trimCatalog } from './video-catalog.js';
|
||||
|
||||
export const MAX_ROWS = Number(process.env.SEARCH_CACHE_MAX_ROWS) || 250_000;
|
||||
export const MAX_BYTES = Number(process.env.SEARCH_CACHE_MAX_BYTES) || 2 * 1024 ** 3;
|
||||
@@ -63,14 +64,7 @@ export async function stats() {
|
||||
export const MAX_VIDEOS = Number(process.env.VIDEO_META_MAX) || 500_000;
|
||||
|
||||
export async function rememberVideos(cards) {
|
||||
const now = Date.now();
|
||||
const rows = (cards || []).filter((c) => c && typeof c.id === 'string' && c.title);
|
||||
if (!rows.length) return;
|
||||
await db.batch(rows.map((c) => ({
|
||||
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(c), `${c.title} ${c.channel || ''}`.toLowerCase(), now],
|
||||
})));
|
||||
await ingest(cards, 'search');
|
||||
if (Math.random() < 0.02) trimVideos().catch(() => {});
|
||||
}
|
||||
|
||||
@@ -84,12 +78,9 @@ export async function searchVideos(q, limit = 60) {
|
||||
ORDER BY seen DESC, updated_at DESC LIMIT ?`,
|
||||
args: [...words.map(like), limit],
|
||||
});
|
||||
return r.rows.map((x) => { try { return JSON.parse(x.card); } catch { return null; } }).filter(Boolean);
|
||||
return withLocalThumbnails(r.rows.map((x) => { try { return JSON.parse(x.card); } catch { return null; } }).filter(Boolean));
|
||||
}
|
||||
|
||||
export async function trimVideos(max = MAX_VIDEOS) {
|
||||
const n = Number((await db.execute('SELECT COUNT(*) AS n FROM video_meta')).rows[0].n);
|
||||
if (n <= max) return 0;
|
||||
await db.execute({ sql: 'DELETE FROM video_meta WHERE id IN (SELECT id FROM video_meta ORDER BY updated_at ASC LIMIT ?)', args: [n - max] });
|
||||
return n - max;
|
||||
return trimCatalog(max);
|
||||
}
|
||||
|
||||
@@ -57,6 +57,8 @@ import { createP2pHub, holdersPayload, createRehydrator } from './p2p-hub.js';
|
||||
import { sha256Range } from './hash.js';
|
||||
import * as innertube from './innertube.js';
|
||||
import * as searchCacheDb from './search-cache.js';
|
||||
import { ingest as collectVideoMetadata, syncListening, linkListening, startThumbnails } from './video-catalog.js';
|
||||
import { registerCatalogRoutes } from './recommendations.js';
|
||||
import QRCode from 'qrcode';
|
||||
import { createYtdlpPool } from './ytdlp-pool.js';
|
||||
import { dirname, join as pathJoin } from 'node:path';
|
||||
@@ -339,6 +341,10 @@ function slimEntry(j) {
|
||||
channelUrl: pickChannelUrl(j),
|
||||
duration: typeof j.duration === 'number' ? j.duration : 0,
|
||||
thumbnail: `https://i.ytimg.com/vi/${id}/mqdefault.jpg`,
|
||||
tags: Array.isArray(j.tags) ? j.tags.slice(0, 30) : [],
|
||||
categories: Array.isArray(j.categories) ? j.categories.slice(0, 30) : [],
|
||||
description: typeof j.description === 'string' ? j.description.slice(0, 1200) : '',
|
||||
viewCount: typeof j.view_count === 'number' ? j.view_count : 0,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -405,6 +411,7 @@ function channelToUrl(c) {
|
||||
|
||||
const app = new Hono();
|
||||
app.use('*', logger());
|
||||
registerCatalogRoutes(app, { resolveListener: recommendationListener });
|
||||
|
||||
// ============================================================================
|
||||
// API routes
|
||||
@@ -458,7 +465,7 @@ function fetchYoutube(q) {
|
||||
const fetchedAt = Date.now();
|
||||
if (searchCache.size >= SEARCH_CACHE_MAX) searchCache.delete(searchCache.keys().next().value);
|
||||
searchCache.set(key, { yt, fetchedAt });
|
||||
if (yt.length) searchCacheDb.rememberVideos(yt).catch(() => {});
|
||||
if (yt.length) await searchCacheDb.rememberVideos(yt);
|
||||
if (yt.length) searchCacheDb.put(q, yt).catch((e) => console.warn(`[search] cache write: ${e.message}`));
|
||||
return { yt, fetchedAt };
|
||||
})().finally(() => inflightSearches.delete(key));
|
||||
@@ -648,6 +655,7 @@ async function resolveStreamsUncached(videoId) {
|
||||
|
||||
const out = await runYtdlpResilient(['-J', '--no-warnings', `https://www.youtube.com/watch?v=${videoId}`], { pooled: true });
|
||||
const info = JSON.parse(out);
|
||||
await collectVideoMetadata([slimEntry(info)], 'streams');
|
||||
setYtChapters(videoId, chaptersFromInfo(info)).catch(() => {});
|
||||
const raw = Array.isArray(info.formats) ? info.formats : [];
|
||||
const formats = [];
|
||||
@@ -1775,6 +1783,16 @@ app.get('/api/download/:videoId', async (c) => {
|
||||
const PROFILE_NAME_RE = /^[A-Za-z0-9][A-Za-z0-9_-]{2,39}$/;
|
||||
const PROFILE_MAX_BYTES = 2_000_000; // full data blob; typical payloads are ~KBs
|
||||
|
||||
async function recommendationListener(c, name, fp) {
|
||||
if (!name) return 'device:' + fp;
|
||||
name = String(name).trim().toLowerCase();
|
||||
if (!PROFILE_NAME_RE.test(name)) return c.json({ ok: false, error: 'invalid profile' }, 400);
|
||||
const row = await getProfile(name);
|
||||
if (!row) return c.json({ ok: false, error: 'profile not found' }, 404);
|
||||
const denied = await profileGate(c, name, row, c.req.header('x-profile-secret'));
|
||||
return denied || 'profile:' + name;
|
||||
}
|
||||
|
||||
const RAND_ADJ = ['amber', 'brave', 'calm', 'coral', 'crimson', 'dusty', 'gentle', 'golden',
|
||||
'hidden', 'ivory', 'jade', 'lunar', 'mellow', 'misty', 'noble', 'quiet',
|
||||
'rapid', 'silver', 'solar', 'stormy', 'swift', 'velvet', 'wild', 'zesty'];
|
||||
@@ -2200,11 +2218,15 @@ app.post('/api/user/sync', async (c) => {
|
||||
let body;
|
||||
try { body = await c.req.json(); } catch { return c.json({ ok: false, error: 'invalid JSON' }, 400); }
|
||||
|
||||
const fp = (body.fingerprint || '').trim();
|
||||
const fp = typeof body.fingerprint === 'string' ? body.fingerprint.trim().slice(0, 200) : '';
|
||||
if (!fp) return c.json({ ok: false, error: 'missing fingerprint' }, 400);
|
||||
|
||||
try {
|
||||
await upsertUser({ fingerprint: fp, appVersion: body.appVersion, playlists: body.playlists });
|
||||
const listener = await recommendationListener(c, body.profileName || '', fp);
|
||||
if (typeof listener !== 'string') return listener;
|
||||
if (body.profileName) await linkListening('device:' + fp, listener);
|
||||
await syncListening(listener, body.stats);
|
||||
if (body.recentVideo && body.recentVideo.id) {
|
||||
await recordVideoAccess(fp, body.recentVideo);
|
||||
}
|
||||
@@ -2508,6 +2530,7 @@ app.get('/*', indexHtml);
|
||||
// ============================================================================
|
||||
async function main() {
|
||||
await initDb();
|
||||
startThumbnails();
|
||||
await p2pDb.initP2pSchema();
|
||||
console.log(`[ytplayer] DB ready`);
|
||||
await media.init();
|
||||
|
||||
218
server/video-catalog.js
Normal file
218
server/video-catalog.js
Normal file
@@ -0,0 +1,218 @@
|
||||
// One catalog for every discovery source. Media-file caching is independent.
|
||||
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 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') {
|
||||
const task = ingestion.then(() => ingestBatch(cards, source));
|
||||
ingestion = task.catch(() => {});
|
||||
return task;
|
||||
}
|
||||
|
||||
async function ingestBatch(cards, source) {
|
||||
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 = { ...(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)',
|
||||
], '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');
|
||||
}
|
||||
Reference in New Issue
Block a user