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 <noreply@paperclip.ing>
This commit is contained in:
DottaandPaperclip authored and GitHub committed 2026-10-06 21:38:15 -05:00
1 parent caf120105c
commit a9a20fb5c6
15 files changed
+2010 -2

No files matched your search

+4
View File
@@ -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.
+85
View File
@@ -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.
@@ -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.
+221
View File
@@ -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;
}
+1
View File
@@ -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";
@@ -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<typeof createServer>[] = [];
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<typeof createServer>[0]) => {
const server = createServer(handler);
servers.push(server);
await new Promise<void>((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<string, unknown>; 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<Buffer>, 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<void>((resolve) => server.close(() => resolve()));
await pool.end();
await tenant.cleanup();
await home.cleanup();
vi.unstubAllEnvs();
await rm(directory, { recursive: true, force: true });
}
}, 60000);
});
@@ -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<ReturnType<typeof startEmbeddedPostgresTestDatabase>>;
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<string, unknown> = {};
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<string, unknown>;
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<string, unknown>; 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();
});
});
@@ -34,6 +34,7 @@ const apiPrefixes: Record<string, string> = {
"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<string>();
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",
+11 -2
View File
@@ -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,
+2
View File
@@ -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
@@ -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
+136
View File
@@ -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<Record<string, unknown>> },
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;
}
@@ -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<RunAuthority> {
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" },
);
}
@@ -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<InspectionResource, Reader> = {
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<string, AnyPgColumn>;
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<string, unknown>;
}
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<unknown> {
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<string, AnyPgColumn>;
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<string, unknown>),
runtimeConfig: redactEventPayload(r.runtimeConfig as Record<string, unknown>),
}));
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<string, unknown> | null;
const routine = snapshot?.routine as Record<string, unknown> | 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<typeof deriveBundleState>[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<Parameters<typeof files.list>[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<string, unknown> {
if (!snapshot || typeof snapshot !== "object" || Array.isArray(snapshot)) return {};
const record = snapshot as Record<string, unknown>;
const object = (value: unknown) =>
value && typeof value === "object" ? (value as Record<string, unknown>) : {};
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),
};
}
@@ -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;
}