Something went wrong. Try again.
A local-first sync engine for atproto
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251import { AtUri, type AtUriString, type DidString, type HandleString, type NsidString, type RecordKeyString } from '@atproto/syntax';import { IDBStorage } from './storage/indexeddb';import type { ATRecord, CollectionURI, RecordPath, RecordURI } from './types';
export interface ATCache { getDid(handle: HandleString): Promise<DidString | null>; getHandle(did: DidString): Promise<HandleString | null>; getRecord<N extends NsidString = NsidString>(uri: RecordURI<N> | AtUri): Promise<ATRecord<N> | null>; getCollection<N extends NsidString = NsidString>(uri: CollectionURI<N> | AtUri): Promise<ATRecord<N>[]>; getBacklinks<N extends NsidString = NsidString>(uri: RecordURI<N> | AtUri, source: NsidString, path: RecordPath): Promise<ATRecord<N>[]>; getBacklinksCount(uri: RecordURI | AtUri, source: NsidString, path: RecordPath): Promise<number>;}
interface BacklinkResponse { total: number; records: Backlink[]; cursor: string;}
interface Backlink { did: DidString; collection: NsidString; rkey: RecordKeyString;}
interface Identity { did: DidString; handle: HandleString; pds: string; signing_key: string;}
export class MicrocosmCache implements ATCache { static slingshotURL = 'https://slingshot.microcosm.blue'; static constellationURL = 'https://constellation.microcosm.blue';
#identityQueue = new Map<DidString | HandleString, Promise<Identity>>(); #cache = new IDBStorage(); // #backlinkCache = new Map<AtUriString, Backlink[]>();
async getDid(handle: HandleString): Promise<DidString | null> { let identity = await this.#handleCache.get(handle); if (identity === undefined) { identity = await this.#fetchIdentity(handle); } return identity?.did || null; }
async getHandle(did: DidString): Promise<HandleString | null> { let identity = await this.#didCache.get(did); if (identity === undefined) { identity = await this.#fetchIdentity(did); } return identity?.handle || null; }
async getPDS(did: DidString): Promise<string | null> { let identity = await this.#didCache.get(did); if (identity === undefined) { identity = await this.#fetchIdentity(did); } return identity?.pds || null; }
async getSigningKey(did: DidString): Promise<string | null> { let identity = await this.#didCache.get(did); if (identity === undefined) { identity = await this.#fetchIdentity(did); } return identity?.signing_key || null; }
async #fetchIdentity(identifier: DidString | HandleString) { const identityPromise = this.#identityQueue.get(identifier);
if (identityPromise !== undefined) { return identityPromise; }
const url = new URL(`${MicrocosmCache.slingshotURL}/xrpc/blue.microcosm.identity.resolveMiniDoc`); url.searchParams.set('identifier', identifier);
const p = await fetch(url, { headers: { Accept: 'application/json', }, }) .then((response) => { if (!response.ok) { console.warn(`Failed to resolve identifier:`, identifier); return null; }
return response.json(); }) .catch(() => null);
this.#identityQueue.set(identifier, p);
const identity = await p; console.log(identity);
if (!identity) return undefined;
this.#didCache.set(identity.did, identity); this.#handleCache.set(identity.handle, identity); return identity as Identity; }
// https://slingshot.microcosm.blue/xrpc/com.atproto.repo.getRecord?repo=did%3Aplc%3Ahdhoaan3xa3jiuq4fg4mefid&collection=app.bsky.feed.like&rkey=3lv4ouczo2b2a async getRecord<N extends NsidString = NsidString>(uri: RecordURI<N> | AtUri): Promise<ATRecord<N> | null> { try { if (typeof uri === 'string') { uri = new AtUri(uri); } const uriString = uri.toString(); const record = await this.#recordCache.get(uriString); if (record !== undefined) return record as ATRecord<N>;
const url = new URL(`${MicrocosmCache.slingshotURL}/xrpc/com.atproto.repo.getRecord`); url.searchParams.set('repo', uri.did); url.searchParams.set('collection', uri.collection); url.searchParams.set('rkey', uri.rkey);
const recordResponse = await fetch(url, { headers: { Accept: 'application/json', }, });
if (!recordResponse.ok) { return null; }
const recordData = await recordResponse.json(); this.#recordCache.set(uriString, recordData);
return recordData; } catch (error) { console.error(uri, error); return null; } }
// https://puffball.us-east.host.bsky.network/xrpc/com.atproto.repo.listRecords?repo=did%3Aplc%3Azcanytzlaumjwgaopolw6wes&collection=app.bsky.feed.repost&limit=100&reverse=false async getCollection<N extends NsidString = NsidString>(uri: CollectionURI<N> | AtUri): Promise<ATRecord<N>[]> { if (typeof uri === 'string') { uri = new AtUri(uri); }
const pds = await this.getPDS(uri.did);
let cursor = ''; const records: ATRecord<N>[] = [];
do { const url = new URL(`${pds}/xrpc/com.atproto.repo.listRecords`); url.searchParams.set('repo', uri.did); url.searchParams.set('collection', uri.collection); url.searchParams.set('limit', '100'); url.searchParams.set('reversed', 'false'); if (cursor) { url.searchParams.set('cursor', cursor); } const response = await fetch(url);
if (!response.ok) break;
const data = await response.json();
if (data.records.length < 100) { cursor = ''; } else { cursor = data.cursor; }
for (const record of data.records as ATRecord[]) { this.#recordCache.set(record.uri, record); records.push(record); } } while (cursor);
return records; }
// https://constellation.microcosm.blue/xrpc/blue.microcosm.links.getBacklinks?subject=at%3A%2F%2Fdid%3Aplc%3Azcanytzlaumjwgaopolw6wes%2Fnetwork.cosmik.collection%2F3mfrzrpx6fw26&source=network.cosmik.collectionLink%3Acollection.uri&limit=100 async getBacklinks<N extends NsidString = NsidString>( uri: RecordURI<N> | AtUri, source: NsidString, path: RecordPath, ): Promise<ATRecord[]> { try { if (typeof uri === 'string') { uri = new AtUri(uri); } const url = new URL(`${MicrocosmCache.constellationURL}/xrpc/blue.microcosm.links.getBacklinks`); url.searchParams.set('subject', uri.toString()); url.searchParams.set('source', `${source}:${path}`); url.searchParams.set('limit', '100');
const recordResponse = await fetch(url, { headers: { Accept: 'application/json', }, });
if (!recordResponse.ok) { return []; }
const response: BacklinkResponse = await recordResponse.json(); const records = await Promise.all(response.records.map((bl) => this.getRecord(AtUri.make(bl.did, bl.collection, bl.rkey)))); return records.filter((record) => record !== null); } catch (error) { console.error(uri, error); return []; } } // https://constellation.microcosm.blue/xrpc/blue.microcosm.links.getBacklinksCount?subject=at%3A%2F%2Fdid%3Aplc%3Azcanytzlaumjwgaopolw6wes%2Fnetwork.cosmik.collection%2F3mfrzrpx6fw26&source=network.cosmik.collectionLink%3Acollection.uri async getBacklinksCount<N extends NsidString = NsidString>( uri: RecordURI<N> | AtUri, source: NsidString, path: RecordPath, ): Promise<number> { try { if (typeof uri === 'string') { uri = new AtUri(uri); } const url = new URL(`${MicrocosmCache.constellationURL}/xrpc/blue.microcosm.links.getBacklinksCount`); url.searchParams.set('subject', uri.toString()); url.searchParams.set('source', `${source}:${path}`);
const recordResponse = await fetch(url, { headers: { Accept: 'application/json', }, });
if (!recordResponse.ok) { return 0; }
const { total } = await recordResponse.json(); return total; } catch (error) { console.error(uri, error); return 0; } }}