diff --git a/package.json b/package.json index ec29dc5..f4838c4 100644 --- a/package.json +++ b/package.json @@ -11,7 +11,7 @@ "knowledge:promote": "tsx scripts/promote-knowledge.ts", "knowledge:sync": "tsx scripts/sync-knowledge-atproto.ts", "test:markdown": "tsx --test src/markdown.test.ts", - "test:content": "tsx --test src/about-content.test.ts src/blog-data.test.ts src/profile-data.test.ts src/knowledge.test.ts src/knowledge-landing.test.tsx", + "test:content": "tsx --test src/about-content.test.ts src/blog-data.test.ts src/pds-data.test.ts src/profile-data.test.ts src/knowledge.test.ts src/knowledge-landing.test.tsx", "test:charts": "tsx --test src/charts/charts.test.ts src/knowledge-charts/knowledge-charts.test.tsx", "test:worker": "bash scripts/test-content-sync-worker.sh", "test": "pnpm test:markdown && pnpm test:content && pnpm test:charts && pnpm test:worker && pnpm knowledge:check", diff --git a/src/data.ts b/src/data.ts index e6a262c..6648506 100644 --- a/src/data.ts +++ b/src/data.ts @@ -38,8 +38,12 @@ export interface ProfileData { } const DEFAULT_PROFILE_FETCH_TIMEOUT_MS = 3_000; +const DEFAULT_PDS_FETCH_TIMEOUT_MS = 5_000; const profileRequests = new Map>(); const lastKnownProfiles = new Map(); +type RecordList = Array<{ uri: string; value: Record }>; +const recordRequests = new Map>(); +const lastKnownRecords = new Map(); function getDid(): string { return process.env.CAMERON_DID || "did:plc:gfrmhdmjvxn2sjedzboeudef"; @@ -47,9 +51,22 @@ function getDid(): string { let pdsEndpointPromise: Promise | undefined; -async function getPdsEndpoint(): Promise { +function pdsFetchTimeoutMs(): number { + const configured = Number(process.env.PDS_FETCH_TIMEOUT_MS); + return Number.isFinite(configured) && configured > 0 + ? configured + : DEFAULT_PDS_FETCH_TIMEOUT_MS; +} + +function remainingTimeoutMs(deadline: number): number { + return Math.max(1, deadline - Date.now()); +} + +async function getPdsEndpoint(timeoutMs = pdsFetchTimeoutMs()): Promise { if (process.env.ATP_SERVICE) return process.env.ATP_SERVICE.replace(/\/$/, ""); - pdsEndpointPromise ??= fetch(`https://plc.directory/${getDid()}`) + pdsEndpointPromise ??= fetch(`https://plc.directory/${getDid()}`, { + signal: AbortSignal.timeout(timeoutMs), + }) .then(async (response) => { if (!response.ok) throw new Error(`DID document lookup failed: HTTP ${response.status}`); const document = await response.json() as { @@ -60,41 +77,80 @@ async function getPdsEndpoint(): Promise { ); if (!service?.serviceEndpoint) throw new Error("DID document has no ATProto PDS endpoint"); return service.serviceEndpoint.replace(/\/$/, ""); + }) + .catch((error) => { + pdsEndpointPromise = undefined; + throw error; }); return pdsEndpointPromise; } +async function fetchRecords( + repo: string, + collection: string, + limit: number, + cacheKey: string, + required: boolean, +): Promise { + const fallback = lastKnownRecords.get(cacheKey); + const deadline = Date.now() + pdsFetchTimeoutMs(); + try { + const pds = await getPdsEndpoint(remainingTimeoutMs(deadline)); + const res = await fetch( + `${pds}/xrpc/com.atproto.repo.listRecords?repo=${encodeURIComponent(repo)}&collection=${encodeURIComponent(collection)}&limit=${limit}`, + { signal: AbortSignal.timeout(remainingTimeoutMs(deadline)) }, + ); + if (!res.ok) { + throw new Error(`PDS listRecords failed for ${collection}: HTTP ${res.status}`); + } + const data = await res.json(); + const records = (data.records ?? []) as RecordList; + lastKnownRecords.set(cacheKey, records); + await cacheSet(cacheKey, records, { life: "minutes" }); + return records; + } catch (error) { + if (fallback !== undefined) { + console.warn( + `[pds] ${collection} unavailable; rendering last-known records: ${error instanceof Error ? error.message : String(error)}`, + ); + await cacheSet(cacheKey, fallback, { life: "seconds" }); + return fallback; + } + if (required) throw error; + console.warn( + `[pds] ${collection} unavailable; rendering an empty optional surface: ${error instanceof Error ? error.message : String(error)}`, + ); + await cacheSet(cacheKey, [], { life: "seconds" }); + return []; + } +} + async function listRecords( repo: string, collection: string, limit = 100, options: { required?: boolean } = {} -): Promise }>> { - const pds = await getPdsEndpoint(); +): Promise { const cacheKey = `records:${repo}:${collection}:${limit}`; const cached = await cacheGet(cacheKey); if (cached !== undefined) { - return cached as Array<{ uri: string; value: Record }>; + return cached as RecordList; } - const res = await fetch( - `${pds}/xrpc/com.atproto.repo.listRecords?repo=${encodeURIComponent(repo)}&collection=${encodeURIComponent(collection)}&limit=${limit}` - ); - if (!res.ok) { - if (options.required) { - throw new Error( - `PDS listRecords failed for ${collection}: HTTP ${res.status}` - ); - } - return []; - } - const data = await res.json(); - const records = (data.records ?? []) as Array<{ - uri: string; - value: Record; - }>; - await cacheSet(cacheKey, records, { life: "minutes" }); - return records; + const existing = recordRequests.get(cacheKey); + if (existing) return existing; + + const request = fetchRecords( + repo, + collection, + limit, + cacheKey, + options.required === true, + ).finally(() => { + if (recordRequests.get(cacheKey) === request) recordRequests.delete(cacheKey); + }); + recordRequests.set(cacheKey, request); + return request; } function slugFromPath(path: string | undefined): string { diff --git a/src/pds-data.test.ts b/src/pds-data.test.ts new file mode 100644 index 0000000..97cb83e --- /dev/null +++ b/src/pds-data.test.ts @@ -0,0 +1,204 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { getEndorsements, listBlogPosts } from "./data.ts"; + +const originalFetch = globalThis.fetch; +const originalDid = process.env.CAMERON_DID; +const originalService = process.env.ATP_SERVICE; +const originalTimeout = process.env.PDS_FETCH_TIMEOUT_MS; + +function restoreEnvironment(): void { + globalThis.fetch = originalFetch; + if (originalDid === undefined) delete process.env.CAMERON_DID; + else process.env.CAMERON_DID = originalDid; + if (originalService === undefined) delete process.env.ATP_SERVICE; + else process.env.ATP_SERVICE = originalService; + if (originalTimeout === undefined) delete process.env.PDS_FETCH_TIMEOUT_MS; + else process.env.PDS_FETCH_TIMEOUT_MS = originalTimeout; +} + +test("shares one in-flight PDS record request across concurrent Blog renders", async () => { + process.env.CAMERON_DID = "did:example:pds-singleflight"; + process.env.ATP_SERVICE = "https://pds.example"; + let fetchCalls = 0; + let releaseFetch: (() => void) | undefined; + globalThis.fetch = (() => { + fetchCalls += 1; + return new Promise((resolve) => { + releaseFetch = () => resolve(Response.json({ records: [{ + uri: "at://did:example:pds-singleflight/site.standard.document/post", + value: { + site: "at://did:plc:gfrmhdmjvxn2sjedzboeudef/site.standard.publication/3md7ylshxzk2y", + title: "Post", + path: "/post", + publishedAt: "2026-07-27T00:00:00.000Z", + tags: ["blog"], + textContent: "Body", + }, + }] })); + }); + }) as typeof fetch; + + try { + const first = listBlogPosts(); + const second = listBlogPosts(); + await new Promise((resolve) => setImmediate(resolve)); + assert.equal(fetchCalls, 1); + releaseFetch?.(); + assert.deepEqual(await first, await second); + } finally { + restoreEnvironment(); + } +}); + +test("bounds required PDS record reads when the service stalls", async () => { + process.env.CAMERON_DID = "did:example:pds-timeout"; + process.env.ATP_SERVICE = "https://pds.example"; + process.env.PDS_FETCH_TIMEOUT_MS = "20"; + globalThis.fetch = ((_input, init) => new Promise((_resolve, reject) => { + const fallback = setTimeout( + () => reject(new Error("PDS fetch was not aborted")), + 500, + ); + init?.signal?.addEventListener( + "abort", + () => { + clearTimeout(fallback); + reject(init.signal?.reason ?? new Error("aborted")); + }, + { once: true }, + ); + })) as typeof fetch; + + try { + const startedAt = Date.now(); + await assert.rejects(listBlogPosts(), /timeout|aborted/i); + assert.ok(Date.now() - startedAt < 250); + } finally { + restoreEnvironment(); + } +}); + +test("serves last-known Blog records after the cache expires during a PDS outage", async () => { + process.env.CAMERON_DID = "did:example:pds-stale"; + process.env.ATP_SERVICE = "https://pds.example"; + process.env.PDS_FETCH_TIMEOUT_MS = "20"; + const originalDateNow = Date.now; + let now = 1_785_192_000_000; + Date.now = () => now; + globalThis.fetch = (async () => Response.json({ records: [{ + uri: "at://did:example:pds-stale/site.standard.document/post", + value: { + site: "at://did:plc:gfrmhdmjvxn2sjedzboeudef/site.standard.publication/3md7ylshxzk2y", + title: "Stale Post", + path: "/stale-post", + publishedAt: "2026-07-27T00:00:00.000Z", + tags: ["blog"], + textContent: "Still useful during an outage.", + }, + }] })) as typeof fetch; + + try { + const fresh = await listBlogPosts(); + assert.equal(fresh[0]?.title, "Stale Post"); + now += 5 * 60 * 1_000 + 1; + globalThis.fetch = ((_input, init) => new Promise((_resolve, reject) => { + const fallback = setTimeout( + () => reject(new Error("PDS fetch was not aborted")), + 500, + ); + init?.signal?.addEventListener( + "abort", + () => { + clearTimeout(fallback); + reject(init.signal?.reason ?? new Error("aborted")); + }, + { once: true }, + ); + })) as typeof fetch; + + const stale = await listBlogPosts(); + assert.deepEqual(stale, fresh); + } finally { + Date.now = originalDateNow; + restoreEnvironment(); + } +}); + +test("does not promote an optional outage sentinel into last-known records", async () => { + process.env.CAMERON_DID = "did:example:pds-optional-empty"; + process.env.ATP_SERVICE = "https://pds.example"; + process.env.PDS_FETCH_TIMEOUT_MS = "20"; + const originalDateNow = Date.now; + const originalWarn = console.warn; + let now = 1_785_192_000_000; + const warnings: string[] = []; + Date.now = () => now; + console.warn = (message) => warnings.push(String(message)); + globalThis.fetch = ((_input, init) => new Promise((_resolve, reject) => { + const fallback = setTimeout( + () => reject(new Error("PDS fetch was not aborted")), + 500, + ); + init?.signal?.addEventListener( + "abort", + () => { + clearTimeout(fallback); + reject(init.signal?.reason ?? new Error("aborted")); + }, + { once: true }, + ); + })) as typeof fetch; + + try { + assert.deepEqual(await getEndorsements(), []); + assert.deepEqual(await getEndorsements(), []); + now += 31_000; + assert.deepEqual(await getEndorsements(), []); + assert.equal(warnings.length, 2); + assert.ok(warnings.every((warning) => warning.includes("empty optional surface"))); + } finally { + Date.now = originalDateNow; + console.warn = originalWarn; + restoreEnvironment(); + } +}); + +test("uses one wall-clock deadline across DID resolution and PDS reads", async () => { + process.env.CAMERON_DID = "did:example:pds-deadline"; + delete process.env.ATP_SERVICE; + process.env.PDS_FETCH_TIMEOUT_MS = "50"; + const startedAt = Date.now(); + globalThis.fetch = ((input, init) => { + if (String(input).startsWith("https://plc.directory/")) { + return new Promise((resolve) => { + setTimeout(() => resolve(Response.json({ service: [{ + id: "#atproto_pds", + type: "AtprotoPersonalDataServer", + serviceEndpoint: "https://pds.example", + }] })), 30); + }); + } + return new Promise((_resolve, reject) => { + const fallback = setTimeout( + () => reject(new Error("PDS fetch was not aborted")), + 500, + ); + init?.signal?.addEventListener( + "abort", + () => { + clearTimeout(fallback); + reject(init.signal?.reason ?? new Error("aborted")); + }, + { once: true }, + ); + }); + }) as typeof fetch; + + try { + await assert.rejects(listBlogPosts(), /timeout|aborted/i); + assert.ok(Date.now() - startedAt < 75); + } finally { + restoreEnvironment(); + } +});