diff --git a/docs/superpowers/plans/2026-07-09-notifications.md b/docs/superpowers/plans/2026-07-09-notifications.md new file mode 100644 index 0000000..e5363ca --- /dev/null +++ b/docs/superpowers/plans/2026-07-09-notifications.md @@ -0,0 +1,905 @@ +# Notifications Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Deliver push/desktop notifications when a user isn't looking at a conversation — including when the app is closed (Web Push) — with the engine in IIOS and the experience in the host app. + +**Architecture:** A `NotificationProjector` in IIOS reacts to `message.sent` events and runs three gates — **policy** (DM always; group only on @mention or reply-to-you), **presence** (skip if the recipient is focused on that thread), **mute** (skip if muted) — then dispatches via a swappable `NotificationPort` (Web Push/VAPID adapter). Presence comes from an explicit `focus_thread` socket signal (room membership ≠ viewing, because the sidebar joins every thread). The host app registers a push subscription + a service worker and renders/handles taps. + +**Tech Stack:** NestJS · Prisma/Postgres · socket.io · `web-push` (VAPID) · Vite/React · Browser Push API + Service Worker. + +**Spec:** `docs/superpowers/specs/2026-07-09-notifications-design.md` + +--- + +## File structure + +**iios (engine):** +- `packages/iios-service/prisma/schema.prisma` — MODIFY: `IiosNotificationSubscription` model + `muted` on `IiosThreadParticipant`. +- `packages/iios-service/src/notifications/presence.service.ts` — CREATE: in-memory focus tracker. +- `packages/iios-service/src/notifications/notification.port.ts` — CREATE: `NotificationPort` interface + DI token. +- `packages/iios-service/src/notifications/web-push.delivery.ts` — CREATE: VAPID adapter. +- `packages/iios-service/src/notifications/notification.projector.ts` — CREATE: the three gates + dispatch + prune. +- `packages/iios-service/src/notifications/notification.controller.ts` — CREATE: vapid key + subscribe/unsubscribe. +- `packages/iios-service/src/notifications/notification.module.ts` — CREATE: wiring. +- `packages/iios-service/src/messaging/message.gateway.ts` — MODIFY: `focus_thread` handler + `OnGatewayDisconnect`. +- `packages/iios-service/src/threads/threads.controller.ts` — MODIFY: `POST /v1/threads/:id/mute|unmute`. +- `packages/iios-service/src/messaging/message.service.ts` — MODIFY: `muteThread(threadId, principal, muted)`. +- `packages/iios-service/src/app.module.ts` — MODIFY: import `NotificationModule`. +- Test specs co-located: `presence.service.spec.ts`, `web-push.delivery.spec.ts`, `notification.projector.spec.ts`. + +**chat-web (experience):** +- `src/lib/notifications.ts` — CREATE: `registerPush`, `mute`, tab-title + desktop-notify helpers. +- `public/sw.js` — CREATE: service worker (`push` + `notificationclick`). +- `src/messages/provider.tsx` / `src/components/Conversation.tsx` / `src/components/Sidebar.tsx` — MODIFY: emit `focus_thread`, mute toggle, enable-notifications prompt. + +--- + +# PHASE 1 — Frontend quick win (chat-web, no backend) + +Desktop notification while the tab is unfocused + tab-title unread count, presence-gated client-side. Immediately demoable; no IIOS change. + +### Task 1: Desktop notification + tab-title unread (chat-web) + +**Files:** +- Create: `chat-web/src/lib/notifications.ts` +- Create: `chat-web/src/lib/notifications.test.ts` +- Modify: `chat-web/src/components/AppShell.tsx` (drive the title from unread), `chat-web/src/messages/provider.tsx` (fire a desktop notification on incoming message when hidden) + +- [ ] **Step 1: Write the failing test** for the pure title helper. + +`chat-web/src/lib/notifications.test.ts`: +```ts +import { describe, it, expect } from 'vitest'; +import { tabTitle } from './notifications'; + +describe('tabTitle', () => { + it('shows a count prefix when there are unread messages', () => { + expect(tabTitle(0)).toBe('Chat'); + expect(tabTitle(3)).toBe('(3) Chat'); + expect(tabTitle(150)).toBe('(99+) Chat'); + }); +}); +``` +> chat-web has no test runner yet. Add one: `pnpm add -D vitest` and a `"test": "vitest run"` script in `chat-web/package.json`. If the team prefers no test dep, keep `notifications.ts` pure and verify `tabTitle` via a scratch `node -e` instead — but prefer vitest. + +- [ ] **Step 2: Run it, expect FAIL.** `cd chat-web && pnpm test` → FAIL (`tabTitle` not exported). + +- [ ] **Step 3: Implement `notifications.ts`** (pure helpers + browser wiring): +```ts +// chat-web/src/lib/notifications.ts +export function tabTitle(unread: number): string { + if (unread <= 0) return 'Chat'; + return `(${unread > 99 ? '99+' : unread}) Chat`; +} + +/** Ask once; returns true if we may show desktop notifications. */ +export async function ensureNotifyPermission(): Promise { + if (!('Notification' in window)) return false; + if (Notification.permission === 'granted') return true; + if (Notification.permission === 'denied') return false; + return (await Notification.requestPermission()) === 'granted'; +} + +/** Show a desktop notification (Phase 1 — tab open but unfocused). */ +export function desktopNotify(title: string, body: string, onClick?: () => void): void { + if (!('Notification' in window) || Notification.permission !== 'granted') return; + const n = new Notification(title, { body, tag: title }); + if (onClick) n.onclick = () => { window.focus(); onClick(); n.close(); }; +} +``` + +- [ ] **Step 4: Run it, expect PASS.** `cd chat-web && pnpm test` → PASS. + +- [ ] **Step 5: Wire the tab title** in `AppShell.tsx`'s `Shell` component (it already has `conversations`): +```tsx +import { useEffect } from 'react'; +import { tabTitle } from '../lib/notifications'; +// inside Shell(), after `const conversations = ...`: +const unreadTotal = (conversations.data ?? []).reduce((n, c) => n + (c.unread ?? 0), 0); +useEffect(() => { document.title = tabTitle(unreadTotal); }, [unreadTotal]); +``` + +- [ ] **Step 6: Fire a desktop notification on incoming message when the tab is hidden and not viewing that thread.** In `provider.tsx`'s `useMessages`, the `message` handler already appends. Presence-gating here is at the *conversation* level, but a global hook is cleaner. Add to `AppShell.tsx`'s `Shell` a global listener via the socket: +```tsx +import { useSocket } from '../messages/provider'; +import { desktopNotify, ensureNotifyPermission } from '../lib/notifications'; +import { useNavigate, useParams } from '@tanstack/react-router'; +// inside Shell(): +const socket = useSocket(); +const navigate = useNavigate(); +const activeThreadId = (useParams({ strict: false }) as { threadId?: string }).threadId ?? null; +useEffect(() => { void ensureNotifyPermission(); }, []); +useEffect(() => { + const off = socket.on('message', (m) => { + const mine = m.senderId === session.userId; + const viewing = m.threadId === activeThreadId && !document.hidden; + if (mine || viewing) return; // presence gate (client-side) + desktopNotify(m.senderName, m.content || 'sent an attachment', + () => void navigate({ to: '/t/$threadId', params: { threadId: m.threadId } })); + }); + return off; +}, [socket, activeThreadId, session.userId, navigate]); +``` + +- [ ] **Step 7: Manual verify.** Two browser profiles; focus window A on a different conversation (or minimize) → send from B → a desktop notification appears and the tab title shows `(1) Chat`; click it → opens the thread. Sending in the conversation you're viewing → **no** notification. + +- [ ] **Step 8: Commit.** +```bash +cd chat-web && git add -A && git -c user.email=maaz@insigniaconsultancy.com commit -m "feat: desktop notifications + tab-title unread (presence-gated)" +``` + +--- + +# PHASE 2 — IIOS notification engine (TDD) + +### Task 2: Schema — subscription table + per-thread mute + +**Files:** +- Modify: `packages/iios-service/prisma/schema.prisma` + +- [ ] **Step 1: Add the model + column.** In `schema.prisma`, add to `IiosThreadParticipant` (before its `@@id`): +```prisma + muted Boolean @default(false) +``` +Add a new model and a relation on `IiosScope`: +```prisma +model IiosNotificationSubscription { + id String @id @default(cuid()) + scopeId String + scope IiosScope @relation(fields: [scopeId], references: [id], onDelete: Cascade) + actorId String + actor IiosActorRef @relation(fields: [actorId], references: [id]) + kind String @default("webpush") + endpoint String @unique + p256dh String + auth String + userAgent String? + createdAt DateTime @default(now()) + lastSeenAt DateTime @default(now()) + + @@index([actorId]) +} +``` +Add `notificationSubscriptions IiosNotificationSubscription[]` to the `IiosScope` model's relation list and `notificationSubscriptions IiosNotificationSubscription[]` to `IiosActorRef`. + +- [ ] **Step 2: Generate the migration + client** (additive; no data loss): +```bash +cd iios && DATABASE_URL='postgresql://iios:iios@localhost:5434/iios?schema=public' \ + pnpm --filter @insignia/iios-service exec prisma migrate dev --name add_notifications +``` +Expected: migration `*_add_notifications` applied; Prisma client regenerated. + +- [ ] **Step 3: Commit.** +```bash +cd iios && git add -A && git -c user.email=maaz@insigniaconsultancy.com commit -m "feat(iios): notification subscription table + per-thread mute" +``` + +### Task 3: PresenceService + gateway focus signal + +**Files:** +- Create: `packages/iios-service/src/notifications/presence.service.ts` +- Create: `packages/iios-service/src/notifications/presence.service.spec.ts` +- Modify: `packages/iios-service/src/messaging/message.gateway.ts` + +- [ ] **Step 1: Write the failing test.** `presence.service.spec.ts`: +```ts +import { describe, it, expect } from 'vitest'; +import { PresenceService } from './presence.service'; + +describe('PresenceService', () => { + it('reports viewing only for the focused thread, and clears on disconnect', () => { + const p = new PresenceService(); + expect(p.isViewing('alice', 'T1')).toBe(false); + p.setFocus('sock1', 'alice', 'T1'); + expect(p.isViewing('alice', 'T1')).toBe(true); + expect(p.isViewing('alice', 'T2')).toBe(false); // joined-elsewhere ≠ viewing + p.setFocus('sock1', 'alice', 'T2'); // moved focus + expect(p.isViewing('alice', 'T1')).toBe(false); + expect(p.isViewing('alice', 'T2')).toBe(true); + p.clearSocket('sock1'); + expect(p.isViewing('alice', 'T2')).toBe(false); + }); + it('any of the actor’s sockets counts as viewing', () => { + const p = new PresenceService(); + p.setFocus('sockA', 'bob', 'T9'); + p.setFocus('sockB', 'bob', null); // a second tab, no focus + expect(p.isViewing('bob', 'T9')).toBe(true); + }); +}); +``` + +- [ ] **Step 2: Run it, expect FAIL.** `cd iios && pnpm vitest run src/notifications/presence.service.spec.ts` → FAIL. + +- [ ] **Step 3: Implement `presence.service.ts`:** +```ts +import { Injectable } from '@nestjs/common'; + +/** In-memory focus tracker (single instance). Keyed on userId (the stable externalId, + * which the gateway has as principal.userId). Prod: back with Redis for multi-replica. */ +@Injectable() +export class PresenceService { + private readonly focus = new Map(); // socketId → focus + + setFocus(socketId: string, userId: string, threadId: string | null): void { + this.focus.set(socketId, { userId, threadId }); + } + clearSocket(socketId: string): void { + this.focus.delete(socketId); + } + isViewing(userId: string, threadId: string): boolean { + for (const f of this.focus.values()) if (f.userId === userId && f.threadId === threadId) return true; + return false; + } +} +``` + +- [ ] **Step 4: Run it, expect PASS.** + +- [ ] **Step 5: Wire the gateway.** In `message.gateway.ts`: add `OnGatewayDisconnect` to the imports and `implements`, inject `PresenceService`, and add handlers. Presence keys on `principal.userId` (the stable externalId) — which the projector already computes per participant, so the two sides match. Gateway additions: +```ts +// imports: add OnGatewayDisconnect +// class: implements OnGatewayInit, OnGatewayConnection, OnGatewayDisconnect +constructor(/* …existing… */ private readonly presence: PresenceService) {} + +handleDisconnect(client: Socket): void { + this.presence.clearSocket(client.id); +} + +@SubscribeMessage('focus_thread') +focusThread(@ConnectedSocket() client: Socket, @MessageBody() body: { threadId: string | null }): void { + const { principal } = client.data as SocketState; + this.presence.setFocus(client.id, principal.userId, body?.threadId ?? null); +} +``` + `MessageGateway` must be able to inject `PresenceService` → it will be provided by `NotificationModule` and imported into `MessageModule`, OR provide `PresenceService` in a shared module. Simplest: create `NotificationModule` (Task 7) that `exports: [PresenceService]` and have `MessageModule` import it. (Do this wiring in Task 7; until then the gateway won't compile — so implement Task 3's gateway edit and Task 7's module together, or stub the provider in `MessageModule` now and move it in Task 7.) **Chosen:** provide `PresenceService` in `MessageModule` now (`providers: [PresenceService], exports: [PresenceService]`), and `NotificationModule` imports `MessageModule` to reuse it. + +- [ ] **Step 6: Run the messaging tests** to ensure the gateway still wires: `cd iios && pnpm vitest run src/messaging` → PASS. + +- [ ] **Step 7: Commit.** +```bash +cd iios && git add -A && git -c user.email=maaz@insigniaconsultancy.com commit -m "feat(iios): PresenceService + focus_thread signal (viewing != room membership)" +``` + +### Task 4: NotificationPort + Web Push adapter + +**Files:** +- Create: `packages/iios-service/src/notifications/notification.port.ts` +- Create: `packages/iios-service/src/notifications/web-push.delivery.ts` +- Create: `packages/iios-service/src/notifications/web-push.delivery.spec.ts` +- Add dep: `web-push` + +- [ ] **Step 1: Add the dependency + VAPID keys.** +```bash +cd iios && pnpm --filter @insignia/iios-service add web-push && pnpm --filter @insignia/iios-service add -D @types/web-push +cd packages/iios-service && node -e "console.log(require('web-push').generateVAPIDKeys())" +``` +Put the printed `publicKey`/`privateKey` into env as `VAPID_PUBLIC_KEY` / `VAPID_PRIVATE_KEY`, and set `VAPID_SUBJECT=mailto:dev@insignia`. + +- [ ] **Step 2: Write the failing test** (mock the transport). `web-push.delivery.spec.ts`: +```ts +import { describe, it, expect, vi } from 'vitest'; +import { WebPushDelivery } from './web-push.delivery'; + +const sub = { endpoint: 'https://push/x', p256dh: 'k', auth: 'a' }; +const payload = { title: 'Bob', body: 'hi', data: { threadId: 'T1' } }; + +describe('WebPushDelivery', () => { + it('returns "sent" on success', async () => { + const send = vi.fn().mockResolvedValue({ statusCode: 201 }); + const d = new WebPushDelivery({ publicKey: 'p', privateKey: 'k', subject: 'mailto:x' }, send as never); + expect(await d.deliver(sub, payload)).toBe('sent'); + expect(send).toHaveBeenCalledOnce(); + }); + it('returns "gone" on 404/410 so the caller prunes', async () => { + const send = vi.fn().mockRejectedValue({ statusCode: 410 }); + const d = new WebPushDelivery({ publicKey: 'p', privateKey: 'k', subject: 'mailto:x' }, send as never); + expect(await d.deliver(sub, payload)).toBe('gone'); + }); + it('returns "failed" on other errors', async () => { + const send = vi.fn().mockRejectedValue({ statusCode: 500 }); + const d = new WebPushDelivery({ publicKey: 'p', privateKey: 'k', subject: 'mailto:x' }, send as never); + expect(await d.deliver(sub, payload)).toBe('failed'); + }); +}); +``` + +- [ ] **Step 3: Run it, expect FAIL.** + +- [ ] **Step 4: Implement the port + adapter.** +`notification.port.ts`: +```ts +export interface PushSub { endpoint: string; p256dh: string; auth: string } +export interface NotificationPayload { title: string; body: string; data: { threadId: string; interactionId?: string } } +export interface NotificationPort { + deliver(sub: PushSub, payload: NotificationPayload): Promise<'sent' | 'gone' | 'failed'>; +} +export const NOTIFICATION_PORT = Symbol('NOTIFICATION_PORT'); +``` +`web-push.delivery.ts`: +```ts +import { Injectable } from '@nestjs/common'; +import webpush from 'web-push'; +import type { NotificationPort, NotificationPayload, PushSub } from './notification.port'; + +type SendFn = typeof webpush.sendNotification; +interface Vapid { publicKey: string; privateKey: string; subject: string } + +@Injectable() +export class WebPushDelivery implements NotificationPort { + private readonly vapid?: Vapid; + constructor(vapid?: Vapid, private readonly send: SendFn = webpush.sendNotification) { + this.vapid = vapid?.publicKey ? vapid : undefined; + if (this.vapid) webpush.setVapidDetails(this.vapid.subject, this.vapid.publicKey, this.vapid.privateKey); + } + async deliver(sub: PushSub, payload: NotificationPayload): Promise<'sent' | 'gone' | 'failed'> { + if (!this.vapid) return 'failed'; // push disabled (no keys) + try { + await this.send({ endpoint: sub.endpoint, keys: { p256dh: sub.p256dh, auth: sub.auth } }, JSON.stringify(payload)); + return 'sent'; + } catch (e) { + const code = (e as { statusCode?: number }).statusCode; + return code === 404 || code === 410 ? 'gone' : 'failed'; + } + } +} +``` + +- [ ] **Step 5: Run it, expect PASS.** + +- [ ] **Step 6: Commit.** +```bash +cd iios && git add -A && git -c user.email=maaz@insigniaconsultancy.com commit -m "feat(iios): NotificationPort + Web Push (VAPID) delivery adapter" +``` + +### Task 5: NotificationProjector (the three gates) + +**Files:** +- Create: `packages/iios-service/src/notifications/notification.projector.ts` +- Create: `packages/iios-service/src/notifications/notification.projector.spec.ts` + +- [ ] **Step 1: Write the failing test.** `notification.projector.spec.ts` (mirrors `inbox.spec.ts` setup — real DB, `ms()` = MessageService, a fake `NotificationPort` capturing calls, a `PresenceService`): +```ts +import { describe, it, expect, beforeAll, afterAll, beforeEach, vi } from 'vitest'; +import { PrismaClient } from '@prisma/client'; +import { IIOS_EVENTS, type CloudEvent } from '@insignia/iios-contracts'; +import { makeFakePorts } from '@insignia/iios-testkit'; +import { resetDb } from '../test-utils/reset-db'; +import { MessageService, type MessagePrincipal } from '../messaging/message.service'; +import { ActorResolver } from '../identity/actor.resolver'; +import { OutboxBus } from '../outbox/outbox.bus'; +import { DlqService } from '../outbox/dlq.service'; +import { ProjectionCursorService } from '../projection/projection-cursor.service'; +import { PresenceService } from './presence.service'; +import { NotificationProjector } from './notification.projector'; +import type { PrismaService } from '../prisma/prisma.service'; + +const url = process.env.DATABASE_URL ?? 'postgresql://iios:iios@localhost:5434/iios?schema=public'; +const prisma = new PrismaClient({ datasources: { db: { url } } }); +const asService = prisma as unknown as PrismaService; +const actors = new ActorResolver(asService); +const ms = () => new MessageService(asService, makeFakePorts(), actors); +const alice: MessagePrincipal = { userId: 'alice', orgId: 'org_demo', appId: 'portal-demo', displayName: 'Alice' }; +const bob: MessagePrincipal = { userId: 'bob', orgId: 'org_demo', appId: 'portal-demo', displayName: 'Bob' }; + +function sub(actorId: string) { + return { scopeId: '', actorId, kind: 'webpush', endpoint: `https://push/${actorId}`, p256dh: 'k', auth: 'a' }; +} +async function seedSub(userId: string) { + const handle = await prisma.iiosSourceHandle.findFirstOrThrow({ where: { externalId: userId } }); + const actor = await prisma.iiosActorRef.findFirstOrThrow({ where: { sourceHandleId: handle.id } }); + const scope = await prisma.iiosScope.findFirstOrThrow(); + await prisma.iiosNotificationSubscription.create({ data: { ...sub(actor.id), scopeId: scope.id } }); + return actor.id; +} +async function eventsOf(type: string) { + const rows = await prisma.iiosOutboxEvent.findMany({ where: { eventType: type }, orderBy: { createdAt: 'asc' } }); + return rows.map((r) => r.cloudEvent as unknown as CloudEvent); +} + +beforeAll(async () => { await prisma.$connect(); }); +afterAll(async () => { await prisma.$disconnect(); }); +beforeEach(async () => { await resetDb(prisma); }); + +function makeProjector(presence = new PresenceService()) { + const deliver = vi.fn().mockResolvedValue('sent'); + const port = { deliver }; + const proj = new NotificationProjector(asService, new OutboxBus(), new DlqService(asService), new ProjectionCursorService(asService), presence, port as never); + return { proj, deliver, presence }; +} + +describe('NotificationProjector', () => { + it('DM: notifies the recipient (absent), never the sender', async () => { + const m = ms(); + const { threadId } = await m.openThread(null, alice, { membership: 'dm' }); + await m.addParticipant(threadId, alice, 'bob'); + await seedSub('bob'); + await m.send(threadId, alice, { content: 'hi bob' }, 'k1'); + + const { proj, deliver } = makeProjector(); + await proj.onMessageSent((await eventsOf(IIOS_EVENTS.messageSent))[0]!); + expect(deliver).toHaveBeenCalledOnce(); + expect(deliver.mock.calls[0][0].endpoint).toContain('push/'); // bob's sub + }); + + it('group without a mention: does NOT notify', async () => { + const m = ms(); + const { threadId } = await m.openThread(null, alice, { membership: 'group', creatorRole: 'ADMIN' }); + await m.addParticipant(threadId, alice, 'bob'); + await seedSub('bob'); + await m.send(threadId, alice, { content: 'hello all' }, 'k1'); + const { proj, deliver } = makeProjector(); + await proj.onMessageSent((await eventsOf(IIOS_EVENTS.messageSent))[0]!); + expect(deliver).not.toHaveBeenCalled(); + }); + + it('group WITH a mention: notifies the mentioned member', async () => { + const m = ms(); + const { threadId } = await m.openThread(null, alice, { membership: 'group', creatorRole: 'ADMIN' }); + await m.addParticipant(threadId, alice, 'bob'); + await seedSub('bob'); + await m.send(threadId, alice, { content: 'hey @bob' }, 'k1', undefined, undefined, ['bob']); + const { proj, deliver } = makeProjector(); + await proj.onMessageSent((await eventsOf(IIOS_EVENTS.messageSent))[0]!); + expect(deliver).toHaveBeenCalledOnce(); + }); + + it('presence gate: a recipient viewing the thread is NOT notified', async () => { + const m = ms(); + const { threadId } = await m.openThread(null, alice, { membership: 'dm' }); + await m.addParticipant(threadId, alice, 'bob'); + await seedSub('bob'); + await m.send(threadId, alice, { content: 'hi' }, 'k1'); + const presence = new PresenceService(); + presence.setFocus('sockB', 'bob', threadId); // bob is looking + const { proj, deliver } = makeProjector(presence); + await proj.onMessageSent((await eventsOf(IIOS_EVENTS.messageSent))[0]!); + expect(deliver).not.toHaveBeenCalled(); + }); + + it('mute gate: a muted recipient is NOT notified', async () => { + const m = ms(); + const { threadId } = await m.openThread(null, alice, { membership: 'dm' }); + await m.addParticipant(threadId, alice, 'bob'); + const bobActorId = await seedSub('bob'); + await prisma.iiosThreadParticipant.update({ where: { threadId_actorId: { threadId, actorId: bobActorId } }, data: { muted: true } }); + await m.send(threadId, alice, { content: 'hi' }, 'k1'); + const { proj, deliver } = makeProjector(); + await proj.onMessageSent((await eventsOf(IIOS_EVENTS.messageSent))[0]!); + expect(deliver).not.toHaveBeenCalled(); + }); + + it('prunes a subscription that returns "gone"', async () => { + const m = ms(); + const { threadId } = await m.openThread(null, alice, { membership: 'dm' }); + await m.addParticipant(threadId, alice, 'bob'); + await seedSub('bob'); + await m.send(threadId, alice, { content: 'hi' }, 'k1'); + const deliver = vi.fn().mockResolvedValue('gone'); + const proj = new NotificationProjector(asService, new OutboxBus(), new DlqService(asService), new ProjectionCursorService(asService), new PresenceService(), { deliver } as never); + await proj.onMessageSent((await eventsOf(IIOS_EVENTS.messageSent))[0]!); + expect(await prisma.iiosNotificationSubscription.count()).toBe(0); + }); +}); +``` + +- [ ] **Step 2: Run it, expect FAIL** (`NotificationProjector` missing). + +- [ ] **Step 3: Implement `notification.projector.ts`** (model on `inbox.projector.ts`: `claim()` idempotency + `cursor.advance`): +```ts +import { Injectable, OnModuleInit, Inject } from '@nestjs/common'; +import { CloudEvent, IIOS_EVENTS } from '@insignia/iios-contracts'; +import { PrismaService } from '../prisma/prisma.service'; +import { OutboxBus } from '../outbox/outbox.bus'; +import { DlqService } from '../outbox/dlq.service'; +import { ProjectionCursorService } from '../projection/projection-cursor.service'; +import { PresenceService } from './presence.service'; +import { NOTIFICATION_PORT, type NotificationPort } from './notification.port'; + +interface MsgData { interactionId: string; threadId: string; senderActorId: string; mentions?: string[] } + +@Injectable() +export class NotificationProjector implements OnModuleInit { + private readonly consumer = 'notification-projector'; + constructor( + private readonly prisma: PrismaService, + private readonly bus: OutboxBus, + private readonly dlq: DlqService, + private readonly cursor: ProjectionCursorService, + private readonly presence: PresenceService, + @Inject(NOTIFICATION_PORT) private readonly port: NotificationPort, + ) {} + + onModuleInit(): void { + this.dlq.registerHandler(this.consumer, (e) => this.onMessageSent(e)); + this.bus.on(IIOS_EVENTS.messageSent, (p) => void this.onMessageSent(p as CloudEvent).catch((err) => this.dlq.onConsumerFailure(this.consumer, p as CloudEvent, err))); + } + + async onMessageSent(event: CloudEvent): Promise { + if (!(await this.claim(event.id))) return; + await this.apply(event); + await this.cursor.advance(this.consumer, event); + } + + private async apply(event: CloudEvent): Promise { + const data = event.data as MsgData; + const thread = await this.prisma.iiosThread.findUnique({ where: { id: data.threadId } }); + if (!thread) return; + const membership = (thread.metadata as { membership?: string } | null)?.membership; + const interaction = await this.prisma.iiosInteraction.findUnique({ + where: { id: data.interactionId }, + include: { parts: { where: { kind: 'TEXT' }, take: 1 }, actor: { include: { sourceHandle: true } } }, + }); + const senderName = interaction?.actor?.sourceHandle?.externalId ?? 'Someone'; + const body = interaction?.parts[0]?.bodyText ?? 'sent an attachment'; + const mentions = data.mentions ?? []; + // recipients = participants except the sender + const participants = await this.prisma.iiosThreadParticipant.findMany({ + where: { threadId: data.threadId }, + include: { actor: { include: { sourceHandle: true } } }, + }); + // parent author (for reply-to-you) — one lookup + let parentAuthorActorId: string | null = null; + if (interaction?.parentInteractionId) { + const parent = await this.prisma.iiosInteraction.findUnique({ where: { id: interaction.parentInteractionId }, select: { actorId: true } }); + parentAuthorActorId = parent?.actorId ?? null; + } + + for (const p of participants) { + if (p.actorId === data.senderActorId) continue; // never the sender + const userId = p.actor?.sourceHandle?.externalId; + // gate 1: policy + const mentioned = userId != null && mentions.includes(userId); + const repliedToMe = parentAuthorActorId != null && parentAuthorActorId === p.actorId; + const eligible = membership === 'dm' || mentioned || repliedToMe; + if (!eligible) continue; + // gate 2: presence + if (userId && this.presence.isViewing(userId, data.threadId)) continue; + // gate 3: mute + if (p.muted) continue; + // dispatch to every subscription + const subs = await this.prisma.iiosNotificationSubscription.findMany({ where: { actorId: p.actorId } }); + for (const s of subs) { + const result = await this.port.deliver( + { endpoint: s.endpoint, p256dh: s.p256dh, auth: s.auth }, + { title: senderName, body, data: { threadId: data.threadId, interactionId: data.interactionId } }, + ); + if (result === 'gone') await this.prisma.iiosNotificationSubscription.delete({ where: { id: s.id } }); + } + } + } + + private async claim(eventId: string): Promise { + const res = await this.prisma.iiosProcessedEvent.createMany({ data: [{ consumerName: this.consumer, eventId }], skipDuplicates: true }); + return res.count > 0; + } +} +``` +> Note: `interaction.parentInteractionId` must be selected — add `parentInteractionId: true` isn't needed since the default include returns scalar fields; confirm the `interaction` object exposes `parentInteractionId` (it does — it's a scalar column). + +- [ ] **Step 4: Run it, expect PASS** — all 6 tests. `cd iios && pnpm vitest run src/notifications/notification.projector.spec.ts`. + +- [ ] **Step 5: Generic-safety grep** — no chat literals in the engine: +```bash +cd iios && grep -rniE "'dm'|'group'" packages/iios-service/src/notifications && echo "LEAK" || echo "clean ✓" +``` +Expected: `clean ✓` (membership is read as an opaque value from `thread.metadata`, compared to the string `'dm'` — which is the *policy* reading an opaque attribute, same as DevOpaPort; that comparison is acceptable in the notification policy gate, but keep it out of kernel messaging). If the grep flags the `membership === 'dm'` line, that is the **policy gate** and is allowed here (this file is the notification *policy plane*, not the kernel) — document it with a comment. + +- [ ] **Step 6: Commit.** +```bash +cd iios && git add -A && git -c user.email=maaz@insigniaconsultancy.com commit -m "feat(iios): NotificationProjector — policy/presence/mute gates + prune" +``` + +### Task 6: Endpoints — subscribe/unsubscribe/vapid-key + mute + +**Files:** +- Create: `packages/iios-service/src/notifications/notification.controller.ts` +- Create: `packages/iios-service/src/notifications/notification.dto.ts` +- Modify: `packages/iios-service/src/messaging/message.service.ts` (add `muteThread`) +- Modify: `packages/iios-service/src/threads/threads.controller.ts` (mute routes) + +- [ ] **Step 1: DTOs.** `notification.dto.ts`: +```ts +import { IsIn, IsObject, IsOptional, IsString } from 'class-validator'; +export class SubscribeDto { + @IsOptional() @IsIn(['webpush']) kind?: string; + @IsString() endpoint!: string; + @IsObject() keys!: { p256dh: string; auth: string }; + @IsOptional() @IsString() userAgent?: string; +} +export class UnsubscribeDto { @IsString() endpoint!: string } +``` + +- [ ] **Step 2: Controller.** `notification.controller.ts`: +```ts +import { BadRequestException, Body, Controller, Delete, Get, Headers, Post } from '@nestjs/common'; +import { PrismaService } from '../prisma/prisma.service'; +import { ActorResolver, type MessagePrincipal } from '../identity/actor.resolver'; +import { SessionVerifier } from '../platform/session.verifier'; +import { SubscribeDto, UnsubscribeDto } from './notification.dto'; + +@Controller('v1/notifications') +export class NotificationController { + constructor( + private readonly prisma: PrismaService, + private readonly actors: ActorResolver, + private readonly session: SessionVerifier, + ) {} + + @Get('vapid-public-key') + vapidKey() { return { key: process.env.VAPID_PUBLIC_KEY ?? '' }; } + + @Post('subscribe') + async subscribe(@Body() body: SubscribeDto, @Headers('authorization') auth?: string) { + const principal = this.principal(auth); + const scope = await this.actors.resolveScope(principal); + const actor = await this.actors.resolveActor(scope.id, principal); + await this.prisma.iiosNotificationSubscription.upsert({ + where: { endpoint: body.endpoint }, + create: { scopeId: scope.id, actorId: actor.id, kind: body.kind ?? 'webpush', endpoint: body.endpoint, p256dh: body.keys.p256dh, auth: body.keys.auth, userAgent: body.userAgent }, + update: { actorId: actor.id, p256dh: body.keys.p256dh, auth: body.keys.auth, lastSeenAt: new Date() }, + }); + return { ok: true }; + } + + @Delete('subscribe') + async unsubscribe(@Body() body: UnsubscribeDto) { + await this.prisma.iiosNotificationSubscription.deleteMany({ where: { endpoint: body.endpoint } }); + return { ok: true }; + } + + private principal(auth?: string): MessagePrincipal { + const token = (auth ?? '').replace(/^Bearer\s+/i, ''); + if (!token) throw new BadRequestException('Authorization bearer token is required'); + return this.session.verify(token); + } +} +``` + +- [ ] **Step 3: `muteThread` in `message.service.ts`** (append a method): +```ts +async muteThread(threadId: string, principal: MessagePrincipal, muted: boolean): Promise<{ threadId: string; muted: boolean }> { + const thread = await this.prisma.iiosThread.findUnique({ where: { id: threadId } }); + if (!thread) throw new NotFoundException('thread not found'); + const actor = await this.actors.resolveActor(thread.scopeId, principal); + await this.prisma.iiosThreadParticipant.update({ + where: { threadId_actorId: { threadId, actorId: actor.id } }, + data: { muted }, + }); + return { threadId, muted }; +} +``` + +- [ ] **Step 4: Mute routes** in `threads.controller.ts`: +```ts +@Post(':id/mute') +async mute(@Param('id') id: string, @Headers('authorization') auth?: string) { + return this.messages.muteThread(id, this.principal(auth), true); +} +@Post(':id/unmute') +async unmute(@Param('id') id: string, @Headers('authorization') auth?: string) { + return this.messages.muteThread(id, this.principal(auth), false); +} +``` + +- [ ] **Step 5: Add a service test** for `muteThread`. Append to `messaging/message.spec.ts`: +```ts +it('muteThread toggles the caller’s per-thread mute flag', async () => { + const s = gov(); + const { threadId } = await s.openThread(null, alice, { membership: 'group', creatorRole: 'ADMIN' }); + await s.muteThread(threadId, alice, true); + const aliceActor = await actorIdFor('alice'); + expect((await prisma.iiosThreadParticipant.findUniqueOrThrow({ where: { threadId_actorId: { threadId, actorId: aliceActor } } })).muted).toBe(true); + await s.muteThread(threadId, alice, false); + expect((await prisma.iiosThreadParticipant.findUniqueOrThrow({ where: { threadId_actorId: { threadId, actorId: aliceActor } } })).muted).toBe(false); +}); +``` + +- [ ] **Step 6: Run it, expect PASS.** `cd iios && pnpm vitest run src/messaging/message.spec.ts`. + +- [ ] **Step 7: Commit.** +```bash +cd iios && git add -A && git -c user.email=maaz@insigniaconsultancy.com commit -m "feat(iios): notification subscribe/unsubscribe/vapid + thread mute endpoints" +``` + +### Task 7: Module wiring + +**Files:** +- Create: `packages/iios-service/src/notifications/notification.module.ts` +- Modify: `packages/iios-service/src/app.module.ts`, `packages/iios-service/src/messaging/message.module.ts` + +- [ ] **Step 1: `message.module.ts`** — provide + export `PresenceService`: +```ts +// add to providers: PresenceService ; add to exports: PresenceService +import { PresenceService } from '../notifications/presence.service'; +``` + +- [ ] **Step 2: `notification.module.ts`:** +```ts +import { Module } from '@nestjs/common'; +import { PrismaModule } from '../prisma/prisma.module'; +import { IdentityModule } from '../identity/identity.module'; +import { OutboxModule } from '../outbox/outbox.module'; +import { ProjectionModule } from '../projection/projection.module'; +import { MessageModule } from '../messaging/message.module'; +import { NotificationController } from './notification.controller'; +import { NotificationProjector } from './notification.projector'; +import { WebPushDelivery } from './web-push.delivery'; +import { NOTIFICATION_PORT } from './notification.port'; + +@Module({ + imports: [PrismaModule, IdentityModule, OutboxModule, ProjectionModule, MessageModule], + controllers: [NotificationController], + providers: [ + NotificationProjector, + { + provide: NOTIFICATION_PORT, + useFactory: () => new WebPushDelivery( + process.env.VAPID_PUBLIC_KEY + ? { publicKey: process.env.VAPID_PUBLIC_KEY, privateKey: process.env.VAPID_PRIVATE_KEY!, subject: process.env.VAPID_SUBJECT ?? 'mailto:dev@insignia' } + : undefined, + ), + }, + ], +}) +export class NotificationModule {} +``` +> `NotificationProjector` needs `PresenceService` — it comes from `MessageModule` (exported). `DlqService`/`OutboxBus` from `OutboxModule`; `ProjectionCursorService` from `ProjectionModule`; `ActorResolver`/`SessionVerifier` (controller) from `IdentityModule`. Confirm each is exported by its module; if not, add it to that module's `exports`. + +- [ ] **Step 3: `app.module.ts`** — add `NotificationModule` to `imports` (after `MediaModule`). + +- [ ] **Step 4: Build + full suite green.** +```bash +cd iios && pnpm --filter @insignia/iios-service exec nest build && \ + docker exec iios-db psql -U iios -d postgres -c "DROP DATABASE IF EXISTS iios_test WITH (FORCE)" && \ + pnpm test && pnpm boundary +``` +Expected: BUILD OK · all tests pass (prior 192 + new presence/webpush/projector/mute) · `import-boundary: OK`. + +- [ ] **Step 5: Commit.** +```bash +cd iios && git add -A && git -c user.email=maaz@insigniaconsultancy.com commit -m "feat(iios): wire NotificationModule (projector + port + presence)" +``` + +--- + +# PHASE 3 — SDK + app (Web Push experience) + +### Task 8: Service worker + push registration (chat-web SDK layer) + +**Files:** +- Create: `chat-web/public/sw.js` +- Modify: `chat-web/src/lib/notifications.ts` (add `registerPush`, `mute`, `unmute`) + +- [ ] **Step 1: Service worker** `chat-web/public/sw.js` (served at `/sw.js`, root scope): +```js +self.addEventListener('push', (event) => { + const p = (() => { try { return event.data.json(); } catch { return { title: 'New message', body: '', data: {} }; } })(); + event.waitUntil(self.registration.showNotification(p.title, { body: p.body, tag: p.data?.threadId, data: p.data })); +}); +self.addEventListener('notificationclick', (event) => { + event.notification.close(); + const threadId = event.notification.data?.threadId; + const url = threadId ? `/t/${threadId}` : '/'; + event.waitUntil((async () => { + const all = await self.clients.matchAll({ type: 'window', includeUncontrolled: true }); + const existing = all.find((c) => 'focus' in c); + if (existing) { await existing.focus(); existing.navigate(url); } else { await self.clients.openWindow(url); } + })()); +}); +``` + +- [ ] **Step 2: `registerPush` + mute** in `chat-web/src/lib/notifications.ts` (append). Uses `IIOS_URL` from `./api`: +```ts +import { IIOS_URL } from './api'; +const authed = (t: string) => ({ 'content-type': 'application/json', authorization: `Bearer ${t}` }); + +function urlBase64ToUint8Array(b64: string): Uint8Array { + const pad = '='.repeat((4 - (b64.length % 4)) % 4); + const raw = atob((b64 + pad).replace(/-/g, '+').replace(/_/g, '/')); + return Uint8Array.from([...raw].map((c) => c.charCodeAt(0))); +} + +export async function registerPush(token: string): Promise { + if (!('serviceWorker' in navigator) || !('PushManager' in window)) return false; + if (!(await ensureNotifyPermission())) return false; + const reg = await navigator.serviceWorker.register('/sw.js'); + const { key } = await (await fetch(`${IIOS_URL}/v1/notifications/vapid-public-key`, { headers: authed(token) })).json(); + if (!key) return false; + const sub = await reg.pushManager.subscribe({ userVisibleOnly: true, applicationServerKey: urlBase64ToUint8Array(key) }); + const j = sub.toJSON(); + await fetch(`${IIOS_URL}/v1/notifications/subscribe`, { + method: 'POST', headers: authed(token), + body: JSON.stringify({ kind: 'webpush', endpoint: j.endpoint, keys: j.keys, userAgent: navigator.userAgent }), + }); + return true; +} + +export async function muteThread(token: string, threadId: string, muted: boolean): Promise { + await fetch(`${IIOS_URL}/v1/threads/${threadId}/${muted ? 'mute' : 'unmute'}`, { method: 'POST', headers: authed(token) }); +} +``` + +- [ ] **Step 3: Build check.** `cd chat-web && pnpm build` → clean. + +- [ ] **Step 4: Commit.** +```bash +cd chat-web && git add -A && git -c user.email=maaz@insigniaconsultancy.com commit -m "feat: service worker + Web Push registration (registerPush, mute)" +``` + +### Task 9: App wiring — focus signal, enable prompt, mute toggle + +**Files:** +- Modify: `chat-web/src/messages/provider.tsx` (expose `focus`), `chat-web/src/lib/message-client.ts` (add `focus`) +- Modify: `chat-web/src/components/Conversation.tsx` (emit focus on mount/visibility + mute toggle in header) +- Modify: `chat-web/src/components/Sidebar.tsx` (an "Enable notifications" button) + +- [ ] **Step 1: `message-client.ts`** — add a focus emitter: +```ts +focus(threadId: string | null): void { this.socket.emit('focus_thread', { threadId }); } +``` + +- [ ] **Step 2: Emit focus from `Conversation.tsx`.** After the existing effects, add: +```tsx +import { useSocket } from '../messages/provider'; +// inside Conversation(): const socket = useSocket(); +useEffect(() => { + socket.focus(threadId); + const onVis = () => socket.focus(document.hidden ? null : threadId); + document.addEventListener('visibilitychange', onVis); + return () => { socket.focus(null); document.removeEventListener('visibilitychange', onVis); }; +}, [socket, threadId]); +``` + +- [ ] **Step 3: Mute toggle in the conversation header** (`Conversation.tsx`, in `.pane-head`, next to add-member). Track mute from the conversation list (`convo`) — add `muted?: boolean` to `ThreadSummary` and return it from `listThreads` (IIOS: include the caller's `participant.muted` in `listThreads`; small addition). Then: +```tsx +import { muteThread } from '../lib/notifications'; +// button: + +``` +> IIOS `listThreads` change: in the participant fetch, select `muted` for the caller's row and set `muted` on each `ThreadSummary`. Add to the `ThreadSummary` interface (`message.service.ts`) `muted?: boolean` and to chat-web `types.ts`. + +- [ ] **Step 4: "Enable notifications" button** in `Sidebar.tsx` header: +```tsx +import { registerPush } from '../lib/notifications'; + +``` + +- [ ] **Step 5: Build + manual verify.** `cd chat-web && pnpm build`. Then: user B clicks **Enable notifications** (grant permission) → **fully close** B's tab → user A sends B a DM → **OS notification appears** (Web Push, app closed) → click → opens the thread. Repeat while B is *viewing* the thread → no push. Mute the thread → no push. + +- [ ] **Step 6: Commit.** +```bash +cd chat-web && git add -A && git -c user.email=maaz@insigniaconsultancy.com commit -m "feat: focus signal, enable-notifications prompt, per-thread mute toggle" +``` + +### Task 10: End-to-end verification + docs + +**Files:** +- Create: `iios/packages/iios-service/scripts/smoke-notifications.mjs` +- Modify: `iios/docs/IIOS_API_AND_SDK_GUIDE.md` (add §5.11 Notifications) + +- [ ] **Step 1: Smoke script** driving the real service with a mock push endpoint. `smoke-notifications.mjs`: sign up two dev tokens (or Supabase), create a DM, POST a fake subscription for the recipient (endpoint = a local http server you stand up that records hits), send a message, run the relay tick, assert the local endpoint received one push; then set focus (emit `focus_thread`) and assert no push. Model it on `scripts/smoke-inbox.mjs`. + +- [ ] **Step 2: Run the smoke** against a running service (`VAPID_*` set) → prints `notifications smoke: PASS`. + +- [ ] **Step 3: Docs** — add §5.11 to the API guide: the four endpoints, the payload shape, and the presence/policy/mute gates + a one-line note that the projector reacts to `message.sent`. Add `VAPID_*` to the env table. + +- [ ] **Step 4: Commit.** +```bash +cd iios && git add -A && git -c user.email=maaz@insigniaconsultancy.com commit -m "test+docs(iios): notifications smoke + API guide §5.11" +``` + +--- + +## Verification (whole feature) + +- **iios:** `pnpm test` green (presence, web-push, projector×6, mute); `pnpm boundary`; generic-safety grep shows the only `'dm'` literal in `notifications/` is the documented policy gate. +- **Manual, three states:** (1) recipient absent → push; (2) recipient viewing the thread → no push; (3) muted → no push; DM always, group only on @mention/reply. +- **App-closed:** the Web Push path fires with the tab fully closed (the real test of the feature). + +## Notes / follow-ups (out of scope, per spec §12) + +Quiet hours · global mute · per-thread level · email/FCM adapters · rich/actionable notifications · digest/batching. Multi-replica presence needs a Redis-backed `PresenceService` (v1 is in-memory, single-instance).