18 Commits

Author SHA1 Message Date
Claude eb385707fa feat(observability): OpenTelemetry tracing + log/trace correlation
CI / build (push) Successful in 3m49s
Deploy iios-service / build-deploy (push) Successful in 4m19s
Adds OTel auto-instrumentation (http, express, prisma/pg, redis, ...) exporting
OTLP to Tempo, so iios-service requests appear as distributed traces in Grafana.

observability/tracing.ts is imported FIRST in main.ts — auto-instrumentation
patches modules as they are require()d, so anything loaded earlier is never
traced. Env-driven (OTEL_SERVICE_NAME / OTEL_EXPORTER_OTLP_ENDPOINT); with the
endpoint unset the SDK never starts, so dev and tests carry zero overhead.
Health/metrics probes are excluded — they would otherwise swamp Tempo and bury
real request traces.

The existing P9 trace middleware now adopts the ACTIVE OTel trace id, so an
http.request log line and its distributed trace share one id (Loki -> Tempo
pivot in Grafana), falling back to x-trace-id / traceparent / generated.

pnpm-lock.yaml regenerated — CI and the Dockerfile both use
'pnpm install --frozen-lockfile', which would hard-fail on an unchanged lock.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-18 13:14:53 +00:00
mcp-bot 0bdae453fc ci(iios): rm -rf /tmp/kp before clone (persistent runner leaves stale dir)
Deploy iios-service / build-deploy (push) Successful in 4m8s
CI / build (push) Successful in 5m7s
The dind-builder runner is host-mode/persistent, so /tmp/kp from a prior run
survives and 'git clone ... /tmp/kp' fails with 'destination path already exists
and is not an empty directory'. Clear it first.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-10 18:14:42 +00:00
maaz519 0a8544b6fc feat(iios): notification engine — presence-gated Web Push (engine phase)
Deploy iios-service / build-deploy (push) Failing after 4m37s
CI / build (push) Successful in 5m25s
Generic notification engine: reacts to message.sent and runs three gates before
delivering — policy (DM always; group only on @mention or reply-to-you), presence
(skip if the recipient is focused on that thread), mute (per-thread). Delivery via a
swappable NotificationPort (Web Push/VAPID adapter); a 'gone' (404/410) prunes the sub.

- PresenceService + gateway `focus_thread` signal + disconnect cleanup (room membership
  != viewing, since the sidebar joins every thread room).
- IiosNotificationSubscription table + `muted` on IiosThreadParticipant (migration).
- Endpoints: GET vapid-public-key, POST/DELETE subscribe, POST threads/:id/mute|unmute;
  listThreads returns the caller's `muted`.
- DM-vs-group read from the opaque `membership` attribute in the notification POLICY only
  — kernel stays generic (grep-verified).

Tests: 13 new (3 gates + reply-to-you + prune + web-push adapter states); full suite 205
green. Verified live: module boots, subscribe stored, mute/unmute round-trips.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-09 19:16:53 +05:30
maaz519 f2ef8922ce chore: remove notifications planning docs (dropping the superpowers workflow)
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-09 18:59:02 +05:30
maaz519 48c2e6589c docs(plan): notifications implementation plan (10 tasks, TDD, 3 phases)
Phase 1 frontend quick win (desktop notif + tab-title), Phase 2 IIOS engine
(schema, PresenceService + focus signal, NotificationPort + Web Push, projector
with policy/presence/mute gates, endpoints, module wiring), Phase 3 SDK+app
(service worker, registerPush, focus emit, mute toggle, deep-link) + smoke/docs.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-09 18:52:20 +05:30
maaz519 ba2d5f4193 docs(spec): notifications design (engine in IIOS + Web Push, per-thread mute)
Approved design: IIOS notification engine (PresenceService, NotificationProjector with
policy/presence/mute gates, swappable NotificationPort → Web Push VAPID, subscription +
participant.muted data) + host-app experience (service worker, registerPush, mute toggle,
deep-link). Trigger policy: DMs always, group only on @mention or reply-to-you. Phased:
frontend quick win → IIOS engine → SDK/app.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-09 18:44:49 +05:30
maaz519 e49a71feaf docs(iios): bring API guide, DEPLOYMENT, and CEO overview up to as-built
API & SDK guide:
- §5.3 Threads & Messaging fully filled in: GET/POST /v1/threads, participants,
  my-annotations; the Message shape (senderId, attachment, annotations); the full
  socket event set (open_thread/send_message/add_participant/annotate/read/typing +
  server message/receipt/typing/annotation); reactions/pins/saves as the generic
  interaction-annotation primitive; mentions → MENTION inbox item.
- §5.4 inbox MENTION kind; SDK reactions/pins/saves + listThreads/createThread;
  vocab adds MENTION, message-part kinds, attachment kind, annotation types.

DEPLOYMENT:
- Real-IdP auth env (SUPABASE_URL/AUTH_ISSUERS via JWKS, no secret) + MEDIA_SECRET.
- Media StoragePort in topology: dev local disk → prod object storage; ⚠ local disk is
  single-instance/non-durable; bytes flow client↔storage direct (scales independently).

CEO overview: refreshed counts (192 tests, 53 tables / 19 migrations) and noted P9 has
begun — real Supabase auth (first stubbed port turned real) + generic reactions/pins/
saves/mentions/media on the kernel.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-09 01:48:29 +05:30
maaz519 ac7f790303 docs(iios): media section (endpoints + governed flow diagram), real-IdP auth, media env
- New §5.10 Media: presign-upload/download + upload/blob endpoints, attachment shape,
  governance notes (25MB + type allowlist via OPA, signed tokens, tenant fence), and a
  sequence diagram of the presign→upload→send→signed-download flow with the OPA gate.
- SDK: uploadMedia/mediaUrl; send() attachment/mentions.
- Auth §3: real OIDC (Supabase/AUTH_ISSUERS) JWKS verification alongside dev HS256.
- Env: MEDIA_DIR/MEDIA_SECRET/PUBLIC_URL/SUPABASE_URL/AUTH_ISSUERS; test count 192.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-09 01:38:37 +05:30
maaz519 24a87f6fb6 feat(iios): media sharing — presigned storage port + attachments on messages
Deploy iios-service / build-deploy (push) Failing after 3m35s
CI / build (push) Successful in 3m47s
Media the industry way: the DB stores a reference, bytes live behind a storage port.
- StoragePort + LocalDiskStorage (dev). MediaService presigns short-lived signed
  upload/download URLs (HS256 tokens) pointing at IIOS's own endpoints; the bytes
  never touch the kernel. Prod swaps STORAGE_PORT to S3/Supabase — same as auth.
- MediaController: presign-upload / PUT upload/:token (raw stream) / GET blob/:token /
  presign-download. Uploads are OPA-governed (iios.media.upload: 25 MB cap +
  image/video/audio/pdf/office allowlist); downloads are tenant-fenced by object key.
- send() + MessageDto carry an attachment (contentRef/mimeType/sizeBytes → generic
  MEDIA_REF/VOICE_REF/FILE_REF part; DTO exposes kind image|video|audio|file).

Tests: media.service.spec (6) — round-trip, oversize/type denied, oversized PUT
refused, tampered token rejected, tenant fence. Full suite 192 green. Verified live:
presign→upload→send→history→signed download round-trips the exact bytes.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-09 01:15:03 +05:30
maaz519 1256664361 feat(iios): expose senderId (stable externalId) on MessageDto
CI / build (push) Successful in 4m53s
Deploy iios-service / build-deploy (push) Failing after 8m37s
The app needs a reliable "is this message mine?" signal. senderActorId worked only
with a lazily-learned actorId, which Supabase token-refresh wipes. senderId is the
sender's externalId (email/username) — always present and comparable to the session
user — so message alignment is deterministic. Verified: senderId round-trips = email.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-07 20:58:10 +05:30
maaz519 0ecc8a4ada feat(iios): multi-issuer session verifier (many IdPs/apps → isolated appId scopes)
SessionVerifier now holds a REGISTRY of trusted OIDC issuers instead of a single
Supabase project. A token is routed by its `iss` claim to that issuer's entry,
verified against that issuer's JWKS (ES256, no secret), and stamped with that
entry's appId/orgId — so two Supabase projects / IdPs map to two isolated app
scopes on one IIOS (chat vs a future support app). App A's tokens can't reach B.

- Config: `AUTH_ISSUERS` (JSON array of { url, appId, orgId? }); the single
  `SUPABASE_URL` (+ SUPABASE_APP_ID) still works as a one-entry shorthand.
- Per-issuer JWKS cache; untrusted issuer → reject; forgery (right issuer claim,
  wrong key) → reject.
- HS256 app-token path (dev/tests) unchanged.

Tests: new session.verifier.spec (4) — routes to correct appId, second issuer →
different appId, untrusted issuer rejected, cross-issuer forgery rejected.

Note: outbox `replay.spec` has a pre-existing clock-skew flake (nextAttemptAt vs
Postgres now()) that surfaces under certain suite orderings — unrelated to this
change (fails in isolation on a clean tree; passes when another spec runs first).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-07 20:38:30 +05:30
maaz519 e956ad3cb9 feat(iios): verify real Supabase access tokens (ES256 via JWKS)
SessionVerifier gains a Supabase mode (set SUPABASE_URL): user SATs are ES256,
verified against the project's public JWKS — no shared secret. Keys are fetched at
startup (onModuleInit) and cached as PEM so verify() stays synchronous; a rotated
kid triggers a background refresh. Claims map to MessagePrincipal with userId=email
(stable, human-readable → mentions/directory keep working; RealMDM canonicalises
later). The legacy HS256 app-token path is unchanged (dev/tests untouched).

senderName now prefers the actor display name over the raw handle, so real names
show on bubbles when identity is an email.

Verified end-to-end against a live project: signup → real ES256 SAT → GET /v1/threads
200 → thread created + owned by the resolved email principal. Full suite 182 green.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-07 20:14:28 +05:30
maaz519 8c814d9b86 feat(iios): @mentions via inbox fan-out + my-annotations query (pins/saves)
Mentions ride the existing event-driven inbox — no chat parsing in the kernel:
- send() carries an OPAQUE mentions[] (userIds) into the message event; the kernel
  never parses "@". The app supplies the notify-list.
- InboxProjector fans out a MENTION inbox item (new generic inbox kind) to each
  mentioned *participant* (never the sender), idempotent per source message; reading
  the thread resolves the reader's MENTION + NEEDS_REPLY items to DONE.
- New MENTION value in the generic IiosInboxItemKind taxonomy (migration).

Pins/saves reuse the annotation primitive; new generic query powers a Saved list:
- MessageService.listMyAnnotated(principal, type) + GET /v1/threads/my-annotations
  returns the caller's annotated messages (type "save") with thread context.

Tests: 182 pass (+4: mention fan-out, resolve-on-read, saved query, opaque mentions).
Verified live: mention→inbox delivery, read→resolve, pin/save persistence.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-07 01:35:30 +05:30
maaz519 f3c4ba72b5 feat(iios): generic interaction-annotation primitive (backs emoji reactions)
Adds IiosInteractionAnnotation — an actor attaches an OPAQUE (annotationType, value)
label to an interaction. The kernel stores + aggregates them but never interprets the
strings (the chat app writes type "reaction" / value = emoji); the same primitive backs
pins/saves/flags/tags later. No chat vocabulary in kernel code — verified by grep.

- MessageService.toggleAnnotation(): governed (participant-only via new OPA rule
  iios.interaction.annotate), idempotent toggle keyed on
  (scope, interaction, actor, type, value); history DTO carries aggregated
  { type, value, users[] } groups.
- Gateway: `annotate` event → broadcasts `annotation` (refreshed user list) to the room.
- DevOpaPort: annotate allowed only for thread members (real OPA can tighten later).
- Migration add_interaction_annotations (additive; dev data untouched).

Tests: 178 pass (+3: toggle/aggregate, coexisting values, non-member denied).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-07 00:20:47 +05:30
maaz519 c4206f9809 fix(iios): thread subject on create, graceful open_thread error, isolated test DB
- openThread/createThread accept a generic `subject` (group name) — stored on the
  thread, echoed via listThreads; kernel never interprets it.
- Gateway open_thread now returns an ACK'd { error } on failure instead of throwing
  (which never acked → clients hung on "loading"). Non-member/missing-thread opens
  fail cleanly.
- Tests run against an isolated `iios_test` database (vitest globalSetup creates +
  migrates it; test.env overrides DATABASE_URL) so `pnpm test` can never TRUNCATE the
  dev database again. Verified: full suite green, dev DB row counts unchanged.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-06 19:39:19 +05:30
maaz519 1af10d0f4a feat(iios): message senderName + dev user directory
- Messages now carry senderName (resolved actorId → sourceHandle.externalId) so
  clients render usernames instead of raw actor cuids.
- GET /v1/dev/users exposes the dev directory (DEV_USERS keys) so a host app can
  validate add-by-username; real apps validate against their user store / MDM.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-06 17:11:20 +05:30
maaz519 31f7682a46 feat(iios): include participant names in listThreads (generic)
GET /v1/threads now returns participants[] (member usernames) so a host app can
render conversation titles / member lists without a second call. Still generic —
it only lists members of threads the caller belongs to.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-06 16:27:47 +05:30
maaz519 ba745bb71a feat(iios): policy-enforced membership + threaded replies + dev IdP (generic)
Enriches the platform PLANE, not the kernel:
- DevOpaPort: the dev OPA stub now evaluates a small policy table (behind the
  same opa.decide) — DM capped at 2, add-participant requires member/group-admin,
  governed self-join for membership threads. Real OPA swaps in unchanged.
- Dev IdP login (POST /v1/dev/login) issues the same JWT claims a real IdP would.
- PolicyDeniedFilter maps fail-closed denials to HTTP 403.

Generic kernel additions (no chat vocabulary — 'dm'/'group' live only as OPA
policy + an opaque thread attribute):
- MessageService.addParticipant (governed membership by userId), governed
  openThread self-join (scoped to threads with a membership attribute),
  parentInteractionId on send (reply link), and a generic listThreads.
- REST: GET /v1/threads, POST /v1/threads, POST /v1/threads/:id/participants;
  socket add_participant + membership/parentInteractionId. ensureParticipant
  gains a role.

