diff --git a/architecture.d/next-private-collections.json b/architecture.d/next-private-collections.json new file mode 100644 index 00000000..91201848 --- /dev/null +++ b/architecture.d/next-private-collections.json @@ -0,0 +1,9 @@ +{ + "reason": "Private collection create and private device enrol (coordinator 2026-10-06: control reviews) are one feature module beside the cloud-copy routes. The proof, current-identity, membership, transaction, enrolment-lookup and refusal helpers that both use move out of cloud-copy-bootstrap.ts into bootstrap-common.ts unchanged, so the two route modules share one audited implementation instead of copying it; most new export declarations are those existing helpers becoming module exports. Both routes go through registerNextCollection/queueNextPolicy and the existing challenge table; no new transport, table or credential store. app.ts mounts the private routes only with MDBASE_NEXT_PRIVATE_BOOTSTRAP=1.", + "growth": { + "productionFiles": 2, + "relativeImports": 14, + "typeScriptExportDeclarations": 23, + "services/server": 2 + } +} diff --git a/changelog.d/next-private-collections.md b/changelog.d/next-private-collections.md new file mode 100644 index 00000000..992aab62 --- /dev/null +++ b/changelog.d/next-private-collections.md @@ -0,0 +1,9 @@ +## Added + +- Create private (end-to-end) collections on the next control plane from an owner's + registered device (`POST /v1/next/collections/private`), behind + `MDBASE_NEXT_PRIVATE_BOOTSTRAP=1`. The genesis enrols only that device, and no + hosted or escrow device is ever enrolled. +- Enrol a registered device of a current member into a private collection with its + SAS commitment (`POST /v1/next/collections/:id/private/devices`). The device holds + no key until an existing keyed device approves it and grants the key. diff --git a/services/server/src/app.ts b/services/server/src/app.ts index 4f67b291..3e364267 100644 --- a/services/server/src/app.ts +++ b/services/server/src/app.ts @@ -24,6 +24,7 @@ import { ProviderRevocationWorker } from "./hosted-capability-lifecycle.js"; import { LogServiceClient } from "./features/next/log-service-client.js"; import { PolicyEmitter } from "./features/next/policy-outbox.js"; import { registerCloudCopyRoutes } from "./features/next/cloud-copy-bootstrap.js"; +import { registerPrivateCollectionRoutes } from "./features/next/private-collections.js"; import { registerNextHostedRoutes } from "./features/next/hosted-routes.js"; import { registerPolicyRecoveryRoutes } from "./features/next/policy-recovery-routes.js"; import { registerLabFixtureRoutes } from "./features/next/lab-fixture-routes.js"; @@ -537,6 +538,7 @@ export async function buildApp(options: BuildOptions) { registerNoisePipeClientRoute(app, { db: options.db, broker: relayBroker }); registerNextHostedRoutes(app, { db: options.db, tokens: options.nextControlPlane.serviceTokens, log: nextLog }); if (options.nextControlPlane.cloudCopyBootstrap) registerCloudCopyRoutes(app, { db: options.db, next: options.nextControlPlane, emitter: nextPolicyEmitter!, log: nextLog, tailscaleAuth: options.tailscaleAuth }); + if (options.nextControlPlane.privateBootstrap) registerPrivateCollectionRoutes(app, { db: options.db, next: options.nextControlPlane, emitter: nextPolicyEmitter!, log: nextLog }); registerPolicyRecoveryRoutes(app, options.db, nextPolicyEmitter!); registerNextRouteRoutes(app, { db: options.db, publicUrl, broker: relayBroker }); if (options.nextControlPlane.labFixtures) registerLabFixtureRoutes(app, { diff --git a/services/server/src/features/next/bootstrap-common.ts b/services/server/src/features/next/bootstrap-common.ts new file mode 100644 index 00000000..544a19ef --- /dev/null +++ b/services/server/src/features/next/bootstrap-common.ts @@ -0,0 +1,159 @@ +// Shared pieces of the next control plane's collection bootstrap routes (cloud copy +// and private): device proofs, current-identity and membership checks under locks, +// bounded transactions, enrolment lookups and refusals. +import { verify } from "node:crypto"; +import type { FastifyReply } from "fastify"; +import type { DatabaseConnection, DatabasePool } from "../../database-types.js"; +import { apiError } from "../../platform/http-errors.js"; +import { ed25519PublicKeyObject } from "./policy-keys.js"; +import type { PolicyOp } from "./policy-wire.js"; + +/** Service devices belong to no account (policy.md: hosted and escrow enrol with the zero account). */ +export const SERVICE_ACCOUNT = "00000000-0000-0000-0000-000000000000"; +export const NIL = SERVICE_ACCOUNT; + +export interface Proof { device_id: string; challenge: string; sig: string } +export interface Device { sign_pk: Buffer; kem_pk: Buffer; noise_pk: Buffer; kind: "desktop" | "cli" } +export type Connector = { id: string; user_id: string }; +/** PostgreSQL lock_timeout or statement_timeout: answer busy (fail closed), never the driver error. */ +export const isLockTimeout = (error: unknown) => ["55P03", "57014"].includes(String((error as { code?: unknown } | null)?.code)); + +export class CreateError extends Error { + constructor(readonly status: number, readonly code: string) { super(code); } +} + +export const lock = (client: DatabaseConnection, collection: string) => + client.query("SELECT pg_advisory_xact_lock(hashtextextended($1::uuid::text, 20261005))", [collection]); + +/** Verify a device's signature over `digest` and consume its challenge. */ +export async function authenticate(client: DatabaseConnection, body: Proof, connector: Connector, digest: (challenge: Uint8Array) => Uint8Array): Promise { + const device = (await client.query( + "SELECT sign_pk, kem_pk, noise_pk, kind FROM next_devices WHERE id = $1 AND connector_id = $2 AND user_id = $3", + [body.device_id, connector.id, connector.user_id] + )).rows[0]; + const challenge = Buffer.from(body.challenge, "hex"); + if (!device || !verify(null, digest(challenge), ed25519PublicKeyObject(device.sign_pk), Buffer.from(body.sig, "hex"))) throw new CreateError(403, "invalid_proof"); + const used = await client.query( + "UPDATE next_device_challenges SET used_at = now() WHERE challenge = $1 AND connector_id = $2 AND used_at IS NULL AND expires_at > now()", + [challenge, connector.id] + ); + if (used.rowCount !== 1) throw new CreateError(403, "invalid_proof"); + return device; +} + +/** + * The connector, account and device are still current, with the exact keys + * authenticated in phase 1; locked until the transaction ends, so a revocation, + * suspension or device removal either happened before (and is refused here) or waits. + */ +export async function currentIdentity(client: DatabaseConnection, connector: Connector, deviceId: string, device: Device): Promise { + const row = await client.query( + `SELECT 1 FROM connectors c JOIN users u ON u.id = c.user_id + JOIN next_devices d ON d.connector_id = c.id AND d.user_id = u.id + WHERE c.id = $1 AND u.id = $2 AND d.id = $3 AND c.revoked_at IS NULL AND u.suspended_at IS NULL + AND d.sign_pk = $4 AND d.kem_pk = $5 AND d.noise_pk = $6 AND d.kind = $7 + FOR SHARE OF c, u, d`, + [connector.id, connector.user_id, deviceId, device.sign_pk, device.kem_pk, device.noise_pk, device.kind] + ); + if (!row.rows.length) throw new CreateError(403, "identity_not_current"); +} + +/** + * The signed-in session is still the account's current credential (not revoked, not + * expired, same session epoch, account not suspended); share-locked until the + * transaction ends, so a sign-out or suspension either happened before or waits. + */ +export async function currentSession(client: DatabaseConnection, session: string, user: string): Promise { + const row = await client.query( + `SELECT 1 FROM sessions s JOIN users u ON u.id = s.user_id + WHERE s.id = $1 AND u.id = $2 AND s.revoked_at IS NULL AND s.expires_at > now() + AND u.suspended_at IS NULL AND s.account_session_epoch = u.session_epoch + FOR SHARE OF s, u`, + [session, user] + ); + if (!row.rows.length) throw new CreateError(403, "identity_not_current"); +} + +/** The account is still active; share-locked until the transaction ends. */ +export async function currentAccount(client: DatabaseConnection, user: string): Promise { + const row = await client.query("SELECT 1 FROM users WHERE id = $1 AND suspended_at IS NULL FOR SHARE", [user]); + if (!row.rows.length) throw new CreateError(403, "identity_not_current"); +} + +export async function inTransaction(db: DatabasePool, run: (client: DatabaseConnection) => Promise): Promise { + const client = await db.connect(); + try { + await client.query("BEGIN"); + // Bounded: no request waits on another's locks for long. Network calls never run + // inside these transactions. + await client.query("SET LOCAL lock_timeout = '5s'"); + // Every statement is bounded too: no history scan holds share locks for long. + await client.query("SET LOCAL statement_timeout = '5s'"); + const result = await run(client); + await client.query("COMMIT"); + return result; + } catch (error) { + await client.query("ROLLBACK").catch(() => undefined); + throw error; + } finally { + client.release(); + } +} + +export const enrolOp = (device: string, account: string, d: { kind: "desktop" | "cli" | "hosted" | "escrow"; sign_pk: Buffer; kem_pk: Buffer; noise_pk: Buffer }): PolicyOp => ({ + op: "device-enrol", device, account, kind: d.kind, signPublicKey: d.sign_pk, kemPublicKey: d.kem_pk, noisePublicKey: d.noise_pk +}); + +/** The outbox row that enrols `device` in `collection`, if any (any keys, any account). */ +export const ENROLMENT = `SELECT o.ops, b.seq, b.item, b.state FROM next_policy_outbox o + LEFT JOIN next_policy_batches b ON b.id = o.batch_id + WHERE o.collection_id = $1 AND o.ops->'ops' @> $2::jsonb ORDER BY o.id LIMIT 1`; +export const enrolmentKey = (device: string) => JSON.stringify([{ op: "device-enrol", device }]); +/** The whole immutable enrolment tuple: device, account, kind and all three keys. */ +export const exactEnrolment = (device: string, account: string, d: Device) => JSON.stringify([{ + op: "device-enrol", device, account, kind: d.kind, + signPublicKey: { $hex: d.sign_pk.toString("hex") }, + kemPublicKey: { $hex: d.kem_pk.toString("hex") }, + noisePublicKey: { $hex: d.noise_pk.toString("hex") } +}]); + +/** A historical enrolment is never current once the device has been revoked. */ +export async function refuseRevoked(client: DatabaseConnection, collection: string, device: string): Promise { + const revoked = await client.query( + "SELECT 1 FROM next_policy_outbox WHERE collection_id = $1 AND ops->'ops' @> $2::jsonb LIMIT 1", + [collection, JSON.stringify([{ op: "device-revoke", device }])] + ); + if (revoked.rows.length) throw new CreateError(409, "device_revoked"); +} + +export function refuse(reply: FastifyReply, error: unknown, message: string) { + if (error instanceof CreateError) { + const status = error.status === 502 ? 503 : error.status; + return reply.code(status).send(apiError(error.code, message)); + } + if (isLockTimeout(error)) return reply.code(503).send(apiError("busy", message)); + throw error; +} + +/** + * The account is a current member of the collection: its latest effective membership + * op (outbox order, then op order within the batch) is a member-set acknowledged by + * the log. A member-remove is effective even while pending; a pending member-set is + * not. One row at most, projected in SQL. + */ +export async function currentMember(client: DatabaseConnection, collection: string, account: string): Promise { + const latest = await client.query<{ op: string }>( + `SELECT e.value->>'op' AS op + FROM next_policy_outbox o + LEFT JOIN next_policy_batches b ON b.id = o.batch_id + CROSS JOIN LATERAL jsonb_array_elements(o.ops->'ops') WITH ORDINALITY AS e(value, ord) + WHERE o.collection_id = $1 + AND (o.ops->'ops' @> $2::jsonb OR o.ops->'ops' @> $3::jsonb) + AND e.value->>'account' = $4 + AND (e.value->>'op' = 'member-remove' OR (e.value->>'op' = 'member-set' AND b.state = 'appended')) + ORDER BY o.id DESC, e.ord DESC + LIMIT 1`, + [collection, JSON.stringify([{ op: "member-set", account }]), JSON.stringify([{ op: "member-remove", account }]), account] + ); + if (latest.rows[0]?.op !== "member-set") throw new CreateError(409, "not_member"); +} diff --git a/services/server/src/features/next/cloud-copy-bootstrap.ts b/services/server/src/features/next/cloud-copy-bootstrap.ts index f039cb0a..364aac4a 100644 --- a/services/server/src/features/next/cloud-copy-bootstrap.ts +++ b/services/server/src/features/next/cloud-copy-bootstrap.ts @@ -16,30 +16,28 @@ // The control plane never holds a collection key. Each deployment generates its own // service device; the first committed record wins, and nothing is ever deleted on a // failure path. -import { verify } from "node:crypto"; import type { FastifyInstance, FastifyReply } from "fastify"; import type { DatabaseConnection, DatabasePool } from "../../database-types.js"; import { apiError } from "../../platform/http-errors.js"; import { requireConnector, requireSessionContext, requireUser } from "../../platform/request-authentication.js"; import { LOG_TOKEN_LIFETIME_MS, type LogServiceClient } from "./log-service-client.js"; -import { ed25519PublicKeyObject, type NextControlPlaneConfig } from "./policy-keys.js"; +import { + authenticate, CreateError, currentAccount, currentIdentity, currentSession, ENROLMENT, enrolmentKey, enrolOp, exactEnrolment, + inTransaction, lock, NIL, refuse as refuseCommon, refuseRevoked, SERVICE_ACCOUNT, type Device, type Proof +} from "./bootstrap-common.js"; +import { type NextControlPlaneConfig } from "./policy-keys.js"; import { queueNextPolicy, registerNextCollection, type PolicyEmitter } from "./policy-outbox.js"; -import { domainHash, encodeCbor, uuidBytes, type PolicyOp } from "./policy-wire.js"; +import { domainHash, encodeCbor, uuidBytes } from "./policy-wire.js"; import { generateServiceDevice, loadServiceDevice, ServiceDeviceError, storeServiceDevice, type ServiceDeviceRecord } from "./service-devices.js"; -/** Service devices belong to no account (policy.md: hosted and escrow enrol with the zero account). */ -const SERVICE_ACCOUNT = "00000000-0000-0000-0000-000000000000"; -const NIL = SERVICE_ACCOUNT; const KINDS = ["hosted", "escrow"] as const; -interface Proof { device_id: string; challenge: string; sig: string } -interface Device { sign_pk: Buffer; kem_pk: Buffer; noise_pk: Buffer; kind: "desktop" | "cli" } -type Connector = { id: string; user_id: string }; -/** PostgreSQL lock_timeout: another request holds the rows; answer busy, never the driver error. */ -const isLockTimeout = (error: unknown) => (error as { code?: unknown } | null)?.code === "55P03"; - -class CreateError extends Error { - constructor(readonly status: number, readonly code: string) { super(code); } +function refuse(reply: FastifyReply, error: unknown, message: string) { + if (error instanceof ServiceDeviceError) { + const status = error.status === 502 ? 503 : error.status; + return reply.code(status).send(apiError(error.code, message)); + } + return refuseCommon(reply, error, message); } /** `H("mdbase/v1/cloud-copy-create", cbor[challenge, connector, device, collection])`, signed by the owner's device. */ @@ -52,9 +50,6 @@ export function cloudCopyJoinDigest(input: { challenge: Uint8Array; connector: s return domainHash("mdbase/v1/cloud-copy-join", encodeCbor([input.challenge, uuidBytes(input.connector), uuidBytes(input.device), uuidBytes(input.collection)])); } -const lock = (client: DatabaseConnection, collection: string) => - client.query("SELECT pg_advisory_xact_lock(hashtextextended($1::uuid::text, 20261005))", [collection]); - /** Whether the collection already exists: false when free, true when this owner's cloud copy, else refused. */ async function existing(client: DatabaseConnection, collection: string, owner: string): Promise { const row = (await client.query<{ owner_user_id: string; sync: string; left: boolean }>( @@ -78,119 +73,11 @@ async function currentCloudCopy(client: DatabaseConnection, collection: string, if (!current.rows.length) throw new CreateError(409, "not_current_cloud_copy"); } -/** Verify a device's signature over `digest` and consume its challenge. */ -async function authenticate(client: DatabaseConnection, body: Proof, connector: Connector, digest: (challenge: Uint8Array) => Uint8Array): Promise { - const device = (await client.query( - "SELECT sign_pk, kem_pk, noise_pk, kind FROM next_devices WHERE id = $1 AND connector_id = $2 AND user_id = $3", - [body.device_id, connector.id, connector.user_id] - )).rows[0]; - const challenge = Buffer.from(body.challenge, "hex"); - if (!device || !verify(null, digest(challenge), ed25519PublicKeyObject(device.sign_pk), Buffer.from(body.sig, "hex"))) throw new CreateError(403, "invalid_proof"); - const used = await client.query( - "UPDATE next_device_challenges SET used_at = now() WHERE challenge = $1 AND connector_id = $2 AND used_at IS NULL AND expires_at > now()", - [challenge, connector.id] - ); - if (used.rowCount !== 1) throw new CreateError(403, "invalid_proof"); - return device; -} - -/** - * The connector, account and device are still current, with the exact keys - * authenticated in phase 1; locked until the transaction ends, so a revocation, - * suspension or device removal either happened before (and is refused here) or waits. - */ -async function currentIdentity(client: DatabaseConnection, connector: Connector, deviceId: string, device: Device): Promise { - const row = await client.query( - `SELECT 1 FROM connectors c JOIN users u ON u.id = c.user_id - JOIN next_devices d ON d.connector_id = c.id AND d.user_id = u.id - WHERE c.id = $1 AND u.id = $2 AND d.id = $3 AND c.revoked_at IS NULL AND u.suspended_at IS NULL - AND d.sign_pk = $4 AND d.kem_pk = $5 AND d.noise_pk = $6 AND d.kind = $7 - FOR SHARE OF c, u, d`, - [connector.id, connector.user_id, deviceId, device.sign_pk, device.kem_pk, device.noise_pk, device.kind] - ); - if (!row.rows.length) throw new CreateError(403, "identity_not_current"); -} - -/** - * The signed-in session is still the account's current credential (not revoked, not - * expired, same session epoch, account not suspended); share-locked until the - * transaction ends, so a sign-out or suspension either happened before or waits. - */ -async function currentSession(client: DatabaseConnection, session: string, user: string): Promise { - const row = await client.query( - `SELECT 1 FROM sessions s JOIN users u ON u.id = s.user_id - WHERE s.id = $1 AND u.id = $2 AND s.revoked_at IS NULL AND s.expires_at > now() - AND u.suspended_at IS NULL AND s.account_session_epoch = u.session_epoch - FOR SHARE OF s, u`, - [session, user] - ); - if (!row.rows.length) throw new CreateError(403, "identity_not_current"); -} - -/** The account is still active; share-locked until the transaction ends. */ -async function currentAccount(client: DatabaseConnection, user: string): Promise { - const row = await client.query("SELECT 1 FROM users WHERE id = $1 AND suspended_at IS NULL FOR SHARE", [user]); - if (!row.rows.length) throw new CreateError(403, "identity_not_current"); -} - -async function inTransaction(db: DatabasePool, run: (client: DatabaseConnection) => Promise): Promise { - const client = await db.connect(); - try { - await client.query("BEGIN"); - // Bounded: no request waits on another's locks for long. Network calls never run - // inside these transactions. - await client.query("SET LOCAL lock_timeout = '5s'"); - const result = await run(client); - await client.query("COMMIT"); - return result; - } catch (error) { - await client.query("ROLLBACK").catch(() => undefined); - throw error; - } finally { - client.release(); - } -} - const publicRecord = (record: ServiceDeviceRecord) => ({ kind: record.kind, device_id: record.device_id, sign_pk: record.sign_pk.toString("hex"), kem_pk: record.kem_pk.toString("hex"), noise_pk: record.noise_pk.toString("hex") }); -const enrolOp = (device: string, account: string, d: { kind: "desktop" | "cli" | "hosted" | "escrow"; sign_pk: Buffer; kem_pk: Buffer; noise_pk: Buffer }): PolicyOp => ({ - op: "device-enrol", device, account, kind: d.kind, signPublicKey: d.sign_pk, kemPublicKey: d.kem_pk, noisePublicKey: d.noise_pk -}); - -/** The outbox row that enrols `device` in `collection`, if any (any keys, any account). */ -const ENROLMENT = `SELECT o.ops, b.seq, b.item, b.state FROM next_policy_outbox o - LEFT JOIN next_policy_batches b ON b.id = o.batch_id - WHERE o.collection_id = $1 AND o.ops->'ops' @> $2::jsonb ORDER BY o.id LIMIT 1`; -const enrolmentKey = (device: string) => JSON.stringify([{ op: "device-enrol", device }]); -/** The whole immutable enrolment tuple: device, account, kind and all three keys. */ -const exactEnrolment = (device: string, account: string, d: Device) => JSON.stringify([{ - op: "device-enrol", device, account, kind: d.kind, - signPublicKey: { $hex: d.sign_pk.toString("hex") }, - kemPublicKey: { $hex: d.kem_pk.toString("hex") }, - noisePublicKey: { $hex: d.noise_pk.toString("hex") } -}]); - -/** A historical enrolment is never current once the device has been revoked. */ -async function refuseRevoked(client: DatabaseConnection, collection: string, device: string): Promise { - const revoked = await client.query( - "SELECT 1 FROM next_policy_outbox WHERE collection_id = $1 AND ops->'ops' @> $2::jsonb LIMIT 1", - [collection, JSON.stringify([{ op: "device-revoke", device }])] - ); - if (revoked.rows.length) throw new CreateError(409, "device_revoked"); -} - -function refuse(reply: FastifyReply, error: unknown, message: string) { - if (error instanceof CreateError || error instanceof ServiceDeviceError) { - const status = error.status === 502 ? 503 : error.status; - return reply.code(status).send(apiError(error.code, message)); - } - if (isLockTimeout(error)) return reply.code(503).send(apiError("busy", message)); - throw error; -} - export function registerCloudCopyRoutes(app: FastifyInstance, options: { db: DatabasePool; next: NextControlPlaneConfig; emitter: PolicyEmitter; log: Pick; fetchImpl?: typeof fetch; now?: () => number; diff --git a/services/server/src/features/next/policy-keys.ts b/services/server/src/features/next/policy-keys.ts index af503f11..b4e0d673 100644 --- a/services/server/src/features/next/policy-keys.ts +++ b/services/server/src/features/next/policy-keys.ts @@ -39,6 +39,8 @@ export interface NextControlPlaneConfig { * inbound `serviceTokens`. */ cloudCopyBootstrap?: { hosted: { url: string; token: string }; escrow: { url: string; token: string } }; + /** Private collection create and device enrol routes; set only by MDBASE_NEXT_PRIVATE_BOOTSTRAP=1. */ + privateBootstrap?: true; labFixtures?: LabFixtureConfig; } @@ -135,6 +137,8 @@ export function parseNextControlPlaneEnv(env: NodeJS.ProcessEnv): NextControlPla throw new Error("Inbound and outbound service tokens must all differ."); } } + const privateBootstrap = env.MDBASE_NEXT_PRIVATE_BOOTSTRAP?.trim() ?? ""; + if (privateBootstrap !== "" && privateBootstrap !== "0" && privateBootstrap !== "1") throw new Error("MDBASE_NEXT_PRIVATE_BOOTSTRAP must be 0 or 1."); return { rootPublicKey: hexBytes(root, 32, "MDBASE_NEXT_ROOT_PUBLIC_KEY"), policyPrivateKeyPem: pem, @@ -142,6 +146,7 @@ export function parseNextControlPlaneEnv(env: NodeJS.ProcessEnv): NextControlPla logService: { url: logServiceUrl, tokenIssuerKeyPem, transportKeyPem }, serviceTokens: { ...(hosted ? { hosted } : {}), ...(escrow ? { escrow } : {}) }, ...(cloudCopyBootstrap ? { cloudCopyBootstrap } : {}), + ...(privateBootstrap === "1" ? { privateBootstrap: true as const } : {}), ...(labFixtures ? { labFixtures } : {}), }; } diff --git a/services/server/src/features/next/private-collections.postgres.test.ts b/services/server/src/features/next/private-collections.postgres.test.ts new file mode 100644 index 00000000..41865e02 --- /dev/null +++ b/services/server/src/features/next/private-collections.postgres.test.ts @@ -0,0 +1,333 @@ +import { generateKeyPairSync, randomUUID, sign } from "node:crypto"; +import cookie from "@fastify/cookie"; +import Fastify from "fastify"; +import pg from "pg"; +import { afterAll, beforeAll, describe, expect, it } from "vitest"; +import { createDatabase, type DatabasePool } from "../../db.js"; +import { tokenHash } from "../../security.js"; +import { deviceRegistrationDigest, issueDeviceChallenge, registerDevice } from "./devices.js"; +import { LogServiceClient } from "./log-service-client.js"; +import { certToJson, ed25519RawPublicKey, loadPolicySigner, parseNextControlPlaneEnv, type NextControlPlaneConfig } from "./policy-keys.js"; +import { PolicyEmitter, queueNextPolicy, registerNextCollection } from "./policy-outbox.js"; +import { certDigest, chainHash, decodeCbor, encodeCbor, keyId, type Cbor, type Decoded, type PolicyOp } from "./policy-wire.js"; +import { inTransaction, refuse } from "./bootstrap-common.js"; +import { privateCreateDigest, privateDeviceEnrolDigest, registerPrivateCollectionRoutes } from "./private-collections.js"; + +const testUrl = process.env.MDBASE_CONNECT_TEST_DATABASE_URL; +const approved = process.env.MDBASE_CONNECT_DESTRUCTIVE_TEST_APPROVAL === "I APPROVE MDBASE CONNECT DESTRUCTIVE POSTGRES TESTS"; +const describePg = testUrl && approved ? describe : describe.skip; +const field = (value: Decoded, key: number) => value instanceof Map ? value.get(key) : undefined; +const hex = (bytes: Uint8Array) => Buffer.from(bytes).toString("hex"); +const rawX = () => (generateKeyPairSync("x25519").publicKey.export({ format: "der", type: "spki" }) as Buffer).subarray(-32); + +describe("private bootstrap configuration", () => { + const base = { + MDBASE_NEXT_CONTROL_PLANE: "1", MDBASE_NEXT_ROOT_PUBLIC_KEY: "00".repeat(32), MDBASE_NEXT_POLICY_SIGNING_KEY: "pem", MDBASE_NEXT_POLICY_KEY_CERT: "{}", + MDBASE_NEXT_LOG_SERVICE_URL: "https://log.example", MDBASE_NEXT_LOG_TOKEN_SIGNING_KEY: "pem", MDBASE_NEXT_LOG_TRANSPORT_KEY: "pem" + }; + it("is off by default, on with 1, and refuses anything else", () => { + expect(parseNextControlPlaneEnv(base)?.privateBootstrap).toBeUndefined(); + expect(parseNextControlPlaneEnv({ ...base, MDBASE_NEXT_PRIVATE_BOOTSTRAP: "0" })?.privateBootstrap).toBeUndefined(); + expect(parseNextControlPlaneEnv({ ...base, MDBASE_NEXT_PRIVATE_BOOTSTRAP: "1" })?.privateBootstrap).toBe(true); + expect(() => parseNextControlPlaneEnv({ ...base, MDBASE_NEXT_PRIVATE_BOOTSTRAP: "yes" })).toThrow(/0 or 1/); + }); +}); + +function configuration(): NextControlPlaneConfig { + const root = generateKeyPairSync("ed25519").privateKey; + const policy = generateKeyPairSync("ed25519").privateKey; + const cert = { policyPublicKey: ed25519RawPublicKey(policy), notBefore: Date.now() - 60_000, notAfter: Date.now() + 30 * 86_400_000, root: keyId(ed25519RawPublicKey(root)) }; + const pem = (key: typeof root) => key.export({ type: "pkcs8", format: "pem" }).toString(); + return { + rootPublicKey: ed25519RawPublicKey(root), policyPrivateKeyPem: pem(policy), policyCert: certToJson({ ...cert, signature: sign(null, certDigest(cert), root) }), + serviceTokens: {}, privateBootstrap: true, + logService: { url: "http://log.test", tokenIssuerKeyPem: pem(generateKeyPairSync("ed25519").privateKey), transportKeyPem: pem(generateKeyPairSync("ed25519").privateKey) } + }; +} + +/** The log service: control items per collection, with create, append, head and read. */ +class Log { + readonly logs = new Map(); + /** Runs once, while the route awaits a read-back. */ + onRead: (() => Promise) | undefined; + readonly fetch: typeof fetch = async (input, init) => { + if (String(input).endsWith("/v1/nonce")) return new Response("ab".repeat(32)); + const frame = decodeCbor(Buffer.from(init!.body as Uint8Array)); + const method = field(frame, 2); + const params = field(frame, 3)!; + const id = hex(field(params, 0) as Uint8Array); + const items = this.logs.get(id); + const chain = () => chainHash(items![items!.length - 1]!); + let result: Cbor; + if (method === "create_log") { + if (!items) this.logs.set(id, [Buffer.from(field(params, 1) as Uint8Array)]); + result = { struct: [[0, 1], [1, chainHash(this.logs.get(id)![0]!)]] }; + } else if (method === "head") { + result = { struct: [[0, items!.length], [1, chain()], [2, 1]] }; + } else if (method === "append") { + const expect = field(params, 1) as number; + const prev = Buffer.from(field(params, 2) as Uint8Array); + if (expect !== items!.length + 1 || !prev.equals(Buffer.from(chain()))) { + result = { struct: [[0, 1], [1, items!.length], [2, chain()]] }; + } else { + const added = (field(params, 3) as Uint8Array[]).map((b) => Buffer.from(b)); + items!.push(...added); + result = { struct: [[0, 0], [1, expect], [2, items!.length]] }; + } + } else if (method === "read") { + const hook = this.onRead; + this.onRead = undefined; + await hook?.(); + if (!items) return new Response(encodeCbor({ struct: [[0, 1], [1, 1], [3, { struct: [[0, "not_found"]] }]] })); + const after = field(params, 1) as number; + const limit = field(params, 2) as number; + result = { struct: [[0, items.slice(after, after + limit).map((item, i) => [after + i + 1, item])]] }; + } else throw new Error("unexpected log operation"); + return new Response(encodeCbor({ struct: [[0, 1], [1, 1], [2, result]] }), { headers: { "content-type": "application/vnd.mdbase.v1+cbor" } }); + }; +} + +describePg("private collections", () => { + let db: DatabasePool; + let admin: pg.Pool; + let schema: string; + const config = configuration(); + const log = new Log(); + const client = new LogServiceClient(config.logService, log.fetch); + const app = Fastify(); + + beforeAll(async () => { + const url = new URL(testUrl!); + if (!["localhost", "127.0.0.1", "::1"].includes(url.hostname) || !/test/i.test(url.pathname)) throw new Error("Private collection tests require dedicated local test Postgres."); + schema = `private_${randomUUID().replaceAll("-", "")}`; + admin = new pg.Pool({ connectionString: url.toString(), max: 2 }); + await admin.query(`CREATE SCHEMA "${schema}"`); + url.searchParams.set("options", `-csearch_path=${schema}`); + db = await createDatabase(url.toString()); + const emitter = new PolicyEmitter(db, client, loadPolicySigner(config, Date.now())); + await app.register(cookie); + registerPrivateCollectionRoutes(app, { db, next: config, emitter, log: client }); + }, 60_000); + afterAll(async () => { + await app.close(); await db?.end(); + if (admin && schema) await admin.query(`DROP SCHEMA "${schema}" CASCADE`); + await admin?.end(); + }); + + async function identity(user = randomUUID()) { + const connector = { id: randomUUID(), user_id: user }; + const token = randomUUID(); + const device = randomUUID(); + const key = generateKeyPairSync("ed25519").privateKey; + const signPk = ed25519RawPublicKey(key); + const kemPk = rawX(); const noisePk = rawX(); + await db.query("INSERT INTO users(id,email,name) VALUES($1,$2,'Owner') ON CONFLICT DO NOTHING", [user, `${user}@example.test`]); + await db.query("INSERT INTO connectors(id,user_id,name,token_hash) VALUES($1,$2,'Daemon',$3)", [connector.id, user, tokenHash(token)]); + const registration = await issueDeviceChallenge(db, connector.id); + await registerDevice(db, connector, { + device_id: device, kind: "desktop", sign_pk: hex(signPk), kem_pk: hex(kemPk), noise_pk: hex(noisePk), challenge: registration.challenge, + sig: hex(sign(null, deviceRegistrationDigest({ challenge: Buffer.from(registration.challenge, "hex"), connectorId: connector.id, deviceId: device, signPk, kemPk, noisePk }), key)) + }); + return { connector, device, key, signPk, headers: { authorization: `Bearer ${token}` } }; + } + type Who = Awaited>; + async function createProof(who: Who, collection: string) { + const { challenge } = await issueDeviceChallenge(db, who.connector.id); + const digest = privateCreateDigest({ challenge: Buffer.from(challenge, "hex"), connector: who.connector.id, device: who.device, collection }); + return { collection_id: collection, device_id: who.device, challenge, sig: hex(sign(null, digest, who.key)) }; + } + async function enrolProof(who: Who, collection: string, commit: string) { + const { challenge } = await issueDeviceChallenge(db, who.connector.id); + const digest = privateDeviceEnrolDigest({ challenge: Buffer.from(challenge, "hex"), connector: who.connector.id, device: who.device, collection, sasCommit: Buffer.from(commit, "hex") }); + return { device_id: who.device, challenge, sig: hex(sign(null, digest, who.key)), sas_commit: commit }; + } + const create = (who: Who, payload: unknown) => app.inject({ method: "POST", url: "/v1/next/collections/private", headers: who.headers, payload }); + const enrol = (who: Who, collection: string, payload: unknown) => + app.inject({ method: "POST", url: `/v1/next/collections/${collection}/private/devices`, headers: who.headers, payload }); + const commit = () => hex(Buffer.from(randomUUID().replaceAll("-", "").repeat(2), "hex")); + const registered = async (collection: string) => (await db.query("SELECT 1 FROM next_collections WHERE collection_id = $1", [collection])).rows.length === 1; + const drain = (collection: string) => new PolicyEmitter(db, client, loadPolicySigner(config, Date.now())).drainCollection(collection); + /** Queue ops without appending them (pending). */ + const queue = async (collection: string, ops: PolicyOp[]) => { + const c = await db.connect(); + try { + await c.query("BEGIN"); + await queueNextPolicy(c, collection, ops); + await c.query("COMMIT"); + } finally { + c.release(); + } + }; + const opsOf = (item: Buffer | Uint8Array) => field(decodeCbor(field(decodeCbor(Buffer.from(item)), 11) as Uint8Array), 3) as Decoded[]; + const created = async (who: Who) => { + const collection = randomUUID(); + const response = await create(who, await createProof(who, collection)); + expect(response.statusCode, response.body).toBe(200); + return collection; + }; + + it("creates an e2e genesis that enrols only the owner's device", async () => { + const who = await identity(); const collection = randomUUID(); + const response = await create(who, await createProof(who, collection)); + expect(response.statusCode, response.body).toBe(200); + expect(response.headers["cache-control"]).toBe("no-store"); + const result = response.json(); + expect(result).toMatchObject({ collection_id: collection, state: "private", owner_account: who.connector.user_id, head: { seq: 1 }, rekey_recipients: [who.device] }); + expect(result.service_devices).toBeUndefined(); + const ops = opsOf(Buffer.from(result.genesis.item, "hex")); + expect(ops.map((op) => field(op, 0))).toEqual([1, 4, 2]); + expect(field(ops[0]!, 3)).toBe(0); // e2e + expect(hex(field(ops[2]!, 1) as Uint8Array)).toBe(who.device.replaceAll("-", "")); + expect(field(ops[2]!, 3)).toBe(0); // desktop + expect(field(ops[2]!, 7)).toBeUndefined(); // the creator needs no approval + expect(result.genesis.item).toBe(hex(log.logs.get(collection.replaceAll("-", ""))![0]!)); + const claims = decodeCbor(Buffer.from(result.device.token.split(".")[0], "hex")); + expect(hex(field(claims, 1) as Uint8Array)).toBe(who.device.replaceAll("-", "")); + expect(hex(field(claims, 5) as Uint8Array)).toBe(collection.replaceAll("-", "")); + const row = (await db.query("SELECT sync, runtime FROM next_collections WHERE collection_id = $1", [collection])).rows[0]; + expect(row).toEqual({ sync: "private", runtime: "next" }); + }); + + it("answers a retry from the creating device with the same genesis", async () => { + const who = await identity(); const collection = randomUUID(); + const first = (await create(who, await createProof(who, collection))).json(); + const again = await create(who, await createProof(who, collection)); + expect(again.statusCode, again.body).toBe(200); + expect(again.json().genesis).toEqual(first.genesis); + expect(log.logs.get(collection.replaceAll("-", ""))!.length).toBe(1); + }); + + it("bounds every statement and lock wait in its transactions, and answers busy on either timeout", async () => { + const settings = await inTransaction(db, async (c) => (await c.query<{ s: string; l: string }>( + "SELECT current_setting('statement_timeout') AS s, current_setting('lock_timeout') AS l" + )).rows[0]); + expect(settings).toEqual({ s: "5s", l: "5s" }); + const timedOut = await inTransaction(db, (c) => c.query("SET LOCAL statement_timeout = '10ms'").then(() => c.query("SELECT pg_sleep(1)"))) + .then(() => undefined, (error: unknown) => error); + expect((timedOut as { code?: string }).code).toBe("57014"); + for (const code of ["55P03", "57014"]) { + const reply = Fastify(); + reply.get("/", (_req, r) => refuse(r, Object.assign(new Error("timeout"), { code }), "retry")); + const res = await reply.inject({ method: "GET", url: "/" }); + expect([res.statusCode, res.json().error.code]).toEqual([503, "busy"]); + await reply.close(); + } + }); + + it("a retry mints nothing once the owner account is removed, even while the removal is pending", async () => { + const who = await identity(); + const collection = await created(who); + await queue(collection, [{ op: "member-remove", account: who.connector.user_id }]); + expect((await create(who, await createProof(who, collection))).json().error.code).toBe("not_member"); + }); + + it("refuses other devices, other owners, cloud copies and collections that left sync", async () => { + const who = await identity(); + const collection = await created(who); + const sibling = await identity(who.connector.user_id); + expect((await create(sibling, await createProof(sibling, collection))).json().error.code).toBe("collection_exists"); + const stranger = await identity(); + expect((await create(stranger, await createProof(stranger, collection))).json().error.code).toBe("collection_exists"); + const cloud = randomUUID(); + await db.query("INSERT INTO next_collections(collection_id, owner_user_id, runtime, sync, root_key_id) VALUES($1,$2,'next','cloud_copy',$3)", [cloud, who.connector.user_id, Buffer.from(config.policyCert.root_key_id, "hex")]); + expect((await create(who, await createProof(who, cloud))).json().error.code).toBe("collection_exists"); + await db.query("UPDATE next_collections SET left_sync_at = now() WHERE collection_id = $1", [collection]); + expect((await create(who, await createProof(who, collection))).statusCode).toBe(409); + }); + + it("needs a fresh proof under its own domain, and refuses nil identifiers", async () => { + const who = await identity(); const collection = randomUUID(); + const payload = await createProof(who, collection); + expect((await app.inject({ method: "POST", url: "/v1/next/collections/private", payload })).statusCode).toBe(401); + expect((await create(who, { ...payload, collection_id: randomUUID() })).statusCode).toBe(403); + expect((await create(who, payload)).statusCode).toBe(200); + expect((await create(who, payload)).statusCode).toBe(403); + // An enrolment proof never creates. + const other = randomUUID(); + const e = await enrolProof(who, other, commit()); + expect((await create(who, { collection_id: other, device_id: e.device_id, challenge: e.challenge, sig: e.sig })).statusCode).toBe(403); + expect(await registered(other)).toBe(false); + const nil = "00000000-0000-0000-0000-000000000000"; + expect((await create(who, await createProof(who, nil))).statusCode).toBe(400); + }); + + it("registers nothing when the device is removed before the proof is checked", async () => { + const who = await identity(); const collection = randomUUID(); + const payload = await createProof(who, collection); + await db.query("DELETE FROM next_devices WHERE id = $1", [who.device]); + expect((await create(who, payload)).statusCode).toBe(403); + expect(await registered(collection)).toBe(false); + }); + + it("enrols the owner's second device with its SAS commitment, approval pending", async () => { + const owner = await identity(); + const collection = await created(owner); + const second = await identity(owner.connector.user_id); + const sas = commit(); + const response = await enrol(second, collection, await enrolProof(second, collection, sas)); + expect(response.statusCode, response.body).toBe(200); + expect(response.json()).toMatchObject({ collection_id: collection, enrolled_at: 2, approval: "pending", device: { device_id: second.device } }); + const items = log.logs.get(collection.replaceAll("-", ""))!; + const [op] = opsOf(items[1]!); + expect(field(op!, 0)).toBe(2); + expect(hex(field(op!, 1) as Uint8Array)).toBe(second.device.replaceAll("-", "")); + expect(hex(field(op!, 7) as Uint8Array)).toBe(sas); + // The same device and commitment again: the same enrolment, nothing appended. + const again = await enrol(second, collection, await enrolProof(second, collection, sas)); + expect(again.statusCode, again.body).toBe(200); + expect(again.json().enrolled_at).toBe(2); + expect(items.length).toBe(2); + // A different commitment is the device's own approval-request, never this route. + expect((await enrol(second, collection, await enrolProof(second, collection, commit()))).json().error.code).toBe("device_enrolled_differently"); + }); + + it("enrols another member account's device, and refuses non-members and removed members", async () => { + const owner = await identity(); + const collection = await created(owner); + const editor = await identity(); + expect((await enrol(editor, collection, await enrolProof(editor, collection, commit()))).json().error.code).toBe("not_member"); + await queue(collection, [{ op: "member-set", account: editor.connector.user_id, role: "editor" }]); + expect((await enrol(editor, collection, await enrolProof(editor, collection, commit()))).json().error.code).toBe("not_member"); + await drain(collection); + expect((await enrol(editor, collection, await enrolProof(editor, collection, commit()))).statusCode).toBe(200); + const removed = await identity(editor.connector.user_id); + await queue(collection, [{ op: "member-remove", account: editor.connector.user_id }]); + expect((await enrol(removed, collection, await enrolProof(removed, collection, commit()))).json().error.code).toBe("not_member"); + }); + + it("refuses revoked devices, cloud copies, collections that left sync and other proofs", async () => { + const owner = await identity(); + const collection = await created(owner); + const second = await identity(owner.connector.user_id); + await queue(collection, [{ op: "device-revoke", device: second.device }]); + expect((await enrol(second, collection, await enrolProof(second, collection, commit()))).json().error.code).toBe("device_revoked"); + + const cloud = randomUUID(); + const c = await db.connect(); + try { + await c.query("BEGIN"); + await registerNextCollection(c, { + collectionId: cloud, ownerUserId: owner.connector.user_id, runtime: "next", sync: "cloud_copy", rootKeyId: Buffer.from(config.policyCert.root_key_id, "hex"), + ops: [{ op: "genesis", owner: owner.connector.user_id, root: Buffer.from(config.policyCert.root_key_id, "hex"), state: "cloud-copy" }, { op: "member-set", account: owner.connector.user_id, role: "owner" }] + }); + await c.query("COMMIT"); + } finally { + c.release(); + } + const third = await identity(owner.connector.user_id); + expect((await enrol(third, cloud, await enrolProof(third, cloud, commit()))).json().error.code).toBe("not_current_private"); + + const create2 = await created(owner); + const p = await createProof(third, create2); + expect((await enrol(third, create2, { device_id: p.device_id, challenge: p.challenge, sig: p.sig, sas_commit: commit() })).json().error.code).toBe("invalid_proof"); + const sas = commit(); + const proofed = await enrolProof(third, create2, sas); + expect((await enrol(third, create2, { ...proofed, sas_commit: commit() })).json().error.code).toBe("invalid_proof"); + await db.query("UPDATE next_collections SET runtime = 'shadow' WHERE collection_id = $1", [create2]); + expect((await enrol(third, create2, await enrolProof(third, create2, sas))).json().error.code).toBe("not_current_private"); + expect((await create(owner, await createProof(owner, create2))).json().error.code).toBe("collection_exists"); + await db.query("UPDATE next_collections SET runtime = 'next' WHERE collection_id = $1", [create2]); + await db.query("UPDATE next_collections SET left_sync_at = now() WHERE collection_id = $1", [create2]); + expect((await enrol(third, create2, await enrolProof(third, create2, sas))).json().error.code).toBe("not_current_private"); + }); +}); diff --git a/services/server/src/features/next/private-collections.ts b/services/server/src/features/next/private-collections.ts new file mode 100644 index 00000000..00a2eaef --- /dev/null +++ b/services/server/src/features/next/private-collections.ts @@ -0,0 +1,239 @@ +// Private (end-to-end) collections on the next control plane. Mounted only with +// MDBASE_NEXT_PRIVATE_BOOTSTRAP=1. +// +// - Create (`POST /v1/next/collections/private`): an owner's registered device creates +// it. Genesis is e2e, sets the owner and enrols that device only. No service device +// is ever enrolled. The device's own initial rekey keys the collection. +// - Device enrol (`POST /v1/next/collections/:id/private/devices`): a registered device +// of a current member account is enrolled, carrying its SAS commitment +// (`device-enrol` key 7). It holds no key until an existing keyed device approves it +// (SAS commit-then-reveal) and appends a `key_grant`. A changed commitment travels +// as the device's own `approval-request`, never through here. +// +// The control plane never approves, never grants and never holds a collection key. It +// enrols with its policy key and mints log credentials, nothing more. +import type { FastifyInstance } from "fastify"; +import type { DatabaseConnection, DatabasePool } from "../../database-types.js"; +import { apiError } from "../../platform/http-errors.js"; +import { requireConnector } from "../../platform/request-authentication.js"; +import { + authenticate, CreateError, currentIdentity, currentMember, ENROLMENT, enrolmentKey, enrolOp, exactEnrolment, inTransaction, lock, NIL, + refuse, refuseRevoked, type Device, type Proof +} from "./bootstrap-common.js"; +import { LOG_TOKEN_LIFETIME_MS, type LogServiceClient } from "./log-service-client.js"; +import type { NextControlPlaneConfig } from "./policy-keys.js"; +import { queueNextPolicy, registerNextCollection, type PolicyEmitter } from "./policy-outbox.js"; +import { domainHash, encodeCbor, uuidBytes } from "./policy-wire.js"; + +/** `H("mdbase/v1/private-create", cbor[challenge, connector, device, collection])`, signed by the owner's device. */ +export function privateCreateDigest(input: { challenge: Uint8Array; connector: string; device: string; collection: string }): Uint8Array { + return domainHash("mdbase/v1/private-create", encodeCbor([input.challenge, uuidBytes(input.connector), uuidBytes(input.device), uuidBytes(input.collection)])); +} + +/** `H("mdbase/v1/private-device-enrol", cbor[challenge, connector, device, collection, sas_commit])`, signed by the enrolling device. */ +export function privateDeviceEnrolDigest(input: { challenge: Uint8Array; connector: string; device: string; collection: string; sasCommit: Uint8Array }): Uint8Array { + return domainHash("mdbase/v1/private-device-enrol", encodeCbor([ + input.challenge, uuidBytes(input.connector), uuidBytes(input.device), uuidBytes(input.collection), input.sasCommit + ])); +} + +/** Whether the collection already exists: false when free, true when this owner's private collection, else refused. */ +async function existing(client: DatabaseConnection, collection: string, owner: string): Promise { + const row = (await client.query<{ owner_user_id: string; sync: string; runtime: string; left: boolean }>( + "SELECT owner_user_id, sync, runtime, left_sync_at IS NOT NULL AS left FROM next_collections WHERE collection_id = $1 FOR UPDATE", [collection] + )).rows[0]; + if (row && (row.owner_user_id !== owner || row.sync !== "private" || row.runtime !== "next" || row.left)) throw new CreateError(409, "collection_exists"); + if (!row) { + // A local collection with this logical ID that belongs to someone else is never adopted. + const other = await client.query("SELECT 1 FROM collections WHERE local_id = $1 AND user_id <> $2 AND removed_at IS NULL", [collection, owner]); + if (other.rows.length) throw new CreateError(409, "collection_exists"); + } + return Boolean(row); +} + +/** A current private collection on the next runtime (any owner); share-locked until the transaction ends. */ +async function currentPrivate(client: DatabaseConnection, collection: string, owner?: string): Promise { + const current = await client.query( + `SELECT 1 FROM next_collections WHERE collection_id = $1 AND sync = 'private' AND runtime = 'next' AND left_sync_at IS NULL + AND ($2::uuid IS NULL OR owner_user_id = $2) FOR SHARE`, + [collection, owner ?? null] + ); + if (!current.rows.length) throw new CreateError(409, "not_current_private"); +} + +/** The enrolment tuple including the SAS commitment. */ +const exactPrivateEnrolment = (device: string, account: string, d: Device, sasCommit: Buffer) => { + const [op] = JSON.parse(exactEnrolment(device, account, d)) as Array>; + return JSON.stringify([{ ...op, sasCommit: { $hex: sasCommit.toString("hex") } }]); +}; + +export function registerPrivateCollectionRoutes(app: FastifyInstance, options: { + db: DatabasePool; next: NextControlPlaneConfig; emitter: PolicyEmitter; + log: Pick; now?: () => number; +}): void { + if (!options.next.privateBootstrap) throw new Error("private collection routes need MDBASE_NEXT_PRIVATE_BOOTSTRAP=1"); + const rootKeyId = Buffer.from(options.next.policyCert.root_key_id, "hex"); + const uuid = { type: "string", pattern: "^[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}$" }; + const proof = { device_id: uuid, challenge: { type: "string", pattern: "^[0-9a-f]{64}$" }, sig: { type: "string", pattern: "^[0-9a-f]{128}$" } }; + const limited = { bodyLimit: 4096, config: { rateLimit: { max: 6, timeWindow: "1 minute" } } }; + + /** The exact bytes of a policy batch, as the log returns them at its position. */ + async function appendedBatch(collection: string, batch: { seq: string | number | null; item: Buffer | null; state: string | null } | undefined) { + const seq = batch?.seq === null || batch?.seq === undefined ? null : Number(batch.seq); + const external = seq !== null && batch?.state === "appended" ? await options.log.controlItemAt(collection, seq) : null; + if (seq === null || !batch?.item || !external || !batch.item.equals(Buffer.from(external))) throw new CreateError(503, "not_ready"); + return { seq, item: batch.item }; + } + + const mint = (device: string, signPk: Buffer, collection: string) => { + const expiresAt = (options.now ?? Date.now)() + LOG_TOKEN_LIFETIME_MS; + return { device_id: device, token: options.log.mintToken({ device, signPublicKey: signPk, collection, expiresAt }), expires_at: expiresAt }; + }; + + // ---- Create: a registered device of the owner. ---- + app.post<{ Body: Proof & { collection_id: string } }>("/v1/next/collections/private", { + ...limited, + schema: { body: { + type: "object", additionalProperties: false, required: ["collection_id", "device_id", "challenge", "sig"], + properties: { collection_id: uuid, ...proof } + } } + }, async (request, reply) => { + reply.header("cache-control", "no-store"); + const connector = await requireConnector(request, reply, options.db); + if (!connector) return reply; + const body = { ...request.body, collection_id: request.body.collection_id.toLowerCase(), device_id: request.body.device_id.toLowerCase() }; + const collection = body.collection_id; + if (collection === NIL || body.device_id === NIL) return reply.code(400).send(apiError("invalid_request", "Nil identifiers are not accepted.")); + const digest = (challenge: Uint8Array) => privateCreateDigest({ challenge, connector: connector.id, device: body.device_id, collection }); + let device: Device; + try { + // 1. Proof, current identity and ownership, then genesis, in one transaction. + // Nothing here calls the network. + device = await inTransaction(options.db, async (client) => { + await lock(client, collection); + const owner = await authenticate(client, body, connector, digest); + await currentIdentity(client, connector, body.device_id, owner); + if (!(await existing(client, collection, connector.user_id))) { + await registerNextCollection(client, { + collectionId: collection, ownerUserId: connector.user_id, runtime: "next", sync: "private", rootKeyId, + ops: [ + { op: "genesis", owner: connector.user_id, root: rootKeyId, state: "e2e" }, + { op: "member-set", account: connector.user_id, role: "owner" }, + enrolOp(body.device_id, connector.user_id, owner) + ] + }); + } + return owner; + }); + } catch (error) { + return refuse(reply, error, "The private collection was not created; retry with a fresh proof."); + } + try { + // 2. Only an appended genesis whose exact bytes the log returns counts as created. + await options.emitter.drainCollection(collection); + const row = (await options.db.query<{ seq: string; item: Buffer; state: string }>( + "SELECT seq, item, state FROM next_policy_batches WHERE collection_id = $1 AND seq = 1 ORDER BY id LIMIT 1", [collection] + )).rows[0]; + const genesis = await appendedBatch(collection, row); + const head = await options.log.head(collection); + // 3. After every await: the identity, the collection and the enrolment are + // current, and stay locked until the token is minted. + return await inTransaction(options.db, async (client) => { + await currentIdentity(client, connector, body.device_id, device); + await currentPrivate(client, collection, connector.user_id); + // Ownership and a historical genesis are not membership: the owner account + // must still be a current member (a member-remove counts even while pending). + await currentMember(client, collection, connector.user_id); + await refuseRevoked(client, collection, body.device_id); + // The requesting device must be the one the genesis enrolled, with the same keys. + const enrolled = await client.query( + `SELECT 1 FROM next_policy_outbox WHERE id = (SELECT min(id) FROM next_policy_outbox WHERE collection_id = $1) + AND ops->'ops' @> $2::jsonb`, + [collection, exactEnrolment(body.device_id, connector.user_id, device)] + ); + if (!enrolled.rows.length) throw new CreateError(409, "collection_exists"); + return { + collection_id: collection, state: "private", owner_account: connector.user_id, + log_url: options.next.logService.url, head: { seq: head.seq, chain: Buffer.from(head.chain).toString("hex") }, + root_public_key: Buffer.from(options.next.rootPublicKey).toString("hex"), policy_cert: options.next.policyCert, + genesis: { seq: 1, item: genesis.item.toString("hex") }, + // Advisory: the creating owner device is an editor user device, the legal + // initial-rekey signer in e2e (never hosted or escrow); at genesis it is + // the only active device, so it wraps the first epoch key to itself. + rekey_recipients: [body.device_id], + device: mint(body.device_id, device.sign_pk, collection) + }; + }); + } catch (error) { + if (error instanceof CreateError && error.status !== 503) return refuse(reply, error, "The private collection is not current for this device."); + return reply.code(503).send(apiError("not_ready", "The private collection outcome is not verified; retry with a fresh proof.")); + } + }); + + // ---- Device enrol: a member account's registered device, with its SAS commitment. ---- + app.post<{ Params: { id: string }; Body: Proof & { sas_commit: string } }>("/v1/next/collections/:id/private/devices", { + ...limited, + schema: { + params: { type: "object", required: ["id"], properties: { id: uuid } }, + body: { + type: "object", additionalProperties: false, required: ["device_id", "challenge", "sig", "sas_commit"], + properties: { ...proof, sas_commit: { type: "string", pattern: "^[0-9a-f]{64}$" } } + } + } + }, async (request, reply) => { + reply.header("cache-control", "no-store"); + const connector = await requireConnector(request, reply, options.db); + if (!connector) return reply; + const collection = request.params.id.toLowerCase(); + const body = { ...request.body, device_id: request.body.device_id.toLowerCase() }; + const sasCommit = Buffer.from(body.sas_commit, "hex"); + if (collection === NIL || body.device_id === NIL) return reply.code(400).send(apiError("invalid_request", "Nil identifiers are not accepted.")); + const digest = (challenge: Uint8Array) => privateDeviceEnrolDigest({ challenge, connector: connector.id, device: body.device_id, collection, sasCommit }); + const checks = async (client: DatabaseConnection, d: Device) => { + // The connector, account and device are current and stay locked; the collection + // is a current private one and the account a current member. A cloud copy, or a + // collection that has left sync, refuses before any policy op exists. + await currentIdentity(client, connector, body.device_id, d); + await currentPrivate(client, collection); + await currentMember(client, collection, connector.user_id); + await refuseRevoked(client, collection, body.device_id); + }; + let device: Device; + try { + device = await inTransaction(options.db, async (client) => { + await lock(client, collection); + const enrolling = await authenticate(client, body, connector, digest); + await checks(client, enrolling); + const prior = (await client.query(ENROLMENT, [collection, enrolmentKey(body.device_id)])).rows[0]; + if (prior) { + // Idempotent only for the identical tuple and commitment. + const same = (await client.query(ENROLMENT, [collection, exactPrivateEnrolment(body.device_id, connector.user_id, enrolling, sasCommit)])).rows.length > 0; + if (!same) throw new CreateError(409, "device_enrolled_differently"); + } else if (!(await queueNextPolicy(client, collection, [{ + op: "device-enrol", device: body.device_id, account: connector.user_id, kind: enrolling.kind, + signPublicKey: enrolling.sign_pk, kemPublicKey: enrolling.kem_pk, noisePublicKey: enrolling.noise_pk, sasCommit + }]))) { + throw new CreateError(409, "not_current_private"); + } + return enrolling; + }); + } catch (error) { + return refuse(reply, error, "The device was not enrolled; retry with a fresh proof."); + } + try { + await options.emitter.drainCollection(collection); + const row = (await options.db.query<{ seq: string | null; item: Buffer | null; state: string | null }>( + ENROLMENT, [collection, exactPrivateEnrolment(body.device_id, connector.user_id, device, sasCommit)] + )).rows[0]; + const batch = await appendedBatch(collection, row); + return await inTransaction(options.db, async (client) => { + await checks(client, device); + // An existing keyed device approves this one (SAS) and grants it the key next. + return { collection_id: collection, enrolled_at: batch.seq, approval: "pending", device: mint(body.device_id, device.sign_pk, collection) }; + }); + } catch (error) { + if (error instanceof CreateError && error.status !== 503) return refuse(reply, error, "The private collection is not current for this device."); + return reply.code(503).send(apiError("not_ready", "The enrolment is not verified; retry with a fresh proof.")); + } + }); +}