Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231import type { ExtractedAction } from './hydrant';import type { UrlEntry } from './types';import { buildUrlsetXml, measureXmlBytes } from './xml';
// static urls have no record behind themexport type Change = ExtractedAction | { type: 'upsert'; uri?: undefined; url: string; entry: UrlEntry };
export interface ShardLimits { maxUrls: number; maxBytes: number;}
// d1 allows 1000 queries a run and 100 params a query, so every batch goes// over as json arrays through json_each instead of a statement per rowconst CHUNK_BYTES = 512 * 1024;
const EMPTY_SHARD_BYTES = measureXmlBytes(buildUrlsetXml([]));
// the newline that joins it to the next url is counted tooexport function urlBytes(entry: UrlEntry): number { return measureXmlBytes(buildUrlsetXml([entry])) - EMPTY_SHARD_BYTES + 1;}
function chunked(items: unknown[]): string[] { const out: string[] = []; let chunk: string[] = []; let size = 0; for (const item of items) { const json = JSON.stringify(item); if (size + json.length > CHUNK_BYTES && chunk.length > 0) { out.push(`[${chunk.join(',')}]`); chunk = []; size = 0; } chunk.push(json); size += json.length + 1; } if (chunk.length > 0) out.push(`[${chunk.join(',')}]`); return out;}
async function lookup<T>(db: D1Database, sql: string, keys: string[]): Promise<T[]> { const unique = [...new Set(keys)]; const results = await Promise.all(chunked(unique).map((json) => db.prepare(sql).bind(json).all<T>())); return results.flatMap((res) => res.results);}
export async function getState(db: D1Database, key: string): Promise<string | undefined> { const row = await db.prepare('SELECT value FROM state WHERE key = ?1').bind(key).first<{ value: string }>(); return row?.value;}
function setState(db: D1Database, key: string, value: string): D1PreparedStatement { return db .prepare('INSERT INTO state (key, value) VALUES (?1, ?2) ON CONFLICT (key) DO UPDATE SET value = excluded.value') .bind(key, value);}
export async function countUrls(db: D1Database): Promise<number> { const row = await db.prepare('SELECT count(*) AS n FROM urls').first<{ n: number }>(); return row?.n ?? 0;}
// the newest shard takes new urls until it's full, older ones only change in placeasync function activeShard(db: D1Database): Promise<{ shard: number; count: number; bytes: number }> { const row = await db .prepare( `SELECT s.shard AS shard, count(u.url) AS count, coalesce(sum(u.bytes), 0) AS bytes FROM (SELECT coalesce(max(shard), 1) AS shard FROM shards) s LEFT JOIN urls u ON u.shard = s.shard`, ) .first<{ shard: number; count: number; bytes: number }>(); return { shard: row?.shard ?? 1, count: row?.count ?? 0, bytes: EMPTY_SHARD_BYTES + (row?.bytes ?? 0) };}
// applies one batch and moves the cursor in the same transaction, so a run// that dies halfway picks up exactly where the last batch landedexport async function applyChanges( db: D1Database, changes: Change[], limits: ShardLimits, cursor: string | undefined, now: Date,): Promise<number> { const uris = changes.flatMap((change) => (change.uri ? [change.uri] : [])); const recordUrl = new Map( (await lookup<{ uri: string; url: string }>( db, 'SELECT uri, url FROM records WHERE uri IN (SELECT value FROM json_each(?1))', uris, )).map((row) => [row.uri, row.url]), ); const urlShard = new Map<string, number | null>( (await lookup<{ url: string; shard: number }>( db, 'SELECT url, shard FROM urls WHERE url IN (SELECT value FROM json_each(?1))', [...changes.map((change) => change.url), ...recordUrl.values()], )).map((row) => [row.url, row.shard]), );
const active = await activeShard(db); const upserts = new Map<string, UrlEntry & { shard: number; bytes: number }>(); const records = new Map<string, string | null>(); const dirty = new Set<number>(); let mutations = 0;
const remove = (url: string): boolean => { const shard = urlShard.get(url); if (shard == null) return false; urlShard.set(url, null); upserts.delete(url); dirty.add(shard); return true; };
for (const change of changes) { const previous = change.uri ? recordUrl.get(change.uri) : undefined;
if (change.type === 'delete') { recordUrl.delete(change.uri); records.set(change.uri, null); if (remove(previous ?? change.url)) mutations++; continue; }
if (change.uri && previous !== change.url) { // a repo rename keeps its at-uri but changes its url if (previous) remove(previous); recordUrl.set(change.uri, change.url); records.set(change.uri, change.url); }
const bytes = urlBytes(change.entry); let shard = urlShard.get(change.url); if (shard == null) { if (active.count >= limits.maxUrls || active.bytes + bytes > limits.maxBytes) { active.shard++; active.count = 0; active.bytes = EMPTY_SHARD_BYTES; } active.count++; active.bytes += bytes; shard = active.shard; urlShard.set(change.url, shard); } upserts.set(change.url, { ...change.entry, shard, bytes }); dirty.add(shard); mutations++; }
const deletedUrls = [...urlShard].filter(([, shard]) => shard === null).map(([url]) => url); const statements: D1PreparedStatement[] = [ ...chunked([...upserts.values()]).map((json) => db .prepare( `INSERT INTO urls (url, shard, lastmod, changefreq, priority, bytes) SELECT value ->> 'loc', value ->> 'shard', value ->> 'lastmod', value ->> 'changefreq', value ->> 'priority', value ->> 'bytes' FROM json_each(?1) WHERE true ON CONFLICT (url) DO UPDATE SET shard = excluded.shard, lastmod = excluded.lastmod, changefreq = excluded.changefreq, priority = excluded.priority, bytes = excluded.bytes`, ) .bind(json), ), ...chunked(deletedUrls).map((json) => db.prepare('DELETE FROM urls WHERE url IN (SELECT value FROM json_each(?1))').bind(json), ), ...chunked([...records].filter(([, url]) => url !== null)).map((json) => db .prepare( `INSERT INTO records (uri, url) SELECT value ->> 0, value ->> 1 FROM json_each(?1) WHERE true ON CONFLICT (uri) DO UPDATE SET url = excluded.url`, ) .bind(json), ), ...chunked([...records].filter(([, url]) => url === null).map(([uri]) => uri)).map((json) => db.prepare('DELETE FROM records WHERE uri IN (SELECT value FROM json_each(?1))').bind(json), ), ...chunked([...dirty]).map((json) => db .prepare( `INSERT INTO shards (shard, lastmod, dirty) SELECT value, ?2, 1 FROM json_each(?1) WHERE true ON CONFLICT (shard) DO UPDATE SET dirty = dirty + 1`, ) .bind(json, now.toISOString().slice(0, 10)), ), ]; if (cursor !== undefined) statements.push(setState(db, 'cursor', cursor)); if (statements.length > 0) await db.batch(statements); return mutations;}
export async function dirtyShards(db: D1Database): Promise<{ shard: number; dirty: number }[]> { const res = await db.prepare('SELECT shard, dirty FROM shards WHERE dirty > 0 ORDER BY shard').all<{ shard: number; dirty: number }>(); return res.results;}
export async function shardEntries(db: D1Database, shard: number): Promise<UrlEntry[]> { const res = await db .prepare('SELECT url AS loc, lastmod, changefreq, priority FROM urls WHERE shard = ?1 ORDER BY url') .bind(shard) .all<{ loc: string; lastmod: string; changefreq: UrlEntry['changefreq'] | null; priority: number | null }>(); return res.results.map((row) => ({ loc: row.loc, lastmod: row.lastmod, changefreq: row.changefreq ?? undefined, priority: row.priority ?? undefined, }));}
// only clears the shard if nothing changed it while it was being renderedexport async function markRendered(db: D1Database, shard: number, dirty: number, lastmod: string): Promise<void> { await db .prepare('UPDATE shards SET dirty = dirty - ?2, lastmod = ?3 WHERE shard = ?1') .bind(shard, dirty, lastmod) .run();}
export async function allShards(db: D1Database): Promise<{ shard: number; lastmod: string }[]> { const res = await db.prepare('SELECT shard, lastmod FROM shards ORDER BY shard').all<{ shard: number; lastmod: string }>(); return res.results;}
export async function markSynced(db: D1Database, now: Date): Promise<void> { await setState(db, 'last_sync_at', now.toISOString()).run();}
export async function clearDb(db: D1Database): Promise<void> { await db.batch(['urls', 'records', 'shards', 'state'].map((table) => db.prepare(`DELETE FROM ${table}`)));}