Tests: dev-opa.port.spec + message.spec (DM cap / group admin / governed join /
listThreads / reply). smoke-membership.mjs; realtime smokes updated for governed
join. 175 unit tests + all smokes green; kernel free of dm/group literals.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-06 16:27:47 +05:30
52 changed files with 4087 additions and 119 deletions
@@ -42,6 +42,7 @@ jobs:
- name: Bump k8s-pods image tag (ArgoCD deploys)
run: |
rm -rf /tmp/kp # dind-builder is a persistent host: clear any stale checkout from a prior run
git clone --depth 1 -b main \
"https://mcp-bot:${{ secrets.K8S_PODS_TOKEN }}@git.lynkedup.cloud/platform-engineering/k8s-pods.git" /tmp/kp
cd /tmp/kp
+14 -2
View File
@@ -46,7 +46,14 @@ Every knob is an environment variable — see [`packages/iios-service/.env.examp
for the full, commented list. Highlights:
- **Secrets** (inject from a vault, never bake into the image): `DATABASE_URL`,
`APP_SECRETS` (per-app JWT signing keys), `ADAPTER_SECRETS` (webhook HMAC keys).
`APP_SECRETS` (per-app HS256 JWT keys), `ADAPTER_SECRETS` (webhook HMAC keys),
`MEDIA_SECRET` (signs presigned media upload/download URLs).
- **Auth (real IdP):** set `SUPABASE_URL` or `AUTH_ISSUERS` to trust real OIDC access
tokens, verified against the issuer's public **JWKS** (ES256) — **no secret to store**.
`AUTH_ISSUERS` is a JSON registry `[{url, appId, orgId?}]` routing many issuers → isolated
app scopes. The dev HS256 path (`APP_SECRETS`) stays for local/tests.
- **Media storage:** `MEDIA_DIR` + `PUBLIC_URL` configure the **dev** local-disk store; for
prod, bind the `StoragePort` to object storage (see topology) — the API/SDK don't change.
- **⚠️ `IIOS_DEV_TOKENS` MUST be `0`/unset in production.** It exposes `/v1/dev/*`
(unauthenticated token minting, webhook injection, chaos, retention sweep). This is the
single most important prod-hardening flag.
@@ -76,7 +83,12 @@ for the full, commented list. Highlights:
> by message `id`** (the payload always carries one).
- **Platform ports** (OPA policy, CMP consent, MDM, CRRE, SAS) — today in-process permissive
stubs (`LocalDevPorts`). For production, point these at real external services; the service
already calls them **fail-closed**.
already calls them **fail-closed**. The **session** port already verifies real OIDC tokens
(Supabase/JWKS) when configured.
- **Media storage** (`StoragePort`) — dev = **local disk** (`MEDIA_DIR`). ⚠️ Local disk is
**single-instance and non-durable**; with N>1 replicas or for persistence, bind it to shared
**object storage** (S3/R2/Supabase Storage). Bytes always flow **client ↔ storage directly**
via presigned URLs — they never transit the service — so this scales independently.
## Scaling & release strategy
+78 -8
View File
@@ -58,6 +58,8 @@ Use it on every call: `-H "authorization: Bearer <token>"`.
- **Different `orgId` = different tenant.** Mint two tokens with `org_A` / `org_B` to test isolation.
- Token TTL: 2h.
**Real IdP tokens (production path).** Beyond the dev HS256 tokens, `SessionVerifier` also verifies **real OIDC access tokens** (ES256) against a trusted issuer's public **JWKS** — set `SUPABASE_URL` (single issuer) or `AUTH_ISSUERS` (a registry mapping many issuers → app scopes). The verifier routes a token by its `iss` claim, so two projects/IdPs map to two isolated `appId` scopes on one service, and app A's tokens can't reach app B. No shared secret needed. The dev token path stays for tests/local.
### Error responses (what QA will see)
| Status | When |
|---|---|
@@ -101,12 +103,24 @@ Grouped by domain. All paths are relative to the base URL. **Auth = Bearer token
### 5.3 Threads & Messaging (native chat)
| Method | Path | Body / Headers | Returns |
|---|---|---|---|
| GET | `/v1/threads` | — | `ThreadSummary[]` — your threads (members, unread, last message, `membership`) |
| POST | `/v1/threads` | `{membership?, creatorRole?, subject?}` | `{threadId, status, history}` (201) — `membership`/`creatorRole`/`subject` are **opaque app attributes** the kernel stores but never interprets |
| POST | `/v1/threads/:id/participants` | `{userId, role?}` | `{threadId, participantCount}` (201) — **governed** (DM cap / group-admin via OPA) |
| GET | `/v1/threads/:id/messages` | — | `{threadId, messages[]}` |
| POST | `/v1/threads/:id/messages` | `{content, contentRef?}` + `idempotency-key?` header | `Message` (201) |
| POST | `/v1/threads/:id/messages` | `{content, contentRef?, mimeType?, sizeBytes?, checksumSha256?, parentInteractionId?}` + `idempotency-key?` header | `Message` (201) |
| GET | `/v1/threads/my-annotations` | `?type=save` | `{message, threadId, threadSubject}[]` — messages you annotated (e.g. saved) |
**Realtime (Socket.IO, namespace `/message`):** connect, then emit client→server events:
- `open_thread` `{threadId?|otherUserId}`, `send_message` `{threadId, content, idempotencyKey?}`, `read` `{threadId}`, `delivered` `{threadId}`.
- Server emits to the thread room: `message` (a Message) and `receipt` `{interactionId, actorId, kind: READ|DELIVERED}`.
A **`Message`** carries `{id, threadId, senderId, senderName, content, attachment?, parentInteractionId?, annotations[], traceId, createdAt}`. `senderId` is the sender's stable externalId (email/username) — the reliable "is this mine?" check. `attachment` = `{contentRef, mimeType, sizeBytes, kind}` (§5.10). `annotations` = `[{type, value, users[]}]` — the reaction/pin/save aggregate.
**Realtime (Socket.IO, namespace `/message`).** Token verified on connect (`auth.token`), then emit client→server:
- `open_thread` `{threadId?, membership?, creatorRole?, subject?}` — opens, or **creates** when no id; a **governed join** for membership threads. Returns `{threadId, status, history}` or `{error}` (acked, never a throw — clients don't hang).
- `send_message` `{threadId, content, contentRef?, mimeType?, sizeBytes?, parentInteractionId?, mentions?}``mentions` is an **opaque userId notify-list** (the app parses `@`; the kernel never does).
- `add_participant` `{threadId, userId, role?}` · `read` `{threadId, interactionId}` · `delivered` `{…}` · `typing` `{threadId}`.
- `annotate` `{threadId, interactionId, type, value}` — toggle a **generic annotation** (chat app uses `type: reaction|pin|save`; `value` = emoji, or empty).
Server → thread room: `message` (a Message), `receipt` `{interactionId, actorId, kind: READ|DELIVERED}`, `typing` `{threadId, userId}`, `annotation` `{interactionId, type, value, op: add|remove, users[], userId}`.
**Governance & primitives (all fail-closed via OPA).** Membership (`iios.thread.participant.add`, `iios.thread.join`), annotating (`iios.interaction.annotate` — participant-only), and send all gate on policy. **Reactions/pins/saves are one generic primitive** — an *interaction annotation* the kernel stores + aggregates as opaque `(type, value)` but never interprets (the same primitive backs pins/saves/flags/tags). **Mentions → Inbox:** `mentions[]` flows into the message event; the inbox projector fans out a `MENTION` inbox item to each mentioned participant (never the sender); reading the thread resolves it.
### 5.4 Inbox (work queue)
| Method | Path | Query / Body | Returns |
@@ -114,6 +128,8 @@ Grouped by domain. All paths are relative to the base URL. **Auth = Bearer token
| GET | `/v1/inbox/items` | `?state=OPEN|SNOOZED|DONE|ARCHIVED|CANCELLED|STALE` | `InboxItem[]` |
| PATCH | `/v1/inbox/items/:id` | `{state, reason?}` | `InboxItem` |
Item **kinds** include `NEEDS_REPLY` (unreplied thread activity, one per owner+thread) and `MENTION` (someone @-mentioned you; `priority: HIGH`, one per source message). Reading the thread resolves both to `DONE`.
### 5.5 Support (tickets, escalation, callbacks, queues, agents)
| Method | Path | Body / Query | Returns |
|---|---|---|---|
@@ -179,6 +195,52 @@ Grouped by domain. All paths are relative to the base URL. **Auth = Bearer token
`ScheduleMeetingDto`: `{meetingType, title, startAt, endAt?, timezone?, attendees?:[{userId, displayName?, role?, visibility?}], requestId?}`. `meetingType ∈ {ZOOM, PHONE, IN_PERSON, CALLBACK, INTERNAL}`; consent `status ∈ {UNKNOWN, GRANTED, DENIED, REVOKED}`. **Transcript/summary is `BLOCKED` until every attendee has `GRANTED` consent.**
### 5.10 Media (attachments)
The DB stores a **reference**; the bytes live behind a **storage port** (dev = local disk `MEDIA_DIR`; prod = swap to S3/R2/Supabase). Bytes go **client ↔ storage directly** via short-lived signed URLs — they never pass through the kernel.
| Method | Path | Auth | Body | Returns |
|---|---|---|---|---|
| POST | `/v1/media/presign-upload` | Bearer | `{mime, sizeBytes}` | `{objectKey, uploadUrl}`**OPA-gated** (size/type) |
| PUT | `/v1/media/upload/:token` | token in URL | raw bytes | `{objectKey, sizeBytes, checksumSha256}` |
| POST | `/v1/media/presign-download` | Bearer | `{contentRef, mime?}` | `{url}`**tenant-fenced**, signed 1h |
| GET | `/v1/media/blob/:token` | token in URL | — | streams the bytes with their `Content-Type` |
**Attaching to a message:** `send`/`send_message` accept `{contentRef, mimeType, sizeBytes, checksumSha256?}`. The stored `Message` then carries `attachment: { contentRef, mimeType, sizeBytes, kind }` where `kind ∈ {image, video, audio, file}` (a friendly view of the generic `MEDIA_REF`/`VOICE_REF`/`FILE_REF` part the kernel writes).
**Governance (fail-closed, built in):**
- **Upload policy** `iios.media.upload`**≤ 25 MB** and an allowlist (`image/*`, `video/*`, `audio/*`, `application/pdf`, common Office/text/zip). A violation → **403** *before* any bytes are sent.
- **Signed tokens** (HS256, `MEDIA_SECRET`) — upload URL lives **5 min**, download URL **1 h**; a tampered/expired token → 403.
- **Tenant fence** — object keys are prefixed with the caller's `scopeId`; a download for an object outside your scope → 403.
**The flow (with governance):**
```mermaid
sequenceDiagram
autonumber
participant C as Client (SDK uploadMedia)
participant API as IIOS Media API
participant OPA as OPA policy
participant ST as StoragePort (disk → S3)
participant K as Message kernel
C->>API: POST /v1/media/presign-upload {mime, sizeBytes}
API->>OPA: decide iios.media.upload (≤25MB? type allowed?)
alt denied
OPA-->>C: 403 (too large / type not allowed)
else allowed
API-->>C: { objectKey, uploadUrl } (signed, 5-min)
C->>ST: PUT bytes → uploadUrl (direct, not via kernel)
ST-->>C: { sizeBytes, checksumSha256 }
C->>K: send_message { contentRef, mime, size }
K-->>C: Message { attachment }
C->>API: POST /v1/media/presign-download { contentRef }
API->>API: tenant-fence (objectKey in my scope?)
API-->>C: { url } (signed, 1-hour)
C->>ST: GET url → bytes (stream, inline render / download)
end
```
---
## 6. SDK reference
@@ -191,7 +253,9 @@ import { RestClient } from '@insignia/iios-kernel-client';
const client = new RestClient({ serviceUrl: 'http://localhost:3200', token });
```
Methods (all return typed promises):
- **Messaging:** `getThreadMessages(threadId)`, `sendMessage(threadId, content, {contentRef?, idempotencyKey?})`
- **Messaging:** `listThreads()`, `createThread({membership?, creatorRole?, subject?})`, `addParticipant(threadId, userId, role?)`, `getThreadMessages(threadId)`, `sendMessage(threadId, content, {attachment?, parentInteractionId?, mentions?, idempotencyKey?})`
- **Reactions / pins / saves (annotations):** `MessageSocket.react(threadId, interactionId, emoji)` · `pin(...)` · `save(...)` (generic `annotate` under the hood); `listMyAnnotated('save')` for a cross-thread saved list. Subscribe to the `annotation` event for live updates.
- **Media:** `uploadMedia(file, {onProgress?}) → {contentRef, mimeType, sizeBytes, checksumSha256, kind}` (presign → PUT-with-progress → normalized ref); `mediaUrl(contentRef) → signed view URL` (cached, short-lived). *This plumbing is identical for every app, so it lives in the SDK; the app only renders by `kind`.*
- **Inbox:** `listInboxItems(state?)`, `patchInboxItem(id, {state, reason?})`
- **Support:** `createTicket({subject, priority?, threadId?})`, `escalate(threadId, subject?)`, `listTickets('mine'|'assigned')`, `patchTicket(id, state)`, `requestCallback({...})`, `createQueue(name)`, `joinQueue(id)`, `joinDefaultQueue()`, `setAvailability(state)`
- **Routing:** `createBinding(input)`, `listBindings()`, `simulateRoute({interactionId, originChannelType, originRef?})`, `listRouteDecisions(state?)`, `approveDecision(id)`, `denyDecision(id)`
@@ -269,7 +333,7 @@ Run against a live service (from `packages/iios-service`, `node scripts/<name>`)
| `smoke-capability.mjs` | governed egress + real HTTP provider (needs `IIOS_PROVIDER_URL_EMAIL`) |
| `smoke-tenant.mjs` | cross-tenant 403 + list isolation |
Automated unit/integration suite: `pnpm test` (114 tests). Import-boundary check: `pnpm boundary`.
Automated unit/integration suite: `pnpm test` (192 tests). Import-boundary check: `pnpm boundary`.
---
@@ -289,7 +353,11 @@ Automated unit/integration suite: `pnpm test` (114 tests). Import-boundary check
|---|---|---|
| `PORT` | `3200` | HTTP port |
| `DATABASE_URL` | `postgresql://iios:iios@localhost:5434/iios?schema=public` | Postgres |
| `APP_SECRETS` | `{"portal-demo":"dev-secret"}` | per-app JWT secrets (`{appId: secret}`) |
| `APP_SECRETS` | `{"portal-demo":"dev-secret"}` | per-app HS256 JWT secrets (`{appId: secret}`) |
| `SUPABASE_URL` / `AUTH_ISSUERS` | — | trusted OIDC issuer(s) — verify real IdP tokens (ES256) via JWKS. `AUTH_ISSUERS` is a JSON array `[{url, appId, orgId?}]` mapping issuers → app scopes; `SUPABASE_URL` is the single-issuer shorthand |
| `MEDIA_DIR` | `<tmp>/iios-media` | local media storage dir (dev `StoragePort`) |
| `MEDIA_SECRET` | `dev-media-secret` | signs media upload/download URLs |
| `PUBLIC_URL` | `http://localhost:$PORT` | base used to build presigned media URLs |
| `IIOS_DEV_TOKENS` | `0` | set `1` to enable `/v1/dev/*` |
| `ADAPTER_SECRETS` | `{}` | per-channel HMAC secrets (default `dev-adapter-secret`) |
| `IIOS_OUTBOUND_LIMIT` / `_WINDOW_MS` | `5` / `60000` | per-(channel,target) rate limit |
@@ -305,7 +373,9 @@ Automated unit/integration suite: `pnpm test` (114 tests). Import-boundary check
- **Interaction kind:** MESSAGE, EMAIL, SYSTEM_NOTICE, INBOX_WORK, SUPPORT_CASE, MEETING_REQUEST, DIGEST, SUMMARY, NOTIFICATION
- **Channel types:** WEBHOOK, EMAIL, WHATSAPP, PORTAL
- **Inbox state:** OPEN, SNOOZED, DONE, ARCHIVED, CANCELLED, STALE · **Inbox kind:** NEEDS_REPLY, NEEDS_REVIEW, NEEDS_APPROVAL, SUPPORT_UPDATE, MEETING_FOLLOWUP, DIGEST, SYSTEM_ALERT, CRM_OWNER_INTEREST
- **Inbox state:** OPEN, SNOOZED, DONE, ARCHIVED, CANCELLED, STALE · **Inbox kind:** NEEDS_REPLY, NEEDS_REVIEW, NEEDS_APPROVAL, **MENTION**, SUPPORT_UPDATE, MEETING_FOLLOWUP, DIGEST, SYSTEM_ALERT, CRM_OWNER_INTEREST
- **Message part kind:** TEXT, HTML, MARKDOWN, MEDIA_REF, FILE_REF, VOICE_REF, LOCATION, STRUCTURED_JSON · **Attachment kind (DTO):** image, video, audio, file
- **Interaction annotation (app-level, opaque to the kernel):** `type` = reaction | pin | save … ; `value` = emoji (reactions) or empty
- **Ticket state:** NEW, OPEN, PENDING, RESOLVED, CLOSED · **priority:** LOW, NORMAL, HIGH, URGENT
- **Route mode:** MANUAL, AUTOMATIC, HYBRID, SIMULATION_ONLY · **output format:** FORWARD, THREADED, DIGEST, SUMMARY, TRANSCRIPT · **decision:** ALLOW, DENY, REVIEW, SUPPRESS, SIMULATED
- **AI job:** CLASSIFY, SUMMARIZE, EXTRACT · **artifact:** CLASSIFICATION, SUMMARY, TRANSCRIPT, DIGEST, EXTRACTION · **artifact status:** PROPOSED, ACCEPTED, REJECTED, SUPERSEDED
+5 -3
View File
@@ -40,7 +40,7 @@ Today the external world (real WhatsApp, real AI models, real calendars, the rea
| **P8** | **Calendar & meetings** — schedule, attendee consent, transcript, summary, action items | `iios-meeting-web` | Support/scheduling; callback→meeting | meeting-studio (5178) |
| **P9** | *(next)* **Production hardening** — real providers, real policy/consent, scale, retention | — | Ops / platform | — |
All of it runs on **one NestJS service** (`iios-service`) with **one Postgres database** (46 tables across 9 migrations), fronted by **one low-level client** (`iios-kernel-client`) that every React SDK is built on.
All of it runs on **one NestJS service** (`iios-service`) with **one Postgres database** (53 tables across 19 migrations), fronted by **one low-level client** (`iios-kernel-client`) that every React SDK is built on.
---
@@ -159,7 +159,9 @@ Each product SDK is a handful of React hooks — a front-end dev wires the UI, t
| Calendar/Zoom providers | Simulated sync | P9 |
| Multi-tenant scale, retention, SLOs | Not yet | P9 |
**Proof it works:** 95 automated tests pass; every capability has a runnable demo and an end-to-end smoke script; the layer-boundary check enforces the architecture.
**Proof it works:** 192 automated tests pass; every capability has a runnable demo and an end-to-end smoke script; the layer-boundary check enforces the architecture.
**P9 has begun (real providers).** A production chat app (`chat-web`) now runs on IIOS with **real Supabase login** (the session port verifies real OIDC tokens via JWKS — the first stubbed port turned real), plus richer chat built generically on the kernel: **emoji reactions, pinned & saved messages, @mentions → inbox, and media sharing** (images/video/audio/docs on a swappable storage port). Each was a thin, generic addition — no chat-specific logic in the kernel — which is the reuse thesis paying off.
---
@@ -175,4 +177,4 @@ Narrative arc for the CEO: **one engine → chat → inbox → support → chann
---
*Appendix — repo facts: 11 packages, 6 demo apps, one NestJS service, 46 Postgres tables across 9 migrations (kernel → messaging → inbox → support → adapters → routing → ai → calendar). Boundary-enforced dependency law; 95 passing tests; 7 end-to-end smoke scripts.*
*Appendix — repo facts: 11 packages, 6 demo apps, one NestJS service, 53 Postgres tables across 19 migrations (kernel → messaging → inbox → support → adapters → routing → ai → calendar → annotations/mentions → media). Boundary-enforced dependency law; 192 passing tests; 8+ end-to-end smoke scripts. First real-provider swap live: Supabase auth (JWKS-verified) + a media storage port.*
+6 -1
View File
@@ -29,7 +29,11 @@
"prisma": "^6.2.1",
"reflect-metadata": "^0.2.2",
"rxjs": "^7.8.2",
"socket.io": "^4.8.3"
"socket.io": "^4.8.3",
"web-push": "^3.6.7",
"@opentelemetry/api": "^1.9.1",
"@opentelemetry/sdk-node": "^0.220.0",
"@opentelemetry/auto-instrumentations-node": "^0.78.0"
},
"devDependencies": {
"@insignia/iios-testkit": "workspace:*",
@@ -38,6 +42,7 @@
"@types/express": "^5.0.6",
"@types/jsonwebtoken": "^9.0.10",
"@types/node": "^26.0.1",
"@types/web-push": "^3.6.4",
"socket.io-client": "^4.8.3",
"typescript": "^5.7.3"
}
@@ -0,0 +1,27 @@
-- CreateTable
CREATE TABLE "IiosInteractionAnnotation" (
"id" TEXT NOT NULL,
"scopeId" TEXT NOT NULL,
"targetInteractionId" TEXT NOT NULL,
"actorId" TEXT NOT NULL,
"annotationType" TEXT NOT NULL,
"value" TEXT NOT NULL,
"createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
CONSTRAINT "IiosInteractionAnnotation_pkey" PRIMARY KEY ("id")
);
-- CreateIndex
CREATE INDEX "IiosInteractionAnnotation_targetInteractionId_idx" ON "IiosInteractionAnnotation"("targetInteractionId");
-- CreateIndex
CREATE UNIQUE INDEX "IiosInteractionAnnotation_scopeId_targetInteractionId_actor_key" ON "IiosInteractionAnnotation"("scopeId", "targetInteractionId", "actorId", "annotationType", "value");
-- AddForeignKey
ALTER TABLE "IiosInteractionAnnotation" ADD CONSTRAINT "IiosInteractionAnnotation_scopeId_fkey" FOREIGN KEY ("scopeId") REFERENCES "IiosScope"("id") ON DELETE CASCADE ON UPDATE CASCADE;
-- AddForeignKey
ALTER TABLE "IiosInteractionAnnotation" ADD CONSTRAINT "IiosInteractionAnnotation_targetInteractionId_fkey" FOREIGN KEY ("targetInteractionId") REFERENCES "IiosInteraction"("id") ON DELETE CASCADE ON UPDATE CASCADE;
-- AddForeignKey
ALTER TABLE "IiosInteractionAnnotation" ADD CONSTRAINT "IiosInteractionAnnotation_actorId_fkey" FOREIGN KEY ("actorId") REFERENCES "IiosActorRef"("id") ON DELETE RESTRICT ON UPDATE CASCADE;
@@ -0,0 +1,2 @@
-- AlterEnum
ALTER TYPE "IiosInboxItemKind" ADD VALUE 'MENTION';
@@ -0,0 +1,30 @@
-- AlterTable
ALTER TABLE "IiosThreadParticipant" ADD COLUMN "muted" BOOLEAN NOT NULL DEFAULT false;
-- CreateTable
CREATE TABLE "IiosNotificationSubscription" (
"id" TEXT NOT NULL,
"scopeId" TEXT NOT NULL,
"actorId" TEXT NOT NULL,
"kind" TEXT NOT NULL DEFAULT 'webpush',
"endpoint" TEXT NOT NULL,
"p256dh" TEXT NOT NULL,
"auth" TEXT NOT NULL,
"userAgent" TEXT,
"createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
"lastSeenAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP,
CONSTRAINT "IiosNotificationSubscription_pkey" PRIMARY KEY ("id")
);
-- CreateIndex
CREATE UNIQUE INDEX "IiosNotificationSubscription_endpoint_key" ON "IiosNotificationSubscription"("endpoint");
-- CreateIndex
CREATE INDEX "IiosNotificationSubscription_actorId_idx" ON "IiosNotificationSubscription"("actorId");
-- AddForeignKey
ALTER TABLE "IiosNotificationSubscription" ADD CONSTRAINT "IiosNotificationSubscription_scopeId_fkey" FOREIGN KEY ("scopeId") REFERENCES "IiosScope"("id") ON DELETE CASCADE ON UPDATE CASCADE;
-- AddForeignKey
ALTER TABLE "IiosNotificationSubscription" ADD CONSTRAINT "IiosNotificationSubscription_actorId_fkey" FOREIGN KEY ("actorId") REFERENCES "IiosActorRef"("id") ON DELETE RESTRICT ON UPDATE CASCADE;
@@ -79,6 +79,7 @@ enum IiosInboxItemKind {
NEEDS_REPLY
NEEDS_REVIEW
NEEDS_APPROVAL
MENTION
SUPPORT_UPDATE
MEETING_FOLLOWUP
DIGEST
@@ -303,10 +304,12 @@ model IiosScope {
channels IiosChannel[]
threads IiosThread[]
interactions IiosInteraction[]
annotations IiosInteractionAnnotation[]
inboxItems IiosInboxItem[]
supportQueues IiosSupportQueue[]
tickets IiosTicket[]
callbacks IiosCallbackRequest[]
notificationSubscriptions IiosNotificationSubscription[]
@@index([orgId, appId, tenantId])
}
@@ -348,6 +351,7 @@ model IiosActorRef {
threadsCreated IiosThread[] @relation("ThreadCreatedBy")
participations IiosThreadParticipant[]
interactions IiosInteraction[]
annotations IiosInteractionAnnotation[]
receipts IiosMessageReceipt[]
unreadCounters IiosUnreadCounter[]
inboxItemsOwned IiosInboxItem[]
@@ -355,6 +359,7 @@ model IiosActorRef {
ticketsRequested IiosTicket[] @relation("TicketRequester")
ticketsAssigned IiosTicket[] @relation("TicketAssignee")
callbacksRequested IiosCallbackRequest[]
notificationSubscriptions IiosNotificationSubscription[]
}
/// A configured channel surface (PORTAL in P1).
@@ -409,12 +414,32 @@ model IiosThreadParticipant {
actorId String
actor IiosActorRef @relation(fields: [actorId], references: [id])
participantRole String @default("MEMBER")
muted Boolean @default(false)
joinedAt DateTime @default(now())
leftAt DateTime?
@@id([threadId, actorId])
}
/// A device/browser push subscription for an actor (Web Push endpoint + keys).
/// The notification engine delivers to these when the actor is absent.
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])
}
/// The semantic envelope. Idempotent per (scope, idempotencyKey).
model IiosInteraction {
id String @id @default(cuid())
@@ -444,6 +469,7 @@ model IiosInteraction {
parts IiosMessagePart[]
receipts IiosMessageReceipt[]
annotations IiosInteractionAnnotation[]
inboxItems IiosInboxItem[]
ticketsCreatedFrom IiosTicket[]
@@ -451,6 +477,26 @@ model IiosInteraction {
@@index([threadId, occurredAt])
}
/// A generic annotation an actor attaches to an interaction. Both `annotationType`
/// and `value` are OPAQUE, app-supplied strings the kernel stores and aggregates
/// but never interprets — e.g. the chat app writes type "reaction" / value "👍".
/// The same primitive backs pins, saves, flags, tags later. No chat vocabulary here.
model IiosInteractionAnnotation {
id String @id @default(cuid())
scopeId String
scope IiosScope @relation(fields: [scopeId], references: [id], onDelete: Cascade)
targetInteractionId String
target IiosInteraction @relation(fields: [targetInteractionId], references: [id], onDelete: Cascade)
actorId String
actor IiosActorRef @relation(fields: [actorId], references: [id])
annotationType String
value String
createdAt DateTime @default(now())
@@unique([scopeId, targetInteractionId, actorId, annotationType, value])
@@index([targetInteractionId])
}
model IiosMessagePart {
id String @id @default(cuid())
interactionId String
@@ -0,0 +1,61 @@
// v1.1 smoke: dev IdP login, policy-enforced DM cap / group membership, listThreads,
// and threaded replies — all over REST. Requires the service with IIOS_DEV_TOKENS=1.
import 'dotenv/config';
const SERVICE = process.env.SMOKE_URL ?? 'http://localhost:3200';
const assert = (c, m) => { if (!c) { console.error('✗', m); process.exit(1); } console.log('✓', m); };
async function login(username, password) {
const r = await fetch(`${SERVICE}/v1/dev/login`, {
method: 'POST', headers: { 'content-type': 'application/json' },
body: JSON.stringify({ username, password }),
});
return { status: r.status, body: r.ok ? await r.json() : null };
}
function call(token, path, method = 'GET', body) {
return fetch(`${SERVICE}${path}`, {
method,
headers: { 'content-type': 'application/json', authorization: `Bearer ${token}` },
body: body ? JSON.stringify(body) : undefined,
});
}
async function json(token, path, method = 'GET', body) {
const r = await call(token, path, method, body);
if (!r.ok) throw new Error(`${method} ${path} ${r.status}: ${await r.text()}`);
return r.json();
}
// 1) dev IdP: correct password issues a token; wrong password is rejected.
const ok = await login('alice', 'alice');
assert(ok.body?.token, 'dev IdP: alice logs in with password');
const bad = await login('alice', 'nope');
assert(bad.status === 401, 'dev IdP: wrong password → 401');
const alice = ok.body.token;
// 2) DM is capped at two (policy).
const dm = await json(alice, '/v1/threads', 'POST', { membership: 'dm' });
assert(dm.threadId, `alice created a DM (${dm.threadId})`);
const add2 = await json(alice, `/v1/threads/${dm.threadId}/participants`, 'POST', { userId: 'bob' });
assert(add2.participantCount === 2, 'added bob → 2 participants');
const add3 = await call(alice, `/v1/threads/${dm.threadId}/participants`, 'POST', { userId: 'carol' });
assert(add3.status === 403, `adding a 3rd to a DM → 403 policy denied (${add3.status})`);
// 3) group: creator (ADMIN) can add several.
const grp = await json(alice, '/v1/threads', 'POST', { membership: 'group', creatorRole: 'ADMIN' });
await json(alice, `/v1/threads/${grp.threadId}/participants`, 'POST', { userId: 'bob' });
const g3 = await json(alice, `/v1/threads/${grp.threadId}/participants`, 'POST', { userId: 'carol' });
assert(g3.participantCount === 3, 'group admin added bob + carol → 3 participants');
// 4) threaded reply round-trips.
const first = await json(alice, `/v1/threads/${grp.threadId}/messages`, 'POST', { content: 'question?' });
const reply = await json(alice, `/v1/threads/${grp.threadId}/messages`, 'POST', { content: 'answer', parentInteractionId: first.id });
assert(reply.parentInteractionId === first.id, 'a reply carries parentInteractionId');
// 5) listThreads shows the caller's threads with membership + last message.
const mine = await json(alice, '/v1/threads');
const g = mine.find((t) => t.threadId === grp.threadId);
assert(g && g.membership === 'group' && g.participantCount === 3, 'GET /v1/threads lists the group with membership + count');
assert(g.lastMessage === 'answer', 'thread summary carries the last message');
console.log('\nv1.1 membership + replies smoke: PASS');
process.exit(0);
@@ -29,10 +29,12 @@ try {
await Promise.all([once(alice, 'connect'), once(bob, 'connect')]);
assert(true, `Alice connected to A (${A_URL}), Bob to B (${B_URL})`);
// Alice opens a thread on instance A; Bob joins the SAME thread on instance B.
// Alice opens a thread on instance A; membership is governed so Alice ADDS Bob, who
// then opens the SAME thread on instance B.
const opened = await alice.emitWithAck('open_thread', {});
const threadId = opened.threadId;
assert(!!threadId, `Alice opened thread ${threadId} on instance A`);
await alice.emitWithAck('add_participant', { threadId, userId: 'bob-cluster' });
await bob.emitWithAck('open_thread', { threadId });
assert(true, 'Bob joined the same thread on instance B');
@@ -44,10 +44,11 @@ const alice = connect('alice');
const bob = connect('bob');
try {
// Alice creates a thread; Bob joins it.
// Alice creates a thread; membership is governed, so Alice ADDS Bob (he can't self-join).
const opened = await alice.emitWithAck('open_thread', {});
const threadId = opened.threadId;
assert(!!threadId, `Alice created thread ${threadId}`);
await alice.emitWithAck('add_participant', { threadId, userId: 'bob' });
await bob.emitWithAck('open_thread', { threadId });
// Alice sends; Bob should receive it live.
+4
View File
@@ -9,6 +9,8 @@ import { OutboxModule } from './outbox/outbox.module';
import { ThreadsModule } from './threads/threads.module';
import { MessageModule } from './messaging/message.module';
import { InboxModule } from './inbox/inbox.module';
import { MediaModule } from './media/media.module';
import { NotificationModule } from './notifications/notification.module';
import { SupportModule } from './support/support.module';
import { AdaptersModule } from './adapters/adapters.module';
import { RoutingModule } from './routing/routing.module';
@@ -33,6 +35,8 @@ import { DevController } from './dev/dev.controller';
ThreadsModule,
MessageModule,
InboxModule,
MediaModule,
NotificationModule,
SupportModule,
AdaptersModule,
RoutingModule,
@@ -1,5 +1,5 @@
import { randomUUID } from 'node:crypto';
import { BadRequestException, Body, Controller, ForbiddenException, Headers, Param, Post } from '@nestjs/common';
import { BadRequestException, Body, Controller, ForbiddenException, Get, Headers, Param, Post, UnauthorizedException } from '@nestjs/common';
import jwt from 'jsonwebtoken';
import { IIOS_EVENTS } from '@insignia/iios-contracts';
import { hmacSign } from '@insignia/iios-adapter-sdk';
@@ -27,20 +27,58 @@ export class DevController {
@Post('token')
token(@Body() body: { appId: string; userId: string; name?: string; orgId?: string }): { token: string } {
this.assertDev();
return { token: this.signToken(body.appId, body.userId, body.name, body.orgId) };
}
/**
* Dev IdP: a credentialed login that mints the SAME JWT claims a real IdP / Session
* Broker would issue (sub/name/appId/orgId). Credentials come from DEV_USERS (JSON
* `{username: password}`; defaults to a few demo users). SessionVerifier is the stable
* verify seam — swapping in a real IdP means issuing these same claims, nothing else.
*/
@Post('login')
login(@Body() body: { username: string; password: string; appId?: string }): { token: string; userId: string } {
this.assertDev();
const appId = body.appId ?? 'portal-demo';
let users: Record<string, string>;
try {
users = JSON.parse(process.env.DEV_USERS ?? '{"alice":"alice","bob":"bob","carol":"carol"}');
} catch {
users = {};
}
const expected = users[body.username];
if (expected === undefined || expected !== body.password) throw new UnauthorizedException('invalid username or password');
return { token: this.signToken(appId, body.username, body.username), userId: body.username };
}
/** Dev directory: the known dev usernames (from DEV_USERS). A real app validates
* add-by-username against its user store / the MDM port; this stands in for it. */
@Get('users')
users(): { users: string[] } {
this.assertDev();
let map: Record<string, string>;
try {
map = JSON.parse(process.env.DEV_USERS ?? '{"alice":"alice","bob":"bob","carol":"carol"}');
} catch {
map = {};
}
return { users: Object.keys(map) };
}
private signToken(appId: string, userId: string, name?: string, orgId?: string): string {
let secrets: Record<string, string>;
try {
secrets = JSON.parse(process.env.APP_SECRETS ?? '{}');
} catch {
secrets = {};
}
const secret = secrets[body.appId];
if (!secret) throw new BadRequestException(`unknown app: ${body.appId}`);
const token = jwt.sign(
{ sub: body.userId, name: body.name ?? body.userId, appId: body.appId, orgId: body.orgId ?? `org_${body.appId}` },
const secret = secrets[appId];
if (!secret) throw new BadRequestException(`unknown app: ${appId}`);
return jwt.sign(
{ sub: userId, name: name ?? userId, appId, orgId: orgId ?? `org_${appId}` },
secret,
{ algorithm: 'HS256', expiresIn: '2h' },
);
return { token };
}
/** Server signs a sample provider payload and feeds it to the inbound pipeline. */
@@ -68,10 +68,10 @@ export class ActorResolver {
);
}
async ensureParticipant(threadId: string, actorId: string): Promise<void> {
async ensureParticipant(threadId: string, actorId: string, participantRole = 'MEMBER'): Promise<void> {
await this.prisma.iiosThreadParticipant.upsert({
where: { threadId_actorId: { threadId, actorId } },
create: { threadId, actorId },
create: { threadId, actorId, participantRole },
update: {},
});
}
@@ -10,6 +10,7 @@ interface MessageSentData {
interactionId: string;
threadId: string;
senderActorId: string;
mentions?: string[]; // opaque app notify-list (userIds); the kernel never parses "@"
}
interface MessageReadData {
threadId: string;
@@ -57,6 +58,60 @@ export class InboxProjector implements OnModuleInit {
if (p.actorId === data.senderActorId) continue;
await this.upsertNeedsReply(thread.scopeId, p.actorId, data.threadId, data.interactionId, traceId, thread.subject);
}
// Targeted mention notifications from the app-supplied opaque notify-list.
const mentions = data.mentions ?? [];
if (mentions.length > 0) {
const src = await this.prisma.iiosInteraction.findUnique({
where: { id: data.interactionId },
include: { parts: { where: { kind: 'TEXT' }, take: 1 }, actor: { include: { sourceHandle: true } } },
});
const senderName = src?.actor?.sourceHandle?.externalId ?? 'Someone';
const text = src?.parts[0]?.bodyText ?? undefined;
for (const userId of mentions) {
await this.createMention(thread.scopeId, userId, data.threadId, data.interactionId, data.senderActorId, senderName, text, traceId);
}
}
}
/** One MENTION item per (owner, source message) for a mentioned participant (never the sender). */
private async createMention(
scopeId: string,
userId: string,
threadId: string,
sourceInteractionId: string,
senderActorId: string,
senderName: string,
summary: string | undefined,
traceId: string | undefined,
): Promise<void> {
const handle = await this.prisma.iiosSourceHandle.findUnique({
where: { scopeId_kind_externalId: { scopeId, kind: 'PORTAL_USER', externalId: userId } },
});
if (!handle) return;
const actor = await this.prisma.iiosActorRef.findFirst({ where: { sourceHandleId: handle.id } });
if (!actor || actor.id === senderActorId) return; // resolvable, and never notify the sender
const isMember = await this.prisma.iiosThreadParticipant.findUnique({ where: { threadId_actorId: { threadId, actorId: actor.id } } });
if (!isMember) return; // only thread participants get mention items
const existing = await this.prisma.iiosInboxItem.findFirst({ where: { ownerActorId: actor.id, sourceInteractionId, kind: 'MENTION' } });
if (existing) return; // idempotent per source message
const item = await this.prisma.iiosInboxItem.create({
data: {
scopeId,
ownerActorId: actor.id,
kind: 'MENTION',
title: `${senderName} mentioned you`,
summary,
priority: 'HIGH',
threadId,
sourceInteractionId,
traceId,
},
});
await this.prisma.iiosInboxItemStateHistory.create({
data: { inboxItemId: item.id, toState: 'OPEN', reasonCode: 'mentioned' },
});
}
async onMessageRead(event: CloudEvent): Promise<void> {
@@ -67,21 +122,24 @@ export class InboxProjector implements OnModuleInit {
private async applyMessageRead(event: CloudEvent): Promise<void> {
const data = event.data as MessageReadData;
const item = await this.prisma.iiosInboxItem.findFirst({
// Reading the thread clears the reader's activity AND mention items for it.
const items = await this.prisma.iiosInboxItem.findMany({
where: {
threadId: data.threadId,
ownerActorId: data.actorId,
kind: 'NEEDS_REPLY',
kind: { in: ['NEEDS_REPLY', 'MENTION'] },
state: { in: ['OPEN', 'SNOOZED'] },
},
});
if (!item) return;
await this.prisma.$transaction([
this.prisma.iiosInboxItem.update({ where: { id: item.id }, data: { state: 'DONE' } }),
this.prisma.iiosInboxItemStateHistory.create({
data: { inboxItemId: item.id, fromState: item.state, toState: 'DONE', reasonCode: 'message_read' },
}),
]);
if (items.length === 0) return;
await this.prisma.$transaction(
items.flatMap((item) => [
this.prisma.iiosInboxItem.update({ where: { id: item.id }, data: { state: 'DONE' } }),
this.prisma.iiosInboxItemStateHistory.create({
data: { inboxItemId: item.id, fromState: item.state, toState: 'DONE', reasonCode: 'message_read' },
}),
]),
);
}
private async upsertNeedsReply(
@@ -88,6 +88,38 @@ describe('InboxProjector (P3 work surface)', () => {
expect(await itemsFor('bob')).toHaveLength(1);
});
it('message.sent with mentions → a MENTION item for the mentioned participant only', async () => {
const m = ms();
const { threadId } = await m.openThread(null, alice, { membership: 'group', creatorRole: 'ADMIN' });
await m.addParticipant(threadId, alice, 'bob');
await m.addParticipant(threadId, alice, 'carol');
await m.send(threadId, alice, { content: 'hey @bob look' }, 'k1', undefined, undefined, ['bob']);
await projector().onMessageSent((await eventsOf(IIOS_EVENTS.messageSent))[0]!);
const bobMentions = (await itemsFor('bob')).filter((i) => i.kind === 'MENTION');
expect(bobMentions).toHaveLength(1);
expect(bobMentions[0]?.state).toBe('OPEN');
expect(bobMentions[0]?.sourceInteractionId).toBeTruthy();
expect((await itemsFor('carol')).filter((i) => i.kind === 'MENTION')).toHaveLength(0); // not mentioned
expect((await itemsFor('alice')).filter((i) => i.kind === 'MENTION')).toHaveLength(0); // never the sender
});
it('reading the thread resolves the readers MENTION item to DONE', async () => {
const m = ms();
const { threadId } = await m.openThread(null, alice, { membership: 'group', creatorRole: 'ADMIN' });
await m.addParticipant(threadId, alice, 'bob');
const sent = await m.send(threadId, alice, { content: '@bob ping' }, 'k1', undefined, undefined, ['bob']);
const proj = projector();
await proj.onMessageSent((await eventsOf(IIOS_EVENTS.messageSent))[0]!);
expect((await itemsFor('bob')).filter((i) => i.kind === 'MENTION')[0]?.state).toBe('OPEN');
await m.markRead(threadId, bob, sent.id);
await proj.onMessageRead((await eventsOf(IIOS_EVENTS.messageRead))[0]!);
expect((await itemsFor('bob')).filter((i) => i.kind === 'MENTION')[0]?.state).toBe('DONE');
});
it('message.read → the recipient item goes DONE', async () => {
const m = ms();
const { threadId } = await m.openThread(null, alice);
+3
View File
@@ -1,3 +1,4 @@
import './observability/tracing'; // MUST be first — auto-instrumentation patches modules on require
import 'dotenv/config';
import 'reflect-metadata';
import { NestFactory } from '@nestjs/core';
@@ -5,6 +6,7 @@ import { ValidationPipe } from '@nestjs/common';
import { AppModule } from './app.module';
import { traceMiddleware } from './observability/trace.middleware';
import { RedisIoAdapter } from './realtime/redis-io.adapter';
import { PolicyDeniedFilter } from './platform/policy-denied.filter';
async function bootstrap(): Promise<void> {
const app = await NestFactory.create(AppModule, { rawBody: true });
@@ -12,6 +14,7 @@ async function bootstrap(): Promise<void> {
app.useGlobalPipes(
new ValidationPipe({ whitelist: true, forbidNonWhitelisted: true, transform: true }),
);
app.useGlobalFilters(new PolicyDeniedFilter()); // fail-closed policy denials → 403
app.enableCors({ origin: true, credentials: true });
// Multi-replica realtime: fan socket.io room emits across instances via Redis.
@@ -0,0 +1,49 @@
import { createHash } from 'node:crypto';
import { promises as fs } from 'node:fs';
import * as os from 'node:os';
import * as path from 'node:path';
import { Injectable } from '@nestjs/common';
import type { StoragePort } from './storage.port';
/**
* Dev storage: bytes on local disk under MEDIA_DIR, with a sidecar `.meta.json`
* holding the mime/size/checksum. Object keys may contain a scope subdirectory.
* A drop-in for an S3/Supabase adapter in prod.
*/
@Injectable()
export class LocalDiskStorage implements StoragePort {
private readonly root = process.env.MEDIA_DIR?.trim() || path.join(os.tmpdir(), 'iios-media');
private full(objectKey: string): string {
// Prevent path traversal: keep everything under root.
const p = path.normalize(path.join(this.root, objectKey));
if (!p.startsWith(path.normalize(this.root))) throw new Error('invalid object key');
return p;
}
async put(objectKey: string, data: Buffer, mime: string): Promise<{ sizeBytes: number; checksumSha256: string }> {
const file = this.full(objectKey);
await fs.mkdir(path.dirname(file), { recursive: true });
const checksumSha256 = createHash('sha256').update(data).digest('hex');
await fs.writeFile(file, data);
await fs.writeFile(`${file}.meta.json`, JSON.stringify({ mime, sizeBytes: data.length, checksumSha256 }));
return { sizeBytes: data.length, checksumSha256 };
}
async get(objectKey: string): Promise<{ data: Buffer; mime: string; sizeBytes: number } | null> {
const file = this.full(objectKey);
try {
const data = await fs.readFile(file);
const meta = JSON.parse(await fs.readFile(`${file}.meta.json`, 'utf8')) as { mime: string };
return { data, mime: meta.mime || 'application/octet-stream', sizeBytes: data.length };
} catch {
return null;
}
}
async remove(objectKey: string): Promise<void> {
const file = this.full(objectKey);
await fs.rm(file, { force: true });
await fs.rm(`${file}.meta.json`, { force: true });
}
}
@@ -0,0 +1,61 @@
import { BadRequestException, Body, Controller, Get, Headers, Param, Post, Put, Req, Res } from '@nestjs/common';
import type { Request, Response } from 'express';
import { MediaService } from './media.service';
import { SessionVerifier } from '../platform/session.verifier';
import { PresignDownloadDto, PresignUploadDto } from './media.dto';
import type { MessagePrincipal } from '../identity/actor.resolver';
@Controller('v1/media')
export class MediaController {
constructor(
private readonly media: MediaService,
private readonly session: SessionVerifier,
) {}
/** Authorize an upload → short-lived signed PUT url + object key. */
@Post('presign-upload')
async presignUpload(@Body() body: PresignUploadDto, @Headers('authorization') auth?: string) {
return this.media.presignUpload(this.principal(auth), { mime: body.mime, sizeBytes: body.sizeBytes });
}
/** Authorize a download → short-lived signed GET url. */
@Post('presign-download')
async presignDownload(@Body() body: PresignDownloadDto, @Headers('authorization') auth?: string) {
return this.media.presignDownload(this.principal(auth), body.contentRef, body.mime);
}
/** The presigned PUT target — token IS the auth. Reads the raw binary body stream. */
@Put('upload/:token')
async upload(@Param('token') token: string, @Req() req: Request) {
const data = await readBody(req);
if (data.length === 0) throw new BadRequestException('empty upload body');
return this.media.put(token, data);
}
/** The presigned GET target — token IS the auth. Streams the bytes with their mime. */
@Get('blob/:token')
async blob(@Param('token') token: string, @Res() res: Response) {
const { data, mime } = await this.media.get(token);
res.setHeader('Content-Type', mime);
res.setHeader('Cache-Control', 'private, max-age=3600');
res.send(data);
}
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);
}
}
/** Collect a request body stream into a Buffer (binary media upload; no body parser touches it). */
function readBody(req: Request): Promise<Buffer> {
const raw = (req as Request & { rawBody?: Buffer }).rawBody;
if (raw && raw.length) return Promise.resolve(raw);
return new Promise<Buffer>((resolve, reject) => {
const chunks: Buffer[] = [];
req.on('data', (c: Buffer) => chunks.push(c));
req.on('end', () => resolve(Buffer.concat(chunks)));
req.on('error', reject);
});
}
@@ -0,0 +1,11 @@
import { IsInt, IsOptional, IsString, Max, Min } from 'class-validator';
export class PresignUploadDto {
@IsString() mime!: string;
@IsInt() @Min(1) @Max(26 * 1024 * 1024) sizeBytes!: number;
}
export class PresignDownloadDto {
@IsString() contentRef!: string;
@IsOptional() @IsString() mime?: string;
}
@@ -0,0 +1,15 @@
import { Module } from '@nestjs/common';
import { PlatformModule } from '../platform/platform.module';
import { IdentityModule } from '../identity/identity.module';
import { MediaController } from './media.controller';
import { MediaService } from './media.service';
import { LocalDiskStorage } from './local-disk.storage';
import { STORAGE_PORT } from './storage.port';
@Module({
imports: [PlatformModule, IdentityModule],
controllers: [MediaController],
// Dev binds local disk; prod swaps STORAGE_PORT to an S3/Supabase adapter.
providers: [MediaService, { provide: STORAGE_PORT, useClass: LocalDiskStorage }],
})
export class MediaModule {}
@@ -0,0 +1,72 @@
import { describe, it, expect, beforeAll, afterAll, beforeEach } from 'vitest';
import { PrismaClient } from '@prisma/client';
import { PolicyDeniedError, type IiosPlatformPorts } from '@insignia/iios-contracts';
import { makeFakePorts } from '@insignia/iios-testkit';
import { resetDb } from '../test-utils/reset-db';
import { ActorResolver, type MessagePrincipal } from '../identity/actor.resolver';
import { DevOpaPort } from '../platform/dev-opa.port';
import { LocalDiskStorage } from './local-disk.storage';
import { MediaService } from './media.service';
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 ports = { ...makeFakePorts(), opa: new DevOpaPort() } as IiosPlatformPorts;
const svc = () => new MediaService(new LocalDiskStorage(), ports, new ActorResolver(asService));
const alice: MessagePrincipal = { userId: 'alice', orgId: 'org_demo', appId: 'portal-demo', displayName: 'Alice' };
// point storage at an isolated temp dir for the test run
process.env.MEDIA_DIR = process.env.MEDIA_DIR ?? '/tmp/iios-media-test';
beforeAll(async () => { await prisma.$connect(); });
afterAll(async () => { await prisma.$disconnect(); });
beforeEach(async () => { await resetDb(prisma); });
function tokenFrom(url: string): string {
return url.split('/').pop()!;
}
describe('MediaService (presigned local storage)', () => {
it('presign → PUT → presign-download → GET round-trips the bytes with mime + checksum', async () => {
const s = svc();
const bytes = Buffer.from('hello-image-bytes');
const { objectKey, uploadUrl } = await s.presignUpload(alice, { mime: 'image/png', sizeBytes: bytes.length });
expect(objectKey).toContain('/'); // scopeId/uuid
const put = await s.put(tokenFrom(uploadUrl), bytes);
expect(put.sizeBytes).toBe(bytes.length);
expect(put.checksumSha256).toHaveLength(64);
const { url } = await s.presignDownload(alice, objectKey, 'image/png');
const got = await s.get(tokenFrom(url));
expect(got.data.equals(bytes)).toBe(true);
expect(got.mime).toBe('image/png');
});
it('rejects an over-the-cap upload (OPA policy)', async () => {
await expect(svc().presignUpload(alice, { mime: 'image/png', sizeBytes: 26 * 1024 * 1024 })).rejects.toBeInstanceOf(PolicyDeniedError);
});
it('rejects a disallowed file type (OPA policy)', async () => {
await expect(svc().presignUpload(alice, { mime: 'application/x-msdownload', sizeBytes: 1000 })).rejects.toBeInstanceOf(PolicyDeniedError);
});
it('a PUT larger than the presigned size is refused', async () => {
const s = svc();
const { uploadUrl } = await s.presignUpload(alice, { mime: 'image/png', sizeBytes: 4 });
await expect(s.put(tokenFrom(uploadUrl), Buffer.from('way too many bytes'))).rejects.toThrow();
});
it('rejects a tampered / non-media token', async () => {
await expect(svc().get('not-a-real-token')).rejects.toThrow();
// an upload token cannot be used to download
const s = svc();
const { uploadUrl } = await s.presignUpload(alice, { mime: 'image/png', sizeBytes: 4 });
await expect(s.get(tokenFrom(uploadUrl))).rejects.toThrow();
});
it('tenant-fences downloads: an object outside my scope is forbidden', async () => {
await expect(svc().presignDownload(alice, 'some-other-scope/abc', 'image/png')).rejects.toThrow();
});
});
@@ -0,0 +1,76 @@
import { randomUUID } from 'node:crypto';
import { BadRequestException, ForbiddenException, Inject, Injectable, PayloadTooLargeException } from '@nestjs/common';
import jwt from 'jsonwebtoken';
import type { IiosPlatformPorts } from '@insignia/iios-contracts';
import { PLATFORM_PORTS } from '../platform/platform-ports';
import { decideOrThrow } from '../platform/fail-closed';
import { ActorResolver, type MessagePrincipal } from '../identity/actor.resolver';
import { STORAGE_PORT, type StoragePort } from './storage.port';
interface UploadToken { op: 'put'; objectKey: string; mime: string; maxBytes: number }
interface DownloadToken { op: 'get'; objectKey: string; mime: string }
/**
* Presigned media. IIOS never trusts the client with storage — it authorizes an
* upload (OPA: size/type/scope), then hands back a short-lived signed URL that
* points at its OWN storage endpoints. The bytes live behind the StoragePort; the
* kernel only ever stores the object key (contentRef) on a MessagePart.
*/
@Injectable()
export class MediaService {
private readonly secret = process.env.MEDIA_SECRET?.trim() || 'dev-media-secret';
private readonly publicUrl = (process.env.PUBLIC_URL?.trim() || `http://localhost:${process.env.PORT ?? 3200}`).replace(/\/$/, '');
constructor(
@Inject(STORAGE_PORT) private readonly storage: StoragePort,
@Inject(PLATFORM_PORTS) private readonly ports: IiosPlatformPorts,
private readonly actors: ActorResolver,
) {}
/** Authorize an upload and return a short-lived signed PUT url + the object key. */
async presignUpload(
principal: MessagePrincipal,
input: { mime: string; sizeBytes: number },
): Promise<{ objectKey: string; uploadUrl: string }> {
const scope = await this.actors.resolveScope(principal);
await decideOrThrow(this.ports, { action: 'iios.media.upload', scopeId: scope.id, mime: input.mime, sizeBytes: input.sizeBytes });
const objectKey = `${scope.id}/${randomUUID()}`;
const token = jwt.sign({ op: 'put', objectKey, mime: input.mime, maxBytes: input.sizeBytes } satisfies UploadToken, this.secret, { expiresIn: '5m' });
return { objectKey, uploadUrl: `${this.publicUrl}/v1/media/upload/${token}` };
}
/** Store bytes for a valid, unexpired upload token (enforcing the declared size cap). */
async put(token: string, data: Buffer): Promise<{ objectKey: string; sizeBytes: number; checksumSha256: string }> {
const t = this.verify<UploadToken>(token, 'put');
if (data.length > t.maxBytes) throw new PayloadTooLargeException('upload exceeds the presigned size');
const { sizeBytes, checksumSha256 } = await this.storage.put(t.objectKey, data, t.mime);
return { objectKey: t.objectKey, sizeBytes, checksumSha256 };
}
/** Authorize a download and return a short-lived signed GET url (tenant-fenced). */
async presignDownload(principal: MessagePrincipal, contentRef: string, mime = 'application/octet-stream'): Promise<{ url: string }> {
const scope = await this.actors.resolveScope(principal);
if (!contentRef.startsWith(`${scope.id}/`)) throw new ForbiddenException('object not in your scope');
const token = jwt.sign({ op: 'get', objectKey: contentRef, mime } satisfies DownloadToken, this.secret, { expiresIn: '1h' });
return { url: `${this.publicUrl}/v1/media/blob/${token}` };
}
/** Fetch bytes for a valid, unexpired download token. */
async get(token: string): Promise<{ data: Buffer; mime: string; sizeBytes: number }> {
const t = this.verify<DownloadToken>(token, 'get');
const obj = await this.storage.get(t.objectKey);
if (!obj) throw new BadRequestException('object not found');
return obj;
}
private verify<T extends { op: string }>(token: string, op: T['op']): T {
let payload: T;
try {
payload = jwt.verify(token, this.secret) as T;
} catch {
throw new ForbiddenException('invalid or expired media token');
}
if (payload.op !== op) throw new ForbiddenException('wrong media token');
return payload;
}
}
@@ -0,0 +1,12 @@
/**
* Byte storage seam. The kernel stores a `contentRef` (an opaque object key) on a
* MessagePart; the actual bytes live behind this port. Dev binds LocalDiskStorage;
* prod swaps to S3/R2/Supabase Storage with zero changes to the media service or app.
*/
export interface StoragePort {
put(objectKey: string, data: Buffer, mime: string): Promise<{ sizeBytes: number; checksumSha256: string }>;
get(objectKey: string): Promise<{ data: Buffer; mime: string; sizeBytes: number } | null>;
remove(objectKey: string): Promise<void>;
}
export const STORAGE_PORT = Symbol('STORAGE_PORT');
@@ -4,6 +4,7 @@ import {
ConnectedSocket,
MessageBody,
OnGatewayConnection,
OnGatewayDisconnect,
OnGatewayInit,
SubscribeMessage,
WebSocketGateway,
@@ -15,6 +16,7 @@ import { logJson } from '../observability/logger';
import { MessageService, type MessagePrincipal } from './message.service';
import { SessionVerifier } from '../platform/session.verifier';
import { OutboxBus } from '../outbox/outbox.bus';
import { PresenceService } from '../notifications/presence.service';
interface SocketState {
principal: MessagePrincipal;
@@ -27,7 +29,7 @@ interface SocketState {
* propagates messages persisted elsewhere (REST/ingest) to connected clients.
*/
@WebSocketGateway({ namespace: '/message', cors: { origin: true } })
export class MessageGateway implements OnGatewayInit, OnGatewayConnection {
export class MessageGateway implements OnGatewayInit, OnGatewayConnection, OnGatewayDisconnect {
private readonly logger = new Logger(MessageGateway.name);
private readonly emittedLocally = new Set<string>();
@@ -37,6 +39,7 @@ export class MessageGateway implements OnGatewayInit, OnGatewayConnection {
private readonly messages: MessageService,
private readonly session: SessionVerifier,
private readonly bus: OutboxBus,
private readonly presence: PresenceService,
) {}
afterInit(): void {
@@ -65,31 +68,82 @@ export class MessageGateway implements OnGatewayInit, OnGatewayConnection {
}
}
handleDisconnect(client: Socket): void {
this.presence.clearSocket(client.id);
}
/** The app reports which thread is in the foreground (or null when blurred) — presence. */
@SubscribeMessage('focus_thread')
focusThread(@ConnectedSocket() client: Socket, @MessageBody() body: { threadId: string | null }): void {
const state = client.data as SocketState | undefined;
if (!state?.principal) return;
this.presence.setFocus(client.id, state.principal.userId, body?.threadId ?? null);
}
@SubscribeMessage('open_thread')
async openThread(@ConnectedSocket() client: Socket, @MessageBody() body: { threadId?: string }) {
async openThread(@ConnectedSocket() client: Socket, @MessageBody() body: { threadId?: string; membership?: string; creatorRole?: string; subject?: string }) {
const { principal } = client.data as SocketState;
const result = await this.messages.openThread(body?.threadId ?? null, principal);
await client.join(result.threadId);
return result;
try {
const result = await this.messages.openThread(body?.threadId ?? null, principal, {
membership: body?.membership,
creatorRole: body?.creatorRole,
subject: body?.subject,
});
await client.join(result.threadId);
return result;
} catch (err) {
// Fail to an ACK'd error instead of throwing (which never acks → client hangs on "loading").
return { error: (err as Error).message ?? 'could not open the conversation' };
}
}
@SubscribeMessage('add_participant')
async addParticipant(@ConnectedSocket() client: Socket, @MessageBody() body: { threadId: string; userId: string; role?: string }) {
const { principal } = client.data as SocketState;
return this.messages.addParticipant(body.threadId, principal, body.userId, body.role);
}
@SubscribeMessage('send_message')
async sendMessage(
@ConnectedSocket() client: Socket,
@MessageBody() body: { threadId: string; content: string; contentRef?: string },
@MessageBody()
body: { threadId: string; content: string; contentRef?: string; mimeType?: string; sizeBytes?: number; checksumSha256?: string; parentInteractionId?: string; mentions?: string[] },
) {
const { principal } = client.data as SocketState;
const msg = await this.messages.send(
body.threadId,
principal,
{ content: body.content, contentRef: body.contentRef },
{ content: body.content, contentRef: body.contentRef, mimeType: body.mimeType, sizeBytes: body.sizeBytes, checksumSha256: body.checksumSha256 },
randomUUID(),
undefined,
body.parentInteractionId,
body.mentions,
);
this.emittedLocally.add(msg.id);
this.server.to(body.threadId).emit('message', msg);
return msg;
}
@SubscribeMessage('annotate')
async annotate(
@ConnectedSocket() client: Socket,
@MessageBody() body: { threadId: string; interactionId: string; type: string; value: string },
) {
const { principal } = client.data as SocketState;
const r = await this.messages.toggleAnnotation(body.interactionId, principal, body.type, body.value);
// Broadcast the refreshed user list for this (interaction,type,value) so every client re-renders.
this.server.to(r.threadId).emit('annotation', {
threadId: r.threadId,
interactionId: r.interactionId,
type: r.type,
value: r.value,
op: r.op,
users: r.users,
userId: principal.userId,
});
return r;
}
@SubscribeMessage('read')
async read(@ConnectedSocket() client: Socket, @MessageBody() body: { threadId: string; interactionId: string }) {
const { principal } = client.data as SocketState;
@@ -2,10 +2,11 @@ import { Module } from '@nestjs/common';
import { MessageService } from './message.service';
import { MessageGateway } from './message.gateway';
import { OutboxModule } from '../outbox/outbox.module';
import { PresenceService } from '../notifications/presence.service';
@Module({
imports: [OutboxModule],
providers: [MessageService, MessageGateway],
exports: [MessageService],
providers: [MessageService, MessageGateway, PresenceService],
exports: [MessageService, PresenceService],
})
export class MessageModule {}
@@ -9,16 +9,55 @@ import { ActorResolver, type MessagePrincipal } from '../identity/actor.resolver
export type { MessagePrincipal };
/** A generic annotation aggregate on a message (opaque type/value + who applied it). */
export interface AnnotationDto {
type: string;
value: string;
users: string[];
}
/** A media attachment on a message (the app renders it by `kind`). */
export interface AttachmentDto {
contentRef: string;
mimeType: string;
sizeBytes: number;
kind: 'image' | 'video' | 'audio' | 'file';
}
function mediaKind(mime: string): AttachmentDto['kind'] {
if (mime.startsWith('image/')) return 'image';
if (mime.startsWith('video/')) return 'video';
if (mime.startsWith('audio/')) return 'audio';
return 'file';
}
export interface MessageDto {
id: string;
threadId: string;
senderActorId: string;
senderId: string; // the sender's stable externalId (email/username) — reliable "is this mine?"
senderName: string;
content: string;
contentRef?: string;
attachment?: AttachmentDto;
parentInteractionId?: string;
annotations: AnnotationDto[];
traceId: string;
createdAt: Date;
}
export interface ThreadSummary {
threadId: string;
subject: string | null;
membership?: string;
participants: string[];
participantCount: number;
unread: number;
muted?: boolean;
lastMessage?: string;
lastAt?: Date;
}
export interface OpenThreadResult {
threadId: string;
status: string;
@@ -39,16 +78,29 @@ export class MessageService {
private readonly actors: ActorResolver,
) {}
/** Open an existing thread, or create one when no id is given. Joins as participant. */
async openThread(threadId: string | null, principal: MessagePrincipal): Promise<OpenThreadResult> {
/**
* Open an existing thread, or create one when no id is given. On create, an optional
* generic `membership` attribute is stamped on the thread's metadata (the app's DM/group
* hint the kernel never branches on it; policy does), and the creator joins as ADMIN
* for a group. Opening an existing thread is a GOVERNED join: only an existing member may
* re-open it (policy `iios.thread.join`) new members enter via addParticipant.
*/
async openThread(threadId: string | null, principal: MessagePrincipal, opts?: { membership?: string; creatorRole?: string; subject?: string }): Promise<OpenThreadResult> {
if (!threadId) {
await decideOrThrow(this.ports, { action: 'iios.thread.create', scope: principal });
const scope = await this.actors.resolveScope(principal);
const actor = await this.actors.resolveActor(scope.id, principal);
// `membership`/`creatorRole`/`subject` are generic, app-supplied thread attributes — the
// kernel stores/echoes them but never branches on their chat meaning (that lives in policy + app).
const thread = await this.prisma.iiosThread.create({
data: { scopeId: scope.id, createdByActorId: actor.id },
data: {
scopeId: scope.id,
createdByActorId: actor.id,
subject: opts?.subject?.trim() || undefined,
metadata: opts?.membership ? ({ membership: opts.membership } as Prisma.InputJsonValue) : undefined,
},
});
await this.actors.ensureParticipant(thread.id, actor.id);
await this.actors.ensureParticipant(thread.id, actor.id, opts?.creatorRole ?? 'MEMBER');
return { threadId: thread.id, status: thread.status, history: [] };
}
@@ -56,16 +108,122 @@ export class MessageService {
if (!thread) throw new NotFoundException('thread not found');
await decideOrThrow(this.ports, { action: 'iios.thread.read', threadId, scopeId: thread.scopeId });
const actor = await this.actors.resolveActor(thread.scopeId, principal);
const alreadyMember =
(await this.prisma.iiosThreadParticipant.findUnique({ where: { threadId_actorId: { threadId, actorId: actor.id } } })) !== null;
// Join is governed ONLY for threads that opted into a membership model (chat dm/group);
// support/inbox/generic threads (no membership attr) keep open-join, unchanged.
const membership = (thread.metadata as { membership?: string } | null)?.membership;
await decideOrThrow(this.ports, { action: 'iios.thread.join', threadId, scopeId: thread.scopeId, membership, alreadyMember });
await this.actors.ensureParticipant(threadId, actor.id);
return { threadId, status: thread.status, history: await this.history(threadId) };
}
/**
* Governed thread membership (a member adds another user by userId). Fail-closed via
* policy: the DM cap / group-admin rule is enforced by OPA (dev stub now, real later);
* the kernel only supplies generic context and adds the participant. The target's actor
* is resolved-or-created so you can add someone who hasn't logged in yet.
*/
async addParticipant(threadId: string, principal: MessagePrincipal, targetUserId: string, role = 'MEMBER'): Promise<{ threadId: string; participantCount: number }> {
const thread = await this.prisma.iiosThread.findUnique({ where: { id: threadId } });
if (!thread) throw new NotFoundException('thread not found');
const caller = await this.actors.resolveActor(thread.scopeId, principal);
const callerP = await this.prisma.iiosThreadParticipant.findUnique({ where: { threadId_actorId: { threadId, actorId: caller.id } } });
const participantCount = await this.prisma.iiosThreadParticipant.count({ where: { threadId } });
const membership = (thread.metadata as { membership?: string } | null)?.membership;
await decideOrThrow(this.ports, {
action: 'iios.thread.participant.add',
threadId,
scopeId: thread.scopeId,
membership,
participantCount,
callerRole: callerP?.participantRole,
targetUserId,
role,
});
const scope = await this.actors.resolveScope(principal);
const target = await this.actors.resolveActor(scope.id, {
userId: targetUserId,
appId: principal.appId,
orgId: principal.orgId,
tenantId: principal.tenantId,
displayName: targetUserId,
});
await this.actors.ensureParticipant(threadId, target.id, role);
return { threadId, participantCount: participantCount + 1 };
}
/** Toggle the caller's per-thread notification mute flag. */
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 };
}
/** Generic "my threads": every thread the caller participates in, with last message + unread. */
async listThreads(principal: MessagePrincipal): Promise<ThreadSummary[]> {
const scope = await this.actors.findScope(principal);
if (!scope) return [];
const actor = await this.actors.resolveActor(scope.id, principal);
const memberships = await this.prisma.iiosThreadParticipant.findMany({ where: { actorId: actor.id }, select: { threadId: true, muted: true } });
const threadIds = memberships.map((m) => m.threadId);
if (threadIds.length === 0) return [];
const mutedBy = new Map(memberships.map((m) => [m.threadId, m.muted]));
const [threads, unreads, allParts] = await Promise.all([
this.prisma.iiosThread.findMany({ where: { id: { in: threadIds } } }),
this.prisma.iiosUnreadCounter.findMany({ where: { threadId: { in: threadIds }, actorId: actor.id } }),
this.prisma.iiosThreadParticipant.findMany({
where: { threadId: { in: threadIds } },
include: { actor: { include: { sourceHandle: true } } },
}),
]);
const unreadBy = new Map(unreads.map((u) => [u.threadId, u.unreadCount]));
const membersBy = new Map<string, string[]>();
for (const p of allParts) {
const name = p.actor?.sourceHandle?.externalId ?? p.actor?.displayName ?? p.actorId;
membersBy.set(p.threadId, [...(membersBy.get(p.threadId) ?? []), name]);
}
const summaries = await Promise.all(
threads.map(async (t) => {
const last = await this.prisma.iiosInteraction.findFirst({
where: { threadId: t.id },
orderBy: { occurredAt: 'desc' },
include: { parts: { where: { kind: 'TEXT' }, take: 1 } },
});
const members = membersBy.get(t.id) ?? [];
return {
threadId: t.id,
subject: t.subject,
membership: (t.metadata as { membership?: string } | null)?.membership,
participants: members,
participantCount: members.length,
unread: unreadBy.get(t.id) ?? 0,
muted: mutedBy.get(t.id) ?? false,
lastMessage: last?.parts[0]?.bodyText ?? undefined,
lastAt: last?.occurredAt,
};
}),
);
return summaries.sort((a, b) => (b.lastAt?.getTime() ?? 0) - (a.lastAt?.getTime() ?? 0));
}
async send(
threadId: string,
principal: MessagePrincipal,
body: { content: string; contentRef?: string },
body: { content: string; contentRef?: string; mimeType?: string; sizeBytes?: number; checksumSha256?: string },
idempotencyKey: string,
traceId: string = randomUUID(),
parentInteractionId?: string,
mentions?: string[],
): Promise<MessageDto> {
const thread = await this.prisma.iiosThread.findUnique({ where: { id: threadId } });
if (!thread) throw new NotFoundException('thread not found');
@@ -75,10 +233,15 @@ export class MessageService {
const actor = await this.actors.resolveActor(thread.scopeId, principal);
await this.actors.ensureParticipant(threadId, actor.id);
// A reply links to its parent — but only if the parent is in the same thread (else ignored).
const parentRef = parentInteractionId
? (await this.prisma.iiosInteraction.findFirst({ where: { id: parentInteractionId, threadId }, select: { id: true } }))?.id
: undefined;
// Idempotency: a repeat key returns the existing message (no re-increment).
const existing = await this.prisma.iiosInteraction.findUnique({
where: { scopeId_idempotencyKey: { scopeId: thread.scopeId, idempotencyKey } },
include: { parts: { orderBy: { partIndex: 'asc' } } },
include: { parts: { orderBy: { partIndex: 'asc' } }, actor: { include: { sourceHandle: true } } },
});
if (existing) return this.toDto(existing, threadId);
@@ -93,6 +256,7 @@ export class MessageService {
idempotencyKey,
status: 'NORMALIZED',
traceId,
parentInteractionId: parentRef,
},
});
@@ -100,7 +264,18 @@ export class MessageService {
{ interactionId: created.id, partIndex: 0, kind: 'TEXT', bodyText: body.content },
];
if (body.contentRef) {
parts.push({ interactionId: created.id, partIndex: 1, kind: 'FILE_REF', contentRef: body.contentRef });
const mime = body.mimeType ?? '';
// Generic part kind from mime — the kernel stores media as opaque parts.
const kind = /^(image|video)\//.test(mime) ? 'MEDIA_REF' : /^audio\//.test(mime) ? 'VOICE_REF' : 'FILE_REF';
parts.push({
interactionId: created.id,
partIndex: 1,
kind,
contentRef: body.contentRef,
mimeType: body.mimeType,
sizeBytes: body.sizeBytes != null ? BigInt(body.sizeBytes) : undefined,
checksumSha256: body.checksumSha256,
});
}
await tx.iiosMessagePart.createMany({ data: parts });
@@ -114,7 +289,9 @@ export class MessageService {
datacontenttype: 'application/json',
traceparent: `00-${traceId.replace(/-/g, '')}-0000000000000000-01`,
insignia: { scopeSnapshotId: thread.scopeId, correlationId: traceId, idempotencyKey, dataClass: 'internal' },
data: { interactionId: created.id, threadId, senderActorId: actor.id },
// `mentions` is an OPAQUE app-supplied notify-list (userIds) — the kernel never parses
// "@"; the inbox projector generically fans out a notification to those actors.
data: { interactionId: created.id, threadId, senderActorId: actor.id, mentions: mentions ?? [] },
};
await tx.iiosOutboxEvent.create({
data: {
@@ -139,7 +316,7 @@ export class MessageService {
return tx.iiosInteraction.findUniqueOrThrow({
where: { id: created.id },
include: { parts: { orderBy: { partIndex: 'asc' } } },
include: { parts: { orderBy: { partIndex: 'asc' } }, actor: { include: { sourceHandle: true } } },
});
});
@@ -148,7 +325,7 @@ export class MessageService {
if (err instanceof Prisma.PrismaClientKnownRequestError && err.code === 'P2002') {
const winner = await this.prisma.iiosInteraction.findUniqueOrThrow({
where: { scopeId_idempotencyKey: { scopeId: thread.scopeId, idempotencyKey } },
include: { parts: { orderBy: { partIndex: 'asc' } } },
include: { parts: { orderBy: { partIndex: 'asc' } }, actor: { include: { sourceHandle: true } } },
});
return this.toDto(winner, threadId);
}
@@ -220,7 +397,7 @@ export class MessageService {
async getMessageById(id: string, principal?: MessagePrincipal): Promise<MessageDto | null> {
const i = await this.prisma.iiosInteraction.findUnique({
where: { id },
include: { parts: { orderBy: { partIndex: 'asc' } } },
include: { parts: { orderBy: { partIndex: 'asc' } }, actor: { include: { sourceHandle: true } } },
});
if (!i || !i.threadId) return null;
if (principal) await this.actors.assertOwns(principal, i.scopeId); // tenant fence (KG-02) when called on behalf of a caller
@@ -231,25 +408,162 @@ export class MessageService {
const interactions = await this.prisma.iiosInteraction.findMany({
where: { threadId },
orderBy: { occurredAt: 'asc' },
include: { parts: { orderBy: { partIndex: 'asc' } } },
include: { parts: { orderBy: { partIndex: 'asc' } }, actor: { include: { sourceHandle: true } } },
});
return interactions.map((i) => this.toDto(i, threadId));
const annotations = await this.annotationsByInteraction(interactions.map((i) => i.id));
return interactions.map((i) => this.toDto(i, threadId, annotations.get(i.id) ?? []));
}
/**
* Toggle a generic annotation (actor attaches/removes an opaque label on an interaction).
* `annotationType`/`value` are app-supplied and never interpreted here (the chat app uses
* type "reaction" + an emoji value). Governed: only a thread participant may annotate.
* Returns the refreshed user list for that (type,value) so callers can broadcast it.
*/
async toggleAnnotation(
interactionId: string,
principal: MessagePrincipal,
annotationType: string,
value: string,
): Promise<{ interactionId: string; threadId: string; type: string; value: string; op: 'add' | 'remove'; users: string[] }> {
const target = await this.prisma.iiosInteraction.findUnique({
where: { id: interactionId },
select: { id: true, threadId: true, scopeId: true },
});
if (!target || !target.threadId) throw new NotFoundException('message not found');
const actor = await this.actors.resolveActor(target.scopeId, principal);
const isMember =
(await this.prisma.iiosThreadParticipant.findUnique({ where: { threadId_actorId: { threadId: target.threadId, actorId: actor.id } } })) !== null;
await decideOrThrow(this.ports, { action: 'iios.interaction.annotate', threadId: target.threadId, scopeId: target.scopeId, isMember });
const key = {
scopeId_targetInteractionId_actorId_annotationType_value: {
scopeId: target.scopeId,
targetInteractionId: interactionId,
actorId: actor.id,
annotationType,
value,
},
};
const existing = await this.prisma.iiosInteractionAnnotation.findUnique({ where: key });
let op: 'add' | 'remove';
if (existing) {
await this.prisma.iiosInteractionAnnotation.delete({ where: { id: existing.id } });
op = 'remove';
} else {
await this.prisma.iiosInteractionAnnotation.create({
data: { scopeId: target.scopeId, targetInteractionId: interactionId, actorId: actor.id, annotationType, value },
});
op = 'add';
}
const users =
(await this.annotationsByInteraction([interactionId])).get(interactionId)?.find((a) => a.type === annotationType && a.value === value)?.users ?? [];
return { interactionId, threadId: target.threadId, type: annotationType, value, op, users };
}
/**
* Generic "my annotated interactions" every message the caller annotated with a given
* type, newest first, with thread context. The app uses type "save" for a personal
* bookmarks list (and could use "pin" for a cross-thread pinned view). Generic: the kernel
* never interprets the type string.
*/
async listMyAnnotated(
principal: MessagePrincipal,
annotationType: string,
): Promise<Array<{ message: MessageDto; threadId: string; threadSubject: string | null }>> {
const scope = await this.actors.findScope(principal);
if (!scope) return [];
const actor = await this.actors.resolveActor(scope.id, principal);
const anns = await this.prisma.iiosInteractionAnnotation.findMany({
where: { actorId: actor.id, annotationType },
orderBy: { createdAt: 'desc' },
select: { targetInteractionId: true },
});
const ids = anns.map((a) => a.targetInteractionId);
if (ids.length === 0) return [];
const [interactions, agg] = await Promise.all([
this.prisma.iiosInteraction.findMany({
where: { id: { in: ids } },
include: {
parts: { orderBy: { partIndex: 'asc' } },
actor: { include: { sourceHandle: true } },
thread: { select: { subject: true } },
},
}),
this.annotationsByInteraction(ids),
]);
const byId = new Map(interactions.map((i) => [i.id, i]));
return ids
.map((id) => byId.get(id))
.filter((i): i is NonNullable<typeof i> => Boolean(i && i.threadId))
.map((i) => ({
message: this.toDto(i, i.threadId as string, agg.get(i.id) ?? []),
threadId: i.threadId as string,
threadSubject: i.thread?.subject ?? null,
}));
}
// ─── helpers ────────────────────────────────────────────────────
/** Aggregate annotations for the given interactions into { type, value, users[] } groups. */
private async annotationsByInteraction(interactionIds: string[]): Promise<Map<string, AnnotationDto[]>> {
const map = new Map<string, AnnotationDto[]>();
if (interactionIds.length === 0) return map;
const rows = await this.prisma.iiosInteractionAnnotation.findMany({
where: { targetInteractionId: { in: interactionIds } },
include: { actor: { include: { sourceHandle: true } } },
orderBy: { createdAt: 'asc' },
});
for (const r of rows) {
const user = r.actor?.sourceHandle?.externalId ?? r.actor?.displayName ?? r.actorId;
const list = map.get(r.targetInteractionId) ?? [];
let entry = list.find((e) => e.type === r.annotationType && e.value === r.value);
if (!entry) {
entry = { type: r.annotationType, value: r.value, users: [] };
list.push(entry);
}
entry.users.push(user);
map.set(r.targetInteractionId, list);
}
for (const list of map.values()) for (const e of list) e.users.sort(); // deterministic order
return map;
}
private toDto(
interaction: { id: string; actorId: string | null; traceId: string | null; occurredAt: Date; parts: Array<{ kind: string; bodyText: string | null; contentRef: string | null }> },
interaction: {
id: string;
actorId: string | null;
actor?: { displayName: string | null; sourceHandle: { externalId: string } | null } | null;
parentInteractionId?: string | null;
traceId: string | null;
occurredAt: Date;
parts: Array<{ kind: string; bodyText: string | null; contentRef: string | null; mimeType?: string | null; sizeBytes?: bigint | null }>;
},
threadId: string,
annotations: AnnotationDto[] = [],
): MessageDto {
const text = interaction.parts.find((p) => p.kind === 'TEXT');
const file = interaction.parts.find((p) => p.contentRef);
const attachment: AttachmentDto | undefined = file?.contentRef
? {
contentRef: file.contentRef,
mimeType: file.mimeType ?? 'application/octet-stream',
sizeBytes: file.sizeBytes != null ? Number(file.sizeBytes) : 0,
kind: mediaKind(file.mimeType ?? ''),
}
: undefined;
return {
id: interaction.id,
threadId,
senderActorId: interaction.actorId ?? '',
senderId: interaction.actor?.sourceHandle?.externalId ?? '',
senderName: interaction.actor?.displayName ?? interaction.actor?.sourceHandle?.externalId ?? 'unknown',
content: text?.bodyText ?? '',
contentRef: file?.contentRef ?? undefined,
attachment,
parentInteractionId: interaction.parentInteractionId ?? undefined,
annotations,
traceId: interaction.traceId ?? '',
createdAt: interaction.occurredAt,
};
@@ -1,15 +1,20 @@
import { describe, it, expect, beforeAll, afterAll, beforeEach } from 'vitest';
import { PrismaClient } from '@prisma/client';
import { PolicyDeniedError, type IiosPlatformPorts } from '@insignia/iios-contracts';
import { resetDb } from '../test-utils/reset-db';
import { makeFakePorts } from '@insignia/iios-testkit';
import { MessageService, type MessagePrincipal } from './message.service';
import { ActorResolver } from '../identity/actor.resolver';
import { DevOpaPort } from '../platform/dev-opa.port';
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 svc = () => new MessageService(asService, makeFakePorts(), new ActorResolver(asService));
// A service whose OPA actually enforces the membership rules (real dev policy plane).
const gov = () =>
new MessageService(asService, { ...makeFakePorts(), opa: new DevOpaPort() } as IiosPlatformPorts, new ActorResolver(asService));
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' };
@@ -101,3 +106,124 @@ describe('MessageService (P2 native messaging)', () => {
expect(ce.insignia.correlationId).toBe(msg.traceId);
});
});
describe('Governed membership + replies (v1.1, policy-enforced)', () => {
it('a direct message is capped at two people (OPA policy)', async () => {
const s = gov();
const { threadId } = await s.openThread(null, alice, { membership: 'dm' });
expect((await s.addParticipant(threadId, alice, 'bob')).participantCount).toBe(2);
await expect(s.addParticipant(threadId, alice, 'carol')).rejects.toBeInstanceOf(PolicyDeniedError);
expect(await prisma.iiosThreadParticipant.count({ where: { threadId } })).toBe(2); // unchanged
});
it('group: the creator is ADMIN and can add; a plain member cannot', async () => {
const s = gov();
const { threadId } = await s.openThread(null, alice, { membership: 'group', creatorRole: 'ADMIN' });
await s.addParticipant(threadId, alice, 'bob'); // admin adds a member
expect(await prisma.iiosThreadParticipant.count({ where: { threadId } })).toBe(2);
await expect(s.addParticipant(threadId, bob, 'carol')).rejects.toBeInstanceOf(PolicyDeniedError); // bob is MEMBER
});
it('self-join is governed: a non-member cannot open a thread by id; after being added, they can', async () => {
const s = gov();
const { threadId } = await s.openThread(null, alice, { membership: 'group', creatorRole: 'ADMIN' });
await expect(s.openThread(threadId, bob)).rejects.toBeInstanceOf(PolicyDeniedError);
await s.addParticipant(threadId, alice, 'bob');
expect((await s.openThread(threadId, bob)).threadId).toBe(threadId);
});
it('listThreads returns the callers threads with membership, count, last message + unread', async () => {
const s = gov();
const { threadId } = await s.openThread(null, alice, { membership: 'group', creatorRole: 'ADMIN' });
await s.addParticipant(threadId, alice, 'bob');
await s.send(threadId, alice, { content: 'hello team' }, 'k1');
const mine = await s.listThreads(alice);
expect(mine).toHaveLength(1);
expect(mine[0]).toMatchObject({ membership: 'group', participantCount: 2, lastMessage: 'hello team' });
expect((await s.listThreads(bob))[0]?.unread).toBe(1);
});
it('a reply stores + returns parentInteractionId; a cross-thread parent is ignored', async () => {
const s = gov();
const { threadId } = await s.openThread(null, alice, { membership: 'dm' });
const first = await s.send(threadId, alice, { content: 'question?' }, 'k1');
const reply = await s.send(threadId, alice, { content: 'answer' }, 'k2', undefined, first.id);
expect(reply.parentInteractionId).toBe(first.id);
const { threadId: other } = await s.openThread(null, alice, { membership: 'dm' });
const cross = await s.send(other, alice, { content: 'x' }, 'k3', undefined, first.id);
expect(cross.parentInteractionId).toBeUndefined(); // parent not in this thread
});
});
describe('Interaction annotations (generic reactions primitive)', () => {
it('toggleAnnotation adds then removes (toggle) and aggregates users into history', async () => {
const s = gov();
const { threadId } = await s.openThread(null, alice, { membership: 'group', creatorRole: 'ADMIN' });
await s.addParticipant(threadId, alice, 'bob');
await s.openThread(threadId, bob);
const msg = await s.send(threadId, alice, { content: 'ship it' }, 'k1');
const add = await s.toggleAnnotation(msg.id, bob, 'reaction', '👍');
expect(add.op).toBe('add');
expect(add.users).toEqual(['bob']);
const add2 = await s.toggleAnnotation(msg.id, alice, 'reaction', '👍'); // alice also 👍
expect(add2.op).toBe('add');
expect(add2.users).toEqual(['alice', 'bob']); // sorted, deterministic
const rem = await s.toggleAnnotation(msg.id, bob, 'reaction', '👍'); // bob toggles off
expect(rem.op).toBe('remove');
expect(rem.users).toEqual(['alice']);
const hist = await s.history(threadId);
const m = hist.find((x) => x.id === msg.id)!;
expect(m.annotations).toEqual([{ type: 'reaction', value: '👍', users: ['alice'] }]);
});
it('different values coexist for the same actor (👍 and 🎉 both stick)', async () => {
const s = gov();
const { threadId } = await s.openThread(null, alice, { membership: 'group', creatorRole: 'ADMIN' });
const msg = await s.send(threadId, alice, { content: 'hi' }, 'k1');
await s.toggleAnnotation(msg.id, alice, 'reaction', '👍');
await s.toggleAnnotation(msg.id, alice, 'reaction', '🎉');
const m = (await s.history(threadId)).find((x) => x.id === msg.id)!;
expect(m.annotations).toEqual(
expect.arrayContaining([
{ type: 'reaction', value: '👍', users: ['alice'] },
{ type: 'reaction', value: '🎉', users: ['alice'] },
]),
);
});
it('a non-member cannot annotate (governed by policy)', async () => {
const s = gov();
const { threadId } = await s.openThread(null, alice, { membership: 'group', creatorRole: 'ADMIN' });
const msg = await s.send(threadId, alice, { content: 'secret' }, 'k1');
await expect(s.toggleAnnotation(msg.id, bob, 'reaction', '👍')).rejects.toBeInstanceOf(PolicyDeniedError);
});
it('listMyAnnotated returns the callers saved messages with thread context; others see none', async () => {
const s = gov();
const { threadId } = await s.openThread(null, alice, { membership: 'group', creatorRole: 'ADMIN', subject: 'Design' });
await s.addParticipant(threadId, alice, 'bob');
const msg = await s.send(threadId, alice, { content: 'save this' }, 'k1');
await s.toggleAnnotation(msg.id, alice, 'save', ''); // save is a personal annotation
const saved = await s.listMyAnnotated(alice, 'save');
expect(saved).toHaveLength(1);
expect(saved[0]).toMatchObject({ threadId, threadSubject: 'Design' });
expect(saved[0]?.message.content).toBe('save this');
expect(await s.listMyAnnotated(bob, 'save')).toHaveLength(0); // bob saved nothing
});
it('send carries an OPAQUE mentions[] into the message event (kernel never parses @)', async () => {
const s = svc();
const { threadId } = await s.openThread(null, alice);
await s.openThread(threadId, bob);
await s.send(threadId, alice, { content: 'hi @bob' }, 'k1', undefined, undefined, ['bob']);
const ev = await prisma.iiosOutboxEvent.findFirstOrThrow({ where: { eventType: 'com.insignia.iios.message.sent.v1' } });
const ce = ev.cloudEvent as { data: { mentions?: string[] } };
expect(ce.data.mentions).toEqual(['bob']);
});
});
@@ -0,0 +1,46 @@
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,
) {}
/** The public VAPID key the client needs to create a push subscription. */
@Get('vapid-public-key')
vapidKey(): { key: string } {
return { key: process.env.VAPID_PUBLIC_KEY ?? '' };
}
/** Store (or refresh) a push subscription for the caller. */
@Post('subscribe')
async subscribe(@Body() body: SubscribeDto, @Headers('authorization') auth?: string): Promise<{ ok: true }> {
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): Promise<{ ok: true }> {
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);
}
}
@@ -0,0 +1,12 @@
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;
}
@@ -0,0 +1,26 @@
import { Module } from '@nestjs/common';
import { OutboxModule } from '../outbox/outbox.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: [OutboxModule, MessageModule], // OutboxModule → bus/dlq; MessageModule → PresenceService. Prisma/Projection/Identity are @Global.
controllers: [NotificationController],
providers: [
NotificationProjector,
{
// Dev binds Web Push (VAPID); prod can swap this to an email/FCM adapter.
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 {}
@@ -0,0 +1,18 @@
/** The delivery seam. Dev = Web Push (VAPID); prod can swap to email/FCM/APNs. */
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');
@@ -0,0 +1,131 @@
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' };
async function actorIdFor(userId: string): Promise<string> {
const handle = await prisma.iiosSourceHandle.findFirstOrThrow({ where: { externalId: userId } });
return (await prisma.iiosActorRef.findFirstOrThrow({ where: { sourceHandleId: handle.id } })).id;
}
async function seedSub(userId: string): Promise<string> {
const actorId = await actorIdFor(userId);
const scope = await prisma.iiosScope.findFirstOrThrow();
await prisma.iiosNotificationSubscription.create({ data: { scopeId: scope.id, actorId, kind: 'webpush', endpoint: `https://push/${userId}`, p256dh: 'k', auth: 'a' } });
return actorId;
}
async function sentEvents(): Promise<CloudEvent[]> {
const rows = await prisma.iiosOutboxEvent.findMany({ where: { eventType: IIOS_EVENTS.messageSent }, orderBy: { createdAt: 'asc' } });
return rows.map((r) => r.cloudEvent as unknown as CloudEvent);
}
function makeProjector(presence = new PresenceService(), deliverResult: 'sent' | 'gone' = 'sent') {
const deliver = vi.fn().mockResolvedValue(deliverResult);
const proj = new NotificationProjector(asService, new OutboxBus(), new DlqService(asService), new ProjectionCursorService(asService), presence, { deliver } as never);
return { proj, deliver, presence };
}
beforeAll(async () => { await prisma.$connect(); });
afterAll(async () => { await prisma.$disconnect(); });
beforeEach(async () => { await resetDb(prisma); });
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 sentEvents())[0]!);
expect(deliver).toHaveBeenCalledOnce();
expect(deliver.mock.calls[0][0].endpoint).toContain('push/bob');
});
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 sentEvents())[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 sentEvents())[0]!);
expect(deliver).toHaveBeenCalledOnce();
});
it('group reply-to-you: notifies the parent author even without a mention', async () => {
const m = ms();
const bob: MessagePrincipal = { userId: 'bob', orgId: 'org_demo', appId: 'portal-demo', displayName: 'Bob' };
const { threadId } = await m.openThread(null, alice, { membership: 'group', creatorRole: 'ADMIN' });
await m.addParticipant(threadId, alice, 'bob');
await seedSub('bob');
const parent = await m.send(threadId, bob, { content: 'question?' }, 'k1'); // bob authors the parent
await m.send(threadId, alice, { content: 'answer' }, 'k2', undefined, parent.id); // alice replies to bob (no mention)
const { proj, deliver } = makeProjector();
const evs = await sentEvents();
await proj.onMessageSent(evs[evs.length - 1]!); // project the reply
expect(deliver).toHaveBeenCalledOnce();
expect(deliver.mock.calls[0][0].endpoint).toContain('push/bob');
});
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);
const { proj, deliver } = makeProjector(presence);
await proj.onMessageSent((await sentEvents())[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 sentEvents())[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 { proj } = makeProjector(new PresenceService(), 'gone');
await proj.onMessageSent((await sentEvents())[0]!);
expect(await prisma.iiosNotificationSubscription.count()).toBe(0);
});
});
@@ -0,0 +1,114 @@
import { Inject, Injectable, OnModuleInit } 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[];
}
/**
* Turns `message.sent` into push notifications for absent recipients. Three gates:
* 1. policy DM always; group only if @mentioned or a reply to that recipient
* 2. presence skip if the recipient is currently focused on that thread
* 3. mute skip if the recipient muted the thread
* Then dispatches to each of the recipient's subscriptions; a 'gone' result prunes it.
* Idempotent per event id (same pattern as InboxProjector). Generic: DM-vs-group is read
* from the opaque `membership` thread attribute here in the notification *policy*, not the kernel.
*/
@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<void> {
if (!(await this.claim(event.id))) return; // duplicate → no apply, no cursor advance
await this.apply(event);
await this.cursor.advance(this.consumer, event);
}
private async apply(event: CloudEvent): Promise<void> {
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 ?? [];
// 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;
}
const participants = await this.prisma.iiosThreadParticipant.findMany({
where: { threadId: data.threadId },
include: { actor: { include: { sourceHandle: true } } },
});
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;
if (!(membership === 'dm' || mentioned || repliedToMe)) continue;
// gate 2 — presence
if (userId && this.presence.isViewing(userId, data.threadId)) continue;
// gate 3 — mute
if (p.muted) continue;
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<boolean> {
const res = await this.prisma.iiosProcessedEvent.createMany({
data: [{ consumerName: this.consumer, eventId }],
skipDuplicates: true,
});
return res.count > 0;
}
}
@@ -0,0 +1,24 @@
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 actors 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);
});
});
@@ -0,0 +1,25 @@
import { Injectable } from '@nestjs/common';
/**
* In-memory focus tracker (single instance). Keyed on userId (the stable externalId the
* gateway has as principal.userId). NOTE: room membership viewing the sidebar joins
* every thread room for live updates, so presence uses an explicit `focus_thread` signal.
* Prod (multi-replica): back this with Redis.
*/
@Injectable()
export class PresenceService {
private readonly focus = new Map<string, { userId: string; threadId: string | null }>(); // 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;
}
}
@@ -0,0 +1,34 @@
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' } };
const vapid = { publicKey: 'p', privateKey: 'k', subject: 'mailto:x' };
describe('WebPushDelivery', () => {
it('returns "sent" on success', async () => {
const send = vi.fn().mockResolvedValue({ statusCode: 201 });
const d = new WebPushDelivery(vapid, 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 d410 = new WebPushDelivery(vapid, vi.fn().mockRejectedValue({ statusCode: 410 }) as never);
expect(await d410.deliver(sub, payload)).toBe('gone');
const d404 = new WebPushDelivery(vapid, vi.fn().mockRejectedValue({ statusCode: 404 }) as never);
expect(await d404.deliver(sub, payload)).toBe('gone');
});
it('returns "failed" on other errors', async () => {
const d = new WebPushDelivery(vapid, vi.fn().mockRejectedValue({ statusCode: 500 }) as never);
expect(await d.deliver(sub, payload)).toBe('failed');
});
it('is disabled (returns "failed", never calls send) without VAPID keys', async () => {
const send = vi.fn();
const d = new WebPushDelivery(undefined, send as never);
expect(await d.deliver(sub, payload)).toBe('failed');
expect(send).not.toHaveBeenCalled();
});
});
@@ -0,0 +1,40 @@
import { Injectable } from '@nestjs/common';
import webpush from 'web-push';
import type { NotificationPayload, NotificationPort, PushSub } from './notification.port';
type SendFn = typeof webpush.sendNotification;
export interface Vapid {
publicKey: string;
privateKey: string;
subject: string;
}
/** Web Push (VAPID) delivery. A dead subscription (404/410) → 'gone' so the caller prunes it. */
@Injectable()
export class WebPushDelivery implements NotificationPort {
private readonly vapid?: Vapid;
constructor(
vapid?: Vapid,
private readonly send: SendFn = webpush.sendNotification,
) {
// Store VAPID; pass it per-send (not global setVapidDetails) so construction never
// validates keys — keeps the adapter unit-testable with a mocked transport.
this.vapid = vapid?.publicKey ? vapid : undefined;
}
async deliver(sub: PushSub, payload: NotificationPayload): Promise<'sent' | 'gone' | 'failed'> {
if (!this.vapid) return 'failed'; // push disabled (no VAPID keys)
try {
await this.send(
{ endpoint: sub.endpoint, keys: { p256dh: sub.p256dh, auth: sub.auth } },
JSON.stringify(payload),
{ vapidDetails: this.vapid },
);
return 'sent';
} catch (e) {
const code = (e as { statusCode?: number }).statusCode;
return code === 404 || code === 410 ? 'gone' : 'failed';
}
}
}
@@ -1,4 +1,5 @@
import type { Request, Response, NextFunction } from 'express';
import { trace } from '@opentelemetry/api';
import { runWithTrace, newTraceId, traceIdFromTraceparent } from './trace-context';
import { logJson } from './logger';
@@ -9,8 +10,13 @@ import { logJson } from './logger';
* the response finishes (method, path, status, ms). Wired via app.use() in main.ts.
*/
export function traceMiddleware(req: Request, res: Response, next: NextFunction): void {
// Prefer the ACTIVE OpenTelemetry trace id: that is what makes this log line
// joinable to its distributed trace in Grafana (Loki -> Tempo). Falls back to
// the inbound header, then traceparent, then a generated id, so a correlation
// id is always present even when tracing is disabled.
const otelTraceId = trace.getActiveSpan()?.spanContext().traceId;
const headerTrace = (req.headers['x-trace-id'] as string | undefined)?.trim();
const traceId = headerTrace || traceIdFromTraceparent(req.headers['traceparent'] as string | undefined) || newTraceId();
const traceId = otelTraceId || headerTrace || traceIdFromTraceparent(req.headers['traceparent'] as string | undefined) || newTraceId();
res.setHeader('x-trace-id', traceId);
const startedAt = Date.now();
runWithTrace(traceId, () => {
@@ -0,0 +1,50 @@
/**
* OpenTelemetry bootstrap. MUST be imported before anything else in main.ts
* auto-instrumentation works by patching modules (http, express, pg, redis, )
* as they are require()d, so any module loaded before this runs is never traced.
*
* Env-driven on purpose: the OTel SDK reads OTEL_SERVICE_NAME,
* OTEL_EXPORTER_OTLP_ENDPOINT and OTEL_EXPORTER_OTLP_PROTOCOL itself, so the
* collector target is a deployment concern rather than a code change. When
* OTEL_EXPORTER_OTLP_ENDPOINT is unset the SDK never starts, so local dev and
* tests run with zero tracing overhead and no exporter errors.
*
* Traces land in Tempo and are viewable in Grafana. The existing P9 trace
* middleware adopts the active OTel trace id, so an `http.request` log line and
* its distributed trace share one id (Loki -> Tempo pivot).
*/
import { NodeSDK } from '@opentelemetry/sdk-node';
import { getNodeAutoInstrumentations } from '@opentelemetry/auto-instrumentations-node';
const endpoint = process.env.OTEL_EXPORTER_OTLP_ENDPOINT;
if (endpoint) {
const sdk = new NodeSDK({
instrumentations: [
getNodeAutoInstrumentations({
// Noisy and low value: every file read becomes a span.
'@opentelemetry/instrumentation-fs': { enabled: false },
// k8s probes hit /health constantly; tracing them would swamp Tempo
// and bury the real request traces.
'@opentelemetry/instrumentation-http': {
ignoreIncomingRequestHook: (req) => {
const url = req.url ?? '';
return ['/health', '/ready', '/healthz', '/readyz', '/metrics'].some((p) =>
url.startsWith(p),
);
},
},
}),
],
});
sdk.start();
const shutdown = (): void => {
// Flush buffered spans before exit, else the last requests before a
// rollout are lost.
void sdk.shutdown().finally(() => process.exit(0));
};
process.on('SIGTERM', shutdown);
process.on('SIGINT', shutdown);
}
@@ -0,0 +1,45 @@
import { describe, it, expect } from 'vitest';
import { DevOpaPort } from './dev-opa.port';
const opa = new DevOpaPort();
describe('DevOpaPort (dev policy plane — membership rules)', () => {
it('allows every existing action by default (thread create/read, message send)', async () => {
for (const action of ['iios.thread.create', 'iios.thread.read', 'iios.message.send', undefined]) {
expect((await opa.decide({ action })).allow).toBe(true);
}
});
it('DM is capped at two: adding a 2nd is allowed, a 3rd is denied', async () => {
const add = (participantCount: number) =>
opa.decide({ action: 'iios.thread.participant.add', membership: 'dm', callerRole: 'MEMBER', participantCount });
expect((await add(1)).allow).toBe(true); // 1 → adding the 2nd
const third = await add(2); // 2 → adding a 3rd
expect(third.allow).toBe(false);
expect(third.obligations[0]?.reason).toMatch(/two people/);
});
it('adding a participant requires the caller be a member', async () => {
const d = await opa.decide({ action: 'iios.thread.participant.add', membership: 'group', callerRole: 'NONE', participantCount: 3 });
expect(d.allow).toBe(false);
expect(d.obligations[0]?.reason).toMatch(/member/);
});
it('group add/remove requires ADMIN', async () => {
const asMember = await opa.decide({ action: 'iios.thread.participant.add', membership: 'group', callerRole: 'MEMBER', participantCount: 3 });
expect(asMember.allow).toBe(false);
expect(asMember.obligations[0]?.reason).toMatch(/admin/);
const asAdmin = await opa.decide({ action: 'iios.thread.participant.add', membership: 'group', callerRole: 'ADMIN', participantCount: 3 });
expect(asAdmin.allow).toBe(true);
});
it('self-join is governed only on membership threads', async () => {
// generic / support thread (no membership attr) → open join, unchanged
expect((await opa.decide({ action: 'iios.thread.join', alreadyMember: false })).allow).toBe(true);
// a membership (chat) thread → existing member allowed, stranger denied
expect((await opa.decide({ action: 'iios.thread.join', membership: 'group', alreadyMember: true })).allow).toBe(true);
const stranger = await opa.decide({ action: 'iios.thread.join', membership: 'group', alreadyMember: false });
expect(stranger.allow).toBe(false);
expect(stranger.obligations[0]?.reason).toMatch(/not a member/);
});
});
@@ -0,0 +1,66 @@
import type { PolicyDecision } from '@insignia/iios-contracts';
/**
* Dev stand-in for the real OPA policy plane. It evaluates a small **policy table**
* against the generic `{ action, ...context }` the kernel passes to `opa.decide`.
* Everything defaults to ALLOW (so existing actions are unchanged); only the rules
* below deny. This is the *policy plane*, NOT the kernel the chat meaning of "dm"
* vs "group" lives here as policy, and a real OPA drops in behind the same interface.
*/
export interface OpaInput {
action?: string;
membership?: string; // app-set generic thread attribute: 'dm' | 'group'
participantCount?: number;
callerRole?: string; // 'MEMBER' | 'ADMIN'
alreadyMember?: boolean;
isMember?: boolean;
mime?: string;
sizeBytes?: number;
[k: string]: unknown;
}
const MEDIA_MAX_BYTES = 25 * 1024 * 1024;
const MEDIA_ALLOWED = /^(image|video|audio)\//;
const MEDIA_ALLOWED_DOCS = new Set([
'application/pdf', 'text/plain', 'application/zip',
'application/msword', 'application/vnd.openxmlformats-officedocument.wordprocessingml.document',
'application/vnd.ms-excel', 'application/vnd.openxmlformats-officedocument.spreadsheetml.sheet',
'application/vnd.ms-powerpoint', 'application/vnd.openxmlformats-officedocument.presentationml.presentation',
]);
const allow = (): PolicyDecision => ({ decisionRef: 'dev-allow', allow: true, obligations: [], ttlSeconds: 60 });
const deny = (reason: string): PolicyDecision => ({ decisionRef: 'dev-deny', allow: false, obligations: [{ kind: 'AUDIT', reason }], ttlSeconds: 0 });
export class DevOpaPort {
async decide(input: unknown): Promise<PolicyDecision> {
const i = (input ?? {}) as OpaInput;
switch (i.action) {
case 'iios.thread.participant.add': {
const count = i.participantCount ?? 0;
if (i.membership === 'dm' && count >= 2) return deny('a direct message is limited to two people');
if (i.callerRole !== 'MEMBER' && i.callerRole !== 'ADMIN') return deny('only a member can add participants');
if (i.membership === 'group' && i.callerRole !== 'ADMIN') return deny('only a group admin can add or remove members');
return allow();
}
case 'iios.thread.join': {
// Only threads that opted into a membership model (chat dm/group) are governed;
// generic/support threads (no membership attr) keep open-join. For a governed
// thread, new members must be added first (governed add-participant).
if (!i.membership) return allow();
return i.alreadyMember ? allow() : deny('you are not a member of this thread');
}
case 'iios.interaction.annotate': {
// Annotating (e.g. reacting to) a message requires being a participant of its thread.
return i.isMember ? allow() : deny('you are not a member of this thread');
}
case 'iios.media.upload': {
const mime = i.mime ?? '';
if ((i.sizeBytes ?? 0) > MEDIA_MAX_BYTES) return deny('file too large (max 25 MB)');
if (!MEDIA_ALLOWED.test(mime) && !MEDIA_ALLOWED_DOCS.has(mime)) return deny(`file type not allowed: ${mime || 'unknown'}`);
return allow();
}
default:
return allow();
}
}
}
@@ -1,10 +1,10 @@
import type {
IiosPlatformPorts,
PlatformPrincipal,
PolicyDecision,
ContextDecisionBundle,
SourceHandleResolution,
} from '@insignia/iios-contracts';
import { DevOpaPort } from './dev-opa.port';
/** DI token for the platform ports (Session/OPA/CMP/MDM/CRRE/SAS/Capability). */
export const PLATFORM_PORTS = Symbol('PLATFORM_PORTS');
@@ -22,11 +22,8 @@ export class LocalDevPorts implements IiosPlatformPorts {
return { principalRef: 'local-dev', orgId: 'org_demo', appId: 'portal-demo', assurance: 'service' };
},
};
opa = {
async decide(_input: unknown): Promise<PolicyDecision> {
return { decisionRef: 'local-allow', allow: true, obligations: [], ttlSeconds: 60 };
},
};
// Dev policy plane: default-allow + the membership rules (real OPA swaps in here).
opa = new DevOpaPort();
cmp = {
async checkPurpose(_input: unknown): Promise<{ receiptRef?: string; status: 'ALLOW' | 'DENY' | 'NOT_REQUIRED' }> {
return { status: 'NOT_REQUIRED' };
@@ -0,0 +1,16 @@
import { ArgumentsHost, Catch, HttpStatus, type ExceptionFilter } from '@nestjs/common';
import { PolicyDeniedError } from '@insignia/iios-contracts';
import type { Response } from 'express';
/**
* Maps a fail-closed `PolicyDeniedError` (from `decideOrThrow`) to HTTP 403 instead of a
* generic 500, so callers can tell "policy denied" (with the reason) from a server fault.
* Applies to every policy-gated REST endpoint (e.g. the DM-cap / membership rules).
*/
@Catch(PolicyDeniedError)
export class PolicyDeniedFilter implements ExceptionFilter {
catch(err: PolicyDeniedError, host: ArgumentsHost): void {
const res = host.switchToHttp().getResponse<Response>();
res.status(HttpStatus.FORBIDDEN).json({ statusCode: 403, code: err.code, message: err.message });
}
}
@@ -0,0 +1,75 @@
import { describe, it, expect, beforeAll, afterAll, vi } from 'vitest';
import { generateKeyPairSync, type KeyObject } from 'node:crypto';
import jwt from 'jsonwebtoken';
import { SessionVerifier } from './session.verifier';
// Two independent fake IdPs (each its own EC signing key + kid).
function makeIssuer(kid: string) {
const { publicKey, privateKey } = generateKeyPairSync('ec', { namedCurve: 'P-256' });
const jwk = { ...(publicKey.export({ format: 'jwk' }) as Record<string, unknown>), kid, use: 'sig', alg: 'ES256' };
return { privateKey, jwks: { keys: [jwk] } };
}
const chat = makeIssuer('chat-kid');
const support = makeIssuer('support-kid');
const CHAT_URL = 'https://chat-proj.example.co';
const SUPPORT_URL = 'https://support-proj.example.co';
function sat(priv: KeyObject, url: string, kid: string, email: string, name: string): string {
return jwt.sign(
{ email, user_metadata: { full_name: name }, role: 'authenticated' },
priv,
{ algorithm: 'ES256', issuer: `${url}/auth/v1`, audience: 'authenticated', subject: 'uuid-x', keyid: kid, expiresIn: '1h' },
);
}
describe('SessionVerifier — multi-issuer (OIDC/JWKS registry)', () => {
let verifier: SessionVerifier;
beforeAll(async () => {
process.env.AUTH_ISSUERS = JSON.stringify([
{ url: CHAT_URL, appId: 'portal-demo' },
{ url: SUPPORT_URL, appId: 'support-app' },
]);
// Serve each issuer's JWKS from its own well-known URL.
vi.stubGlobal(
'fetch',
vi.fn(async (u: string) => {
const s = String(u);
const body = s.includes('chat-proj') ? chat.jwks : s.includes('support-proj') ? support.jwks : { keys: [] };
return { ok: true, json: async () => body } as unknown as Response;
}),
);
verifier = new SessionVerifier();
await verifier.onModuleInit();
});
afterAll(() => {
vi.unstubAllGlobals();
delete process.env.AUTH_ISSUERS;
});
it('routes a token to the appId of its own issuer', () => {
expect(verifier.verify(sat(chat.privateKey, CHAT_URL, 'chat-kid', 'Alice@x.com', 'Alice'))).toMatchObject({
userId: 'alice@x.com',
appId: 'portal-demo',
displayName: 'Alice',
});
});
it('a second issuer maps to a DIFFERENT appId — isolated scope', () => {
expect(verifier.verify(sat(support.privateKey, SUPPORT_URL, 'support-kid', 'Bob@x.com', 'Bob'))).toMatchObject({
userId: 'bob@x.com',
appId: 'support-app',
});
});
it('rejects a token from an untrusted issuer', () => {
const rogue = makeIssuer('rogue-kid');
expect(() => verifier.verify(sat(rogue.privateKey, 'https://evil.example.co', 'rogue-kid', 'e@x.com', 'E'))).toThrow();
});
it('rejects a forgery: signed by issuer B but claiming issuer A + As kid', () => {
const forged = sat(support.privateKey, CHAT_URL, 'chat-kid', 'x@x.com', 'X'); // wrong key for the claimed issuer
expect(() => verifier.verify(forged)).toThrow();
});
});
@@ -1,30 +1,148 @@
import { Injectable, UnauthorizedException } from '@nestjs/common';
import { createPublicKey } from 'node:crypto';
import { Injectable, OnModuleInit, UnauthorizedException } from '@nestjs/common';
import jwt from 'jsonwebtoken';
import type { MessagePrincipal } from '../messaging/message.service';
/** One trusted OIDC issuer (a Supabase project / IdP) → the app scope it maps to. */
interface IssuerEntry {
issuer: string; // the token `iss` claim, e.g. https://<ref>.supabase.co/auth/v1
jwksUrl: string;
appId: string;
orgId: string;
audience: string;
kidToPem: Map<string, string>;
}
/**
* The real `session` port for P2: verifies a host app's HS256 token (reuses the
* support-service AppTokenVerifier pattern). Per-app secrets come from the
* APP_SECRETS env (JSON map keyed by appId). Expected claims: { sub, name?,
* appId, orgId?, tenantId? }.
* The `session` port. Verifies whatever token a caller presents:
*
* 1. OIDC / JWKS (real IdPs) a REGISTRY of trusted issuers (env `AUTH_ISSUERS`,
* or the single `SUPABASE_URL` shorthand). A token is routed by its `iss` claim
* to that issuer's entry, verified against that issuer's public JWKS (ES256, no
* secret), and stamped with that entry's `appId`/`orgId`. So two projects/IdPs
* map to two isolated app scopes on one IIOS app A's tokens can't reach app B.
*
* 2. App token (dev / HS256) legacy per-app secrets from `APP_SECRETS`, keyed by
* the `appId` claim. Unchanged; used by the dev IdP + tests.
*
* A Session Broker would later collapse case 1 to a single issuer (the Broker's PAT).
*/
@Injectable()
export class SessionVerifier {
private readonly secrets: Record<string, string>;
export class SessionVerifier implements OnModuleInit {
private readonly appSecrets: Record<string, string>;
private readonly issuers = new Map<string, IssuerEntry>(); // keyed by issuer string
constructor() {
try {
this.secrets = JSON.parse(process.env.APP_SECRETS ?? '{}');
this.appSecrets = JSON.parse(process.env.APP_SECRETS ?? '{}');
} catch {
this.secrets = {};
this.appSecrets = {};
}
for (const cfg of this.readIssuerConfig()) {
const url = this.normalize(cfg.url);
const appId = cfg.appId?.trim() || 'portal-demo';
const issuer = cfg.issuer?.trim() || `${url}/auth/v1`;
const jwksUrl = cfg.jwksUrl?.trim() || `${url}/auth/v1/.well-known/jwks.json`;
this.issuers.set(issuer, {
issuer,
jwksUrl,
appId,
orgId: cfg.orgId?.trim() || `org_${appId}`,
audience: cfg.audience?.trim() || 'authenticated',
kidToPem: new Map(),
});
}
}
async onModuleInit(): Promise<void> {
await Promise.all([...this.issuers.values()].map((e) => this.refreshJwks(e)));
}
/** Assemble the issuer registry from AUTH_ISSUERS (multi) or SUPABASE_URL (single, back-compat). */
private readIssuerConfig(): Array<{ url: string; appId?: string; orgId?: string; issuer?: string; jwksUrl?: string; audience?: string }> {
const out: Array<{ url: string; appId?: string; orgId?: string; issuer?: string; jwksUrl?: string; audience?: string }> = [];
try {
const raw = process.env.AUTH_ISSUERS;
if (raw) {
const parsed = JSON.parse(raw) as Array<{ url: string; appId?: string; orgId?: string; issuer?: string; jwksUrl?: string; audience?: string }>;
if (Array.isArray(parsed)) out.push(...parsed.filter((e) => e && (e.url || e.issuer)));
}
} catch {
/* ignore malformed AUTH_ISSUERS */
}
const single = process.env.SUPABASE_URL?.trim();
if (single && !out.length) {
out.push({ url: single, appId: process.env.SUPABASE_APP_ID?.trim(), orgId: process.env.SUPABASE_ORG_ID?.trim() });
}
return out;
}
/** Fetch one issuer's JWKS and cache each key as PEM (so verify() stays synchronous). */
private async refreshJwks(entry: IssuerEntry): Promise<void> {
try {
const res = await fetch(entry.jwksUrl);
if (!res.ok) return;
const { keys } = (await res.json()) as { keys: Array<Record<string, unknown>> };
const next = new Map<string, string>();
for (const jwk of keys ?? []) {
const kid = jwk.kid as string | undefined;
if (!kid) continue;
try {
next.set(kid, createPublicKey({ key: jwk as never, format: 'jwk' }).export({ type: 'spki', format: 'pem' }) as string);
} catch {
/* skip a key we can't import */
}
}
if (next.size) entry.kidToPem = next;
} catch {
/* keep whatever keys we already have */
}
}
verify(token: string): MessagePrincipal {
const decoded = jwt.decode(token) as jwt.JwtPayload | null;
const appId = decoded?.appId ? String(decoded.appId) : undefined;
const decoded = jwt.decode(token, { complete: true }) as { header?: { alg?: string; kid?: string }; payload?: jwt.JwtPayload } | null;
if (!decoded?.payload) throw new UnauthorizedException('invalid token');
if (this.issuers.size && decoded.header?.alg === 'ES256') {
const iss = decoded.payload.iss ? String(decoded.payload.iss) : '';
const entry = this.issuers.get(iss);
if (!entry) throw new UnauthorizedException(`untrusted issuer: ${iss || '(none)'}`);
return this.verifyOidc(token, decoded.header.kid, entry);
}
return this.verifyAppToken(token, decoded.payload);
}
/** Verify an ES256 SAT against its issuer's JWKS key, mapping to that issuer's app scope. */
private verifyOidc(token: string, kid: string | undefined, entry: IssuerEntry): MessagePrincipal {
const pem = kid ? entry.kidToPem.get(kid) : undefined;
if (!pem) {
void this.refreshJwks(entry); // rotated / not loaded — pull fresh for next time
throw new UnauthorizedException('unknown signing key — retry');
}
let payload: jwt.JwtPayload;
try {
payload = jwt.verify(token, pem, { algorithms: ['ES256'], issuer: entry.issuer, audience: entry.audience }) as jwt.JwtPayload;
} catch {
throw new UnauthorizedException('invalid token');
}
const email = payload.email ? String(payload.email).toLowerCase() : undefined;
const meta = (payload.user_metadata ?? {}) as { full_name?: string; name?: string };
// userId = email (stable + human-readable → mentions/directory work); RealMDM canonicalises `sub` later.
const userId = email ?? String(payload.sub);
return {
userId,
appId: entry.appId,
orgId: entry.orgId,
tenantId: undefined,
displayName: meta.full_name ?? meta.name ?? email ?? userId,
};
}
/** Legacy per-app HS256 token (dev IdP / host apps). */
private verifyAppToken(token: string, decodedPayload: jwt.JwtPayload): MessagePrincipal {
const appId = decodedPayload.appId ? String(decodedPayload.appId) : undefined;
if (!appId) throw new UnauthorizedException('missing appId claim');
const secret = this.secrets[appId];
const secret = this.appSecrets[appId];
if (!secret) throw new UnauthorizedException(`unknown app: ${appId}`);
let payload: jwt.JwtPayload;
@@ -43,4 +161,9 @@ export class SessionVerifier {
displayName: payload.name ? String(payload.name) : undefined,
};
}
/** Accept a project URL in any shape (…/rest/v1, …/auth/v1, trailing slash). */
private normalize(url: string): string {
return url.trim().replace(/\/+$/, '').replace(/\/(rest|auth)\/v1$/, '');
}
}
@@ -1,6 +1,10 @@
import { IsOptional, IsString } from 'class-validator';
import { IsInt, IsOptional, IsString } from 'class-validator';
export class SendMessageDto {
@IsString() content!: string;
@IsOptional() @IsString() contentRef?: string;
@IsOptional() @IsString() mimeType?: string;
@IsOptional() @IsInt() sizeBytes?: number;
@IsOptional() @IsString() checksumSha256?: string;
@IsOptional() @IsString() parentInteractionId?: string;
}
@@ -8,6 +8,7 @@ import {
HttpCode,
Param,
Post,
Query,
} from '@nestjs/common';
import { ThreadsService } from './threads.service';
import { MessageService } from '../messaging/message.service';
@@ -22,11 +23,50 @@ export class ThreadsController {
private readonly session: SessionVerifier,
) {}
/** Generic "my threads" — every thread the caller participates in (last message + unread). */
@Get()
async listThreads(@Headers('authorization') auth?: string) {
return this.messages.listThreads(this.principal(auth));
}
/** Generic "my annotated messages" (e.g. ?type=save for a personal bookmarks list). */
@Get('my-annotations')
async myAnnotations(@Headers('authorization') auth?: string, @Query('type') type = 'save') {
return this.messages.listMyAnnotated(this.principal(auth), type);
}
/** Create a thread; `membership`/`creatorRole`/`subject` are opaque, app-supplied attributes the kernel stores but never interprets. */
@Post()
@HttpCode(201)
async createThread(@Body() body: { membership?: string; creatorRole?: string; subject?: string }, @Headers('authorization') auth?: string) {
return this.messages.openThread(null, this.principal(auth), { membership: body?.membership, creatorRole: body?.creatorRole, subject: body?.subject });
}
/** Governed membership: add a user (by userId) to a thread — policy enforces DM cap / roles. */
@Post(':id/participants')
@HttpCode(201)
async addParticipant(@Param('id') id: string, @Body() body: { userId: string; role?: string }, @Headers('authorization') auth?: string) {
return this.messages.addParticipant(id, this.principal(auth), body.userId, body.role);
}
@Get(':id/messages')
async listMessages(@Param('id') id: string) {
return this.threads.getMessages(id);
}
/** Mute / unmute notifications from this thread for the caller. */
@Post(':id/mute')
@HttpCode(200)
async mute(@Param('id') id: string, @Headers('authorization') auth?: string) {
return this.messages.muteThread(id, this.principal(auth), true);
}
@Post(':id/unmute')
@HttpCode(200)
async unmute(@Param('id') id: string, @Headers('authorization') auth?: string) {
return this.messages.muteThread(id, this.principal(auth), false);
}
/** REST/polling fallback for native send (same write as the socket path). */
@Post(':id/messages')
@HttpCode(201)
@@ -36,14 +76,19 @@ export class ThreadsController {
@Headers('authorization') authorization?: string,
@Headers('idempotency-key') idempotencyKey?: string,
) {
const token = (authorization ?? '').replace(/^Bearer\s+/i, '');
if (!token) throw new BadRequestException('Authorization bearer token is required');
const principal = this.session.verify(token);
return this.messages.send(
id,
principal,
{ content: body.content, contentRef: body.contentRef },
this.principal(authorization),
{ content: body.content, contentRef: body.contentRef, mimeType: body.mimeType, sizeBytes: body.sizeBytes, checksumSha256: body.checksumSha256 },
idempotencyKey ?? randomUUID(),
undefined,
body.parentInteractionId,
);
}
private principal(authorization?: string) {
const token = (authorization ?? '').replace(/^Bearer\s+/i, '');
if (!token) throw new BadRequestException('Authorization bearer token is required');
return this.session.verify(token);
}
}
+1897 -35
View File
File diff suppressed because it is too large Load Diff
+20
View File
@@ -0,0 +1,20 @@
import { execSync } from 'node:child_process';
// Tests TRUNCATE between cases, so they MUST NOT touch the dev database. Run them
// against an isolated `iios_test` DB (created + migrated here, once, before the suite).
const TEST_URL = 'postgresql://iios:iios@localhost:5434/iios_test?schema=public';
export default function setup() {
try {
execSync(`docker exec iios-db psql -U iios -d postgres -c "CREATE DATABASE iios_test"`, { stdio: 'pipe' });
} catch (e) {
const msg = String(e.stderr ?? e.stdout ?? e);
if (!/already exists/i.test(msg)) {
console.warn(`[vitest] could not create iios_test (is the iios-db container up?): ${msg.slice(0, 160)}`);
}
}
execSync('pnpm --filter @insignia/iios-service exec prisma migrate deploy', {
stdio: 'inherit',
env: { ...process.env, DATABASE_URL: TEST_URL },
});
}
+7 -3
View File
@@ -15,9 +15,13 @@ export default defineConfig({
},
test: {
include: ['packages/**/src/**/*.{test,spec}.ts', 'test/**/*.{test,spec}.ts'],
// DB-backed specs share one Postgres and TRUNCATE between tests. Run every
// file in a single worker process, sequentially, so there is no cross-file
// race on the shared database.
// DB-backed specs TRUNCATE between tests, so they run against an ISOLATED
// `iios_test` database (created + migrated by the global setup) — never the dev
// DB. This env overrides any DATABASE_URL from the shell.
env: { DATABASE_URL: 'postgresql://iios:iios@localhost:5434/iios_test?schema=public' },
globalSetup: ['./scripts/vitest-global-setup.mjs'],
// Run every file in a single worker process, sequentially, so there is no
// cross-file race on the shared database.
fileParallelism: false,
sequence: { concurrent: false },
pool: 'forks',