// Bluesky API layer — reads live state (notifications, subscriptions, mutes), // builds per-actor ManagementTargets for the model, and provides before/after // diffing for verification. No mutations here — those happen via code-mode.ts. import { AppBskyFeedPost, type Agent, type AppBskyNotificationListNotifications, } from '@atproto/api' import type {ActorRef, EvidenceRow, ManagementTarget, NormalizedNotification} from './types' function chunk(items: T[], size: number): T[][] { const out: T[][] = [] for (let i = 0; i < items.length; i += size) out.push(items.slice(i, i + size)) return out } function preview(text: string | undefined, limit = 120) { if (!text) return '' const collapsed = text.replace(/\s+/g, ' ').trim() if (collapsed.length <= limit) return collapsed return `${collapsed.slice(0, limit - 1).trimEnd()}…` } export function actorProfileUrl(identifier: string) { return `https://bsky.app/profile/${encodeURIComponent(identifier)}` } function actorRefFromProfile(profile: {did: string; handle?: string; displayName?: string}) { const handle = profile.handle || profile.did return { did: profile.did, handle, name: profile.displayName || handle, profileUrl: actorProfileUrl(handle), } satisfies ActorRef } // When a notification is about someone else's post (e.g. a like on their post), // the notification has a reasonSubject URI but no text. This fetches the actual // post text so we can show a snippet in evidence rows. async function fetchSubjectTexts( agent: Agent, notifications: AppBskyNotificationListNotifications.Notification[], ) { const uris = [ ...new Set( notifications .map(n => n.reasonSubject) .filter((uri): uri is string => Boolean(uri?.includes('app.bsky.feed.post'))), ), ] const texts = new Map() for (const batch of chunk(uris, 25)) { const res = await agent.app.bsky.feed.getPosts({uris: batch}) for (const post of res.data.posts) { if (AppBskyFeedPost.isRecord(post.record) && post.record.text) { texts.set(String(post.uri), String(post.record.text)) } } } return texts } // Fetch the 50 most recent notifications, normalize into our flat shape, // enrich with subject text. Persistence is handled by the route runtime so // Bun can use SQLite and Cloudflare Workers can use D1. export async function listNotifications(agent: Agent): Promise { const res = await agent.listNotifications({limit: 50}) const subjectTexts = await fetchSubjectTexts(agent, res.data.notifications) const normalized = res.data.notifications.map(notification => { const record = notification.record const text = AppBskyFeedPost.isRecord(record) && typeof record.text === 'string' ? record.text : undefined const isReply = Boolean(AppBskyFeedPost.isRecord(record) && record.reply) return { uri: notification.uri, reason: notification.reason, reasonSubject: notification.reasonSubject, authorDid: notification.author.did, authorHandle: notification.author.handle, authorName: notification.author.displayName || notification.author.handle, indexedAt: notification.indexedAt, isRead: notification.isRead, text, subjectText: notification.reasonSubject ? subjectTexts.get(notification.reasonSubject) : undefined, isReply, } satisfies NormalizedNotification }) return normalized } export type ActivitySubscription = { post: boolean reply: boolean actorHandle: string actorName: string } export async function getActorProfiles(agent: Agent, actorDids: string[]): Promise> { const unique = [...new Set(actorDids.filter(Boolean))] const out = new Map() for (const batch of chunk(unique, 25)) { const res = await agent.getProfiles({actors: batch}) for (const profile of res.data.profiles) { out.set(profile.did, actorRefFromProfile(profile)) } } for (const did of unique) { if (out.has(did)) continue out.set(did, { did, handle: did, name: did, profileUrl: actorProfileUrl(did), }) } return out } const TYPEAHEAD_URL = process.env.TYPEAHEAD_URL || 'https://typeahead.waow.tech' type TypeaheadActor = { did: string handle: string displayName?: string avatar?: string } // Community-run actor search (https://typeahead.waow.tech) — drop-in compatible // with app.bsky.actor.searchActorsTypeahead. Public endpoint, no auth required, // so this no longer burns login-session budget on every keystroke. export async function searchActorSuggestions( query: string, limit = 6, baseUrl = TYPEAHEAD_URL, ): Promise { const trimmed = query.trim().replace(/^@/, '') if (!trimmed) return [] const url = `${baseUrl}/xrpc/app.bsky.actor.searchActorsTypeahead?q=${encodeURIComponent(trimmed)}&limit=${limit}` try { const res = await fetch(url, {headers: {'X-Client': 'noti'}}) if (!res.ok) return [] const body = (await res.json()) as {actors?: TypeaheadActor[]} const actors = Array.isArray(body.actors) ? body.actors : [] return actors.map(actor => ({ did: actor.did, handle: actor.handle, name: actor.displayName || actor.handle, profileUrl: actorProfileUrl(actor.handle), })) } catch { return [] } } export async function listActivitySubscriptions(agent: Agent) { const rows = new Map() let cursor: string | undefined do { const res = await agent.app.bsky.notification.listActivitySubscriptions({cursor}) for (const profile of res.data.subscriptions) { const sub = profile.viewer?.activitySubscription if (sub?.post || sub?.reply) { rows.set(profile.did, { post: Boolean(sub.post), reply: Boolean(sub.reply), actorHandle: profile.handle, actorName: profile.displayName || profile.handle, }) } } cursor = res.data.cursor } while (cursor) return rows } export async function listMutedActors(agent: Agent) { const rows = new Set() let cursor: string | undefined do { const res = await agent.app.bsky.graph.getMutes({cursor, limit: 100}) for (const profile of res.data.mutes) rows.add(profile.did) cursor = res.data.cursor } while (cursor) return rows } export type BlockedActor = { actorHandle: string actorName: string blockingUri?: string } export async function listBlockedActors(agent: Agent) { const rows = new Map() let cursor: string | undefined do { const res = await agent.app.bsky.graph.getBlocks({cursor, limit: 100}) for (const profile of res.data.blocks) { rows.set(profile.did, { actorHandle: profile.handle, actorName: profile.displayName || profile.handle, blockingUri: profile.viewer?.blocking, }) } cursor = res.data.cursor } while (cursor) return rows } export type ActorState = { subscription: {post: boolean; reply: boolean} | null muted: boolean blocked: boolean blockingUri?: string } // Snapshot current subscription + mute state for a set of actors. // Used before and after apply to detect what actually changed. export async function getActorStates(agent: Agent, actorDids: string[]): Promise> { const unique = [...new Set(actorDids.filter(Boolean))] const [subscriptions, mutedActors, blockedActors] = await Promise.all([ listActivitySubscriptions(agent), listMutedActors(agent), listBlockedActors(agent), ]) const out = new Map() for (const did of unique) { const subscription = subscriptions.get(did) const blocked = blockedActors.get(did) out.set(did, { subscription: subscription ? {post: subscription.post, reply: subscription.reply} : null, muted: mutedActors.has(did), blocked: Boolean(blocked), blockingUri: blocked?.blockingUri, }) } return out } export type ActorStateChange = | {actorDid: string; kind: 'mute'; muted: boolean} | {actorDid: string; kind: 'block'; blocked: boolean} | {actorDid: string; kind: 'subscription'; subscription: {post: boolean; reply: boolean} | null} // Compare before/after snapshots — returns only fields that actually changed. // This is what verification reports to the user. export function diffActorStates( before: Map, after: Map, ): ActorStateChange[] { const changes: ActorStateChange[] = [] for (const [did, afterState] of after) { const beforeState = before.get(did) if (!beforeState) continue if (beforeState.muted !== afterState.muted) { changes.push({actorDid: did, kind: 'mute', muted: afterState.muted}) } if (beforeState.blocked !== afterState.blocked) { changes.push({actorDid: did, kind: 'block', blocked: afterState.blocked}) } const bSub = beforeState.subscription const aSub = afterState.subscription if (JSON.stringify(bSub) !== JSON.stringify(aSub)) { changes.push({actorDid: did, kind: 'subscription', subscription: aSub}) } } return changes } function evidenceRows(notifications: NormalizedNotification[]): EvidenceRow[] { return notifications.slice(0, 3).map(notification => ({ uri: notification.uri, actor: notification.authorName, actorProfileUrl: actorProfileUrl(notification.authorHandle || notification.authorDid || notification.authorName), reason: notification.reason, snippet: preview( notification.text || notification.subjectText || notification.reasonSubject || notification.reason, ), })) } // The core state-building function. Groups unread notifications by actor, // enriches with subscription/mute state and historical pressure from storage, // and produces the ManagementTarget[] that gets sent to the model. // Sorted: subscribed actors first, then by unread count, then by recency. export function buildManagementTargets( notifications: NormalizedNotification[], subscriptionModes: Map, mutedActors: Set, blockedActors: Map, history: Map, ): ManagementTarget[] { const unread = notifications.filter(n => !n.isRead && n.authorDid) const grouped = new Map() for (const notification of unread) { const actorDid = notification.authorDid if (!actorDid) continue grouped.set(actorDid, [...(grouped.get(actorDid) || []), notification]) } const actorDids = [...new Set([...grouped.keys(), ...subscriptionModes.keys(), ...blockedActors.keys()])] return actorDids .map(actorDid => { const actorNotifications = grouped.get(actorDid) || [] const latest = [...actorNotifications].sort((a, b) => b.indexedAt.localeCompare(a.indexedAt))[0] const modes = subscriptionModes.get(actorDid) const block = blockedActors.get(actorDid) const hist = history.get(actorDid) || {allTimeCount: 0, recent24hCount: 0, recent7dCount: 0} const reasons: Record = {} for (const notification of actorNotifications) { reasons[notification.reason] = (reasons[notification.reason] || 0) + 1 } const actorHandle = latest?.authorHandle || modes?.actorHandle || block?.actorHandle || actorDid const actorName = latest?.authorName || modes?.actorName || block?.actorName || actorHandle return { actorDid, actorHandle, actorName, sourceUris: actorNotifications.map(n => n.uri), currentUnreadCount: actorNotifications.length, recent24hCount: hist.recent24hCount, recent7dCount: hist.recent7dCount, allTimeCount: hist.allTimeCount, reasons, subscriptionPosts: Boolean(modes?.post), subscriptionReplies: Boolean(modes?.reply), muted: mutedActors.has(actorDid), blocked: Boolean(block), blockingUri: block?.blockingUri, examples: evidenceRows(actorNotifications), } satisfies ManagementTarget }) .sort((a, b) => { return ( Number(b.subscriptionPosts || b.subscriptionReplies) - Number(a.subscriptionPosts || a.subscriptionReplies) || Number(b.blocked) - Number(a.blocked) || b.currentUnreadCount - a.currentUnreadCount || b.recent24hCount - a.recent24hCount || b.allTimeCount - a.allTimeCount ) }) }