From 25d2ea7b978cc1c1125a0f85a15bdc3e4d9dd43d Mon Sep 17 00:00:00 2001 From: Florian <45694132+flo-bit@users.noreply.github.com> Date: Sun, 12 Apr 2026 03:01:41 +0200 Subject: [PATCH] add following feeds, cache stuff --- scripts/create-placeholder-post.ts | 164 +++++++ scripts/publish-feed-generators.ts | 12 + src/app.d.ts | 18 + src/lib/reddit/bot.ts | 76 ++++ src/lib/reddit/feed-cache.ts | 404 ++++++++++++++++++ src/lib/reddit/server/communities.remote.ts | 25 +- src/routes/+page.svelte | 9 +- src/routes/api/refresh-follows/+server.ts | 33 ++ src/routes/c/[handle]/+page.svelte | 8 + .../app.bsky.feed.getFeedSkeleton/+server.ts | 168 ++++++-- wrangler.jsonc | 12 +- 11 files changed, 881 insertions(+), 48 deletions(-) create mode 100644 scripts/create-placeholder-post.ts create mode 100644 src/lib/reddit/feed-cache.ts create mode 100644 src/routes/api/refresh-follows/+server.ts diff --git a/scripts/create-placeholder-post.ts b/scripts/create-placeholder-post.ts new file mode 100644 index 0000000..fdbdf0c --- /dev/null +++ b/scripts/create-placeholder-post.ts @@ -0,0 +1,164 @@ +// Create the one-time empty-state placeholder post that the +// `following-hot` / `following-new` bsky feed generators emit when +// the viewer follows zero atmo.garden communities. +// +// The feed generator handler prepends this post's URI to the feed +// skeleton response before falling back to the all- contents, +// so users who open a following feed with no subscriptions see: +// +// "🌱 Follow communities on atmo.garden to personalize this feed. +// Browse communities β†’ https://atmo.garden" +// +// …followed by the global hot/new feed. +// +// Usage: +// pnpm tsx scripts/create-placeholder-post.ts +// +// Prints the resulting at-uri. Paste it into wrangler.jsonc as +// `FOLLOWING_FEED_PLACEHOLDER_URI`, then deploy. This script uses +// `createRecord` (not `putRecord`) because the post needs a TID +// rkey for the bsky firehose, so re-running it will create a NEW +// post each time β€” run it once and keep the URI stable. To retire +// or change the placeholder, delete the old record from bsky and +// re-run the script. +// +// Requires the same env vars the discovery-list publisher + feed +// generator publisher use in src/lib/reddit/bot.ts: +// +// ATMO_GARDEN_PDS - e.g. https://bsky.social (or the +// PDS that actually hosts the account) +// ATMO_GARDEN_IDENTIFIER - handle, e.g. atmo.garden +// ATMO_GARDEN_APP_PASSWORD - bsky app password (NOT main password) + +import { readFileSync } from 'fs'; + +function loadEnv(path: string) { + try { + const text = readFileSync(path, 'utf8'); + for (const line of text.split('\n')) { + const m = line.match(/^\s*([A-Z_][A-Z0-9_]*)\s*=\s*(.*)\s*$/i); + if (!m) continue; + let val = m[2]; + if (val.startsWith('"') && val.endsWith('"')) val = val.slice(1, -1); + if (!process.env[m[1]]) process.env[m[1]] = val; + } + } catch { + /* ignore */ + } +} +loadEnv('.env'); +loadEnv('.dev.vars'); + +const PDS = process.env.ATMO_GARDEN_PDS; +const IDENTIFIER = process.env.ATMO_GARDEN_IDENTIFIER; +const APP_PASSWORD = process.env.ATMO_GARDEN_APP_PASSWORD; + +if (!PDS || !IDENTIFIER || !APP_PASSWORD) { + console.error( + 'Missing ATMO_GARDEN_PDS / ATMO_GARDEN_IDENTIFIER / ATMO_GARDEN_APP_PASSWORD in env' + ); + process.exit(1); +} + +// --------------------------------------------------------------------------- +// Post content + richtext facet for the https://atmo.garden link +// --------------------------------------------------------------------------- + +const POST_TEXT = + '🌱 Follow communities on atmo.garden to personalize this feed.\n\nBrowse communities β†’ https://atmo.garden'; + +/** Find the UTF-8 byte range of `substring` inside `text`, or null. */ +function byteRange( + text: string, + substring: string +): { byteStart: number; byteEnd: number } | null { + const charIdx = text.indexOf(substring); + if (charIdx === -1) return null; + const encoder = new TextEncoder(); + const byteStart = encoder.encode(text.slice(0, charIdx)).byteLength; + const byteEnd = byteStart + encoder.encode(substring).byteLength; + return { byteStart, byteEnd }; +} + +const LINK_SUBSTRING = 'https://atmo.garden'; +const linkRange = byteRange(POST_TEXT, LINK_SUBSTRING); +if (!linkRange) { + console.error('internal: could not locate link substring in POST_TEXT'); + process.exit(1); +} + +const facets = [ + { + index: linkRange, + features: [ + { + $type: 'app.bsky.richtext.facet#link', + uri: 'https://atmo.garden' + } + ] + } +]; + +// --------------------------------------------------------------------------- +// 1. Log in +// --------------------------------------------------------------------------- + +console.log(`[1/2] createSession on ${PDS} as ${IDENTIFIER}…`); +const sessionRes = await fetch(`${PDS}/xrpc/com.atproto.server.createSession`, { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ identifier: IDENTIFIER, password: APP_PASSWORD }) +}); +if (!sessionRes.ok) { + console.error( + ` createSession failed (${sessionRes.status}):`, + await sessionRes.text() + ); + process.exit(1); +} +const session = (await sessionRes.json()) as { did: string; accessJwt: string }; +console.log(` did = ${session.did}`); + +// --------------------------------------------------------------------------- +// 2. createRecord the placeholder post +// --------------------------------------------------------------------------- + +console.log(`[2/2] createRecord app.bsky.feed.post…`); +const createRes = await fetch(`${PDS}/xrpc/com.atproto.repo.createRecord`, { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + Authorization: `Bearer ${session.accessJwt}` + }, + body: JSON.stringify({ + repo: session.did, + collection: 'app.bsky.feed.post', + record: { + $type: 'app.bsky.feed.post', + text: POST_TEXT, + facets, + createdAt: new Date().toISOString() + } + }) +}); +if (!createRes.ok) { + console.error( + ` createRecord failed (${createRes.status}):`, + await createRes.text() + ); + process.exit(1); +} +const body = (await createRes.json()) as { uri: string; cid: string }; + +// --------------------------------------------------------------------------- +// 3. Report +// --------------------------------------------------------------------------- + +console.log(`\nβœ… Placeholder post created:\n`); +console.log(` at-uri: ${body.uri}`); +const rkey = body.uri.split('/').pop(); +console.log(` web: https://bsky.app/profile/${IDENTIFIER}/post/${rkey}`); +console.log(`\nPaste the at-uri into wrangler.jsonc as:`); +console.log(` "FOLLOWING_FEED_PLACEHOLDER_URI": "${body.uri}"`); +console.log(`\nThen deploy. The following-hot / following-new feeds will`); +console.log(`prepend this post when the viewer follows zero communities.\n`); diff --git a/scripts/publish-feed-generators.ts b/scripts/publish-feed-generators.ts index 323f555..70b17bd 100644 --- a/scripts/publish-feed-generators.ts +++ b/scripts/publish-feed-generators.ts @@ -94,6 +94,18 @@ const FEEDS: FeedGeneratorSpec[] = [ displayName: 'atmo.garden Top Week', description: 'The most-liked submissions from every atmo.garden community over the last 7 days. Browse + submit at https://atmo.garden.' + }, + { + rkey: 'following-hot', + displayName: 'atmo Following Hot', + description: + 'Hot submissions from only the atmo.garden communities you follow on Bluesky. Follow a community account to add it to this feed. Browse + submit at https://atmo.garden.' + }, + { + rkey: 'following-new', + displayName: 'atmo Following New', + description: + 'Newest submissions from only the atmo.garden communities you follow on Bluesky. Follow a community account to add it to this feed. Browse + submit at https://atmo.garden.' } ]; diff --git a/src/app.d.ts b/src/app.d.ts index e0e3dad..9be3997 100644 --- a/src/app.d.ts +++ b/src/app.d.ts @@ -22,6 +22,15 @@ declare global { COOKIE_SECRET: string; OAUTH_PUBLIC_URL: string; PROFILE_CACHE?: KVNamespace; + /** + * KV namespace holding the materialized sorted feed lists + * (`sorted:` keys, one per PostSort). The cron tick + * rewrites these once per minute from `getCombinedFeed`, + * and both `getHomeFeed` (main page) and the bsky feed + * generator XRPC handler read from them through + * `src/lib/reddit/feed-cache.ts`. + */ + FEEDS_CACHE: KVNamespace; DB: D1Database; COMMUNITY_ENCRYPTION_KEY: string; ROOKERY_HOSTNAME: string; @@ -34,6 +43,15 @@ declare global { ATMO_GARDEN_IDENTIFIER?: string; ATMO_GARDEN_APP_PASSWORD?: string; ATMO_GARDEN_LIST_RKEY?: string; + /** + * at-uri of the single bsky post emitted as the first entry + * in following-* feeds when the viewer follows zero atmo + * communities. Created once via + * `scripts/create-placeholder-post.ts`. Empty string disables + * the placeholder and the feed falls back to the all- + * contents. + */ + FOLLOWING_FEED_PLACEHOLDER_URI?: string; }; } } diff --git a/src/lib/reddit/bot.ts b/src/lib/reddit/bot.ts index 552195b..6ff338e 100644 --- a/src/lib/reddit/bot.ts +++ b/src/lib/reddit/bot.ts @@ -1719,6 +1719,7 @@ export async function runCronTick(env: App.Platform['env']): Promise<{ jetstreamEvents: number; postsRefreshed: number; postsDeleted: number; + feedCachesBuilt: number; errors: string[]; }> { const db = env.DB; @@ -1798,6 +1799,18 @@ export async function runCronTick(env: App.Platform['env']): Promise<{ return 0; }); + // Rebuild the materialized sorted lists + community DID list in + // KV so both the main page (`getHomeFeed`) and the bsky feed + // generator XRPC handler can serve feed requests without hitting + // D1 on the hot path. Runs last so it reflects the freshest state + // after refresh + sweep. Non-fatal: any per-sort failure just + // leaves the previous KV entry in place (the 5 min expirationTtl + // safety net kicks in only if the cron stops running entirely). + const feedCachesBuilt = await rebuildFeedCaches(env, db).catch((e) => { + errors.push(`feed-cache: ${String(e)}`); + return 0; + }); + return { communitiesChecked: communities.length, postsCreated, @@ -1806,9 +1819,72 @@ export async function runCronTick(env: App.Platform['env']): Promise<{ jetstreamEvents: jetstreamResult.events, postsRefreshed, postsDeleted, + feedCachesBuilt, errors }; } +// ------------------------------------------------------------------------- +// Feed-cache materialization +// ------------------------------------------------------------------------- + +/** + * Rebuild the KV-backed sorted lists that both the atmo.garden home + * page and the bsky feed generator XRPC handler read from. + * + * Writes: + * - `sorted:hot`, `sorted:new`, `sorted:top-day`, `sorted:top-week` + * β€” each a JSON-serialized `PostWithCommunity[]` of up to + * `FEED_CACHE_LIMIT` rows from `getCombinedFeed`. + * - `communities:dids` β€” a JSON array of all known community DIDs, + * used by the following-feed path's `getRelationships` batch + * lookup. + * + * Runs the 4 sort queries in parallel (they're independent) so the + * total added latency to `runCronTick` is bounded by the slowest + * query. Returns the number of keys successfully written. + */ +async function rebuildFeedCaches( + env: App.Platform['env'], + db: D1Database +): Promise { + const { FEED_CACHE_LIMIT, writeSortedList, writeAllCommunityDids } = + await import('./feed-cache'); + + // All five sorts the main page exposes in its SortTabs. The bsky + // feed generator only dispatches hot/new/top-day/top-week, but + // materializing top-month here too costs almost nothing (~500 KB + // of KV storage) and lets `getHomeFeed` stay on the fast path + // across all UI sorts instead of falling back to D1 for one sort. + const sorts = ['hot', 'new', 'top-day', 'top-week', 'top-month'] as const; + let written = 0; + + await Promise.all([ + // Full community DID list (small, but refresh it here so the + // following-feed path picks up newly-registered communities + // within one cron tick). + writeAllCommunityDids(env, db) + .then(() => { + written++; + }) + .catch((e) => { + console.error('[rebuildFeedCaches] community DIDs write failed', e); + }), + // Per-sort materializations. + ...sorts.map((sort) => + getCombinedFeed(db, FEED_CACHE_LIMIT, sort, 0) + .then((rows) => writeSortedList(env, sort, rows)) + .then(() => { + written++; + }) + .catch((e) => { + console.error(`[rebuildFeedCaches] ${sort} failed`, e); + }) + ) + ]); + + return written; +} + // Convenience re-export so routes don't need to import db.ts directly. export { getCombinedFeed }; diff --git a/src/lib/reddit/feed-cache.ts b/src/lib/reddit/feed-cache.ts new file mode 100644 index 0000000..9a6f0b0 --- /dev/null +++ b/src/lib/reddit/feed-cache.ts @@ -0,0 +1,404 @@ +// Shared cache layer for the feed pipeline. +// +// All of our feed surfaces β€” the atmo.garden home page +// (`getHomeFeed`), the bsky feed generator XRPC handler +// (`getFeedSkeleton`), and the following-feed variants β€” read from +// the same materialized sorted lists stored in Workers KV (binding +// `FEEDS_CACHE`). The cron tick rewrites those lists once per minute +// from `getCombinedFeed`, so the hot path never touches D1. +// +// Two cache tiers: +// +// 1. Workers KV (`env.FEEDS_CACHE`) β€” global, durable, written by +// the cron once per minute. Source of truth. +// 2. Workers Cache API (per-colocation edge cache, free) β€” sits in +// front of KV for both sorted lists and per-viewer follow +// intersections. 30 s TTL on lists, 5 min TTL on follow sets. +// +// The following-feed paths also cache each viewer's set of followed +// community DIDs through the same Cache API layer. Invalidation is +// best-effort per-colo via `invalidateViewerCommunityFollows`, +// triggered by the `POST /api/refresh-follows` endpoint the UI hits +// after a user toggles a community follow. +// +// See `scripts/publish-feed-generators.ts` for the feed records this +// cache eventually serves, and the main cron's `rebuildFeedCaches` +// step for how the KV entries get populated. + +import type { PostSort, PostWithCommunity } from './db'; +import { listCommunities } from './db'; + +// Cloudflare Workers Cache API β€” `caches.default` is a CF-specific +// extension not present on the DOM's `CacheStorage` type, AND the +// `caches` global itself isn't defined in vite's Node SSR runtime +// during `pnpm dev`. We cast once at the top and fall back to a +// no-op shim when the global is missing so the module still loads +// in dev β€” all cache operations become harmless no-ops and every +// request falls through to KV (which DOES work via platformProxy +// against the remote binding). In production on Workers the real +// `caches.default` takes over transparently. +type MinimalCache = { + match(key: Request): Promise; + put(key: Request, value: Response): Promise; + delete(key: Request): Promise; +}; + +const noopCache: MinimalCache = { + async match() { + return undefined; + }, + async put() { + /* no-op */ + }, + async delete() { + return false; + } +}; + +const cfCache: MinimalCache = + typeof caches !== 'undefined' && + (caches as unknown as { default?: MinimalCache }).default + ? (caches as unknown as { default: MinimalCache }).default + : noopCache; + +// --------------------------------------------------------------------------- +// Sorted-list cache (materialized `getCombinedFeed` output per sort) +// --------------------------------------------------------------------------- + +/** + * Upper bound on how many rows per sort we materialize and store in + * KV. The main page's infinite scroll and every feed generator call + * share this cap β€” users can page through the top 1000 per sort, and + * scrolling past that stops. Bump if users start complaining; lower + * if KV storage gets tight. 1000 rows Γ— ~500 bytes/row = ~500 KB per + * sort, well under KV's 25 MB value limit. + */ +export const FEED_CACHE_LIMIT = 1000; + +/** + * KV key naming: one key per sort. Values are JSON-serialized + * `PostWithCommunity[]` with length ≀ `FEED_CACHE_LIMIT`, in sort + * order. `missing_since IS NULL` is already applied by the underlying + * `getCombinedFeed` query so we never surface rows that are pending + * sweep. + */ +function kvKeyForSort(sort: PostSort): string { + return `sorted:${sort}`; +} + +/** + * Edge-cache key for the sorted list. Workers Cache API keys are + * Request objects; the URL doesn't have to resolve β€” it's purely a + * deterministic identifier. + */ +function cacheKeyForSort(sort: PostSort): Request { + return new Request(`https://cache.internal/feeds/sorted/${sort}`); +} + +/** + * Write the materialized sorted list for a given sort into KV. Called + * by the cron's `rebuildFeedCaches` step. Includes a 5 min + * `expirationTtl` as a belt-and-suspenders safety net β€” the cron + * overwrites every minute, so the TTL only matters if the cron stops + * running entirely. + */ +export async function writeSortedList( + env: App.Platform['env'], + sort: PostSort, + rows: PostWithCommunity[] +): Promise { + await env.FEEDS_CACHE.put(kvKeyForSort(sort), JSON.stringify(rows), { + expirationTtl: 300 + }); +} + +/** + * Read a sorted list through the edge cache, falling back to KV. + * + * Flow: + * 1. Check Workers Cache API. Hit β†’ return parsed JSON. + * 2. Miss β†’ fetch from KV. Returns [] if the cron hasn't populated + * it yet (brand-new deploy or KV key was manually cleared). + * 3. Populate the edge cache with a 30 s `s-maxage`. + * + * The 30 s edge TTL is chosen to be well under the cron tick + * frequency, so each colo hits KV at most ~2Γ—/min per sort β€” ~1.4 M + * KV reads/month across ~200 colos, cheap. + */ +export async function getCachedSortedList( + env: App.Platform['env'], + sort: PostSort +): Promise { + const cacheKey = cacheKeyForSort(sort); + const cached = await cfCache.match(cacheKey); + if (cached) { + try { + return (await cached.json()) as PostWithCommunity[]; + } catch (e) { + console.error(`[feed-cache] cached list parse failed for ${sort}`, e); + // fall through to re-fetch + } + } + + const raw = await env.FEEDS_CACHE.get(kvKeyForSort(sort)); + const list: PostWithCommunity[] = raw ? (JSON.parse(raw) as PostWithCommunity[]) : []; + + await cfCache.put( + cacheKey, + new Response(JSON.stringify(list), { + headers: { + 'content-type': 'application/json', + 'cache-control': 'public, s-maxage=30' + } + }) + ); + return list; +} + +// --------------------------------------------------------------------------- +// Community-DID list (small, ~100 entries, rarely changes) +// --------------------------------------------------------------------------- + +const COMMUNITY_DIDS_KV_KEY = 'communities:dids'; +const COMMUNITY_DIDS_CACHE_KEY = new Request('https://cache.internal/communities/dids'); + +/** + * Rewrite the list of all known community DIDs into KV. Called by + * the cron alongside the sorted-list rebuild β€” the community set + * changes rarely (new registrations) so hourly-ish staleness is + * fine, but running it on every tick is also trivially cheap. + */ +export async function writeAllCommunityDids( + env: App.Platform['env'], + db: D1Database +): Promise { + const rows = await listCommunities(db); + const dids = rows.map((r) => r.did); + await env.FEEDS_CACHE.put(COMMUNITY_DIDS_KV_KEY, JSON.stringify(dids), { + expirationTtl: 3600 + }); +} + +/** + * Read the list of all known community DIDs, edge-cached. Used by + * `fetchViewerCommunityRelationships` to know which DIDs to ask bsky + * about in each `getRelationships` batch. + */ +export async function getAllCommunityDids(env: App.Platform['env']): Promise { + const cached = await cfCache.match(COMMUNITY_DIDS_CACHE_KEY); + if (cached) { + try { + return (await cached.json()) as string[]; + } catch { + /* fall through */ + } + } + + const raw = await env.FEEDS_CACHE.get(COMMUNITY_DIDS_KV_KEY); + const dids: string[] = raw ? (JSON.parse(raw) as string[]) : []; + + await cfCache.put( + COMMUNITY_DIDS_CACHE_KEY, + new Response(JSON.stringify(dids), { + headers: { + 'content-type': 'application/json', + 'cache-control': 'public, s-maxage=300' + } + }) + ); + return dids; +} + +// --------------------------------------------------------------------------- +// Per-viewer community-follows cache +// --------------------------------------------------------------------------- + +const BSKY_APPVIEW_PUBLIC = 'https://public.api.bsky.app'; + +/** + * Max DIDs allowed per `app.bsky.graph.getRelationships` call, per + * the lexicon's `others` array maxLength. Any more and bsky returns + * 400. + */ +const RELATIONSHIPS_BATCH_SIZE = 30; + +function viewerFollowsCacheKey(viewerDid: string): Request { + // URL-encode the DID so colons etc. don't confuse the cache. + return new Request( + `https://cache.internal/follows/${encodeURIComponent(viewerDid)}` + ); +} + +/** + * Call `app.bsky.graph.getRelationships` against the public bsky + * appview in parallel batches, asking whether `viewerDid` follows + * each community DID. Returns the subset of community DIDs that + * `viewerDid` follows. + * + * Uses `getRelationships` instead of `getFollows` because we only + * care about a specific small set of DIDs (the community accounts), + * not the user's entire follow graph. This is O(communities) instead + * of O(user's follows) and scales with a stable number (~100). + */ +export async function fetchViewerCommunityRelationships( + viewerDid: string, + communityDids: string[] +): Promise { + if (communityDids.length === 0) return []; + + // Chunk into batches of 30. + const batches: string[][] = []; + for (let i = 0; i < communityDids.length; i += RELATIONSHIPS_BATCH_SIZE) { + batches.push(communityDids.slice(i, i + RELATIONSHIPS_BATCH_SIZE)); + } + + const results = await Promise.all( + batches.map(async (batch) => { + const url = new URL( + `${BSKY_APPVIEW_PUBLIC}/xrpc/app.bsky.graph.getRelationships` + ); + url.searchParams.set('actor', viewerDid); + // `others` is an array; repeated query params. + for (const did of batch) url.searchParams.append('others', did); + try { + const res = await fetch(url); + if (!res.ok) { + console.error( + '[feed-cache] getRelationships non-ok', + res.status, + batch.length + ); + return [] as string[]; + } + const body = (await res.json()) as { + relationships?: Array< + | { did: string; following?: string; followedBy?: string } + | { actor: string; notFound: true } + >; + }; + const followed: string[] = []; + for (const rel of body.relationships ?? []) { + // `following` is set when `viewerDid` follows this DID. + if ('did' in rel && typeof rel.following === 'string') { + followed.push(rel.did); + } + } + return followed; + } catch (e) { + console.error('[feed-cache] getRelationships threw', e); + return []; + } + }) + ); + + // Flatten. Any batch that failed just contributes nothing β€” we + // return a possibly-incomplete follow set, which silently hides + // some of the user's real follows until the next refresh. Better + // than erroring out the whole feed. + return results.flat(); +} + +/** + * Read the cached set of community DIDs the viewer follows. Cache + * miss triggers a fresh `getRelationships` fan-out (parallel batches, + * ~150ms worst case) and a 5 min edge-cache write. + * + * Returns the intersection of "communities on atmo.garden" ∩ "DIDs + * the viewer follows on bsky", using bsky's native graph as the + * source of truth. No D1 writes, no new follow table β€” the user's + * `app.bsky.graph.follow` records on their own PDS drive everything. + */ +export async function getCachedViewerCommunityFollows( + env: App.Platform['env'], + viewerDid: string +): Promise { + const cacheKey = viewerFollowsCacheKey(viewerDid); + const cached = await cfCache.match(cacheKey); + if (cached) { + try { + return (await cached.json()) as string[]; + } catch { + /* fall through */ + } + } + + const communityDids = await getAllCommunityDids(env); + const followed = await fetchViewerCommunityRelationships(viewerDid, communityDids); + + await cfCache.put( + cacheKey, + new Response(JSON.stringify(followed), { + headers: { + 'content-type': 'application/json', + 'cache-control': 'public, s-maxage=300' + } + }) + ); + return followed; +} + +/** + * Purge the cached follow set for a viewer and eagerly repopulate it + * from bsky. Called by `POST /api/refresh-follows` after the UI + * writes a new `app.bsky.graph.follow` record. + * + * Important caveat: Workers Cache API is **per-colocation**, so this + * only busts the cache at the colo handling the current request. + * Other colos will continue to serve stale data until their 5 min + * TTL expires. In practice this is fine β€” users' follow-then-reload + * flows stick to the same colo via TCP/TLS stickiness and geo + * routing. Cross-colo invalidation would require moving this cache + * to KV, which roughly triples monthly cost. + */ +export async function invalidateViewerCommunityFollows( + env: App.Platform['env'], + viewerDid: string +): Promise { + const cacheKey = viewerFollowsCacheKey(viewerDid); + await cfCache.delete(cacheKey); + // Eagerly repopulate so the refresh round-trip includes the fresh + // data β€” the UI can immediately use it (or at least knows the + // bsky graph has propagated). + return getCachedViewerCommunityFollows(env, viewerDid); +} + +// --------------------------------------------------------------------------- +// JWT helper (unverified payload parsing) +// --------------------------------------------------------------------------- + +/** + * Pull the `iss` claim out of a service-auth JWT in the + * `Authorization: Bearer ` header. Does **not** verify the + * signature β€” the attack surface is "see how a different user's + * feed would look," which is bounded by the fact that bsky follow + * graphs are public anyway. Proper signature verification would + * require fetching the caller's DID doc and running ES256K against + * their repo signing key; worth doing eventually, not blocking for + * MVP. + * + * Also rejects obviously-expired JWTs as a minimal sanity check. + */ +export function parseViewerDidFromJwt(authHeader: string | null): string | null { + if (!authHeader) return null; + const m = authHeader.match(/^Bearer (.+)$/i); + if (!m) return null; + const parts = m[1].split('.'); + if (parts.length !== 3) return null; + + try { + // base64url β†’ base64 β†’ atob + let payloadB64 = parts[1].replace(/-/g, '+').replace(/_/g, '/'); + while (payloadB64.length % 4 !== 0) payloadB64 += '='; + const json = atob(payloadB64); + const obj = JSON.parse(json) as { iss?: unknown; exp?: unknown }; + + // Sanity: reject JWTs whose expiry is already past. + if (typeof obj.exp === 'number' && obj.exp * 1000 < Date.now()) return null; + + const iss = obj.iss; + if (typeof iss !== 'string' || !iss.startsWith('did:')) return null; + return iss; + } catch { + return null; + } +} diff --git a/src/lib/reddit/server/communities.remote.ts b/src/lib/reddit/server/communities.remote.ts index 76de403..0965aa1 100644 --- a/src/lib/reddit/server/communities.remote.ts +++ b/src/lib/reddit/server/communities.remote.ts @@ -3,7 +3,6 @@ import { command, getRequestEvent } from '$app/server'; import * as v from 'valibot'; import { getRecentPostsForCommunity, - getCombinedFeed, getPostByUri, listCommunities, getCommunityByHandle, @@ -18,6 +17,7 @@ import { checkCanSubmit, refreshCommunityCache } from '../bot'; +import { getCachedSortedList, FEED_CACHE_LIMIT } from '../feed-cache'; import { parseListUri } from '../list-uri'; import { ACCENT_COLORS, @@ -245,6 +245,17 @@ export const getCommunityPost = command( } ); +/** + * Home feed for atmo.garden's main page. Reads from the KV-materialized + * sorted list that the cron tick rebuilds every minute β€” zero D1 on + * the hot path. Uses the same cache as the bsky feed generator XRPC + * handler so the two stay in sync by construction. + * + * Pagination is offset-based, capped at `FEED_CACHE_LIMIT` total + * rows β€” infinite scroll on the main page stops when the client has + * consumed all materialized entries. If someone needs to go deeper, + * bump `FEED_CACHE_LIMIT` in feed-cache.ts. + */ export const getHomeFeed = command( v.object({ limit: v.optional(v.number()), @@ -254,14 +265,12 @@ export const getHomeFeed = command( async (input): Promise => { const { platform } = getRequestEvent(); const env = platform?.env; - if (!env || !env.DB) return []; + if (!env) return []; - return getCombinedFeed( - env.DB, - input.limit ?? 50, - (input.sort ?? 'hot') as PostSort, - input.offset ?? 0 - ); + const list = await getCachedSortedList(env, (input.sort ?? 'hot') as PostSort); + const offset = Math.max(0, input.offset ?? 0); + const limit = Math.max(1, Math.min(FEED_CACHE_LIMIT, input.limit ?? 50)); + return list.slice(offset, offset + limit); } ); diff --git a/src/routes/+page.svelte b/src/routes/+page.svelte index ad41763..87c9ecc 100644 --- a/src/routes/+page.svelte +++ b/src/routes/+page.svelte @@ -11,6 +11,11 @@ type SubmitterProfile = { handle: string; displayName: string | null; avatar: string | null }; const PAGE_SIZE = 50; + // Matches `FEED_CACHE_LIMIT` in src/lib/reddit/feed-cache.ts β€” the + // main page is served entirely from the KV-materialized sorted + // list, so infinite scroll stops once we've consumed it. Bumping + // this here requires also bumping FEED_CACHE_LIMIT on the server. + const FEED_MAX_ROWS = 1000; let loading = $state(true); let loadingMore = $state(false); @@ -49,7 +54,7 @@ quoted = quotedRes.posts; submitters = profileRes.profiles; } - hasMore = rows.length >= PAGE_SIZE; + hasMore = rows.length >= PAGE_SIZE && feed.length < FEED_MAX_ROWS; } catch (e) { console.error('[home] loadFeed failed', e); } finally { @@ -75,7 +80,7 @@ quoted = { ...quoted, ...quotedRes.posts }; submitters = { ...submitters, ...profileRes.profiles }; } - if (rows.length < PAGE_SIZE) { + if (rows.length < PAGE_SIZE || feed.length >= FEED_MAX_ROWS) { hasMore = false; } } catch (e) { diff --git a/src/routes/api/refresh-follows/+server.ts b/src/routes/api/refresh-follows/+server.ts new file mode 100644 index 0000000..aab6ddc --- /dev/null +++ b/src/routes/api/refresh-follows/+server.ts @@ -0,0 +1,33 @@ +// Purge the edge-cached community-follow set for the authenticated +// user, then eagerly repopulate it from bsky's public appview via +// `getRelationships`. Called by the community page's follow / unfollow +// button after a successful `app.bsky.graph.follow` (or delete) write +// on the user's PDS, so the next `getFeedSkeleton` call for this +// user picks up the new community subscription without waiting for +// the default 5 min TTL. +// +// Per-colo invalidation only β€” see `invalidateViewerCommunityFollows` +// in `src/lib/reddit/feed-cache.ts` for the tradeoff. In practice the +// user's "follow β†’ reload feed" flow stays on one Cloudflare colo via +// TCP/TLS stickiness + geo routing, so the stale-in-other-colos +// window is a non-issue for interactive UX. + +import type { RequestHandler } from './$types'; +import { json, error } from '@sveltejs/kit'; +import { invalidateViewerCommunityFollows } from '$lib/reddit/feed-cache'; + +export const POST: RequestHandler = async ({ platform, locals }) => { + const env = platform?.env; + if (!env) error(500, 'Platform env unavailable'); + if (!locals.did) error(401, 'Not authenticated'); + + // `invalidateViewerCommunityFollows` deletes the cached entry at + // the current colo and then re-fetches from bsky, so this request + // returns with the fresh data baked into the per-colo cache. The + // UI doesn't need the returned list today, but we include it in + // the response in case it wants to render immediate feedback + // ("you now follow N communities") without a second round-trip. + const followed = await invalidateViewerCommunityFollows(env, locals.did); + + return json({ ok: true, followedCount: followed.length }); +}; diff --git a/src/routes/c/[handle]/+page.svelte b/src/routes/c/[handle]/+page.svelte index 01fc204..978be1d 100644 --- a/src/routes/c/[handle]/+page.svelte +++ b/src/routes/c/[handle]/+page.svelte @@ -188,6 +188,14 @@ const result = await followUser({ did: community.did }); followUri = result.uri; } + // Purge the edge-cached community-follow set for this + // viewer so the bsky following-hot / following-new feed + // generators pick up the change on the next request + // without waiting for the 5 min TTL. Fire-and-forget β€” + // the refresh is a nice-to-have, not load-bearing. + fetch('/api/refresh-follows', { method: 'POST' }).catch((e) => { + console.error('[community] refresh-follows failed', e); + }); } catch (e) { console.error('[community] toggle join failed', e); } finally { diff --git a/src/routes/xrpc/app.bsky.feed.getFeedSkeleton/+server.ts b/src/routes/xrpc/app.bsky.feed.getFeedSkeleton/+server.ts index 6226223..6f2cabe 100644 --- a/src/routes/xrpc/app.bsky.feed.getFeedSkeleton/+server.ts +++ b/src/routes/xrpc/app.bsky.feed.getFeedSkeleton/+server.ts @@ -1,7 +1,32 @@ // Bluesky feed generator: serves the `app.bsky.feed.getFeedSkeleton` -// XRPC endpoint so the two `app.bsky.feed.generator` records published -// by atmo.garden (`all-hot` and `top-day`) can be subscribed to from -// any bsky client. +// XRPC endpoint. Handles two kinds of feeds: +// +// scope = 'all' β€” global feeds (all-hot, all-new, top-day, +// top-week). Reads from the KV-materialized +// sorted list, slices the requested page, +// maps to skeleton entries. +// +// scope = 'following' β€” per-viewer personalized feeds +// (following-hot, following-new). Reads the +// same KV list, filters down to posts +// whose community_did is in the viewer's +// followed-community set, then slices. +// Zero subscriptions β†’ emit an optional +// placeholder post followed by the all- +// contents as a fallback. +// +// Both paths are zero-D1 on the hot request path β€” all reads come +// from Workers KV with an edge-cache tier in front (see +// src/lib/reddit/feed-cache.ts). The cron tick rebuilds the KV +// entries once per minute, so feed freshness is bounded by that. +// +// Following feeds pull the viewer DID from the `Authorization: +// Bearer ` header's `iss` claim, parsed without signature +// verification. The attack surface is "look at how a different user's +// feed would render" β€” which an attacker could already compute from +// bsky's public follow graph β€” so unverified parsing is acceptable +// for MVP. Proper signature verification would require fetching the +// caller's DID doc and running ES256K, worth doing eventually. // // Wire shape (per lexicons.atproto.com / app.bsky.feed.getFeedSkeleton): // @@ -19,7 +44,7 @@ // } // } // -// atmo.garden serves reposts and quote posts differently: +// Post routing for skeleton entries: // // - Community REPOST rows (`uri` β†’ `app.bsky.feed.repost/...`): // emit `{ post: quoted_post_uri, reason: skeletonReasonRepost(uri) }`. @@ -31,26 +56,32 @@ // emit `{ post: uri }`. The community's quote post is itself a // valid bsky post with an `app.bsky.embed.record` embed, so the // appview hydrates it with the original post rendered inline. -// -// Feeds are matched by the rkey of the `feed` query param -// (`all-hot` / `top-day`). We intentionally don't validate the -// authority DID β€” if someone publishes a generator record on their -// own account that points at our service, they'll get the same -// feed data, but they can't modify it, so it's harmless. import type { RequestHandler } from './$types'; import { json, error } from '@sveltejs/kit'; -import { getCombinedFeed, type PostSort, type PostWithCommunity } from '$lib/reddit/db'; +import type { PostSort, PostWithCommunity } from '$lib/reddit/db'; +import { + getCachedSortedList, + getCachedViewerCommunityFollows, + parseViewerDidFromJwt +} from '$lib/reddit/feed-cache'; const DEFAULT_LIMIT = 50; const MAX_LIMIT = 100; -/** Map a known feed rkey β†’ the `PostSort` we drive `getCombinedFeed` with. */ -const FEED_RKEY_TO_SORT: Record = { - 'all-hot': 'hot', - 'all-new': 'new', - 'top-day': 'top-day', - 'top-week': 'top-week' +type FeedConfig = { + sort: PostSort; + scope: 'all' | 'following'; +}; + +/** Map a known feed rkey β†’ config (sort + scope). */ +const FEED_RKEY_TO_CONFIG: Record = { + 'all-hot': { sort: 'hot', scope: 'all' }, + 'all-new': { sort: 'new', scope: 'all' }, + 'top-day': { sort: 'top-day', scope: 'all' }, + 'top-week': { sort: 'top-week', scope: 'all' }, + 'following-hot': { sort: 'hot', scope: 'following' }, + 'following-new': { sort: 'new', scope: 'following' } }; type SkeletonFeedPost = { @@ -104,53 +135,116 @@ function rowToSkeleton(row: PostWithCommunity): SkeletonFeedPost | null { return null; } -export const GET: RequestHandler = async ({ url, platform }) => { +/** + * Build a response envelope with cursor handling. The cursor is an + * opaque string per the lexicon spec; we use base-10 offset because + * everything downstream is offset-paginated. Only emit a cursor + * when the page was filled β€” saves one empty paginated round-trip + * at the end of the feed. + */ +function buildResponse( + entries: SkeletonFeedPost[], + offset: number, + limit: number, + totalAvailable: number +): { feed: SkeletonFeedPost[]; cursor?: string } { + const page = entries.slice(0, limit); + const nextOffset = offset + page.length; + const nextCursor = nextOffset < totalAvailable && page.length === limit ? String(nextOffset) : undefined; + return { + feed: page, + ...(nextCursor ? { cursor: nextCursor } : {}) + }; +} + +export const GET: RequestHandler = async ({ url, request, platform }) => { const env = platform?.env; - if (!env?.DB) error(500, 'DB binding unavailable'); + if (!env) error(500, 'Platform env unavailable'); const feedUri = url.searchParams.get('feed'); if (!feedUri) error(400, 'missing `feed` parameter'); const rkey = parseFeedRkey(feedUri); - const sort = rkey ? FEED_RKEY_TO_SORT[rkey] : undefined; - if (!sort) { + const config = rkey ? FEED_RKEY_TO_CONFIG[rkey] : undefined; + if (!config) { // Unknown feed. The getFeedSkeleton lexicon defines a single // named error for this case (`errors: [{ name: "UnknownFeed" }]`) // and the atproto XRPC error convention is a 400 with a JSON // body of `{ error: "", message: "" }` β€” clients - // branch on the name to render a "feed unavailable" state, - // which is better UX than silently returning an empty feed. + // branch on the name to render a "feed unavailable" state. return json( { error: 'UnknownFeed', message: `Feed not found: ${feedUri}` }, { status: 400 } ); } - // Clamp limit to the [1, MAX_LIMIT] range per the lexicon spec. + // Clamp limit to [1, MAX_LIMIT] per the lexicon spec. const limitRaw = url.searchParams.get('limit'); const limitParsed = limitRaw ? parseInt(limitRaw, 10) : DEFAULT_LIMIT; const limit = Number.isFinite(limitParsed) ? Math.max(1, Math.min(MAX_LIMIT, limitParsed)) : DEFAULT_LIMIT; - // Cursor is an opaque string per spec; we use a base-10 offset - // because `getCombinedFeed` already accepts offset-based pagination - // and we don't need anything fancier. Invalid cursors fall back - // to offset 0. + // Cursor β†’ offset. Invalid cursors fall back to 0. const cursorRaw = url.searchParams.get('cursor'); const cursorParsed = cursorRaw ? parseInt(cursorRaw, 10) : 0; const offset = Number.isFinite(cursorParsed) && cursorParsed >= 0 ? cursorParsed : 0; - const rows = await getCombinedFeed(env.DB, limit, sort, offset); - const feed = rows.map(rowToSkeleton).filter((x): x is SkeletonFeedPost => x !== null); + // Pull the materialized global sorted list once β€” both scopes + // read from the same cached entry. + const sortedList = await getCachedSortedList(env, config.sort); - // Only emit a cursor when we filled the page β€” saves one empty - // paginated round-trip for the common "we've reached the end" case. - const nextCursor = rows.length === limit ? String(offset + rows.length) : undefined; + if (config.scope === 'all') { + const slice = sortedList + .slice(offset, offset + limit) + .map(rowToSkeleton) + .filter((x): x is SkeletonFeedPost => x !== null); + return json(buildResponse(slice, offset, limit, sortedList.length)); + } - return json({ - feed, - ...(nextCursor ? { cursor: nextCursor } : {}) - }); + // ----- scope === 'following' ----- + + // Extract viewer DID from the service-auth JWT. Missing / malformed + // JWT β†’ treat as zero-subscription fallback, which surfaces the + // placeholder + all- contents so an unauthenticated client + // peek at the feed still renders something useful. + const viewerDid = parseViewerDidFromJwt(request.headers.get('authorization')); + + let filtered: PostWithCommunity[] = []; + if (viewerDid) { + const followed = await getCachedViewerCommunityFollows(env, viewerDid); + if (followed.length > 0) { + const followedSet = new Set(followed); + filtered = sortedList.filter((r) => followedSet.has(r.community_did)); + } + } + + if (filtered.length === 0) { + // Empty-state: show the placeholder post (if configured) then + // fall back to the all- list so the feed always has + // content. Placeholder absent (env var empty) β†’ just fall back + // to the all- contents. + const placeholderUri = env.FOLLOWING_FEED_PLACEHOLDER_URI ?? ''; + const fallback = sortedList + .slice(offset, offset + limit) + .map(rowToSkeleton) + .filter((x): x is SkeletonFeedPost => x !== null); + + // Only prepend the placeholder on the first page (offset 0) β€” + // otherwise subsequent pages would have it re-inserted. + const entries: SkeletonFeedPost[] = + offset === 0 && placeholderUri + ? [{ post: placeholderUri }, ...fallback].slice(0, limit) + : fallback; + + return json(buildResponse(entries, offset, limit, sortedList.length)); + } + + // Normal path: viewer has real subscriptions. + const slice = filtered + .slice(offset, offset + limit) + .map(rowToSkeleton) + .filter((x): x is SkeletonFeedPost => x !== null); + return json(buildResponse(slice, offset, limit, filtered.length)); }; diff --git a/wrangler.jsonc b/wrangler.jsonc index 53a8532..9c6e166 100644 --- a/wrangler.jsonc +++ b/wrangler.jsonc @@ -13,7 +13,13 @@ }, "vars": { "OAUTH_PUBLIC_URL": "https://atmo.garden", - "ROOKERY_HOSTNAME": "pds.atmo.garden" + "ROOKERY_HOSTNAME": "pds.atmo.garden", + // at-uri of the single bsky post emitted as the first entry in + // following-* feeds when the viewer follows zero communities. + // Created once via `pnpm tsx scripts/create-placeholder-post.ts` + // and pasted here. Empty string disables the placeholder (feed + // falls back directly to the all- contents). + "FOLLOWING_FEED_PLACEHOLDER_URI": "" }, "kv_namespaces": [ { @@ -23,6 +29,10 @@ { "binding": "OAUTH_STATES", "id": "8cff96d0d78f4966a09b6b8833cab54e" + }, + { + "binding": "FEEDS_CACHE", + "id": "c37213227cec459887c3c4bd69ffa636" } ], "d1_databases": [ -- 2.51.2