Persistent server search cache: 200 results per query, 250k queries, 15-day stale-while-revalidate, refresh=1
This commit is contained in:
11
server/db.js
11
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.
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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: '' });
|
||||
});
|
||||
|
||||
60
server/search-cache.js
Normal file
60
server/search-cache.js
Normal file
@@ -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 };
|
||||
}
|
||||
20
server/search-cache.test.js
Normal file
20
server/search-cache.test.js
Normal file
@@ -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);
|
||||
});
|
||||
@@ -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,26 +397,22 @@ 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);
|
||||
yt = await innertube.searchDeep(q, { limit: searchCacheDb.PER_QUERY });
|
||||
} catch (e) {
|
||||
console.warn(`[search] innertube failed, using yt-dlp: ${e.message}`);
|
||||
}
|
||||
@@ -428,37 +425,50 @@ function fetchSearch(q) {
|
||||
], { pooled: true });
|
||||
yt = parseCards(out);
|
||||
}
|
||||
const results = [...mine, ...yt];
|
||||
const fetchedAt = Date.now();
|
||||
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;
|
||||
}
|
||||
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=<query>
|
||||
// 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=<query>[&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);
|
||||
}
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user