From a9a20fb5c6c080884ced7ee0a4e3b04f9cbf3a6e Mon Sep 17 00:00:00 2001 From: Dotta <34892728+cryppadotta@users.noreply.github.com> Date: Tue, 6 Oct 2026 21:38:15 -0500 Subject: [PATCH] feat(security): add read-only customer-success inspection APIs (#15405) ## Thinking Path > - Paperclip is the open source app people use to manage AI agents for work. > - Agents need a persistent identity and a verified active run for governed access. > - Customer-success inspection needs broad reads without tenant writes or secret access. > - Ordinary board login and database credentials give more authority than this task needs. > - This pull request adds a dedicated inspection API and strict managed-run authority. > - Cloud owns short grants, human approval, replay protection, and audit records. > - The benefit is inspectable access that an operator can disable immediately. ## Linked Issues or Issue Description Refs: #15352. This change reuses the persistent Ed25519 identity from that PR. **Subsystem affected** Server authentication, pure resource readers, and the shared wire contract. **Problem or motivation** One internal Paperclip agent must inspect customer onboarding work. It must not receive owner login or database credentials. Reads must not create customer sessions, memberships, activity, or read receipts. **Proposed solution** Add disabled-by-default run authority and versioned tenant inspection endpoints. Require strict instance-bound managed-run JWTs on the home instance. Require exact-operation, single-use Cloud permits on tenants. Execute a reviewed company-scoped catalog in read-only transactions. Cloud applies seven-day stack-age eligibility and human exceptions. **Roadmap alignment** This is access support for Cloud deployments and governed agent identities. Bot creation, scheduling, scoring, and reports are separate work. The maintainer requested this implementation. ## What Changed - Reuse existing public identity reads and managed private-key injection. Reject unprovisioned keys, paused agents, ended runs, legacy signatures, and wrong instances. - Mount `/api/customer-success/v1` before actor/session synchronization. Verify Cloud permits and consume them centrally before reading. - Add explicit company-scoped database readers and bounded instruction, skill snapshot, run log, workspace, and asset reads. Preserve existing redactions and file protections. - Add protocol, security, database immutability, and managed-agent qualification tests. Add deployment and rollback documentation. ## Verification - Full `pnpm -r typecheck` and `pnpm build` passed. Server typecheck passed after review fixes. - The broad local `pnpm test:run` recorded 14,277 passes and four failures in unchanged suites: two timeouts and two PR-metadata mock assertions. All three affected suites passed on isolated reruns (36 tests). The complete CI matrix passes at the final head, including every test lane, typecheck, build, runner checks, canary dry run, and the security scan. - Focused inspection, JWT, and existing identity tests pass. The catalog test compares every public database table before and after reads. - Inspection and route-contract tests: 22 passed. The coordinated test runs a real managed process agent against separate home/customer PostgreSQL databases and a PostgreSQL broker over HTTP. It proves wake through the existing controller, bounded binary file reads, single challenge consumption across replicas, concurrent grants with a two-connection pool, scoped SQL audits, append-only runtime auditing, one-year retention, and unchanged tenant data/files. - Run the coordinated test with `PAPERCLIP_INSPECTION_CLOUD_DIST` pointing at the sibling Cloud build. Normal unit runs skip that optional private integration. - Final-head Greptile is 5/5 with no unresolved findings. - No production deployment or customer inspection occurred. ## Risks - This adds an authentication boundary. Keep both feature flags disabled until coordinated staging and canary qualification. - Cloud support must deploy after this API. Unsupported tenants fail closed. There is no owner-login or database fallback. - Existing redactions remain the content boundary. Arbitrary pasted secrets in readable prose or files may remain. - Remote files and suppressed provider traces remain unavailable. Wake can cause normal startup/background writes; test those separately. - Disable Cloud policy first during rollback. Preserve existing identity material and Cloud audit history. ## Model Used OpenAI GPT-6 (Codex), with reasoning, code execution, and browser testing. The session does not expose a more specific deployment ID or context-window size. ## Checklist - [x] I have included a thinking path that traces from project context to this change - [x] I have specified the model used (with version and capability details) - [x] I have checked ROADMAP.md and confirmed this PR does not duplicate planned core work - [x] I have searched GitHub for duplicate or related PRs and linked them above - [x] I have either (a) linked existing issues with `Fixes: #` / `Closes #` / `Refs #` OR (b) described the issue in-PR following the relevant issue template - [x] I have not referenced internal/instance-local Paperclip issues or links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip` URLs) - [x] My branch name describes the change (e.g. `docs/...`, `fix/...`) and contains no internal Paperclip ticket id or instance-derived details - [x] I have run tests locally and they pass (focused checks and isolated reruns; broad-run flakes are documented above) - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation to reflect my changes - [x] I have considered and documented any risks above - [x] All Paperclip CI gates are green - [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups - [x] I will address all Greptile and reviewer comments before requesting merge --------- Co-authored-by: Paperclip --- doc/AGENT-IDENTITY.md | 4 + doc/CUSTOMER-SUCCESS-INSPECTION.md | 85 +++ .../2026-10-06-customer-success-inspection.md | 48 ++ packages/shared/src/customer-success.ts | 221 +++++++ packages/shared/src/index.ts | 1 + .../__tests__/customer-success-cloud.test.ts | 407 ++++++++++++ server/src/__tests__/customer-success.test.ts | 605 ++++++++++++++++++ server/src/__tests__/openapi-routes.test.ts | 6 + server/src/agent-auth-jwt.ts | 13 +- server/src/app.ts | 2 + server/src/middleware/http-log-redaction.ts | 1 + server/src/routes/customer-success.ts | 136 ++++ .../services/customer-success-authority.ts | 56 ++ .../services/customer-success-inspection.ts | 422 ++++++++++++ .../src/services/workspace-file-resources.ts | 5 + 15 files changed, 2010 insertions(+), 2 deletions(-) create mode 100644 doc/CUSTOMER-SUCCESS-INSPECTION.md create mode 100644 doc/plans/2026-10-06-customer-success-inspection.md create mode 100644 packages/shared/src/customer-success.ts create mode 100644 server/src/__tests__/customer-success-cloud.test.ts create mode 100644 server/src/__tests__/customer-success.test.ts create mode 100644 server/src/routes/customer-success.ts create mode 100644 server/src/services/customer-success-authority.ts create mode 100644 server/src/services/customer-success-inspection.ts 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; }