From 710f90c96454cac46833a97e6ceff1cc43b30009 Mon Sep 17 00:00:00 2001 From: Florian <45694132+flo-bit@users.noreply.github.com> Date: Sat, 23 May 2026 12:09:26 +0200 Subject: [PATCH] inbox --- apps/relay/migrations/0004_inbox.sql | 17 +++ apps/relay/src/db/queries.ts | 106 +++++++++++++++++ apps/relay/src/rpc/entrypoint.ts | 8 ++ apps/relay/src/rpc/ops.ts | 47 ++++++++ apps/relay/src/xrpc/send.ts | 15 +++ apps/relay/test/inbox.test.ts | 48 ++++++++ apps/relay/test/send.test.ts | 19 +++ apps/web/src/lib/remote/notifs.remote.ts | 7 ++ apps/web/src/lib/server/relay.ts | 5 +- .../src/routes/(app)/inbox/+page.server.ts | 15 +++ apps/web/src/routes/(app)/inbox/+page.svelte | 109 +++++++++++++++--- .../lexicons/tools/atmo/notifs/send.json | 16 +++ packages/lexicons/src/rpc.ts | 29 +++++ 13 files changed, 427 insertions(+), 14 deletions(-) create mode 100644 apps/relay/migrations/0004_inbox.sql create mode 100644 apps/relay/test/inbox.test.ts create mode 100644 apps/web/src/routes/(app)/inbox/+page.server.ts diff --git a/apps/relay/migrations/0004_inbox.sql b/apps/relay/migrations/0004_inbox.sql new file mode 100644 index 0000000..ccbd018 --- /dev/null +++ b/apps/relay/migrations/0004_inbox.sql @@ -0,0 +1,17 @@ +-- Inbox (Phase 3). Every accepted `send` is recorded here as the canonical +-- history; per-category routing (Phase 4) only decides which alert channels also +-- fire. `read_at` null = unread. `actors` is a JSON array of handles/DIDs. + +CREATE TABLE notifications ( + id TEXT PRIMARY KEY, + recipient_did TEXT NOT NULL, + sender_did TEXT NOT NULL, + category TEXT, + title TEXT NOT NULL, + body TEXT NOT NULL, + uri TEXT, + actors TEXT, + created_at INTEGER NOT NULL, + read_at INTEGER +); +CREATE INDEX notifications_by_recipient ON notifications (recipient_did, created_at DESC); diff --git a/apps/relay/src/db/queries.ts b/apps/relay/src/db/queries.ts index b1c9e9d..5eab091 100644 --- a/apps/relay/src/db/queries.ts +++ b/apps/relay/src/db/queries.ts @@ -509,3 +509,109 @@ export async function deletePushSubscription(db: D1Database, endpoint: string): .run(); return changed(result); } + +// --------------------------------------------------------------------------- +// notifications (inbox) +// --------------------------------------------------------------------------- + +export interface NotificationRow { + id: string; + recipient_did: Did; + sender_did: Did; + category: string | null; + title: string; + body: string; + uri: string | null; + actors: string | null; // JSON array + created_at: number; + read_at: number | null; +} + +export interface InsertNotificationInput { + id: string; + recipientDid: Did; + senderDid: Did; + category: string | null; + title: string; + body: string; + uri: string | null; + actors: string[] | null; + createdAt: number; +} + +export async function insertNotification(db: D1Database, input: InsertNotificationInput): Promise { + await db + .prepare( + `INSERT INTO notifications + (id, recipient_did, sender_did, category, title, body, uri, actors, created_at, read_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, NULL)`, + ) + .bind( + input.id, + input.recipientDid, + input.senderDid, + input.category, + input.title, + input.body, + input.uri, + input.actors ? JSON.stringify(input.actors) : null, + input.createdAt, + ) + .run(); +} + +/** Page the inbox newest-first; `before` is a `created_at` cursor (exclusive). */ +export async function listNotificationsForRecipient( + db: D1Database, + recipientDid: Did, + limit: number, + before?: number, +): Promise { + const sql = + before !== undefined + ? 'SELECT * FROM notifications WHERE recipient_did = ? AND created_at < ? ORDER BY created_at DESC LIMIT ?' + : 'SELECT * FROM notifications WHERE recipient_did = ? ORDER BY created_at DESC LIMIT ?'; + const stmt = + before !== undefined + ? db.prepare(sql).bind(recipientDid, before, limit) + : db.prepare(sql).bind(recipientDid, limit); + const { results } = await stmt.all(); + return results; +} + +export function countUnreadNotifications(db: D1Database, recipientDid: Did): Promise<{ c: number } | null> { + return db + .prepare('SELECT COUNT(*) AS c FROM notifications WHERE recipient_did = ? AND read_at IS NULL') + .bind(recipientDid) + .first<{ c: number }>(); +} + +export async function markNotificationsRead( + db: D1Database, + recipientDid: Did, + ids: string[], + readAt: number, +): Promise { + if (ids.length === 0) return 0; + const placeholders = ids.map(() => '?').join(','); + const result = await db + .prepare( + `UPDATE notifications SET read_at = ? + WHERE recipient_did = ? AND read_at IS NULL AND id IN (${placeholders})`, + ) + .bind(readAt, recipientDid, ...ids) + .run(); + return result.meta.changes ?? 0; +} + +export async function markAllNotificationsRead( + db: D1Database, + recipientDid: Did, + readAt: number, +): Promise { + const result = await db + .prepare('UPDATE notifications SET read_at = ? WHERE recipient_did = ? AND read_at IS NULL') + .bind(readAt, recipientDid) + .run(); + return result.meta.changes ?? 0; +} diff --git a/apps/relay/src/rpc/entrypoint.ts b/apps/relay/src/rpc/entrypoint.ts index a301934..0e4a4d6 100644 --- a/apps/relay/src/rpc/entrypoint.ts +++ b/apps/relay/src/rpc/entrypoint.ts @@ -2,6 +2,8 @@ import { WorkerEntrypoint } from 'cloudflare:workers'; import type { Did } from '@atcute/lexicons'; import type { + ListNotificationsResult, + MarkReadInput, NotifsRpc, PushSubscriptionInput, ToolsAtmoNotifsDenyPending, @@ -70,4 +72,10 @@ export class RelayRpc extends WorkerEntrypoint implements NotifsRpc { unregisterWebPush(did: Did, endpoint: string) { return ops.unregisterWebPush(this.env, did, endpoint); } + listNotifications(did: Did, cursor?: string): Promise { + return ops.listNotifications(this.env, did, cursor); + } + markRead(did: Did, input: MarkReadInput) { + return ops.markRead(this.env, did, input); + } } diff --git a/apps/relay/src/rpc/ops.ts b/apps/relay/src/rpc/ops.ts index e90998c..212d130 100644 --- a/apps/relay/src/rpc/ops.ts +++ b/apps/relay/src/rpc/ops.ts @@ -6,6 +6,9 @@ import type { Did } from '@atcute/lexicons'; import type { + ListNotificationsResult, + MarkReadInput, + NotificationView, PushSubscriptionInput, ToolsAtmoNotifsDenyPending, ToolsAtmoNotifsGetSettings, @@ -233,3 +236,47 @@ export async function unregisterWebPush( const unregistered = await q.deletePushSubscriptionForDid(env.DB, did, endpoint); return { unregistered }; } + +const INBOX_PAGE_SIZE = 30; + +function toNotificationView(row: q.NotificationRow): NotificationView { + return { + id: row.id, + sender: row.sender_did, + category: row.category ?? undefined, + title: row.title, + body: row.body, + uri: row.uri ?? undefined, + actors: row.actors ? (JSON.parse(row.actors) as string[]) : [], + createdAt: toIsoDatetime(row.created_at), + read: row.read_at !== null, + }; +} + +export async function listNotifications( + env: Env, + did: Did, + cursor?: string, +): Promise { + const before = cursor !== undefined ? Number(cursor) : undefined; + const rows = await q.listNotificationsForRecipient(env.DB, did, INBOX_PAGE_SIZE, before); + const unread = (await q.countUnreadNotifications(env.DB, did))?.c ?? 0; + const last = rows.at(-1); + return { + notifications: rows.map(toNotificationView), + unread, + cursor: rows.length === INBOX_PAGE_SIZE && last ? String(last.created_at) : undefined, + }; +} + +export async function markRead( + env: Env, + did: Did, + input: MarkReadInput, +): Promise<{ marked: number }> { + const readAt = now(); + const marked = input.all + ? await q.markAllNotificationsRead(env.DB, did, readAt) + : await q.markNotificationsRead(env.DB, did, input.ids ?? [], readAt); + return { marked }; +} diff --git a/apps/relay/src/xrpc/send.ts b/apps/relay/src/xrpc/send.ts index 27c19c3..902a589 100644 --- a/apps/relay/src/xrpc/send.ts +++ b/apps/relay/src/xrpc/send.ts @@ -30,6 +30,21 @@ export function makeSend(app: AppContext): ProcedureConfig { + await q.insertNotification(env.DB, { + id, + recipientDid: user, + senderDid: SENDER, + category: 'mention', + title: id, + body: 'body', + uri: null, + actors, + createdAt, + }); +} + +it('listNotifications returns newest-first with actors and an unread count', async () => { + const user = 'did:plc:inbox1' as Did; + await seed(user, 'n-a', 1000, ['alice.test']); + await seed(user, 'n-b', 2000, null); + + const res = await ops.listNotifications(env, user); + + expect(res.notifications.map((n) => n.id)).toEqual(['n-b', 'n-a']); + expect(res.notifications[1]?.actors).toEqual(['alice.test']); + expect(res.notifications[0]?.read).toBe(false); + expect(res.unread).toBe(2); +}); + +it('markRead({ ids }) marks those; markRead({ all }) clears the rest', async () => { + const user = 'did:plc:inbox2' as Did; + await seed(user, 'm1', 1000, null); + await seed(user, 'm2', 2000, null); + await seed(user, 'm3', 3000, null); + + expect((await ops.markRead(env, user, { ids: ['m1', 'm2'] })).marked).toBe(2); + expect((await ops.listNotifications(env, user)).unread).toBe(1); + + expect((await ops.markRead(env, user, { all: true })).marked).toBe(1); + expect((await ops.listNotifications(env, user)).unread).toBe(0); +}); diff --git a/apps/relay/test/send.test.ts b/apps/relay/test/send.test.ts index 5112bd6..fcac17b 100644 --- a/apps/relay/test/send.test.ts +++ b/apps/relay/test/send.test.ts @@ -86,6 +86,25 @@ it('enqueues and reports delivered=1 with a linked channel', async () => { expect(log?.delivered_count).toBe(1); }); +it('records the notification in the inbox', async () => { + const sender = await makeIdentity('did:plc:sendinbox'); + mockPlc(sender); + await q.upsertGrant(env.DB, { + recipientDid: RECIPIENT, + senderDid: sender.did, + grantedAt: Date.now(), + title: null, + description: null, + iconUrl: null + }); + const jwt = await makeJwt(sender, { lxm: SEND }); + + await call(send(jwt)); + + const rows = await q.listNotificationsForRecipient(env.DB, RECIPIENT, 50); + expect(rows.some((r) => r.sender_did === sender.did && r.title === 'Hello')).toBe(true); +}); + it('accepts silently with delivered=0 when the grant is muted', async () => { const sender = await makeIdentity('did:plc:sendmuted'); mockPlc(sender); diff --git a/apps/web/src/lib/remote/notifs.remote.ts b/apps/web/src/lib/remote/notifs.remote.ts index a26f898..33e35d3 100644 --- a/apps/web/src/lib/remote/notifs.remote.ts +++ b/apps/web/src/lib/remote/notifs.remote.ts @@ -65,3 +65,10 @@ export const registerPush = command( export const unregisterPush = command(v.object({ endpoint: v.string() }), async ({ endpoint }) => { await requireRelay().unregisterWebPush(endpoint); }); + +export const markNotificationsRead = command( + v.object({ ids: v.optional(v.array(v.string())), all: v.optional(v.boolean()) }), + async (input) => { + await requireRelay().markRead(input); + } +); diff --git a/apps/web/src/lib/server/relay.ts b/apps/web/src/lib/server/relay.ts index de3801e..dddd36a 100644 --- a/apps/web/src/lib/server/relay.ts +++ b/apps/web/src/lib/server/relay.ts @@ -10,6 +10,7 @@ // `relayFor` throws a clear error rather than failing cryptically. import type { Did } from '@atcute/lexicons'; import type { + MarkReadInput, PushSubscriptionInput, ToolsAtmoNotifsDenyPending, ToolsAtmoNotifsGrant, @@ -50,6 +51,8 @@ export function relayFor(platform: App.Platform | undefined, did: Did | null) { updateSettings: (input: ToolsAtmoNotifsUpdateSettings.$input) => svc.updateSettings(did, input), registerWebPush: (sub: PushSubscriptionInput) => svc.registerWebPush(did, sub), - unregisterWebPush: (endpoint: string) => svc.unregisterWebPush(did, endpoint) + unregisterWebPush: (endpoint: string) => svc.unregisterWebPush(did, endpoint), + listNotifications: (cursor?: string) => svc.listNotifications(did, cursor), + markRead: (input: MarkReadInput) => svc.markRead(did, input) }; } diff --git a/apps/web/src/routes/(app)/inbox/+page.server.ts b/apps/web/src/routes/(app)/inbox/+page.server.ts new file mode 100644 index 0000000..01d5abf --- /dev/null +++ b/apps/web/src/routes/(app)/inbox/+page.server.ts @@ -0,0 +1,15 @@ +import { relayFor } from '$lib/server/relay'; + +import type { PageServerLoad } from './$types'; + +// Inbox: the recipient's notification history. The (app) layout guard guarantees +// `locals.did` is present. +export const load: PageServerLoad = async ({ locals, platform }) => { + const relay = relayFor(platform, locals.did); + const res = await relay.listNotifications(); + + return { + notifications: res?.notifications ?? [], + unread: res?.unread ?? 0 + }; +}; diff --git a/apps/web/src/routes/(app)/inbox/+page.svelte b/apps/web/src/routes/(app)/inbox/+page.svelte index da4f93c..11c2bbe 100644 --- a/apps/web/src/routes/(app)/inbox/+page.svelte +++ b/apps/web/src/routes/(app)/inbox/+page.svelte @@ -1,24 +1,107 @@ Inbox · atmo.pub
-
-

