diff --git a/apps/web/src/app.d.ts b/apps/web/src/app.d.ts index c62ab7e..40a16ec 100644 --- a/apps/web/src/app.d.ts +++ b/apps/web/src/app.d.ts @@ -72,6 +72,12 @@ declare global { OAUTH_PUBLIC_URL: string; DB: D1Database; CRON_SECRET: string; + /** Handle or DID of the reply bot, e.g. `going.atmo.rsvp`. */ + BOT_IDENTIFIER?: string; + /** App password for the bot account (set via `wrangler secret put`). */ + BOT_APP_PASSWORD?: string; + /** PDS the bot authenticates against; defaults to bsky.social. */ + BOT_PDS_URL?: string; }; /** Cloudflare Worker execution context. Use `ctx.waitUntil(promise)` to * let the worker keep a fire-and-forget task alive after the response diff --git a/apps/web/src/lib/bot/actions.ts b/apps/web/src/lib/bot/actions.ts new file mode 100644 index 0000000..60f3dcb --- /dev/null +++ b/apps/web/src/lib/bot/actions.ts @@ -0,0 +1,66 @@ +import type { Client } from '@atcute/client'; +import type { Did } from '@atcute/lexicons'; +import { LIKE_COLLECTION, POST_COLLECTION } from './config'; + +export type PostRef = { uri: string; cid: string }; + +const encoder = new TextEncoder(); +const URL_RE = /https?:\/\/[^\s)\]]+/g; + +/** Build `#link` facets for any URLs in `text`, using UTF-8 byte offsets. */ +function buildLinkFacets(text: string) { + const facets: Array<{ + index: { byteStart: number; byteEnd: number }; + features: Array<{ $type: 'app.bsky.richtext.facet#link'; uri: string }>; + }> = []; + + for (const match of text.matchAll(URL_RE)) { + // Trim trailing punctuation that's part of the sentence, not the URL. + const uri = match[0].replace(/[.,;:!?]+$/, ''); + const byteStart = encoder.encode(text.slice(0, match.index)).length; + const byteEnd = byteStart + encoder.encode(uri).length; + facets.push({ + index: { byteStart, byteEnd }, + features: [{ $type: 'app.bsky.richtext.facet#link', uri }] + }); + } + + return facets.length ? facets : undefined; +} + +/** Like a post as the bot (confirms a successful RSVP). */ +export async function likePost(client: Client, botDid: Did, subject: PostRef): Promise { + await client.post('com.atproto.repo.createRecord', { + input: { + repo: botDid, + collection: LIKE_COLLECTION, + record: { + $type: LIKE_COLLECTION, + subject, + createdAt: new Date().toISOString() + } + } + }); +} + +/** Reply to a post as the bot. `root`/`parent` thread the reply correctly. */ +export async function replyToPost( + client: Client, + botDid: Did, + opts: { root: PostRef; parent: PostRef; text: string } +): Promise { + const facets = buildLinkFacets(opts.text); + await client.post('com.atproto.repo.createRecord', { + input: { + repo: botDid, + collection: POST_COLLECTION, + record: { + $type: POST_COLLECTION, + text: opts.text, + createdAt: new Date().toISOString(), + reply: { root: opts.root, parent: opts.parent }, + ...(facets ? { facets } : {}) + } + } + }); +} diff --git a/apps/web/src/lib/bot/config.ts b/apps/web/src/lib/bot/config.ts new file mode 100644 index 0000000..efc72bc --- /dev/null +++ b/apps/web/src/lib/bot/config.ts @@ -0,0 +1,51 @@ +/** + * Config + copy for the `going.atmo.rsvp` reply bot. + * + * The bot lets people RSVP to an event by tagging it in a reply to a post that + * links to an `atmo.rsvp` event. It uses the replier's own OAuth token (stored + * at login) to write the RSVP record, then likes their reply to confirm — or + * replies with a hint when it can't act. + */ + +/** Public host of atmo.rsvp links we resolve events from. */ +export const ATMO_HOSTS = ['atmo.rsvp', 'www.atmo.rsvp']; + +/** Default PDS used to authenticate the bot when BOT_PDS_URL is unset. */ +export const DEFAULT_BOT_PDS = 'https://bsky.social'; + +/** Lexicon NSIDs. */ +export const RSVP_COLLECTION = 'community.lexicon.calendar.rsvp'; +export const EVENT_COLLECTION = 'community.lexicon.calendar.event'; +export const LIKE_COLLECTION = 'app.bsky.feed.like'; +export const POST_COLLECTION = 'app.bsky.feed.post'; + +/** Always RSVP as "going" — that's what the handle promises. */ +export const RSVP_GOING_STATUS = 'community.lexicon.calendar.rsvp#going'; + +/** How many notifications to pull per cron tick. */ +export const NOTIF_LIMIT = 50; +/** Cap mentions handled per tick so one busy minute can't blow the CPU budget. */ +export const MAX_MENTIONS_PER_RUN = 25; + +/** Don't re-nag the same person to sign in more than once per day. */ +export const NO_TOKEN_COOLDOWN_MS = 24 * 60 * 60 * 1000; +/** Don't repeat "no event here" to the same person more than once an hour. */ +export const NO_EVENT_COOLDOWN_MS = 60 * 60 * 1000; + +/** Reply copy. Kept short; URLs become clickable via link facets. */ +export const REPLIES = { + noToken: + "To RSVP I need you to sign in once at https://atmo.rsvp — then tag me again and I'll add you 🎟️", + noEvent: + "I couldn't find an atmo.rsvp event in this thread. Reply to a post that links to one and tag me again 👋", + external: (url: string) => `RSVPs for this event are handled on the original site: ${url}` +} as const; + +/** Outcomes recorded per processed mention (for dedup + cooldowns + auditing). */ +export type BotOutcome = + | 'liked' + | 'replied_no_token' + | 'replied_no_event' + | 'replied_external' + | 'skipped_cooldown' + | 'error'; diff --git a/apps/web/src/lib/bot/db.ts b/apps/web/src/lib/bot/db.ts new file mode 100644 index 0000000..5f36cc1 --- /dev/null +++ b/apps/web/src/lib/bot/db.ts @@ -0,0 +1,81 @@ +import type { BotOutcome } from './config'; +import { NO_EVENT_COOLDOWN_MS, NO_TOKEN_COOLDOWN_MS } from './config'; + +/** + * Records every mention the bot has acted on. Doubles as the dedup guard (so a + * mention is never processed twice across overlapping cron ticks) and the + * cooldown source (so we don't repeatedly nag the same person). + */ +let schemaReady = false; + +export async function ensureBotSchema(db: D1Database): Promise { + if (schemaReady) return; + await db.batch([ + db.prepare( + `CREATE TABLE IF NOT EXISTS rsvp_bot_processed ( + notif_uri TEXT PRIMARY KEY, + author_did TEXT NOT NULL, + outcome TEXT NOT NULL, + event_uri TEXT, + created_at INTEGER NOT NULL + )` + ), + db.prepare( + `CREATE INDEX IF NOT EXISTS idx_rsvp_bot_author + ON rsvp_bot_processed (author_did, created_at)` + ) + ]); + schemaReady = true; +} + +/** Returns the subset of `uris` already processed (so we skip them). */ +export async function getProcessed(db: D1Database, uris: string[]): Promise> { + if (uris.length === 0) return new Set(); + const placeholders = uris.map(() => '?').join(','); + const { results } = await db + .prepare(`SELECT notif_uri FROM rsvp_bot_processed WHERE notif_uri IN (${placeholders})`) + .bind(...uris) + .all<{ notif_uri: string }>(); + return new Set(results.map((r) => r.notif_uri)); +} + +export async function recordOutcome( + db: D1Database, + input: { notifUri: string; authorDid: string; outcome: BotOutcome; eventUri?: string | null } +): Promise { + await db + .prepare( + `INSERT OR REPLACE INTO rsvp_bot_processed + (notif_uri, author_did, outcome, event_uri, created_at) + VALUES (?, ?, ?, ?, ?)` + ) + .bind(input.notifUri, input.authorDid, input.outcome, input.eventUri ?? null, Date.now()) + .run(); +} + +async function repliedWithin( + db: D1Database, + authorDid: string, + outcome: BotOutcome, + windowMs: number +): Promise { + const row = await db + .prepare( + `SELECT 1 FROM rsvp_bot_processed + WHERE author_did = ? AND outcome = ? AND created_at > ? + LIMIT 1` + ) + .bind(authorDid, outcome, Date.now() - windowMs) + .first(); + return row !== null; +} + +/** True if we've already told this person to sign in within the cooldown window. */ +export function repliedNoTokenRecently(db: D1Database, authorDid: string): Promise { + return repliedWithin(db, authorDid, 'replied_no_token', NO_TOKEN_COOLDOWN_MS); +} + +/** True if we've already said "no event here" to this person within the window. */ +export function repliedNoEventRecently(db: D1Database, authorDid: string): Promise { + return repliedWithin(db, authorDid, 'replied_no_event', NO_EVENT_COOLDOWN_MS); +} diff --git a/apps/web/src/lib/bot/process-mentions.ts b/apps/web/src/lib/bot/process-mentions.ts new file mode 100644 index 0000000..56b1fdb --- /dev/null +++ b/apps/web/src/lib/bot/process-mentions.ts @@ -0,0 +1,153 @@ +import type { Did } from '@atcute/lexicons'; +import { getServerClient } from '$lib/contrail'; +import { MAX_MENTIONS_PER_RUN, NOTIF_LIMIT, REPLIES } from './config'; +import { getBotHandle, type BotHandle } from './session'; +import { + ensureBotSchema, + getProcessed, + recordOutcome, + repliedNoEventRecently, + repliedNoTokenRecently +} from './db'; +import { likePost, replyToPost, type PostRef } from './actions'; +import { findEventLink, loadEvent, resolveActorDid, type BotPostRecord } from './resolve-event'; +import { rsvpOnBehalf } from './rsvp'; + +/** The notification fields we rely on (structurally compatible with the lexicon type). */ +type Mention = { + uri: string; + cid: string; + reason: string; + author: { did: Did }; + record: unknown; +}; + +/** + * Poll the bot's mentions and act on each new one: RSVP the replier and like + * their post, or reply with a hint. Wrapped per-mention so one failure can't + * stall the batch; unrecorded failures are retried on the next tick. + */ +export async function processBotMentions(env: App.Platform['env'], db: D1Database): Promise { + const bot = await getBotHandle(env); + if (!bot) return; + + await ensureBotSchema(db); + + try { + const res = await bot.client.get('app.bsky.notification.listNotifications', { + params: { limit: NOTIF_LIMIT } + }); + if (!res.ok) { + console.warn('[bot] listNotifications failed:', res.status); + return; + } + + const mentions = (res.data.notifications ?? []).filter( + (n) => n.reason === 'mention' && n.author.did !== bot.did + ) as Mention[]; + if (mentions.length === 0) return; + + const processed = await getProcessed( + db, + mentions.map((n) => n.uri) + ); + const todo = mentions.filter((n) => !processed.has(n.uri)).slice(0, MAX_MENTIONS_PER_RUN); + + for (const mention of todo) { + try { + await handleMention(env, db, bot, mention); + } catch (e) { + // Leave unrecorded so the next tick retries this mention. + console.error('[bot] failed to handle mention', mention.uri, e); + } + } + } finally { + await bot.flush(); + } +} + +async function handleMention( + env: App.Platform['env'], + db: D1Database, + bot: BotHandle, + mention: Mention +): Promise { + const serverClient = getServerClient(db); + const record = mention.record as BotPostRecord; + const authorDid = mention.author.did; + + const parent: PostRef = { uri: mention.uri, cid: mention.cid }; + const root: PostRef = record.reply?.root ?? parent; + const reply = (text: string) => replyToPost(bot.client, bot.did, { root, parent, text }); + + const replyNoEvent = async () => { + if (await repliedNoEventRecently(db, authorDid)) { + await recordOutcome(db, { notifUri: mention.uri, authorDid, outcome: 'skipped_cooldown' }); + return; + } + await reply(REPLIES.noEvent); + await recordOutcome(db, { notifUri: mention.uri, authorDid, outcome: 'replied_no_event' }); + }; + + // 1. Find the event link (mention → parent → root). + const link = await findEventLink(bot.client, record); + if (!link) return replyNoEvent(); + + // 2. Resolve it to an indexed event (with cid + external-RSVP signal). + const did = await resolveActorDid(bot.client, link.actor); + const event = did ? await loadEvent(serverClient, did, link.rkey) : null; + if (!event) return replyNoEvent(); + + // 3. External-only events can't be RSVP'd here — point them at the source. + if (event.externalOnly) { + await reply( + event.externalRsvpUrl + ? REPLIES.external(event.externalRsvpUrl) + : 'RSVPs for this event are handled on the original site.' + ); + await recordOutcome(db, { + notifUri: mention.uri, + authorDid, + outcome: 'replied_external', + eventUri: event.uri + }); + return; + } + + // 4. RSVP on the replier's behalf. + const result = await rsvpOnBehalf(env, serverClient, authorDid, event); + + if (result === 'ok') { + await likePost(bot.client, bot.did, parent); + await recordOutcome(db, { + notifUri: mention.uri, + authorDid, + outcome: 'liked', + eventUri: event.uri + }); + return; + } + + if (result === 'no-token') { + if (await repliedNoTokenRecently(db, authorDid)) { + await recordOutcome(db, { + notifUri: mention.uri, + authorDid, + outcome: 'skipped_cooldown', + eventUri: event.uri + }); + return; + } + await reply(REPLIES.noToken); + await recordOutcome(db, { + notifUri: mention.uri, + authorDid, + outcome: 'replied_no_token', + eventUri: event.uri + }); + return; + } + + // Transient error — throw so this mention stays unrecorded and is retried. + throw new Error(`transient RSVP failure for ${mention.uri}`); +} diff --git a/apps/web/src/lib/bot/resolve-event.ts b/apps/web/src/lib/bot/resolve-event.ts new file mode 100644 index 0000000..b065626 --- /dev/null +++ b/apps/web/src/lib/bot/resolve-event.ts @@ -0,0 +1,161 @@ +import type { Client } from '@atcute/client'; +import type { Did } from '@atcute/lexicons'; +import type { Handle, ResourceUri } from '@atcute/lexicons/syntax'; +import { getEventRecordFromContrail } from '$lib/contrail'; +import { ATMO_HOSTS } from './config'; + +/** The fields of a Bluesky post record we care about for link extraction. */ +export type BotPostRecord = { + text?: string; + facets?: Array<{ features?: Array<{ $type?: string; uri?: string }> }>; + embed?: { + $type?: string; + external?: { uri?: string }; + media?: { $type?: string; external?: { uri?: string } }; + }; + reply?: { + root?: { uri: string; cid: string }; + parent?: { uri: string; cid: string }; + }; +}; + +export type AtmoEventLink = { actor: string; rkey: string }; + +/** Parse `https://atmo.rsvp/p/{actor}/e/{rkey}` into its parts (public events only). */ +export function parseAtmoEventUrl(raw: string): AtmoEventLink | null { + let url: URL; + try { + url = new URL(raw); + } catch { + return null; + } + if (!ATMO_HOSTS.includes(url.hostname.toLowerCase())) return null; + // Public event path only. Space URLs (`/p/.../e/.../s/...`) won't match the + // `$` anchor, so private-space RSVP is intentionally out of scope for v1. + const m = url.pathname.match(/^\/p\/([^/]+)\/e\/([^/]+)\/?$/); + if (!m) return null; + return { actor: decodeURIComponent(m[1]), rkey: m[2] }; +} + +const URL_RE = /https?:\/\/[^\s)\]]+/g; + +/** Pull the first atmo.rsvp event link out of a single post record. */ +function linkFromRecord(record: BotPostRecord | undefined): AtmoEventLink | null { + if (!record) return null; + + const candidates: string[] = []; + + for (const facet of record.facets ?? []) { + for (const feature of facet.features ?? []) { + if (feature.$type?.endsWith('richtext.facet#link') && feature.uri) { + candidates.push(feature.uri); + } + } + } + + const embed = record.embed; + if (embed?.$type === 'app.bsky.embed.external' && embed.external?.uri) { + candidates.push(embed.external.uri); + } + if (embed?.$type === 'app.bsky.embed.recordWithMedia' && embed.media?.external?.uri) { + candidates.push(embed.media.external.uri); + } + + if (record.text) { + for (const m of record.text.matchAll(URL_RE)) candidates.push(m[0]); + } + + for (const candidate of candidates) { + const parsed = parseAtmoEventUrl(candidate); + if (parsed) return parsed; + } + return null; +} + +/** Fetch raw post records for the given AT-URIs (best-effort). */ +async function fetchPostRecords( + client: Client, + uris: string[] +): Promise> { + const out = new Map(); + if (uris.length === 0) return out; + const res = await client.get('app.bsky.feed.getPosts', { + params: { uris: uris as ResourceUri[] } + }); + if (!res.ok) return out; + for (const post of res.data.posts ?? []) { + out.set(post.uri, post.record as BotPostRecord); + } + return out; +} + +/** + * Lenient lookup: the mention reply itself, then its parent, then the thread + * root. First atmo.rsvp event link wins. + */ +export async function findEventLink( + client: Client, + mention: BotPostRecord +): Promise { + const direct = linkFromRecord(mention); + if (direct) return direct; + + const parentUri = mention.reply?.parent?.uri; + const rootUri = mention.reply?.root?.uri; + const uris = [...new Set([parentUri, rootUri].filter((u): u is string => !!u))]; + if (uris.length === 0) return null; + + const records = await fetchPostRecords(client, uris); + // Prefer the parent over the root when both contain a link. + for (const uri of uris) { + const link = linkFromRecord(records.get(uri)); + if (link) return link; + } + return null; +} + +/** Resolve an actor (handle or DID) from an atmo.rsvp link to a DID. */ +export async function resolveActorDid(client: Client, actor: string): Promise { + if (actor.startsWith('did:')) return actor as Did; + const res = await client.get('com.atproto.identity.resolveHandle', { + params: { handle: actor as Handle } + }); + if (!res.ok) return null; + return res.data.did; +} + +export type ResolvedEvent = { + uri: string; + cid: string; + /** True when RSVPs are handled off-platform (imported event, `external_only`). */ + externalOnly: boolean; + /** The off-platform RSVP URL, when known. */ + externalRsvpUrl: string | null; +}; + +/** + * Load the event referenced by a link via contrail's in-process index, returning + * the strong-ref pieces (uri + cid) and the external-RSVP signal. Returns null + * when the event isn't indexed (treated as "no event"). + */ +export async function loadEvent( + serverClient: Client, + did: Did, + rkey: string +): Promise { + const event = await getEventRecordFromContrail(serverClient, { did, rkey }); + if (!event?.cid) return null; + + const value = event.value as unknown as { + additionalData?: { externalSource?: { rsvpMode?: string; url?: string } }; + }; + const externalSource = value?.additionalData?.externalSource; + const externalOnly = externalSource?.rsvpMode === 'external_only'; + + return { + uri: event.uri, + cid: event.cid, + externalOnly, + externalRsvpUrl: externalOnly ? (externalSource?.url ?? null) : null + }; +} diff --git a/apps/web/src/lib/bot/rsvp.ts b/apps/web/src/lib/bot/rsvp.ts new file mode 100644 index 0000000..a14636c --- /dev/null +++ b/apps/web/src/lib/bot/rsvp.ts @@ -0,0 +1,107 @@ +import { Client } from '@atcute/client'; +import type { Did } from '@atcute/lexicons'; +import type { ResourceUri } from '@atcute/lexicons/syntax'; +import { + TokenInvalidError, + TokenRefreshError, + TokenRevokedError, + AuthMethodUnsatisfiableError +} from '@atcute/oauth-node-client'; +import { createOAuthClient } from '$lib/atproto/server/oauth'; +import { getRsvpStatus, getViewerRsvpFromContrail } from '$lib/contrail'; +import { RSVP_COLLECTION, RSVP_GOING_STATUS } from './config'; +import type { ResolvedEvent } from './resolve-event'; + +/** + * - `ok`: RSVP exists/created as going — the bot should like the reply. + * - `no-token`: the user's session is gone or lacks RSVP scope — nudge them. + * - `error`: transient failure — leave unprocessed so the next tick retries. + */ +export type RsvpResult = 'ok' | 'no-token' | 'error'; + +const RSVP_NSID = RSVP_COLLECTION as `${string}.${string}.${string}`; + +/** Whether an OAuth restore error means the session is genuinely unrecoverable. */ +function isSessionGone(e: unknown): boolean { + return ( + e instanceof TokenInvalidError || + e instanceof TokenRevokedError || + e instanceof TokenRefreshError || + e instanceof AuthMethodUnsatisfiableError + ); +} + +/** + * RSVP `userDid` to `event` as "going", using their stored OAuth token. + * Idempotent: if they're already going we no-op; if they previously RSVP'd + * interested/notgoing we flip the existing record to going. + */ +export async function rsvpOnBehalf( + env: App.Platform['env'], + serverClient: Client, + userDid: Did, + event: ResolvedEvent +): Promise { + // No stored session at all ⇒ they've never signed in ⇒ nudge them (don't retry). + if ((await env.OAUTH_SESSIONS.get(userDid)) === null) return 'no-token'; + + let session; + try { + session = await createOAuthClient(env).restore(userDid); + } catch (e) { + if (isSessionGone(e)) return 'no-token'; + console.error('[bot] session restore failed (transient):', e); + return 'error'; + } + + const client = new Client({ handler: session }); + + // Look up any existing RSVP so we don't duplicate, and can flip status. + let existing = null; + try { + existing = await getViewerRsvpFromContrail(serverClient, { + eventUri: event.uri, + actor: userDid + }); + } catch { + // Index miss is fine — we'll just create a fresh record. + } + + if (existing && getRsvpStatus(existing.value?.status) === 'going') return 'ok'; + + const record = { + $type: RSVP_COLLECTION, + createdAt: new Date().toISOString(), + status: RSVP_GOING_STATUS, + subject: { uri: event.uri, cid: event.cid } + }; + + const response = existing + ? await client.post('com.atproto.repo.putRecord', { + input: { repo: userDid, collection: RSVP_NSID, rkey: existing.rkey, record } + }) + : await client.post('com.atproto.repo.createRecord', { + input: { repo: userDid, collection: RSVP_NSID, record } + }); + + if (!response.ok) { + // 401/403 ⇒ token revoked or missing the RSVP write scope ⇒ ask them to re-auth. + if (response.status === 401 || response.status === 403) return 'no-token'; + console.error('[bot] RSVP write failed:', response.status, response.data); + return 'error'; + } + + const createdUri = (response.data as { uri?: string })?.uri; + if (createdUri) { + // Best-effort: nudge contrail to re-index now instead of waiting for the firehose. + try { + await serverClient.post('rsvp.atmo.notifyOfUpdate', { + input: { uris: [createdUri as ResourceUri] } + }); + } catch { + // Harmless — the next ingest tick will pick it up. + } + } + + return 'ok'; +} diff --git a/apps/web/src/lib/bot/session.ts b/apps/web/src/lib/bot/session.ts new file mode 100644 index 0000000..dc15061 --- /dev/null +++ b/apps/web/src/lib/bot/session.ts @@ -0,0 +1,81 @@ +import { Client, CredentialManager, type AtpSessionData } from '@atcute/client'; +import type { Did } from '@atcute/lexicons'; +import { DEFAULT_BOT_PDS } from './config'; + +/** + * Authenticated handle for the bot account. `client` makes API calls; + * `flush()` persists the (possibly refreshed) session back to KV — call it once + * at the end of a run so rotated tokens survive to the next tick. + */ +export type BotHandle = { + client: Client; + did: Did; + flush: () => Promise; +}; + +/** KV key for the bot's persisted credential session (namespaced away from DIDs). */ +const sessionKey = (identifier: string) => `bot:session:${identifier}`; + +/** + * Log the bot in (or resume a persisted session) using an app password. + * + * Reuses the OAUTH_SESSIONS KV namespace for storage. Returns `null` when the + * bot isn't configured (no identifier/password) so the cron can no-op cleanly + * in dev or before the secret is set. + */ +export async function getBotHandle(env: App.Platform['env']): Promise { + const identifier = env.BOT_IDENTIFIER; + const password = env.BOT_APP_PASSWORD; + if (!identifier || !password) { + console.warn('[bot] BOT_IDENTIFIER / BOT_APP_PASSWORD not set — skipping bot run'); + return null; + } + + const kv = env.OAUTH_SESSIONS; + const key = sessionKey(identifier); + + const manager = new CredentialManager({ + service: env.BOT_PDS_URL || DEFAULT_BOT_PDS + }); + + const saved = await loadSession(kv, key); + let session: AtpSessionData | undefined; + if (saved) { + try { + session = await manager.resume(saved); + } catch (e) { + console.warn('[bot] resume failed, logging in fresh:', e); + } + } + if (!session) { + session = await manager.login({ identifier, password }); + } + + // Persist immediately so a fresh login is reusable even if the run later throws. + await saveSession(kv, key, manager.session); + + return { + client: new Client({ handler: manager }), + did: session.did, + flush: () => saveSession(kv, key, manager.session) + }; +} + +async function loadSession(kv: KVNamespace, key: string): Promise { + const raw = await kv.get(key, 'text'); + if (!raw) return undefined; + try { + return JSON.parse(raw) as AtpSessionData; + } catch { + return undefined; + } +} + +async function saveSession( + kv: KVNamespace, + key: string, + session: AtpSessionData | undefined +): Promise { + if (!session) return; + await kv.put(key, JSON.stringify(session)); +} diff --git a/apps/web/src/routes/api/cron/+server.ts b/apps/web/src/routes/api/cron/+server.ts index e80801f..4bae09f 100644 --- a/apps/web/src/routes/api/cron/+server.ts +++ b/apps/web/src/routes/api/cron/+server.ts @@ -1,4 +1,5 @@ import { contrail, ensureInit } from '$lib/contrail/index'; +import { processBotMentions } from '$lib/bot/process-mentions'; import type { RequestHandler } from './$types'; export const POST: RequestHandler = async ({ request, platform }) => { @@ -11,5 +12,12 @@ export const POST: RequestHandler = async ({ request, platform }) => { await ensureInit(db); await contrail.ingest({}, db); + // Reply bot — isolated so its failures never break firehose ingest. + try { + await processBotMentions(platform!.env, db); + } catch (e) { + console.error('[bot] processBotMentions failed:', e); + } + return new Response('OK'); }; diff --git a/apps/web/wrangler.jsonc b/apps/web/wrangler.jsonc index 4f79ebd..017b450 100644 --- a/apps/web/wrangler.jsonc +++ b/apps/web/wrangler.jsonc @@ -16,7 +16,10 @@ } }, "vars": { - "OAUTH_PUBLIC_URL": "https://atmo.rsvp" + "OAUTH_PUBLIC_URL": "https://atmo.rsvp", + // Reply bot. BOT_APP_PASSWORD is a secret: `wrangler secret put BOT_APP_PASSWORD`. + "BOT_IDENTIFIER": "going.atmo.rsvp", + "BOT_PDS_URL": "https://eurosky.social" }, "d1_databases": [ { @@ -39,4 +42,4 @@ "id": "46cfb5c0bb8c41378e757b46afd7dd35" } ] -} \ No newline at end of file +}