Experimental Bluesky client for agents
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067206820692070207120722073207420752076207720782079208020812082208320842085208620872088208920902091209220932094209520962097209820992100210121022103210421052106210721082109211021112112211321142115211621172118211921202121212221232124212521262127212821292130213121322133213421352136213721382139214021412142214321442145214621472148214921502151215221532154215521562157215821592160216121622163216421652166216721682169217021712172217321742175217621772178217921802181218221832184218521862187218821892190219121922193219421952196219721982199220022012202220322042205220622072208220922102211221222132214221522162217221822192220222122222223222422252226222722282229223022312232223322342235223622372238223922402241224222432244224522462247224822492250225122522253225422552256225722582259226022612262226322642265226622672268226922702271227222732274227522762277227822792280228122822283228422852286228722882289229022912292229322942295229622972298229923002301230223032304230523062307230823092310231123122313231423152316231723182319232023212322232323242325232623272328232923302331233223332334233523362337233823392340234123422343234423452346234723482349235023512352235323542355235623572358235923602361236223632364236523662367236823692370237123722373237423752376237723782379238023812382238323842385238623872388238923902391239223932394239523962397239823992400240124022403240424052406240724082409241024112412241324142415241624172418241924202421242224232424242524262427242824292430243124322433243424352436243724382439244024412442244324442445244624472448244924502451245224532454245524562457245824592460246124622463246424652466246724682469247024712472247324742475247624772478247924802481248224832484248524862487248824892490249124922493249424952496249724982499250025012502250325042505250625072508250925102511251225132514251525162517251825192520252125222523252425252526252725282529253025312532253325342535253625372538253925402541254225432544254525462547254825492550255125522553255425552556255725582559256025612562256325642565256625672568256925702571257225732574257525762577257825792580258125822583258425852586258725882589259025912592259325942595259625972598259926002601260226032604260526062607260826092610261126122613261426152616261726182619262026212622262326242625262626272628262926302631263226332634263526362637263826392640264126422643264426452646264726482649265026512652265326542655265626572658265926602661266226632664266526662667266826692670267126722673267426752676267726782679268026812682268326842685268626872688268926902691269226932694269526962697269826992700270127022703270427052706270727082709271027112712271327142715271627172718271927202721272227232724272527262727272827292730273127322733273427352736273727382739274027412742274327442745274627472748274927502751275227532754275527562757275827592760276127622763276427652766import { AppBskyEmbedExternal, AppBskyEmbedRecord, AppBskyFeedPost, AtUri, BskyAgent, RichText, jsonToLex, lexToJson,} from '@atproto/api';import { TID } from '@atproto/common-web';import { createHash, randomUUID } from 'node:crypto';import { lookup } from 'node:dns/promises';import { readFile } from 'node:fs/promises';import { request as httpsRequest } from 'node:https';import { isIP, type LookupFunction } from 'node:net';import { homedir } from 'node:os';import type { Author, Facet, FeedOptions, FeedReason, FeedSnapshot, LinkCard, Media, PostNode, ProfileEnrichment, ProfileLabel, ProfileSnapshot, Reader, ReaderOptions, Source, Target, ThreadOptions, ThreadSnapshot,} from './model.ts';import type { FrozenAttachment, LoadedDraft, PreparedPublication, PreparedPublicationPost, PublicationClaim, PublicationRecordInspection, PublicationTargetInput, SocialClient, SocialMutationReceipt, StrongReference,} from './social.ts';
type ObjectValue = { [key: string]: unknown };
const PUBLIC_SERVICE = 'https://public.api.bsky.app';const DEFAULT_PDS = 'https://pds.mlf.one';const REQUEST_TIMEOUT_MS = 15_000;const TID_RECORD_KEY_PATTERN = /^[234567abcdefghij][234567abcdefghijklmnopqrstuvwxyz]{12}$/u;const DEFAULT_THREAD_OPTIONS: Required<ThreadOptions> = { depth: 3, parentHeight: 40, maxNodes: 60,};const DEFAULT_FEED_LIMIT = 20;const MAX_FEED_NODES = 600;
function objectValue(value: unknown): ObjectValue | undefined { return value !== null && typeof value === 'object' && !Array.isArray(value) ? value as ObjectValue : undefined;}
function stringValue(value: unknown): string | undefined { return typeof value === 'string' ? value : undefined;}
function numberValue(value: unknown): number | undefined { return typeof value === 'number' && Number.isFinite(value) ? value : undefined;}
function typeName(value: unknown): string | undefined { return stringValue(objectValue(value)?.$type);}
/** ATProto record-key baseline; AtUri parses path segments but does not enforce this grammar. */function isValidRecordKey(value: string): boolean { return value !== '.' && value !== '..' && /^[a-zA-Z0-9_~.:-]{1,512}$/u.test(value);}
function canonicalPostRecordUri(value: unknown): string | undefined { const uri = stringValue(value); if (!uri?.startsWith('at://')) return undefined; try { const parsed = new AtUri(uri); if ( !parsed.hostname.startsWith('did:') || parsed.collection !== 'app.bsky.feed.post' || !parsed.rkey || !isValidRecordKey(parsed.rkey) || parsed.pathname !== `/${parsed.collection}/${parsed.rkey}` || parsed.search || parsed.hash ) return undefined; return uri; } catch { return undefined; }}
function unique(values: readonly string[]): string[] { return [...new Set(values)];}
function uniqueObjects<T>(values: readonly T[]): T[] { const seen = new Set<string>(); return values.filter((value) => { const key = JSON.stringify(value); if (seen.has(key)) return false; seen.add(key); return true; });}
function normalizeService(value: string, label: string): string { let url: URL; try { url = new URL(value); } catch { throw new Error(`${label} must be an absolute HTTP(S) URL`); } if ((url.protocol !== 'https:' && url.protocol !== 'http:') || url.username || url.password) { throw new Error(`${label} must be an HTTP(S) URL without embedded credentials`); } if (url.pathname !== '/' || url.search || url.hash) { throw new Error(`${label} must identify an origin, not a path, query, or fragment`); } return url.origin;}
function requestSignal(): AbortSignal { return AbortSignal.timeout(REQUEST_TIMEOUT_MS);}
const timedFetch: typeof fetch = Object.assign( (input: Parameters<typeof fetch>[0], init?: Parameters<typeof fetch>[1]): Promise<Response> => { const timeout = requestSignal(); const signal = init?.signal ? AbortSignal.any([init.signal, timeout]) : timeout; return fetch(input, { ...init, signal }); }, { preconnect( url: Parameters<typeof fetch.preconnect>[0], options?: Parameters<typeof fetch.preconnect>[1], ): void { fetch.preconnect(url, options); }, },);
function redactSecrets(value: string, secrets: readonly string[]): string { const needles = unique(secrets.flatMap((secret) => { if (!secret) return []; try { return [secret, encodeURIComponent(secret)]; } catch { return [secret]; } })).filter((needle) => needle.length > 0); const spans: Array<{ start: number; end: number }> = []; for (const needle of needles) { let searchFrom = 0; while (searchFrom < value.length) { const start = value.indexOf(needle, searchFrom); if (start < 0) break; spans.push({ start, end: start + needle.length }); searchFrom = start + 1; } } if (spans.length === 0) return value; spans.sort((left, right) => left.start - right.start || right.end - left.end); const merged: Array<{ start: number; end: number }> = []; for (const span of spans) { const previous = merged.at(-1); if (previous && span.start <= previous.end) { previous.end = Math.max(previous.end, span.end); } else { merged.push({ ...span }); } } let redacted = ''; let copiedThrough = 0; for (const span of merged) { redacted += `${value.slice(copiedThrough, span.start)}[redacted]`; copiedThrough = span.end; } return redacted + value.slice(copiedThrough);}
function sanitizedError( action: string, error: unknown, secrets: readonly string[] = [], safeMessage?: string,): Error { const data = objectValue(error); const status = numberValue(data?.status); let code = safeMessage === undefined ? stringValue(data?.error) : undefined; let message = safeMessage ?? (error instanceof Error ? error.message : String(error)); message = redactSecrets(message, secrets); if (code !== undefined) code = redactSecrets(code, secrets); message = message.replace(/[\r\n\t]+/g, ' ').slice(0, 300).trim(); const detail = [status !== undefined ? `status ${status}` : undefined, code] .filter((part): part is string => Boolean(part)) .join(', '); return new Error(`${action} failed${detail ? ` (${detail})` : ''}: ${message || 'unknown error'}`);}
function authorFrom(value: unknown): Author | undefined { const actor = objectValue(value); const did = stringValue(actor?.did); const handle = stringValue(actor?.handle); if (!actor || !did || !handle) return undefined;
const viewer = objectValue(actor.viewer); const author: Author = { did, handle, displayName: stringValue(actor.displayName) ?? '', }; const avatar = stringValue(actor.avatar); const pronouns = stringValue(actor.pronouns); if (avatar !== undefined) author.avatar = avatar; if (pronouns !== undefined) author.pronouns = pronouns; if (viewer) { author.following = typeof viewer.following === 'string'; author.followedBy = typeof viewer.followedBy === 'string'; } return author;}
function labelValues(...sources: unknown[]): string[] { const labels: string[] = []; for (const source of sources) { const sourceObject = objectValue(source); const entries = Array.isArray(source) ? source : sourceObject?.values; if (!Array.isArray(entries)) continue; for (const entry of entries) { const val = stringValue(objectValue(entry)?.val); if (val !== undefined) labels.push(val); } } return unique(labels);}
function serializableRecord(value: unknown): Record<string, unknown> | undefined { const record = objectValue(value); if (!record) return undefined; try { return JSON.parse(JSON.stringify(lexToJson(record as Parameters<typeof lexToJson>[0]))) as Record<string, unknown>; } catch { return undefined; }}
function profileLabelValues(value: unknown): { labels: ProfileLabel[]; warnings: string[] } { if (!Array.isArray(value)) return { labels: [], warnings: [] }; const labels: ProfileLabel[] = []; const warnings: string[] = []; value.forEach((rawLabel, index) => { const label = objectValue(rawLabel); const value = stringValue(label?.val); const issuerDid = stringValue(label?.src); const uri = stringValue(label?.uri); if (!label || value === undefined || !issuerDid || !uri) { warnings.push(`Profile label ${index + 1} omitted its value, issuer DID, or target URI and was not retained.`); return; } const normalized: ProfileLabel = { value, issuerDid, uri }; const cid = stringValue(label.cid); const createdAt = stringValue(label.cts); const expiresAt = stringValue(label.exp); if (cid !== undefined) normalized.cid = cid; if (typeof label.neg === 'boolean') normalized.negated = label.neg; if (createdAt !== undefined) normalized.createdAt = createdAt; if (expiresAt !== undefined) normalized.expiresAt = expiresAt; if (typeof label.sig === 'string') { normalized.signature = label.sig; } else if (label.sig instanceof Uint8Array) { normalized.signature = Buffer.from(label.sig).toString('base64'); } else { const encoded = stringValue(objectValue(label.sig)?.$bytes); if (encoded !== undefined) normalized.signature = encoded; } labels.push(normalized); }); return { labels, warnings };}
function facetValues(value: unknown, text: string, ownerUri: string): { facets: Facet[]; warnings: string[] } { if (!Array.isArray(value)) return { facets: [], warnings: [] };
const encoded = new TextEncoder().encode(text); const decoder = new TextDecoder('utf-8', { fatal: true }); const facets: Facet[] = []; const warnings: string[] = [];
value.forEach((rawFacet, facetIndex) => { const facet = objectValue(rawFacet); const index = objectValue(facet?.index); const byteStart = numberValue(index?.byteStart); const byteEnd = numberValue(index?.byteEnd); if ( byteStart === undefined || byteEnd === undefined || !Number.isInteger(byteStart) || !Number.isInteger(byteEnd) ) { warnings.push(`Post ${ownerUri} has a facet without integer UTF-8 byte bounds.`); return; }
const unsupported: string[] = []; let facetText = ''; if (byteStart < 0 || byteEnd < byteStart || byteEnd > encoded.length) { unsupported.push('invalid-byte-range'); } else { try { facetText = decoder.decode(encoded.slice(byteStart, byteEnd)); } catch { facetText = new TextDecoder().decode(encoded.slice(byteStart, byteEnd)); unsupported.push('invalid-utf8-boundary'); } }
const targets: Target[] = []; const features = Array.isArray(facet?.features) ? facet.features : []; for (const rawFeature of features) { const feature = objectValue(rawFeature); const featureType = typeName(feature) ?? 'unknown-facet-feature'; if (featureType === 'app.bsky.richtext.facet#mention') { const did = stringValue(feature?.did); if (did !== undefined) { targets.push({ kind: 'profile', value: did, label: facetText, ownerUri, index: facetIndex }); } else { unsupported.push(`${featureType}:malformed`); } } else if (featureType === 'app.bsky.richtext.facet#link') { const uri = stringValue(feature?.uri); if (uri !== undefined) { targets.push({ kind: 'link', value: uri, label: facetText, ownerUri, index: facetIndex }); } else { unsupported.push(`${featureType}:malformed`); } } else if (featureType === 'app.bsky.richtext.facet#tag') { const tag = stringValue(feature?.tag); if (tag !== undefined) { targets.push({ kind: 'tag', value: tag, label: facetText, ownerUri, index: facetIndex }); } else { unsupported.push(`${featureType}:malformed`); } } else { unsupported.push(featureType); } }
facets.push({ byteStart, byteEnd, text: facetText, targets, unsupported: unique(unsupported), }); });
return { facets, warnings };}
function nodeScore(node: PostNode): number { return (node.availability === 'available' ? 100 : 0) + (node.author ? 8 : 0) + (node.text !== undefined ? 4 : 0) + (node.createdAt !== undefined ? 2 : 0) + (node.indexedAt !== undefined ? 2 : 0) + node.facets.length * 2 + node.media.length * 2 + node.quotes.length * 2 + node.cards.length * 2 + node.replies.length;}
function mergeAuthors(primary: Author | undefined, secondary: Author | undefined): Author | undefined { if (!primary) return secondary; if (!secondary) return primary; return { did: primary.did, handle: primary.handle, displayName: primary.displayName, ...(primary.avatar !== undefined || secondary.avatar !== undefined ? { avatar: primary.avatar ?? secondary.avatar } : {}), ...(primary.pronouns !== undefined || secondary.pronouns !== undefined ? { pronouns: primary.pronouns ?? secondary.pronouns } : {}), ...(primary.following !== undefined || secondary.following !== undefined ? { following: primary.following ?? secondary.following } : {}), ...(primary.followedBy !== undefined || secondary.followedBy !== undefined ? { followedBy: primary.followedBy ?? secondary.followedBy } : {}), };}
function mergeNodes(existing: PostNode, candidate: PostNode): PostNode { const [primary, secondary] = nodeScore(candidate) >= nodeScore(existing) ? [candidate, existing] : [existing, candidate]; const unsupported = unique([...primary.unsupported, ...secondary.unsupported]); return { ...primary, cid: primary.cid ?? secondary.cid, author: mergeAuthors(primary.author, secondary.author), text: primary.text ?? secondary.text, createdAt: primary.createdAt ?? secondary.createdAt, indexedAt: primary.indexedAt ?? secondary.indexedAt, parentUri: primary.parentUri ?? secondary.parentUri, rootUri: primary.rootUri ?? secondary.rootUri, replies: unique([...primary.replies, ...secondary.replies]), repliesFetched: primary.repliesFetched || secondary.repliesFetched, facets: uniqueObjects([...primary.facets, ...secondary.facets]), media: uniqueObjects([...primary.media, ...secondary.media]), quotes: unique([...primary.quotes, ...secondary.quotes]), cards: uniqueObjects([...primary.cards, ...secondary.cards]), labels: unique([...primary.labels, ...secondary.labels]), unsupported: primary.availability === 'available' ? unsupported.filter((item) => item !== 'unhydrated-record-reference') : unsupported, counts: { replies: primary.counts.replies ?? secondary.counts.replies, likes: primary.counts.likes ?? secondary.counts.likes, reposts: primary.counts.reposts ?? secondary.counts.reposts, quotes: primary.counts.quotes ?? secondary.counts.quotes, }, };}
function unavailableNode(value: unknown, fallbackUri?: string): PostNode | undefined { const raw = objectValue(value); const uri = stringValue(raw?.uri) ?? fallbackUri; if (!uri) return undefined; const rawType = typeName(raw); const availability: PostNode['availability'] = raw?.notFound === true ? 'not-found' : raw?.blocked === true ? 'blocked' : 'unsupported'; return { uri, availability, replies: [], repliesFetched: false, facets: [], media: [], quotes: [], cards: [], labels: [], unsupported: availability === 'unsupported' ? [rawType ?? 'unknown-thread-node'] : [], counts: {}, };}
function availableNode(value: unknown, repliesSource?: unknown): { node?: PostNode; warnings: string[] } { const post = objectValue(value); const uri = stringValue(post?.uri); const record = objectValue(post?.record) ?? objectValue(post?.value); if (!post || !uri || !record) return { warnings: [] };
const text = stringValue(record.text); const facets = facetValues(record.facets, text ?? '', uri); const reply = objectValue(record.reply); const parentUri = stringValue(objectValue(reply?.parent)?.uri); const rootUri = stringValue(objectValue(reply?.root)?.uri); const rawReplies = objectValue(repliesSource)?.replies; const replies = Array.isArray(rawReplies) ? rawReplies.flatMap((entry) => { const entryObject = objectValue(entry); const entryPost = objectValue(entryObject?.post); const replyUri = stringValue(entryPost?.uri) ?? stringValue(entryObject?.uri); return replyUri === undefined ? [] : [replyUri]; }) : []; const replyCount = numberValue(post.replyCount); const node: PostNode = { uri, availability: 'available', replies, repliesFetched: Array.isArray(rawReplies), facets: facets.facets, media: [], quotes: [], cards: [], labels: labelValues(post.labels, record.labels), unsupported: [], counts: { replies: replyCount, likes: numberValue(post.likeCount), reposts: numberValue(post.repostCount), quotes: numberValue(post.quoteCount), }, }; const cid = stringValue(post.cid); const author = authorFrom(post.author); const createdAt = stringValue(record.createdAt); const indexedAt = stringValue(post.indexedAt); if (cid !== undefined) node.cid = cid; if (author !== undefined) node.author = author; if (text !== undefined) node.text = text; if (createdAt !== undefined) node.createdAt = createdAt; if (indexedAt !== undefined) node.indexedAt = indexedAt; if (parentUri !== undefined) node.parentUri = parentUri; if (rootUri !== undefined) node.rootUri = rootUri; const recordType = typeName(record); if (recordType && recordType !== 'app.bsky.feed.post') { node.unsupported.push(`record:${recordType}`); } if (record.embed !== undefined && post.embed === undefined && !Array.isArray(post.embeds)) { node.unsupported.push(`unhydrated:${typeName(record.embed) ?? 'unknown-embed'}`); } return { node, warnings: facets.warnings };}
interface NormalizationContext { nodes: Map<string, PostNode>; quoteExpansions: Map<string, number>; replyExpansions: Map<string, number>; options: Required<ThreadOptions>; warnings: string[]; locallyTruncated: boolean;}
function addNode(context: NormalizationContext, node: PostNode): boolean { const existing = context.nodes.get(node.uri); if (existing) { context.nodes.set(node.uri, mergeNodes(existing, node)); return true; } if (context.nodes.size >= context.options.maxNodes) { context.locallyTruncated = true; return false; } context.nodes.set(node.uri, node); return true;}
function appendOwnerData( context: NormalizationContext, ownerUri: string, values: { media?: Media[]; quotes?: string[]; cards?: LinkCard[]; unsupported?: string[] },): void { const owner = context.nodes.get(ownerUri); if (!owner) return; context.nodes.set(ownerUri, { ...owner, media: values.media ? uniqueObjects([...owner.media, ...values.media]) : owner.media, quotes: values.quotes ? unique([...owner.quotes, ...values.quotes]) : owner.quotes, cards: values.cards ? uniqueObjects([...owner.cards, ...values.cards]) : owner.cards, unsupported: values.unsupported ? unique([...owner.unsupported, ...values.unsupported]) : owner.unsupported, });}
function ingestRawQuoteReference(context: NormalizationContext, ownerUri: string, value: unknown): void { const embed = objectValue(value); const embedType = typeName(embed); const strongRef = embedType === 'app.bsky.embed.record' ? objectValue(embed?.record) : embedType === 'app.bsky.embed.recordWithMedia' ? objectValue(objectValue(embed?.record)?.record) : undefined; if (!strongRef) return; const uri = canonicalPostRecordUri(strongRef.uri); if (!uri) { appendOwnerData(context, ownerUri, { unsupported: [`${embedType ?? 'record-embed'}:invalid-uri`] }); return; } if (context.nodes.get(ownerUri)?.quotes.includes(uri)) return; appendOwnerData(context, ownerUri, { quotes: [uri] }); const placeholder = unavailableNode({ $type: 'unhydrated-record-reference', uri, }); if (placeholder) addNode(context, placeholder);}
function ingestQuoteRecord(context: NormalizationContext, ownerUri: string, value: unknown): void { const quote = objectValue(value); const rawQuoteUri = stringValue(quote?.uri); if (!quote || !rawQuoteUri) { appendOwnerData(context, ownerUri, { unsupported: [`${typeName(quote) ?? 'record-embed'}:malformed`] }); return; } const quoteUri = canonicalPostRecordUri(rawQuoteUri); if (!quoteUri) { appendOwnerData(context, ownerUri, { unsupported: [`${typeName(quote) ?? 'record-embed'}:invalid-uri`] }); return; }
appendOwnerData(context, ownerUri, { quotes: [quoteUri] }); const normalized = objectValue(quote?.value) ? availableNode(quote) : { node: unavailableNode(quote), warnings: [] }; if (!normalized.node || !addNode(context, normalized.node)) return; context.warnings.push(...normalized.warnings);
const quoteType = typeName(quote); if (quoteType === 'app.bsky.embed.record#viewRecord' || objectValue(quote.value)) { const embedded = Array.isArray(quote.embeds) ? quote.embeds : []; const expansionScore = nodeScore(normalized.node) + embedded.length * 20; const previousExpansion = context.quoteExpansions.get(quoteUri) ?? -1; if (expansionScore <= previousExpansion) return; context.quoteExpansions.set(quoteUri, expansionScore); ingestEmbeds(context, quoteUri, embedded); ingestRawQuoteReference(context, quoteUri, objectValue(quote.value)?.embed); } else if ( quoteType !== 'app.bsky.embed.record#viewNotFound' && quoteType !== 'app.bsky.embed.record#viewBlocked' && quoteType !== 'app.bsky.embed.record#viewDetached' && quote?.notFound !== true && quote?.blocked !== true && quote?.detached !== true ) { appendOwnerData(context, ownerUri, { unsupported: [quoteType ?? 'unknown-record-embed'] }); }}
function ingestEmbed(context: NormalizationContext, ownerUri: string, value: unknown): void { const embed = objectValue(value); const embedType = typeName(embed) ?? 'unknown-embed'; if (!embed) { appendOwnerData(context, ownerUri, { unsupported: [embedType] }); return; }
if (embedType === 'app.bsky.embed.images#view') { if (!Array.isArray(embed.images)) { appendOwnerData(context, ownerUri, { unsupported: [`${embedType}:malformed`] }); return; } const media: Media[] = []; const unsupported: string[] = []; embed.images.forEach((rawImage, index) => { const image = objectValue(rawImage); if (!image) { unsupported.push(`${embedType}:image-${index + 1}-malformed`); return; } const aspect = objectValue(image.aspectRatio); const item: Media = { kind: 'image' }; const url = stringValue(image.fullsize); const thumbnail = stringValue(image.thumb); const alt = stringValue(image.alt); const width = numberValue(aspect?.width); const height = numberValue(aspect?.height); if (url !== undefined) item.url = url; if (thumbnail !== undefined) item.thumbnail = thumbnail; if (alt !== undefined) item.alt = alt; if (width !== undefined) item.width = width; if (height !== undefined) item.height = height; media.push(item); }); appendOwnerData(context, ownerUri, { media, unsupported }); return; }
if (embedType === 'app.bsky.embed.video#view') { const aspect = objectValue(embed.aspectRatio); const item: Media = { kind: 'video' }; const url = stringValue(embed.playlist); const thumbnail = stringValue(embed.thumbnail); const alt = stringValue(embed.alt); const width = numberValue(aspect?.width); const height = numberValue(aspect?.height); if (url !== undefined) item.url = url; if (thumbnail !== undefined) item.thumbnail = thumbnail; if (alt !== undefined) item.alt = alt; if (width !== undefined) item.width = width; if (height !== undefined) item.height = height; appendOwnerData(context, ownerUri, { media: [item] }); return; }
if (embedType === 'app.bsky.embed.external#view') { const external = objectValue(embed.external); const uri = stringValue(external?.uri); const title = stringValue(external?.title); const description = stringValue(external?.description); if (uri === undefined || title === undefined || description === undefined) { appendOwnerData(context, ownerUri, { unsupported: [`${embedType}:malformed`] }); return; } const card: LinkCard = { uri, title, description }; const thumbnail = stringValue(external?.thumb); if (thumbnail !== undefined) card.thumbnail = thumbnail; appendOwnerData(context, ownerUri, { cards: [card] }); return; }
if (embedType === 'app.bsky.embed.record#view') { ingestQuoteRecord(context, ownerUri, embed.record); return; }
if (embedType === 'app.bsky.embed.recordWithMedia#view') { ingestEmbed(context, ownerUri, embed.media); const record = objectValue(embed.record); if (typeName(record) === 'app.bsky.embed.record#view' || record?.record !== undefined) { ingestQuoteRecord(context, ownerUri, record?.record); } else { appendOwnerData(context, ownerUri, { unsupported: [`${embedType}:malformed-record`] }); } return; }
appendOwnerData(context, ownerUri, { unsupported: [embedType] });}
function ingestEmbeds(context: NormalizationContext, ownerUri: string, values: unknown[]): void { for (const value of values) ingestEmbed(context, ownerUri, value);}
function normalizeThreadEntry(value: unknown, fallbackUri?: string): { node?: PostNode; warnings: string[] } { const entry = objectValue(value); const post = objectValue(entry?.post); if (post) return availableNode(post, entry); return { node: unavailableNode(entry, fallbackUri), warnings: [] };}
function entryUri(value: unknown): string | undefined { const entry = objectValue(value); return stringValue(objectValue(entry?.post)?.uri) ?? stringValue(entry?.uri);}
function ingestThreadEntry( context: NormalizationContext, value: unknown, fallbackUri?: string, ingestOwnEmbeds = true,): string | undefined { const normalized = normalizeThreadEntry(value, fallbackUri); if (!normalized.node || !addNode(context, normalized.node)) return normalized.node?.uri; context.warnings.push(...normalized.warnings); if (ingestOwnEmbeds && normalized.node.availability === 'available') { const post = objectValue(objectValue(value)?.post); const embed = post?.embed; if (embed !== undefined) ingestEmbed(context, normalized.node.uri, embed); ingestRawQuoteReference(context, normalized.node.uri, objectValue(post?.record)?.embed); } return normalized.node.uri;}
function ingestReplies(context: NormalizationContext, value: unknown, level: number): void { if (level > context.options.depth) return; const entry = objectValue(value); if (!Array.isArray(entry?.replies)) return; for (const reply of entry.replies) { const admittedUri = ingestThreadEntry(context, reply); const admittedNode = admittedUri ? context.nodes.get(admittedUri) : undefined; if (!admittedUri || !admittedNode) continue;
const replyEntry = objectValue(reply); const childCount = Array.isArray(replyEntry?.replies) ? replyEntry.replies.length : 0; const expansionScore = nodeScore(admittedNode) + childCount * 20; const previousExpansion = context.replyExpansions.get(admittedUri) ?? -1; if (expansionScore <= previousExpansion) continue; context.replyExpansions.set(admittedUri, expansionScore); ingestReplies(context, reply, level + 1); }}
/** * Losslessly converts a standard app.bsky.feed.getPostThread response union into * Perch's serializable graph. It deliberately accepts unknown for fixture-driven * contract tests and for forward-compatible API unions. */export function normalizeThread( thread: unknown, context: { requested: string; focusUri: string; fetchedAt: string; source: Source; options: Required<ThreadOptions>; },): ThreadSnapshot { const state: NormalizationContext = { nodes: new Map(), quoteExpansions: new Map(), replyExpansions: new Map(), options: context.options, warnings: [], locallyTruncated: false, };
const focusUri = ingestThreadEntry(state, thread, context.focusUri, false) ?? context.focusUri; if (!state.nodes.has(focusUri)) { const fallbackFocus = unavailableNode(thread, focusUri); if (fallbackFocus) addNode(state, fallbackFocus); }
const parentsNearestFirst: unknown[] = []; const parentUris = new Set<string>([focusUri]); let parent = objectValue(thread)?.parent; while (parent !== undefined && parentsNearestFirst.length < context.options.parentHeight) { const uri = entryUri(parent); if (uri && parentUris.has(uri)) { state.warnings.push(`A parent cycle at ${uri} was ignored.`); break; } if (uri) parentUris.add(uri); parentsNearestFirst.push(parent); parent = objectValue(parent)?.parent; } if (parent !== undefined) { state.warnings.push(`Parent traversal stopped at the configured parentHeight of ${context.options.parentHeight}.`); state.locallyTruncated = true; }
// Reserve the focus, then prefer its nearest ancestors when maxNodes is small. for (const rawParent of parentsNearestFirst) ingestThreadEntry(state, rawParent, undefined, false);
const ancestorChain = [...parentsNearestFirst].reverse(); const admittedAncestorUris: string[] = []; for (const rawParent of ancestorChain) { const uri = entryUri(rawParent); const node = uri ? state.nodes.get(uri) : undefined; if (uri && node?.availability === 'available') admittedAncestorUris.push(uri); }
// Connect the parent chain even when the API only expresses it on each child. const chainNearestFirst = [thread, ...parentsNearestFirst]; for (let index = 0; index + 1 < chainNearestFirst.length; index += 1) { const childUri = entryUri(chainNearestFirst[index]) ?? (index === 0 ? focusUri : undefined); const parentUri = entryUri(chainNearestFirst[index + 1]); const parentNode = parentUri ? state.nodes.get(parentUri) : undefined; if (childUri && parentUri && parentNode) { state.nodes.set(parentUri, { ...parentNode, replies: unique([...parentNode.replies, childUri]) }); } }
for (const rawParent of ancestorChain) { const uri = entryUri(rawParent); const node = uri ? state.nodes.get(uri) : undefined; const post = objectValue(objectValue(rawParent)?.post); if (uri && node?.availability === 'available') { if (post?.embed !== undefined) ingestEmbed(state, uri, post.embed); ingestRawQuoteReference(state, uri, objectValue(post?.record)?.embed); } } const focusPost = objectValue(objectValue(thread)?.post); if (focusPost?.embed !== undefined) ingestEmbed(state, focusUri, focusPost.embed); ingestRawQuoteReference(state, focusUri, objectValue(focusPost?.record)?.embed); ingestReplies(state, thread, 1);
const focus = state.nodes.get(focusUri); if (focus) { const reply = objectValue(objectValue(focusPost?.record)?.reply); const parentUri = stringValue(objectValue(reply?.parent)?.uri); const rootUri = stringValue(objectValue(reply?.root)?.uri); state.nodes.set(focusUri, { ...focus, ...(parentUri !== undefined ? { parentUri } : {}), ...(rootUri !== undefined ? { rootUri } : {}), }); }
const remoteUnknown = [...state.nodes.values()].filter((node) => node.availability === 'available' && !node.repliesFetched ).length; if (remoteUnknown > 0) { state.warnings.push( `Reply coverage is unknown for ${remoteUnknown} post(s); the API supplied no cursor or certain remaining count.`, ); } if (state.locallyTruncated) { state.warnings.push('Local maxNodes or safety bounds omitted one or more supplied nodes.'); }
const orderedUris = unique([ ...admittedAncestorUris, focusUri, ...state.nodes.keys(), ]); return { kind: 'thread', requested: context.requested, focusUri, fetchedAt: context.fetchedAt, source: { ...context.source }, nodes: orderedUris.flatMap((uri) => { const node = state.nodes.get(uri); return node === undefined ? [] : [node]; }), ancestors: admittedAncestorUris, replyOrder: 'source', limits: { ...context.options }, locallyTruncated: state.locallyTruncated, warnings: unique(state.warnings), };}
function feedReason(value: unknown, warnings: string[], occurrenceId: string): FeedReason | undefined { if (value === undefined) return undefined; const raw = objectValue(value); if (!raw) { warnings.push(`Feed occurrence ${occurrenceId} supplied a malformed unknown reason.`); return { kind: 'unknown' }; }
const rawType = typeName(raw); const kind: FeedReason['kind'] = rawType === 'app.bsky.feed.defs#reasonRepost' ? 'repost' : rawType === 'app.bsky.feed.defs#reasonPin' ? 'pin' : 'unknown'; const reason: FeedReason = { kind }; const by = authorFrom(raw.by); const indexedAt = stringValue(raw.indexedAt); if (by !== undefined) reason.by = by; if (indexedAt !== undefined) reason.indexedAt = indexedAt; if (kind === 'unknown' && rawType !== undefined) reason.type = rawType; if (raw.by !== undefined && by === undefined) { warnings.push(`Feed occurrence ${occurrenceId} supplied a reason actor without a usable DID and handle.`); } return reason;}
function normalizedFeedNode(value: unknown, fallbackUri?: string): { node?: PostNode; warnings: string[] } { const normalized = availableNode(value); if (normalized.node) return normalized; return { node: unavailableNode(value, fallbackUri), warnings: [] };}
function ingestFeedContextPost(context: NormalizationContext, value: unknown): string | undefined { const uri = canonicalPostRecordUri(objectValue(value)?.uri); if (!uri) return undefined; const normalized = normalizedFeedNode(value, uri); if (!normalized.node || !addNode(context, normalized.node)) return uri; context.warnings.push(...normalized.warnings); if (normalized.node.availability === 'available') { const post = objectValue(value); if (post?.embed !== undefined) ingestEmbed(context, uri, post.embed); ingestRawQuoteReference(context, uri, objectValue(post?.record)?.embed); } return uri;}
/** Losslessly normalizes one source-ordered following or author-feed page. */export function normalizeFeed( data: unknown, context: { fetchedAt: string; source: Source; limit: number; cursor?: string; feed?: 'following' | 'author'; actor?: Author; },): FeedSnapshot { const requested = boundedFeedOptions({ limit: context.limit, cursor: context.cursor }); const feed = context.feed ?? 'following'; if (feed === 'author' && !context.actor) { throw new Error('Author-feed normalization requires its canonical actor'); } const response = objectValue(data); if (!response || !Array.isArray(response.feed)) { throw new Error('Feed response omitted its feed array'); }
const state: NormalizationContext = { nodes: new Map(), quoteExpansions: new Map(), replyExpansions: new Map(), options: { depth: 0, parentHeight: 0, maxNodes: MAX_FEED_NODES }, warnings: [], locallyTruncated: false, }; const entries: FeedSnapshot['entries'] = []; const frames: Array<{ post: ObjectValue; postUri: string; parent?: unknown; root?: unknown; parentUri?: string; rootUri?: string; }> = [];
// Admit every distinct primary post before optional context or quote nodes can consume the page cap. response.feed.forEach((rawEntry, index) => { const occurrenceId = String(index + 1); const entry = objectValue(rawEntry); const post = objectValue(entry?.post); const postUri = canonicalPostRecordUri(post?.uri); if (!entry || !post || !postUri) { throw new Error(`Feed occurrence ${occurrenceId} omitted a usable canonical post URI`); }
const normalized = normalizedFeedNode(post, postUri); if (!normalized.node) { throw new Error(`Feed occurrence ${occurrenceId} could not retain its primary post`); } addNode(state, normalized.node); state.warnings.push(...normalized.warnings);
const reply = objectValue(entry.reply); const parent = reply?.parent; const root = reply?.root; const suppliedParentUri = canonicalPostRecordUri(objectValue(parent)?.uri); const suppliedRootUri = canonicalPostRecordUri(objectValue(root)?.uri); if (parent !== undefined && suppliedParentUri === undefined) { state.warnings.push(`Feed occurrence ${occurrenceId} supplied parent context without a usable canonical post URI.`); } if (root !== undefined && suppliedRootUri === undefined) { state.warnings.push(`Feed occurrence ${occurrenceId} supplied root context without a usable canonical post URI.`); }
const retainedPrimary = state.nodes.get(postUri) ?? normalized.node; const parentUri = suppliedParentUri ?? retainedPrimary.parentUri; const rootUri = suppliedRootUri ?? retainedPrimary.rootUri; const reason = feedReason(entry.reason, state.warnings, occurrenceId); entries.push({ id: occurrenceId, postUri, ...(reason !== undefined ? { reason } : {}), ...(parentUri !== undefined ? { parentUri } : {}), ...(rootUri !== undefined ? { rootUri } : {}), }); frames.push({ post, postUri, ...(parent !== undefined ? { parent } : {}), ...(root !== undefined ? { root } : {}), ...(parentUri !== undefined ? { parentUri } : {}), ...(rootUri !== undefined ? { rootUri } : {}), }); });
for (const frame of frames) { const primary = state.nodes.get(frame.postUri); if (primary) { state.nodes.set(frame.postUri, { ...primary, parentUri: primary.parentUri ?? frame.parentUri, rootUri: primary.rootUri ?? frame.rootUri, }); }
if (frame.parent !== undefined) ingestFeedContextPost(state, frame.parent); if (frame.root !== undefined) ingestFeedContextPost(state, frame.root); const parentUri = frame.parentUri; if (parentUri) { const parentNode = state.nodes.get(parentUri); if (parentNode) { state.nodes.set(parentUri, { ...parentNode, replies: unique([...parentNode.replies, frame.postUri]), }); } }
if (frame.post.embed !== undefined) ingestEmbed(state, frame.postUri, frame.post.embed); ingestRawQuoteReference(state, frame.postUri, objectValue(frame.post.record)?.embed); }
let cursor: string | undefined; if (typeof response.cursor === 'string' && response.cursor.length > 0) { cursor = response.cursor; } else if (response.cursor === '') { state.warnings.push('Feed response supplied an empty cursor; continuation is unavailable.'); } else if (Object.prototype.hasOwnProperty.call(response, 'cursor') && response.cursor !== undefined) { state.warnings.push('Feed response supplied a non-string cursor; continuation is unavailable.'); } if (state.locallyTruncated) { state.warnings.push('Local feed maxNodes omitted one or more supplied context or quote nodes.'); }
return { kind: 'feed', feed, fetchedAt: context.fetchedAt, source: { ...context.source }, ...(context.actor !== undefined ? { actor: { ...context.actor } } : {}), entries, nodes: [...state.nodes.values()], limit: requested.limit, ...(cursor !== undefined ? { cursor } : {}), ...(requested.cursor !== undefined ? { requestCursor: requested.cursor } : {}), maxNodes: MAX_FEED_NODES, locallyTruncated: state.locallyTruncated, warnings: unique(state.warnings), };}
function boundedFeedOptions(options: FeedOptions = {}): { limit: number; cursor?: string } { const limit = options.limit ?? DEFAULT_FEED_LIMIT; if (!Number.isInteger(limit) || limit < 1 || limit > 100) { throw new Error('limit must be an integer from 1 through 100'); } if (options.cursor !== undefined && (typeof options.cursor !== 'string' || options.cursor.length === 0)) { throw new Error('cursor must be a nonempty opaque string when supplied'); } return { limit, ...(options.cursor !== undefined ? { cursor: options.cursor } : {}), };}
function boundedOptions(options: ThreadOptions = {}): Required<ThreadOptions> { const resolved = { depth: options.depth ?? DEFAULT_THREAD_OPTIONS.depth, parentHeight: options.parentHeight ?? DEFAULT_THREAD_OPTIONS.parentHeight, maxNodes: options.maxNodes ?? DEFAULT_THREAD_OPTIONS.maxNodes, }; for (const [name, value] of [ ['depth', resolved.depth], ['parentHeight', resolved.parentHeight], ] as const) { if (!Number.isInteger(value) || value < 0 || value > 100) { throw new Error(`${name} must be an integer from 0 through 100`); } } if (!Number.isInteger(resolved.maxNodes) || resolved.maxNodes < 1 || resolved.maxNodes > 200) { throw new Error('maxNodes must be an integer from 1 through 200'); } return resolved;}
function postInputParts(input: string): { actor: string; rkey: string } { const trimmed = input.trim(); if (trimmed.startsWith('at://')) { let uri: AtUri; try { uri = new AtUri(trimmed); } catch { throw new Error('Post target is not a valid AT URI'); } if ( uri.collection !== 'app.bsky.feed.post' || !uri.rkey || !isValidRecordKey(uri.rkey) || uri.pathname !== `/${uri.collection}/${uri.rkey}` || uri.search || uri.hash ) { throw new Error('Post AT URI must identify one app.bsky.feed.post record'); } return { actor: uri.hostname.startsWith('did:') ? uri.hostname : uri.hostname.toLowerCase(), rkey: uri.rkey, }; }
let url: URL; try { url = new URL(trimmed); } catch { throw new Error('Post target must be a bsky.app post URL or AT URI'); } if (url.protocol !== 'https:' || (url.hostname !== 'bsky.app' && url.hostname !== 'www.bsky.app')) { throw new Error('Post URL must use https://bsky.app'); } const parts = url.pathname.split('/').filter(Boolean); if (parts.length !== 4 || parts[0] !== 'profile' || parts[2] !== 'post') { throw new Error('Post URL must have the form https://bsky.app/profile/<actor>/post/<rkey>'); } let actor: string; let rkey: string; try { actor = decodeURIComponent(parts[1]); rkey = decodeURIComponent(parts[3]); } catch { throw new Error('Post URL contains invalid percent encoding'); } // AtUri.make validates the authority and record-key grammar without guessing from display text. try { const authority = actor.startsWith('did:') ? actor : actor.toLowerCase(); const validated = AtUri.make(authority, 'app.bsky.feed.post', rkey); return { actor: validated.hostname, rkey: validated.rkey }; } catch { throw new Error('Post URL contains an invalid actor or record key'); }}
function profileInputActor(input: string): string { const trimmed = input.trim(); if (trimmed.startsWith('at://')) { try { const authority = new AtUri(trimmed).hostname; return authority.startsWith('did:') ? authority : authority.toLowerCase(); } catch { throw new Error('Profile target contains an invalid AT identifier'); } } if (trimmed.startsWith('https://') || trimmed.startsWith('http://')) { let url: URL; try { url = new URL(trimmed); } catch { throw new Error('Profile target is not a valid URL'); } const parts = url.pathname.split('/').filter(Boolean); if ( url.protocol !== 'https:' || (url.hostname !== 'bsky.app' && url.hostname !== 'www.bsky.app') || parts.length !== 2 || parts[0] !== 'profile' ) { throw new Error('Profile URL must have the form https://bsky.app/profile/<actor>'); } try { const decoded = decodeURIComponent(parts[1]); const authority = decoded.startsWith('did:') ? decoded : decoded.toLowerCase(); return new AtUri(`at://${authority}`).hostname; } catch { throw new Error('Profile URL contains an invalid actor'); } } const actor = trimmed.startsWith('@') ? trimmed.slice(1) : trimmed; try { const authority = actor.startsWith('did:') ? actor : actor.toLowerCase(); return new AtUri(`at://${authority}`).hostname; } catch { throw new Error('Profile target must be a valid handle, DID, AT URI, or bsky.app profile URL'); }}
export function normalizeProfile( value: unknown, context: { requested: string; fetchedAt: string; source: Source; enrichments?: ProfileEnrichment[]; warnings?: string[]; },): ProfileSnapshot { const profile = objectValue(value); const author = authorFrom(profile); if (!profile || !author) throw new Error('Profile response omitted canonical DID or handle'); const normalizedLabels = profileLabelValues(profile.labels); const warnings = unique([...(context.warnings ?? []), ...normalizedLabels.warnings]); const snapshot: ProfileSnapshot = { kind: 'profile', requested: context.requested, fetchedAt: context.fetchedAt, source: { ...context.source }, author, counts: { followers: numberValue(profile.followersCount), follows: numberValue(profile.followsCount), posts: numberValue(profile.postsCount), }, labels: normalizedLabels.labels, };
const description = stringValue(profile.description); const banner = stringValue(profile.banner); const website = stringValue(profile.website); const createdAt = stringValue(profile.createdAt); const indexedAt = stringValue(profile.indexedAt); if (description !== undefined) snapshot.description = description; else if (profile.description !== undefined) warnings.push('Profile field description was supplied with an unsupported non-string value.'); if (banner !== undefined) snapshot.banner = banner; else if (profile.banner !== undefined) warnings.push('Profile field banner was supplied with an unsupported non-string value.'); if (website !== undefined) snapshot.website = website; else if (profile.website !== undefined) warnings.push('Profile field website was supplied with an unsupported non-string value.'); if (createdAt !== undefined) snapshot.createdAt = createdAt; else if (profile.createdAt !== undefined) warnings.push('Profile field createdAt was supplied with an unsupported non-string value.'); if (indexedAt !== undefined) snapshot.indexedAt = indexedAt; else if (profile.indexedAt !== undefined) warnings.push('Profile field indexedAt was supplied with an unsupported non-string value.');
const associated = serializableRecord(profile.associated); const verification = serializableRecord(profile.verification); if (associated !== undefined) snapshot.associated = associated; else if (profile.associated !== undefined) warnings.push('Profile associated metadata was not a serializable object.'); if (verification !== undefined) snapshot.verification = verification; else if (profile.verification !== undefined) warnings.push('Profile verification metadata was not a serializable object.');
if (profile.pinnedPost !== undefined) { const pinned = objectValue(profile.pinnedPost); const uri = canonicalPostRecordUri(pinned?.uri); const cid = stringValue(pinned?.cid); if (uri && cid) snapshot.pinnedPost = { uri, cid }; else warnings.push('Profile pinned post omitted a canonical post URI or CID.'); } if (context.enrichments !== undefined) snapshot.enrichments = context.enrichments; if (warnings.length > 0) snapshot.warnings = unique(warnings); return snapshot;}
async function canonicalPostUri( agent: BskyAgent, input: string, authenticationSecrets: () => readonly string[],): Promise<string> { const { actor, rkey } = postInputParts(input); let did = actor; if (!actor.startsWith('did:')) { try { const response = await agent.resolveHandle({ handle: actor }, { signal: requestSignal() }); did = response.data.did; } catch (error) { throw sanitizedError(`Resolving handle ${actor}`, error, authenticationSecrets()); } } try { return AtUri.make(did, 'app.bsky.feed.post', rkey).toString(); } catch { throw new Error('Resolved post identity did not produce a valid canonical AT URI'); }}
const PROFILE_COLLECTIONS: Record<string, string> = { 'sh.tangled.actor.profile': 'Tangled profile', 'sh.tangled.feed.comment': 'Tangled comments', 'sh.tangled.feed.reaction': 'Tangled reactions', 'sh.tangled.feed.star': 'Tangled stars', 'sh.tangled.knot': 'Tangled knots', 'sh.tangled.knot.member': 'Tangled knot members', 'sh.tangled.label.definition': 'Tangled label definitions', 'sh.tangled.label.op': 'Tangled label operations', 'sh.tangled.repo': 'Tangled repositories', 'sh.tangled.repo.collaborator': 'Tangled collaborators', 'sh.tangled.repo.issue': 'Tangled issues', 'sh.tangled.repo.issue.state': 'Tangled issue state', 'sh.tangled.repo.pull': 'Tangled pull requests', 'sh.tangled.repo.pull.status': 'Tangled pull-request status', 'sh.tangled.spindle': 'Tangled spindles', 'sh.tangled.spindle.member': 'Tangled spindle members', 'app.greengale.document': 'GreenGale documents', 'app.greengale.publication': 'GreenGale publication', 'site.standard.document': 'Standard Site documents', 'site.standard.publication': 'Standard Site publication',};const MAX_DISCOVERED_RESPONSE_BYTES = 1_000_000;
const NON_GLOBAL_IPV4_PREFIXES = [ [0x00000000, 8], // Current network. [0x0a000000, 8], // Private use. [0x64400000, 10], // Shared address space. [0x7f000000, 8], // Loopback. [0xa9fe0000, 16], // Link-local. [0xac100000, 12], // Private use. [0xc0000000, 24], // IETF protocol assignments. [0xc0000200, 24], // Documentation. [0xc0a80000, 16], // Private use. [0xc0586300, 24], // Deprecated 6to4 relay anycast (192.88.99.0/24). [0xc6120000, 15], // Benchmarking. [0xc6336400, 24], // Documentation. [0xcb007100, 24], // Documentation. [0xe0000000, 4], // Multicast. [0xf0000000, 4], // Reserved.] as const;
function ipv4Number(address: string): number | undefined { const parts = address.split('.').map(Number); if (parts.length !== 4 || parts.some((part) => !Number.isInteger(part) || part < 0 || part > 255)) return undefined; return ((parts[0] * 0x1000000) + (parts[1] * 0x10000) + (parts[2] * 0x100) + parts[3]) >>> 0;}
function globallyRoutableIPv4(address: string): boolean { const value = ipv4Number(address); if (value === undefined) return false; return !NON_GLOBAL_IPV4_PREFIXES.some(([network, length]) => ( value >>> (32 - length) ) === ( network >>> (32 - length) ));}
type IPv6Words = [number, number, number, number, number, number, number, number];const NON_GLOBAL_IPV6_PREFIXES: ReadonlyArray<readonly [IPv6Words, number]> = [ [[0, 0, 0, 0, 0, 0, 0, 0], 96], // Unspecified, loopback, and deprecated IPv4-compatible space. [[0, 0, 0, 0, 0xffff, 0, 0, 0], 96], // IPv4-translated space. [[0x0064, 0xff9b, 0x0001, 0, 0, 0, 0, 0], 48], // Local-use translation prefix. [[0x0100, 0, 0, 0, 0, 0, 0, 0], 64], // Discard-only. [[0x2001, 0, 0, 0, 0, 0, 0, 0], 23], // IETF special-purpose assignments. [[0x2001, 0x0db8, 0, 0, 0, 0, 0, 0], 32], // Documentation. [[0x2002, 0, 0, 0, 0, 0, 0, 0], 16], // Deprecated 6to4. [[0x3ffe, 0, 0, 0, 0, 0, 0, 0], 16], // Former 6bone space. [[0x3fff, 0, 0, 0, 0, 0, 0, 0], 20], // Documentation. [[0x5f00, 0, 0, 0, 0, 0, 0, 0], 16], // Segment-routing SIDs. [[0xfc00, 0, 0, 0, 0, 0, 0, 0], 7], // Unique-local. [[0xfe80, 0, 0, 0, 0, 0, 0, 0], 10], // Link-local. [[0xff00, 0, 0, 0, 0, 0, 0, 0], 8], // Multicast.];
function ipv6Words(address: string): IPv6Words | undefined { const halves = address.toLowerCase().split('%')[0].split('::'); if (halves.length > 2) return undefined; const parseHalf = (half: string): number[] | undefined => { if (half === '') return []; const tokens = half.split(':'); const words: number[] = []; for (let index = 0; index < tokens.length; index += 1) { const token = tokens[index]; if (token.includes('.')) { if (index !== tokens.length - 1) return undefined; const embedded = ipv4Number(token); if (embedded === undefined) return undefined; words.push(embedded >>> 16, embedded & 0xffff); } else { if (!/^[0-9a-f]{1,4}$/u.test(token)) return undefined; words.push(Number.parseInt(token, 16)); } } return words; }; const left = parseHalf(halves[0]); const right = parseHalf(halves[1] ?? ''); if (!left || !right) return undefined; const omitted = 8 - left.length - right.length; if ((halves.length === 1 && omitted !== 0) || (halves.length === 2 && omitted < 1)) return undefined; const words = [...left, ...Array.from({ length: omitted }, () => 0), ...right]; return words.length === 8 ? words as IPv6Words : undefined;}
function ipv6PrefixMatches(address: IPv6Words, network: IPv6Words, length: number): boolean { const wholeWords = Math.floor(length / 16); for (let index = 0; index < wholeWords; index += 1) { if (address[index] !== network[index]) return false; } const remaining = length % 16; if (remaining === 0) return true; const mask = (0xffff << (16 - remaining)) & 0xffff; return (address[wholeWords] & mask) === (network[wholeWords] & mask);}
export function globallyRoutableAddress(address: string): boolean { const normalized = address.toLowerCase().replace(/^\[|\]$/gu, ''); const family = isIP(normalized); if (family === 4) return globallyRoutableIPv4(normalized); if (family !== 6) return false; const words = ipv6Words(normalized); if (!words) return false; if (words.slice(0, 5).every((word) => word === 0) && words[5] === 0xffff) { const embedded = `${words[6] >>> 8}.${words[6] & 0xff}.${words[7] >>> 8}.${words[7] & 0xff}`; return globallyRoutableIPv4(embedded); } if (words[0] === 0x0064 && words[1] === 0xff9b && words.slice(2, 6).every((word) => word === 0)) { const embedded = `${words[6] >>> 8}.${words[6] & 0xff}.${words[7] >>> 8}.${words[7] & 0xff}`; return globallyRoutableIPv4(embedded); } if (NON_GLOBAL_IPV6_PREFIXES.some(([network, length]) => ipv6PrefixMatches(words, network, length))) return false; return words[0] >= 0x2000 && words[0] <= 0x3fff;}
export interface PublicAddressResolution { address: string; family: number;}
export interface DiscoveredFetchDependencies { resolveAddress?: ( hostname: string, options: { all: true; verbatim: true }, ) => Promise<PublicAddressResolution[]>; request?: typeof httpsRequest; timeoutMs?: number;}
async function vettedPublicAddress( url: URL, signal: AbortSignal = requestSignal(), resolveAddress: NonNullable<DiscoveredFetchDependencies['resolveAddress']> = async (hostname) => ( lookup(hostname, { all: true, verbatim: true }) ),): Promise<{ address: string; family: 4 | 6 }> { if (signal.aborted) throw new Error('Discovered PDS request was aborted before transport dispatch'); const hostname = url.hostname.toLowerCase().replace(/^\[|\]$/gu, ''); if ( hostname === 'localhost' || hostname.endsWith('.localhost') || hostname.endsWith('.local') || hostname.endsWith('.internal') || hostname.endsWith('.home') || hostname.endsWith('.lan') || hostname.endsWith('.test') || hostname.endsWith('.invalid') || hostname.endsWith('.example') ) throw new Error(`Discovered PDS hostname ${hostname} is not a public destination`);
const literalFamily = isIP(hostname); if (literalFamily !== 0) { if (!globallyRoutableAddress(hostname)) throw new Error(`Discovered PDS address ${hostname} is not globally routable`); return { address: hostname, family: literalFamily as 4 | 6 }; } const abort = Promise.withResolvers<never>(); const abortResolution = () => abort.reject( new Error('Discovered PDS address resolution exceeded the request deadline'), ); signal.addEventListener('abort', abortResolution, { once: true }); if (signal.aborted) abortResolution(); let addresses: PublicAddressResolution[]; try { addresses = await Promise.race([ resolveAddress(hostname, { all: true, verbatim: true }), abort.promise, ]); } finally { signal.removeEventListener('abort', abortResolution); } if (signal.aborted) throw new Error('Discovered PDS request was aborted before transport dispatch'); if (addresses.length === 0) throw new Error(`Discovered PDS hostname ${hostname} resolved to no addresses`); if (addresses.some((address) => !globallyRoutableAddress(address.address))) { throw new Error(`Discovered PDS hostname ${hostname} resolved to a non-global address`); } const selected = addresses[0]; const selectedFamily = isIP(selected.address); if (selectedFamily !== 4 && selectedFamily !== 6) { throw new Error(`Discovered PDS hostname ${hostname} resolved to an invalid address`); } return { address: selected.address, family: selectedFamily };}
export function discoveredFetch( service: string, dependencies: DiscoveredFetchDependencies = {},): typeof fetch { const serviceUrl = new URL(service); const resolveAddress = dependencies.resolveAddress ?? (async (hostname: string) => ( lookup(hostname, { all: true, verbatim: true }) )); const dispatchRequest = dependencies.request ?? httpsRequest; const timeoutMs = dependencies.timeoutMs ?? REQUEST_TIMEOUT_MS; if (!Number.isFinite(timeoutMs) || timeoutMs <= 0) throw new Error('Discovery timeout must be a positive number'); const implementation = async ( input: Parameters<typeof fetch>[0], init?: Parameters<typeof fetch>[1], ): Promise<Response> => { const request = new Request(input, init); const url = new URL(request.url); if (url.origin !== serviceUrl.origin || url.protocol !== 'https:' || url.username || url.password) { throw new Error('Discovered PDS request escaped its vetted HTTPS origin'); } if (request.method !== 'GET' && request.method !== 'HEAD') { throw new Error('Discovered PDS enrichment permits only read requests'); } const deadline = new AbortController(); const timer = setTimeout(() => { deadline.abort(new Error('Discovered PDS request deadline exceeded')); }, timeoutMs); const signal = request.signal.aborted ? request.signal : AbortSignal.any([request.signal, deadline.signal]); try { const selected = await vettedPublicAddress(url, signal, resolveAddress); if (signal.aborted) throw new Error('Discovered PDS request was aborted before transport dispatch'); const pinnedLookup = (( _hostname: string, options: { all?: boolean }, callback: ( error: Error | null, address: string | Array<{ address: string; family: number }>, family?: number, ) => void, ): void => { if (options.all === true) { callback(null, [{ address: selected.address, family: selected.family }]); } else { callback(null, selected.address, selected.family); } }) as LookupFunction;
return await new Promise<Response>((resolveResponse, rejectResponse) => { const nodeRequest = dispatchRequest({ protocol: 'https:', hostname: url.hostname.replace(/^\[|\]$/gu, ''), port: url.port || 443, path: `${url.pathname}${url.search}`, method: request.method, servername: isIP(url.hostname.replace(/^\[|\]$/gu, '')) === 0 ? url.hostname.replace(/^\[|\]$/gu, '') : undefined, lookup: pinnedLookup, headers: { accept: 'application/json', 'accept-encoding': 'identity', 'user-agent': 'perch-reader/0.1', }, signal, }, (response) => { const status = response.statusCode ?? 0; if (status >= 300 && status < 400) { response.destroy(); rejectResponse(new Error(`Discovered PDS redirected with status ${status}; redirects are refused`)); return; } const declaredLength = Number(response.headers['content-length']); if (Number.isFinite(declaredLength) && declaredLength > MAX_DISCOVERED_RESPONSE_BYTES) { response.destroy(); rejectResponse(new Error('Discovered PDS response exceeded the 1,000,000-byte limit')); return; } const chunks: Buffer[] = []; let byteLength = 0; response.on('data', (chunk: Buffer | string) => { const bytes = typeof chunk === 'string' ? Buffer.from(chunk) : chunk; byteLength += bytes.length; if (byteLength > MAX_DISCOVERED_RESPONSE_BYTES) { response.destroy(new Error('Discovered PDS response exceeded the 1,000,000-byte limit')); return; } chunks.push(bytes); }); response.on('error', rejectResponse); response.on('end', () => { const headers = new Headers(); const contentType = response.headers['content-type']; if (typeof contentType === 'string') headers.set('content-type', contentType); headers.set('content-length', String(byteLength)); resolveResponse(new Response(Buffer.concat(chunks), { status, headers })); }); }); nodeRequest.on('error', rejectResponse); nodeRequest.end(); }); } finally { clearTimeout(timer); } }; return Object.assign(implementation, { preconnect(): void { // Deliberately no-op: preconnection would bypass the pinned-address policy. }, });}
function didResolutionUrl(did: string): URL { if (/^did:plc:[a-z2-7]{24}$/u.test(did)) { return new URL(`https://plc.directory/${did}`); } if (!did.startsWith('did:web:')) { throw new Error(`Public profile enrichment does not support DID method ${did.split(':', 3).slice(0, 2).join(':') || did}`); }
const segments = did.slice('did:web:'.length).split(':'); const encodedDomain = segments.shift(); if ( !encodedDomain || !/^[a-z0-9.-]+(?:%3a[0-9]+)?$/iu.test(encodedDomain) || encodedDomain.startsWith('.') || encodedDomain.endsWith('.') || encodedDomain.includes('..') ) { throw new Error('did:web identifier omitted a safe fully-qualified domain name'); } const domain = encodedDomain.replace(/%3a/iu, ':'); const origin = new URL(`https://${domain}`); if ( origin.protocol !== 'https:' || origin.username || origin.password || isIP(origin.hostname) !== 0 || !origin.hostname.includes('.') || origin.port === '0' ) { throw new Error('did:web identifier omitted a safe public HTTPS domain'); }
const path = segments.length === 0 ? '/.well-known/did.json' : `/${segments.map((segment) => { let decoded: string; try { decoded = decodeURIComponent(segment); } catch { throw new Error('did:web identifier contained invalid percent encoding'); } if (!decoded || decoded === '.' || decoded === '..' || /[\\/\u0000-\u001f\u007f]/u.test(decoded)) { throw new Error('did:web identifier contained an unsafe path segment'); } return encodeURIComponent(decoded); }).join('/')}/did.json`; return new URL(path, origin);}
async function readPublicJson(url: URL, context: string): Promise<unknown> { const response = await discoveredFetch(url.origin)(new Request(url, { method: 'GET', signal: requestSignal(), })); const body = await response.text(); let parsed: unknown; try { parsed = JSON.parse(body); } catch { throw new Error(`${context} returned malformed JSON (status ${response.status})`); } if (!response.ok) { const error = objectValue(parsed); const code = stringValue(error?.error); const rawMessage = stringValue(error?.message); const message = rawMessage?.replace(/[\r\n\t]+/gu, ' ').slice(0, 200); throw new Error(`${context} failed (status ${response.status}${code ? `, ${code}` : ''})${message ? `: ${message}` : ''}`); } return parsed;}
export function pdsEndpointFromDidDocument(value: unknown, expectedDid: string): string { const document = objectValue(value); if (stringValue(document?.id) !== expectedDid) { throw new Error('Resolved DID document id did not match the requested canonical DID'); } const rawServices = document?.service; const services = Array.isArray(rawServices) ? rawServices : []; const candidates = unique(services.flatMap((rawService) => { const service = objectValue(rawService); const id = stringValue(service?.id); const type = stringValue(service?.type); const endpoint = stringValue(service?.serviceEndpoint); if (!id || !type || !endpoint || !id.endsWith('#atproto_pds') || type !== 'AtprotoPersonalDataServer') return []; try { return [normalizeService(endpoint, 'Discovered PDS')]; } catch { return []; } })); if (candidates.length !== 1) throw new Error('Resolved DID document did not supply exactly one safe ATProto PDS origin'); const endpoint = new URL(candidates[0]); if (endpoint.protocol !== 'https:') throw new Error('Discovered PDS must use HTTPS'); return endpoint.origin;}
function unavailableEnrichments( actorDid: string, service: string, fetchedAt: string, status: 'unavailable' | 'unsupported', note: string,): ProfileEnrichment[] { return [ ['repository-collections', 'Known public repository collections'], ['tangled-profile', 'Tangled profile'], ['greengale-publication', 'GreenGale publication'], ].map(([id, title]) => ({ id, title, status, items: [], source: { service, actorDid, fetchedAt }, note, }));}
function previewItems( record: ObjectValue, scalarFields: readonly string[], arrayFields: readonly string[],): { items: Array<{ label: string; value: string; uri?: string }>; notes: string[] } { const items: Array<{ label: string; value: string; uri?: string }> = []; const notes: string[] = []; for (const field of scalarFields) { if (!Object.prototype.hasOwnProperty.call(record, field)) continue; const value = record[field]; if (typeof value === 'string' || typeof value === 'boolean' || typeof value === 'number') { const item: { label: string; value: string; uri?: string } = { label: field, value: String(value) }; if (field === 'url' && typeof value === 'string') { try { const url = new URL(value); if (url.protocol === 'https:' || url.protocol === 'http:') item.uri = url.toString(); } catch { notes.push(`${field} was present but was not a valid HTTP(S) URL.`); } } items.push(item); if (value === '') notes.push(`${field} was present-empty and retained verbatim.`); } else { notes.push(`${field} was present with an unsupported value shape.`); } } for (const field of arrayFields) { if (!Object.prototype.hasOwnProperty.call(record, field)) continue; const values = record[field]; if (!Array.isArray(values)) { notes.push(`${field} was present but was not an array.`); continue; } values.slice(0, 10).forEach((value, index) => { if (typeof value === 'string') { items.push({ label: `${field} ${index + 1}`, value }); if (value === '') notes.push(`${field} item ${index + 1} was present-empty and retained verbatim.`); } else { notes.push(`${field} item ${index + 1} had an unsupported value shape.`); } }); if (values.length > 10) notes.push(`${field} has ${values.length - 10} additional item(s) outside the bounded preview.`); } return { items, notes };}
async function recordEnrichment( agent: BskyAgent, service: string, actorDid: string, fetchedAt: string, collections: ReadonlySet<string>, options: { id: string; title: string; collection: string; scalarFields: readonly string[]; arrayFields: readonly string[]; },): Promise<ProfileEnrichment> { const baseSource = { service, actorDid, fetchedAt }; if (!collections.has(options.collection)) { return { id: options.id, title: options.title, status: 'absent', items: [], source: baseSource }; } try { const response = await agent.com.atproto.repo.getRecord( { repo: actorDid, collection: options.collection, rkey: 'self' }, { signal: requestSignal() }, ); const uri = stringValue(response.data.uri); const cid = stringValue(response.data.cid); const record = objectValue(response.data.value); if (!uri || !cid || !record) { return { id: options.id, title: options.title, status: 'unsupported', items: [], source: baseSource, note: 'The self record response omitted a URI, CID, or object value.', }; } const preview = previewItems(record, options.scalarFields, options.arrayFields); const previewedFields = new Set(['$type', ...options.scalarFields, ...options.arrayFields]); const unexpanded = Object.keys(record).filter((field) => !previewedFields.has(field)); if (unexpanded.length > 0) { preview.notes.push(`Additional record fields were not expanded by this bounded provider: ${unexpanded.join(', ')}.`); } return { id: options.id, title: options.title, status: 'available', items: preview.items, source: { ...baseSource, uri, cid }, ...(preview.notes.length > 0 ? { note: preview.notes.join(' ') } : {}), }; } catch (error) { const code = stringValue(objectValue(error)?.error); if (code === 'RecordNotFound') { return { id: options.id, title: options.title, status: 'absent', items: [], source: baseSource, note: 'The collection exists, but its self record is absent.', }; } return { id: options.id, title: options.title, status: 'unavailable', items: [], source: baseSource, note: sanitizedError(`Reading ${options.collection}/self`, error).message, }; }}
async function profileEnrichments( actorDid: string,): Promise<{ enrichments: ProfileEnrichment[]; warnings: string[] }> { const fetchedAt = new Date().toISOString(); let identityUrl: URL; try { identityUrl = didResolutionUrl(actorDid); } catch (error) { const note = error instanceof Error ? error.message : String(error); return { enrichments: unavailableEnrichments(actorDid, actorDid, fetchedAt, 'unsupported', note), warnings: [note], }; }
let didDocument: unknown; try { didDocument = await readPublicJson(identityUrl, `Resolving DID ${actorDid} for public enrichment`); } catch (error) { const note = error instanceof Error ? error.message : String(error); return { enrichments: unavailableEnrichments(actorDid, identityUrl.toString(), fetchedAt, 'unavailable', note), warnings: [note], }; }
let service: string; try { service = pdsEndpointFromDidDocument(didDocument, actorDid); await vettedPublicAddress(new URL(service)); } catch (error) { const note = error instanceof Error ? error.message : String(error); return { enrichments: unavailableEnrichments(actorDid, identityUrl.toString(), fetchedAt, 'unsupported', note), warnings: [`Public profile enrichment refused the discovered endpoint: ${note}`], }; }
const publicAgent = new BskyAgent({ service, fetch: discoveredFetch(service) }); let collections: string[]; try { const described = await publicAgent.com.atproto.repo.describeRepo( { repo: actorDid }, { signal: requestSignal() }, ); if (described.data.did !== actorDid) { throw new Error('describeRepo DID did not match the requested canonical DID'); } if (!Array.isArray(described.data.collections) || described.data.collections.some((entry) => typeof entry !== 'string')) { throw new Error('describeRepo omitted a usable collection inventory'); } collections = described.data.collections; } catch (error) { const note = sanitizedError(`Describing public repository ${actorDid}`, error).message; return { enrichments: unavailableEnrichments(actorDid, service, fetchedAt, 'unavailable', note), warnings: [note], }; }
const collectionSet = new Set(collections); const inventoryItems = collections.flatMap((collection) => { const title = PROFILE_COLLECTIONS[collection]; return title === undefined ? [] : [{ label: title, value: collection }]; }); const inventory: ProfileEnrichment = { id: 'repository-collections', title: 'Known public repository collections', status: inventoryItems.length > 0 ? 'available' : 'absent', items: inventoryItems, source: { service, actorDid, fetchedAt, uri: `at://${actorDid}` }, }; const [tangled, greengale] = await Promise.all([ recordEnrichment(publicAgent, service, actorDid, fetchedAt, collectionSet, { id: 'tangled-profile', title: 'Tangled profile', collection: 'sh.tangled.actor.profile', scalarFields: ['description', 'pronouns', 'location', 'bluesky'], arrayFields: ['links', 'stats', 'pinnedRepositories'], }), recordEnrichment(publicAgent, service, actorDid, fetchedAt, collectionSet, { id: 'greengale-publication', title: 'GreenGale publication', collection: 'app.greengale.publication', scalarFields: ['url', 'name', 'description', 'hideBlueskyBio', 'enableSiteStandard'], arrayFields: ['pinnedPosts'], }), ]); return { enrichments: [inventory, tangled, greengale], warnings: [] };}
function expandedPath(path: string): string { if (path === '~') return homedir(); if (path.startsWith('~/')) return `${homedir()}/${path.slice(2)}`; return path;}
async function authenticationConfig(path: string): Promise<{ identifier: string; password: string }> { let parsed: unknown; const resolved = expandedPath(path); try { parsed = JSON.parse(await readFile(resolved, 'utf8')); } catch (error) { const message = error instanceof Error ? error.message.replace(/[\r\n\t]+/g, ' ').slice(0, 200) : 'unknown error'; throw new Error(`Unable to read authentication config ${resolved}: ${message}`); } const env = objectValue(objectValue(objectValue(parsed)?.mcpServers)?.atproto)?.env; const identifier = stringValue(objectValue(env)?.BSKY_IDENTIFIER); const password = stringValue(objectValue(env)?.BSKY_APP_PASSWORD); if (!identifier || !password) { throw new Error('Authentication config is missing mcpServers.atproto.env BSKY_IDENTIFIER or BSKY_APP_PASSWORD'); } return { identifier, password };}
interface ClientContext { agent: BskyAgent; source: Source; authenticated: boolean; authenticationSecrets: () => readonly string[];}
async function createClientContext(options: ReaderOptions): Promise<ClientContext> { const authConfig = options.authConfig; const authenticated = authConfig !== undefined; const service = normalizeService( authenticated ? (options.pds ?? DEFAULT_PDS) : (options.service ?? PUBLIC_SERVICE), authenticated ? 'PDS' : 'service', ); const secrets: string[] = []; const agent = new BskyAgent({ service, fetch: timedFetch, persistSession(_event, session) { if (session?.accessJwt) secrets.push(session.accessJwt); if (session?.refreshJwt) secrets.push(session.refreshJwt); }, }); const authenticationSecrets = (): string[] => unique([ ...secrets, agent.session?.accessJwt, agent.session?.refreshJwt, ].filter((value): value is string => typeof value === 'string' && value.length > 0));
if (authConfig === undefined) { return { agent, authenticated, authenticationSecrets, source: Object.freeze({ service, mode: 'public' as const }), }; }
const credentials = await authenticationConfig(authConfig); secrets.push(credentials.identifier, credentials.password); try { await agent.login({ identifier: credentials.identifier, password: credentials.password }); } catch (error) { throw sanitizedError( 'Authentication', error, authenticationSecrets(), 'credentials were rejected or the authentication service was unavailable', ); } const session = agent.session; if (!session?.did) throw new Error('Authentication response omitted the viewer DID'); if (!session.accessJwt || !session.refreshJwt) { throw new Error('Authentication response omitted usable session credentials'); } return { agent, authenticated, authenticationSecrets, source: Object.freeze({ service, mode: 'authenticated' as const, viewerDid: session.did }), };}
async function fetchPrimaryProfile( agent: BskyAgent, actor: string,): Promise<{ value: unknown; author: Author }> { const response = await agent.getProfile({ actor }, { signal: requestSignal() }); const author = authorFrom(response.data); if (!author) throw new Error('Profile response omitted canonical DID or handle'); if (actor.startsWith('did:') && author.did !== actor) { throw new Error(`Profile changed canonical DID from ${actor} to ${author.did}; refusing target substitution`); } return { value: response.data, author };}
/** Creates a public reader by default, or an in-memory authenticated reader when authConfig is supplied. */export async function createReader(options: ReaderOptions = {}): Promise<Reader> { const context = await createClientContext(options); const { agent, source, authenticated, authenticationSecrets } = context; return { source, async thread(input: string, threadOptions: ThreadOptions = {}): Promise<ThreadSnapshot> { const resolvedOptions = boundedOptions(threadOptions); const focusUri = await canonicalPostUri(agent, input, authenticationSecrets); try { const response = await agent.getPostThread( { uri: focusUri, depth: resolvedOptions.depth, parentHeight: resolvedOptions.parentHeight, }, { signal: requestSignal() }, ); const fetchedAt = new Date().toISOString(); return normalizeThread(response.data.thread, { requested: input, focusUri, fetchedAt, source, options: resolvedOptions, }); } catch (error) { throw sanitizedError(`Fetching thread ${focusUri}`, error, authenticationSecrets()); } }, async profile(actorInput: string): Promise<ProfileSnapshot> { const actor = profileInputActor(actorInput); try { const primary = await fetchPrimaryProfile(agent, actor); const fetchedAt = new Date().toISOString(); const enrichment = await profileEnrichments(primary.author.did); return normalizeProfile(primary.value, { requested: actorInput, fetchedAt, source, enrichments: enrichment.enrichments, warnings: enrichment.warnings, }); } catch (error) { throw sanitizedError(`Fetching profile ${actor}`, error, authenticationSecrets()); } }, async following(feedOptions: FeedOptions = {}): Promise<FeedSnapshot> { if (!authenticated || source.mode !== 'authenticated') { throw new Error('Following feed requires an authenticated reader'); } const session = agent.session; if (!session?.accessJwt || !session.refreshJwt || session.did !== source.viewerDid) { throw new Error('Following feed requires a usable authenticated session'); } const resolvedOptions = boundedFeedOptions(feedOptions); try { const response = await agent.getTimeline( { limit: resolvedOptions.limit, ...(resolvedOptions.cursor !== undefined ? { cursor: resolvedOptions.cursor } : {}), }, { signal: requestSignal() }, ); const snapshot = normalizeFeed(response.data, { fetchedAt: new Date().toISOString(), source, limit: resolvedOptions.limit, ...(resolvedOptions.cursor !== undefined ? { cursor: resolvedOptions.cursor } : {}), }); if (resolvedOptions.cursor !== undefined && snapshot.cursor === resolvedOptions.cursor) { throw new Error('Timeline returned the unchanged continuation cursor; refusing a silent loop'); } return snapshot; } catch (error) { throw sanitizedError('Fetching following feed', error, authenticationSecrets()); } }, async authorFeed(actorInput: string, feedOptions: FeedOptions = {}): Promise<FeedSnapshot> { const actor = profileInputActor(actorInput); const resolvedOptions = boundedFeedOptions(feedOptions); try { const primary = await fetchPrimaryProfile(agent, actor); const response = await agent.getAuthorFeed( { actor: primary.author.did, limit: resolvedOptions.limit, ...(resolvedOptions.cursor !== undefined ? { cursor: resolvedOptions.cursor } : {}), }, { signal: requestSignal() }, ); const snapshot = normalizeFeed(response.data, { feed: 'author', actor: primary.author, fetchedAt: new Date().toISOString(), source, limit: resolvedOptions.limit, ...(resolvedOptions.cursor !== undefined ? { cursor: resolvedOptions.cursor } : {}), }); if (resolvedOptions.cursor !== undefined && snapshot.cursor === resolvedOptions.cursor) { throw new Error('Author feed returned the unchanged continuation cursor; refusing a silent loop'); } return snapshot; } catch (error) { throw sanitizedError(`Fetching author feed ${actor}`, error, authenticationSecrets()); } }, };}
function strongReference(value: unknown, location: string, preserve: boolean = false): StrongReference { const reference = objectValue(value); const uri = canonicalPostRecordUri(reference?.uri); const cid = stringValue(reference?.cid); if (!uri || !cid) throw new Error(`${location} omitted a canonical post URI or CID`); if (!preserve) return { uri, cid }; const type = stringValue(reference?.$type); if (reference?.$type !== undefined && type !== 'com.atproto.repo.strongRef') { throw new Error(`${location} supplied an invalid strong-reference type`); } return { ...structuredClone(reference), uri, cid };}
function stableJson(value: unknown): string { const normalize = (entry: unknown): unknown => { if (Array.isArray(entry)) return entry.map(normalize); const object = objectValue(entry); if (!object) return entry; return Object.fromEntries(Object.keys(object).sort().map((key) => [key, normalize(object[key])])); }; return JSON.stringify(normalize(value));}
async function resolvedPublicationTarget( agent: BskyAgent, target: PublicationTargetInput, authenticationSecrets: () => readonly string[],): Promise<{ reference: StrongReference; record: ObjectValue }> { const sourceInput = target.stored?.uri ?? target.input; const uri = await canonicalPostUri(agent, sourceInput, authenticationSecrets); try { const response = await agent.getPosts({ uris: [uri] }, { signal: requestSignal() }); const post = response.data.posts.find((candidate) => candidate.uri === uri); if (!post) throw new Error('the target was missing, blocked, or unavailable'); const reference = strongReference(post, `Post target ${target.input}`); const record = objectValue(post.record); if (!record) throw new Error(`Post target ${target.input} omitted its record value`); if (target.stored && (target.stored.uri !== reference.uri || target.stored.cid !== reference.cid)) { throw new Error(`Stored post reference ${target.stored.ref} is stale; its URI or CID no longer matches the live record`); } return { reference, record }; } catch (error) { throw sanitizedError(`Resolving post target ${target.input}`, error, authenticationSecrets()); }}
async function preparePublication( context: ClientContext & { source: Source & { mode: 'authenticated'; viewerDid: string } }, draft: LoadedDraft,): Promise<PreparedPublication> { if (draft.posts.length === 0) throw new Error('Publication requires at least one explicit post'); const richTexts: RichText[] = []; for (let index = 0; index < draft.posts.length; index += 1) { const draftPost = draft.posts[index]; const richText = new RichText({ text: draftPost.text }); if (richText.graphemeLength > 300 || richText.length > 3_000) { throw new Error( `Post ${index + 1} exceeds Bluesky limits: ${richText.graphemeLength}/300 graphemes and ${richText.length}/3000 UTF-8 bytes.`, ); } await richText.detectFacets(context.agent); const preliminary: Record<string, unknown> = { $type: 'app.bsky.feed.post', text: richText.text, createdAt: new Date().toISOString(), ...(richText.facets !== undefined && richText.facets.length > 0 ? { facets: richText.facets } : {}), ...(draftPost.langs !== undefined ? { langs: draftPost.langs } : {}), }; const postValidation = AppBskyFeedPost.validateRecord(preliminary); if (!postValidation.success) { throw new Error(`Post ${index + 1} failed installed protocol validation: ${postValidation.error.message}`); } if (draftPost.external !== undefined) { const externalValidation = AppBskyEmbedExternal.validateMain({ $type: 'app.bsky.embed.external', external: { uri: draftPost.external.uri, title: draftPost.external.title, description: draftPost.external.description, }, }); if (!externalValidation.success) { throw new Error(`Post ${index + 1} external card failed installed protocol validation: ${externalValidation.error.message}`); } } richTexts.push(richText); }
const quoteTargets = await Promise.all(draft.posts.map(async (post) => post.quote === undefined ? undefined : resolvedPublicationTarget(context.agent, post.quote, context.authenticationSecrets) )); quoteTargets.forEach((target, index) => { if (target === undefined) return; const validation = AppBskyEmbedRecord.validateMain({ $type: 'app.bsky.embed.record', record: target.reference, }); if (!validation.success) { throw new Error(`Post ${index + 1} quote failed installed protocol validation: ${validation.error.message}`); } }); const initialParent = draft.replyTo === undefined ? undefined : await resolvedPublicationTarget(context.agent, draft.replyTo, context.authenticationSecrets); let initialReply: { parent: StrongReference; root: StrongReference } | undefined; if (initialParent) { const reply = objectValue(initialParent.record.reply); const root = reply === undefined ? initialParent.reference : strongReference(reply.root, `Reply root for ${initialParent.reference.uri}`, true); initialReply = { parent: initialParent.reference, root }; }
const publicationUuid = randomUUID().replaceAll('-', ''); const id = `pub-${publicationUuid}`; const preparedAt = new Date().toISOString(); let previousRkey: string | undefined; const posts: PreparedPublicationPost[] = draft.posts.map((post, index) => { const rkey = TID.nextStr(previousRkey); previousRkey = rkey; const facets = richTexts[index].facets; return { index: index + 1, rkey, plannedUri: AtUri.make(context.source.viewerDid, 'app.bsky.feed.post', rkey).toString(), text: post.text, ...(facets !== undefined && facets.length > 0 ? { facets: lexToJson(facets as Parameters<typeof lexToJson>[0]) as unknown[] } : {}), ...(quoteTargets[index] !== undefined ? { quote: quoteTargets[index]!.reference } : {}), ...(post.images !== undefined ? { images: post.images.map((image) => ({ ...image })) } : {}), ...(post.external !== undefined ? { external: { uri: post.external.uri, title: post.external.title, description: post.external.description, ...(post.external.thumb !== undefined ? { thumb: { ...post.external.thumb } } : {}), }, } : {}), ...(post.langs !== undefined ? { langs: [...post.langs] } : {}), createdAt: preparedAt, }; }); return { plan: { version: 1, id, ownerDid: context.source.viewerDid, service: context.source.service, preparedAt, sourceFile: draft.sourceFile, ...(initialReply !== undefined ? { initialReply } : {}), posts, attachments: draft.attachments.map((attachment) => ({ ...attachment.metadata })), }, attachments: draft.attachments.map((attachment) => ({ metadata: { ...attachment.metadata }, bytes: new Uint8Array(attachment.bytes), })), };}
function publicationReply( claim: PublicationClaim,): { root: StrongReference; parent: StrongReference } | undefined { const index = claim.post.index; if (index === 1) return claim.plan.initialReply; if (claim.confirmed.length !== index - 1) { throw new Error(`Post ${index} requires ${index - 1} confirmed predecessor(s), but ${claim.confirmed.length} were supplied`); } const parent = claim.confirmed[index - 2]; const root = claim.plan.initialReply?.root ?? claim.confirmed[0]; if (!parent || !root) throw new Error(`Post ${index} is missing its confirmed chain references`); return { root, parent };}
async function uploadPublicationAttachments( context: ClientContext, attachments: FrozenAttachment[],): Promise<Map<string, unknown>> { const blobs = new Map<string, unknown>(); for (const attachment of attachments) { if (attachment.bytes.length !== attachment.metadata.byteLength) { throw new Error(`Frozen attachment ${attachment.metadata.id} byte length no longer matches its plan`); } const digest = createHash('sha256').update(attachment.bytes).digest('hex'); if (digest !== attachment.metadata.sha256) { throw new Error(`Frozen attachment ${attachment.metadata.id} failed its durable SHA-256 check`); } try { const response = await context.agent.uploadBlob(attachment.bytes, { encoding: attachment.metadata.mimeType, signal: requestSignal(), }); if (!response.data.blob) throw new Error('upload response omitted its blob reference'); blobs.set(attachment.metadata.id, response.data.blob); } catch (error) { throw sanitizedError(`Uploading frozen attachment ${attachment.metadata.id}`, error, context.authenticationSecrets()); } } return blobs;}
async function publishClaimedPost( context: ClientContext & { source: Source & { mode: 'authenticated'; viewerDid: string } }, claim: PublicationClaim, beforeDispatch: (record: ObjectValue) => void,): Promise<StrongReference> { if (claim.plan.ownerDid !== context.source.viewerDid || claim.plan.service !== context.source.service) { throw new Error('Publication plan belongs to a different authenticated DID or service'); } const expectedUri = AtUri.make(context.source.viewerDid, 'app.bsky.feed.post', claim.post.rkey).toString(); if (!TID_RECORD_KEY_PATTERN.test(claim.post.rkey) || claim.post.plannedUri !== expectedUri) { throw new Error(`Prepared post ${claim.post.index} does not have its immutable app.bsky.feed.post TID identity`); } const expectedAttachments = new Set([ ...(claim.post.images?.map((image) => image.attachmentId) ?? []), ...(claim.post.external?.thumb ? [claim.post.external.thumb.attachmentId] : []), ]); if ( claim.attachments.length !== expectedAttachments.size || claim.attachments.some((attachment) => !expectedAttachments.has(attachment.metadata.id)) ) throw new Error(`Post ${claim.post.index} did not receive its exact frozen attachment set`);
let candidate: unknown; if (claim.record !== undefined) { candidate = jsonToLex(claim.record as Parameters<typeof jsonToLex>[0]); } else { const blobs = await uploadPublicationAttachments(context, claim.attachments); let media: ObjectValue | undefined; if (claim.post.images !== undefined && claim.post.images.length > 0) { media = { $type: 'app.bsky.embed.images', images: claim.post.images.map((image) => ({ image: blobs.get(image.attachmentId), alt: image.alt, })), }; } else if (claim.post.external !== undefined) { media = { $type: 'app.bsky.embed.external', external: { uri: claim.post.external.uri, title: claim.post.external.title, description: claim.post.external.description, ...(claim.post.external.thumb !== undefined ? { thumb: blobs.get(claim.post.external.thumb.attachmentId) } : {}), }, }; } let embed: ObjectValue | undefined; if (claim.post.quote && media) { embed = { $type: 'app.bsky.embed.recordWithMedia', record: { $type: 'app.bsky.embed.record', record: claim.post.quote }, media, }; } else if (claim.post.quote) { embed = { $type: 'app.bsky.embed.record', record: claim.post.quote }; } else { embed = media; } const reply = publicationReply(claim); candidate = { $type: 'app.bsky.feed.post', text: claim.post.text, createdAt: claim.post.createdAt, ...(claim.post.facets !== undefined ? { facets: jsonToLex(claim.post.facets as Parameters<typeof jsonToLex>[0]) } : {}), ...(claim.post.langs !== undefined ? { langs: claim.post.langs } : {}), ...(reply !== undefined ? { reply } : {}), ...(embed !== undefined ? { embed } : {}), }; } const validation = AppBskyFeedPost.validateRecord(candidate); if (!validation.success) { throw new Error(`Prepared post ${claim.post.index} failed final protocol validation: ${validation.error.message}`); } const record = claim.record === undefined ? lexToJson(validation.value as Parameters<typeof lexToJson>[0]) as ObjectValue : structuredClone(claim.record); beforeDispatch(record); try { const response = await context.agent.com.atproto.repo.createRecord( { repo: context.source.viewerDid, collection: 'app.bsky.feed.post', rkey: claim.post.rkey, validate: true, record: validation.value, }, { signal: requestSignal() }, ); if (response.data.uri !== claim.post.plannedUri || !response.data.cid) { throw new Error('createRecord response omitted or changed the planned URI/CID'); } return { uri: response.data.uri, cid: response.data.cid }; } catch (error) { throw sanitizedError(`Publishing post ${claim.post.index} as ${claim.post.plannedUri}`, error, context.authenticationSecrets()); }}
async function inspectPublicationRecord( context: ClientContext & { source: Source & { mode: 'authenticated'; viewerDid: string } }, rkey: string, plannedUri: string, record: ObjectValue,): Promise<PublicationRecordInspection> { const observedAt = new Date().toISOString(); try { const response = await context.agent.com.atproto.repo.getRecord( { repo: context.source.viewerDid, collection: 'app.bsky.feed.post', rkey }, { signal: requestSignal() }, ); const actual = lexToJson(response.data.value as Parameters<typeof lexToJson>[0]); if (response.data.uri !== plannedUri || !response.data.cid) { return { state: 'mismatch', observedAt, uri: response.data.uri, ...(response.data.cid !== undefined ? { cid: response.data.cid } : {}), detail: 'Exact record key returned an unexpected URI or omitted its CID.', }; } if (stableJson(actual) !== stableJson(record)) { return { state: 'mismatch', observedAt, uri: response.data.uri, cid: response.data.cid, detail: 'Exact record key exists, but its value does not match the durable prepared record.', }; } return { state: 'matching', observedAt, uri: response.data.uri, cid: response.data.cid, }; } catch (error) { if (stringValue(objectValue(error)?.error) === 'RecordNotFound') return { state: 'absent', observedAt }; throw sanitizedError(`Reconciling exact post record ${plannedUri}`, error, context.authenticationSecrets()); }}
async function canonicalMutationTarget( context: ClientContext & { source: Source & { mode: 'authenticated'; viewerDid: string } }, input: string,): Promise<{ did: string; handle?: string }> { const trimmed = input.trim(); if (trimmed.startsWith('at://')) { let parsed: AtUri; try { parsed = new AtUri(trimmed); } catch { throw new Error('Social target contains an invalid AT identifier'); } if ( (parsed.collection !== '' || parsed.rkey !== '') && (parsed.collection !== 'app.bsky.actor.profile' || parsed.rkey !== 'self') ) { throw new Error('Social target AT identifier must be a bare actor or app.bsky.actor.profile/self record'); } } const actor = profileInputActor(trimmed); try { const response = await context.agent.getProfile({ actor }, { signal: requestSignal() }); const did = stringValue(response.data.did); const handle = stringValue(response.data.handle); if (!did?.startsWith('did:')) throw new Error('target profile omitted its canonical DID'); if (did === context.source.viewerDid) throw new Error('authenticated account cannot target itself for this operation'); return { did, ...(handle !== undefined ? { handle } : {}) }; } catch (error) { throw sanitizedError(`Resolving social target ${input}`, error, context.authenticationSecrets()); }}
async function relationshipFor( context: ClientContext & { source: Source & { mode: 'authenticated'; viewerDid: string } }, targetDid: string,): Promise<ObjectValue> { try { const response = await context.agent.app.bsky.graph.getRelationships( { actor: context.source.viewerDid, others: [targetDid] }, { signal: requestSignal() }, ); if (response.data.actor !== undefined && response.data.actor !== context.source.viewerDid) { throw new Error('relationship response actor did not match the authenticated DID'); } const relationship = response.data.relationships.find((candidate) => objectValue(candidate)?.did === targetDid); const value = objectValue(relationship); if (!value || value.notFound === true) throw new Error('relationship response did not retain the canonical target'); return value; } catch (error) { throw sanitizedError(`Reading relationship to ${targetDid}`, error, context.authenticationSecrets()); }}
async function verifiedOwnedGraphRecord( context: ClientContext & { source: Source & { mode: 'authenticated'; viewerDid: string } }, uri: string, collection: 'app.bsky.graph.follow' | 'app.bsky.graph.block', targetDid: string,): Promise<{ uri: string; cid: string; rkey: string }> { let parsed: AtUri; try { parsed = new AtUri(uri); } catch { throw new Error(`Relationship supplied malformed record URI ${uri}`); } if ( parsed.hostname !== context.source.viewerDid || parsed.collection !== collection || !parsed.rkey || !isValidRecordKey(parsed.rkey) || parsed.pathname !== `/${collection}/${parsed.rkey}` ) throw new Error(`Relationship record ${uri} is not an owned ${collection} record`); try { const response = await context.agent.com.atproto.repo.getRecord( { repo: context.source.viewerDid, collection, rkey: parsed.rkey }, { signal: requestSignal() }, ); if (response.data.uri !== uri || !response.data.cid || objectValue(response.data.value)?.subject !== targetDid) { throw new Error(`Owned relationship record ${uri} did not match target ${targetDid}`); } return { uri, cid: response.data.cid, rkey: parsed.rkey }; } catch (error) { throw sanitizedError(`Verifying owned relationship record ${uri}`, error, context.authenticationSecrets()); }}
function graphRecordKey( operation: 'follow' | 'block', viewerDid: string, targetDid: string,): string { const digest = createHash('sha256').update(`${operation}\u0000${viewerDid}\u0000${targetDid}`).digest(); const stableTimestamp = 1_500_000_000_000_000 + digest.readUIntBE(0, 6); return TID.fromTime(stableTimestamp, digest[6] % 32).toString();}
async function createGraphRecord( context: ClientContext & { source: Source & { mode: 'authenticated'; viewerDid: string } }, operation: 'follow' | 'block', target: { did: string; handle?: string },): Promise<SocialMutationReceipt> { const collection = operation === 'follow' ? 'app.bsky.graph.follow' : 'app.bsky.graph.block'; const rkey = graphRecordKey(operation, context.source.viewerDid, target.did); const plannedUri = AtUri.make(context.source.viewerDid, collection, rkey).toString(); try { const response = await context.agent.com.atproto.repo.createRecord( { repo: context.source.viewerDid, collection, rkey, validate: true, record: { $type: collection, subject: target.did, createdAt: new Date().toISOString() }, }, { signal: requestSignal() }, ); if (response.data.uri !== plannedUri || !response.data.cid) { throw new Error(`createRecord response did not confirm planned relationship record ${plannedUri}`); } return { operation, targetDid: target.did, ...(target.handle !== undefined ? { targetHandle: target.handle } : {}), state: 'created', recordUri: response.data.uri, }; } catch (error) { try { await verifiedOwnedGraphRecord(context, plannedUri, collection, target.did); return { operation, targetDid: target.did, ...(target.handle !== undefined ? { targetHandle: target.handle } : {}), state: 'already-present', recordUri: plannedUri, }; } catch { throw sanitizedError(`Creating ${operation} record for ${target.did}`, error, context.authenticationSecrets()); } }}
/** Creates the authenticated mutation client without exposing its SDK agent or session tokens. */export async function createSocialClient(options: ReaderOptions): Promise<SocialClient> { if (options.authConfig === undefined) throw new Error('Social mutations require explicit authentication'); const baseContext = await createClientContext(options); if (baseContext.source.mode !== 'authenticated' || !baseContext.source.viewerDid) { throw new Error('Social mutations require a usable authenticated session'); } const context = baseContext as ClientContext & { source: Source & { mode: 'authenticated'; viewerDid: string }; }; return { source: context.source, async preparePublication(draft: LoadedDraft): Promise<PreparedPublication> { try { return await preparePublication(context, draft); } catch (error) { throw sanitizedError('Preparing publication', error, context.authenticationSecrets()); } }, publishPost(claim: PublicationClaim, beforeDispatch: (record: ObjectValue) => void): Promise<StrongReference> { return publishClaimedPost(context, claim, beforeDispatch); }, inspectPublicationRecord( rkey: string, plannedUri: string, record: ObjectValue, ): Promise<PublicationRecordInspection> { return inspectPublicationRecord(context, rkey, plannedUri, record); }, async follow(input: string): Promise<SocialMutationReceipt> { const target = await canonicalMutationTarget(context, input); const relationship = await relationshipFor(context, target.did); const existing = stringValue(relationship.following); if (existing) { await verifiedOwnedGraphRecord(context, existing, 'app.bsky.graph.follow', target.did); return { operation: 'follow', targetDid: target.did, ...(target.handle !== undefined ? { targetHandle: target.handle } : {}), state: 'already-present', recordUri: existing, }; } return createGraphRecord(context, 'follow', target); }, async unfollow(input: string): Promise<SocialMutationReceipt> { const target = await canonicalMutationTarget(context, input); const relationship = await relationshipFor(context, target.did); const existing = stringValue(relationship.following); if (!existing) { return { operation: 'unfollow', targetDid: target.did, ...(target.handle !== undefined ? { targetHandle: target.handle } : {}), state: 'absent', }; } const record = await verifiedOwnedGraphRecord( context, existing, 'app.bsky.graph.follow', target.did, ); try { await context.agent.com.atproto.repo.deleteRecord( { repo: context.source.viewerDid, collection: 'app.bsky.graph.follow', rkey: record.rkey, swapRecord: record.cid, }, { signal: requestSignal() }, ); } catch (error) { throw sanitizedError(`Deleting verified follow record ${existing}`, error, context.authenticationSecrets()); } return { operation: 'unfollow', targetDid: target.did, ...(target.handle !== undefined ? { targetHandle: target.handle } : {}), state: 'deleted', recordUri: existing, }; }, async block(input: string): Promise<SocialMutationReceipt> { const target = await canonicalMutationTarget(context, input); const relationship = await relationshipFor(context, target.did); const existing = stringValue(relationship.blocking); if (existing) { await verifiedOwnedGraphRecord(context, existing, 'app.bsky.graph.block', target.did); return { operation: 'block', targetDid: target.did, ...(target.handle !== undefined ? { targetHandle: target.handle } : {}), state: 'already-present', recordUri: existing, }; } return createGraphRecord(context, 'block', target); }, };}