import { AppBskyFeedPost } from "@atcute/bluesky"; import { type JetstreamEvent, JetstreamSubscription } from "@atcute/jetstream"; import { is } from "@atcute/lexicons"; import type { AtprotoDid } from "@atcute/lexicons/syntax"; import { eq } from "drizzle-orm"; import db, { jetstreamState, selfDeletingPosts, sessions } from "../db"; import { loggers } from "../logger"; const logger = loggers.jetstream; let subscription: JetstreamSubscription | null = null; let subscriptionIterator: AsyncIterator | null = null; function getDids() { const allSessions = db.select().from(sessions).all(); return allSessions.map((session) => session.did as AtprotoDid); } function getStoredCursor() { const state = db .select() .from(jetstreamState) .where(eq(jetstreamState.id, "singleton")) .get(); logger.info(state, "retrieved cursor"); return state?.cursor; } function updateStoredCursor(cursor: number) { db.insert(jetstreamState) .values({ id: "singleton", cursor, updatedAt: new Date(), }) .onConflictDoUpdate({ target: jetstreamState.id, set: { cursor, updatedAt: new Date(), }, }) .execute(); } async function handleEvent(event: JetstreamEvent) { const { did, kind } = event; if (kind !== "commit") return; const { commit } = event; if (commit.operation !== "create") return; if (commit.collection !== "app.bsky.feed.post") return; const { record: post, rkey } = commit; if (!is(AppBskyFeedPost.mainSchema, post)) return; const uri = `at://${did}/app.bsky.feed.post/${rkey}`; const postId = crypto.randomUUID(); const labels = post.labels; const marker = labels?.values.find((label) => label.val.startsWith("cat.vt3e.bbell#marker"), ); if (!marker) return; const expiresAtEpoch = marker.val.split(":")[1]; if (!expiresAtEpoch || Number.isNaN(Number(expiresAtEpoch))) return; const date = new Date(Number(expiresAtEpoch)); if (date < new Date()) return; const base = { id: postId, did, uri, }; logger.info(base, "creating post deletion"); const [result] = await db .insert(selfDeletingPosts) .values({ id: postId, createdAt: new Date(), deletesAt: date, did: did, uri: uri, status: "active", }) .returning() .execute(); if (!result) { logger.error(base, "failed to create post deletion"); throw new Error("failed to create post deletion"); } logger.info(base, "post deletion created"); } export async function startListening() { if (subscription) return; const wantedDids = getDids(); if (!wantedDids || wantedDids.length === 0) { logger.info("no sessions found, not starting jetstream"); return; } const storedCursor = getStoredCursor(); subscription = new JetstreamSubscription({ url: "wss://jetstream2.us-east.bsky.network", cursor: storedCursor, wantedDids, wantedCollections: ["app.bsky.feed.post"], onConnectionOpen: (event) => { const wanted = subscription?.getOptions().wantedDids || []; logger.info( event, `jetstream connected, listening to ${wanted?.length || 0} repo(s)`, ); logger.trace(wanted, "listening to repos"); }, }); subscriptionIterator = subscription[Symbol.asyncIterator](); try { while (true) { const { value: event, done } = await subscriptionIterator.next(); if (done) break; await handleEvent(event); updateStoredCursor(subscription.cursor); } } catch (err) { logger.error(err, "jetstream loop error"); } finally { try { await subscriptionIterator?.return?.(); } catch {} subscription = null; subscriptionIterator = null; logger.info("jetstream subscription ended"); } } export function updateListeningDids() { const wanted = getDids(); if (subscription) { if (!wanted || wanted.length === 0) { logger.info("no dids left, stopping jetstream subscription"); try { subscriptionIterator?.return?.(); } catch (err) { logger.warn({ err }, "error while returning iterator"); } subscriptionIterator = null; return; } subscription.updateOptions({ wantedDids: wanted, }); logger.info("updated jetstream DIDs"); return; } if (wanted && wanted.length > 0) { logger.info("dids present, starting jetstream"); startListening().catch((err) => { logger.error(err, "failed to start jetstream"); }); return; } logger.info("no dids to listen to"); }