Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions apps/control-plane/src/app.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -103,6 +104,7 @@ export async function buildControlPlane(config: ControlPlaneConfig): Promise<Fas
const requireSession = accessGuard(database);

app.decorateRequest("user", null);
registerReviews(app, database, requireSession);
registerTeam(app, database, requireSession, config.oidc?.issuer);
registerOidc(app, database, config);
registerIntegrations(app, database, config.controlPlaneSecret, requireSession);
Expand Down
77 changes: 77 additions & 0 deletions apps/control-plane/src/reviews.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
import { randomUUID } from "node:crypto";
import { z } from "zod";
import type { FastifyInstance, FastifyReply, FastifyRequest } from "fastify";
import type { ClassificationEvent, Delivery, StoredIntegration, UserRecord } from "@pyro/contracts";
import type { Database } from "@pyro/storage";
import { decisionDeliveries, deliveryFor } from "@pyro/integrations";
import { appScope, canAccessApp, visibleEvent } from "./access.js";
export interface ReviewRecord {
id: string; appId: string; revision: number; status: "open" | "resolved";
severity: "low" | "medium" | "high"; assignedTo?: string;
disposition?: "true_positive" | "false_positive" | "uncertain";
updatedAt: string; resolvedAt?: string; resolvedBy?: string;
comments: Array<{ id: string; actorId: string; at: string; text: string }>;
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<unknown>) {
const store = database.document<ReviewRecord[]>("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<UserRecord[]>("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<UserRecord[]>("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<StoredIntegration[]>("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<void> | 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; });
}
21 changes: 21 additions & 0 deletions apps/control-plane/test/reviews.test.ts
Original file line number Diff line number Diff line change
@@ -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);
});
5 changes: 4 additions & 1 deletion apps/dashboard/src/App.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -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: <Activity className="size-3.5" /> },
{ value: "usage", label: "Usage", icon: <ChartColumn className="size-3.5" /> },
{ value: "reviews", label: "Review inbox", icon: <BookOpenCheck className="size-3.5" /> },
{ value: "activity", label: "Activity", icon: <BookOpenCheck className="size-3.5" /> },
{ value: "playground", label: "Playground", icon: <TerminalSquare className="size-3.5" /> },
] },
Expand Down Expand Up @@ -107,6 +109,7 @@ export default function App() {
playground: <PlaygroundPage onDecision={() => setRefreshKey((value) => value + 1)} />,
profiles: <ProfilesPage />,
history: <PolicyHistoryPage />,
reviews: <ReviewsPage canReview={user.role !== "viewer"} />,
activity: <ActivityPage refreshKey={refreshKey} />,
keys: <ApiKeysPage />,
settings: <SettingsPage />,
Expand Down
2 changes: 1 addition & 1 deletion apps/dashboard/src/pages/IntegrationsPage.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ export function IntegrationsPage() {
return <>
<PageHeader title="Webhooks" description="Send filtered decisions to your HTTP receiver with signed JSON payloads. Delivery runs in the background with retries and a persistent history." actions={<Button onClick={() => { setError(undefined); setEditing(createWebhookDraft()); }}>Add webhook</Button>} />
{error && <p role="alert" className="mb-4 text-sm text-danger">{error}</p>}{notice && <p role="status" className="mb-4 text-sm">{notice}</p>}
{integrations.length > 0 && <div className="mb-6 grid gap-4 lg:grid-cols-2">{integrations.map((item) => <Card key={item.id}><CardHeader><div className="flex items-center justify-between gap-3"><CardTitle>{item.name}</CardTitle><Switch disabled={busy} aria-label={`Enable ${item.name}`} checked={item.enabled} onCheckedChange={(enabled) => void act(async () => { await api.put(`/api/integrations/${item.id}`, { enabled }); })} /></div><p className="text-sm text-muted">Signed webhook · {item.destinationHost}</p></CardHeader><CardContent><p className="text-sm">{item.actions.join(" / ")} · minimum risk {item.minimumRisk} · {item.appIds.length ? `${item.appIds.length} selected ${item.appIds.length === 1 ? "application" : "applications"}` : "all applications"}</p><div className="mt-4 flex flex-wrap gap-2"><Button size="sm" variant="outline" onClick={() => { setError(undefined); setEditing(createWebhookDraft(item)); }}>Edit</Button><Button size="sm" variant="outline" disabled={busy || !item.enabled} onClick={() => void act(async () => { await api.post(`/api/integrations/${item.id}/test`); setNotice("Test queued. Check delivery history below."); })}>Send test</Button>{item.type === "webhook" && <Button size="sm" variant="outline" disabled={busy} onClick={() => { if (window.confirm("Rotate the signing secret? Update your receiver to use the new secret.")) void act(async () => { const result = await api.post<{ signingSecret: string }>(`/api/integrations/${item.id}/rotate-secret`); setSecret(result.signingSecret); }); }}>Rotate secret</Button>}<Button size="sm" variant="ghost" disabled={busy} onClick={() => { if (window.confirm(`Delete ${item.name}? Pending deliveries will stop.`)) void act(async () => { await api.delete(`/api/integrations/${item.id}`); }); }}>Delete</Button></div></CardContent></Card>)}</div>}
{integrations.length > 0 && <div className="mb-6 grid gap-4 lg:grid-cols-2">{integrations.map((item) => <Card key={item.id}><CardHeader><div className="flex items-center justify-between gap-3"><CardTitle>{item.name}</CardTitle><Switch disabled={busy} aria-label={`Enable ${item.name}`} checked={item.enabled} onCheckedChange={(enabled) => void act(async () => { await api.put(`/api/integrations/${item.id}`, { enabled }); })} /></div><p className="text-sm text-muted">Signed webhook · {item.destinationHost}</p></CardHeader><CardContent><p className="text-sm">{item.actions.join(" / ")} · minimum risk {item.minimumRisk} · {item.appIds.length ? `${item.appIds.length} selected ${item.appIds.length === 1 ? "application" : "applications"}` : "all applications"}</p><label className="mt-3 flex items-center gap-3 text-sm"><Switch disabled={busy} checked={item.reviewResolutions ?? false} onCheckedChange={(reviewResolutions) => void act(async () => { await api.put(`/api/integrations/${item.id}`, { reviewResolutions }); })} />Signed review resolution callbacks</label><div className="mt-4 flex flex-wrap gap-2"><Button size="sm" variant="outline" onClick={() => { setError(undefined); setEditing(createWebhookDraft(item)); }}>Edit</Button><Button size="sm" variant="outline" disabled={busy || !item.enabled} onClick={() => void act(async () => { await api.post(`/api/integrations/${item.id}/test`); setNotice("Test queued. Check delivery history below."); })}>Send test</Button>{item.type === "webhook" && <Button size="sm" variant="outline" disabled={busy} onClick={() => { if (window.confirm("Rotate the signing secret? Update your receiver to use the new secret.")) void act(async () => { const result = await api.post<{ signingSecret: string }>(`/api/integrations/${item.id}/rotate-secret`); setSecret(result.signingSecret); }); }}>Rotate secret</Button>}<Button size="sm" variant="ghost" disabled={busy} onClick={() => { if (window.confirm(`Delete ${item.name}? Pending deliveries will stop.`)) void act(async () => { await api.delete(`/api/integrations/${item.id}`); }); }}>Delete</Button></div></CardContent></Card>)}</div>}
{!integrations.length && <p className="mb-6 border border-dashed p-6 text-sm text-muted">No webhooks yet. Add your receiver, then send a test.</p>}
<WebhookDeliveries deliveries={deliveries} webhooks={integrations} busy={busy} onRefresh={() => void act(async () => {})} onRetry={(id) => void act(async () => { await api.post(`/api/integration-deliveries/${id}/retry`); })} />
<Dialog open={Boolean(editing)} onOpenChange={(open) => { if (!open && !busy) setEditing(undefined); }}>
Expand Down
Loading
Loading