diff --git a/apps/web/package.json b/apps/web/package.json index ec35fd8..a5e8d73 100644 --- a/apps/web/package.json +++ b/apps/web/package.json @@ -17,6 +17,7 @@ "format": "prettier --write .", "lint": "prettier --check . && eslint .", "env:generate-key": "npx tsx src/lib/atproto/scripts/generate-key.ts", + "env:generate-notify-key": "npx tsx src/lib/atproto/scripts/generate-notify-key.ts", "env:generate-secret": "npx tsx src/lib/atproto/scripts/generate-secret.ts", "env:setup-dev": "npx tsx src/lib/atproto/scripts/setup-dev.ts", "tunnel": "npx tsx src/lib/atproto/scripts/tunnel.ts", @@ -26,11 +27,13 @@ "@atcute/atproto": "^3.1.10", "@atcute/bluesky": "^3.3.0", "@atcute/client": "^4.2.1", + "@atcute/crypto": "^2.4.1", "@atcute/identity-resolver": "^1.2.2", "@atcute/lex-cli": "^2.5.3", "@atcute/lexicons": "^1.2.9", "@atcute/oauth-node-client": "^1.1.0", "@atcute/tid": "^1.1.2", + "@atcute/xrpc-server": "^0.1.12", "@atmo-dev/contrail-lexicons": "^0.4.4", "@cloudflare/workers-types": "^4.20260317.1", "@eslint/compat": "^2.0.3", diff --git a/apps/web/src/app.d.ts b/apps/web/src/app.d.ts index 40a16ec..82f04b8 100644 --- a/apps/web/src/app.d.ts +++ b/apps/web/src/app.d.ts @@ -78,6 +78,10 @@ declare global { BOT_APP_PASSWORD?: string; /** PDS the bot authenticates against; defaults to bsky.social. */ BOT_PDS_URL?: string; + /** P-256 private key (multikey) used to sign atmo.pub app tokens and + * publish the `did:web` DID document. Set via `wrangler secret put`. + * When unset, notifications are disabled (the feature no-ops). */ + ATMO_NOTIFY_PRIVATE_KEY?: 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/hooks.server.ts b/apps/web/src/hooks.server.ts index 98b4fae..b0f8de4 100644 --- a/apps/web/src/hooks.server.ts +++ b/apps/web/src/hooks.server.ts @@ -1,11 +1,20 @@ import type { Handle } from '@sveltejs/kit'; import { restoreSession } from '$lib/atproto/server/session'; +import { buildDidDocument } from '$lib/notify/did-document'; export const handle: Handle = async ({ event, resolve }) => { - const { session, client, did } = await restoreSession( - event.cookies, - event.platform?.env - ); + // Publish our did:web document (served via a hook because SvelteKit ignores + // route directories that start with a dot). Public endpoint — no session. + if (event.url.pathname === '/.well-known/did.json' && event.platform?.env) { + const doc = await buildDidDocument(event.platform.env); + if (doc) { + return new Response(JSON.stringify(doc, null, 2), { + headers: { 'content-type': 'application/json', 'cache-control': 'public, max-age=300' } + }); + } + } + + const { session, client, did } = await restoreSession(event.cookies, event.platform?.env); event.locals.session = session; event.locals.client = client; diff --git a/apps/web/src/lib/atproto/scripts/generate-notify-key.ts b/apps/web/src/lib/atproto/scripts/generate-notify-key.ts new file mode 100644 index 0000000..2143c6b --- /dev/null +++ b/apps/web/src/lib/atproto/scripts/generate-notify-key.ts @@ -0,0 +1,21 @@ +import { P256PrivateKeyExportable } from '@atcute/crypto'; + +// Generates the P-256 signing key used for atmo.pub notifications. The private +// multikey is the server secret (ATMO_NOTIFY_PRIVATE_KEY); the public key is +// derived from it automatically when serving /.well-known/did.json. + +const keypair = await P256PrivateKeyExportable.createKeypair(); +const privateMultikey = await keypair.exportPrivateKey('multikey'); +const publicMultikey = await keypair.exportPublicKey('multikey'); +const didKey = await keypair.exportPublicKey('did'); + +console.log('\natmo.pub notification signing key (P-256)\n'); +console.log('Private key — set as a secret, keep it server-side only:'); +console.log(` ATMO_NOTIFY_PRIVATE_KEY=${privateMultikey}\n`); +console.log('Public key multibase (reference — the DID doc derives this automatically):'); +console.log(` ${publicMultikey}`); +console.log(`did:key (reference): ${didKey}\n`); +console.log('Next steps:'); +console.log(' • prod: wrangler secret put ATMO_NOTIFY_PRIVATE_KEY'); +console.log(' • dev: add ATMO_NOTIFY_PRIVATE_KEY=… to apps/web/.env'); +console.log(' • verify https:///.well-known/did.json resolves with your key\n'); diff --git a/apps/web/src/lib/atproto/settings.ts b/apps/web/src/lib/atproto/settings.ts index 20961a2..a185415 100644 --- a/apps/web/src/lib/atproto/settings.ts +++ b/apps/web/src/lib/atproto/settings.ts @@ -24,7 +24,12 @@ export const scopes = [ scope.repo({ collection: [...collections] }), scope.blob({ accept: ['image/*'] }), 'include:rsvp.atmo.permissionSet', - 'include:app.bsky.authCreatePosts' + 'include:app.bsky.authCreatePosts', + // atmo.pub notifications: lets us mint user tokens (via + // com.atproto.server.getServiceAuth) to ask for/revoke notification consent. + // `send` itself needs no user scope (it's signed with our app key). + 'rpc?lxm=pub.atmo.notify.requestPermission&aud=*', + 'rpc?lxm=pub.atmo.notify.revokeSelf&aud=*' ]; // set to false to disable signup diff --git a/apps/web/src/lib/components/RecentActivity.svelte b/apps/web/src/lib/components/RecentActivity.svelte index 1a41120..cba08da 100644 --- a/apps/web/src/lib/components/RecentActivity.svelte +++ b/apps/web/src/lib/components/RecentActivity.svelte @@ -5,13 +5,25 @@ let { activities }: { activities: ActivityCluster[] } = $props(); + function dedupeByDid(attendees: AttendeeInfo[]): AttendeeInfo[] { + const seen = new Set(); + return attendees.filter((a) => { + if (seen.has(a.did)) return false; + seen.add(a.did); + return true; + }); + } + function visibleAttendees(cluster: ActivityCluster): { shown: AttendeeInfo[]; status: 'going' | 'interested'; } { - const going = cluster.attendees.filter((a) => a.status === 'going'); + // Dedupe by DID: the same actor can surface from multiple sources when a + // cluster is assembled, which would otherwise produce duplicate `{#each}` + // keys (each_key_duplicate) and inflate counts. + const going = dedupeByDid(cluster.attendees.filter((a) => a.status === 'going')); if (going.length > 0) return { shown: going, status: 'going' }; - return { shown: cluster.attendees, status: 'interested' }; + return { shown: dedupeByDid(cluster.attendees), status: 'interested' }; } function namesSentence(attendees: AttendeeInfo[]): string { @@ -60,9 +72,7 @@ class="hover:bg-base-100 dark:hover:bg-base-900/50 -mx-2 block rounded-lg px-2 py-3 transition-colors" >
-

+

{eventTitle}

diff --git a/apps/web/src/lib/notify/config.ts b/apps/web/src/lib/notify/config.ts new file mode 100644 index 0000000..5c93a70 --- /dev/null +++ b/apps/web/src/lib/notify/config.ts @@ -0,0 +1,43 @@ +// atmo.pub notification relay integration — shared constants. +// Docs: https://docs.atmo.pub/llms.txt + +/** Relay XRPC base + DID (the `aud` for every app/user service-auth token). */ +export const RELAY_ORIGIN = 'https://relay.atmo.pub'; +export const RELAY_DID = 'did:web:relay.atmo.pub'; + +// Lexicon method names (also used as the JWT `lxm` claim). +export const LXM_SEND = 'pub.atmo.notify.send'; +export const LXM_REQUEST_PERMISSION = 'pub.atmo.notify.requestPermission'; +export const LXM_REVOKE_SELF = 'pub.atmo.notify.revokeSelf'; +export const LXM_SUBSCRIBER_CHANGED = 'pub.atmo.notify.subscriberChanged'; + +/** Shown on the atmo.pub approval screen for our sender DID. */ +export const APP_NAME = 'atmo.rsvp'; +export const APP_DESCRIPTION = + 'Reminders for events you RSVP to, plus a ping when someone RSVPs to your events.'; + +/** RSVP statuses that count as "attending" for reminders and host alerts. + * Matched as a suffix on `community.lexicon.calendar.rsvp#`. */ +export const ATTENDING_SUFFIXES = ['#going', '#interested'] as const; + +// Reminder milestones, measured as ms before the event's `startsAt`. +export const REMINDER_24H_MS = 24 * 60 * 60 * 1000; +export const REMINDER_1H_MS = 60 * 60 * 1000; + +/** How long after a milestone passes we'll still fire it (tolerates worker + * downtime / ingest backlog). Kept small so a late "24h" reminder never lands + * hours off; the per-(recipient,event,milestone) dedup key prevents repeats. */ +export const REMINDER_CATCHUP_MS = 30 * 60 * 1000; + +/** Caps per cron tick — bounds D1 CPU per cycle and the relay's 1/sec-per-pair + * rate limit. Leftover work is picked up on subsequent ticks (every minute). */ +export const MAX_SENDS_PER_TICK = 40; +export const MAX_NEW_RSVPS_PER_TICK = 100; + +export type ReminderKind = 'r24' | 'r1' | 'r0'; + +export const REMINDER_MILESTONES: { kind: ReminderKind; offsetMs: number }[] = [ + { kind: 'r24', offsetMs: REMINDER_24H_MS }, + { kind: 'r1', offsetMs: REMINDER_1H_MS }, + { kind: 'r0', offsetMs: 0 } +]; diff --git a/apps/web/src/lib/notify/db.ts b/apps/web/src/lib/notify/db.ts new file mode 100644 index 0000000..a5e0ceb --- /dev/null +++ b/apps/web/src/lib/notify/db.ts @@ -0,0 +1,124 @@ +// D1 tables backing atmo.pub notifications. Created lazily on first use (same +// pattern as the reply bot's schema in $lib/bot/db.ts), so no migration step. + +let schemaReady = false; + +export async function ensureNotifySchema(db: D1Database): Promise { + if (schemaReady) return; + await db.batch([ + // Who has granted us permission to notify them. Kept authoritative by the + // `subscriberChanged` webhook + the in-app enable/disable flow. + db.prepare( + `CREATE TABLE IF NOT EXISTS notify_subscribers ( + did TEXT PRIMARY KEY, + enabled INTEGER NOT NULL DEFAULT 1, + changed_at INTEGER NOT NULL + )` + ), + // Dedup ledger: one row per notification we've already sent (or claimed). + // Prevents duplicates across overlapping cron ticks. + db.prepare( + `CREATE TABLE IF NOT EXISTS notify_sent ( + key TEXT PRIMARY KEY, + recipient TEXT NOT NULL, + kind TEXT NOT NULL, + created_at INTEGER NOT NULL + )` + ), + db.prepare(`CREATE INDEX IF NOT EXISTS idx_notify_sent_created ON notify_sent (created_at)`), + // Single-row cursor tracking the last RSVP we've considered for host + // alerts (records_rsvp.indexed_at is in microseconds). + db.prepare( + `CREATE TABLE IF NOT EXISTS notify_state ( + id INTEGER PRIMARY KEY CHECK (id = 1), + rsvp_cursor_us INTEGER NOT NULL + )` + ) + ]); + schemaReady = true; +} + +export async function setSubscriber(db: D1Database, did: string, enabled: boolean): Promise { + await db + .prepare( + `INSERT INTO notify_subscribers (did, enabled, changed_at) + VALUES (?, ?, ?) + ON CONFLICT(did) DO UPDATE SET enabled = excluded.enabled, changed_at = excluded.changed_at` + ) + .bind(did, enabled ? 1 : 0, Date.now()) + .run(); +} + +export async function getSubscriber( + db: D1Database, + did: string +): Promise<{ enabled: boolean } | null> { + const row = await db + .prepare(`SELECT enabled FROM notify_subscribers WHERE did = ?`) + .bind(did) + .first<{ enabled: number }>(); + return row ? { enabled: row.enabled === 1 } : null; +} + +export async function countEnabledSubscribers(db: D1Database): Promise { + const row = await db + .prepare(`SELECT COUNT(*) AS n FROM notify_subscribers WHERE enabled = 1`) + .first<{ n: number }>(); + return Number(row?.n ?? 0); +} + +/** + * Atomically claim a notification key. Returns true if this caller won the claim + * (and should send), false if it was already claimed. Insert-or-ignore makes + * this safe under concurrent/overlapping ticks. + */ +export async function claimNotification( + db: D1Database, + key: string, + recipient: string, + kind: string +): Promise { + const res = await db + .prepare( + `INSERT OR IGNORE INTO notify_sent (key, recipient, kind, created_at) VALUES (?, ?, ?, ?)` + ) + .bind(key, recipient, kind, Date.now()) + .run(); + return (res.meta?.changes ?? 0) > 0; +} + +/** Release a previously-claimed key so it can be retried (e.g. after a 429). */ +export async function releaseNotification(db: D1Database, key: string): Promise { + await db.prepare(`DELETE FROM notify_sent WHERE key = ?`).bind(key).run(); +} + +export async function getRsvpCursor(db: D1Database): Promise { + const row = await db + .prepare(`SELECT rsvp_cursor_us FROM notify_state WHERE id = 1`) + .first<{ rsvp_cursor_us: number }>(); + return row ? Number(row.rsvp_cursor_us) : null; +} + +export async function setRsvpCursor(db: D1Database, us: number): Promise { + await db + .prepare( + `INSERT INTO notify_state (id, rsvp_cursor_us) VALUES (1, ?) + ON CONFLICT(id) DO UPDATE SET rsvp_cursor_us = excluded.rsvp_cursor_us` + ) + .bind(us) + .run(); +} + +/** Resolve DIDs to handles from contrail's `identities` index (best-effort). */ +export async function getHandles(db: D1Database, dids: string[]): Promise> { + const unique = [...new Set(dids)]; + if (unique.length === 0) return new Map(); + const placeholders = unique.map(() => '?').join(','); + const { results } = await db + .prepare(`SELECT did, handle FROM identities WHERE did IN (${placeholders})`) + .bind(...unique) + .all<{ did: string; handle: string | null }>(); + const map = new Map(); + for (const r of results) if (r.handle) map.set(r.did, r.handle); + return map; +} diff --git a/apps/web/src/lib/notify/did-document.ts b/apps/web/src/lib/notify/did-document.ts new file mode 100644 index 0000000..8673d8f --- /dev/null +++ b/apps/web/src/lib/notify/did-document.ts @@ -0,0 +1,35 @@ +import { getPublicMultikey, getSenderDid, notifyConfigured } from './keypair'; +import { SERVICE_URL } from '../spaces/config'; + +type Env = App.Platform['env']; + +/** + * The `did:web` document served at `/.well-known/did.json`, so the atmo + * relay can resolve our DID and verify the app tokens we sign. Returns null when + * notifications aren't configured (no signing key) — the caller should 404. + */ +export async function buildDidDocument(env: Env): Promise | null> { + if (!notifyConfigured(env)) return null; + + const did = getSenderDid(env); + const origin = env.OAUTH_PUBLIC_URL; + const publicKeyMultibase = await getPublicMultikey(env); + + const service: Record[] = [ + { id: '#atmo_notify', type: 'AtmoNotifsSender', serviceEndpoint: origin } + ]; + // Keep the spaces service entry when this origin is also the spaces host + // (e.g. a dev tunnel), so publishing the DID doc here doesn't drop it. + if (SERVICE_URL && new URL(SERVICE_URL).host === new URL(origin).host) { + service.push({ id: '#event_space', type: 'AtmoSpaceService', serviceEndpoint: SERVICE_URL }); + } + + return { + '@context': ['https://www.w3.org/ns/did/v1', 'https://w3id.org/security/multikey/v1'], + id: did, + verificationMethod: [ + { id: `${did}#atproto`, type: 'Multikey', controller: did, publicKeyMultibase } + ], + service + }; +} diff --git a/apps/web/src/lib/notify/keypair.ts b/apps/web/src/lib/notify/keypair.ts new file mode 100644 index 0000000..4498775 --- /dev/null +++ b/apps/web/src/lib/notify/keypair.ts @@ -0,0 +1,39 @@ +import { P256PrivateKey, parsePrivateMultikey } from '@atcute/crypto'; +import type { Did } from '@atcute/lexicons'; + +type Env = App.Platform['env']; + +/** True when the deployment is configured to send/receive atmo.pub notifications. */ +export function notifyConfigured(env: Env): boolean { + return !!env.ATMO_NOTIFY_PRIVATE_KEY && !!env.OAUTH_PUBLIC_URL; +} + +/** Our app's `did:web`, derived from the public origin so it always matches the + * DID document served at `/.well-known/did.json`. */ +export function getSenderDid(env: Env): Did { + if (!env.OAUTH_PUBLIC_URL) throw new Error('OAUTH_PUBLIC_URL is not set'); + return `did:web:${new URL(env.OAUTH_PUBLIC_URL).host}` as Did; +} + +// The imported keypair is cached per process (re-imported only if the secret +// rotates). importRaw is async (WebCrypto), so callers await getKeypair(). +let cached: { multikey: string; keypair: P256PrivateKey } | null = null; + +export async function getKeypair(env: Env): Promise { + const multikey = env.ATMO_NOTIFY_PRIVATE_KEY; + if (!multikey) throw new Error('ATMO_NOTIFY_PRIVATE_KEY secret is not set'); + if (cached?.multikey === multikey) return cached.keypair; + + const parsed = parsePrivateMultikey(multikey); + if (parsed.type !== 'p256') { + throw new Error(`ATMO_NOTIFY_PRIVATE_KEY must be a P-256 multikey, got ${parsed.type}`); + } + const keypair = await P256PrivateKey.importRaw(parsed.privateKeyBytes); + cached = { multikey, keypair }; + return keypair; +} + +/** Public key as a multibase string for the DID document's `publicKeyMultibase`. */ +export async function getPublicMultikey(env: Env): Promise { + return (await getKeypair(env)).exportPublicKey('multikey'); +} diff --git a/apps/web/src/lib/notify/process.ts b/apps/web/src/lib/notify/process.ts new file mode 100644 index 0000000..b65d6bb --- /dev/null +++ b/apps/web/src/lib/notify/process.ts @@ -0,0 +1,261 @@ +import { + MAX_NEW_RSVPS_PER_TICK, + MAX_SENDS_PER_TICK, + REMINDER_CATCHUP_MS, + REMINDER_24H_MS, + REMINDER_MILESTONES, + type ReminderKind +} from './config'; +import { + claimNotification, + countEnabledSubscribers, + ensureNotifySchema, + getHandles, + getRsvpCursor, + releaseNotification, + setRsvpCursor, + setSubscriber +} from './db'; +import { notifyConfigured } from './keypair'; +import { RelayError, relaySend, type SendPayload } from './relay'; + +type Env = App.Platform['env']; + +// json_extract($.status) is the full token, e.g. +// `community.lexicon.calendar.rsvp#going`. Match the going/interested suffixes. +const ATTENDING_SQL = `( + json_extract(r.record,'$.status') LIKE '%#going' + OR json_extract(r.record,'$.status') LIKE '%#interested' +)`; + +/** Entry point, called from the cron after firehose ingest. No-ops cleanly when + * notifications aren't configured (dev, or before the signing key is set). */ +export async function runNotifications(env: Env, db: D1Database): Promise { + if (!notifyConfigured(env)) return; + await ensureNotifySchema(db); + + let budget = MAX_SENDS_PER_TICK; + // New-RSVP alerts must advance the cursor even when there are no subscribers, + // so the backlog never grows; reminders are a pure no-op without subscribers. + if ((await countEnabledSubscribers(db)) > 0) { + budget = await processReminders(env, db, budget); + } + await processNewRsvps(env, db, budget); +} + +type SendOutcome = 'sent' | 'ratelimited' | 'nogrant' | 'error'; + +/** Send, mapping relay errors to an outcome the callers act on. A 403 means the + * recipient has no active grant (revoked outside our webhook) → mark them + * disabled so we stop trying. */ +async function trySend(env: Env, db: D1Database, payload: SendPayload): Promise { + try { + await relaySend(env, payload); + return 'sent'; + } catch (e) { + if (e instanceof RelayError) { + if (e.status === 429) return 'ratelimited'; + if (e.status === 403) { + await setSubscriber(db, payload.recipient, false); + return 'nogrant'; + } + } + console.error('[notify] send failed:', e); + return 'error'; + } +} + +function eventWebUrl(env: Env, eventUri: string): string | undefined { + // at:///community.lexicon.calendar.event/ + const m = eventUri.match(/^at:\/\/([^/]+)\/[^/]+\/([^/]+)$/); + if (!m) return undefined; + return `${env.OAUTH_PUBLIC_URL}/p/${m[1]}/e/${m[2]}`; +} + +function reminderCopy(kind: ReminderKind, name: string): { title: string; body: string } { + const ev = name || 'Your event'; + switch (kind) { + case 'r24': + return { title: 'Event tomorrow', body: `${ev} starts in 24 hours` }; + case 'r1': + return { title: 'Event soon', body: `${ev} starts in 1 hour` }; + case 'r0': + return { title: 'Starting now', body: `${ev} is starting now` }; + } +} + +type ReminderRow = { + recipient: string; + event_uri: string; + host: string; + name: string | null; + starts_at: string | null; +}; + +/** + * Reminders (24h / 1h / at start) for events a subscriber is going to or + * interested in. We fetch each subscriber's upcoming attending events (joined in + * SQL so non-subscribers never enter the set), then fire whichever milestones + * are currently due. + */ +async function processReminders(env: Env, db: D1Database, budget: number): Promise { + if (budget <= 0) return budget; + const now = Date.now(); + + // Candidate window on startsAt as ISO strings. Padded by 36h on each side so + // non-UTC offsets (max ±14h) can't slip a relevant event past the string + // comparison; precise filtering happens in JS via Date.parse below. + const PAD = 36 * 60 * 60 * 1000; + const lo = new Date(now - REMINDER_CATCHUP_MS - PAD).toISOString(); + const hi = new Date(now + REMINDER_24H_MS + PAD).toISOString(); + + const { results } = await db + .prepare( + `SELECT r.did AS recipient, + e.uri AS event_uri, + e.did AS host, + json_extract(e.record,'$.name') AS name, + json_extract(e.record,'$.startsAt') AS starts_at + FROM records_rsvp r + JOIN records_event e ON json_extract(r.record,'$.subject.uri') = e.uri + JOIN notify_subscribers s ON s.did = r.did AND s.enabled = 1 + WHERE json_extract(e.record,'$.startsAt') >= ? + AND json_extract(e.record,'$.startsAt') <= ? + AND ${ATTENDING_SQL}` + ) + .bind(lo, hi) + .all(); + + const hostHandles = await getHandles( + db, + results.map((r) => r.host) + ); + + for (const row of results) { + if (!row.starts_at) continue; + const startMs = Date.parse(row.starts_at); + if (!Number.isFinite(startMs)) continue; + const url = eventWebUrl(env, row.event_uri); + const hostHandle = hostHandles.get(row.host); + + for (const m of REMINDER_MILESTONES) { + if (budget <= 0) return budget; + const fireAt = startMs - m.offsetMs; + // Due if we're within [fireAt, fireAt + catchup]. Past events (startMs + // far behind) and far-future ones fall outside every milestone window. + if (now < fireAt || now > fireAt + REMINDER_CATCHUP_MS) continue; + + const key = `${m.kind}:${row.recipient}:${row.event_uri}`; + if (!(await claimNotification(db, key, row.recipient, m.kind))) continue; + + const { title, body } = reminderCopy(m.kind, row.name ?? ''); + const outcome = await trySend(env, db, { + recipient: row.recipient, + title, + body, + uri: url, + category: 'reminder', + categoryDescription: 'Event reminders', + threadKey: row.event_uri, + actors: hostHandle ? [hostHandle] : undefined + }); + + if (outcome === 'ratelimited') { + await releaseNotification(db, key); + return 0; // stop this tick; pick up next minute + } + budget--; + } + } + + return budget; +} + +type NewRsvpRow = { + rsvp_uri: string; + rsvp_author: string; + indexed_at: number; + status: string | null; + event_uri: string; + host: string; + name: string | null; +}; + +/** + * Host alerts: notify an event's owner when someone (else) RSVPs going/interested + * to it. Walks RSVPs indexed since our cursor; only events whose owner is an + * enabled subscriber are joined in. The cursor advances every tick so the + * backlog stays bounded; on first run it's seeded to "now" to avoid replaying + * history. + */ +async function processNewRsvps(env: Env, db: D1Database, budget: number): Promise { + let cursor = await getRsvpCursor(db); + if (cursor === null) { + const row = await db + .prepare(`SELECT MAX(indexed_at) AS m FROM records_rsvp`) + .first<{ m: number | null }>(); + cursor = Number(row?.m) || Date.now() * 1000; + await setRsvpCursor(db, cursor); + return; // baseline only — don't alert on pre-existing RSVPs + } + if (budget <= 0) return; + + const { results } = await db + .prepare( + `SELECT r.uri AS rsvp_uri, + r.did AS rsvp_author, + r.indexed_at AS indexed_at, + json_extract(r.record,'$.status') AS status, + e.uri AS event_uri, + e.did AS host, + json_extract(e.record,'$.name') AS name + FROM records_rsvp r + JOIN records_event e ON json_extract(r.record,'$.subject.uri') = e.uri + JOIN notify_subscribers s ON s.did = e.did AND s.enabled = 1 + WHERE r.indexed_at > ? + AND r.did <> e.did + AND ${ATTENDING_SQL} + ORDER BY r.indexed_at ASC + LIMIT ?` + ) + .bind(cursor, MAX_NEW_RSVPS_PER_TICK) + .all(); + + if (results.length === 0) return; + + const authorHandles = await getHandles( + db, + results.map((r) => r.rsvp_author) + ); + + let advanceTo = cursor; + for (const row of results) { + if (budget <= 0) break; // leave cursor before this row; retry next tick + + const key = `rsvp:${row.rsvp_uri}`; + if (await claimNotification(db, key, row.host, 'rsvp')) { + const who = authorHandles.get(row.rsvp_author) ?? row.rsvp_author; + const verb = row.status?.endsWith('#interested') ? 'is interested in' : 'is going to'; + const outcome = await trySend(env, db, { + recipient: row.host, + title: 'New RSVP', + body: `${who} ${verb} ${row.name || 'your event'}`, + uri: eventWebUrl(env, row.event_uri), + category: 'rsvp', + categoryDescription: 'RSVPs to your events', + threadKey: row.event_uri, + actors: [row.rsvp_author] + }); + + if (outcome === 'ratelimited') { + await releaseNotification(db, key); + break; // leave cursor before this row + } + budget--; + } + + advanceTo = Number(row.indexed_at); + } + + if (advanceTo > cursor) await setRsvpCursor(db, advanceTo); +} diff --git a/apps/web/src/lib/notify/relay.ts b/apps/web/src/lib/notify/relay.ts new file mode 100644 index 0000000..6e1ea15 --- /dev/null +++ b/apps/web/src/lib/notify/relay.ts @@ -0,0 +1,97 @@ +import { createServiceJwt } from '@atcute/xrpc-server/auth'; +import type { Did, Nsid } from '@atcute/lexicons'; +import { + APP_DESCRIPTION, + APP_NAME, + LXM_REQUEST_PERMISSION, + LXM_REVOKE_SELF, + LXM_SEND, + RELAY_DID, + RELAY_ORIGIN +} from './config'; +import { getKeypair, getSenderDid } from './keypair'; + +type Env = App.Platform['env']; + +/** Thrown when the relay returns a non-2xx. `status` drives caller behaviour: + * 403 = no active grant, 429 = rate limited (back off), 401 = bad token. */ +export class RelayError extends Error { + constructor( + readonly status: number, + readonly lxm: string, + readonly data: unknown + ) { + super(`relay ${lxm} -> ${status}`); + this.name = 'RelayError'; + } +} + +/** Mint a short-lived app token (atproto service-auth JWT, ES256) signed with + * our server-side key. Tokens are method-scoped (`lxm`) — mint one per call. */ +export async function mintAppToken(env: Env, lxm: string): Promise { + return createServiceJwt({ + keypair: await getKeypair(env), + issuer: getSenderDid(env), + audience: RELAY_DID as Did, + lxm: lxm as Nsid, + expiresIn: 60 + }); +} + +async function relayCall>( + lxm: string, + bearer: string, + body: object +): Promise { + const res = await fetch(`${RELAY_ORIGIN}/xrpc/${lxm}`, { + method: 'POST', + headers: { authorization: `Bearer ${bearer}`, 'content-type': 'application/json' }, + body: JSON.stringify(body) + }); + const data = res.headers.get('content-type')?.includes('json') + ? ((await res.json().catch(() => ({}))) as T) + : ({} as T); + if (!res.ok) throw new RelayError(res.status, lxm, data); + return data; +} + +export type SendPayload = { + recipient: string; + title: string; + body: string; + uri?: string; + category?: string; + categoryDescription?: string; + threadKey?: string; + actors?: string[]; +}; + +/** Deliver a notification. App token only — requires an active grant from the + * recipient. `delivered: 0` is success-with-no-channels, not an error. */ +export async function relaySend( + env: Env, + payload: SendPayload +): Promise<{ id: string; delivered: number }> { + return relayCall(LXM_SEND, await mintAppToken(env, LXM_SEND), payload); +} + +/** Ask a user to allow our app to notify them. Uses the user's token (minted by + * the caller from their OAuth session via `com.atproto.server.getServiceAuth`). */ +export async function relayRequestPermission( + env: Env, + userToken: string +): Promise<{ id: string; status: 'pending' | 'alreadyGranted' }> { + return relayCall(LXM_REQUEST_PERMISSION, userToken, { + senderDid: getSenderDid(env), + title: APP_NAME, + description: APP_DESCRIPTION, + iconUrl: `${env.OAUTH_PUBLIC_URL}/og.png` + }); +} + +/** Remove our app's grant for the user. Dual-auth: app token in the header, + * fresh user token in the body. */ +export async function relayRevokeSelf(env: Env, userToken: string): Promise<{ ok: boolean }> { + const appToken = await mintAppToken(env, LXM_REVOKE_SELF); + return relayCall(LXM_REVOKE_SELF, appToken, { userToken }); +} diff --git a/apps/web/src/lib/notify/settings.remote.ts b/apps/web/src/lib/notify/settings.remote.ts new file mode 100644 index 0000000..e4fab39 --- /dev/null +++ b/apps/web/src/lib/notify/settings.remote.ts @@ -0,0 +1,63 @@ +import { error } from '@sveltejs/kit'; +import { command, getRequestEvent } from '$app/server'; +import { Client } from '@atcute/client'; +import type { Did, Nsid } from '@atcute/lexicons'; +import { LXM_REQUEST_PERMISSION, LXM_REVOKE_SELF, RELAY_DID } from './config'; +import { notifyConfigured } from './keypair'; +import { relayRequestPermission, relayRevokeSelf } from './relay'; +import { ensureNotifySchema, setSubscriber } from './db'; + +/** Mint a user service-auth token (issued by the user's PDS) for a relay method. + * Requires the matching `rpc?lxm=…&aud=*` OAuth scope (see settings.ts). */ +async function mintUserToken(client: Client, lxm: string): Promise { + const res = await client.get('com.atproto.server.getServiceAuth', { + params: { aud: RELAY_DID as Did, lxm: lxm as Nsid } + }); + if (!res.ok) throw new Error('could not mint user token'); + return res.data.token; +} + +/** Ask atmo.pub to let us notify the signed-in user. `alreadyGranted` means we + * can send right away; `pending` means they must approve it on atmo.pub (the + * subscriberChanged webhook then flips them to enabled). */ +export const enableNotifications = command(async (): Promise<{ status: 'enabled' | 'pending' }> => { + const { locals, platform } = getRequestEvent(); + const env = platform?.env; + if (!env || !notifyConfigured(env)) error(400, 'Notifications are not configured'); + if (!locals.client || !locals.did) error(401, 'Not signed in'); + + try { + const userToken = await mintUserToken(locals.client, LXM_REQUEST_PERMISSION); + const res = await relayRequestPermission(env, userToken); + await ensureNotifySchema(env.DB); + + if (res.status === 'alreadyGranted') { + await setSubscriber(env.DB, locals.did, true); + return { status: 'enabled' }; + } + return { status: 'pending' }; + } catch (e) { + if (e && typeof e === 'object' && 'status' in e) throw e; // re-throw SvelteKit errors + error(400, e instanceof Error ? e.message : 'Could not enable notifications'); + } +}); + +/** Revoke our grant at the relay and stop sending locally. */ +export const disableNotifications = command(async (): Promise<{ status: 'disabled' }> => { + const { locals, platform } = getRequestEvent(); + const env = platform?.env; + if (!env || !notifyConfigured(env)) error(400, 'Notifications are not configured'); + if (!locals.client || !locals.did) error(401, 'Not signed in'); + + try { + const userToken = await mintUserToken(locals.client, LXM_REVOKE_SELF); + await relayRevokeSelf(env, userToken); + } catch (e) { + // Even if the relay call fails, mark disabled locally so we stop sending. + console.error('[notify] revokeSelf failed:', e); + } + + await ensureNotifySchema(env.DB); + await setSubscriber(env.DB, locals.did, false); + return { status: 'disabled' }; +}); diff --git a/apps/web/src/routes/(app)/+layout.svelte b/apps/web/src/routes/(app)/+layout.svelte index 26e47b9..89f5b42 100644 --- a/apps/web/src/routes/(app)/+layout.svelte +++ b/apps/web/src/routes/(app)/+layout.svelte @@ -4,11 +4,10 @@ import { Head, Navbar, Button, Avatar } from '@foxui/core'; import { resolve } from '$app/paths'; import { page } from '$app/state'; + import { dev } from '$app/environment'; import { ModeWatcher } from 'mode-watcher'; import LoginModal from '$lib/components/LoginModal.svelte'; - import CreateEventModal, { - createEventModalState - } from '$lib/components/CreateEventModal.svelte'; + import CreateEventModal, { createEventModalState } from '$lib/components/CreateEventModal.svelte'; let { children } = $props(); @@ -68,6 +67,35 @@ create event {/if} + + {#if dev} + + + + + + + {/if} { + const env = platform?.env; + const configured = !!env && notifyConfigured(env); + + if (!locals.did || !configured || !env) { + return { configured, loggedIn: !!locals.did, state: 'none' as NotifyState }; + } + + await ensureNotifySchema(env.DB); + const sub = await getSubscriber(env.DB, locals.did); + const state: NotifyState = !sub ? 'none' : sub.enabled ? 'enabled' : 'disabled'; + + return { configured, loggedIn: true, state }; +}; diff --git a/apps/web/src/routes/(app)/settings/+page.svelte b/apps/web/src/routes/(app)/settings/+page.svelte new file mode 100644 index 0000000..01ff155 --- /dev/null +++ b/apps/web/src/routes/(app)/settings/+page.svelte @@ -0,0 +1,132 @@ + + +
+

+ Settings +

+ +
+
+

Notifications

+

+ Delivered through atmo.pub, on the channels you choose there (web push, Telegram, …). +

+ + {#if !data.loggedIn} +

+ Sign in to turn on notifications. +

+
+ +
+ {:else if !data.configured} +

+ Notifications aren’t available on this server yet. +

+ {:else} +
    + {#each perks as perk (perk)} +
  • + • + {perk} +
  • + {/each} +
+ + {#if notifyState === 'enabled'} +
+ + + Notifications are on + + +
+ {:else} +
+ +
+ + {#if pending} +

+ Almost there — approve atmo.rsvp on + atmo.pub to start receiving notifications. +

+ {/if} + {/if} + + {#if errorMsg} +

{errorMsg}

+ {/if} + {/if} +
+
+
diff --git a/apps/web/src/routes/api/cron/+server.ts b/apps/web/src/routes/api/cron/+server.ts index 229df91..69bf917 100644 --- a/apps/web/src/routes/api/cron/+server.ts +++ b/apps/web/src/routes/api/cron/+server.ts @@ -1,5 +1,6 @@ import { contrail, ensureInit } from '$lib/contrail/index'; import { processBotMentions } from '$lib/bot/process-mentions'; +import { runNotifications } from '$lib/notify/process'; import type { RequestHandler } from './$types'; export const POST: RequestHandler = async ({ request, platform }) => { @@ -29,5 +30,14 @@ export const POST: RequestHandler = async ({ request, platform }) => { console.error('[cron] contrail.ingest failed:', e); } + // atmo.pub notifications: event reminders + host RSVP alerts. Runs after + // ingest so it sees the freshest records; isolated so a failure can't 500 + // the tick. No-ops when notifications aren't configured. + try { + await runNotifications(platform!.env, db); + } catch (e) { + console.error('[cron] runNotifications failed:', e); + } + return new Response('OK'); }; diff --git a/apps/web/src/routes/xrpc/pub.atmo.notify.subscriberChanged/+server.ts b/apps/web/src/routes/xrpc/pub.atmo.notify.subscriberChanged/+server.ts new file mode 100644 index 0000000..78536f8 --- /dev/null +++ b/apps/web/src/routes/xrpc/pub.atmo.notify.subscriberChanged/+server.ts @@ -0,0 +1,61 @@ +import { json } from '@sveltejs/kit'; +import { ServiceJwtVerifier } from '@atcute/xrpc-server/auth'; +import { + CompositeDidDocumentResolver, + PlcDidDocumentResolver, + WebDidDocumentResolver +} from '@atcute/identity-resolver'; +import type { Did, Nsid } from '@atcute/lexicons'; +import { LXM_SUBSCRIBER_CHANGED, RELAY_DID } from '$lib/notify/config'; +import { getSenderDid, notifyConfigured } from '$lib/notify/keypair'; +import { ensureNotifySchema, setSubscriber } from '$lib/notify/db'; +import type { RequestHandler } from './$types'; + +// Relay -> us callback: fires when a user enables/disables notifications from +// our app (e.g. on atmo.pub). Lets us keep `notify_subscribers` accurate without +// polling. This route is more specific than /xrpc/[...path], so it wins. + +let verifier: ServiceJwtVerifier | null = null; +function getVerifier(serviceDid: Did): ServiceJwtVerifier { + if (!verifier) { + verifier = new ServiceJwtVerifier({ + serviceDid, + resolver: new CompositeDidDocumentResolver({ + methods: { plc: new PlcDidDocumentResolver(), web: new WebDidDocumentResolver() } + }) + }); + } + return verifier; +} + +export const POST: RequestHandler = async ({ request, platform }) => { + const env = platform!.env; + if (!notifyConfigured(env)) return json({ error: 'NotConfigured' }, { status: 404 }); + + const authz = request.headers.get('authorization') ?? ''; + const token = authz.startsWith('Bearer ') ? authz.slice(7) : ''; + if (!token) return json({ error: 'NotAuthorized' }, { status: 401 }); + + // Verify signature + aud (our DID) + lxm, then the critical check: the issuer + // MUST be the relay — otherwise anyone could forge enrollment changes. + const result = await getVerifier(getSenderDid(env)).verify(token, { + lxm: LXM_SUBSCRIBER_CHANGED as Nsid + }); + if (!result.ok || result.value.issuer !== RELAY_DID) { + return json({ error: 'NotAuthorized' }, { status: 401 }); + } + + let body: { recipient?: unknown; enabled?: unknown }; + try { + body = await request.json(); + } catch { + return json({ error: 'InvalidRequest' }, { status: 400 }); + } + if (typeof body.recipient !== 'string' || typeof body.enabled !== 'boolean') { + return json({ error: 'InvalidRequest' }, { status: 400 }); + } + + await ensureNotifySchema(env.DB); + await setSubscriber(env.DB, body.recipient, body.enabled); + return json({ ok: true }); +}; diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index dd12d99..9e4913e 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -102,6 +102,9 @@ importers: '@atcute/client': specifier: ^4.2.1 version: 4.2.1 + '@atcute/crypto': + specifier: ^2.4.1 + version: 2.4.1 '@atcute/identity-resolver': specifier: ^1.2.2 version: 1.2.2(@atcute/identity@1.1.4) @@ -117,6 +120,9 @@ importers: '@atcute/tid': specifier: ^1.1.2 version: 1.1.2 + '@atcute/xrpc-server': + specifier: ^0.1.12 + version: 0.1.12 '@atmo-dev/contrail-lexicons': specifier: ^0.4.4 version: 0.4.4(wrangler@4.77.0(@cloudflare/workers-types@4.20260317.1))