diff --git a/sitemap/src/hydrant.ts b/sitemap/src/hydrant.ts index 9e1e532ca..559b32fa5 100644 --- a/sitemap/src/hydrant.ts +++ b/sitemap/src/hydrant.ts @@ -1,11 +1,10 @@ import type { HydrantEvent, HydrantRecordValue, UrlEntry } from './types'; import { sanitizeLastmod } from './xml'; -export interface ExtractedAction { - type: 'upsert' | 'delete'; - url: string; - entry?: UrlEntry; -} +// uri is the record's at-uri, deletes carry no record so it's what finds their url +export type ExtractedAction = + | { type: 'delete'; uri: string; url: string } + | { type: 'upsert'; uri: string; url: string; entry: UrlEntry }; interface CollectionSpec { buildPath: (owner: string, rkey: string, record?: HydrantRecordValue) => string; @@ -48,17 +47,19 @@ export function extractUrlAction( return null; } + const uri = `at://${did}/${collection}/${rkey}`; const normalizedBase = baseUrl.replace(/\/+$/, ''); const path = spec.buildPath(did, rkey, record); const url = `${normalizedBase}/${path}`; if (action === 'delete') { - return { type: 'delete', url }; + return { type: 'delete', uri, url }; } const lastmod = sanitizeLastmod(record?.createdAt as string | undefined, now); return { type: 'upsert', + uri, url, entry: { loc: url, diff --git a/sitemap/src/index.ts b/sitemap/src/index.ts index 306f3e38a..931224b54 100644 --- a/sitemap/src/index.ts +++ b/sitemap/src/index.ts @@ -1,5 +1,5 @@ import type { Env, ShardInfo, UrlEntry, WorkerState } from './types'; -import { extractUrlAction } from './hydrant'; +import { extractUrlAction, type ExtractedAction } from './hydrant'; import { flushShardToR2, formatShardFilename, @@ -85,6 +85,107 @@ export default { }, }; +const KV_CONCURRENCY = 32; + +async function eachLimit(items: T[], limit: number, fn: (item: T) => Promise): Promise { + let next = 0; + const workers = Array.from({ length: Math.min(limit, items.length) }, async () => { + while (next < items.length) await fn(items[next++]); + }); + await Promise.all(workers); +} + +// one kv prefix with a read cache and the writes held back until flush +class CachedKv { + private cache = new Map(); + private dirty = new Map(); + + constructor( + private kv: KVNamespace, + private prefix: string, + ) {} + + async get(key: string): Promise { + if (!this.cache.has(key)) { + this.cache.set(key, await this.kv.get(this.prefix + key)); + } + return this.cache.get(key) ?? null; + } + + set(key: string, value: string | null): void { + this.cache.set(key, value); + this.dirty.set(key, value); + } + + // a backfill brings thousands of unknown keys, and reading them one per + // event takes longer than the cron gets to run + async prime(keys: Iterable): Promise { + const missing = [...new Set(keys)].filter((key) => !this.cache.has(key)); + await eachLimit(missing, KV_CONCURRENCY, (key) => this.get(key)); + } + + hasDirty(): boolean { + return this.dirty.size > 0; + } + + async flush(): Promise { + await eachLimit([...this.dirty], KV_CONCURRENCY, ([key, value]) => + value === null ? this.kv.delete(this.prefix + key) : this.kv.put(this.prefix + key, value), + ); + this.dirty.clear(); + } +} + +// u: is the shard a url sits in, r: is the url a record last made +export class SitemapKvStore { + private shards: CachedKv; + private records: CachedKv; + + constructor(kv: KVNamespace) { + this.shards = new CachedKv(kv, 'u:'); + this.records = new CachedKv(kv, 'r:'); + } + + async getUrlShard(url: string): Promise { + const val = await this.shards.get(url); + return val === null ? undefined : parseInt(val, 10); + } + + setUrlShard(url: string, shardIndex: number): void { + this.shards.set(url, String(shardIndex)); + } + + deleteUrl(url: string): void { + this.shards.set(url, null); + } + + async getRecordUrl(uri: string): Promise { + return (await this.records.get(uri)) ?? undefined; + } + + setRecordUrl(uri: string, url: string | null): void { + this.records.set(uri, url); + } + + async prime(actions: ExtractedAction[]): Promise { + await this.records.prime(actions.map((action) => action.uri)); + const previous = await Promise.all(actions.map((action) => this.getRecordUrl(action.uri))); + await this.shards.prime([ + ...actions.map((action) => action.url), + ...previous.filter((url): url is string => url !== undefined), + ]); + } + + hasDirty(): boolean { + return this.shards.hasDirty() || this.records.hasDirty(); + } + + async flush(): Promise { + await this.shards.flush(); + await this.records.flush(); + } +} + export async function loadWorkerState(kv: KVNamespace): Promise { const raw = await kv.get(KEY_WORKER_STATE); if (!raw) { @@ -93,7 +194,6 @@ export async function loadWorkerState(kv: KVNamespace): Promise { lastSyncAt: 'never', totalUrls: 0, shards: [], - urlLocator: {}, }; } @@ -104,7 +204,6 @@ export async function loadWorkerState(kv: KVNamespace): Promise { 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 { @@ -112,7 +211,6 @@ export async function loadWorkerState(kv: KVNamespace): Promise { lastSyncAt: 'never', totalUrls: 0, shards: [], - urlLocator: {}, }; } } @@ -130,6 +228,7 @@ export async function syncFromHydrant(env: Env): Promise<{ const now = new Date(); const state = await loadWorkerState(env.SITEMAP_KV); + const kvStore = new SitemapKvStore(env.SITEMAP_KV); // Initialize shard 1 if needed if (state.shards.length === 0) { @@ -166,7 +265,7 @@ export async function syncFromHydrant(env: Env): Promise<{ const staticUrls = getStaticUrls(baseUrl, now); for (const entry of staticUrls) { activeEntries.set(entry.loc, entry); - state.urlLocator[entry.loc] = activeShard.index; + kvStore.setUrlShard(entry.loc, activeShard.index); } mutated = true; } @@ -186,7 +285,7 @@ export async function syncFromHydrant(env: Env): Promise<{ const hydrantEvents = fetchResult.events; let highestCursor = state.cursor; - let eventCount = 0; + const eventCount = hydrantEvents.length; let mutationCount = 0; // Track historical shards modified by in-place updates or deletes: shardIndex -> { upserts, deletes } @@ -204,49 +303,54 @@ export async function syncFromHydrant(env: Env): Promise<{ return mod; } + async function removeUrl(url: string): Promise { + const shardIndex = await kvStore.getUrlShard(url); + if (!shardIndex) return false; + if (shardIndex === activeShard.index) { + activeEntries.delete(url); + } else { + const mod = getHistoricalMod(shardIndex); + mod.deletes.add(url); + mod.upserts.delete(url); + } + kvStore.deleteUrl(url); + mutated = true; + return true; + } + for (const evt of hydrantEvents) { highestCursor = String(Math.max(Number(highestCursor), evt.id)); - eventCount++; + } - const action = extractUrlAction(evt, baseUrl, now); - if (!action) continue; + const actions = hydrantEvents.flatMap((evt) => extractUrlAction(evt, baseUrl, now) ?? []); + await kvStore.prime(actions); + + for (const action of actions) { + const previous = await kvStore.getRecordUrl(action.uri); 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++; - } + kvStore.setRecordUrl(action.uri, null); + if (await removeUrl(previous ?? action.url)) mutationCount++; + continue; + } + + if (previous !== action.url) { + // a repo rename keeps its at-uri but changes its url + if (previous) await removeUrl(previous); + kvStore.setRecordUrl(action.uri, action.url); + } + + const existingShardIndex = await kvStore.getUrlShard(action.url); + if (!existingShardIndex || existingShardIndex === activeShard.index) { + activeEntries.set(action.url, action.entry); + if (!existingShardIndex) kvStore.setUrlShard(action.url, activeShard.index); + } else { + 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 @@ -294,10 +398,6 @@ export async function syncFromHydrant(env: Env): Promise<{ chunk, ); - for (const entry of chunk) { - state.urlLocator[entry.loc] = activeShard.index; - } - // Open next shard const nextIndex = activeShard.index + 1; activeShard = { @@ -309,6 +409,9 @@ export async function syncFromHydrant(env: Env): Promise<{ lastmod: sanitizeLastmod(undefined, now), }; state.shards.push(activeShard); + for (const entry of remainder) { + kvStore.setUrlShard(entry.loc, nextIndex); + } entriesArray = remainder; mutated = true; } @@ -323,10 +426,6 @@ export async function syncFromHydrant(env: Env): Promise<{ entriesArray, ); - for (const entry of entriesArray) { - state.urlLocator[entry.loc] = activeShard.index; - } - await updateSitemapIndex(env.SITEMAP_BUCKET, state.shards, baseUrl); } @@ -334,8 +433,10 @@ export async function syncFromHydrant(env: Env): Promise<{ state.totalUrls = state.shards.reduce((acc, s) => acc + s.count, 0); state.cursor = highestCursor; - if (mutated || eventCount > 0) { + // Commit ordering: granular KV entries saved first, then worker state is canonical commit marker + if (mutated || eventCount > 0 || kvStore.hasDirty()) { state.lastSyncAt = now.toISOString(); + await kvStore.flush(); await env.SITEMAP_KV.put(KEY_WORKER_STATE, JSON.stringify(state)); } diff --git a/sitemap/src/types.ts b/sitemap/src/types.ts index f34abbf6a..823ec234d 100644 --- a/sitemap/src/types.ts +++ b/sitemap/src/types.ts @@ -63,5 +63,4 @@ export interface WorkerState { lastSyncAt: string; totalUrls: number; shards: ShardInfo[]; - urlLocator: Record; } diff --git a/sitemap/test/sitemap.test.ts b/sitemap/test/sitemap.test.ts index 4203d33e8..8715d2fa0 100644 --- a/sitemap/test/sitemap.test.ts +++ b/sitemap/test/sitemap.test.ts @@ -101,8 +101,8 @@ describe('Hydrant Event Processing', () => { expect(action).not.toBeNull(); expect(action?.type).toBe('upsert'); expect(action?.url).toBe('https://tangled.org/did:plc:alice123/cool-tool'); - expect(action?.entry?.lastmod).toBe('2026-07-15'); - expect(action?.entry?.priority).toBe(0.8); + expect(action?.type === 'upsert' && action.entry.lastmod).toBe('2026-07-15'); + expect(action?.type === 'upsert' && action.entry.priority).toBe(0.8); }); it('extracts profile creation event', () => { @@ -143,10 +143,10 @@ describe('Hydrant Event Processing', () => { const action = extractUrlAction(event, baseUrl, now); expect(action?.type).toBe('upsert'); expect(action?.url).toBe('https://tangled.org/did:plc:alice123/strings/tid-3xyz'); - expect(action?.entry?.changefreq).toBe('weekly'); + expect(action?.type === 'upsert' && action.entry.changefreq).toBe('weekly'); }); - it('emits delete action when a repo is deleted', () => { + it('emits delete action with the record uri', () => { const event: HydrantEvent = { id: 104, type: 'record', @@ -160,7 +160,7 @@ describe('Hydrant Event Processing', () => { const action = extractUrlAction(event, baseUrl, now); expect(action?.type).toBe('delete'); - expect(action?.url).toBe('https://tangled.org/did:plc:alice123/deleted-repo'); + expect(action?.uri).toBe('at://did:plc:alice123/sh.tangled.repo/deleted-repo'); }); it('ignores pull requests, comments, identity events and unknown collections', () => { diff --git a/sitemap/test/sync.test.ts b/sitemap/test/sync.test.ts index 19e84ceba..229346ff2 100644 --- a/sitemap/test/sync.test.ts +++ b/sitemap/test/sync.test.ts @@ -15,6 +15,10 @@ class MockKV { this.putCount++; this.store.set(key, value); } + + async delete(key: string): Promise { + this.store.delete(key); + } } class MockR2Object { @@ -136,8 +140,8 @@ describe('Full syncFromHydrant workflow with Sharding & Multi-Sync', () => { expect(state1.shards[2].count).toBe(11); expect(state1.shards[2].sealed).toBe(false); - // Assert that repo-1 landed in shard 1 - expect(state1.urlLocator['https://tangled.org/did:plc:dawn/repo-1']).toBe(1); + // Assert that repo-1 landed in shard 1 under u: prefix + expect(await kv.get('u:https://tangled.org/did:plc:dawn/repo-1')).toBe('1'); // ========================================== // SYNC 2: @@ -183,6 +187,7 @@ describe('Full syncFromHydrant workflow with Sharding & Multi-Sync', () => { ]; currentMockEvents = sync2Events; + const kvPutsBeforeSync2 = kv.putCount; const res2 = await syncFromHydrant(env); expect(res2.success).toBe(true); expect(res2.syncedEvents).toBe(3); @@ -190,16 +195,19 @@ describe('Full syncFromHydrant workflow with Sharding & Multi-Sync', () => { expect(res2.newCursor).toBe('102'); expect(res2.totalUrls).toBe(61); // -1 deleted, +1 in-place updated, +1 new in active shard + // new-tool's shard and record url and the worker state, not every url again + expect(kv.putCount - kvPutsBeforeSync2).toBe(3); + const state2Raw = await kv.get('worker_state'); const state2 = JSON.parse(state2Raw!); - expect(state2.urlLocator['https://tangled.org/did:plc:dawn/new-tool']).toBe(3); + expect(await kv.get('u:https://tangled.org/did:plc:dawn/new-tool')).toBe('3'); // repo-2 remains in shard 1 (in-place update) - expect(state2.urlLocator['https://tangled.org/did:plc:dawn/repo-2']).toBe(1); + expect(await kv.get('u:https://tangled.org/did:plc:dawn/repo-2')).toBe('1'); // Deleted URL is completely gone from locator! - expect(state2.urlLocator['https://tangled.org/did:plc:dawn/repo-1']).toBeUndefined(); + expect(await kv.get('u:https://tangled.org/did:plc:dawn/repo-1')).toBeNull(); // Shard 1 in R2 was updated to 24 items and contains updated repo-2 but NO repo-1 expect(state2.shards[0].count).toBe(24); @@ -235,6 +243,43 @@ describe('Full syncFromHydrant workflow with Sharding & Multi-Sync', () => { expect(res4.newCursor).toBe('102'); // Cursor not falsely advanced }); + it('finds deleted and renamed repos through their record uri', async () => { + const kv = new MockKV(); + const bucket = new MockR2Bucket(); + const repo = (id: number, rkey: string, action: 'create' | 'update' | 'delete', name?: string): HydrantEvent => ({ + id, + type: 'record', + record: { + did: 'did:plc:alice', + collection: 'sh.tangled.repo', + rkey, + action, + record: name ? { name } : undefined, + }, + }); + let events = [repo(1, '3mtid1', 'create', 'old-name'), repo(2, '3mtid2', 'create', 'doomed')]; + const env: Env = { + SITEMAP_KV: kv as any, + SITEMAP_BUCKET: bucket as any, + HYDRANT: { fetch: async () => new Response(JSON.stringify(events)) } as any, + BASE_URL: 'https://tangled.org', + }; + + await syncFromHydrant(env); + expect(await kv.get('r:at://did:plc:alice/sh.tangled.repo/3mtid2')).toBe('https://tangled.org/did:plc:alice/doomed'); + + events = [repo(20, '3mtid1', 'update', 'new-name'), repo(21, '3mtid2', 'delete')]; + await syncFromHydrant(env); + + const shard = await gzipDecompress(await (await bucket.get('sitemaps/sitemap-0001.xml.gz'))!.arrayBuffer()); + expect(shard).toContain('https://tangled.org/did:plc:alice/new-name'); + expect(shard).not.toContain('old-name'); + expect(shard).not.toContain('doomed'); + expect(await kv.get('u:https://tangled.org/did:plc:alice/old-name')).toBeNull(); + expect(await kv.get('u:https://tangled.org/did:plc:alice/doomed')).toBeNull(); + expect(await kv.get('r:at://did:plc:alice/sh.tangled.repo/3mtid2')).toBeNull(); + }); + it('enforces byte-size cap rollover', async () => { const kv = new MockKV(); const bucket = new MockR2Bucket(); diff --git a/sitemap/wrangler.dev.jsonc b/sitemap/wrangler.dev.jsonc index ffbf5b339..7b55544db 100644 --- a/sitemap/wrangler.dev.jsonc +++ b/sitemap/wrangler.dev.jsonc @@ -8,6 +8,10 @@ "observability": { "enabled": true, }, + // a backfill reads and writes kv for every url and record, way past the default 10k + "limits": { + "subrequests": 500000, + }, "triggers": { "crons": ["0 0 * * *"], }, diff --git a/sitemap/wrangler.jsonc b/sitemap/wrangler.jsonc index b3844e641..cac7ed846 100644 --- a/sitemap/wrangler.jsonc +++ b/sitemap/wrangler.jsonc @@ -6,6 +6,10 @@ "observability": { "enabled": true, }, + // a backfill reads and writes kv for every url and record, way past the default 10k + "limits": { + "subrequests": 500000, + }, "triggers": { "crons": ["0 0 * * *"], },