// Data fetching layer for protocol-native dynamic surfaces. Blog records are // read directly from Cameron's Leaflet publication on the PDS. About is the // intentionally separate Git-backed document. import { AtUri } from "@atproto/syntax"; import { cacheGet, cacheSet } from "./cache.ts"; import { leafletContentToMarkdown } from "./leaflet-reader.ts"; export { getAbout, } from "./about-content.ts"; const BLOG_COLLECTION = "site.standard.document"; const BLOG_PUBLICATION_URI = process.env.PUBLICATION_URI || "at://did:plc:gfrmhdmjvxn2sjedzboeudef/site.standard.publication/3md7ylshxzk2y"; export interface BlogPost { uri: string; rkey: string; slug: string; title: string; body: string; description?: string; tags?: string[]; publishedAt: string; updatedAt?: string; } export interface ProfileData { did: string; handle: string; displayName?: string; description?: string; avatar?: string; followersCount?: number; followsCount?: number; postsCount?: number; } 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 }>; type RecordPage = { records?: RecordList; cursor?: string }; const recordRequests = new Map>(); const lastKnownRecords = new Map(); function getDid(): string { return process.env.CAMERON_DID || "did:plc:gfrmhdmjvxn2sjedzboeudef"; } let pdsEndpointPromise: Promise | undefined; 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()}`, { 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 { service?: Array<{ id?: string; type?: string; serviceEndpoint?: string }>; }; const service = document.service?.find((item) => item.id === "#atproto_pds" || item.type === "AtprotoPersonalDataServer" ); 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, allPages: boolean, ): Promise { const fallback = lastKnownRecords.get(cacheKey); const deadline = Date.now() + pdsFetchTimeoutMs(); try { const pds = await getPdsEndpoint(remainingTimeoutMs(deadline)); const records: RecordList = []; const seenCursors = new Set(); let cursor: string | undefined; do { const url = new URL(`${pds}/xrpc/com.atproto.repo.listRecords`); url.searchParams.set("repo", repo); url.searchParams.set("collection", collection); url.searchParams.set("limit", String(limit)); if (cursor) url.searchParams.set("cursor", cursor); const res = await fetch(url, { signal: AbortSignal.timeout(remainingTimeoutMs(deadline)), }); if (!res.ok) { throw new Error(`PDS listRecords failed for ${collection}: HTTP ${res.status}`); } const data = await res.json() as RecordPage; const page = data.records ?? []; records.push(...page); if (!allPages || !data.cursor) break; if (seenCursors.has(data.cursor)) { throw new Error(`PDS listRecords repeated cursor for ${collection}`); } seenCursors.add(data.cursor); cursor = data.cursor; } while (true); 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; allPages?: boolean } = {} ): Promise { const allPages = options.allPages === true; const cacheKey = `records:${repo}:${collection}:${limit}:${allPages ? "all" : "page"}`; const cached = await cacheGet(cacheKey); if (cached !== undefined) { return cached as RecordList; } const existing = recordRequests.get(cacheKey); if (existing) return existing; const request = fetchRecords( repo, collection, limit, cacheKey, options.required === true, allPages, ).finally(() => { if (recordRequests.get(cacheKey) === request) recordRequests.delete(cacheKey); }); recordRequests.set(cacheKey, request); return request; } function slugFromPath(path: string | undefined): string { if (!path) return ""; const parts = path.split("/").filter(Boolean); return parts.at(-1) ?? ""; } function stripDuplicateTitle(body: string, title: string): string { const match = body.match(/^#\s+(.+)\n*/); if (match && match[1].trim() === title.trim()) return body.slice(match[0].length); return body; } export function parseBlogDocument( uri: string, value: Record, ): BlogPost { const parsed = new AtUri(uri); const title = String(value.title ?? ""); const content = value.content as Record | undefined; const rawBody = content?.$type === "pub.leaflet.content" ? leafletContentToMarkdown(content, getDid()) : typeof content?.value === "string" ? content.value : String(value.textContent ?? ""); return { uri, rkey: parsed.rkey, slug: slugFromPath(value.path as string | undefined), title, body: stripDuplicateTitle(rawBody, title), description: value.description as string | undefined, tags: value.tags as string[] | undefined, publishedAt: String(value.publishedAt ?? ""), updatedAt: value.updatedAt as string | undefined, }; } function profileFetchTimeoutMs(): number { const configured = Number(process.env.PROFILE_FETCH_TIMEOUT_MS); return Number.isFinite(configured) && configured > 0 ? configured : DEFAULT_PROFILE_FETCH_TIMEOUT_MS; } async function fetchProfile(did: string, cacheKey: string): Promise { const fallback = lastKnownProfiles.get(did) ?? null; try { const res = await fetch( `https://public.api.bsky.app/xrpc/app.bsky.actor.getProfile?actor=${encodeURIComponent(did)}`, { signal: AbortSignal.timeout(profileFetchTimeoutMs()) }, ); if (!res.ok) { await cacheSet(cacheKey, fallback, { life: "seconds" }); return fallback; } const data = await res.json(); const profile = { did: data.did, handle: data.handle, displayName: data.displayName, description: data.description, avatar: data.avatar, followersCount: data.followersCount, followsCount: data.followsCount, postsCount: data.postsCount, } satisfies ProfileData; lastKnownProfiles.set(did, profile); await cacheSet(cacheKey, profile, { life: "minutes" }); return profile; } catch (error) { console.warn( `[profile] Bluesky profile unavailable; rendering without fresh profile data: ${error instanceof Error ? error.message : String(error)}`, ); await cacheSet(cacheKey, fallback, { life: "seconds" }); return fallback; } } export async function getProfile(): Promise { const did = getDid(); const cacheKey = `profile:${did}`; const cached = await cacheGet(cacheKey); if (cached !== undefined) { const profile = cached as ProfileData | null; if (profile) lastKnownProfiles.set(did, profile); return profile; } const existing = profileRequests.get(did); if (existing) return existing; const request = fetchProfile(did, cacheKey).finally(() => { if (profileRequests.get(did) === request) profileRequests.delete(did); }); profileRequests.set(did, request); return request; } export async function listBlogPosts(): Promise { const did = getDid(); const records = await listRecords(did, BLOG_COLLECTION, 100, { required: true, allPages: true, }); return records .filter((record) => record.value.site === BLOG_PUBLICATION_URI) .map((record) => parseBlogDocument(record.uri, record.value)) .filter((post) => post.slug && post.title && post.publishedAt) .sort( (left, right) => new Date(right.publishedAt).getTime() - new Date(left.publishedAt).getTime(), ); } export async function getBlogPost(slug: string): Promise { return (await listBlogPosts()).find((post) => post.slug === slug) ?? null; } // --- Margin annotations --- export interface MarginAnnotation { id: string; body?: { value?: string; format?: string }; target: { source: string; selector?: any; title?: string }; creator?: { name?: string; id?: string }; created?: string; motivation?: string; } export async function getAnnotations( pageUrl: string ): Promise { if (process.env.ENABLE_MARGIN !== "true") return []; try { const res = await fetch( `https://margin.at/api/annotations?url=${encodeURIComponent(pageUrl)}` ); if (!res.ok) return []; const data = await res.json(); return (data.items ?? []) as MarginAnnotation[]; } catch { return []; } } // --- Endorsements (fund.at.endorse) --- export interface Endorsement { uri: string; // endorsed entity (DID or hostname) createdAt: string; } export async function getEndorsements(): Promise { const did = getDid(); const records = await listRecords(did, "fund.at.endorse"); return records .map((r) => ({ uri: r.value.uri as string, createdAt: r.value.createdAt as string, })) .sort( (a, b) => new Date(b.createdAt).getTime() - new Date(a.createdAt).getTime() ); } // --- Bluesky feed --- export interface FeedPost { uri: string; cid: string; text: string; createdAt: string; author: { did: string; handle: string; displayName?: string; avatar?: string; }; likeCount: number; repostCount: number; replyCount: number; // Embedded images images?: Array<{ thumb: string; fullsize: string; alt: string; }>; // Embedded link card external?: { uri: string; title: string; description: string; thumb?: string; }; // Quote post quotedPost?: { uri: string; text: string; author: { handle: string; displayName?: string; avatar?: string; }; }; // Is this a reply? isReply: boolean; } // --- Semble (network.cosmik) --- const SEMBLE_HANDLE = process.env.SEMBLE_HANDLE || "cameron.stream"; export interface SembleCardContent { url: string; title?: string; description?: string; imageUrl?: string; siteName?: string; author?: string; type?: string; } export interface SembleCard { id: string; type: string; url: string; uri?: string; cardContent: SembleCardContent; note?: { id: string; text: string }; createdAt: string; updatedAt: string; collections: Array<{ id: string; name: string }>; author: { id: string; name?: string; handle: string; avatarUrl?: string; }; } export interface SembleCollection { id: string; uri?: string; name: string; description?: string; accessType?: string; cardCount: number; createdAt: string; updatedAt: string; author: { id: string; name?: string; handle: string; avatarUrl?: string; }; } export interface SembleData { cards: SembleCard[]; collections: SembleCollection[]; available: boolean; } interface SembleRecord { uri: string; value: Record; } function objectValue(value: unknown): Record | undefined { if (!value || typeof value !== "object" || Array.isArray(value)) { return undefined; } return value as Record; } function stringValue(value: unknown): string | undefined { return typeof value === "string" ? value : undefined; } function recordId(uri: string): string { return uri.split("/").pop() ?? uri; } function refUri(value: unknown): string | undefined { return stringValue(objectValue(value)?.uri); } export async function getSembleData(limit = 20): Promise { const did = getDid(); let records: [SembleRecord[], SembleRecord[], SembleRecord[]]; try { records = (await Promise.all([ listRecords(did, "network.cosmik.card", 100, { required: true }), listRecords(did, "network.cosmik.collection", 100, { required: true }), listRecords(did, "network.cosmik.collectionLink", 100, { required: true, }), ])) as [SembleRecord[], SembleRecord[], SembleRecord[]]; } catch (error) { console.error("Unable to load Semble records from the PDS", error); return { cards: [], collections: [], available: false }; } const [cardRecords, collectionRecords, collectionLinkRecords] = records; const linksByCard = new Map(); const linkCountByCollection = new Map(); for (const linkRecord of collectionLinkRecords) { const cardUri = refUri(linkRecord.value.card); const collectionUri = refUri(linkRecord.value.collection); if (!cardUri || !collectionUri) continue; const collectionUris = linksByCard.get(cardUri) ?? []; collectionUris.push(collectionUri); linksByCard.set(cardUri, collectionUris); linkCountByCollection.set( collectionUri, (linkCountByCollection.get(collectionUri) ?? 0) + 1 ); } const collections = collectionRecords .map((record): SembleCollection | null => { const name = stringValue(record.value.name); const createdAt = stringValue(record.value.createdAt); if (!name || !createdAt) return null; return { id: recordId(record.uri), uri: record.uri, name, description: stringValue(record.value.description), accessType: stringValue(record.value.accessType), cardCount: linkCountByCollection.get(record.uri) ?? 0, createdAt, updatedAt: stringValue(record.value.updatedAt) ?? createdAt, author: { id: did, handle: SEMBLE_HANDLE, }, }; }) .filter((collection): collection is SembleCollection => collection !== null) .sort( (a, b) => new Date(b.updatedAt).getTime() - new Date(a.updatedAt).getTime() ); const collectionsByUri = new Map( collections .filter((collection) => collection.uri) .map((collection) => [collection.uri!, collection]) ); const notesByParentUri = new Map(); for (const record of cardRecords) { if (record.value.type !== "NOTE") continue; const parentUri = refUri(record.value.parentCard); const text = stringValue(objectValue(record.value.content)?.text); if (!parentUri || !text) continue; notesByParentUri.set(parentUri, { id: recordId(record.uri), text, }); } const cards = cardRecords .map((record): SembleCard | null => { if (record.value.type !== "URL") return null; const content = objectValue(record.value.content); const metadata = objectValue(content?.metadata); const url = stringValue(content?.url) ?? stringValue(record.value.url); const createdAt = stringValue(record.value.createdAt); if (!url || !createdAt) return null; const cardCollections = (linksByCard.get(record.uri) ?? []) .map((collectionUri) => collectionsByUri.get(collectionUri)) .filter( (collection): collection is SembleCollection => collection !== undefined ) .map((collection) => ({ id: collection.id, name: collection.name, })); return { id: recordId(record.uri), type: "URL", url, uri: record.uri, cardContent: { url, title: stringValue(metadata?.title), description: stringValue(metadata?.description), imageUrl: stringValue(metadata?.imageUrl), siteName: stringValue(metadata?.siteName), author: stringValue(metadata?.author), type: stringValue(metadata?.type), }, note: notesByParentUri.get(record.uri), createdAt, updatedAt: stringValue(record.value.updatedAt) ?? createdAt, collections: cardCollections, author: { id: did, handle: SEMBLE_HANDLE, }, }; }) .filter((card): card is SembleCard => card !== null) .sort( (a, b) => new Date(b.createdAt).getTime() - new Date(a.createdAt).getTime() ) .slice(0, limit); return { cards, collections, available: true }; } export async function getSembleCollections(): Promise { return (await getSembleData()).collections; } export async function getSembleCards(limit = 20): Promise { return (await getSembleData(limit)).cards; } export async function getAuthorFeed(limit = 20): Promise { const did = getDid(); const res = await fetch( `https://public.api.bsky.app/xrpc/app.bsky.feed.getAuthorFeed?actor=${encodeURIComponent(did)}&limit=${limit}&filter=posts_no_replies` ); if (!res.ok) return []; const data = await res.json(); return (data.feed ?? []).map((item: any) => { const post = item.post; const record = post.record; const embed = post.embed; // Extract images from embed let images: FeedPost["images"]; if (embed?.$type === "app.bsky.embed.images#view") { images = embed.images.map((img: any) => ({ thumb: img.thumb, fullsize: img.fullsize, alt: img.alt ?? "", })); } else if (embed?.$type === "app.bsky.embed.recordWithMedia#view") { if (embed.media?.$type === "app.bsky.embed.images#view") { images = embed.media.images.map((img: any) => ({ thumb: img.thumb, fullsize: img.fullsize, alt: img.alt ?? "", })); } } // Extract external link card let external: FeedPost["external"]; if (embed?.$type === "app.bsky.embed.external#view") { external = { uri: embed.external.uri, title: embed.external.title, description: embed.external.description, thumb: embed.external.thumb, }; } // Extract quote post let quotedPost: FeedPost["quotedPost"]; const quotedRecord = embed?.$type === "app.bsky.embed.record#view" ? embed.record : embed?.$type === "app.bsky.embed.recordWithMedia#view" ? embed.record?.record : null; if (quotedRecord?.value?.text) { quotedPost = { uri: quotedRecord.uri, text: quotedRecord.value.text, author: { handle: quotedRecord.author?.handle ?? "", displayName: quotedRecord.author?.displayName, avatar: quotedRecord.author?.avatar, }, }; } return { uri: post.uri, cid: post.cid, text: record.text ?? "", createdAt: record.createdAt, author: { did: post.author.did, handle: post.author.handle, displayName: post.author.displayName, avatar: post.author.avatar, }, likeCount: post.likeCount ?? 0, repostCount: post.repostCount ?? 0, replyCount: post.replyCount ?? 0, images, external, quotedPost, isReply: !!item.reply, }; }); }