diff --git a/server/db.js b/server/db.js index 88637a4..2cfdbdb 100644 --- a/server/db.js +++ b/server/db.js @@ -112,6 +112,17 @@ export async function initDb() { CREATE INDEX IF NOT EXISTS idx_media_lru ON media_cache (status, last_access); + -- Persistent search cache: the YouTube part of a search (up to 200 cards, + -- JSON), keyed by the lowercased query. Trimmed LRU by count and bytes. + CREATE TABLE IF NOT EXISTS search_cache ( + q TEXT PRIMARY KEY, + results TEXT NOT NULL, + size INTEGER NOT NULL DEFAULT 0, + fetched_at INTEGER NOT NULL DEFAULT 0, + last_access INTEGER NOT NULL DEFAULT 0 + ); + CREATE INDEX IF NOT EXISTS idx_search_cache_lru ON search_cache (last_access); + -- 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. diff --git a/server/innertube.js b/server/innertube.js index df1fcc0..929f1c3 100644 --- a/server/innertube.js +++ b/server/innertube.js @@ -8,10 +8,8 @@ export function lengthToSeconds(s) { return s.trim().split(':').map(Number).reduce((acc, n) => acc * 60 + n, 0); } -export function parseSearch(json) { - const sections = json?.contents?.twoColumnSearchResultsRenderer?.primaryContents - ?.sectionListRenderer?.contents; - if (!Array.isArray(sections)) throw new Error('innertube: unexpected response shape'); +// Cards from a list of itemSectionRenderer-style sections. +function cardsFromSections(sections) { const out = []; for (const s of sections) { for (const it of s?.itemSectionRenderer?.contents || []) { @@ -34,6 +32,37 @@ export function parseSearch(json) { return out; } +const tokenOf = (sections) => { + for (const s of sections || []) { + const t = s?.continuationItemRenderer?.continuationEndpoint?.continuationCommand?.token; + if (t) return t; + } + return ''; +}; + +export function parseSearch(json) { + const sections = json?.contents?.twoColumnSearchResultsRenderer?.primaryContents + ?.sectionListRenderer?.contents; + if (!Array.isArray(sections)) throw new Error('innertube: unexpected response shape'); + return cardsFromSections(sections); +} + +// First page plus the token for the next one ('' when there is none). +export function parseSearchPage(json) { + const sections = json?.contents?.twoColumnSearchResultsRenderer?.primaryContents + ?.sectionListRenderer?.contents; + if (!Array.isArray(sections)) throw new Error('innertube: unexpected response shape'); + return { cards: cardsFromSections(sections), next: tokenOf(sections) }; +} + +// A continuation response: more cards plus the next token. +export function parseContinuation(json) { + const items = (json?.onResponseReceivedCommands || []) + .flatMap((c) => c?.appendContinuationItemsAction?.continuationItems || []); + if (!Array.isArray(items) || !items.length) return { cards: [], next: '' }; + return { cards: cardsFromSections(items), next: tokenOf(items) }; +} + export async function search(q, { fetchImpl = fetch, timeoutMs = 6000 } = {}) { const res = await fetchImpl('https://www.youtube.com/youtubei/v1/search?prettyPrint=false', { method: 'POST', @@ -44,3 +73,34 @@ export async function search(q, { fetchImpl = fetch, timeoutMs = 6000 } = {}) { if (!res.ok) throw new Error(`innertube: HTTP ${res.status}`); return parseSearch(await res.json()); } + +// Up to `limit` results (default 200) by following continuation tokens. The pages +// are a chain, so they are fetched one after another; every page is optional — +// a failure part-way returns what was gathered. `onPage(cards)` lets a caller +// save progress. +export async function searchDeep(q, { limit = 200, fetchImpl = fetch, timeoutMs = 6000, maxPages = 14, onPage } = {}) { + const post = async (body) => { + const res = await fetchImpl('https://www.youtube.com/youtubei/v1/search?prettyPrint=false', { + method: 'POST', headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ context: { client: CLIENT }, ...body }), signal: AbortSignal.timeout(timeoutMs), + }); + if (!res.ok) throw new Error(`innertube: HTTP ${res.status}`); + return res.json(); + }; + const first = parseSearchPage(await post({ query: q })); + const seen = new Set(); + const all = []; + const add = (cards) => { for (const c of cards) if (!seen.has(c.id)) { seen.add(c.id); all.push(c); } }; + add(first.cards); + if (onPage) onPage(all.slice()); + let next = first.next; + for (let i = 0; i < maxPages && next && all.length < limit; i++) { + let page; + try { page = parseContinuation(await post({ continuation: next })); } catch { break; } + if (!page.cards.length) break; + add(page.cards); + if (onPage) onPage(all.slice()); + next = page.next; + } + return all.slice(0, limit); +} diff --git a/server/innertube.test.js b/server/innertube.test.js index 1f383a1..62dc2ce 100644 --- a/server/innertube.test.js +++ b/server/innertube.test.js @@ -1,6 +1,6 @@ import { test, expect } from 'bun:test'; import { readFileSync } from 'node:fs'; -import { parseSearch, lengthToSeconds, search } from './innertube.js'; +import { parseSearch, lengthToSeconds, search, searchDeep, parseContinuation } from './innertube.js'; const fixture = JSON.parse(readFileSync(new URL('./fixtures/innertube-search.json', import.meta.url))); test('parses video renderers into slim cards', () => { @@ -27,3 +27,23 @@ test('search() throws on HTTP error', async () => { const fetchImpl = async () => new Response('no', { status: 429 }); await expect(search('x', { fetchImpl })).rejects.toThrow('429'); }); + +const card = (id) => ({ videoRenderer: { videoId: id, title: { runs: [{ text: 't' + id }] }, ownerText: { runs: [{ text: 'c' }] } } }); +test('searchDeep follows continuation tokens, de-duplicates and stops at the limit', async () => { + const calls = []; + const fetchImpl = async (_u, init) => { + const body = JSON.parse(init.body); + calls.push(body.continuation || 'first'); + const mk = (ids, token) => [{ itemSectionRenderer: { contents: ids.map(card) } }, ...(token ? [{ continuationItemRenderer: { continuationEndpoint: { continuationCommand: { token } } } }] : [])]; + if (!body.continuation) return Response.json({ contents: { twoColumnSearchResultsRenderer: { primaryContents: { sectionListRenderer: { contents: mk(['a', 'b', 'c'], 'T1') } } } } }); + if (body.continuation === 'T1') return Response.json({ onResponseReceivedCommands: [{ appendContinuationItemsAction: { continuationItems: mk(['c', 'd', 'e'], 'T2') } }] }); + return Response.json({ onResponseReceivedCommands: [{ appendContinuationItemsAction: { continuationItems: mk(['f', 'g'], '') } }] }); + }; + const all = await searchDeep('x', { fetchImpl, limit: 200 }); + expect(all.map((c) => c.id)).toEqual(['a', 'b', 'c', 'd', 'e', 'f', 'g']); + expect(calls).toEqual(['first', 'T1', 'T2']); + expect((await searchDeep('x', { fetchImpl, limit: 4 })).length).toBe(4); +}); +test('parseContinuation of an empty response is empty', () => { + expect(parseContinuation({})).toEqual({ cards: [], next: '' }); +}); diff --git a/server/search-cache.js b/server/search-cache.js new file mode 100644 index 0000000..7fc5472 --- /dev/null +++ b/server/search-cache.js @@ -0,0 +1,60 @@ +// Persistent search cache (libsql). Holds the YouTube cards of a search so a +// 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'; + +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; +export const PER_QUERY = 200; + +let writes = 0; + +export async function get(q) { + const key = q.toLowerCase(); + const r = await db.execute({ sql: 'SELECT results, fetched_at FROM search_cache WHERE q = ?', args: [key] }); + const row = r.rows[0]; + if (!row) return null; + let results; + try { results = JSON.parse(row.results); } catch { return null; } + // Touch at most once a minute per row — reads must stay cheap. + db.execute({ sql: 'UPDATE search_cache SET last_access = ? WHERE q = ? AND last_access < ?', + args: [Date.now(), key, Date.now() - 60_000] }).catch(() => {}); + return { results, fetchedAt: Number(row.fetched_at) }; +} + +export async function put(q, results) { + const json = JSON.stringify(results.slice(0, PER_QUERY)); + const now = Date.now(); + await db.execute({ + sql: `INSERT INTO search_cache (q, results, size, fetched_at, last_access) VALUES (?,?,?,?,?) + ON CONFLICT(q) DO UPDATE SET results=excluded.results, size=excluded.size, + fetched_at=excluded.fetched_at, last_access=excluded.last_access`, + args: [q.toLowerCase(), json, json.length, now, now], + }); + if (++writes % 25 === 0) trim().catch(() => {}); +} + +// Drop least-recently-used rows until under both budgets. +export async function trim({ maxRows = MAX_ROWS, maxBytes = MAX_BYTES } = {}) { + const t = (await db.execute('SELECT COUNT(*) AS n, COALESCE(SUM(size),0) AS b FROM search_cache')).rows[0]; + let n = Number(t.n), b = Number(t.b); + if (n <= maxRows && b <= maxBytes) return 0; + let dropped = 0; + while (n > maxRows || b > maxBytes) { + const batch = (await db.execute('SELECT q, size FROM search_cache ORDER BY last_access ASC LIMIT 500')).rows; + if (!batch.length) break; + const drop = []; + for (const r of batch) { + if (n <= maxRows && b <= maxBytes) break; + drop.push(r.q); n -= 1; b -= Number(r.size); + } + await db.batch(drop.map((q) => ({ sql: 'DELETE FROM search_cache WHERE q = ?', args: [q] }))); + dropped += drop.length; + } + return dropped; +} + +export async function stats() { + const t = (await db.execute('SELECT COUNT(*) AS n, COALESCE(SUM(size),0) AS b FROM search_cache')).rows[0]; + return { queries: Number(t.n), bytes: Number(t.b), maxQueries: MAX_ROWS, maxBytes: MAX_BYTES }; +} diff --git a/server/search-cache.test.js b/server/search-cache.test.js new file mode 100644 index 0000000..7bafb5b --- /dev/null +++ b/server/search-cache.test.js @@ -0,0 +1,20 @@ +import { test, expect } from 'bun:test'; +import * as sc from './search-cache.js'; +import { initDb } from './db.js'; + +await initDb(); +const cards = (n) => Array.from({ length: n }, (_, i) => ({ id: 'v' + i, title: 't' + i })); + +test('put/get round-trips, case-insensitive, capped at 200', async () => { + await sc.put('Hello World', cards(300)); + const r = await sc.get('hello world'); + expect(r.results.length).toBe(200); + expect(await sc.get('nope')).toBeNull(); +}); + +test('trim evicts least-recently-used rows over the row cap', async () => { + for (let i = 0; i < 6; i++) await sc.put('q' + i, cards(3)); + const dropped = await sc.trim({ maxRows: 3 }); + expect(dropped).toBeGreaterThan(0); + expect((await sc.stats()).queries).toBe(3); +}); diff --git a/server/server.js b/server/server.js index 0d07ffe..a4dc4d2 100644 --- a/server/server.js +++ b/server/server.js @@ -56,6 +56,7 @@ import { registerIntakeRoutes } from './p2p-intake.js'; 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 QRCode from 'qrcode'; import { createYtdlpPool } from './ytdlp-pool.js'; import { dirname, join as pathJoin } from 'node:path'; @@ -396,69 +397,78 @@ app.get('/api/version', (c) => // resolveStreams()' streamCache below, just keyed by normalized query text // instead of videoId. Only a full success is cached, so a transient yt-dlp // failure still gets retried on the next request. -const SEARCH_CACHE_MAX = 200; +const SEARCH_CACHE_MAX = 200; // hot in-memory layer const SEARCH_CACHE_TTL_MS = 10 * 60_000; // fresh: answered from memory -const SEARCH_CACHE_STALE_MS = 6 * 3600_000; // stale-while-revalidate window -const searchCache = new Map(); // lowercased q -> { results, expiresAt } -const inflightSearches = new Map(); // lowercased q -> Promise<{ results, youtubeError? }> +const SEARCH_CACHE_STALE_MS = 15 * 24 * 3600_000; // persistent copy: stale-while-revalidate for 15 days +const searchCache = new Map(); // lowercased q -> { yt, fetchedAt } (YouTube cards only) +const inflightSearches = new Map(); // lowercased q -> Promise<{ yt, youtubeError? }> -// One search, shared by every caller asking the same thing at the same time. -function fetchSearch(q) { +// The YouTube side of one search (up to 200 cards), shared by every caller asking +// the same thing at the same time, and written to memory + the persistent cache. +function fetchYoutube(q) { const key = q.toLowerCase(); if (inflightSearches.has(key)) return inflightSearches.get(key); const p = (async () => { - // This server's own library first — and it still answers when YouTube - // (yt-dlp) is unreachable or rate-limited. - let mine = []; - try { mine = (await notesDb.listUploads({ q, limit: 20 })).map(uploads.card); } catch { /* library optional */ } - try { - let yt = []; - if (process.env.SEARCH_INNERTUBE !== '0') { - try { - yt = await innertube.search(q); - } catch (e) { - console.warn(`[search] innertube failed, using yt-dlp: ${e.message}`); - } + let yt = []; + if (process.env.SEARCH_INNERTUBE !== '0') { + try { + yt = await innertube.searchDeep(q, { limit: searchCacheDb.PER_QUERY }); + } catch (e) { + console.warn(`[search] innertube failed, using yt-dlp: ${e.message}`); } - if (!yt.length) { - const out = await runYtdlpResilient([ - `ytsearch${SEARCH_LIMIT}:${q}`, - '--dump-json', '--flat-playlist', - '--no-warnings', '--ignore-errors', - ], { pooled: true }); - yt = parseCards(out); - } - const results = [...mine, ...yt]; - if (searchCache.size >= SEARCH_CACHE_MAX) searchCache.delete(searchCache.keys().next().value); - searchCache.set(key, { results, expiresAt: Date.now() + SEARCH_CACHE_TTL_MS }); - return { results }; - } catch (err) { - if (mine.length) return { results: mine, youtubeError: err.message }; - throw err; } + if (!yt.length) { + const out = await runYtdlpResilient([ + `ytsearch${SEARCH_LIMIT}:${q}`, + '--dump-json', '--flat-playlist', + '--no-warnings', '--ignore-errors', + ], { pooled: true }); + yt = parseCards(out); + } + 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.put(q, yt).catch((e) => console.warn(`[search] cache write: ${e.message}`)); + return { yt, fetchedAt }; })().finally(() => inflightSearches.delete(key)); inflightSearches.set(key, p); return p; } -// GET /api/search?q= -// Fresh cache → instant. Stale (up to 6 h) → answered instantly from the old -// copy while a refresh runs in the background, so a repeat search never waits -// on yt-dlp/InnerTube. Only a full success is cached, so a transient failure -// still gets retried on the next request. +const searchLibrary = async (q) => { + // This server's own library first — it still answers when YouTube is unreachable. + try { return (await notesDb.listUploads({ q, limit: 20 })).map(uploads.card); } catch { return []; } +}; + +// GET /api/search?q=[&refresh=1] +// Memory → persistent cache (fresh 10 min, then stale-while-revalidate up to 15 +// days: answered instantly, refreshed in the background) → YouTube. `refresh=1` +// (the app's "Update search" button) skips every cache. Only a full success is +// cached, so a transient failure gets retried on the next request. app.get('/api/search', async (c) => { const q = (c.req.query('q') || '').trim(); if (!q) return c.json({ ok: false, error: 'empty query' }, 400); - const cached = searchCache.get(q.toLowerCase()); - const age = cached ? Date.now() - (cached.expiresAt - SEARCH_CACHE_TTL_MS) : Infinity; - if (cached && Date.now() < cached.expiresAt) return c.json({ ok: true, results: cached.results }); - if (cached && age < SEARCH_CACHE_STALE_MS) { - fetchSearch(q).catch(() => { /* keep serving the stale copy */ }); - return c.json({ ok: true, results: cached.results, stale: true }); + const key = q.toLowerCase(); + const mine = await searchLibrary(q); + const answer = (yt, extra = {}) => c.json({ ok: true, results: [...mine, ...yt], ...extra }); + if (c.req.query('refresh') !== '1') { + let hit = searchCache.get(key); + if (!hit) { + const row = await searchCacheDb.get(q).catch(() => null); + if (row) { hit = { yt: row.results, fetchedAt: row.fetchedAt }; searchCache.set(key, hit); } + } + const age = hit ? Date.now() - hit.fetchedAt : Infinity; + if (hit && age < SEARCH_CACHE_TTL_MS) return answer(hit.yt, { fetchedAt: hit.fetchedAt }); + if (hit && age < SEARCH_CACHE_STALE_MS) { + fetchYoutube(q).catch(() => { /* keep serving the stale copy */ }); + return answer(hit.yt, { stale: true, fetchedAt: hit.fetchedAt }); + } } try { - return c.json({ ok: true, ...(await fetchSearch(q)) }); + const { yt, fetchedAt } = await fetchYoutube(q); + return answer(yt, { fetchedAt }); } catch (err) { + if (mine.length) return c.json({ ok: true, results: mine, youtubeError: err.message }); return c.json({ ok: false, error: err.message }, 500); } });