import { describe, expect, it } from 'vitest'; import { parseConfigLimits, syncFromHydrant } from '../src/index'; import type { Env, HydrantEvent } from '../src/types'; import { gzipDecompress } from '../src/xml'; class MockKV { store = new Map(); putCount = 0; async get(key: string): Promise { return this.store.get(key) ?? null; } async put(key: string, value: string): Promise { this.putCount++; this.store.set(key, value); } } class MockR2Object { constructor(private buffer: ArrayBuffer, public httpMetadata?: any) {} async arrayBuffer(): Promise { return this.buffer; } } class MockR2Bucket { store = new Map(); putCount = 0; async get(key: string): Promise { const item = this.store.get(key); if (!item) return null; return new MockR2Object(item.buffer, item.httpMetadata); } async put( key: string, value: ArrayBuffer | Uint8Array | string, options?: { httpMetadata?: any }, ): Promise { this.putCount++; let buffer: ArrayBuffer; if (typeof value === 'string') { buffer = new TextEncoder().encode(value).buffer as ArrayBuffer; } else if (value instanceof Uint8Array) { buffer = value.buffer.slice(value.byteOffset, value.byteOffset + value.byteLength) as ArrayBuffer; } else { buffer = value; } this.store.set(key, { buffer, httpMetadata: options?.httpMetadata }); } } describe('Full syncFromHydrant workflow with Sharding & Multi-Sync', () => { it('handles multi-sync, preserves identity, updates historical shards in-place, and enforces caps', async () => { const kv = new MockKV(); const bucket = new MockR2Bucket(); // ========================================== // SYNC 1: Initial creation + multi-shard rollover // ========================================== const sync1Events: HydrantEvent[] = [ { id: 1, type: 'identity', identity: { did: 'did:plc:dawn', handle: 'dawn' }, }, { id: 2, type: 'record', record: { did: 'did:plc:dawn', collection: 'sh.tangled.actor.profile', rkey: 'self', action: 'create', record: { createdAt: '2026-08-01T00:00:00Z' }, }, }, ]; // 58 repo events -> 61 total URLs (2 static + 1 profile + 58 repos) for (let i = 1; i <= 58; i++) { sync1Events.push({ id: 10 + i, type: 'record', record: { did: 'did:plc:dawn', collection: 'sh.tangled.repo', rkey: `repo-${i}`, action: 'create', record: { name: `repo-${i}`, createdAt: '2026-08-15T00:00:00Z', }, }, }); } let currentMockEvents = sync1Events; const mockHydrantFetcher = { fetch: async (url: string | Request) => { const reqUrl = typeof url === 'string' ? url : url.url; const u = new URL(reqUrl); const cursor = Number(u.searchParams.get('cursor') || '0'); const filtered = currentMockEvents.filter((e) => e.id > cursor); return new Response(JSON.stringify(filtered), { status: 200, headers: { 'Content-Type': 'application/json' }, }); }, }; const env: Env = { SITEMAP_KV: kv as any, SITEMAP_BUCKET: bucket as any, HYDRANT: mockHydrantFetcher as any, BASE_URL: 'https://tangled.org', SHARD_MAX_URLS: '25', // 25 max per shard }; const res1 = await syncFromHydrant(env); expect(res1.success).toBe(true); expect(res1.syncedEvents).toBe(60); expect(res1.totalUrls).toBe(61); expect(res1.newCursor).toBe('68'); // Verify identity was persisted in KV const identityRaw = await kv.get('identity_cache'); expect(identityRaw).toContain('"did:plc:dawn":"dawn"'); // Shards: 25 + 25 + 11 = 61 const state1Raw = await kv.get('worker_state'); const state1 = JSON.parse(state1Raw!); expect(state1.shards).toHaveLength(3); expect(state1.shards[0].count).toBe(25); expect(state1.shards[0].sealed).toBe(true); expect(state1.shards[1].count).toBe(25); expect(state1.shards[1].sealed).toBe(true); 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/dawn/repo-1']).toBe(1); // ========================================== // SYNC 2: Second sync without identity event // - Asserts no did: URLs are generated (proves handle persistence) // - Asserts in-place update in historical shard 1 (repo-2) // - Asserts delete against historical sealed shard 1 (repo-1) // ========================================== const sync2Events: HydrantEvent[] = [ // New repo from same user { id: 100, type: 'record', record: { did: 'did:plc:dawn', collection: 'sh.tangled.repo', rkey: 'new-tool', action: 'create', record: { createdAt: '2026-09-01T00:00:00Z' }, }, }, // Delete repo-1 which is inside sealed shard 1 { id: 101, type: 'record', record: { did: 'did:plc:dawn', collection: 'sh.tangled.repo', rkey: 'repo-1', action: 'delete', }, }, // In-place update for repo-2 inside shard 1 { id: 102, type: 'record', record: { did: 'did:plc:dawn', collection: 'sh.tangled.repo', rkey: 'repo-2', action: 'update', record: { createdAt: '2026-09-05T00:00:00Z' }, }, }, ]; currentMockEvents = sync2Events; const res2 = await syncFromHydrant(env); expect(res2.success).toBe(true); expect(res2.syncedEvents).toBe(3); expect(res2.mutatedUrls).toBe(3); expect(res2.newCursor).toBe('102'); expect(res2.totalUrls).toBe(61); // -1 deleted, +1 in-place updated, +1 new in active shard const state2Raw = await kv.get('worker_state'); const state2 = JSON.parse(state2Raw!); // New URL has resolved handle dawn, NOT did:plc:dawn! expect(state2.urlLocator['https://tangled.org/dawn/new-tool']).toBe(3); expect(state2.urlLocator['https://tangled.org/did:plc:dawn/new-tool']).toBeUndefined(); // repo-2 remains in shard 1 (in-place update) expect(state2.urlLocator['https://tangled.org/dawn/repo-2']).toBe(1); // Deleted URL is completely gone from locator! expect(state2.urlLocator['https://tangled.org/dawn/repo-1']).toBeUndefined(); // 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); const shard1Obj = await bucket.get('sitemaps/sitemap-0001.xml.gz'); const shard1Xml = await gzipDecompress(await shard1Obj!.arrayBuffer()); expect(shard1Xml).not.toContain('https://tangled.org/dawn/repo-1'); expect(shard1Xml).toContain('https://tangled.org/dawn/repo-2'); expect(shard1Xml).toContain('2026-09-05'); // ========================================== // SYNC 3: No-op sync pins zero R2 / KV puts! // ========================================== const r2PutsBefore = bucket.putCount; const kvPutsBefore = kv.putCount; currentMockEvents = []; const res3 = await syncFromHydrant(env); expect(res3.success).toBe(true); expect(res3.syncedEvents).toBe(0); expect(res3.mutatedUrls).toBe(0); // Explicitly assert zero R2 and KV puts occurred expect(bucket.putCount).toBe(r2PutsBefore); expect(kv.putCount).toBe(kvPutsBefore); // ========================================== // SYNC 4: Fetch failure propagation // ========================================== mockHydrantFetcher.fetch = async () => new Response('Internal error', { status: 500 }); const res4 = await syncFromHydrant(env); expect(res4.success).toBe(false); expect(res4.error).toContain('HTTP 500'); expect(res4.newCursor).toBe('102'); // Cursor not falsely advanced }); it('enforces byte-size cap rollover', async () => { const kv = new MockKV(); const bucket = new MockR2Bucket(); // Tiny 350-byte max cap -> will force rollover after ~2 URLs const env: Env = { SITEMAP_KV: kv as any, SITEMAP_BUCKET: bucket as any, HYDRANT: { fetch: async () => new Response(JSON.stringify([ { id: 1, type: 'record', record: { did: 'did:plc:alice', handle: 'alice', collection: 'sh.tangled.repo', rkey: 'long-repository-name-to-consume-bytes', action: 'create', }, }, ])), } as any, BASE_URL: 'https://tangled.org', SHARD_MAX_BYTES: '350', // very low byte limit }; const res = await syncFromHydrant(env); expect(res.success).toBe(true); const stateRaw = await kv.get('worker_state'); const state = JSON.parse(stateRaw!); // Shard 1 sealed due to byte cap, shard 2 opened expect(state.shards.length).toBeGreaterThanOrEqual(2); expect(state.shards[0].sealed).toBe(true); }); it('validates config limits and handles NaN safely', () => { expect(parseConfigLimits({ SITEMAP_KV: {} as any, SITEMAP_BUCKET: {} as any, SHARD_MAX_URLS: 'abc' })).toEqual({ maxUrls: 50000, maxBytes: 52428800, }); expect(parseConfigLimits({ SITEMAP_KV: {} as any, SITEMAP_BUCKET: {} as any, SHARD_MAX_URLS: '-5' })).toEqual({ maxUrls: 50000, maxBytes: 52428800, }); expect(parseConfigLimits({ SITEMAP_KV: {} as any, SITEMAP_BUCKET: {} as any, SHARD_MAX_URLS: '1000' })).toEqual({ maxUrls: 1000, maxBytes: 52428800, }); }); });