// Bluesky ingest over the public AT Protocol — no auth, no external binary, // no browser. Verified live against public.api.bsky.app: // com.atproto.identity.resolveHandle -> { did } // app.bsky.feed.getAuthorFeed -> { feed: [...], cursor } // both 200 unauthenticated, with full post text, ISO-8601 ms `createdAt` and // cursor paging. // // Bluesky leads the whole posts feature deliberately: it is free, complete and // verified working, so the model / storage / index / search / AI pipeline can // be proven end-to-end before touching anything as fragile as X. // // Note on PDS resolution: bsky.social 401s for accounts hosted on other PDSes, // so any request that must hit the account's own PDS (the getRepo CAR backfill) // resolves the service endpoint from the DID document at plc.directory first. // The read-only AppView (public.api.bsky.app) needs no such resolution. import { postPermalink, uploadDateFromCreatedAt, type Post, type PostAvailability, type PostRef, } from "../lib/posts"; import { registerSocialFetcher, type PostFetchInput, type PostFetchResult, type SocialFetcher, type SocialFetcherProbe, } from "./fetchers"; const APPVIEW = "https://public.api.bsky.app"; const PLC_DIRECTORY = "https://plc.directory"; // getAuthorFeed's documented ceiling. const PAGE_LIMIT = 100; // Safety valve so a first-run backfill of a prolific account can't page // forever inside one job. const MAX_PAGES = 200; // --------------------------------------------------------------------------- // Wire types (only the fields we consume) // --------------------------------------------------------------------------- type BskyAuthor = { did?: string; handle?: string; displayName?: string; }; type BskyFacetFeature = { $type?: string; uri?: string; }; type BskyFacet = { features?: BskyFacetFeature[]; }; type BskyStrongRef = { uri?: string; cid?: string }; type BskyRecord = { text?: string; createdAt?: string; langs?: string[]; facets?: BskyFacet[]; reply?: { root?: BskyStrongRef; parent?: BskyStrongRef }; embed?: BskyEmbed; }; type BskyEmbed = { $type?: string; images?: unknown[]; media?: { images?: unknown[] }; external?: { uri?: string }; record?: BskyStrongRef & { record?: BskyStrongRef }; }; type BskyPostView = { uri?: string; cid?: string; author?: BskyAuthor; record?: BskyRecord; embed?: BskyEmbed; replyCount?: number; repostCount?: number; likeCount?: number; quoteCount?: number; indexedAt?: string; }; type BskyReplyRef = { root?: BskyPostView; parent?: BskyPostView; }; type BskyReason = { $type?: string; by?: BskyAuthor; indexedAt?: string; }; type BskyFeedItem = { post?: BskyPostView; reply?: BskyReplyRef; reason?: BskyReason; }; type BskyFeedResponse = { feed?: BskyFeedItem[]; cursor?: string; }; // --------------------------------------------------------------------------- // Helpers // --------------------------------------------------------------------------- // at:///app.bsky.feed.post/ -> rkey. The rkey is the post's native // id, the atproto analogue of a tweet id. export function rkeyFromAtUri(uri: string | undefined): string | null { if (!uri) return null; const idx = uri.lastIndexOf("/"); if (idx < 0 || idx === uri.length - 1) return null; return uri.slice(idx + 1); } export function didFromAtUri(uri: string | undefined): string | null { if (!uri) return null; const match = /^at:\/\/([^/]+)\//.exec(uri); return match ? match[1] : null; } function refFromPostView(view: BskyPostView | undefined): PostRef | null { if (!view) return null; const id = rkeyFromAtUri(view.uri); if (!id) return null; const author = view.author?.handle ?? ""; const ref: PostRef = { platform: "bluesky", id }; if (author) { ref.author = author; ref.url = postPermalink("bluesky", author, id); } return ref; } // A strong ref carries only a URI, so the handle is unavailable; the DID // works in a bsky.app profile URL, so the link still resolves. function refFromStrongRef(ref: BskyStrongRef | undefined): PostRef | null { if (!ref?.uri) return null; const id = rkeyFromAtUri(ref.uri); if (!id) return null; const did = didFromAtUri(ref.uri); const out: PostRef = { platform: "bluesky", id }; if (did) { out.author = did; out.url = postPermalink("bluesky", did, id); } return out; } // Outbound URLs: richtext link facets plus an external-embed card's target. function linksFrom(record: BskyRecord | undefined, embed: BskyEmbed | undefined): string[] { const links = new Set(); for (const facet of record?.facets ?? []) { for (const feature of facet.features ?? []) { if (feature.$type === "app.bsky.richtext.facet#link" && feature.uri) { links.add(feature.uri); } } } const external = record?.embed?.external?.uri ?? embed?.external?.uri; if (external) links.add(external); return [...links]; } // Media is counted, not archived (v1). Images can hang off either a plain // image embed or the media half of a recordWithMedia embed. function mediaCountFrom(embed: BskyEmbed | undefined): number { if (!embed) return 0; if (Array.isArray(embed.images)) return embed.images.length; if (Array.isArray(embed.media?.images)) return embed.media!.images!.length; return 0; } function quotedRefFrom( record: BskyRecord | undefined, embed: BskyEmbed | undefined, ): PostRef | null { // app.bsky.embed.record -> record is the strong ref; // app.bsky.embed.recordWithMedia -> record.record is. const candidate = record?.embed?.record ?? embed?.record; if (!candidate) return null; const inner = (candidate as { record?: BskyStrongRef }).record; return refFromStrongRef(inner ?? (candidate as BskyStrongRef)); } // Normalize one feed item into a Post. Returns null for items we can't // identify (missing uri/text/createdAt). // // Reposts: for a repost the feed item's `post` IS the original post, and // `reason.by` is the account that reposted it. We archive the original record // with isRepost + repostOf set — i.e. "this post entered the archive because // the tracked account reposted it" — which keeps attribution honest and lets // the archive dedupe by the original id. export function normalizeBlueskyItem( item: BskyFeedItem, channelSlug: string, ): Post | null { const view = item.post; if (!view) return null; const id = rkeyFromAtUri(view.uri); if (!id) return null; const record = view.record; const createdAt = record?.createdAt ?? view.indexedAt; if (!createdAt) return null; const text = typeof record?.text === "string" ? record.text : ""; const author = view.author?.handle ?? ""; const isRepost = item.reason?.$type === "app.bsky.feed.defs#reasonRepost"; const replyParent = refFromPostView(item.reply?.parent) ?? refFromStrongRef(record?.reply?.parent); const replyRoot = refFromPostView(item.reply?.root) ?? refFromStrongRef(record?.reply?.root); const post: Post = { id, slug: `${channelSlug}/${id}`, channelSlug, author, createdAt, uploadDate: uploadDateFromCreatedAt(createdAt), text, url: postPermalink("bluesky", author, id), platform: "bluesky", isReply: Boolean(replyParent ?? replyRoot), isRepost, links: linksFrom(record, view.embed), }; const displayName = view.author?.displayName; if (displayName) post.authorName = displayName; const lang = record?.langs?.[0]; if (lang) post.lang = lang; // A root post's thread id is its own id, so every post carries a usable // grouping key without a second lookup. post.threadId = replyRoot?.id ?? id; if (replyParent) post.replyTo = replyParent; const quoted = quotedRefFrom(record, view.embed); if (quoted) post.quoted = quoted; if (isRepost) { post.repostOf = { platform: "bluesky", id, ...(author ? { author, url: postPermalink("bluesky", author, id) } : {}), }; } const media = mediaCountFrom(view.embed); if (media > 0) post.mediaCount = media; const engagement: NonNullable = {}; if (typeof view.likeCount === "number") engagement.likes = view.likeCount; if (typeof view.repostCount === "number") engagement.reposts = view.repostCount; if (typeof view.replyCount === "number") engagement.replies = view.replyCount; if (typeof view.quoteCount === "number") engagement.quotes = view.quoteCount; if (Object.keys(engagement).length > 0) post.engagement = engagement; return post; } // --------------------------------------------------------------------------- // API calls // --------------------------------------------------------------------------- async function getJson( url: string, signal: AbortSignal | undefined, ): Promise { const res = await fetch(url, { signal, headers: { accept: "application/json" }, }); if (!res.ok) { throw new Error(`${res.status} ${res.statusText} for ${url}`); } return (await res.json()) as T; } export async function resolveHandle( handle: string, signal?: AbortSignal, ): Promise { const url = `${APPVIEW}/xrpc/com.atproto.identity.resolveHandle?handle=${encodeURIComponent(handle)}`; const body = await getJson<{ did?: string }>(url, signal); if (!body.did) throw new Error(`No DID for handle ${handle}`); return body.did; } type BskyProfile = { did?: string; handle?: string; displayName?: string; description?: string; avatar?: string; postsCount?: number; }; export async function getProfile( actor: string, signal?: AbortSignal, ): Promise { const url = `${APPVIEW}/xrpc/app.bsky.actor.getProfile?actor=${encodeURIComponent(actor)}`; return getJson(url, signal); } // The account's own PDS endpoint, from its DID document. Required — and NOT // optional — for any getRepo CAR backfill: bsky.social 401s for accounts hosted // elsewhere. export async function resolvePdsEndpoint( did: string, signal?: AbortSignal, ): Promise { if (!did.startsWith("did:plc:")) return null; type DidDoc = { service?: Array<{ id?: string; type?: string; serviceEndpoint?: string }>; }; const doc = await getJson( `${PLC_DIRECTORY}/${encodeURIComponent(did)}`, signal, ); for (const service of doc.service ?? []) { if ( service.type === "AtprotoPersonalDataServer" && typeof service.serviceEndpoint === "string" ) { return service.serviceEndpoint.replace(/\/$/, ""); } } return null; } async function getAuthorFeed( actor: string, cursor: string | undefined, signal: AbortSignal | undefined, ): Promise { const params = new URLSearchParams({ actor, limit: String(PAGE_LIMIT), filter: "posts_with_replies", }); if (cursor) params.set("cursor", cursor); return getJson( `${APPVIEW}/xrpc/app.bsky.feed.getAuthorFeed?${params.toString()}`, signal, ); } // --------------------------------------------------------------------------- // Fetcher // --------------------------------------------------------------------------- export const blueskyFetcher: SocialFetcher = { id: "bluesky-atproto", label: "Bluesky (public AT Protocol)", platform: "bluesky", fields: { limit: true }, detect(url: string): boolean { try { const host = new URL(url).hostname.toLowerCase(); return host === "bsky.app" || host.endsWith(".bsky.app"); } catch { return false; } }, async probe(url: string, signal?: AbortSignal): Promise { const handle = handleFromUrl(url); if (!handle) { return { ok: false, error: "Could not read a handle from that URL" }; } try { const profile = await getProfile(handle, signal); const resolved = profile.handle ?? handle; return { ok: true, name: profile.displayName || resolved, handle: resolved, url: `https://bsky.app/profile/${resolved}`, avatar: profile.avatar, description: profile.description, postCount: profile.postsCount, }; } catch (err) { return { ok: false, error: (err as Error).message }; } }, async fetch(input: PostFetchInput): Promise { const { handle, channelSlug, since, seenIds, limit, signal, onLog } = input; const posts: Post[] = []; // Resuming an unfinished backfill: continue from the stored cursor and // ignore the watermark, otherwise the first page (all already-archived) // would stop the run before it reached the missing history. let cursor: string | undefined = input.cursor; const resuming = Boolean(input.cursor); const watermark = resuming ? undefined : since; let complete = false; if (resuming) { onLog?.(`Resuming backfill from cursor ${cursor}`); } for (let page = 0; page < MAX_PAGES; page++) { if (signal.aborted) break; const body = await getAuthorFeed(handle, cursor, signal); const items = body.feed ?? []; if (items.length === 0) { complete = true; break; } let sawKnown = false; let sawOlderThanWatermark = false; let newOnPage = 0; for (const item of items) { const post = normalizeBlueskyItem(item, channelSlug); if (!post) continue; if (seenIds.has(post.id)) { sawKnown = true; continue; } if (watermark && post.createdAt <= watermark) { sawOlderThanWatermark = true; continue; } posts.push(post); newOnPage++; } onLog?.( `Bluesky page ${page + 1}: ${items.length} items, ${newOnPage} new` + (sawKnown ? ", reached already-archived content" : "") + (sawOlderThanWatermark ? ", reached the watermark" : ""), ); cursor = body.cursor; // Stop conditions, in the same spirit as sync()'s "stop at the first // page containing an already-archived entry": we still keep the new // posts found on that page, we just don't page further back. if (sawKnown || sawOlderThanWatermark) { complete = true; break; } if (!cursor) { complete = true; break; } if (limit && posts.length >= limit) { posts.length = limit; break; } } return { posts, complete, cursor }; }, }; function handleFromUrl(url: string): string | null { const trimmed = url.trim(); if (!trimmed) return null; if (!/^https?:\/\//i.test(trimmed)) return trimmed.replace(/^@/, ""); try { const parsed = new URL(trimmed); const segments = parsed.pathname.split("/").filter(Boolean); if (segments[0] === "profile" && segments[1]) { return decodeURIComponent(segments[1]).replace(/^@/, ""); } return null; } catch { return null; } } // app.bsky.feed.getPosts takes up to 25 AT-URIs and returns ONLY the posts // that still exist — so an id absent from the response is a deletion. That // makes the whole sweep a handful of unauthenticated batch calls, which is why // Bluesky can afford to check an entire channel routinely. const GET_POSTS_BATCH = 25; blueskyFetcher.checkAvailability = async ({ ids, handle, signal, onLog }) => { const out = new Map(); if (ids.length === 0) return out; // The AT-URI needs the account's DID, not its handle. let did: string; try { did = await resolveHandle(handle, signal); } catch (err) { // The account itself is unreachable — that is NOT evidence any individual // post was deleted, so report it as such rather than mass-marking. onLog?.(`Could not resolve @${handle}: ${(err as Error).message}`); for (const id of ids) out.set(id, "account_unavailable"); return out; } for (let i = 0; i < ids.length; i += GET_POSTS_BATCH) { if (signal.aborted) break; const batch = ids.slice(i, i + GET_POSTS_BATCH); const params = new URLSearchParams(); for (const id of batch) { params.append("uris", `at://${did}/app.bsky.feed.post/${id}`); } try { const body = await getJson<{ posts?: BskyPostView[] }>( `${APPVIEW}/xrpc/app.bsky.feed.getPosts?${params.toString()}`, signal, ); const alive = new Set( (body.posts ?? []) .map((p) => rkeyFromAtUri(p.uri)) .filter((r): r is string => !!r), ); for (const id of batch) { out.set(id, alive.has(id) ? "available" : "deleted"); } } catch (err) { // A failed batch must not read as 25 deletions. onLog?.(`Availability batch failed: ${(err as Error).message}`); for (const id of batch) out.set(id, "error"); } } onLog?.( `Checked ${out.size} post(s): ${[...out.values()].filter((v) => v === "deleted").length} deleted.`, ); return out; }; registerSocialFetcher(blueskyFetcher);