From d628645e47bbd1227ebfa3cc42ad11f9c32c5866 Mon Sep 17 00:00:00 2001 From: Hardik Bhatia Date: Sat, 26 Sep 2026 21:09:48 +0530 Subject: [PATCH] feat(reviews): add scoped review inbox and signed resolution callbacks --- apps/control-plane/src/app.ts | 2 + apps/control-plane/src/reviews.ts | 77 ++++++++++++++++ apps/control-plane/test/reviews.test.ts | 21 +++++ apps/dashboard/src/App.tsx | 5 +- apps/dashboard/src/pages/IntegrationsPage.tsx | 2 +- apps/dashboard/src/pages/ReviewsPage.tsx | 31 +++++++ docs/control-plane.openapi.yaml | 90 +++++++++++++++++++ docs/deployment.md | 16 ++++ packages/cli/test/contract.test.ts | 2 +- packages/contracts/src/index.ts | 4 +- 10 files changed, 246 insertions(+), 4 deletions(-) create mode 100644 apps/control-plane/src/reviews.ts create mode 100644 apps/control-plane/test/reviews.test.ts create mode 100644 apps/dashboard/src/pages/ReviewsPage.tsx diff --git a/apps/control-plane/src/app.ts b/apps/control-plane/src/app.ts index 92fc5e2..28f1c1b 100644 --- a/apps/control-plane/src/app.ts +++ b/apps/control-plane/src/app.ts @@ -23,6 +23,7 @@ import { createSession, ensureAdmin, sessionUserId, sha256, verifyAdminPassword import type { ControlPlaneConfig } from "./config.js"; import { accessGuard, appScope, canAccessApp, visibleEvent, visibleUser, allowedProfiles } from "./access.js"; +import { registerReviews } from "./reviews.js"; import { registerTeam, verifyPassword } from "./team.js"; import { registerOidc } from "./oidc.js"; import { registerPolicyHistory } from "./policies.js"; @@ -103,6 +104,7 @@ export async function buildControlPlane(config: ControlPlaneConfig): Promise; + pendingCallbacks?: Delivery[]; +} +const update = z.object({ expectedRevision: z.number().int().nonnegative(), assignedTo: z.string().max(100).nullable().optional(), severity: z.enum(["low", "medium", "high"]).optional(), disposition: z.enum(["true_positive", "false_positive", "uncertain"]).optional(), comment: z.string().trim().min(1).max(2000).optional(), reopen: z.boolean().optional() }); +const initial = (event: ClassificationEvent): ReviewRecord => ({ id: event.id, appId: event.appId ?? "default", revision: 0, status: "open", severity: event.risk >= .8 ? "high" : "medium", updatedAt: event.createdAt, comments: [] }); +const visible = ({ pendingCallbacks, ...record }: ReviewRecord) => record; +export function registerReviews(app: FastifyInstance, database: Database, guard: (r: FastifyRequest, p: FastifyReply) => Promise) { + const store = database.document("reviews", () => []); + app.get("/api/reviews", { preHandler: guard }, async (request) => { + const query = request.query as { appId?: string; status?: string; assignedTo?: string; offset?: string }; + const states = await store.read(); + const events = await database.events.query({ action: "review", appId: query.appId, appIds: appScope(request.user!), limit: 500, offset: Math.max(0, Number(query.offset) || 0) }); + const reviews = events.events.map((event) => ({ ...visible(states.find((r) => r.id === event.id) ?? initial(event)), event: visibleEvent(request.user!, event), ageMs: Date.now() - Date.parse(event.createdAt) })).filter((r) => (!query.status || r.status === query.status) && (!query.assignedTo || r.assignedTo === query.assignedTo)); + const users = (await database.document("users", () => []).read()).filter((u) => !u.disabled && ["admin", "reviewer", "operator"].includes(u.role ?? "") && (request.user!.role === "admin" || u.id === request.user!.id)); + return { reviews, totalEvents: events.total, nextOffset: events.events.length === 500 ? (Number(query.offset) || 0) + 500 : null, assignees: users.map((u) => ({ id: u.id, username: u.username })) }; + }); + app.get<{ Params: { id: string } }>("/api/reviews/:id", { preHandler: guard }, async (request, reply) => { + const event = await database.events.findById(request.params.id); + if (!event || event.action !== "review" || !canAccessApp(request.user!, event.appId)) return reply.code(404).send({ error: "Review not found." }); + return { review: visible((await store.read()).find((r) => r.id === event.id) ?? initial(event)), event: visibleEvent(request.user!, event) }; + }); + app.put<{ Params: { id: string } }>("/api/reviews/:id", { preHandler: guard }, async (request, reply) => { + const parsed = update.safeParse(request.body); + if (!parsed.success) return reply.code(400).send({ error: parsed.error.issues[0]?.message }); + const event = await database.events.findById(request.params.id); + if (!event || event.action !== "review" || !canAccessApp(request.user!, event.appId)) return reply.code(404).send({ error: "Review not found." }); + const body = parsed.data; + if (body.assignedTo) { + const user = (await database.document("users", () => []).read()).find((u) => u.id === body.assignedTo && !u.disabled); + if (!user || !canAccessApp(user, event.appId) || !["admin", "operator", "reviewer"].includes(user.role ?? "")) return reply.code(400).send({ error: "Assignee must be an authorized reviewer for this application." }); + } + if (body.reopen && body.disposition) return reply.code(400).send({ error: "Choose reopen or a resolution." }); + const integrations = (await database.document("integrations", () => []).read()).filter((i) => i.reviewResolutions && i.signingSecret); + let saved!: ReviewRecord; let conflict = false; + await store.update((rows) => { + const current = rows.find((r) => r.id === event.id) ?? initial(event); + if (current.revision !== body.expectedRevision) { conflict = true; return rows; } + const now = new Date().toISOString(); + saved = { ...current, revision: current.revision + 1, updatedAt: now, severity: body.severity ?? current.severity, assignedTo: body.assignedTo === null ? undefined : body.assignedTo ?? current.assignedTo, comments: body.comment ? [...current.comments, { id: randomUUID(), actorId: request.user!.id, at: now, text: body.comment }] : current.comments }; + if (body.reopen) { saved.status = "open"; delete saved.disposition; delete saved.resolvedAt; delete saved.resolvedBy; } + if (body.disposition) { + saved.status = "resolved"; saved.disposition = body.disposition; saved.resolvedBy = request.user!.id; saved.resolvedAt = now; + const callbacks = decisionDeliveries(event, integrations).map((d) => deliveryFor(d.integrationId, { ...d.payload, id: `${event.id}:review:${saved.revision}`, type: "review.resolved", createdAt: now, data: { ...d.payload.data, review: { revision: saved.revision, disposition: body.disposition!, actorId: request.user!.id } } })); + saved.pendingCallbacks = [...current.pendingCallbacks ?? [], ...callbacks]; + } + return [...rows.filter((r) => r.id !== event.id), saved]; + }); + if (conflict) return reply.code(409).send({ error: "This review changed. Reload before saving." }); + return { review: visible(saved), warning: "Feedback does not change the original decision or execute a tool." }; + }); + // Persist callback intent with the review, then hand off to the existing signed + // delivery outbox. Stable IDs make retries safe after an interrupted handoff. + let active: Promise | undefined; + const flush = async () => { + const records = await store.read(); + for (const record of records) for (const delivery of record.pendingCallbacks ?? []) { + await database.deliveries.enqueue(delivery); + await store.update((rows) => rows.map((r) => r.id === record.id ? { ...r, pendingCallbacks: r.pendingCallbacks?.filter((d) => d.id !== delivery.id) } : r)); + } + const cutoff = new Date(Date.now() - (Math.max(1, Number(process.env.EVENT_RETENTION_DAYS) || 30)) * 86400_000).toISOString(); + await store.update((rows) => rows.filter((r) => r.updatedAt >= cutoff || r.pendingCallbacks?.length)); + }; + const timer = setInterval(() => { if (!active) { active = flush().catch((e) => app.log.error(e, "Review callback handoff failed")).finally(() => { active = undefined; }); } }, 1000); timer.unref(); + app.addHook("onClose", async () => { clearInterval(timer); await active; }); +} diff --git a/apps/control-plane/test/reviews.test.ts b/apps/control-plane/test/reviews.test.ts new file mode 100644 index 0000000..a0f4288 --- /dev/null +++ b/apps/control-plane/test/reviews.test.ts @@ -0,0 +1,21 @@ +import assert from "node:assert/strict"; +import test from "node:test"; +import { randomUUID } from "node:crypto"; +import { openDatabase } from "@pyro/storage"; +import { buildControlPlane } from "../src/app.js"; +test("review resolution is attributed, immutable, concurrency-safe and does not change a decision", async (t) => { + const config = { host: "127.0.0.1", port: 0, databaseUrl: `memory://reviews-${randomUUID()}`, adminPassword: "correct-horse-battery-staple", controlPlaneSecret: "control-plane-test-secret", gatewayInternalUrl: "http://127.0.0.1:1", gatewayApiKey: "test-key", typesafeEndpoint: "https://api.typesafe.ai/v1/systemone", typesafeModel: "jev-latest" }; + const app = await buildControlPlane(config); t.after(() => app.close()); const db = await openDatabase(config.databaseUrl); + const login = await app.inject({ method: "POST", url: "/api/auth/login", payload: { password: config.adminPassword } }); const cookie = login.headers["set-cookie"]!.split(";")[0]!; + const event = { id: "review-one", appId: "default", createdAt: new Date().toISOString(), profileId: "default", action: "review" as const, verdict: "suspicious" as const, risk: .7, confidence: .7, reason: "review", detectors: [], model: "local", provider: "local-rules", latencyMs: 0, queueMs: 0, inputHash: "hash" }; + await db.events.append(event); + const edit = (disposition: string) => app.inject({ method: "PUT", url: "/api/reviews/review-one", headers: { cookie }, payload: { expectedRevision: 0, disposition, comment: "Redacted feedback" } }); + const responses = await Promise.all([edit("true_positive"), edit("false_positive")]); + assert.deepEqual(responses.map((r) => r.statusCode).sort(), [200, 409]); + assert.deepEqual(await db.events.findById(event.id), event); + const review = (await app.inject({ method: "GET", url: "/api/reviews/review-one", headers: { cookie } })).json().review; + assert.equal(review.revision, 1); assert.equal(review.status, "resolved"); assert.equal(review.resolvedBy, login.json().user.id); assert.equal(review.comments.length, 1); + const viewer = (await app.inject({ method: "POST", url: "/api/team", headers: { cookie }, payload: { username: "viewer", role: "viewer", appIds: ["default"] } })).json(); + const session = await app.inject({ method: "POST", url: "/api/auth/login", payload: { username: "viewer", password: viewer.password } }); + assert.equal((await app.inject({ method: "PUT", url: "/api/reviews/review-one", headers: { cookie: session.headers["set-cookie"]!.split(";")[0]! }, payload: { expectedRevision: 1, disposition: "uncertain" } })).statusCode, 403); +}); diff --git a/apps/dashboard/src/App.tsx b/apps/dashboard/src/App.tsx index 574c82b..dd7aa34 100644 --- a/apps/dashboard/src/App.tsx +++ b/apps/dashboard/src/App.tsx @@ -16,17 +16,19 @@ import { PlaygroundPage } from "@/pages/PlaygroundPage"; import { PolicyHistoryPage } from "@/pages/PolicyHistoryPage"; import { ProfilesPage } from "@/pages/ProfilesPage"; import { IntegrationsPage } from "@/pages/IntegrationsPage"; +import { ReviewsPage } from "@/pages/ReviewsPage"; import { TeamPage } from "@/pages/TeamPage"; import { SettingsPage } from "@/pages/SettingsPage"; import { UsagePage } from "@/pages/UsagePage"; -type Page = "team" | "history" | "overview" | "apps" | "usage" | "playground" | "profiles" | "activity" | "keys" | "settings" | "integrations"; +type Page = "reviews" | "team" | "history" | "overview" | "apps" | "usage" | "playground" | "profiles" | "activity" | "keys" | "settings" | "integrations"; type User = UserRecord; const NAV: BranchedMenuItem[] = [ { label: "Observe", children: [ { value: "overview", label: "Overview", icon: }, { value: "usage", label: "Usage", icon: }, + { value: "reviews", label: "Review inbox", icon: }, { value: "activity", label: "Activity", icon: }, { value: "playground", label: "Playground", icon: }, ] }, @@ -107,6 +109,7 @@ export default function App() { playground: setRefreshKey((value) => value + 1)} />, profiles: , history: , + reviews: , activity: , keys: , settings: , diff --git a/apps/dashboard/src/pages/IntegrationsPage.tsx b/apps/dashboard/src/pages/IntegrationsPage.tsx index 66dd68b..f681c5e 100644 --- a/apps/dashboard/src/pages/IntegrationsPage.tsx +++ b/apps/dashboard/src/pages/IntegrationsPage.tsx @@ -37,7 +37,7 @@ export function IntegrationsPage() { return <> { setError(undefined); setEditing(createWebhookDraft()); }}>Add webhook} /> {error &&

{error}

}{notice &&

{notice}

} - {integrations.length > 0 &&
{integrations.map((item) =>
{item.name} void act(async () => { await api.put(`/api/integrations/${item.id}`, { enabled }); })} />

Signed webhook · {item.destinationHost}

{item.actions.join(" / ")} · minimum risk {item.minimumRisk} · {item.appIds.length ? `${item.appIds.length} selected ${item.appIds.length === 1 ? "application" : "applications"}` : "all applications"}

{item.type === "webhook" && }
)}
} + {integrations.length > 0 &&
{integrations.map((item) =>
{item.name} void act(async () => { await api.put(`/api/integrations/${item.id}`, { enabled }); })} />

Signed webhook · {item.destinationHost}

{item.actions.join(" / ")} · minimum risk {item.minimumRisk} · {item.appIds.length ? `${item.appIds.length} selected ${item.appIds.length === 1 ? "application" : "applications"}` : "all applications"}

{item.type === "webhook" && }
)}
} {!integrations.length &&

No webhooks yet. Add your receiver, then send a test.

} void act(async () => {})} onRetry={(id) => void act(async () => { await api.post(`/api/integration-deliveries/${id}/retry`); })} /> { if (!open && !busy) setEditing(undefined); }}> diff --git a/apps/dashboard/src/pages/ReviewsPage.tsx b/apps/dashboard/src/pages/ReviewsPage.tsx new file mode 100644 index 0000000..dca76a6 --- /dev/null +++ b/apps/dashboard/src/pages/ReviewsPage.tsx @@ -0,0 +1,31 @@ +import { useEffect, useState } from "react"; +import type { ClassificationEvent } from "@pyro/contracts"; +import { api } from "@/lib/api"; +import { PageHeader } from "@/components/shared"; +import { Card, CardContent } from "@/components/ui/card"; +import { Button } from "@/components/ui/button"; +import { Textarea } from "@/components/ui/textarea"; +interface Review { id: string; revision: number; status: string; appId: string; assignedTo?: string; severity: string; disposition?: string; ageMs: number; comments: Array<{ id: string; actorId: string; at: string; text: string }>; event: ClassificationEvent } +export function ReviewsPage({ canReview }: { canReview: boolean }) { + const [reviews, setReviews] = useState([]); const [assignees, setAssignees] = useState>([]); + const [selected, setSelected] = useState(""); const [status, setStatus] = useState(() => localStorage.getItem("pyro-review-filter") || "open"); + const [comment, setComment] = useState(""); const [message, setMessage] = useState(""); const [busy, setBusy] = useState(false); + const [sample, setSample] = useState(""); const [expected, setExpected] = useState("allow"); + const [offset, setOffset] = useState(0); const [nextOffset, setNextOffset] = useState(null); + const load = async () => { const result = await api.get<{ reviews: Review[]; assignees: typeof assignees; nextOffset: number | null }>(`/api/reviews?status=${status}&offset=${offset}`); setReviews(result.reviews); setAssignees(result.assignees); setNextOffset(result.nextOffset); }; + useEffect(() => { void load().catch((e) => setMessage(e.message)); localStorage.setItem("pyro-review-filter", status); }, [status, offset]); + const current = reviews.find((r) => r.id === selected); + const save = async (changes: Record) => { if (!current) return; setBusy(true); setMessage(""); try { await api.put(`/api/reviews/${current.id}`, { expectedRevision: current.revision, ...changes }); setComment(""); await load(); setMessage("Review saved. The original decision is unchanged."); } catch (e) { setMessage(e instanceof Error ? e.message : "Save failed."); } finally { setBusy(false); } }; + const exportSample = () => { const blob = new Blob([JSON.stringify({ id: current!.id, input: sample, expected, category: "review-feedback" }) + "\n"], { type: "application/x-ndjson" }); const url = URL.createObjectURL(blob); const link = document.createElement("a"); link.href = url; link.download = "review-sample.jsonl"; link.click(); URL.revokeObjectURL(url); }; + return
+
+ {message &&

{message}

} + {reviews.map((r) => )}
DecisionApplication / traceSeverityAgeStatus
{r.appId}
{r.event.traceId?.slice(0, 12) ?? r.event.requestId?.slice(0, 12)}
{r.severity}{Math.floor(r.ageMs / 60000)} min{r.disposition ?? r.status}
{!reviews.length &&

No matching reviews on this page.

}
+ {current &&

Review {current.id}

Policy {current.event.profileId} · revision {current.event.policyRevision ?? "legacy"} · original action {current.event.action}

+ {canReview &&
} +
{current.comments.map((c) =>

{c.actorId} · {new Date(c.at).toLocaleString()}

{c.text}

)}
+ {canReview && <>