diff --git a/.env.example b/.env.example index 7185b8e..78255dd 100644 --- a/.env.example +++ b/.env.example @@ -3,3 +3,5 @@ DB=./data.db JETSTREAM="wss://jetstream1.us-east.bsky.network/subscribe" HOSTNAME="feeds.example.com" PUBLISHER="did:plc:fasdv..." + +INGEST_BASE=http://localhost:4001 diff --git a/feedgen.ts b/feedgen.ts index 17b97ca..0397e22 100644 --- a/feedgen.ts +++ b/feedgen.ts @@ -114,6 +114,36 @@ const requireAuth = async ( const getAuthor = db.prepare("SELECT pds FROM authors WHERE did = ?"); +const ingestEndpoint = Deno.env.get("INGEST_BASE") ?? "http://localhost:4001"; + +class CachedDIDs extends Set { + private timeouts: Map> = new Map(); + + override add(key: V) { + const out = super.add(key); + const handle = setTimeout( + () => { + this.delete(key); + }, + 15 * 60 * 1000, + ); + this.timeouts.set(key, handle); + return out; + } + + override delete(key: V) { + const out = super.delete(key); + const handle = this.timeouts.get(key); + if (handle) { + clearTimeout(handle); + this.timeouts.delete(key); + } + return out; + } +} + +const cachedDIDs = new CachedDIDs(); + app.add(AppBskyFeedGetFeedSkeleton.mainSchema, { async handler({ request, params: { feed, limit, cursor } }) { const feedUri = parseResourceUri(feed); @@ -140,23 +170,20 @@ app.add(AppBskyFeedGetFeedSkeleton.mainSchema, { const jwt = await requireAuth(request, "app.bsky.feed.getFeedSkeleton"); const author = getAuthor.get(jwt.issuer); - if (author) { + if (cachedDIDs.has(jwt.issuer) && author) { pds = author.pds; } else { - const resolved = await didResolver.resolve(jwt.issuer as DID); - for (const service of resolved.service ?? []) { - if ( - service.type == "AtprotoPersonalDataServer" && - typeof service.serviceEndpoint === "string" - ) { - pds = service.serviceEndpoint; - } - } - if (typeof pds !== "string") + const refreshReq = await fetch( + `${ingestEndpoint}/refresh?id=${jwt.issuer}`, + ); + if (!refreshReq.ok) { throw new InvalidRequestError({ error: "NoServiceEndpoint", description: "No service endpoint", }); + } + pds = await refreshReq.text(); + cachedDIDs.add(jwt.issuer); } } diff --git a/ingest.ts b/ingest.ts index c402634..ec67fdc 100644 --- a/ingest.ts +++ b/ingest.ts @@ -95,12 +95,13 @@ Deno.serve({ port: Number(Deno.env.get("PORT")) || 4001 }, async (request) => { status: 400, }); } - if (!(await backfillUser(did as ActorIdentifier))) { + const pds = await backfillUser(did as ActorIdentifier); + if (!pds) { return new Response(`Failed to refresh ${did}`, { status: 500, }); } - return new Response(`Refreshed ${did}`); + return new Response(pds); } return new Response("Pong!"); }); @@ -153,7 +154,7 @@ for await (const event of jetstream) { async function backfillUser(did: ActorIdentifier) { const cached = getAuthor.get(did); const pds = await getPDS(did as DID, true); - if (!pds || cached?.pds === pds) return false; + if (!pds || cached?.pds === pds) return pds; const handler = simpleFetchHandler({ service: pds }); const rpc = new Client({ handler }); try { @@ -174,5 +175,5 @@ async function backfillUser(did: ActorIdentifier) { } catch (e) { console.error(`Failed to backfill posts: ${e}`); } - return true; + return pds; }