Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460import type { Env, ShardInfo, UrlEntry, WorkerState } from './types';import { extractUrlAction } from './hydrant';import { flushShardToR2, formatShardFilename, getStaticUrls, loadShardEntries, splitEntriesByLimits, updateSitemapIndex,} from './sitemap';import { buildUrlsetXml, measureXmlBytes, PROTOCOL_MAX_BYTES, PROTOCOL_MAX_URLS, sanitizeLastmod,} from './xml';
const KEY_WORKER_STATE = 'worker_state';const KEY_IDENTITY_CACHE = 'identity_cache';
export function parseConfigLimits(env: Env): { maxUrls: number; maxBytes: number } { const parsedUrls = parseInt(env.SHARD_MAX_URLS || '', 10); const parsedBytes = parseInt(env.SHARD_MAX_BYTES || '', 10);
const maxUrls = !isNaN(parsedUrls) && parsedUrls > 0 ? Math.min(parsedUrls, PROTOCOL_MAX_URLS) : PROTOCOL_MAX_URLS;
const maxBytes = !isNaN(parsedBytes) && parsedBytes > 0 ? Math.min(parsedBytes, PROTOCOL_MAX_BYTES) : PROTOCOL_MAX_BYTES;
return { maxUrls, maxBytes };}
export default { async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> { const url = new URL(request.url); const pathname = url.pathname;
if (pathname === '/sitemap.xml') { const obj = await env.SITEMAP_BUCKET.get('sitemap.xml'); if (!obj) { return new Response('Sitemap index not found. Has sync run yet?', { status: 404 }); }
return new Response(obj.body, { headers: { 'Content-Type': 'application/xml; charset=utf-8', 'Cache-Control': 'public, max-age=86400, s-maxage=86400', }, }); }
if (pathname.startsWith('/sitemaps/')) { const filename = pathname.replace(/^\/sitemaps\//, ''); const obj = await env.SITEMAP_BUCKET.get(`sitemaps/${filename}`); if (!obj) { return new Response('Sitemap shard not found', { status: 404 }); }
return new Response(obj.body, { headers: { 'Content-Type': 'application/x-gzip', 'Content-Encoding': 'gzip', 'Cache-Control': 'public, max-age=3600, stale-while-revalidate=86400', }, }); }
return new Response('Not Found', { status: 404 }); },
async scheduled(event: ScheduledEvent, env: Env, ctx: ExecutionContext): Promise<void> { console.log(`[sitemap] Cron trigger fired at ${new Date().toISOString()} (cron: ${event.cron})`); ctx.waitUntil((async () => { try { const res = await syncFromHydrant(env); console.log(`[sitemap] Sync finished: success=${res.success}, events=${res.syncedEvents}, mutated=${res.mutatedUrls}, totalUrls=${res.totalUrls}, cursor=${res.newCursor}${res.error ? `, error=${res.error}` : ''}`); } catch (err) { console.error('[sitemap] Unhandled error during scheduled sync:', err); } })()); },};
export async function loadWorkerState(kv: KVNamespace): Promise<WorkerState> { const raw = await kv.get(KEY_WORKER_STATE); if (!raw) { return { cursor: '0', lastSyncAt: 'never', totalUrls: 0, shards: [], urlLocator: {}, }; }
try { const parsed = JSON.parse(raw); return { cursor: typeof parsed.cursor === 'string' ? parsed.cursor : '0', lastSyncAt: typeof parsed.lastSyncAt === 'string' ? parsed.lastSyncAt : 'never', totalUrls: typeof parsed.totalUrls === 'number' ? parsed.totalUrls : 0, shards: Array.isArray(parsed.shards) ? parsed.shards : [], urlLocator: typeof parsed.urlLocator === 'object' && parsed.urlLocator !== null ? parsed.urlLocator : {}, }; } catch { return { cursor: '0', lastSyncAt: 'never', totalUrls: 0, shards: [], urlLocator: {}, }; }}
export async function loadIdentityMap(kv: KVNamespace): Promise<Map<string, string>> { const raw = await kv.get(KEY_IDENTITY_CACHE); if (!raw) return new Map(); try { const record = JSON.parse(raw); return new Map(Object.entries(record)); } catch { return new Map(); }}
export async function saveIdentityMap(kv: KVNamespace, map: Map<string, string>): Promise<void> { const obj = Object.fromEntries(map.entries()); await kv.put(KEY_IDENTITY_CACHE, JSON.stringify(obj));}
export async function syncFromHydrant(env: Env): Promise<{ success: boolean; syncedEvents: number; mutatedUrls: number; newCursor: string; totalUrls: number; error?: string;}> { const baseUrl = env.BASE_URL || 'https://tangled.org'; const { maxUrls, maxBytes } = parseConfigLimits(env); const now = new Date();
const state = await loadWorkerState(env.SITEMAP_KV); const identityMap = await loadIdentityMap(env.SITEMAP_KV);
// Initialize shard 1 if needed if (state.shards.length === 0) { state.shards.push({ index: 1, filename: formatShardFilename(1), count: 0, uncompressedBytes: 0, sealed: false, lastmod: sanitizeLastmod(undefined, now), }); }
let activeShard = state.shards[state.shards.length - 1]; if (activeShard.sealed) { const nextIndex = activeShard.index + 1; activeShard = { index: nextIndex, filename: formatShardFilename(nextIndex), count: 0, uncompressedBytes: 0, sealed: false, lastmod: sanitizeLastmod(undefined, now), }; state.shards.push(activeShard); }
// Load active shard entries const activeEntries = await loadShardEntries(env.SITEMAP_BUCKET, activeShard.filename);
// Seed static URLs on initial empty setup let mutated = false; if (activeShard.index === 1 && activeEntries.size === 0) { const staticUrls = getStaticUrls(baseUrl, now); for (const entry of staticUrls) { activeEntries.set(entry.loc, entry); state.urlLocator[entry.loc] = activeShard.index; } mutated = true; }
// Fetch events with bounded pagination const fetchResult = await fetchHydrantStreamPages(env, state.cursor); if (!fetchResult.ok) { return { success: false, syncedEvents: 0, mutatedUrls: 0, newCursor: state.cursor, totalUrls: state.totalUrls, error: fetchResult.error, }; }
const hydrantEvents = fetchResult.events; let highestCursor = state.cursor; let eventCount = 0; let mutationCount = 0;
// Track historical shards modified by in-place updates or deletes: shardIndex -> { upserts, deletes } const historicalModifications = new Map< number, { upserts: Map<string, UrlEntry>; deletes: Set<string> } >();
function getHistoricalMod(shardIndex: number) { let mod = historicalModifications.get(shardIndex); if (!mod) { mod = { upserts: new Map(), deletes: new Set() }; historicalModifications.set(shardIndex, mod); } return mod; }
for (const evt of hydrantEvents) { highestCursor = String(Math.max(Number(highestCursor), evt.id)); eventCount++;
const action = extractUrlAction(evt, baseUrl, identityMap, now); if (!action) continue;
if (action.type === 'delete') { const targetShardIndex = state.urlLocator[action.url]; if (!targetShardIndex) continue;
if (targetShardIndex === activeShard.index) { if (activeEntries.delete(action.url)) { delete state.urlLocator[action.url]; mutated = true; mutationCount++; } } else { const mod = getHistoricalMod(targetShardIndex); mod.deletes.add(action.url); mod.upserts.delete(action.url); delete state.urlLocator[action.url]; mutated = true; mutationCount++; } } else if (action.type === 'upsert' && action.entry) { const existingShardIndex = state.urlLocator[action.url];
if (!existingShardIndex || existingShardIndex === activeShard.index) { // New URL or already in active shard activeEntries.set(action.url, action.entry); state.urlLocator[action.url] = activeShard.index; mutated = true; mutationCount++; } else { // In-place update in historical shard! (finding #4) const mod = getHistoricalMod(existingShardIndex); mod.upserts.set(action.url, action.entry); mod.deletes.delete(action.url); mutated = true; mutationCount++; } } }
// Apply in-place modifications to historical shards for (const [shardIndex, mod] of historicalModifications.entries()) { const shardMeta = state.shards.find((s) => s.index === shardIndex); if (!shardMeta) continue;
const shardEntries = await loadShardEntries(env.SITEMAP_BUCKET, shardMeta.filename); for (const url of mod.deletes) { shardEntries.delete(url); } for (const [url, entry] of mod.upserts) { shardEntries.set(url, entry); }
const updatedBytes = await flushShardToR2( env.SITEMAP_BUCKET, shardMeta.filename, Array.from(shardEntries.values()), ); shardMeta.count = shardEntries.size; shardMeta.uncompressedBytes = updatedBytes; shardMeta.lastmod = sanitizeLastmod(undefined, now); }
// Enforce both maxUrls and maxBytes rollover limits (finding #1) let entriesArray = Array.from(activeEntries.values()); while (true) { const currentBytes = measureXmlBytes(buildUrlsetXml(entriesArray)); if (entriesArray.length <= maxUrls && currentBytes <= maxBytes) { break; }
const { chunk, remainder } = splitEntriesByLimits(entriesArray, maxUrls, maxBytes); if (chunk.length === 0 || remainder.length === 0) { break; }
activeShard.count = chunk.length; activeShard.sealed = true; activeShard.lastmod = sanitizeLastmod(undefined, now); activeShard.uncompressedBytes = await flushShardToR2( env.SITEMAP_BUCKET, activeShard.filename, chunk, );
for (const entry of chunk) { state.urlLocator[entry.loc] = activeShard.index; }
// Open next shard const nextIndex = activeShard.index + 1; activeShard = { index: nextIndex, filename: formatShardFilename(nextIndex), count: remainder.length, uncompressedBytes: 0, sealed: false, lastmod: sanitizeLastmod(undefined, now), }; state.shards.push(activeShard); entriesArray = remainder; mutated = true; }
// Only flush active shard and index if changes actually occurred if (mutated) { activeShard.count = entriesArray.length; activeShard.lastmod = sanitizeLastmod(undefined, now); activeShard.uncompressedBytes = await flushShardToR2( env.SITEMAP_BUCKET, activeShard.filename, entriesArray, );
for (const entry of entriesArray) { state.urlLocator[entry.loc] = activeShard.index; }
await updateSitemapIndex(env.SITEMAP_BUCKET, state.shards, baseUrl); }
// Compute total URL count state.totalUrls = state.shards.reduce((acc, s) => acc + s.count, 0); state.cursor = highestCursor;
// Commit ordering: identity map saved FIRST, then worker state is canonical commit marker (finding #3) if (mutated || eventCount > 0) { state.lastSyncAt = now.toISOString(); await saveIdentityMap(env.SITEMAP_KV, identityMap); await env.SITEMAP_KV.put(KEY_WORKER_STATE, JSON.stringify(state)); }
return { success: true, syncedEvents: eventCount, mutatedUrls: mutationCount, newCursor: state.cursor, totalUrls: state.totalUrls, };}
export async function fetchHydrantStreamPages( env: Env, startCursor: string, maxEvents = 10000, idleTimeoutMs = 3000,): Promise<{ ok: boolean; events: any[]; error?: string }> { const client = env.HYDRANT ?? { fetch }; const endpoint = env.HYDRANT ? `http://hydrant.internal/stream?cursor=${startCursor}` : `${(env.HYDRANT_URL || 'https://api.tangled.org').replace(/\/+$/, '')}/stream?cursor=${startCursor}`;
try { const res = await client.fetch(endpoint, { headers: { Upgrade: 'websocket' }, });
const ws = (res as any).webSocket; if (!ws) { // Fallback for mocks / non-websocket responses in tests if (!res.ok) { const errText = await res.text(); return { ok: false, events: [], error: `Hydrant returned HTTP ${res.status}: ${errText}` }; } const body = await res.json(); const events: any[] = Array.isArray(body) ? body : (body as any)?.events ?? []; return { ok: true, events }; }
ws.accept();
return new Promise((resolve) => { const events: any[] = []; let timer: ReturnType<typeof setTimeout> | null = null;
function resetIdleTimer() { if (timer) clearTimeout(timer); timer = setTimeout(() => { cleanup(); resolve({ ok: true, events }); }, idleTimeoutMs); }
function cleanup() { if (timer) clearTimeout(timer); try { ws.close(1000, 'batch finished'); } catch {} }
ws.addEventListener('message', (event: any) => { try { const raw = typeof event.data === 'string' ? event.data : new TextDecoder().decode(event.data); const parsed = JSON.parse(raw); if (parsed && typeof parsed === 'object') { if (parsed.type === 'error') { cleanup(); resolve({ ok: false, events, error: parsed.error || 'Hydrant stream error' }); return; } events.push(parsed); if (events.length >= maxEvents) { cleanup(); resolve({ ok: true, events }); return; } } resetIdleTimer(); } catch (err: any) { console.error('[sitemap] Error parsing ws frame:', err); } });
ws.addEventListener('error', (err: any) => { cleanup(); resolve({ ok: events.length > 0, events, error: err?.message || 'WebSocket error' }); });
ws.addEventListener('close', () => { if (timer) clearTimeout(timer); resolve({ ok: true, events }); });
resetIdleTimer(); }); } catch (err: any) { return { ok: false, events: [], error: err?.message || 'Network error' }; }}