import { WorkflowEntrypoint, type WorkflowEvent, type WorkflowStep, type WorkflowStepConfig } from 'cloudflare:workers'; import { applyChanges, countUrls, getState, markSynced } from './db'; import { COLLECTION_SPECS, extractUrlAction } from './hydrant'; import { readJetstream, subscribeJetstream, type Subscribe } from './jetstream'; import { getStaticUrls, renderShards, resetSitemap } from './sitemap'; import type { Env } from './types'; import { PROTOCOL_MAX_BYTES, PROTOCOL_MAX_URLS } from './xml'; const DEFAULT_HYDRANT_URL = 'https://index.tangled.network'; function hydrantUrl(env: Env): string { return (env.HYDRANT_URL || DEFAULT_HYDRANT_URL).replace(/\/+$/, ''); } export interface SyncDeps { subscribe: Subscribe; idleMs: number; closeMs: number; } const defaultDeps: SyncDeps = { subscribe: subscribeJetstream, idleMs: 3000, closeMs: 10_000, }; 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 { 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/gzip', 'Cache-Control': 'public, max-age=3600, stale-while-revalidate=86400', }, }); } return new Response('Not Found', { status: 404 }); }, }; export interface Steps { do(name: string, fn: () => Promise): Promise; } // the workflow checkpoints every step, tests and local runs just call them const inline: Steps = { do: (_name, fn) => fn() }; export interface SyncOptions { reset?: boolean; startedAt?: Date; } export interface SyncResult { syncedEvents: number; mutatedUrls: number; newCursor: string; totalUrls: number; } const MAX_PAGES = 10; export async function syncFromHydrant( env: Env, deps: SyncDeps = defaultDeps, steps: Steps = inline, options: SyncOptions = {}, ): Promise { const db = env.SITEMAP_DB; const baseUrl = env.BASE_URL || 'https://tangled.org'; const limits = parseConfigLimits(env); // a workflow replays this on resume, so the time has to come from the instance const now = options.startedAt ?? new Date(); const parsedBatch = parseInt(env.SYNC_BATCH_SIZE || '', 10); const batchSize = !isNaN(parsedBatch) && parsedBatch > 0 ? parsedBatch : 10000; const seen = new Set(); let syncedEvents = 0; let mutatedUrls = 0; if (options.reset) { await steps.do('reset', () => resetSitemap(db, env.SITEMAP_BUCKET)); } mutatedUrls += await steps.do('seed', async () => { if ((await countUrls(db)) > 0) return 0; const seeds = getStaticUrls(baseUrl, now).map((entry) => ({ type: 'upsert' as const, url: entry.loc, entry })); return applyChanges(db, seeds, limits, undefined, now); }); for (let page = 1; page <= MAX_PAGES; page++) { const batch = await steps.do(`batch ${page}`, async () => { // read here rather than carried over, so a retried step starts from what landed const cursor = Number(await getState(db, 'cursor')) || 0; const read = await readJetstream( { url: hydrantUrl(env), cursor, collections: Object.keys(COLLECTION_SPECS), maxEvents: batchSize, idleMs: deps.idleMs, closeMs: deps.closeMs, untilUs: now.getTime() * 1000, maxReconnects: 3, seen, }, deps.subscribe, ); if (read.events.length === 0) { if (read.error) throw new Error(read.error); return { events: 0, mutations: 0, more: false }; } const changes = read.events.flatMap((evt) => extractUrlAction(evt, baseUrl, now) ?? []); const mutations = await applyChanges(db, changes, limits, String(read.cursor), now); return { events: read.events.length, mutations, more: read.full && !read.error }; }); syncedEvents += batch.events; mutatedUrls += batch.mutations; if (!batch.more) break; } await steps.do('render', () => renderShards(db, env.SITEMAP_BUCKET, baseUrl, now)); return steps.do('finish', async () => { if (syncedEvents > 0) await markSynced(db, now); return { syncedEvents, mutatedUrls, newCursor: (await getState(db, 'cursor')) ?? '0', totalUrls: await countUrls(db), }; }); } const STEP_CONFIG: WorkflowStepConfig = { retries: { limit: 3, delay: '30 seconds', backoff: 'exponential' }, timeout: '15 minutes', }; // started by its own schedule at midnight, or by hand with // wrangler workflows trigger '{"reset":true}' export class SyncWorkflow extends WorkflowEntrypoint { async run(event: Readonly>, step: WorkflowStep): Promise { const steps: Steps = { do: (name, fn) => step.do(name, STEP_CONFIG, fn as () => Promise) }; const res = await syncFromHydrant(this.env, defaultDeps, steps, { reset: event.payload?.reset, startedAt: event.timestamp, }); console.log(`[sitemap] sync finished: ${JSON.stringify(res)}`); return res; } }