diff --git a/doc/AGENT-IDENTITY.md b/doc/AGENT-IDENTITY.md index f8dabb3c54..0e379b905d 100644 --- a/doc/AGENT-IDENTITY.md +++ b/doc/AGENT-IDENTITY.md @@ -98,3 +98,7 @@ have distinct identities. Native failure records and reports redact the assigned key before truncating diagnostics. Streaming output buffers settle at item or turn completion: short structural prefixes (such as a trailing dash) are preserved, longer interrupted key fragments become redaction markers, and pending buffers are cleared. Terminal events carry these settled `outputTails`; transcript projection displays them as deltas without changing the source event receipt. Codex shell delivery preserves its default `KEY`/`SECRET`/`TOKEN` name exclusions for other configured credentials. The exceptions are the three identity variables and `PAPERCLIP_API_KEY`: the latter is the short-lived, scoped run/bridge credential required by the Paperclip agent skill’s Bash/curl API calls. Disabling Codex’s automatic exclusions is paired with this explicit filtered allowlist; it does not admit arbitrary host environment variables. Provider authentication secrets can still reach the provider process without being newly exposed to shell commands by this feature. + +Cloud customer-success inspection can consume this existing identity together +with strict active-run authority. See [inspection support](CUSTOMER-SUCCESS-INSPECTION.md). +No additional agent keypair or private-key distribution is introduced. diff --git a/doc/CUSTOMER-SUCCESS-INSPECTION.md b/doc/CUSTOMER-SUCCESS-INSPECTION.md new file mode 100644 index 0000000000..e20a2c807f --- /dev/null +++ b/doc/CUSTOMER-SUCCESS-INSPECTION.md @@ -0,0 +1,85 @@ +# Cloud customer-success inspection support + +Paperclip exposes `/api/customer-success/v1/run-authority` and `/read` for the +coordinated Cloud inspection broker. This foundation supports one registered +Paperclip agent; bot creation, scheduling, scoring and reporting are separate. + +The authority endpoint is disabled unless +`PAPERCLIP_CUSTOMER_SUCCESS_AUTHORITY_ENABLED=true`. It accepts only a strict, +current instance/company-derived managed-run JWT. It requires the authenticated +agent's existing provisioned Ed25519 identity and an active running heartbeat; +paused, cancelled/stopped, finished and unsupported-runtime agents are rejected. +No identity GET provisions keys. Normal API JWT compatibility remains unchanged; +legacy signatures are rejected specifically at this authority boundary. + +Tenant reads are disabled unless +`PAPERCLIP_CUSTOMER_SUCCESS_INSPECTION_ENABLED=true`. Configure the public-only +`PAPERCLIP_CUSTOMER_SUCCESS_INSPECTION_JWKS` and the stable HTTPS +`PAPERCLIP_CUSTOMER_SUCCESS_CLOUD_ORIGIN`. Cloud provision/roll delivers the trust +configuration. The persisted runtime stack identity is authoritative; an explicit +`PAPERCLIP_CLOUD_STACK_ID` is usable on operator-configured qualification stacks. +Cloud's dedicated inspection key is distinct from the inspecting agent's key. + +The tenant verifies a maximum-sixty-second, audience/type-bound Ed25519 permit +with exact stack, query, grant, registered key/agent/run, request and binding +version. It calls Cloud's permit-consumption endpoint before reading, so replay +fencing and every inspection record live in Cloud. Cloud remains responsible for +stack eligibility, seven days from creation/pooled claim, human older-stack +approval, active source-run checks, limits, revocation and audits. The source run +bearer is never forwarded to tenants. + +The namespace terminates before actor middleware. It never creates users, +memberships, sessions, run-identity snapshots, activity records or read receipts. +All database readers run inside repeatable-read, read-only transactions. An +explicit catalog avoids existing GET side effects such as skill reconciliation, +instruction recovery/adoption and provider-trace cleanup. There is no arbitrary +GET proxy, write operation or credential-resolution operation. + +The catalog in `services/customer-success-inspection.ts` covers companies, +directory/membership/permission metadata, agents/config/instruction revisions and +public identity, tasks/comments/documents/interactions, approvals/decisions, +routines/schedules, projects/workspaces/stored repository assignments, skill +sources/versions, execution events/trace metadata, activity/costs, work products/ +assets/attachments, and existing readable connections/installs/grants/catalogs. +`packages/shared/src/customer-success.ts` is the versioned wire contract and is +mirrored byte-for-byte in the private Cloud broker; update both together. + +`list`/`get` readers enforce company scope and bounded pagination. Special +operations read instruction files without repair, installed skill snapshots, +existing run logs, workspace files through the current file-resource service, +and company-scoped storage assets. Workspace context uses an owning task and +existing project/workspace checks. Binary envelopes and inclusive byte ranges +pass through Cloud; ranges cap at 1 MiB, JSON at 2 MiB, and lists at 100 rows. +Remote or unsnapshotted resources report unavailable; large content requires +pagination or download ranges. No storage credentials or bypass URLs are issued. + +Existing config/event/run redactions and path/symlink/secret-file/size restrictions +remain in use. Authentication account/session, secret, private-key/encrypted-key +and raw provider-trace tables/readers are excluded. Identity-enabled trace +suppression stays intact. No new prose/file scanner is added; arbitrary pasted +secrets in otherwise readable content may be present. + +Deploy this support first, then Cloud's migration/broker/admin UI. Keep both +inspection and Cloud policy disabled until staging and internal canary checks +pass. Wake only idle sleeps through Cloud's existing controller; normal startup +and background writes after wake are separate from the inspection read. Existing +human login/Slack alerts are unchanged. + +Focused verification: + +```sh +pnpm exec vitest run server/src/__tests__/customer-success.test.ts \ + server/src/__tests__/agent-auth-jwt.test.ts server/src/__tests__/agent-identity.test.ts +# Coordinated qualification, with the sibling Cloud build: +PAPERCLIP_INSPECTION_CLOUD_DIST=/absolute/path/paperclip-cloud/dist \ + pnpm exec vitest run server/src/__tests__/customer-success-cloud.test.ts +``` + +The coordinated test uses two disposable local PostgreSQL databases, a real +managed process with the existing identity/env injection, a PostgreSQL-backed +Cloud broker, HTTP permit consumption, the existing wake controller with a local +provider, and bounded binary file reads. It tests durable replay fencing and +never touches a customer. Cloud's operations document covers enrollment, +capability grants, exclusions, approval/expiry/revocation, audit retention and +rollback. During rollback, disable Cloud policy first, disable tenant inspection, +and preserve identity material and Cloud audit history. diff --git a/doc/plans/2026-10-06-customer-success-inspection.md b/doc/plans/2026-10-06-customer-success-inspection.md new file mode 100644 index 0000000000..37e52e5a50 --- /dev/null +++ b/doc/plans/2026-10-06-customer-success-inspection.md @@ -0,0 +1,48 @@ +# Customer-success inspection foundation + +Approved scope: one designated Paperclip agent uses its existing persistent +Ed25519 identity plus a verified, active managed-run credential. Cloud owns +short-lived grants, older-stack human approvals, durable replay fencing and audit. +Customer stacks expose an explicit read-only API before ordinary auth middleware. +Automatic eligibility is seven days from stack creation (or recorded pooled +claim), with no user/account/signup-age condition. The eventual bot and morning +schedule are separate work. + +## Implementation and qualification + +- Paperclip: `codex/customer-success-inspection`, updated to master `ac8f3eb14`. +- Cloud: same branch name, isolated worktree `/private/tmp/paperclip-cloud-customer-success`, updated to master `044160d0`. +- Core strict JWT/run-authority, versioned permit routes, explicit resource catalog, + read-only transactions, instruction/workspace/asset/snapshot readers implemented. +- Cloud broker, dedicated signer, durable store/migration, admin-session/capability + routes, registration/kill switch/approval/audit UI and agent client implemented. +- Core catalog, file restrictions, active-run authority and existing identity tests pass. + The complete database snapshot remains equal before and after catalog reads. +- Coordinated qualification passes with a real managed process agent, separate + source/customer PostgreSQL databases, two Cloud broker replicas, HTTP permits, + runtime-role append-only enforcement, one-year retention, the existing wake + controller with a local provider, and binary file reads through Cloud. +- Cloud root suite: 2,544 pass, 73 expected skips. Admin web: 779 pass. The 23 + inspection checks cover scoped approval review, replay, revocation, bounded + anonymous audit traffic, and database-pool concurrency. Routing/wake smoke passes. +- Browser acceptance verified native admin login, navigation, exact target review, + approval, revocation, immediate disablement and persistence after refresh. +- Full core typecheck and build pass. The local broad test run recorded 14,277 + passes and four failures in existing suites (two timeouts and two PR-metadata + mock assertions); all three affected suites pass on isolated reruns. The full + CI matrix passes at the implementation head. Public PR #15405 and its companion + Cloud PR carry current review/check status. No deployment or customer enablement. + +## Required invariants + +No database credentials, owner-login fallback, private identity material, tenant +sessions/memberships/activity/read receipts or routine Slack alerts. Preserve +existing content redactions, file restrictions and raw-trace suppression. Audit +before dispatch and before release; recheck run, policy, grant and age exception. +The source bearer never reaches customers. Tenant permits are short-lived, +operation-bound and consumed atomically in Cloud. Human approvals override age +only, freeze exact targets, and can authorize later runs of the same identity. + +Deploy core support before Cloud migration/broker/UI. Disabled by default; staging +and internal canary qualification precede customer enablement. Default retention +one year; rollback/reenrollment/exclusions/retention runbooks accompany both PRs. diff --git a/packages/shared/src/customer-success.ts b/packages/shared/src/customer-success.ts new file mode 100644 index 0000000000..a2d02aaa82 --- /dev/null +++ b/packages/shared/src/customer-success.ts @@ -0,0 +1,221 @@ +/** Version 1 inspection wire contract. Keep Cloud's protocol.ts byte-identical. */ +export const INSPECTION_VERSION = 1; +export const INSPECTION_AUDIENCE = "paperclip-customer-success/v1"; +export const INSPECTION_PERMIT_TYPE = "paperclip-inspection+jwt"; +export const INSPECTION_HEADER = "x-paperclip-cloud-inspection"; +export const CHALLENGE_DOMAIN = "paperclip-customer-success-challenge/v1\n"; +export const MAX_INSPECTION_BYTES = 1024 * 1024; +export const MAX_INSPECTION_ROWS = 100; +export const INSPECTION_RESOURCES = [ + "companies", + "goals", + "users", + "memberships", + "permissions", + "agents", + "agentIdentities", + "agentInstructions", + "agentConfigRevisions", + "tasks", + "comments", + "documents", + "documentRevisions", + "taskDocuments", + "interactions", + "approvals", + "approvalComments", + "decisions", + "routines", + "routineTriggers", + "routineRuns", + "routineRevisions", + "routineDocuments", + "projects", + "projectWorkspaces", + "projectMemberships", + "skills", + "skillVersions", + "skillSources", + "runs", + "runEvents", + "traceMetadata", + "activity", + "costs", + "workProducts", + "artifacts", + "attachments", + "connections", + "connectionInstalls", + "connectionGrants", + "connectionGrantMembers", + "connectionCatalog", + "toolProfiles", + "toolProfileEntries", + "agentMemberships", + "agentCommentary", + "skillPolicies", + "projectGoals", + "labels", + "taskLabels", + "taskRelations", + "taskApprovals", + "documentMemberships", + "documentThreads", + "documentComments", + "threads", +] as const; +export type InspectionResource = (typeof INSPECTION_RESOURCES)[number]; +export interface InspectionQuery { + operation: + | "list" + | "get" + | "instructions.file" + | "skills.file" + | "runs.log" + | "files.list" + | "files.read" + | "files.download" + | "assets.content"; + resource?: InspectionResource; + companyId?: string; + resourceId?: string; + parentId?: string; + path?: string; + versionId?: string; + limit?: number; + offset?: number; + range?: { start: number; end: number }; + context?: { + projectId?: string; + workspaceId?: string; + workspace?: "auto" | "execution" | "project"; + }; +} +export interface RunAuthority { + version: 1; + instanceId: string; + companyId: string; + agentId: string; + runId: string; + keyId: string; + active: true; +} +export interface InspectionPermit { + v: 1; + iss: "paperclip-cloud"; + aud: typeof INSPECTION_AUDIENCE; + sub: string; + jti: string; + iat: number; + exp: number; + bindingVersion: number; + grantId: string; + requestId: string; + agentId: string; + keyId: string; + runId: string; + query: InspectionQuery; +} +export function canonicalJson(value: unknown): string { + if (value === null || typeof value === "string" || typeof value === "boolean") + return JSON.stringify(value); + if (typeof value === "number" && Number.isFinite(value)) return JSON.stringify(value); + if (Array.isArray(value)) return `[${value.map(canonicalJson).join(",")}]`; + if (value && typeof value === "object") { + const entries = Object.entries(value) + .filter(([, v]) => v !== undefined) + .sort(([a], [b]) => (a < b ? -1 : a > b ? 1 : 0)); + return `{${entries.map(([k, v]) => `${JSON.stringify(k)}:${canonicalJson(v)}`).join(",")}}`; + } + throw new Error("Invalid JSON value"); +} +export function parseInspectionQuery(value: unknown): InspectionQuery { + if (!value || typeof value !== "object" || Array.isArray(value)) + throw new Error("Invalid inspection query"); + const q = value as InspectionQuery; + const allowed = new Set([ + "operation", + "resource", + "companyId", + "resourceId", + "parentId", + "path", + "versionId", + "limit", + "offset", + "range", + "context", + ]); + if (Object.keys(q).some((k) => !allowed.has(k))) + throw new Error("Unsupported inspection parameter"); + if ( + ![ + "list", + "get", + "instructions.file", + "skills.file", + "runs.log", + "files.list", + "files.read", + "files.download", + "assets.content", + ].includes(q.operation) + ) + throw new Error("Unsupported inspection operation"); + for (const k of ["companyId", "resourceId", "parentId", "versionId"] as const) { + if (q[k] !== undefined && (typeof q[k] !== "string" || !/^[a-zA-Z0-9_-]{1,128}$/.test(q[k]!))) + throw new Error(`Invalid ${k}`); + } + if (q.resource !== undefined && !INSPECTION_RESOURCES.includes(q.resource)) + throw new Error("Unsupported inspection resource"); + if (q.operation === "get" && !q.resourceId) throw new Error("resourceId required"); + if (["get", "list"].includes(q.operation) && !q.resource) throw new Error("resource required"); + if (!(q.operation === "list" && q.resource === "companies") && !q.companyId) + throw new Error("companyId required"); + if (!["get", "list"].includes(q.operation) && !q.resourceId) + throw new Error("resourceId required"); + if ( + q.path !== undefined && + (typeof q.path !== "string" || q.path.length > 1024 || q.path.includes("\0")) + ) + throw new Error("Invalid path"); + if ( + q.range !== undefined && + !["files.download", "assets.content", "runs.log"].includes(q.operation) + ) + throw new Error("This operation does not support byte ranges"); + for (const k of ["limit", "offset"] as const) + if ( + q[k] !== undefined && + (!Number.isSafeInteger(q[k]) || q[k]! < 0 || q[k]! > (k === "limit" ? 100 : 1000000)) + ) + throw new Error(`Invalid ${k}`); + if ( + q.range !== undefined && + (!q.range || + typeof q.range !== "object" || + Array.isArray(q.range) || + Object.keys(q.range).sort().join(",") !== "end,start" || + !Number.isSafeInteger(q.range.start) || + !Number.isSafeInteger(q.range.end) || + q.range.start < 0 || + q.range.end < q.range.start || + q.range.end - q.range.start + 1 > MAX_INSPECTION_BYTES) + ) + throw new Error("Invalid byte range"); + if (q.context !== undefined) { + if ( + !q.context || + typeof q.context !== "object" || + Array.isArray(q.context) || + Object.keys(q.context).some((k) => !["projectId", "workspaceId", "workspace"].includes(k)) || + Object.values(q.context).some( + (v) => typeof v !== "string" || !/^[a-zA-Z0-9_-]{1,128}$/.test(v), + ) || + (q.context.workspace !== undefined && + !["auto", "execution", "project"].includes(q.context.workspace)) + ) + throw new Error("Invalid file context"); + } + return JSON.parse(canonicalJson(q)) as InspectionQuery; +} diff --git a/packages/shared/src/index.ts b/packages/shared/src/index.ts index 4c01d1b482..0e40261c5a 100644 --- a/packages/shared/src/index.ts +++ b/packages/shared/src/index.ts @@ -2832,3 +2832,4 @@ export { aiConnectionRouterSlug, aiConnectionRouterAppDefinition, aiConnectionRo export { isAppAggregator, aggregatorManagementUrl, aggregatorAppsSyncSchema, aggregatorAppsRefreshSchema, arcadeDiscoverySetupSchema, type AggregatorAppSnapshot, type AggregatorAppsResponse, type ArcadeDiscoverySetupInput } from "./aggregator-apps.js"; export * from "./connection-instructions.js"; +export * from "./customer-success.js"; diff --git a/server/src/__tests__/customer-success-cloud.test.ts b/server/src/__tests__/customer-success-cloud.test.ts new file mode 100644 index 0000000000..3fcbb12530 --- /dev/null +++ b/server/src/__tests__/customer-success-cloud.test.ts @@ -0,0 +1,407 @@ +import { createServer } from "node:http"; +import type { AddressInfo } from "node:net"; +import { generateKeyPairSync, randomUUID } from "node:crypto"; +import { mkdtemp, readFile, rm, stat, writeFile } from "node:fs/promises"; +import { join, resolve } from "node:path"; +import { tmpdir } from "node:os"; +import { pathToFileURL } from "node:url"; +import express from "express"; +import { describe, it, expect, vi } from "vitest"; +import { + createDb, + agents, + companies, + heartbeatRuns, + issues, + projects, + projectWorkspaces, +} from "@paperclipai/db"; +import { customerSuccessRoutes } from "../routes/customer-success.js"; +import { agentIdentityService } from "../services/agent-identity.js"; +import { createLocalAgentJwt } from "../agent-auth-jwt.js"; +import { execute } from "../adapters/process/execute.js"; +import { errorHandler } from "../middleware/error-handler.js"; +import { startEmbeddedPostgresTestDatabase } from "./helpers/embedded-postgres.js"; +import type { StorageService } from "../storage/types.js"; + +// Coordinated qualification against the sibling Cloud build, never a customer +// stack or hosted provider. Default unit runs have no cross-repository dependency. +const cloudDist = process.env.PAPERCLIP_INSPECTION_CLOUD_DIST; +describe.skipIf(!cloudDist)("customer-success managed agent / Cloud / tenant qualification", () => { + it("uses the existing managed-runtime identity, active JWT, PostgreSQL broker, HTTP permits, and audited tenant reads", async () => { + const load = (module: string) => import(pathToFileURL(resolve(cloudDist!, module)).href); + const [ + { CustomerSuccessInspection }, + { InspectionSigner }, + { PostgresInspectionStore }, + { inspectionRoute }, + { InMemoryCloudHarnessRegistry }, + { CloudStackSleepController }, + pg, + ] = await Promise.all([ + load("customer-success/service.js"), + load("customer-success/crypto.js"), + load("customer-success/store.js"), + load("customer-success/routes.js"), + load("provisioner/memory.js"), + load("sleep/controller.js"), + load("../node_modules/pg/lib/index.js"), + ]); + const home = await startEmbeddedPostgresTestDatabase("inspection-home-"); + const tenant = await startEmbeddedPostgresTestDatabase("inspection-customer-"); + const directory = await mkdtemp(join(tmpdir(), "inspection-journey-")); + const servers: ReturnType[] = []; + const serverErrors: string[] = []; + const pool = new pg.default.Pool({ + connectionString: home.connectionString, + max: 2, + connectionTimeoutMillis: 1000, + }); + try { + vi.stubEnv("PAPERCLIP_HOME", directory); + vi.stubEnv("PAPERCLIP_INSTANCE_ID", "qualification-home"); + vi.stubEnv("PAPERCLIP_AGENT_JWT_SECRET", "isolated-qualification-jwt"); + vi.stubEnv("PAPERCLIP_SECRETS_MASTER_KEY_FILE", join(directory, "master.key")); + vi.stubEnv("PAPERCLIP_SECRETS_MASTER_KEY", ""); + vi.stubEnv("PAPERCLIP_CUSTOMER_SUCCESS_AUTHORITY_ENABLED", "true"); + vi.stubEnv("PAPERCLIP_CUSTOMER_SUCCESS_INSPECTION_ENABLED", "true"); + vi.stubEnv("PAPERCLIP_CLOUD_STACK_ID", "qualification-stack"); + const homeDb = createDb(home.connectionString); + const tenantDb = createDb(tenant.connectionString); + const [homeCompany] = await homeDb + .insert(companies) + .values({ name: "Internal", issuePrefix: "INT" }) + .returning(); + const [agent] = await homeDb + .insert(agents) + .values({ + companyId: homeCompany.id, + name: "Success qualification", + adapterType: "process", + }) + .returning(); + const identity = await agentIdentityService(homeDb).ensureAgentIdentity( + homeCompany.id, + agent.id, + ); + const [run] = await homeDb + .insert(heartbeatRuns) + .values({ + companyId: homeCompany.id, + agentId: agent.id, + status: "running", + startedAt: new Date(), + }) + .returning(); + const [customer] = await tenantDb + .insert(companies) + .values({ name: "Synthetic customer", issuePrefix: "SYN" }) + .returning(); + const [project] = await tenantDb + .insert(projects) + .values({ companyId: customer.id, name: "Qualification project" }) + .returning(); + const [workspace] = await tenantDb + .insert(projectWorkspaces) + .values({ + companyId: customer.id, + projectId: project.id, + name: "Qualification files", + cwd: directory, + }) + .returning(); + const filePath = join(directory, "result.bin"); + await writeFile(filePath, Buffer.from([0, 1, 2, 3, 4, 5])); + const fileBefore = await stat(filePath); + const [task] = await tenantDb + .insert(issues) + .values({ companyId: customer.id, projectId: project.id, title: "First successful task" }) + .returning(); + const storage = { provider: "local_disk" } as StorageService; + const start = async (handler: Parameters[0]) => { + const server = createServer(handler); + servers.push(server); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + return `http://127.0.0.1:${(server.address() as AddressInfo).port}`; + }; + const homeApp = express(); + homeApp.use(express.json()); + homeApp.use("/api/customer-success/v1", customerSuccessRoutes(homeDb, storage)); + homeApp.use(errorHandler); + const sourceOrigin = await start(homeApp); + const tenantApp = express(); + tenantApp.use(express.json()); + tenantApp.use("/api/customer-success/v1", customerSuccessRoutes(tenantDb, storage)); + tenantApp.use(errorHandler); + const tenantOrigin = await start(tenantApp); + await pool.query("CREATE SCHEMA cloud_harness"); + await pool.query( + await readFile( + resolve(cloudDist!, "../migrations/0058_customer_success_inspection.sql"), + "utf8", + ), + ); + const registry = new InMemoryCloudHarnessRegistry(); + let transactionReads = 0; + const registryForClient = (client: any) => ({ + getStack: async (id: string) => { + await client.query("SELECT 1"); + transactionReads++; + return registry.getStack(id); + }, + getAccount: async (id: string) => { + await client.query("SELECT 1"); + transactionReads++; + return registry.getCustomerSuccessAccount(id); + }, + listStackIds: async () => registry.listCustomerSuccessStackIds(), + }); + const store = new PostgresInspectionStore(pool, registryForClient); + const now = new Date(); + registry.accountGroups.set("fixture-account", { + id: "fixture-account", + kind: "customer", + billingStatus: "unconfigured", + }); + registry.stacks.set("qualification-stack", { + id: "qualification-stack", + accountGroupId: "fixture-account", + displayName: "Synthetic customer", + primaryHost: "fixture.example.test", + kind: "customer_production", + lifecycleState: "sleeping", + sleepState: "sleeping", + sleepReason: "idle", + version: 1, + updatedAt: now, + rolloutMetadata: { track: "stable" }, + providerRefs: [ + { + provider: "railway", + resourceType: "service_domain:web", + externalId: "fixture", + metadata: { upstreamOrigin: tenantOrigin }, + }, + ], + createdAt: now, + }); + const signer = new InspectionSigner( + "qualification", + generateKeyPairSync("ed25519").privateKey.export({ format: "jwk" }), + ); + vi.stubEnv("PAPERCLIP_CUSTOMER_SUCCESS_INSPECTION_JWKS", signer.publicJwks); + // Use the production wake controller with a disposable local provider. + // The tenant HTTP process is already listening; no hosted stack is woken. + const sleep = new CloudStackSleepController({ + registry, + provider: { + name: "fake", + wakeStack: async ( + _id: string, + context: { stack: Record; providerRefs: unknown[] }, + ) => ({ + stack: { ...context.stack, lifecycleState: "active", sleepState: "awake" }, + resources: context.providerRefs, + secretRefs: { providerAdminCredentials: {} }, + operations: [], + }), + inspectStack: async () => ({ + stackId: "qualification-stack", + lifecycleState: "active", + sleepState: "awake", + resources: registry.stacks.get("qualification-stack").providerRefs, + }), + }, + }); + const broker = new CustomerSuccessInspection({ + store, + registry, + signer, + allowLoopback: true, + wakeStack: async (stackId: string) => { + await sleep.wakeStack({ stackId }); + }, + }); + const cloudOrigin = await start(async (req, res) => { + try { + const handled = await inspectionRoute({ + req, + res, + url: new URL(req.url!, "http://localhost"), + service: broker, + authenticateHuman: async () => { + throw new Error("Qualification never uses human access"); + }, + readJson: async (request: AsyncIterable, max: number) => { + const parts: Buffer[] = []; + let bytes = 0; + for await (const part of request) { + bytes += part.length; + if (bytes > max) throw new Error("Too large"); + parts.push(part); + } + return JSON.parse(Buffer.concat(parts).toString()); + }, + }); + if (!handled) { + res.statusCode = 404; + res.end(); + } + } catch (error) { + serverErrors.push(String(error)); + res.statusCode = 503; + res.end("{}"); + } + }); + vi.stubEnv("PAPERCLIP_CUSTOMER_SUCCESS_CLOUD_ORIGIN", cloudOrigin); + await broker.configure("qualification-operator", { + binding: { + sourceOrigin, + instanceId: "qualification-home", + companyId: homeCompany.id, + agentId: agent.id, + keyId: identity.keyId, + publicKeyPem: identity.publicKeyPem, + }, + }); + await broker.configure("qualification-operator", { enabled: true }); + const program = ` + const {inspectionAgentClient}=await import(${JSON.stringify(pathToFileURL(resolve(cloudDist!, "customer-success/client.js")).href)}); + const call=inspectionAgentClient(); + const discovery=await call('discover',{}); + const grant=await call('grant',{stackId:'qualification-stack'}); + const tasks=await call('read',{stackId:'qualification-stack',grantToken:grant.token,query:{operation:'list',resource:'tasks',companyId:${JSON.stringify(customer.id)}}}); + const file=await call('read',{stackId:'qualification-stack',grantToken:grant.token,query:{operation:'files.download',companyId:${JSON.stringify(customer.id)},resourceId:${JSON.stringify(task.id)},context:{projectId:${JSON.stringify(project.id)},workspaceId:${JSON.stringify(workspace.id)}},path:'result.bin',range:{start:2,end:5}}}); + console.log(JSON.stringify({stackCount:discovery.items.length,titles:tasks.items.map(t=>t.title),fileBytes:[...Buffer.from(file.content,'base64')]})); + `; + const programFile = join(directory, "qualify.mjs"); + await writeFile(programFile, program); + const output: string[] = []; + const result = await execute({ + runId: run.id, + agent, + agentIdentity: identity, + authToken: createLocalAgentJwt(agent.id, homeCompany.id, "process", run.id)!, + config: { + command: process.execPath, + args: [programFile], + cwd: directory, + env: { PAPERCLIP_CUSTOMER_SUCCESS_CLOUD_ORIGIN: cloudOrigin }, + }, + context: {}, + runtime: { sessionId: null, sessionParams: null, taskKey: null }, + onLog: async (_stream, chunk) => { + output.push(chunk); + }, + onMeta: async () => {}, + }); + expect(result.exitCode, [...output, ...serverErrors].join("\n")).toBe(0); + expect(output.join("")).toContain("First successful task"); + expect(output.join("")).toContain('"fileBytes":[2,3,4,5]'); + expect((await stat(filePath)).mtimeMs).toBe(fileBefore.mtimeMs); + expect(await readFile(filePath)).toEqual(Buffer.from([0, 1, 2, 3, 4, 5])); + const audit = await pool.query( + "SELECT event FROM cloud_harness.customer_success_audit ORDER BY occurred_at", + ); + expect(audit.rows.map((r: { event: string }) => r.event)).toContain("permit.consumed"); + expect(audit.rows.map((r: { event: string }) => r.event)).toContain("read.completed"); + expect(audit.rows.map((r: { event: string }) => r.event)).toContain("stack.woken"); + expect(await tenantDb.select().from(issues)).toEqual([task]); + // The production SQL reader scopes before pagination and hides mixed + // approval cohorts and global policy records from tenant approvers. + await store.transaction(async (tx: any) => { + const localId = randomUUID(); + const mixedId = randomUUID(); + const foreignId = randomUUID(); + for (const [id, detail, stackId] of [ + [localId, { stackIds: ["qualification-stack"] }, undefined], + [mixedId, { stackIds: ["qualification-stack", "foreign-stack"] }, undefined], + [foreignId, {}, "foreign-stack"], + ]) + await tx.audit({ id, at: tx.now, event: "scope.fixture", detail, stackId }); + const scoped = await tx.audits(0, 100, ["qualification-stack"]); + expect(scoped.some((event: any) => event.id === localId)).toBe(true); + expect(scoped.some((event: any) => event.id === mixedId || event.id === foreignId)).toBe( + false, + ); + expect(scoped.some((event: any) => event.event === "policy.changed")).toBe(false); + expect(await tx.audits(0, 100, [])).toEqual([]); + const last = await tx.audits(scoped.length - 1, 100, ["qualification-stack"]); + expect(last).toHaveLength(1); + }); + // Two independent pool connections prove consumption is durable, not a + // process-local Set: exactly one replica can consume a signed challenge. + const token = createLocalAgentJwt(agent.id, homeCompany.id, "process", run.id)!; + const c = await broker.challenge(token, "discover", {}); + const { sign } = await import("node:crypto"); + const sig = sign(null, Buffer.from(c.bytes), identity.privateKeyPem).toString("base64url"); + const replica = new CustomerSuccessInspection({ + store: new PostgresInspectionStore(pool, registryForClient), + registry, + signer, + allowLoopback: true, + }); + const competing = await Promise.allSettled([ + broker.execute(token, c.challengeId, sig), + replica.execute(token, c.challengeId, sig), + ]); + expect(competing.filter((r) => r.status === "fulfilled")).toHaveLength(1); + const grants = await Promise.all( + Array.from({ length: 4 }, () => + broker.challenge(token, "grant", { stackId: "qualification-stack" }), + ), + ); + const concurrentGrants = await Promise.all( + grants.map((proof) => + broker.execute( + token, + proof.challengeId, + sign(null, Buffer.from(proof.bytes), identity.privateKeyPem).toString("base64url"), + ), + ), + ); + expect(concurrentGrants).toHaveLength(4); + expect(transactionReads).toBeGreaterThan(0); + // Prove append-only enforcement with a runtime role, even if someone + // accidentally grants it UPDATE/DELETE privileges on the audit table. + const client = await pool.connect(); + const runtimeRole = `inspection_runtime_${randomUUID().replaceAll("-", "")}`; + try { + await client.query(`CREATE ROLE ${runtimeRole}`); + await client.query(`GRANT USAGE ON SCHEMA cloud_harness TO ${runtimeRole}`); + await client.query( + `GRANT SELECT, UPDATE, DELETE ON cloud_harness.customer_success_audit TO ${runtimeRole}`, + ); + await client.query(`SET ROLE ${runtimeRole}`); + await expect( + client.query("UPDATE cloud_harness.customer_success_audit SET event = 'rewritten'"), + ).rejects.toThrow("append-only"); + await expect( + client.query("DELETE FROM cloud_harness.customer_success_audit"), + ).rejects.toThrow("append-only"); + } finally { + await client.query("RESET ROLE"); + client.release(); + } + const { applyInspectionRetention } = await load("customer-success/retention.js"); + await expect(applyInspectionRetention(pool, 364)).rejects.toThrow("at least one year"); + await pool.query( + "INSERT INTO cloud_harness.customer_success_audit VALUES ($1,clock_timestamp() - interval '366 days','old.fixture','{}')", + [randomUUID()], + ); + expect((await applyInspectionRetention(pool)).auditDeleted).toBe(1); + expect( + (await pool.query("SELECT count(*) FROM cloud_harness.customer_success_audit")).rows[0] + .count, + ).not.toBe("0"); + } finally { + for (const server of servers.reverse()) + await new Promise((resolve) => server.close(() => resolve())); + await pool.end(); + await tenant.cleanup(); + await home.cleanup(); + vi.unstubAllEnvs(); + await rm(directory, { recursive: true, force: true }); + } + }, 60000); +}); diff --git a/server/src/__tests__/customer-success.test.ts b/server/src/__tests__/customer-success.test.ts new file mode 100644 index 0000000000..30ace46a0e --- /dev/null +++ b/server/src/__tests__/customer-success.test.ts @@ -0,0 +1,605 @@ +import { randomUUID, generateKeyPairSync, createHash, createHmac, sign } from "node:crypto"; +import { mkdtemp, writeFile, rm, symlink, stat } from "node:fs/promises"; +import { join } from "node:path"; +import { tmpdir } from "node:os"; +import { Readable } from "node:stream"; +import express from "express"; +import request from "supertest"; +import { afterAll, beforeAll, describe, expect, it, vi } from "vitest"; +import { eq, sql } from "drizzle-orm"; +import * as schema from "@paperclipai/db"; +import { + INSPECTION_AUDIENCE, + INSPECTION_HEADER, + INSPECTION_PERMIT_TYPE, + INSPECTION_RESOURCES, + parseInspectionQuery, + type InspectionQuery, + type InspectionPermit, +} from "@paperclipai/shared"; +import { customerSuccessRunAuthority } from "../services/customer-success-authority.js"; +import { + readCustomerSuccessResource, + inspectionReaders, +} from "../services/customer-success-inspection.js"; +import { customerSuccessRoutes, verifyInspectionPermit } from "../routes/customer-success.js"; +import { agentIdentityService } from "../services/agent-identity.js"; +import { createLocalAgentJwt, verifyLocalAgentJwt } from "../agent-auth-jwt.js"; +import { errorHandler } from "../middleware/error-handler.js"; +import { redactAgentAdapterConfig } from "../redaction.js"; +import type { StorageService } from "../storage/types.js"; +import { startEmbeddedPostgresTestDatabase } from "./helpers/embedded-postgres.js"; + +describe("customer-success read authority", () => { + let db: schema.Db; + let database: Awaited>; + let directory: string; + let companyId: string; + let otherCompanyId: string; + let agentId: string; + let runId: string; + let bearer: string; + const storage: StorageService = { + provider: "local_disk", + putFile: async () => { + throw new Error("write forbidden"); + }, + deleteObject: async () => { + throw new Error("write forbidden"); + }, + headObject: async () => ({ exists: true }), + getObject: async (_companyId, _objectKey, opts) => ({ + stream: Readable.from([ + Buffer.from([0, 1, 2, 3, 4]).subarray( + opts?.range?.start ?? 0, + opts?.range ? opts.range.end + 1 : undefined, + ), + ]), + contentType: "application/octet-stream", + }), + }; + beforeAll(async () => { + directory = await mkdtemp(join(tmpdir(), "inspection-test-")); + vi.stubEnv("PAPERCLIP_HOME", directory); + vi.stubEnv("PAPERCLIP_INSTANCE_ID", "inspection-home"); + vi.stubEnv("PAPERCLIP_AGENT_JWT_SECRET", "inspection-local-fixture"); + vi.stubEnv("PAPERCLIP_SECRETS_MASTER_KEY_FILE", join(directory, "master.key")); + vi.stubEnv("PAPERCLIP_SECRETS_MASTER_KEY", ""); + database = await startEmbeddedPostgresTestDatabase("inspection-db-"); + db = schema.createDb(database.connectionString); + const companies = await db + .insert(schema.companies) + .values([ + { name: "Customer", issuePrefix: "CUS" }, + { name: "Other", issuePrefix: "OTH" }, + ]) + .returning(); + companyId = companies[0].id; + otherCompanyId = companies[1].id; + const [agent] = await db + .insert(schema.agents) + .values({ companyId, name: "Managed inspection fixture", adapterType: "process" }) + .returning(); + agentId = agent.id; + const [run] = await db + .insert(schema.heartbeatRuns) + .values({ companyId, agentId, status: "running", startedAt: new Date() }) + .returning(); + runId = run.id; + bearer = createLocalAgentJwt(agentId, companyId, "process", runId)!; + }); + afterAll(async () => { + await database?.cleanup(); + vi.unstubAllEnvs(); + vi.unstubAllGlobals(); + await rm(directory, { recursive: true, force: true }); + }); + + it("unprovisioned reads never generate a key, strict authority requires a current managed run", async () => { + await expect(customerSuccessRunAuthority(db, bearer)).rejects.toThrow("unprovisioned"); + expect(await db.select().from(schema.agentIdentityKeys)).toHaveLength(0); + const key = await agentIdentityService(db).ensureAgentIdentity(companyId, agentId); + expect(await customerSuccessRunAuthority(db, bearer)).toEqual({ + version: 1, + active: true, + instanceId: "inspection-home", + companyId, + agentId, + runId, + keyId: key.keyId, + }); + await expect(customerSuccessRunAuthority(db, "ordinary-api-key")).rejects.toThrow( + "managed-run", + ); + await db.update(schema.agents).set({ status: "paused" }).where(eq(schema.agents.id, agentId)); + await expect(customerSuccessRunAuthority(db, bearer)).rejects.toThrow("not active"); + await db.update(schema.agents).set({ status: "idle" }).where(eq(schema.agents.id, agentId)); + await db + .update(schema.heartbeatRuns) + .set({ resultJson: { executionCancellation: { state: "requested" } } }) + .where(eq(schema.heartbeatRuns.id, runId)); + await expect(customerSuccessRunAuthority(db, bearer)).rejects.toThrow("not active"); + await db + .update(schema.heartbeatRuns) + .set({ resultJson: null, status: "succeeded" }) + .where(eq(schema.heartbeatRuns.id, runId)); + await expect(customerSuccessRunAuthority(db, bearer)).rejects.toThrow("not active"); + await db + .update(schema.heartbeatRuns) + .set({ status: "running" }) + .where(eq(schema.heartbeatRuns.id, runId)); + }); + it("rejects legacy signatures, missing instance claims and foreign-instance derived tokens", async () => { + const [header, payload] = bearer.split("."); + const input = `${header}.${payload}`; + const legacy = `${input}.${createHmac("sha256", "inspection-local-fixture").update(input).digest("base64url")}`; + expect(verifyLocalAgentJwt(legacy)).not.toBeNull(); + expect(verifyLocalAgentJwt(legacy, { strictRunAuthority: true })).toBeNull(); + vi.stubEnv("PAPERCLIP_INSTANCE_ID", "development-clone"); + const foreign = createLocalAgentJwt(agentId, companyId, "process", runId)!; + vi.stubEnv("PAPERCLIP_INSTANCE_ID", "inspection-home"); + await expect(customerSuccessRunAuthority(db, foreign)).rejects.toThrow("managed-run"); + }); + it("every catalog reader is company scoped, paginated, and leaves the complete database unchanged", async () => { + const snapshot = async () => { + const tables = await db.execute( + sql`select tablename from pg_tables where schemaname = 'public' order by tablename`, + ); + const data: Record = {}; + for (const table of tables) { + const name = String(table.tablename).replaceAll('"', '""'); + data[name] = await db.execute( + sql.raw( + `SELECT coalesce(jsonb_agg(to_jsonb(t) ORDER BY to_jsonb(t)::text),'[]'::jsonb) AS rows FROM "${name}" t`, + ), + ); + } + return JSON.stringify(data); + }; + const before = await snapshot(); + for (const resource of INSPECTION_RESOURCES) { + const result = (await readCustomerSuccessResource(db, storage, { + operation: "list", + resource, + companyId, + limit: 1, + })) as { items: unknown[] }; + expect(result.items.length).toBeLessThanOrEqual(1); + } + expect(await snapshot()).toEqual(before); + await expect( + db.transaction( + (tx) => + tx.update(schema.agents).set({ name: "forbidden" }).where(eq(schema.agents.id, agentId)), + { accessMode: "read only" }, + ), + ).rejects.toThrow(); + expect(inspectionReaders.agentIdentities.omit).toContain("privateKeyMaterial"); + const identity = (await readCustomerSuccessResource(db, storage, { + operation: "get", + resource: "agentIdentities", + companyId, + resourceId: agentId, + })) as Record; + expect(identity).not.toHaveProperty("privateKeyMaterial"); + await expect( + readCustomerSuccessResource(db, storage, { + operation: "get", + resource: "agentIdentities", + companyId: otherCompanyId, + resourceId: agentId, + }), + ).rejects.toThrow("not found"); + }); + it("returns tasks/documents verbatim and uses existing adapter config redaction", async () => { + await db + .update(schema.agents) + .set({ + adapterConfig: { + env: { API_KEY: "hidden", NOTE: "same existing env policy" }, + promptTemplate: "Do useful work", + }, + }) + .where(eq(schema.agents.id, agentId)); + const agent = (await readCustomerSuccessResource(db, storage, { + operation: "get", + resource: "agents", + companyId, + resourceId: agentId, + })) as { adapterConfig: { env: Record; promptTemplate: string } }; + expect(JSON.stringify(agent)).not.toContain("hidden"); + expect(agent.adapterConfig.promptTemplate).toEqual("Do useful work"); + const [task] = await db + .insert(schema.issues) + .values({ + companyId, + title: "A pasted instruction", + description: "Keep prose unchanged: arbitrary pasted material", + }) + .returning(); + const taskResult = (await readCustomerSuccessResource(db, storage, { + operation: "get", + resource: "tasks", + companyId, + resourceId: task.id, + })) as { description: string }; + expect(taskResult.description).toEqual(task.description); + await expect( + readCustomerSuccessResource(db, storage, { + operation: "get", + resource: "tasks", + companyId: otherCompanyId, + resourceId: task.id, + }), + ).rejects.toThrow("not found"); + }); + it("reads existing instruction files without recovery, adoption, or filesystem writes", async () => { + const file = join(directory, "AGENTS.md"); + await writeFile(file, "# Investigate carefully\n"); + await db + .update(schema.agents) + .set({ adapterConfig: { instructionsFilePath: file } }) + .where(eq(schema.agents.id, agentId)); + const before = await stat(file); + const count = (await db.select().from(schema.agentInstructionRevisions)).length; + const result = (await readCustomerSuccessResource(db, storage, { + operation: "instructions.file", + companyId, + resourceId: agentId, + })) as { content: string }; + expect(Buffer.from(result.content, "base64").toString()).toEqual("# Investigate carefully\n"); + expect((await stat(file)).mtimeMs).toEqual(before.mtimeMs); + expect((await db.select().from(schema.agentInstructionRevisions)).length).toEqual(count); + await writeFile(join(directory, ".env"), "CREDENTIAL=do-not-return"); + await expect( + readCustomerSuccessResource(db, storage, { + operation: "instructions.file", + companyId, + resourceId: agentId, + path: ".env", + }), + ).rejects.toThrow("denied by policy"); + await symlink(file, join(directory, "linked.md")); + await expect( + readCustomerSuccessResource(db, storage, { + operation: "instructions.file", + companyId, + resourceId: agentId, + path: "linked.md", + }), + ).rejects.toThrow("symlink"); + }); + it("reuses environment redaction for project and routine configuration and revisions", async () => { + const env = { + CONFIG: { type: "plain" as const, value: "configuration-credential" }, + REFERENCE: { + type: "secret_ref" as const, + secretId: randomUUID(), + version: "latest" as const, + }, + }; + const [project] = await db + .insert(schema.projects) + .values({ companyId, name: "Configured project", env }) + .returning(); + const [routine] = await db + .insert(schema.routines) + .values({ companyId, title: "Configured routine", env }) + .returning(); + const [revision] = await db + .insert(schema.routineRevisions) + .values({ + companyId, + routineId: routine.id, + revisionNumber: 1, + title: routine.title, + snapshot: { + version: 1, + routine: { ...routine }, + triggers: [], + } as typeof schema.routineRevisions.$inferInsert.snapshot, + }) + .returning(); + const expected = redactAgentAdapterConfig({ env }).env; + for (const [resource, resourceId] of [ + ["projects", project.id], + ["routines", routine.id], + ["routineRevisions", revision.id], + ] as const) { + const row = (await readCustomerSuccessResource(db, storage, { + operation: "get", + resource, + companyId, + resourceId, + })) as { env?: unknown; snapshot?: { routine: { env: unknown } } }; + expect(row.snapshot?.routine.env ?? row.env).toEqual(expected); + expect(JSON.stringify(row)).not.toContain("configuration-credential"); + } + await expect( + readCustomerSuccessResource(db, storage, { + operation: "get", + resource: "users", + companyId, + resourceId: "missing", + }), + ).rejects.toThrow("not found"); + }); + it("asset bytes and bounded byte ranges are company scoped, without bypass URLs", async () => { + const [asset] = await db + .insert(schema.assets) + .values({ + companyId, + provider: "local_disk", + objectKey: "test/object", + contentType: "application/octet-stream", + byteSize: 5, + sha256: "fixture", + }) + .returning(); + const result = (await readCustomerSuccessResource(db, storage, { + operation: "assets.content", + companyId, + resourceId: asset.id, + range: { start: 1, end: 3 }, + })) as { content: string; bytes: number }; + expect(Buffer.from(result.content, "base64")).toEqual(Buffer.from([1, 2, 3])); + expect(result.bytes).toEqual(3); + await expect( + readCustomerSuccessResource(db, storage, { + operation: "assets.content", + companyId: otherCompanyId, + resourceId: asset.id, + }), + ).rejects.toThrow("not found"); + }); + it("tenant requests verify exact Cloud permit parameters, consume centrally, reject cookies and unknown paths", async () => { + vi.stubEnv("PAPERCLIP_CUSTOMER_SUCCESS_INSPECTION_ENABLED", "true"); + vi.stubEnv("PAPERCLIP_CLOUD_STACK_ID", "fixture-stack"); + vi.stubEnv("PAPERCLIP_CUSTOMER_SUCCESS_CLOUD_ORIGIN", "https://cloud.example.test"); + vi.stubEnv("PAPERCLIP_CUSTOMER_SUCCESS_AUTHORITY_ENABLED", "true"); + const keys = generateKeyPairSync("ed25519"); + const jwks = { keys: [{ ...keys.publicKey.export({ format: "jwk" }), kid: "cloud" }] }; + vi.stubEnv("PAPERCLIP_CUSTOMER_SUCCESS_INSPECTION_JWKS", JSON.stringify(jwks)); + const query: InspectionQuery = { + operation: "get", + resource: "companies", + companyId, + resourceId: companyId, + }; + const now = Math.floor(Date.now() / 1000); + const claims: InspectionPermit = { + v: 1, + iss: "paperclip-cloud", + aud: INSPECTION_AUDIENCE, + sub: "fixture-stack", + iat: now, + exp: now + 60, + jti: randomUUID(), + bindingVersion: 1, + grantId: randomUUID(), + requestId: randomUUID(), + agentId, + keyId: "sha256:fixture", + runId, + query, + }; + function token(value = claims, privateKey = keys.privateKey) { + const header = Buffer.from( + JSON.stringify({ alg: "EdDSA", typ: INSPECTION_PERMIT_TYPE, kid: "cloud" }), + ).toString("base64url"); + const body = Buffer.from(JSON.stringify(value)).toString("base64url"); + const input = `${header}.${body}`; + return `${input}.${sign(null, Buffer.from(input), privateKey).toString("base64url")}`; + } + const fetcher = vi.fn().mockResolvedValue(Response.json({ consumed: true })); + vi.stubGlobal("fetch", fetcher); + const app = express(); + app.use(express.json()); + app.use("/api/customer-success/v1", customerSuccessRoutes(db, storage)); + app.use(errorHandler); + const response = await request(app) + .post("/api/customer-success/v1/read") + .set(INSPECTION_HEADER, token()) + .send(query); + expect(response.status).toEqual(200); + expect(response.headers["cache-control"]).toEqual("no-store"); + expect(fetcher).toHaveBeenCalledOnce(); + expect(fetcher.mock.calls[0][0]).toEqual( + "https://cloud.example.test/v1/customer-success/permits/consume", + ); + expect( + ( + await request(app) + .post("/api/customer-success/v1/read") + .set(INSPECTION_HEADER, token()) + .send({ ...query, resource: "agents" }) + ).status, + ).toEqual(403); + expect( + ( + await request(app) + .post("/api/customer-success/v1/read") + .set(INSPECTION_HEADER, token()) + .set("Cookie", "session=ignored") + .send(query) + ).status, + ).toEqual(403); + expect( + ( + await request(app) + .get("/api/customer-success/v1/run-authority") + .set("Authorization", `Bearer ${bearer}`) + ).status, + ).toEqual(200); + expect((await request(app).post("/api/customer-success/v1/unknown").send({})).status).toEqual( + 404, + ); + fetcher.mockResolvedValue(new Response("", { status: 503 })); + expect( + ( + await request(app) + .post("/api/customer-success/v1/read") + .set(INSPECTION_HEADER, token()) + .send(query) + ).status, + ).toEqual(403); + expect(() => + verifyInspectionPermit(token({ ...claims, sub: "other-stack" }), jwks, "fixture-stack"), + ).toThrow(); + expect(() => + verifyInspectionPermit(token({ ...claims, exp: now }), jwks, "fixture-stack"), + ).toThrow(); + expect(() => + verifyInspectionPermit(token({ ...claims, exp: now + 61 }), jwks, "fixture-stack"), + ).toThrow(); + expect(() => + verifyInspectionPermit( + token(claims, generateKeyPairSync("ed25519").privateKey), + jwks, + "fixture-stack", + ), + ).toThrow(); + }); + it("reads assigned repository URLs, snapshotted skill files, and workspace downloads with existing file restrictions", async () => { + const root = await mkdtemp(join(directory, "workspace-")); + await writeFile(join(root, "result.bin"), Buffer.alloc(2 * 1024 * 1024, 7)); + await writeFile(join(root, "empty.txt"), ""); + await writeFile(join(root, ".env"), "TOKEN=fixture"); + await symlink(join(root, "result.bin"), join(root, "link.bin")); + const [project] = await db + .insert(schema.projects) + .values({ companyId, name: "Onboarding repo" }) + .returning(); + const [workspace] = await db + .insert(schema.projectWorkspaces) + .values({ + companyId, + projectId: project.id, + name: "Local checkout", + cwd: root, + repoUrl: "https://github.com/example/onboarding", + }) + .returning(); + const [issue] = await db + .insert(schema.issues) + .values({ companyId, projectId: project.id, title: "Workspace outputs" }) + .returning(); + const base = { + companyId, + resourceId: issue.id, + context: { projectId: project.id, workspaceId: workspace.id }, + }; + const repo = (await readCustomerSuccessResource(db, storage, { + operation: "get", + resource: "projectWorkspaces", + companyId, + resourceId: workspace.id, + })) as { repoUrl: string }; + expect(repo.repoUrl).toEqual("https://github.com/example/onboarding"); + const before = await stat(join(root, "result.bin")); + const result = (await readCustomerSuccessResource(db, storage, { + ...base, + operation: "files.download", + path: "result.bin", + range: { start: 2, end: 7 }, + })) as { content: string; bytes: number }; + expect(Buffer.from(result.content, "base64")).toEqual(Buffer.alloc(6, 7)); + expect(result.bytes).toEqual(6); + expect((await stat(join(root, "result.bin"))).mtimeMs).toEqual(before.mtimeMs); + for (const path of [".env", "../result.bin", "link.bin"]) + await expect( + readCustomerSuccessResource(db, storage, { ...base, operation: "files.download", path }), + ).rejects.toThrow(); + await expect( + readCustomerSuccessResource(db, storage, { + ...base, + operation: "files.download", + path: "result.bin", + }), + ).rejects.toThrow("bounded byte range"); + const listed = await readCustomerSuccessResource(db, storage, { + ...base, + operation: "files.list", + limit: 5, + }); + expect(JSON.stringify(listed)).toContain("result.bin"); + expect( + await readCustomerSuccessResource(db, storage, { + ...base, + operation: "files.download", + path: "empty.txt", + }), + ).toMatchObject({ available: true, bytes: 0, content: "" }); + const [unassigned] = await db + .insert(schema.issues) + .values({ companyId, title: "No workspace" }) + .returning(); + expect( + await readCustomerSuccessResource(db, storage, { + companyId, + resourceId: unassigned.id, + operation: "files.read", + path: "result.txt", + }), + ).toMatchObject({ available: false }); + const [skill] = await db + .insert(schema.companySkills) + .values({ + companyId, + key: "example", + slug: "example", + name: "Example", + markdown: "# Example", + fileInventory: [{ path: "SKILL.md", kind: "markdown", sizeBytes: 9 }], + }) + .returning(); + const [version] = await db + .insert(schema.companySkillVersions) + .values({ + companyId, + companySkillId: skill.id, + revisionNumber: 1, + fileInventory: [{ path: "SKILL.md", kind: "markdown", sizeBytes: 9, content: "# Example" }], + }) + .returning(); + await db + .update(schema.companySkills) + .set({ currentVersionId: version.id }) + .where(eq(schema.companySkills.id, skill.id)); + const snapshot = (await readCustomerSuccessResource(db, storage, { + operation: "skills.file", + companyId, + resourceId: skill.id, + })) as { content: string }; + expect(Buffer.from(snapshot.content, "base64").toString()).toEqual("# Example"); + await expect( + readCustomerSuccessResource(db, storage, { + operation: "skills.file", + companyId, + resourceId: skill.id, + path: ".env", + }), + ).rejects.toThrow("denied by policy"); + const versions = (await readCustomerSuccessResource(db, storage, { + operation: "list", + resource: "skillVersions", + companyId, + parentId: skill.id, + })) as { items: unknown[] }; + expect(versions.items).toHaveLength(1); + }); + it("the operation contract has no writes, credentials or unknown parameters", () => { + for (const value of [ + { operation: "write", companyId }, + { operation: "list", resource: "secrets", companyId }, + { operation: "get", resource: "agents", companyId, resourceId: agentId, key: "replacement" }, + { + operation: "assets.content", + companyId, + resourceId: agentId, + range: { start: 0, end: 1048576 }, + }, + ]) + expect(() => parseInspectionQuery(value)).toThrow(); + }); +}); diff --git a/server/src/__tests__/openapi-routes.test.ts b/server/src/__tests__/openapi-routes.test.ts index 212e26eb5a..ac83f76c55 100644 --- a/server/src/__tests__/openapi-routes.test.ts +++ b/server/src/__tests__/openapi-routes.test.ts @@ -34,6 +34,7 @@ const apiPrefixes: Record = { "slack-tools.ts": "/api", "email.ts": "/api", "cloud.ts": "/api/cloud", + "customer-success.ts": "/api/customer-success/v1", "companies.ts": "/api/companies", "company-skills.ts": "/api", "company-skill-policy.ts": "/api", @@ -94,6 +95,11 @@ const HTTP_METHODS = new Set([ const explicitOpenApiCoverageExclusions = new Set(); const explicitOpenApiOperationCoverageExclusions = new Set([ + // Inspection uses its own versioned Cloud-permit/managed-run protocol, + // documented in CUSTOMER-SUCCESS-INSPECTION.md and the shared contract. + // Ordinary board sessions and agent API keys cannot invoke these endpoints. + "GET /api/customer-success/v1/run-authority", + "POST /api/customer-success/v1/read", // OAuth discovery and protocol endpoints have their own metadata contract; // browser connection-management operations remain documented in the board API. "GET /.well-known/oauth-authorization-server", diff --git a/server/src/agent-auth-jwt.ts b/server/src/agent-auth-jwt.ts index a4f42b9d77..dd55611585 100644 --- a/server/src/agent-auth-jwt.ts +++ b/server/src/agent-auth-jwt.ts @@ -160,7 +160,7 @@ export function createLocalAgentJwt( return `${signingInput}.${signature}`; } -export function verifyLocalAgentJwt(token: string): LocalAgentJwtClaims | null { +export function verifyLocalAgentJwt(token: string, options: { strictRunAuthority?: boolean } = {}): LocalAgentJwtClaims | null { if (!token) return null; const config = jwtConfig(); if (!config) return null; @@ -199,7 +199,7 @@ export function verifyLocalAgentJwt(token: string): LocalAgentJwtClaims | null { const perCompanyKey = deriveCompanySigningKey(config.secret, claimedCompanyId, config.instanceId); const perCompanySig = signPayload(perCompanyKey, signingInput); let signatureOk = safeCompare(signature, perCompanySig); - if (!signatureOk && !config.disableLegacyFallback) { + if (!signatureOk && !config.disableLegacyFallback && !options.strictRunAuthority) { const legacySig = signPayload(config.secret, signingInput); signatureOk = safeCompare(signature, legacySig); } @@ -237,6 +237,15 @@ export function verifyLocalAgentJwt(token: string): LocalAgentJwtClaims | null { // enforcement is conditional — matching how iss/aud are handled above. const instanceClaim = typeof claims.instance_id === "string" ? claims.instance_id : undefined; if (instanceClaim && instanceClaim !== config.instanceId) return null; + // Cross-instance inspection cannot inherit the compatibility exceptions of + // ordinary API authentication. Only a current, standard managed-run token + // issued by this instance can attest live authority. + if (options.strictRunAuthority && ( + header.typ !== "JWT" || issuer !== config.issuer || audience !== config.audience + || instanceClaim !== config.instanceId || exp <= now || iat > now + || !Number.isInteger(iat) || !Number.isInteger(exp) + || (Object.hasOwn(claims, "key_scope") && (!claims.key_scope || typeof claims.key_scope !== "object" || (claims.key_scope as { kind?: unknown }).kind !== "standard")) + )) return null; return { sub, diff --git a/server/src/app.ts b/server/src/app.ts index 4b70d76cd5..b9054137ce 100644 --- a/server/src/app.ts +++ b/server/src/app.ts @@ -1,3 +1,4 @@ +import { customerSuccessRoutes } from "./routes/customer-success.js"; import { cloudWarmStandbyMiddleware } from "./middleware/cloud-warm-standby.js"; import type { CloudWarmStandby } from "./services/cloud-warm-standby.js"; import { browserUseRoutes } from "./routes/browser-use.js"; @@ -583,6 +584,7 @@ export async function createApp( // A signed claim above commits identity before any normal request can seed // company data. Unclaimed probes bypass session resolution as well as SQL. app.use(cloudWarmStandbyMiddleware(isWarmStandby, health, staticUi)); + app.use("/api/customer-success/v1", customerSuccessRoutes(db, opts.storageService)); app.use(publicMcpIngress); // Connection-intent tools carry their own short-lived, run-bound bearer and // must be reachable by remote adapters that intentionally do not receive an diff --git a/server/src/middleware/http-log-redaction.ts b/server/src/middleware/http-log-redaction.ts index 0e97df2bda..21acd2a633 100644 --- a/server/src/middleware/http-log-redaction.ts +++ b/server/src/middleware/http-log-redaction.ts @@ -16,6 +16,7 @@ export const HTTP_LOG_REDACT_PATHS = [ 'req.headers["x-paperclip-cloud-session-id"]', 'req.headers["x-paperclip-cloud-runtime-identity"]', 'req.headers["x-paperclip-cloud-control"]', + 'req.headers["x-paperclip-cloud-inspection"]', // Runtime GitHub capabilities authorize credential acquisition for a live run. 'req.headers["x-paperclip-github-capability"]', // Telegram's optional webhook verification header is a reusable bearer diff --git a/server/src/routes/customer-success.ts b/server/src/routes/customer-success.ts new file mode 100644 index 0000000000..a28a0d0eed --- /dev/null +++ b/server/src/routes/customer-success.ts @@ -0,0 +1,136 @@ +import { createPublicKey, verify } from "node:crypto"; +import { Router } from "express"; +import type { Db } from "@paperclipai/db"; +import { + canonicalJson, + parseInspectionQuery, + INSPECTION_AUDIENCE, + INSPECTION_PERMIT_TYPE, + INSPECTION_HEADER, + type InspectionPermit, +} from "@paperclipai/shared"; +import type { StorageService } from "../storage/types.js"; +import { getCloudRuntimeIdentity } from "../services/cloud-runtime-identity.js"; +import { customerSuccessRunAuthority } from "../services/customer-success-authority.js"; +import { readCustomerSuccessResource } from "../services/customer-success-inspection.js"; +import { forbidden, unprocessable } from "../errors.js"; + +export function verifyInspectionPermit( + token: string, + keys: { keys: Array> }, + stackId: string, + now = Date.now(), +): InspectionPermit { + const parts = token.split("."); + if (parts.length !== 3 || token.length > 16384) throw forbidden("Invalid inspection permit"); + try { + const header = JSON.parse(Buffer.from(parts[0], "base64url").toString()); + const claims = JSON.parse(Buffer.from(parts[1], "base64url").toString()) as InspectionPermit; + const jwk = keys.keys.find( + (k) => k.kid === header.kid && k.kty === "OKP" && k.crv === "Ed25519" && !k.d, + ); + if ( + header.alg !== "EdDSA" || + header.typ !== INSPECTION_PERMIT_TYPE || + !jwk || + !verify( + null, + Buffer.from(`${parts[0]}.${parts[1]}`), + createPublicKey({ key: jwk, format: "jwk" }), + Buffer.from(parts[2], "base64url"), + ) + ) + throw new Error(); + const seconds = Math.floor(now / 1000); + if ( + claims.v !== 1 || + claims.iss !== "paperclip-cloud" || + claims.aud !== INSPECTION_AUDIENCE || + claims.sub !== stackId || + !Number.isInteger(claims.iat) || + !Number.isInteger(claims.exp) || + claims.exp <= seconds || + claims.iat > seconds + 5 || + claims.exp - claims.iat > 60 || + claims.exp <= claims.iat || + !claims.jti || + !claims.grantId || + !claims.requestId || + !claims.agentId || + !claims.keyId || + !claims.runId || + !Number.isSafeInteger(claims.bindingVersion) + ) + throw new Error(); + parseInspectionQuery(claims.query); + return claims; + } catch { + throw forbidden("Invalid inspection permit"); + } +} + +/** Mounted before actor middleware: none of these reads creates tenant state. */ +export function customerSuccessRoutes(db: Db, storage: StorageService) { + const router = Router(); + router.use((_req, res, next) => { + res.setHeader("Cache-Control", "no-store"); + next(); + }); + router.get("/run-authority", async (req, res) => { + if (process.env.PAPERCLIP_CUSTOMER_SUCCESS_AUTHORITY_ENABLED !== "true") + throw forbidden("Run authority is disabled"); + if (req.headers.cookie) throw forbidden("Browser credentials are not run authority"); + const match = req.headers.authorization?.match(/^Bearer ([A-Za-z0-9_.-]+)$/); + if (!match) throw forbidden("Managed run bearer required"); + res.json(await customerSuccessRunAuthority(db, match[1])); + }); + router.post("/read", async (req, res) => { + if (process.env.PAPERCLIP_CUSTOMER_SUCCESS_INSPECTION_ENABLED !== "true") + throw forbidden("Inspection is disabled"); + if (req.headers.authorization || req.headers.cookie) + throw forbidden("Only a Cloud inspection permit is accepted"); + const stackId = getCloudRuntimeIdentity()?.stackId ?? process.env.PAPERCLIP_CLOUD_STACK_ID; + const rawKeys = process.env.PAPERCLIP_CUSTOMER_SUCCESS_INSPECTION_JWKS; + const origin = process.env.PAPERCLIP_CUSTOMER_SUCCESS_CLOUD_ORIGIN; + const token = req.get(INSPECTION_HEADER); + if (!stackId || !rawKeys || !origin || !token) + throw forbidden("Inspection trust is unavailable"); + const url = new URL(origin); + if ( + (url.protocol !== "https:" && + !( + url.protocol === "http:" && ["127.0.0.1", "localhost", "[::1]"].includes(url.hostname) + )) || + url.origin !== origin + ) + throw forbidden("Inspection trust is invalid"); + const permit = verifyInspectionPermit(token, JSON.parse(rawKeys), stackId); + let query; + try { + query = parseInspectionQuery(req.body); + } catch { + throw unprocessable("Invalid inspection query"); + } + if (canonicalJson(query) !== canonicalJson(permit.query)) + throw forbidden("Inspection parameters do not match permit"); + // Replay fencing and all inspection state live in Cloud. No tenant receipt. + const consumed = await fetch(`${origin}/v1/customer-success/permits/consume`, { + method: "POST", + headers: { "Content-Type": "application/json", [INSPECTION_HEADER]: token }, + body: "{}", + redirect: "error", + signal: AbortSignal.timeout(10000), + }); + if (!consumed.ok) throw forbidden("Inspection permit could not be consumed"); + const result = await readCustomerSuccessResource(db, storage, query); + const body = JSON.stringify(result); + if (Buffer.byteLength(body) > 2 * 1024 * 1024) + throw unprocessable("Inspection response is too large; use pagination or a byte range"); + res.type("application/json").send(body); + }); + // Unknown methods/paths terminate here rather than entering normal auth. + router.use((_req, res) => { + res.status(404).json({ error: "Unsupported inspection endpoint" }); + }); + return router; +} diff --git a/server/src/services/customer-success-authority.ts b/server/src/services/customer-success-authority.ts new file mode 100644 index 0000000000..2081902ae0 --- /dev/null +++ b/server/src/services/customer-success-authority.ts @@ -0,0 +1,56 @@ +import { and, eq } from "drizzle-orm"; +import { agents, heartbeatRuns, type Db } from "@paperclipai/db"; +import type { RunAuthority } from "@paperclipai/shared"; +import { verifyLocalAgentJwt } from "../agent-auth-jwt.js"; +import { agentRunWritesRevoked } from "../agent-run-cancellation.js"; +import { forbidden } from "../errors.js"; +import { agentIdentityService, supportsManagedAgentIdentity } from "./agent-identity.js"; + +/** No actor/session resolution, key provisioning, or run identity capture. */ +export async function customerSuccessRunAuthority(db: Db, bearer: string): Promise { + const claims = verifyLocalAgentJwt(bearer, { strictRunAuthority: true }); + if (!claims) throw forbidden("A current managed-run credential is required"); + return db.transaction( + async (tx) => { + const [agent] = await tx + .select() + .from(agents) + .where(and(eq(agents.id, claims.sub), eq(agents.companyId, claims.company_id))); + const [run] = await tx + .select() + .from(heartbeatRuns) + .where( + and( + eq(heartbeatRuns.id, claims.run_id), + eq(heartbeatRuns.agentId, claims.sub), + eq(heartbeatRuns.companyId, claims.company_id), + ), + ); + if ( + !agent || + !supportsManagedAgentIdentity(agent.adapterType, agent.adapterConfig, run?.driverKind) || + !["idle", "running", "active"].includes(agent.status) || + !run || + run.status !== "running" || + run.finishedAt || + agentRunWritesRevoked(run) + ) + throw forbidden("The managed run is not active"); + const identity = await agentIdentityService(tx as unknown as Db).getPublicIdentity( + agent.companyId, + agent.id, + ); + if (!identity) throw forbidden("The agent identity is unprovisioned"); + return { + version: 1, + instanceId: claims.instance_id!, + companyId: agent.companyId, + agentId: agent.id, + runId: run.id, + keyId: identity.keyId, + active: true, + }; + }, + { accessMode: "read only" }, + ); +} diff --git a/server/src/services/customer-success-inspection.ts b/server/src/services/customer-success-inspection.ts new file mode 100644 index 0000000000..2f9a031e82 --- /dev/null +++ b/server/src/services/customer-success-inspection.ts @@ -0,0 +1,422 @@ +import fs from "node:fs/promises"; +import { constants } from "node:fs"; +import { and, eq, getTableColumns, type SQL } from "drizzle-orm"; +import type { AnyPgTable, AnyPgColumn } from "drizzle-orm/pg-core"; +import * as schema from "@paperclipai/db"; +import type { Db } from "@paperclipai/db"; +import { + MAX_INSPECTION_BYTES, + type InspectionQuery, + type InspectionResource, +} from "@paperclipai/shared"; +import type { StorageService } from "../storage/types.js"; +import { HttpError, notFound, unprocessable } from "../errors.js"; +import { redactAgentAdapterConfig, redactEventPayload } from "../redaction.js"; +import { createRunSecretRedactionRegistry } from "./run-secret-redaction.js"; +import { + assertWorkspaceFilePathAllowed, + workspaceFileResourceService, +} from "./workspace-file-resources.js"; +import { deriveBundleState } from "./agent-instructions.js"; +import { readInstructionBytes } from "./agent-instruction-files.js"; +import { getRunLogStore } from "./run-log-store.js"; + +// This is a reviewed catalog of readers, not a proxy for existing GET routes. +// Select credentials out of queries; never load encrypted identity records. +interface Reader { + table: AnyPgTable; + omit?: string[]; + parent?: string; + id?: string; +} +export const inspectionReaders: Record = { + companies: { table: schema.companies }, + goals: { table: schema.goals }, + users: { table: schema.authUsers }, + memberships: { table: schema.companyMemberships }, + permissions: { table: schema.principalPermissionGrants }, + agents: { table: schema.agents }, + agentIdentities: { table: schema.agentIdentityKeys, id: "agentId", omit: ["privateKeyMaterial"] }, + agentInstructions: { table: schema.agentInstructionRevisions, parent: "agentId" }, + agentConfigRevisions: { table: schema.agentConfigRevisions, parent: "agentId" }, + tasks: { table: schema.issues }, + comments: { table: schema.issueComments, parent: "issueId" }, + documents: { table: schema.documents }, + documentRevisions: { table: schema.documentRevisions, parent: "documentId" }, + taskDocuments: { table: schema.issueDocuments, parent: "issueId" }, + interactions: { table: schema.issueThreadInteractions, parent: "issueId" }, + approvals: { table: schema.approvals }, + approvalComments: { table: schema.approvalComments, parent: "approvalId" }, + decisions: { table: schema.decisions }, + routines: { table: schema.routines }, + routineTriggers: { + table: schema.routineTriggers, + parent: "routineId", + omit: ["webhookSecretHash", "webhookSecret"], + }, + routineRuns: { table: schema.routineRuns, parent: "routineId" }, + routineRevisions: { table: schema.routineRevisions, parent: "routineId" }, + routineDocuments: { table: schema.routineDocuments, parent: "routineId" }, + projects: { table: schema.projects }, + projectWorkspaces: { table: schema.projectWorkspaces, parent: "projectId" }, + projectMemberships: { table: schema.projectMemberships, parent: "projectId" }, + skills: { table: schema.companySkills, omit: ["publicShareToken"] }, + skillVersions: { table: schema.companySkillVersions, parent: "companySkillId" }, + skillSources: { table: schema.companySkillSources, omit: ["leaseToken"] }, + runs: { table: schema.heartbeatRuns, omit: ["logRef"] }, + runEvents: { table: schema.heartbeatRunEvents, parent: "runId" }, + traceMetadata: { table: schema.providerTraceRecords, parent: "runId", omit: ["traceRef"] }, + activity: { table: schema.activityLog }, + costs: { table: schema.costEvents }, + workProducts: { table: schema.issueWorkProducts, parent: "issueId" }, + artifacts: { table: schema.assets, omit: ["objectKey"] }, + attachments: { table: schema.issueAttachments, parent: "issueId" }, + connections: { + table: schema.toolConnections, + omit: ["externalCredential", "credentialRefs", "credentialSecretRefs"], + }, + connectionInstalls: { table: schema.toolConnectionInstalls, parent: "connectionId" }, + connectionGrants: { + table: schema.connectionGrants, + parent: "connectionId", + omit: ["credentialSecretRefs", "externalCredential"], + }, + connectionGrantMembers: { table: schema.connectionGrantMembers, parent: "grantId" }, + connectionCatalog: { table: schema.toolCatalogEntries, parent: "connectionId" }, + toolProfiles: { table: schema.toolProfiles }, + toolProfileEntries: { table: schema.toolProfileEntries, parent: "profileId" }, + agentMemberships: { table: schema.agentMemberships, parent: "agentId" }, + agentCommentary: { table: schema.agentCommentary, parent: "runId" }, + skillPolicies: { table: schema.companySkillPolicies, id: "companyId" }, + projectGoals: { table: schema.projectGoals, parent: "projectId" }, + labels: { table: schema.labels }, + taskLabels: { table: schema.issueLabels, parent: "issueId" }, + taskRelations: { table: schema.issueRelations }, + taskApprovals: { table: schema.issueApprovals, parent: "issueId" }, + documentMemberships: { table: schema.documentMemberships, parent: "documentId" }, + documentThreads: { table: schema.documentAnnotationThreads, parent: "documentId" }, + documentComments: { table: schema.documentAnnotationComments, parent: "threadId" }, + threads: { table: schema.chatConversations, parent: "endpointId" }, +}; + +async function ownRow(tx: Db, table: AnyPgTable, companyId: string, id: string) { + const columns = getTableColumns(table) as Record; + const [row] = await tx + .select() + .from(table) + .where(and(eq(columns.companyId, companyId), eq(columns.id, id))) + .limit(1); + if (!row) throw notFound("Inspection resource not found"); + return row as Record; +} +function unavailable(reason: string) { + return { available: false, reason }; +} +function binary( + bytes: Buffer, + totalBytes: number, + start = 0, + contentType = "application/octet-stream", +) { + return { + available: true, + encoding: "base64", + contentType, + totalBytes, + start, + end: start + bytes.length - 1, + bytes: bytes.length, + content: bytes.toString("base64"), + }; +} + +/** All SQL, including nested serializers, runs in one read-only transaction. */ +export async function readCustomerSuccessResource( + db: Db, + storage: StorageService, + q: InspectionQuery, +): Promise { + return db.transaction( + async (rawTx) => { + const tx = rawTx as unknown as Db; + if (q.companyId) { + const [company] = await tx + .select({ id: schema.companies.id }) + .from(schema.companies) + .where(eq(schema.companies.id, q.companyId)); + if (!company) throw notFound("Inspection company not found"); + } + if (q.operation === "list" || q.operation === "get") { + const reader = inspectionReaders[q.resource!]; + const all = getTableColumns(reader.table) as Record; + const columns = Object.fromEntries( + Object.entries(all).filter(([key]) => !reader.omit?.includes(key)), + ); + const predicates: SQL[] = []; + if (q.resource === "companies") { + if (q.companyId) predicates.push(eq(all.id, q.companyId)); + } else if (q.resource === "users") { + // Directory only, never the authentication account/session tables. + const memberships = await tx + .select({ userId: schema.companyMemberships.principalId }) + .from(schema.companyMemberships) + .where( + and( + eq(schema.companyMemberships.companyId, q.companyId!), + eq(schema.companyMemberships.principalType, "user"), + ), + ); + const { inArray } = await import("drizzle-orm"); + if (!memberships.length) { + if (q.operation === "get") throw notFound("Inspection resource not found"); + return { items: [], nextOffset: null }; + } + predicates.push( + inArray( + all.id, + memberships.map((m) => m.userId), + ), + ); + } else { + if (!all.companyId) throw unprocessable("Reader has no company boundary"); + predicates.push(eq(all.companyId, q.companyId!)); + } + if (q.operation === "get" && !all[reader.id ?? "id"]) + throw unprocessable("This resource supports listing only"); + if (q.operation === "get") predicates.push(eq(all[reader.id ?? "id"], q.resourceId!)); + if (q.parentId) { + if (!reader.parent || !all[reader.parent]) + throw unprocessable("This reader does not accept parentId"); + predicates.push(eq(all[reader.parent], q.parentId)); + } + const limit = q.operation === "get" ? 1 : Math.max(1, q.limit ?? 25); + const order = all[reader.id ?? "id"] ?? all.createdAt ?? all.companyId; + let rows = await tx + .select(columns) + .from(reader.table) + .where(and(...predicates)) + .orderBy(order) + .limit(limit + (q.operation === "get" ? 0 : 1)) + .offset(q.offset ?? 0); + if (q.resource === "agents") + rows = rows.map((r) => ({ + ...r, + adapterConfig: redactAgentAdapterConfig(r.adapterConfig as Record), + runtimeConfig: redactEventPayload(r.runtimeConfig as Record), + })); + if (q.resource === "agentConfigRevisions") + rows = rows.map((r) => ({ + ...r, + beforeConfig: redactRevisionSnapshot(r.beforeConfig), + afterConfig: redactRevisionSnapshot(r.afterConfig), + })); + if (q.resource === "projects" || q.resource === "routines") + rows = rows.map((r) => ({ ...r, env: redactAgentAdapterConfig({ env: r.env }).env })); + if (q.resource === "routineRevisions") + rows = rows.map((r) => { + const snapshot = r.snapshot as Record | null; + const routine = snapshot?.routine as Record | undefined; + return routine + ? { + ...r, + snapshot: { + ...snapshot, + routine: { + ...routine, + env: redactAgentAdapterConfig({ env: routine.env }).env, + }, + }, + } + : r; + }); + if (q.resource === "runs") + rows = await createRunSecretRedactionRegistry(tx).redactForRuns( + q.companyId!, + rows as Array<{ id: string }>, + ); + // Match the existing structured event/config redaction, without scanning + // prose, documents, instructions, skill contents, or workspace files. + if ( + [ + "runEvents", + "activity", + "connections", + "connectionGrants", + "routineTriggers", + "routineRevisions", + ].includes(q.resource!) + ) + rows = rows.map((r) => redactEventPayload(r) ?? {}); + if (q.resource === "runEvents") { + const redactions = createRunSecretRedactionRegistry(tx); + rows = await Promise.all( + rows.map((r) => redactions.redactForRun(q.companyId!, String(r.runId), r)), + ); + } + if (q.operation === "get") { + if (!rows.length) throw notFound("Inspection resource not found"); + return rows[0]; + } + return { + items: rows.slice(0, limit), + nextOffset: rows.length > limit ? (q.offset ?? 0) + limit : null, + }; + } + if (q.operation === "instructions.file") { + const agent = await ownRow(tx, schema.agents, q.companyId!, q.resourceId!); + const state = deriveBundleState( + agent as unknown as Parameters[0], + ); + if (!state.rootPath) return unavailable("instructions_not_configured"); + const relative = q.path ?? state.entryFile; + assertWorkspaceFilePathAllowed(relative); + const bytes = await readInstructionBytes(state.rootPath, relative); + return bytes + ? binary(bytes, bytes.length, 0, "text/plain; charset=utf-8") + : unavailable("instructions_file_missing"); + } + if (q.operation === "skills.file") { + const skill = await ownRow(tx, schema.companySkills, q.companyId!, q.resourceId!); + const relative = q.path ?? "SKILL.md"; + assertWorkspaceFilePathAllowed(relative); + const versionId = q.versionId ?? skill.currentVersionId; + // Installed version snapshots are authoritative. Reading must not invoke + // inventory reconciliation, checkout, fetching, adoption or materialization. + if (versionId) { + const version = await ownRow( + tx, + schema.companySkillVersions, + q.companyId!, + String(versionId), + ); + if (version.companySkillId !== skill.id) throw notFound("Skill version not found"); + const files = version.fileInventory as Array<{ + path: string; + content?: string; + encoding?: string; + }>; + const file = files.find((f) => f.path === relative); + if (file?.content !== undefined) { + const bytes = Buffer.from(file.content, file.encoding === "base64" ? "base64" : "utf8"); + if (bytes.length > MAX_INSPECTION_BYTES) + throw unprocessable("Skill snapshot exceeds inspection limit"); + return binary(bytes, bytes.length); + } + } + if (relative === "SKILL.md") + return binary( + Buffer.from(String(skill.markdown), "utf8"), + Buffer.byteLength(String(skill.markdown)), + ); + return unavailable("skill_file_not_snapshotted"); + } + if (q.operation === "runs.log") { + const run = await ownRow(tx, schema.heartbeatRuns, q.companyId!, q.resourceId!); + if (!run.logRef || run.logStore !== "local_file") return unavailable("run_log_unavailable"); + const result = await getRunLogStore().read( + { store: "local_file", logRef: String(run.logRef) }, + { + offset: q.range?.start ?? q.offset ?? 0, + limitBytes: q.range ? q.range.end - q.range.start + 1 : 256000, + }, + ); + return createRunSecretRedactionRegistry(tx).redactForRun( + q.companyId!, + q.resourceId!, + result, + ); + } + if (q.operation.startsWith("files.")) { + try { + const issue = await ownRow(tx, schema.issues, q.companyId!, q.resourceId!); + const files = workspaceFileResourceService(tx); + const input = { path: q.path ?? "", ...q.context, limit: q.limit, offset: q.offset }; + const opts = { + issue: issue as unknown as NonNullable[2]>["issue"], + }; + if (q.operation === "files.list") return await files.list(q.resourceId!, input, opts); + if (q.operation === "files.read") + return await files.readContent(q.resourceId!, input, opts); + const resolved = await files.prepareDownload(q.resourceId!, input, opts); + const file = await fs.open(resolved.realPath, constants.O_RDONLY | constants.O_NOFOLLOW); + try { + const stat = await file.stat(); + if (!stat.isFile()) throw unprocessable("Regular file required"); + if (stat.size === 0 && !q.range) return binary(Buffer.alloc(0), 0); + const start = q.range?.start ?? 0; + const end = q.range?.end ?? stat.size - 1; + if (start >= stat.size || end >= stat.size || end - start + 1 > MAX_INSPECTION_BYTES) + throw unprocessable("Request a bounded byte range"); + const bytes = Buffer.alloc(end - start + 1); + const result = await file.read(bytes, 0, bytes.length, start); + const after = await file.stat(); + if ( + after.size !== stat.size || + after.mtimeMs !== stat.mtimeMs || + after.ino !== stat.ino + ) + return unavailable("file_changed_during_read"); + if (result.bytesRead !== bytes.length) return unavailable("file_changed_during_read"); + return binary(bytes, stat.size, start); + } finally { + await file.close(); + } + } catch (error) { + const code = + error instanceof HttpError + ? (error.details as { code?: string } | undefined)?.code + : undefined; + if (code && ["remote_workspace", "no_workspace", "no_local_workspace"].includes(code)) + return unavailable(code); + throw error; + } + } + if (q.operation === "assets.content") { + const asset = await ownRow(tx, schema.assets, q.companyId!, q.resourceId!); + const range = q.range; + if (range && (range.start >= Number(asset.byteSize) || range.end >= Number(asset.byteSize))) + throw unprocessable("Asset byte range is outside the resource"); + if (!range && Number(asset.byteSize) > MAX_INSPECTION_BYTES) + throw unprocessable("Request a bounded byte range"); + const object = await storage.getObject(q.companyId!, String(asset.objectKey), { range }); + const chunks: Buffer[] = []; + let length = 0; + try { + for await (const chunk of object.stream) { + const bytes = Buffer.from(chunk); + length += bytes.length; + if (length > MAX_INSPECTION_BYTES) + throw unprocessable("Asset exceeds inspection limit"); + chunks.push(bytes); + } + } finally { + object.stream.destroy(); + } + return binary( + Buffer.concat(chunks), + Number(asset.byteSize), + q.range?.start ?? 0, + object.contentType, + ); + } + throw unprocessable("Unsupported inspection operation"); + }, + { accessMode: "read only", isolationLevel: "repeatable read" }, + ); +} + +function redactRevisionSnapshot(snapshot: unknown): Record { + if (!snapshot || typeof snapshot !== "object" || Array.isArray(snapshot)) return {}; + const record = snapshot as Record; + const object = (value: unknown) => + value && typeof value === "object" ? (value as Record) : {}; + return { + ...record, + adapterConfig: redactAgentAdapterConfig(object(record.adapterConfig)), + runtimeConfig: redactEventPayload(object(record.runtimeConfig)), + metadata: + record.metadata && typeof record.metadata === "object" + ? redactEventPayload(object(record.metadata)) + : (record.metadata ?? null), + }; +} diff --git a/server/src/services/workspace-file-resources.ts b/server/src/services/workspace-file-resources.ts index 5f11afa72d..90b8534ce6 100644 --- a/server/src/services/workspace-file-resources.ts +++ b/server/src/services/workspace-file-resources.ts @@ -333,6 +333,11 @@ function throwIfDenied(segments: string[]) { } } +/** Apply the existing file-path policy to other read-only file surfaces. */ +export function assertWorkspaceFilePathAllowed(relativePath: string): void { + throwIfDenied(normalizeWorkspaceRelativePath(relativePath).segments); +} + function shouldPruneSegments(segments: string[]) { return denyReasonForPathSegments(segments) != null; }