Inbox

+
+
+

Inbox

+ {#if data.unread > 0} + + {data.unread} + + {/if} +
+ {#if data.unread > 0} + + {/if}
-
- + {:else} +
    + {#each data.notifications as n (n.id)} +
  • + +
  • + {/each} +
+ {/if}
diff --git a/packages/lexicons/lexicons/tools/atmo/notifs/send.json b/packages/lexicons/lexicons/tools/atmo/notifs/send.json index 3ea5d9a..8d61c8b 100644 --- a/packages/lexicons/lexicons/tools/atmo/notifs/send.json +++ b/packages/lexicons/lexicons/tools/atmo/notifs/send.json @@ -21,6 +21,22 @@ "threadKey": { "type": "string", "description": "Optional opaque key for grouping related notifications." + }, + "category": { + "type": "string", + "maxLength": 64, + "description": "Optional free-form category (e.g. 'mention', 'reply') used for per-category routing." + }, + "categoryDescription": { + "type": "string", + "maxLength": 200, + "description": "Optional human description of the category, shown in the routing UI." + }, + "actors": { + "type": "array", + "maxLength": 8, + "items": { "type": "string" }, + "description": "Optional handles/DIDs of the people behind this notification (for avatar stacks)." } } } diff --git a/packages/lexicons/src/rpc.ts b/packages/lexicons/src/rpc.ts index 2a061ab..913b1d1 100644 --- a/packages/lexicons/src/rpc.ts +++ b/packages/lexicons/src/rpc.ts @@ -34,6 +34,31 @@ export interface PushSubscriptionInput { auth: string; } +/** A notification as shown in the inbox (binding-only; no public lexicon). */ +export interface NotificationView { + id: string; + sender: Did; + category?: string; + title: string; + body: string; + uri?: string; + actors: string[]; + createdAt: string; + read: boolean; +} + +export interface ListNotificationsResult { + notifications: NotificationView[]; + /** `created_at` cursor for the next (older) page; absent when there are no more. */ + cursor?: string; + unread: number; +} + +export interface MarkReadInput { + ids?: string[]; + all?: boolean; +} + export interface NotifsRpc { grant(did: Did, input: ToolsAtmoNotifsGrant.$input): Promise; revoke(did: Did, input: ToolsAtmoNotifsRevoke.$input): Promise; @@ -65,4 +90,8 @@ export interface NotifsRpc { // Web push (binding-only; no public lexicon). registerWebPush(did: Did, sub: PushSubscriptionInput): Promise<{ registered: boolean }>; unregisterWebPush(did: Did, endpoint: string): Promise<{ unregistered: boolean }>; + + // Inbox (binding-only; no public lexicon). + listNotifications(did: Did, cursor?: string): Promise; + markRead(did: Did, input: MarkReadInput): Promise<{ marked: number }>; } -- 2.51.2