feat(api): deterministic scan worker + AUD0 wire + tenant binding
CI (SIAX Cloud) / security (pull_request) Successful in 14s
CI (SIAX Cloud) / security (push) Successful in 14s
CI (SIAX Cloud) / contracts (pull_request) Successful in 15s
CI (SIAX Cloud) / contracts (push) Successful in 17s
CI (SIAX Cloud) / quality (push) Successful in 1m9s
CI (SIAX Cloud) / quality (pull_request) Successful in 1m10s

- 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
This commit is contained in:
2026-09-16 22:53:55 +02:00
parent 7fb6043e34
commit 12f6a65b80
13 changed files with 757 additions and 52 deletions
+2 -1
View File
@@ -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
+8
View File
@@ -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);
+1
View File
@@ -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": {
+62
View File
@@ -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<string, unknown>;
source?: Record<string, unknown>;
}
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 = () => {};
+16 -1
View File
@@ -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<string, unknown>;
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,
};
};
}
+209
View File
@@ -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<string, number>;
stylesheets: string[];
images: string[];
computedStyle: Record<string, Record<string, string | null>>;
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<T>(
targetId: string,
route: string,
kind: EvidenceRecord["kind"],
payload: T,
tool: string,
): EvidenceRecord<T> {
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<CapturedPage> {
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(() => {});
}
}
+73 -14
View File
@@ -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<string, string>;
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<string, unknown>) ?? {};
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<string, string | undefined>;
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<string, string>;
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();
+10 -8
View File
@@ -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<string, unknown>): {
};
}
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;
}
+65 -14
View File
@@ -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<string, unknown>): {
export async function createScan(
db: PoolLike,
ownerSub: string,
tenantId: string,
registryId: string,
): Promise<ReturnType<typeof toScan> | 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<PendingScan | null> {
// 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<string | null> {
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<void> {
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<void> {
await db.query(
"UPDATE scans SET status = 'failed', result = $2::jsonb, updated_at = now() WHERE id = $1 AND status = 'running'",
[scanId, JSON.stringify({ error })],
);
}
+107
View File
@@ -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<boolean>;
stop: () => Promise<void>;
} {
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<boolean> {
let scan: Awaited<ReturnType<typeof claimPendingScan>> = 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<void> {
while (!stopped) {
const didWork = await processOne();
if (!didWork) {
await new Promise<void>((resolve) => {
timer = setTimeout(resolve, interval);
});
}
}
}
if (deps.enabled !== false) {
void loop();
}
return {
processOne,
stop: async () => {
stopped = true;
if (timer) clearTimeout(timer);
},
};
}
+176
View File
@@ -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<string, unknown>[] }[]) {
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);
});
});
+18 -14
View File
@@ -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();
});
});
+10
View File
@@ -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