diff --git a/packages/iios-service/prisma/migrations/20260703070438_retention_sweep/migration.sql b/packages/iios-service/prisma/migrations/20260703070438_retention_sweep/migration.sql new file mode 100644 index 0000000..9c89020 --- /dev/null +++ b/packages/iios-service/prisma/migrations/20260703070438_retention_sweep/migration.sql @@ -0,0 +1,29 @@ +-- AlterTable +ALTER TABLE "IiosInteraction" ADD COLUMN "dataClass" TEXT NOT NULL DEFAULT 'internal'; + +-- CreateTable +CREATE TABLE "IiosRetentionPolicySnapshot" ( + "id" TEXT NOT NULL, + "policyKey" TEXT NOT NULL, + "scopeSnapshotId" TEXT NOT NULL, + "targetType" TEXT NOT NULL, + "targetId" TEXT NOT NULL, + "dataClass" TEXT NOT NULL, + "retainUntil" TIMESTAMP(3), + "archiveAfter" TIMESTAMP(3) NOT NULL, + "deleteAfter" TIMESTAMP(3) NOT NULL, + "sourceVersion" TEXT NOT NULL DEFAULT 'v1', + "status" TEXT NOT NULL DEFAULT 'ACTIVE', + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + + CONSTRAINT "IiosRetentionPolicySnapshot_pkey" PRIMARY KEY ("id") +); + +-- CreateIndex +CREATE INDEX "IiosRetentionPolicySnapshot_status_deleteAfter_idx" ON "IiosRetentionPolicySnapshot"("status", "deleteAfter"); + +-- CreateIndex +CREATE INDEX "IiosRetentionPolicySnapshot_scopeSnapshotId_status_idx" ON "IiosRetentionPolicySnapshot"("scopeSnapshotId", "status"); + +-- CreateIndex +CREATE UNIQUE INDEX "IiosRetentionPolicySnapshot_targetType_targetId_key" ON "IiosRetentionPolicySnapshot"("targetType", "targetId"); diff --git a/packages/iios-service/prisma/schema.prisma b/packages/iios-service/prisma/schema.prisma index 6cb4413..4d2f480 100644 --- a/packages/iios-service/prisma/schema.prisma +++ b/packages/iios-service/prisma/schema.prisma @@ -435,6 +435,7 @@ model IiosInteraction { occurredAt DateTime @default(now()) receivedAt DateTime @default(now()) status IiosDeliveryState @default(RECEIVED) + dataClass String @default("internal") // envelope data class; drives retention window (P9) policyDecisionRef String? consentReceiptRef String? crreBundleRef String? @@ -493,6 +494,27 @@ model IiosProcessedEvent { @@id([consumerName, eventId]) } +/// Per-resource retention snapshot (P9). Absolute lifecycle timestamps frozen from the +/// resource's age + policy window; read by the retention sweep. delete = redact-in-place. +model IiosRetentionPolicySnapshot { + id String @id @default(cuid()) + policyKey String + scopeSnapshotId String + targetType String // 'interaction' + targetId String + dataClass String + retainUntil DateTime? + archiveAfter DateTime + deleteAfter DateTime + sourceVersion String @default("v1") + status String @default("ACTIVE") // ACTIVE | ARCHIVED | REDACTED + createdAt DateTime @default(now()) + + @@unique([targetType, targetId]) + @@index([status, deleteAfter]) + @@index([scopeSnapshotId, status]) +} + /// Legal/compliance hold (P9 DSR). An ACTIVE, unexpired hold over a target BLOCKS /// erasure/retention deletion until an authorised release. Never auto-released. model IiosComplianceHold { diff --git a/packages/iios-service/src/app.module.ts b/packages/iios-service/src/app.module.ts index 64d9707..9e16610 100644 --- a/packages/iios-service/src/app.module.ts +++ b/packages/iios-service/src/app.module.ts @@ -16,6 +16,7 @@ import { AiModule } from './ai/ai.module'; import { CalendarModule } from './calendar/calendar.module'; import { CapabilityModule } from './capability/capability.module'; import { DsrModule } from './dsr/dsr.module'; +import { RetentionModule } from './retention/retention.module'; import { HealthController } from './health.controller'; import { MetricsController } from './observability/metrics.controller'; import { DevController } from './dev/dev.controller'; @@ -39,6 +40,7 @@ import { DevController } from './dev/dev.controller'; CalendarModule, CapabilityModule, DsrModule, + RetentionModule, ], controllers: [HealthController, MetricsController, DevController], }) diff --git a/packages/iios-service/src/retention/retention.module.ts b/packages/iios-service/src/retention/retention.module.ts new file mode 100644 index 0000000..20b7956 --- /dev/null +++ b/packages/iios-service/src/retention/retention.module.ts @@ -0,0 +1,10 @@ +import { Global, Module } from '@nestjs/common'; +import { RetentionService } from './retention.service'; + +/** Data-retention sweep, global so the dev controller + /metrics can reach it. */ +@Global() +@Module({ + providers: [RetentionService], + exports: [RetentionService], +}) +export class RetentionModule {} diff --git a/packages/iios-service/src/retention/retention.service.spec.ts b/packages/iios-service/src/retention/retention.service.spec.ts new file mode 100644 index 0000000..3c51f75 --- /dev/null +++ b/packages/iios-service/src/retention/retention.service.spec.ts @@ -0,0 +1,119 @@ +import { randomUUID } from 'node:crypto'; +import { describe, it, expect, beforeAll, afterAll, beforeEach } from 'vitest'; +import { PrismaClient } from '@prisma/client'; +import { makeFakePorts } from '@insignia/iios-testkit'; +import { resetDb } from '../test-utils/reset-db'; +import { ActorResolver } from '../identity/actor.resolver'; +import { IngestService } from '../interactions/ingest.service'; +import { RetentionService } from './retention.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 actors = new ActorResolver(asService); +const svc = () => new RetentionService(asService); + +async function ingest(text: string, orgId = 'org_demo'): Promise<{ interactionId: string; scopeId: string }> { + const res = await new IngestService(asService, makeFakePorts(), actors).ingest( + { + scope: { orgId, appId: 'portal-demo' }, + channel: { type: 'WEBHOOK', externalChannelId: 'webhook' }, + source: { handleKind: 'EXTERNAL_VISITOR', externalId: 'ext-1' }, + kind: 'MESSAGE', + thread: { externalThreadId: `wh-${randomUUID().slice(0, 6)}` }, + parts: [{ kind: 'TEXT', bodyText: text }], + occurredAt: new Date().toISOString(), + providerEventId: `pe-${randomUUID()}`, + }, + `ing-${randomUUID()}`, + ); + const it = await prisma.iiosInteraction.findUniqueOrThrow({ where: { id: res.interactionId } }); + return { interactionId: res.interactionId, scopeId: it.scopeId }; +} + +const past = new Date(Date.now() - 86_400_000); +const future = new Date(Date.now() + 86_400_000); +const setSnapshot = (targetId: string, data: { archiveAfter?: Date; deleteAfter?: Date }) => + prisma.iiosRetentionPolicySnapshot.update({ where: { targetType_targetId: { targetType: 'interaction', targetId } }, data }); +const bodyOf = async (interactionId: string) => (await prisma.iiosMessagePart.findFirstOrThrow({ where: { interactionId } })).bodyText; + +beforeAll(async () => { await prisma.$connect(); }); +afterAll(async () => { await prisma.$disconnect(); }); +beforeEach(async () => { await resetDb(prisma); }); + +describe('RetentionService (P9 — retention sweep)', () => { + it('ensureSnapshots captures one ACTIVE snapshot per interaction with archive < delete', async () => { + const { interactionId } = await ingest('hello'); + await svc().ensureSnapshots(); + const snap = await prisma.iiosRetentionPolicySnapshot.findUniqueOrThrow({ where: { targetType_targetId: { targetType: 'interaction', targetId: interactionId } } }); + expect(snap.status).toBe('ACTIVE'); + expect(snap.dataClass).toBe('internal'); + expect(snap.archiveAfter.getTime()).toBeLessThan(snap.deleteAfter.getTime()); + }); + + it('redacts an interaction whose delete window has passed', async () => { + const { interactionId } = await ingest('my secret'); + await svc().ensureSnapshots(); + await setSnapshot(interactionId, { archiveAfter: past, deleteAfter: past }); + + const res = await svc().applySweep(); + expect(res.redacted).toBe(1); + expect(await bodyOf(interactionId)).toBe('[redacted]'); + expect((await prisma.iiosInteraction.findUniqueOrThrow({ where: { id: interactionId } })).status).toBe('REDACTED'); + expect((await prisma.iiosRetentionPolicySnapshot.findUniqueOrThrow({ where: { targetType_targetId: { targetType: 'interaction', targetId: interactionId } } })).status).toBe('REDACTED'); + expect(await prisma.iiosAuditLink.count({ where: { action: 'retention.redacted' } })).toBe(1); + }); + + it('archives (not deletes) when only the archive window has passed, then deletes later', async () => { + const { interactionId } = await ingest('still needed'); + await svc().ensureSnapshots(); + await setSnapshot(interactionId, { archiveAfter: past, deleteAfter: future }); + + let res = await svc().applySweep(); + expect(res).toMatchObject({ archived: 1, redacted: 0 }); + expect(await bodyOf(interactionId)).toBe('still needed'); // content untouched by archive + + await setSnapshot(interactionId, { deleteAfter: past }); + res = await svc().applySweep(); + expect(res.redacted).toBe(1); + expect(await bodyOf(interactionId)).toBe('[redacted]'); + }); + + it('an active compliance hold blocks the sweep until released', async () => { + const { interactionId, scopeId } = await ingest('held content'); + await svc().ensureSnapshots(); + await setSnapshot(interactionId, { archiveAfter: past, deleteAfter: past }); + const hold = await prisma.iiosComplianceHold.create({ data: { scopeId, targetType: 'message', targetId: interactionId, holdReason: 'legal' } }); + + let res = await svc().applySweep(); + expect(res).toMatchObject({ redacted: 0, skippedHeld: 1 }); + expect(await bodyOf(interactionId)).toBe('held content'); + + await prisma.iiosComplianceHold.update({ where: { id: hold.id }, data: { status: 'RELEASED' } }); + res = await svc().applySweep(); + expect(res.redacted).toBe(1); + expect(await bodyOf(interactionId)).toBe('[redacted]'); + }); + + it('is idempotent — a REDACTED snapshot is terminal', async () => { + const { interactionId } = await ingest('x'); + await svc().ensureSnapshots(); + await setSnapshot(interactionId, { archiveAfter: past, deleteAfter: past }); + await svc().applySweep(); + expect(await svc().applySweep()).toMatchObject({ archived: 0, redacted: 0, skippedHeld: 0 }); + }); + + it('sweep is tenant-scoped', async () => { + const a = await ingest('tenant a', 'org_A'); + const b = await ingest('tenant b', 'org_B'); + await svc().ensureSnapshots(); + await setSnapshot(a.interactionId, { archiveAfter: past, deleteAfter: past }); + await setSnapshot(b.interactionId, { archiveAfter: past, deleteAfter: past }); + + const res = await svc().applySweep(a.scopeId); + expect(res.redacted).toBe(1); + expect(await bodyOf(a.interactionId)).toBe('[redacted]'); + expect(await bodyOf(b.interactionId)).toBe('tenant b'); // scope B untouched + }); +}); diff --git a/packages/iios-service/src/retention/retention.service.ts b/packages/iios-service/src/retention/retention.service.ts new file mode 100644 index 0000000..b5831a7 --- /dev/null +++ b/packages/iios-service/src/retention/retention.service.ts @@ -0,0 +1,127 @@ +import { Injectable, Logger, OnModuleDestroy, OnModuleInit } from '@nestjs/common'; +import { PrismaService } from '../prisma/prisma.service'; +import { recordAudit } from '../observability/audit'; + +const DAY_MS = 86_400_000; + +export interface SweepResult { + archived: number; + redacted: number; + skippedHeld: number; +} + +/** + * Data retention (P9). Each interaction gets a per-resource retention snapshot with + * frozen archive/delete timestamps; a scheduled sweep archives (status change) then + * deletes = redacts-in-place (reusing the DSR tombstone semantics) once a resource ages + * past its window. An active compliance hold blocks both. Never a hard delete. + */ +@Injectable() +export class RetentionService implements OnModuleInit, OnModuleDestroy { + private readonly logger = new Logger(RetentionService.name); + private timer?: ReturnType; + + constructor(private readonly prisma: PrismaService) {} + + onModuleInit(): void { + const ms = Number(process.env.IIOS_RETENTION_SWEEP_INTERVAL_MS ?? 0); + if (ms > 0) { + this.timer = setInterval(() => { + void this.sweep().catch((err) => this.logger.warn(`retention sweep failed: ${(err as Error).message}`)); + }, ms); + } + } + + onModuleDestroy(): void { + if (this.timer) clearInterval(this.timer); + } + + private windowDays(dataClass: string, kind: 'ARCHIVE' | 'DELETE'): number { + const cls = process.env[`IIOS_RETENTION_${kind}_DAYS_${dataClass.toUpperCase()}`]; + const glob = process.env[`IIOS_RETENTION_${kind}_DAYS`]; + return Number(cls ?? glob ?? (kind === 'ARCHIVE' ? 90 : 365)); + } + + /** Lazily capture a per-resource snapshot for any interaction that lacks one. */ + async ensureSnapshots(scopeId?: string): Promise { + const interactions = await this.prisma.iiosInteraction.findMany({ + where: scopeId ? { scopeId } : {}, + select: { id: true, scopeId: true, dataClass: true, receivedAt: true }, + }); + let created = 0; + for (const i of interactions) { + const exists = await this.prisma.iiosRetentionPolicySnapshot.findUnique({ + where: { targetType_targetId: { targetType: 'interaction', targetId: i.id } }, + }); + if (exists) continue; + const base = i.receivedAt.getTime(); + await this.prisma.iiosRetentionPolicySnapshot.create({ + data: { + policyKey: `default:${i.dataClass}:v1`, + scopeSnapshotId: i.scopeId, + targetType: 'interaction', + targetId: i.id, + dataClass: i.dataClass, + archiveAfter: new Date(base + this.windowDays(i.dataClass, 'ARCHIVE') * DAY_MS), + deleteAfter: new Date(base + this.windowDays(i.dataClass, 'DELETE') * DAY_MS), + sourceVersion: process.env.IIOS_RETENTION_POLICY_VERSION ?? 'v1', + }, + }); + created++; + } + return created; + } + + /** Act on due snapshots: archive, or delete (redact-in-place), honoring compliance holds. */ + async applySweep(scopeId?: string): Promise { + const now = new Date(); + const due = await this.prisma.iiosRetentionPolicySnapshot.findMany({ + where: { status: { in: ['ACTIVE', 'ARCHIVED'] }, archiveAfter: { lte: now }, ...(scopeId ? { scopeSnapshotId: scopeId } : {}) }, + }); + const result: SweepResult = { archived: 0, redacted: 0, skippedHeld: 0 }; + + for (const s of due) { + const held = await this.prisma.iiosComplianceHold.findFirst({ + where: { targetId: s.targetId, status: 'ACTIVE', OR: [{ expiresAt: null }, { expiresAt: { gt: now } }] }, + }); + if (held) { + result.skippedHeld++; + continue; + } + + if (s.deleteAfter <= now) { + await this.prisma.$transaction([ + this.prisma.iiosMessagePart.updateMany({ where: { interactionId: s.targetId }, data: { bodyText: '[redacted]', contentRef: null } }), + this.prisma.iiosInboundRawEvent.updateMany({ where: { interactionId: s.targetId }, data: { payload: { redacted: true } } }), + this.prisma.iiosInteraction.update({ where: { id: s.targetId }, data: { status: 'REDACTED' } }), + this.prisma.iiosRetentionPolicySnapshot.update({ where: { id: s.id }, data: { status: 'REDACTED' } }), + ]); + await recordAudit(this.prisma, { action: 'retention.redacted', resourceType: 'interaction', resourceId: s.targetId, scopeId: s.scopeSnapshotId }); + result.redacted++; + } else if (s.status === 'ACTIVE') { + await this.prisma.iiosRetentionPolicySnapshot.update({ where: { id: s.id }, data: { status: 'ARCHIVED' } }); + await recordAudit(this.prisma, { action: 'retention.archived', resourceType: 'interaction', resourceId: s.targetId, scopeId: s.scopeSnapshotId }); + result.archived++; + } + } + return result; + } + + /** Full sweep: capture missing snapshots, then act on due ones. */ + async sweep(scopeId?: string): Promise { + await this.ensureSnapshots(scopeId); + return this.applySweep(scopeId); + } + + /** Snapshot lifecycle counts for /metrics (global ops). */ + async summary(): Promise<{ total: number; byStatus: Record }> { + const groups = await this.prisma.iiosRetentionPolicySnapshot.groupBy({ by: ['status'], _count: true }); + const byStatus: Record = {}; + let total = 0; + for (const g of groups) { + byStatus[g.status] = g._count; + total += g._count; + } + return { total, byStatus }; + } +} diff --git a/packages/iios-service/src/test-utils/reset-db.ts b/packages/iios-service/src/test-utils/reset-db.ts index 6b2f86c..7a01cc0 100644 --- a/packages/iios-service/src/test-utils/reset-db.ts +++ b/packages/iios-service/src/test-utils/reset-db.ts @@ -8,7 +8,7 @@ import type { PrismaClient } from '@prisma/client'; export async function resetDb(prisma: PrismaClient): Promise { await prisma.$executeRawUnsafe( `TRUNCATE TABLE - "IiosComplianceHold","IiosProjectionCursor","IiosIdempotencyCommand","IiosDlqItem","IiosAuditLink", + "IiosRetentionPolicySnapshot","IiosComplianceHold","IiosProjectionCursor","IiosIdempotencyCommand","IiosDlqItem","IiosAuditLink", "IiosMeetingActionItem","IiosActionItem","IiosTranscriptSegment","IiosMeetingTranscript","IiosMeetingSummary","IiosMeetingParticipant","IiosMeeting","IiosMeetingRequest", "IiosCalendarSyncCursor","IiosCalendarEvent","IiosCalendarProviderAccount","IiosAvailabilityWindow", "IiosAiEvidenceLink","IiosAiClaim","IiosAiToolCall","IiosAiArtifact","IiosAiModelRun","IiosAiJob","IiosEmbeddingRef",