From 12f6a65b804c9a908adeb448c6f452d31fbd81cf Mon Sep 17 00:00:00 2001 From: Simon Date: Wed, 16 Sep 2026 22:53:55 +0200 Subject: [PATCH] feat(api): deterministic scan worker + AUD0 wire + tenant binding MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - scan worker: atomic claim (FOR UPDATE SKIP LOCKED), Playwright-core + system chromium capture → typed EvidenceRecords (runtime-html/dom/computed-style/ stylesheet/asset/network-request/screenshot), confidence=measured - AUD0 emitter (fire-and-forget, advisory): c0py.registry_created/.scan_created/ .scan_completed/.scan_failed, tenant_id = Zitadel resourceowner (org) - tenant binding: resourceowner claim from introspection (deny-by-default 403), migration 002 tenant_id on registries+scans, all queries tenant+owner scoped - evidence gaps stay explicit (screenshot miss → no screenshot record) - 30/30 tests, canonical validator OK --- Dockerfile | 3 +- apps/api/migrations/002_tenant.sql | 8 + apps/api/package.json | 1 + apps/api/src/audit/aud0.ts | 62 ++++++++ apps/api/src/auth/introspection.ts | 17 ++- apps/api/src/capture/capture.ts | 209 ++++++++++++++++++++++++++ apps/api/src/index.ts | 87 +++++++++-- apps/api/src/services/registries.ts | 18 ++- apps/api/src/services/scans.ts | 79 ++++++++-- apps/api/src/worker/scan-worker.ts | 107 +++++++++++++ apps/api/tests/capture-worker.test.ts | 176 ++++++++++++++++++++++ apps/api/tests/persistence.test.ts | 32 ++-- pnpm-lock.yaml | 10 ++ 13 files changed, 757 insertions(+), 52 deletions(-) create mode 100644 apps/api/migrations/002_tenant.sql create mode 100644 apps/api/src/audit/aud0.ts create mode 100644 apps/api/src/capture/capture.ts create mode 100644 apps/api/src/worker/scan-worker.ts create mode 100644 apps/api/tests/capture-worker.test.ts diff --git a/Dockerfile b/Dockerfile index a4e0d87..6fb564f 100644 --- a/Dockerfile +++ b/Dockerfile @@ -13,7 +13,8 @@ RUN pnpm install --frozen-lockfile FROM node:22-alpine AS runtime WORKDIR /app -RUN apk add --no-cache wget +# chromium: deterministic browser capture (scan worker, playwright-core) +RUN apk add --no-cache wget chromium COPY package.json pnpm-workspace.yaml pnpm-lock.yaml ./ COPY packages ./packages diff --git a/apps/api/migrations/002_tenant.sql b/apps/api/migrations/002_tenant.sql new file mode 100644 index 0000000..0df9189 --- /dev/null +++ b/apps/api/migrations/002_tenant.sql @@ -0,0 +1,8 @@ +-- C0PY 002 — tenant binding (resourceowner id from Zitadel introspection; +-- resolved server-side, NEVER from client input). + +ALTER TABLE registries ADD COLUMN IF NOT EXISTS tenant_id TEXT NOT NULL DEFAULT ''; +ALTER TABLE scans ADD COLUMN IF NOT EXISTS tenant_id TEXT NOT NULL DEFAULT ''; + +CREATE INDEX IF NOT EXISTS registries_tenant_idx ON registries(tenant_id); +CREATE INDEX IF NOT EXISTS scans_tenant_idx ON scans(tenant_id); diff --git a/apps/api/package.json b/apps/api/package.json index 486ca68..a6989af 100644 --- a/apps/api/package.json +++ b/apps/api/package.json @@ -16,6 +16,7 @@ "@siax/c0py-types": "workspace:*", "fastify": "^5.12.0", "pg": "^8.23.0", + "playwright-core": "^1.63.0", "zod": "^3.24.0" }, "devDependencies": { diff --git a/apps/api/src/audit/aud0.ts b/apps/api/src/audit/aud0.ts new file mode 100644 index 0000000..aa4f842 --- /dev/null +++ b/apps/api/src/audit/aud0.ts @@ -0,0 +1,62 @@ +import type { FastifyBaseLogger } from "fastify"; + +export interface Aud0Event { + event_type: string; + tenant_id: string; + project_id?: string; + actor_id?: string; + app_id: string; + capability_id: string; + resource_type: string; + resource_id: string; + risk_level: "L0" | "L1" | "L2" | "L3" | "L4" | "L5"; + payload?: Record; + source?: Record; +} + +export interface Aud0Emitter { + (event: Aud0Event): void; +} + +export function createAud0Emitter(deps: { + baseUrl?: string; + authToken?: string; + logger?: FastifyBaseLogger; + fetchImpl?: typeof fetch; +}): Aud0Emitter { + const fetchImpl = deps.fetchImpl ?? fetch; + + // Fire-and-forget per AUD0 doctrine: an AUD0 outage must never block the + // caller's real work. Failures are logged and dropped (advisory). + return (event) => { + if (!deps.baseUrl || !deps.authToken) { + deps.logger?.debug({ event_type: event.event_type }, "AUD0 not configured — event dropped"); + return; + } + const url = `${deps.baseUrl.replace(/\/$/, "")}/v1/aud0/events`; + void fetchImpl(url, { + method: "POST", + headers: { + Authorization: `Bearer ${deps.authToken}`, + "x-tenant-id": event.tenant_id, + "x-actor-id": event.actor_id ?? "", + "Content-Type": "application/json", + }, + body: JSON.stringify(event), + signal: AbortSignal.timeout(5000), + }) + .then((res) => { + if (!res.ok) { + deps.logger?.warn( + { status: res.status, event_type: event.event_type }, + "AUD0 ingest non-2xx (advisory)", + ); + } + }) + .catch((err) => { + deps.logger?.warn({ err, event_type: event.event_type }, "AUD0 ingest failed (advisory)"); + }); + }; +} + +export const noopAud0: Aud0Emitter = () => {}; diff --git a/apps/api/src/auth/introspection.ts b/apps/api/src/auth/introspection.ts index 667a4f8..a6a9041 100644 --- a/apps/api/src/auth/introspection.ts +++ b/apps/api/src/auth/introspection.ts @@ -5,6 +5,7 @@ export interface IntrospectionResult { sub?: string; aud?: string[] | string; scope?: string; + resourceOwnerId?: string; reason?: string; } @@ -75,6 +76,20 @@ export function createIntrospector(deps: IntrospectionDeps) { return { active: false, reason: "audience_mismatch" }; } - return { active: true, sub: body.sub, aud: body.aud, scope: body.scope }; + // Tenant binding: the resource owner (org) claim is asserted by Zitadel, + // never taken from any client input. Missing claim is surfaced so the + // caller can deny-by-default. + const raw = body as unknown as Record; + const resourceOwnerId = typeof raw["urn:zitadel:iam:user:resourceowner:id"] === "string" + ? String(raw["urn:zitadel:iam:user:resourceowner:id"]) + : undefined; + + return { + active: true, + sub: body.sub, + aud: body.aud, + scope: body.scope, + resourceOwnerId, + }; }; } diff --git a/apps/api/src/capture/capture.ts b/apps/api/src/capture/capture.ts new file mode 100644 index 0000000..52aed65 --- /dev/null +++ b/apps/api/src/capture/capture.ts @@ -0,0 +1,209 @@ +import { randomUUID } from "node:crypto"; +import type { EvidenceRecord } from "@siax/c0py-types"; + +export interface CapturedPage { + finalUrl: string; + status: number; + contentType: string | null; + title: string | null; + description: string | null; + lang: string | null; + charset: string | null; + viewportMeta: string | null; + domCounts: Record; + stylesheets: string[]; + images: string[]; + computedStyle: Record>; + networkRequests: { url: string; method: string; resourceType: string; status: number | null }[]; + screenshotBase64: string | null; +} + +export interface CaptureEvidence { + evidence: EvidenceRecord[]; + summary: { + url: string; + status: number; + title: string | null; + domElementCount: number; + stylesheetCount: number; + imageCount: number; + networkRequestCount: number; + screenshotCaptured: boolean; + }; +} + +export const DEFAULT_PROFILE_ID = "c0py-default-public"; +export const DEFAULT_VIEWPORT = { width: 1280, height: 800, deviceScaleFactor: 1 }; + +function record( + targetId: string, + route: string, + kind: EvidenceRecord["kind"], + payload: T, + tool: string, +): EvidenceRecord { + return { + id: randomUUID(), + targetId, + kind, + route, + capturedAt: new Date().toISOString(), + method: "deterministic-browser-capture", + confidence: "measured", + captureProfileId: DEFAULT_PROFILE_ID, + viewport: DEFAULT_VIEWPORT, + source: { tool, toolVersion: process.env.CAPTURE_TOOL_VERSION ?? "unknown" }, + payload, + }; +} + +// Pure builder — unit-testable without a browser. +export function buildEvidence(targetId: string, page: CapturedPage): CaptureEvidence { + const tool = "playwright-core+chromium"; + const evidence: EvidenceRecord[] = [ + record(targetId, page.finalUrl, "runtime-html", { + finalUrl: page.finalUrl, + status: page.status, + contentType: page.contentType, + title: page.title, + description: page.description, + lang: page.lang, + charset: page.charset, + viewportMeta: page.viewportMeta, + }, tool), + record(targetId, page.finalUrl, "dom", { + counts: page.domCounts, + }, tool), + record(targetId, page.finalUrl, "computed-style", page.computedStyle, tool), + record(targetId, page.finalUrl, "stylesheet", { + stylesheets: page.stylesheets, + }, tool), + record(targetId, page.finalUrl, "asset", { + images: page.images, + }, tool), + record(targetId, page.finalUrl, "network-request", { + requests: page.networkRequests, + }, tool), + ]; + if (page.screenshotBase64) { + evidence.push( + record(targetId, page.finalUrl, "screenshot", { + base64: page.screenshotBase64, + contentType: "image/jpeg", + }, tool), + ); + } + return { + evidence, + summary: { + url: page.finalUrl, + status: page.status, + title: page.title, + domElementCount: Object.values(page.domCounts).reduce((a, b) => a + b, 0), + stylesheetCount: page.stylesheets.length, + imageCount: page.images.length, + networkRequestCount: page.networkRequests.length, + screenshotCaptured: page.screenshotBase64 !== null, + }, + }; +} + +// Browser-backed capture. Playwright is the canonical runtime per AGENTS.md; +// chromium comes from the system package (apk chromium) via executablePath. +export async function capturePage( + url: string, + opts: { executablePath?: string; timeoutMs: number; maxRequests?: number }, +): Promise { + const { chromium } = await import("playwright-core"); + const executablePath = + opts.executablePath ?? process.env.CHROMIUM_EXECUTABLE_PATH ?? "/usr/bin/chromium-browser"; + const browser = await chromium.launch({ + executablePath, + args: ["--no-sandbox", "--disable-dev-shm-usage"], + }); + const maxRequests = opts.maxRequests ?? 100; + try { + const context = await browser.newContext({ viewport: DEFAULT_VIEWPORT }); + const page = await context.newPage(); + const networkRequests: CapturedPage["networkRequests"] = []; + page.on("response", (res) => { + if (networkRequests.length < maxRequests) { + networkRequests.push({ + url: res.url(), + method: res.request().method(), + resourceType: res.request().resourceType(), + status: res.status(), + }); + } + }); + const response = await page.goto(url, { timeout: opts.timeoutMs, waitUntil: "domcontentloaded" }); + await page.waitForTimeout(500); + const captured = await page.evaluate(() => { + const count = (sel: string) => document.querySelectorAll(sel).length; + const pick = (el: Element | null, props: string[]) => { + if (!el) return Object.fromEntries(props.map((p) => [p, null])); + const cs = getComputedStyle(el); + return Object.fromEntries(props.map((p) => [p, cs.getPropertyValue(p) || null])); + }; + const body = pick(document.body, ["font-family", "color", "background-color", "font-size"]); + const h1 = pick(document.querySelector("h1"), ["font-family", "font-size", "font-weight", "color", "margin"]); + return { + title: document.title, + description: document.querySelector('meta[name="description"]')?.getAttribute("content") ?? null, + lang: document.documentElement.getAttribute("lang"), + charset: document.characterSet, + viewportMeta: document.querySelector('meta[name="viewport"]')?.getAttribute("content") ?? null, + domCounts: { + script: count("script"), + link: count("link"), + style: count("style"), + img: count("img"), + svg: count("svg"), + form: count("form"), + input: count("input"), + button: count("button"), + nav: count("nav"), + header: count("header"), + footer: count("footer"), + main: count("main"), + section: count("section"), + h1: count("h1"), + h2: count("h2"), + h3: count("h3"), + a: count("a"), + table: count("table"), + iframe: count("iframe"), + video: count("video"), + canvas: count("canvas"), + }, + stylesheets: Array.from(document.querySelectorAll('link[rel="stylesheet"]')).map((l) => + (l as HTMLLinkElement).href, + ), + images: Array.from(document.querySelectorAll("img")).slice(0, 50).map((i) => i.src), + computedStyle: { body, h1 }, + }; + }); + let screenshotBase64: string | null = null; + try { + screenshotBase64 = await page.screenshot({ + fullPage: true, + type: "jpeg", + quality: 60, + timeout: Math.min(opts.timeoutMs, 10_000), + }).then((b) => b.toString("base64")); + } catch { + // Screenshot failure is an evidence gap, not a capture failure. + screenshotBase64 = null; + } + return { + finalUrl: page.url(), + status: response?.status() ?? 0, + contentType: response?.headers()["content-type"] ?? null, + ...captured, + networkRequests, + screenshotBase64, + }; + } finally { + await browser.close().catch(() => {}); + } +} diff --git a/apps/api/src/index.ts b/apps/api/src/index.ts index c0dc95f..fde4845 100644 --- a/apps/api/src/index.ts +++ b/apps/api/src/index.ts @@ -8,6 +8,8 @@ import { getRegistry, } from "./services/registries.js"; import { createScan, listScans, getScan } from "./services/scans.js"; +import { createAud0Emitter, noopAud0 } from "./audit/aud0.js"; +import { startScanWorker } from "./worker/scan-worker.js"; async function main() { const env = envSchema.parse(process.env); @@ -54,16 +56,26 @@ async function main() { .send({ error: "Unauthorized" }); return; } - (request as never as { user?: { sub?: string } }).user = { sub: result.sub }; + // Tenant binding: org claim asserted by Zitadel; deny-by-default when absent. + if (!result.resourceOwnerId) { + await reply.code(403).send({ error: "Forbidden" }); + return; + } + (request as unknown as { user?: { sub?: string; tenantId?: string } }).user = { + sub: result.sub, + tenantId: result.resourceOwnerId, + }; } }); - const requireSub = (request: FastifyRequest): string => { - const user = (request as unknown as { user?: { sub?: string } }).user; - if (!user?.sub) { - throw Object.assign(new Error("missing subject"), { statusCode: 401 }); + const requireUser = ( + request: FastifyRequest, + ): { sub: string; tenantId: string } => { + const user = (request as unknown as { user?: { sub?: string; tenantId?: string } }).user; + if (!user?.sub || !user?.tenantId) { + throw Object.assign(new Error("missing subject or tenant"), { statusCode: 401 }); } - return user.sub; + return { sub: user.sub, tenantId: user.tenantId }; }; const dbError = (reply: FastifyReply, err: unknown) => { @@ -71,10 +83,11 @@ async function main() { return reply.code(500).send({ error: "Internal Server Error" }); }; - // Registry endpoints (owner-scoped by introspected subject) + // Registry endpoints (owner + tenant scoped by introspected identity) server.get("/v1/c0py/registries", async (request, reply) => { try { - return { registries: await listRegistries(getPool(), requireSub(request)) }; + const u = requireUser(request); + return { registries: await listRegistries(getPool(), u.sub, u.tenantId) }; } catch (err) { return dbError(reply, err); } @@ -89,11 +102,23 @@ async function main() { return; } try { - const registry = await createRegistry(getPool(), requireSub(request), { + const u = requireUser(request); + const registry = await createRegistry(getPool(), u.sub, u.tenantId, { name, url, config: body.config, }); + aud0({ + event_type: "c0py.registry_created", + tenant_id: u.tenantId, + actor_id: u.sub, + app_id: "c0py", + capability_id: "c0py.manage_registries", + resource_type: "registry", + resource_id: registry.id, + risk_level: "L1", + payload: { name, url }, + }); await reply.code(201).send({ registry }); } catch (err) { return dbError(reply, err); @@ -103,7 +128,8 @@ async function main() { server.get("/v1/c0py/registries/:id", async (request, reply) => { const { id } = request.params as Record; try { - const registry = await getRegistry(getPool(), requireSub(request), id); + const u = requireUser(request); + const registry = await getRegistry(getPool(), u.sub, u.tenantId, id); if (!registry) { await reply.code(404).send({ error: "Not Found" }); return; @@ -114,7 +140,7 @@ async function main() { } }); - // Scan endpoints (owner-scoped, registry ownership enforced in SQL) + // Scan endpoints (owner + tenant scoped, registry ownership enforced in SQL) server.post("/v1/c0py/scans", async (request, reply) => { const body = (request.body as Record) ?? {}; const registryId = typeof body.registryId === "string" ? body.registryId.trim() : ""; @@ -123,11 +149,23 @@ async function main() { return; } try { - const scan = await createScan(getPool(), requireSub(request), registryId); + const u = requireUser(request); + const scan = await createScan(getPool(), u.sub, u.tenantId, registryId); if (!scan) { await reply.code(404).send({ error: "Not Found" }); return; } + aud0({ + event_type: "c0py.scan_created", + tenant_id: u.tenantId, + actor_id: u.sub, + app_id: "c0py", + capability_id: "c0py.scan_run", + resource_type: "scan", + resource_id: scan.id, + risk_level: "L1", + payload: { registryId }, + }); await reply.code(201).send({ scan }); } catch (err) { return dbError(reply, err); @@ -137,7 +175,8 @@ async function main() { server.get("/v1/c0py/scans", async (request, reply) => { const query = request.query as Record; try { - return { scans: await listScans(getPool(), requireSub(request), query.registryId) }; + const u = requireUser(request); + return { scans: await listScans(getPool(), u.sub, u.tenantId, query.registryId) }; } catch (err) { return dbError(reply, err); } @@ -146,7 +185,8 @@ async function main() { server.get("/v1/c0py/scans/:id", async (request, reply) => { const { id } = request.params as Record; try { - const scan = await getScan(getPool(), requireSub(request), id); + const u = requireUser(request); + const scan = await getScan(getPool(), u.sub, u.tenantId, id); if (!scan) { await reply.code(404).send({ error: "Not Found" }); return; @@ -160,6 +200,25 @@ async function main() { await server.listen({ port: env.PORT, host: "0.0.0.0" }); console.log(`c0py API listening on port ${env.PORT}`); + const aud0 = + env.AUD0_BASE_URL && env.AUD0_API_KEY + ? createAud0Emitter({ + baseUrl: env.AUD0_BASE_URL, + authToken: env.AUD0_API_KEY, + logger: server.log, + }) + : noopAud0; + + if (env.NODE_ENV !== "test") { + startScanWorker({ + db: getPool(), + logger: server.log, + aud0, + crawlTimeoutMs: env.CRAWL_TIMEOUT_MS, + }); + console.log("c0py scan worker started"); + } + for (const signal of ["SIGTERM", "SIGINT"] as const) { process.on(signal, async () => { await server.close(); diff --git a/apps/api/src/services/registries.ts b/apps/api/src/services/registries.ts index 28e7a62..54ea681 100644 --- a/apps/api/src/services/registries.ts +++ b/apps/api/src/services/registries.ts @@ -1,6 +1,7 @@ export interface RegistryRow { id: string; owner_sub: string; + tenant_id: string; name: string; url: string; config: unknown; @@ -27,10 +28,10 @@ function toRegistry(row: Record): { }; } -export async function listRegistries(db: PoolLike, ownerSub: string) { +export async function listRegistries(db: PoolLike, ownerSub: string, tenantId: string) { const res = await db.query( - "SELECT id, owner_sub, name, url, config, created_at FROM registries WHERE owner_sub = $1 ORDER BY created_at DESC", - [ownerSub], + "SELECT id, owner_sub, tenant_id, name, url, config, created_at FROM registries WHERE owner_sub = $1 AND tenant_id = $2 ORDER BY created_at DESC", + [ownerSub, tenantId], ); return res.rows.map(toRegistry); } @@ -38,19 +39,20 @@ export async function listRegistries(db: PoolLike, ownerSub: string) { export async function createRegistry( db: PoolLike, ownerSub: string, + tenantId: string, input: { name: string; url: string; config?: unknown }, ) { const res = await db.query( - "INSERT INTO registries (owner_sub, name, url, config) VALUES ($1, $2, $3, $4::jsonb) RETURNING id, owner_sub, name, url, config, created_at", - [ownerSub, input.name, input.url, JSON.stringify(input.config ?? {})], + "INSERT INTO registries (owner_sub, tenant_id, name, url, config) VALUES ($1, $2, $3, $4, $5::jsonb) RETURNING id, owner_sub, tenant_id, name, url, config, created_at", + [ownerSub, tenantId, input.name, input.url, JSON.stringify(input.config ?? {})], ); return toRegistry(res.rows[0]); } -export async function getRegistry(db: PoolLike, ownerSub: string, id: string) { +export async function getRegistry(db: PoolLike, ownerSub: string, tenantId: string, id: string) { const res = await db.query( - "SELECT id, owner_sub, name, url, config, created_at FROM registries WHERE owner_sub = $1 AND id = $2", - [ownerSub, id], + "SELECT id, owner_sub, tenant_id, name, url, config, created_at FROM registries WHERE owner_sub = $1 AND tenant_id = $2 AND id = $3", + [ownerSub, tenantId, id], ); return res.rows[0] ? toRegistry(res.rows[0]) : null; } diff --git a/apps/api/src/services/scans.ts b/apps/api/src/services/scans.ts index 6e0f178..01ebfd1 100644 --- a/apps/api/src/services/scans.ts +++ b/apps/api/src/services/scans.ts @@ -2,7 +2,8 @@ export interface ScanRow { id: string; registry_id: string; owner_sub: string; - status: string; + tenant_id: string; + status: "pending" | "running" | "completed" | "failed"; result: unknown; created_at: Date; updated_at: Date; @@ -33,38 +34,88 @@ function toScan(row: Record): { export async function createScan( db: PoolLike, ownerSub: string, + tenantId: string, registryId: string, ): Promise | null> { // Ownership check happens in SQL: the scan may only reference a registry - // owned by the same subject. Inserting otherwise returns no rows (fail-closed). + // owned by the same subject and tenant. Otherwise no rows return (fail-closed). const res = await db.query( - `INSERT INTO scans (registry_id, owner_sub, status) - SELECT id, $1, 'pending' FROM registries WHERE id = $2 AND owner_sub = $1 - RETURNING id, registry_id, owner_sub, status, result, created_at, updated_at`, - [ownerSub, registryId], + `INSERT INTO scans (registry_id, owner_sub, tenant_id, status) + SELECT id, $1, $2, 'pending' FROM registries WHERE id = $3 AND owner_sub = $1 AND tenant_id = $2 + RETURNING id, registry_id, owner_sub, tenant_id, status, result, created_at, updated_at`, + [ownerSub, tenantId, registryId], ); return res.rows[0] ? toScan(res.rows[0]) : null; } -export async function listScans(db: PoolLike, ownerSub: string, registryId?: string) { +export async function listScans( + db: PoolLike, + ownerSub: string, + tenantId: string, + registryId?: string, +) { if (registryId !== undefined) { const res = await db.query( - "SELECT id, registry_id, owner_sub, status, result, created_at, updated_at FROM scans WHERE owner_sub = $1 AND registry_id = $2 ORDER BY created_at DESC", - [ownerSub, registryId], + "SELECT id, registry_id, owner_sub, tenant_id, status, result, created_at, updated_at FROM scans WHERE owner_sub = $1 AND tenant_id = $2 AND registry_id = $3 ORDER BY created_at DESC", + [ownerSub, tenantId, registryId], ); return res.rows.map(toScan); } const res = await db.query( - "SELECT id, registry_id, owner_sub, status, result, created_at, updated_at FROM scans WHERE owner_sub = $1 ORDER BY created_at DESC", - [ownerSub], + "SELECT id, registry_id, owner_sub, tenant_id, status, result, created_at, updated_at FROM scans WHERE owner_sub = $1 AND tenant_id = $2 ORDER BY created_at DESC", + [ownerSub, tenantId], ); return res.rows.map(toScan); } -export async function getScan(db: PoolLike, ownerSub: string, id: string) { +export async function getScan(db: PoolLike, ownerSub: string, tenantId: string, id: string) { const res = await db.query( - "SELECT id, registry_id, owner_sub, status, result, created_at, updated_at FROM scans WHERE owner_sub = $1 AND id = $2", - [ownerSub, id], + "SELECT id, registry_id, owner_sub, tenant_id, status, result, created_at, updated_at FROM scans WHERE owner_sub = $1 AND tenant_id = $2 AND id = $3", + [ownerSub, tenantId, id], ); return res.rows[0] ? toScan(res.rows[0]) : null; } + +// --- worker-facing (subject is the scan's owner, resolved server-side) --- + +export interface PendingScan { + id: string; + registry_id: string; + owner_sub: string; + tenant_id: string; +} + +export async function claimPendingScan(db: PoolLike): Promise { + // Atomic claim: flip one pending scan to running and return it. Fail-closed + // against double-processing (UPDATE ... WHERE status='pending' races resolve + // to a single winner per scan). + const res = await db.query( + `UPDATE scans SET status = 'running', updated_at = now() + WHERE id = (SELECT id FROM scans WHERE status = 'pending' ORDER BY created_at LIMIT 1 FOR UPDATE SKIP LOCKED) + RETURNING id, registry_id, owner_sub, tenant_id`, + ); + return (res.rows[0] as unknown as PendingScan) ?? null; +} + +export async function getScanTargetUrl(db: PoolLike, registryId: string): Promise { + const res = await db.query("SELECT url FROM registries WHERE id = $1", [registryId]); + return res.rows[0] ? String(res.rows[0].url) : null; +} + +export async function completeScan( + db: PoolLike, + scanId: string, + result: unknown, +): Promise { + await db.query( + "UPDATE scans SET status = 'completed', result = $2::jsonb, updated_at = now() WHERE id = $1 AND status = 'running'", + [scanId, JSON.stringify(result)], + ); +} + +export async function failScan(db: PoolLike, scanId: string, error: string): Promise { + await db.query( + "UPDATE scans SET status = 'failed', result = $2::jsonb, updated_at = now() WHERE id = $1 AND status = 'running'", + [scanId, JSON.stringify({ error })], + ); +} diff --git a/apps/api/src/worker/scan-worker.ts b/apps/api/src/worker/scan-worker.ts new file mode 100644 index 0000000..07952c1 --- /dev/null +++ b/apps/api/src/worker/scan-worker.ts @@ -0,0 +1,107 @@ +import type { FastifyBaseLogger } from "fastify"; +import { claimPendingScan, getScanTargetUrl, completeScan, failScan } from "../services/scans.js"; +import type { PoolLike } from "../services/registries.js"; +import { capturePage, buildEvidence } from "../capture/capture.js"; +import type { Aud0Emitter } from "../audit/aud0.js"; + +export interface ScanWorkerDeps { + db: PoolLike; + logger: FastifyBaseLogger; + aud0: Aud0Emitter; + capture?: (url: string, timeoutMs: number) => Promise<{ evidence: unknown[]; summary: unknown }>; + pollIntervalMs?: number; + crawlTimeoutMs?: number; + enabled?: boolean; +} + +export function startScanWorker(deps: ScanWorkerDeps): { + processOne: () => Promise; + stop: () => Promise; +} { + const interval = deps.pollIntervalMs ?? 2000; + const crawlTimeoutMs = deps.crawlTimeoutMs ?? 30_000; + const capture = + deps.capture ?? + (async (url: string, timeoutMs: number) => { + const page = await capturePage(url, { timeoutMs }); + return buildEvidence("scan", page); + }); + + let stopped = false; + let timer: NodeJS.Timeout | null = null; + + async function processOne(): Promise { + let scan: Awaited> = null; + try { + scan = await claimPendingScan(deps.db); + } catch (err) { + deps.logger.error({ err }, "scan worker claim failed"); + return false; + } + if (!scan) return false; + + const started = Date.now(); + deps.logger.info({ scanId: scan.id, url: scan.registry_id }, "scan running"); + try { + const url = await getScanTargetUrl(deps.db, scan.registry_id); + if (!url) throw new Error("registry url not found"); + const out = await capture(url, crawlTimeoutMs); + await completeScan(deps.db, scan.id, out); + deps.aud0({ + event_type: "c0py.scan_completed", + tenant_id: scan.tenant_id, + actor_id: scan.owner_sub, + app_id: "c0py", + capability_id: "c0py.scan_run", + resource_type: "scan", + resource_id: scan.id, + risk_level: "L2", + payload: { summary: out.summary, durationMs: Date.now() - started }, + }); + deps.logger.info({ scanId: scan.id, durationMs: Date.now() - started }, "scan completed"); + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + try { + await failScan(deps.db, scan.id, message); + } catch (err2) { + deps.logger.error({ err: err2 }, "failScan write failed"); + } + deps.aud0({ + event_type: "c0py.scan_failed", + tenant_id: scan.tenant_id, + actor_id: scan.owner_sub, + app_id: "c0py", + capability_id: "c0py.scan_run", + resource_type: "scan", + resource_id: scan.id, + risk_level: "L2", + payload: { error: message }, + }); + deps.logger.warn({ scanId: scan.id, err: message }, "scan failed"); + } + return true; + } + + async function loop(): Promise { + while (!stopped) { + const didWork = await processOne(); + if (!didWork) { + await new Promise((resolve) => { + timer = setTimeout(resolve, interval); + }); + } + } + } + + if (deps.enabled !== false) { + void loop(); + } + + return { + processOne, + stop: async () => { + stopped = true; + if (timer) clearTimeout(timer); + }, + }; +} diff --git a/apps/api/tests/capture-worker.test.ts b/apps/api/tests/capture-worker.test.ts new file mode 100644 index 0000000..4b2e0e5 --- /dev/null +++ b/apps/api/tests/capture-worker.test.ts @@ -0,0 +1,176 @@ +import { describe, it, expect, vi } from "vitest"; +import { buildEvidence, type CapturedPage } from "../src/capture/capture.js"; +import { createAud0Emitter, noopAud0 } from "../src/audit/aud0.js"; +import { startScanWorker } from "../src/worker/scan-worker.js"; +import type { PoolLike } from "../src/services/registries.js"; +import type { FastifyBaseLogger } from "fastify"; + +const PAGE: CapturedPage = { + finalUrl: "https://example.com/", + status: 200, + contentType: "text/html", + title: "Example", + description: "desc", + lang: "en", + charset: "UTF-8", + viewportMeta: "width=device-width", + domCounts: { script: 3, img: 2, a: 10 }, + stylesheets: ["https://example.com/app.css"], + images: ["https://example.com/logo.png"], + computedStyle: { body: { "font-family": "sans-serif" }, h1: null }, + networkRequests: [{ url: "https://example.com/", method: "GET", resourceType: "document", status: 200 }], + screenshotBase64: "abc", +}; + +describe("capture evidence builder", () => { + it("produces measured EvidenceRecords per kind", () => { + const { evidence, summary } = buildEvidence("reg-1", PAGE); + const kinds = evidence.map((e) => e.kind); + expect(kinds).toEqual( + expect.arrayContaining(["runtime-html", "dom", "computed-style", "stylesheet", "asset", "network-request", "screenshot"]), + ); + for (const e of evidence) { + expect(e.confidence).toBe("measured"); + expect(e.targetId).toBe("reg-1"); + expect(e.source.tool).toContain("playwright-core"); + expect(e.method).toBe("deterministic-browser-capture"); + expect(e.capturedAt).toBeTruthy(); + } + expect(summary).toMatchObject({ + url: "https://example.com/", + status: 200, + domElementCount: 15, + stylesheetCount: 1, + screenshotCaptured: true, + }); + }); + + it("omits screenshot evidence when capture missed it (evidence gap explicit)", () => { + const { evidence, summary } = buildEvidence("reg-1", { ...PAGE, screenshotBase64: null }); + expect(evidence.some((e) => e.kind === "screenshot")).toBe(false); + expect(summary.screenshotCaptured).toBe(false); + }); +}); + +describe("aud0 emitter", () => { + const base = { + event_type: "c0py.scan_completed", + tenant_id: "org-1", + actor_id: "u1", + app_id: "c0py", + capability_id: "c0py.scan_run", + resource_type: "scan", + resource_id: "s1", + risk_level: "L2" as const, + }; + + it("no-op when unconfigured", () => { + const fetchImpl = vi.fn(); + const emit = createAud0Emitter({ fetchImpl: fetchImpl as unknown as typeof fetch }); + emit(base); + expect(fetchImpl).not.toHaveBeenCalled(); + }); + + it("posts event with tenant header (fire-and-forget)", async () => { + const fetchImpl = vi.fn().mockResolvedValue(new Response("{}", { status: 200 })); + const emit = createAud0Emitter({ + baseUrl: "https://aud0.siax.io", + authToken: "tok", + fetchImpl: fetchImpl as unknown as typeof fetch, + }); + emit(base); + await new Promise((r) => setTimeout(r, 10)); + expect(fetchImpl).toHaveBeenCalledWith( + "https://aud0.siax.io/v1/aud0/events", + expect.objectContaining({ + method: "POST", + headers: expect.objectContaining({ "x-tenant-id": "org-1" }), + }), + ); + }); + + it("never throws on network failure (advisory)", async () => { + const fetchImpl = vi.fn().mockRejectedValue(new Error("down")); + const emit = createAud0Emitter({ + baseUrl: "https://aud0.siax.io", + authToken: "tok", + fetchImpl: fetchImpl as unknown as typeof fetch, + }); + expect(() => emit(base)).not.toThrow(); + await new Promise((r) => setTimeout(r, 10)); + }); + + it("noop emitter accepts anything", () => { + expect(() => noopAud0(base)).not.toThrow(); + }); +}); + +describe("scan worker", () => { + const logger = { info: vi.fn(), warn: vi.fn(), error: vi.fn(), debug: vi.fn() } as unknown as FastifyBaseLogger; + + function dbWith(queries: { text: string; rows: Record[] }[]) { + const calls: string[] = []; + const db: PoolLike = { + query: (text) => { + calls.push(text); + const row = queries.find((q) => text.includes(q.text)); + return Promise.resolve({ rows: row?.rows ?? [] }); + }, + }; + return { db, calls }; + } + + it("claims, captures, completes and emits AUD0", async () => { + const { db } = dbWith([ + { text: "UPDATE scans SET status = 'running'", rows: [{ id: "s1", registry_id: "r1", owner_sub: "u1", tenant_id: "org-1" }] }, + { text: "SELECT url FROM registries", rows: [{ url: "https://example.com" }] }, + { text: "UPDATE scans SET status = 'completed'", rows: [] }, + ]); + const aud0 = vi.fn(); + const worker = startScanWorker({ + db, + logger, + aud0: aud0 as never, + enabled: false, + capture: async (url) => { + expect(url).toBe("https://example.com"); + return buildEvidence("scan", PAGE); + }, + }); + const did = await worker.processOne(); + expect(did).toBe(true); + expect(aud0).toHaveBeenCalledWith( + expect.objectContaining({ event_type: "c0py.scan_completed", tenant_id: "org-1", risk_level: "L2" }), + ); + }); + + it("marks scan failed and emits c0py.scan_failed on capture error", async () => { + const { db } = dbWith([ + { text: "UPDATE scans SET status = 'running'", rows: [{ id: "s1", registry_id: "r1", owner_sub: "u1", tenant_id: "org-1" }] }, + { text: "SELECT url FROM registries", rows: [] }, + { text: "UPDATE scans SET status = 'failed'", rows: [] }, + ]); + const aud0 = vi.fn(); + const worker = startScanWorker({ + db, + logger, + aud0: aud0 as never, + enabled: false, + capture: async () => { + throw new Error("boom"); + }, + }); + const did = await worker.processOne(); + expect(did).toBe(true); + expect(aud0).toHaveBeenCalledWith( + expect.objectContaining({ event_type: "c0py.scan_failed" }), + ); + }); + + it("returns false when no pending scan", async () => { + const { db } = dbWith([{ text: "UPDATE scans SET status = 'running'", rows: [] }]); + const worker = startScanWorker({ db, logger, aud0: noopAud0, enabled: false, capture: async () => buildEvidence("s", PAGE) }); + const did = await worker.processOne(); + expect(did).toBe(false); + }); +}); diff --git a/apps/api/tests/persistence.test.ts b/apps/api/tests/persistence.test.ts index 8c39756..f3ae198 100644 --- a/apps/api/tests/persistence.test.ts +++ b/apps/api/tests/persistence.test.ts @@ -26,26 +26,28 @@ const ROW = { }; describe("registries service", () => { - it("listRegistries filters by owner", async () => { + it("listRegistries filters by owner and tenant", async () => { const db = fakePool({ query: () => ({ rows: [ROW] }) }); - const out = await listRegistries(db, "svc-user"); + const out = await listRegistries(db, "svc-user", "org-1"); expect(out).toHaveLength(1); expect(out[0]).toMatchObject({ id: ROW.id, name: "example", createdAt: "2026-09-16T10:00:00.000Z" }); - expect(db.queries[0].values).toEqual(["svc-user"]); + expect(db.queries[0].values).toEqual(["svc-user", "org-1"]); expect(db.queries[0].text).toContain("owner_sub = $1"); + expect(db.queries[0].text).toContain("tenant_id = $2"); }); - it("createRegistry inserts with owner and config jsonb", async () => { + it("createRegistry inserts with owner, tenant and config jsonb", async () => { const db = fakePool({ query: () => ({ rows: [ROW] }) }); - const out = await createRegistry(db, "svc-user", { name: "example", url: "https://example.com" }); + const out = await createRegistry(db, "svc-user", "org-1", { name: "example", url: "https://example.com" }); expect(out.id).toBe(ROW.id); expect(db.queries[0].values?.[0]).toBe("svc-user"); + expect(db.queries[0].values?.[1]).toBe("org-1"); expect(db.queries[0].text).toContain("RETURNING"); }); it("getRegistry returns null when not owner (fail-closed)", async () => { const db = fakePool({ query: () => ({ rows: [] }) }); - const out = await getRegistry(db, "other-user", ROW.id); + const out = await getRegistry(db, "other-user", "org-1", ROW.id); expect(out).toBeNull(); }); }); @@ -54,34 +56,36 @@ describe("scans service", () => { const SCAN_ROW = { ...ROW, registry_id: ROW.id, + tenant_id: "org-1", status: "pending", result: {}, updated_at: ROW.created_at, }; - it("createScan enforces registry ownership in SQL and returns null otherwise", async () => { + it("createScan enforces registry ownership+tenant in SQL and returns null otherwise", async () => { const db = fakePool({ query: () => ({ rows: [] }) }); - const out = await createScan(db, "svc-user", ROW.id); + const out = await createScan(db, "svc-user", "org-1", ROW.id); expect(out).toBeNull(); expect(db.queries[0].text).toContain("owner_sub = $1"); + expect(db.queries[0].text).toContain("tenant_id = $2"); }); it("createScan returns scan for owned registry", async () => { const db = fakePool({ query: () => ({ rows: [SCAN_ROW] }) }); - const out = await createScan(db, "svc-user", ROW.id); + const out = await createScan(db, "svc-user", "org-1", ROW.id); expect(out).toMatchObject({ id: SCAN_ROW.id, registryId: SCAN_ROW.id, status: "pending" }); }); it("listScans supports optional registryId filter", async () => { const db = fakePool({ query: () => ({ rows: [SCAN_ROW] }) }); - await listScans(db, "svc-user"); - expect(db.queries[0].text).not.toContain("registry_id = $2"); - await listScans(db, "svc-user", SCAN_ROW.id); - expect(db.queries[1].values).toEqual(["svc-user", SCAN_ROW.id]); + await listScans(db, "svc-user", "org-1"); + expect(db.queries[0].text).not.toContain("registry_id = $3"); + await listScans(db, "svc-user", "org-1", SCAN_ROW.id); + expect(db.queries[1].values).toEqual(["svc-user", "org-1", SCAN_ROW.id]); }); it("getScan returns null for foreign scan", async () => { const db = fakePool({ query: () => ({ rows: [] }) }); - expect(await getScan(db, "other-user", SCAN_ROW.id)).toBeNull(); + expect(await getScan(db, "other-user", "org-1", SCAN_ROW.id)).toBeNull(); }); }); diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index e0880e5..31b785a 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -50,6 +50,9 @@ importers: pg: specifier: ^8.23.0 version: 8.23.0 + playwright-core: + specifier: ^1.63.0 + version: 1.63.0 zod: specifier: ^3.24.0 version: 3.25.76 @@ -1299,6 +1302,11 @@ packages: resolution: {integrity: sha512-r34yH/GlQpKZbU1BvFFqOjhISRo1MNx1tWYsYvmj6KIRHSPMT2+yHOEb1SG6NMvRoHRF0a07kCOox/9yakl1vg==} hasBin: true + playwright-core@1.63.0: + resolution: {integrity: sha512-rYCsBF/M5HjUch52bbtVONEFjv6Xu8sm8h72dNlR5bzIE1fvC/bxgspzkjSfU+MweEMmPM8KJebG6nnyxo5mCg==} + engines: {node: '>=20'} + hasBin: true + postcss@8.4.31: resolution: {integrity: sha512-PS08Iboia9mts/2ygV3eLpY5ghnUcfLV/EXTOW1E2qYxJKGGBUtNjN76FYHnMs36RmARn41bC0AZmn+rR0OVpQ==} engines: {node: ^10 || ^12 || >=14} @@ -2680,6 +2688,8 @@ snapshots: sonic-boom: 4.2.1 thread-stream: 4.2.0 + playwright-core@1.63.0: {} + postcss@8.4.31: dependencies: nanoid: 5.1.16