Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
3797c63
fix(executor): enforce workspace isolation for every tool call and sc…
hutusi Jul 17, 2026
e97263a
fix(worker): contain artifact ingestion against symlink escapes
hutusi Jul 17, 2026
6749bcb
fix(authz): pin an authorized run manifest and scope repo connections…
hutusi Jul 17, 2026
173d673
fix(engine): atomic run transitions, DB-allocated event seq, recovera…
hutusi Jul 17, 2026
2facd4b
fix(engine): recover crashed no-retry steps and resume their session
hutusi Jul 17, 2026
b514faa
fix(engine): align quota window and stop double-counting run spend on…
hutusi Jul 17, 2026
9f001a1
fix(artifacts): enforce the step output contract on ingestion
hutusi Jul 17, 2026
ad48700
fix(api): close the SSE replay/subscribe gap by subscribing first
hutusi Jul 17, 2026
dbc379f
docs: record the security & correctness hardening
hutusi Jul 17, 2026
300df2c
fix: address review — atomic approvals, safe sweeper, tighter lifecyc…
hutusi Jul 17, 2026
28c319e
fix(security): address review — symlink writes, env allowlist, /app o…
hutusi Jul 17, 2026
775db9d
docs: correct the enqueue window, SSE polling, and optional-resource …
hutusi Jul 17, 2026
f8985a2
fix(security): confine reads + redact event secrets; make seq/finaliz…
hutusi Jul 17, 2026
8eeed20
fix: attempt-safe/binary artifacts, requires.skills, per-phase budget…
hutusi Jul 17, 2026
509a5c7
fix: safe artifact-attempt cleanup + unified run-lifecycle finalization
hutusi Jul 17, 2026
780cd17
fix: symmetric resource resolution, contiguous SSE cursor, split comp…
hutusi Jul 17, 2026
1d444f6
fix(security): allow-list the agent subprocess env; harden the manife…
hutusi Jul 17, 2026
da8f08b
fix: guard the artifact-size env and mount the API artifact store rea…
hutusi Jul 17, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 6 additions & 3 deletions ARCHITECTURE.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
# Architecture

This document is the orientation map for contributors. The authoritative design lives in `docs/design/` (10 documents) and `docs/adr/` (8 decision records); this file tells you what exists, where it lives, and which invariants hold it together.
This document is the orientation map for contributors. The authoritative design lives in `docs/design/` (10 documents) and `docs/adr/` (9 decision records); this file tells you what exists, where it lives, and which invariants hold it together.

## Bird's-eye view

Expand Down Expand Up @@ -35,7 +35,7 @@ Dependency direction is enforced by `scripts/check-deps.ts` (runtime deps only).

## Key flows

**Submit → run.** `POST /projects/:id/tasks` validates params against the compiled input schema (the same schema the SPA rendered the form from), verifies skill/MCP grants, resolves model roles → concrete granted models (frozen into `runs.model_resolution`), checks hard-stop quota headroom, then inserts task+run in one transaction and sends a pg-boss job keyed by run id. A worker sweeper re-enqueues stragglers, so a run can never be stranded nor double-queued.
**Submit → run.** `POST /projects/:id/tasks` validates params against the compiled input schema (the same schema the SPA rendered the form from), verifies each `repoRef` is owned by the project, resolves model roles → concrete granted models (frozen into `runs.model_resolution`), pins the authorized skills/MCP into `runs.resource_manifest` (required grants enforced, optional included only when granted), checks hard-stop quota headroom, then inserts task+run in one transaction and sends a pg-boss job keyed by run id. A worker sweeper re-enqueues stragglers (and runs left paused by a lost approval-resume), so a run can never be stranded nor double-queued.

**Engine loop.** The engine (worker-side) walks phases/steps: `when:` conditions, optional-resource skips, per-step retries, budget checks at every boundary. Agent steps stream normalized executor events which the engine persists to `run_events` (per-run monotonic `seq`), mirrors to the bus, and projects into `run_steps`/`token_usage`/`artifacts`. At the end it enforces the artifact contract — *succeeded* always means the contracted outputs exist.

Expand All @@ -52,7 +52,10 @@ Dependency direction is enforced by `scripts/check-deps.ts` (runtime deps only).
3. **Steps are the idempotency unit** — restart-safe by template rule, resumable by session id where the executor supports it.
4. **Usage rows are keyed `(run, step, attempt)`** so retries re-incur cost without ever double-counting.
5. **Executors are stateless I/O**: all inputs in the request, all outputs as events; they never touch the database.
6. **Secrets never leave as plaintext**: encrypted at rest (AES-256-GCM), write-only in the API, scrubbed from git remotes before agent code runs.
6. **Secrets never leave as plaintext**: encrypted at rest (AES-256-GCM), write-only in the API, scrubbed from git remotes before agent code runs, stripped from the agent subprocess environment (the master key and datastore URLs never reach a tool call), and redacted from event payloads before they persist or stream.
7. **Containment goes through one seam**: every tool call (reads and writes confined to the workspace) and the subprocess env are decided by `packages/executor-core/isolation.ts`; the adapter never reimplements it (ADR-0009). OS-level isolation between runs and keeping the provider key out of the subprocess are deferred to the container layer.
8. **The worker trusts only the pinned manifest**: repos are project-scoped and skills/MCP resolve solely from `runs.resource_manifest`, never the mutable global registry.
9. **Lifecycle mutations are atomic**: run status transitions are compare-and-swap on the expected status and event `seq` is allocated by the database, so concurrent writers can't clobber a status or collide on a seq (`run-lifecycle.ts`).

## Where to look

Expand Down
17 changes: 17 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,23 @@ All notable changes to Agrippa are documented here. The format follows

## [Unreleased]

### Security

- **Executor isolation seam** — one enforceable place (`packages/executor-core/isolation.ts`) decides every tool call and scrubs the subprocess environment. Read-only workspaces now actually deny shell and confine writes to the artifact directory; read-write workspaces confine writes to the workspace with a boundary-safe check (the previous `startsWith` let a sibling `<workspace>-evil` path through, and `Bash` bypassed the check entirely). **Reads (Read/Grep/Glob) are confined to the workspace too**, so the agent can't read `/proc/self/environ`, another run's directory, or the shared artifact store. The SDK subprocess environment is **allow-listed** — only the SDK auth variables and a fixed set of system essentials pass through, so platform secrets, DSNs, and injection vectors like `NODE_OPTIONS` are all dropped — with the OS `sandbox` enabled where available (bubblewrap installed in the worker image), `strictMcpConfig`, and repo-supplied `.claude`/`.mcp.json` removed after checkout so a checked-out repository can't inject hooks or permission overrides. The worker image runs as a non-root user with `/app` kept root-owned.
- **Event-payload secret redaction** — known secret values (the provider key, resolved MCP tokens) are redacted from every event before it is persisted or streamed over SSE, so a secret the agent echoes into output can't leak through the timeline.
- **Artifact path containment** — artifact ingestion resolves sources through `realpath` and rejects any that escape the workspace, closing a symlink disclosure (e.g. `ln -s /proc/self/environ`) that could exfiltrate secrets or other runs' files through the download endpoint.
- **Cross-tenant resource authorization** — submission rejects a `repoConnectionId` that isn't owned by the project, and the worker loads repo connections scoped to the run's project. Optional skills/MCP servers are now grant-checked: an authorized-resource manifest is pinned onto the run at submit (required grants enforced, optional resources included only when granted) and the worker resolves resources only from it, so a project without a grant can no longer receive the platform's global credential (e.g. the shared GitHub token).

### Fixed

- **Crash recovery for no-retry steps** — a worker that died mid-step no longer silently skips (or spuriously fails) a step without template retries; the crashed attempt no longer consumes the retry budget, and the executor session is carried onto the recovery attempt so resume works.
- **Atomic run lifecycle** — one `finalizeRun` owns every terminal transition: a single transaction does the status CAS (requiring `cancel_requested = false` for a success, so a late cancel wins atomically), finishedAt/totals, and the terminal event. Both the engine and the worker's retry-exhaustion path use it, so a run is never half-finalized and a queued run whose setup threw transitions `queued → failed` (now legal) with a terminal event instead of stranding. Event sequence numbers come from an atomic per-run counter (`runs.next_event_seq`), and approval decisions are CAS on `pending` with the sweeper re-enqueuing any run left paused by a lost resume enqueue.
- **Quota accounting** — the engine now counts the same monthly window as the submit gate, excludes the run's own spend from the headroom it checks (no double-count on resume), re-reads project usage at each step boundary so concurrent runs can't jointly overspend, and restores per-phase spend on resume so per-phase budgets aren't reset by a crash.
- **Artifact output contract** — patch steps no longer hand-write the diff (the engine generates it from `git diff`); collected artifacts are validated against the step's declared keys/kinds and matched by exact filename; only the current step's own expected files are cleared before an attempt (never a recursive `rm` a `.agrippa -> /work` symlink could redirect at the shared store, and `.agrippa` is stripped at checkout); missing/empty sources don't create zero-byte rows; binary (`file`-kind) artifacts are stored byte-exact; and ingestion size-caps and streams files instead of buffering them whole.
- **Authorized resources** — `requires.skills` is validated by the compiler and enforced by the engine; `skills` resolution returns `{ resolved, missing }` like `mcpServers`, so an unavailable required skill fails the step and an unavailable optional one is skipped (rather than throwing). `scripts/backfill-manifest.ts` reconstructs the manifest for pre-migration runs from **project grants** (not the full template).
- **SSE live gap** — the events stream subscribes and *awaits* the subscription being live before replaying, and treats the bus purely as a **wake-up** that triggers an ordered Postgres replay — so the cursor advances contiguously and a dropped event can't be skipped past (even on `Last-Event-ID` reconnect) (ADR-0007).
- **Production Compose** — split into a worker-only `workspaces` volume and a shared `artifacts` volume; the `api` runs as the non-root `bun` user and mounts only `artifacts` (consistent ownership whichever service initializes it) and receives `AGRIPPA_EXECUTOR` (the API chooses the executor at submit).

## [0.1.0] — 2026-07-17

The M1 milestone: all three layers of the platform, working end to end.
Expand Down
4 changes: 2 additions & 2 deletions apps/api/src/lib/audit.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { auditLogs, type Db } from "@agrippa/db";
import { auditLogs, type DbOrTx } from "@agrippa/db";
import type { Context } from "hono";
import type { AppEnv } from "../context";

Expand All @@ -14,7 +14,7 @@ type AuditEntry = {
* Every mutating handler records an audit row (docs/design/05-api-and-auth.md).
* Accepts an explicit tx so creations can be atomic with their mutation.
*/
export async function audit(c: Context<AppEnv>, entry: AuditEntry, tx?: Db): Promise<void> {
export async function audit(c: Context<AppEnv>, entry: AuditEntry, tx?: DbOrTx): Promise<void> {
const db = tx ?? c.var.db;
await db.insert(auditLogs).values({
orgId: c.var.user.orgId,
Expand Down
6 changes: 4 additions & 2 deletions apps/api/src/lib/usage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -45,8 +45,10 @@ export async function projectUsage(db: Db, projectId: string): Promise<ProjectUs

/**
* Submit-time quota gate (docs/design/04): hard-stop quotas reject new work
* once the period's spend has reached either limit. The engine re-checks at
* every step boundary for mid-run enforcement.
* once the current month's spend has reached either limit. The engine re-reads
* the same month-scoped project usage at every step boundary (excluding the
* run's own spend, which its budget meter already carries) for mid-run
* enforcement, so this gate and the engine agree on the accounting window.
*/
export async function assertQuotaHeadroom(db: Db, projectId: string): Promise<void> {
const [quota] = await db
Expand Down
139 changes: 72 additions & 67 deletions apps/api/src/routes/execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,13 +18,16 @@ import {
templateVersions,
} from "@agrippa/db";
import {
appendRunEvent as allocateRunEvent,
authorizeResources,
buildParamsValidator,
decideApproval,
resolveModelRoles,
SubmitError,
type TemplateDoc,
verifyResourceGrants,
verifyRepoRefs,
} from "@agrippa/orchestration";
import { and, asc, desc, eq, gt, max } from "drizzle-orm";
import { and, asc, desc, eq, gt } from "drizzle-orm";
import { Hono } from "hono";
import { streamSSE } from "hono/streaming";
import type { AppEnv } from "../context";
Expand All @@ -46,28 +49,6 @@ async function loadRunScoped(
return run;
}

/** Appends an API-originated event to the run log (seq = max + 1) and the bus. */
async function appendRunEvent(
c: { var: AppEnv["Variables"] },
runId: string,
type: string,
payload: Record<string, unknown>,
): Promise<void> {
const [maxSeq] = await c.var.db
.select({ v: max(runEvents.seq) })
.from(runEvents)
.where(eq(runEvents.runId, runId));
const seq = (maxSeq?.v ?? 0) + 1;
await c.var.db.insert(runEvents).values({ runId, seq, type, payload });
await c.var.bus?.publish({
runId,
seq,
type,
payload,
createdAt: new Date().toISOString(),
});
}

export const executionRoutes = new Hono<AppEnv>()
// ── Submit ──────────────────────────────────────────────────────────────────
.post(
Expand Down Expand Up @@ -108,11 +89,15 @@ export const executionRoutes = new Hono<AppEnv>()
await assertQuotaHeadroom(db, projectId);

try {
// every repoRef must reference a connection owned by this project
await verifyRepoRefs(db, projectId, compiled.spec.inputs, parsed.data);

const skillRows = await db.select({ id: skills.id, slug: skills.slug }).from(skills);
const mcpRows = await db
.select({ id: mcpServers.id, slug: mcpServers.slug })
.from(mcpServers);
await verifyResourceGrants(db, projectId, compiled, {
// required grants enforced; optional resources pinned only when granted
const resourceManifest = await authorizeResources(db, projectId, compiled, {
skillIdBySlug: new Map(skillRows.map((s) => [s.slug, s.id])),
mcpIdBySlug: new Map(mcpRows.map((m) => [m.slug, m.id])),
});
Expand Down Expand Up @@ -142,6 +127,7 @@ export const executionRoutes = new Hono<AppEnv>()
executorId: DEFAULT_EXECUTOR,
paramsSnapshot: parsed.data,
modelResolution,
resourceManifest,
budget: compiled.spec.budgets as unknown as Record<string, unknown>,
createdBy: user.id,
})
Expand Down Expand Up @@ -239,6 +225,7 @@ export const executionRoutes = new Hono<AppEnv>()
executorId: latest.executorId,
paramsSnapshot: latest.paramsSnapshot,
modelResolution: latest.modelResolution,
resourceManifest: latest.resourceManifest,
budget: latest.budget,
createdBy: c.var.user.id,
})
Expand Down Expand Up @@ -303,33 +290,56 @@ export const executionRoutes = new Hono<AppEnv>()
.from(approvals)
.where(and(eq(approvals.id, c.req.param("approvalId")), eq(approvals.runId, run.id)));
if (!approval) throw AppError.notFound("Approval");
if (approval.status !== "pending") {
throw AppError.conflict("already_decided", `Approval is ${approval.status}`);
}
const [updated] = await c.var.db
.update(approvals)
.set({
status: input.decision,
decidedBy: c.var.user.id,
decidedAt: new Date(),
comment: input.comment,
})
.where(eq(approvals.id, approval.id))
.returning();
await appendRunEvent(c, run.id, "approval.decided", {
const eventPayload = {
approvalId: approval.id,
checkpointId: approval.checkpointId,
decision: input.decision,
};
// decision (CAS on status='pending'), the approval.decided event, and the
// audit row commit together — a partial write can't leave the timeline or
// audit log missing the decision that a later retry would then skip
const result = await c.var.db.transaction(async (tx) => {
const updated = await decideApproval(tx, approval.id, {
status: input.decision,
decidedBy: c.var.user.id,
comment: input.comment,
});
if (!updated) return { updated: null as null };
const event = await allocateRunEvent(tx, {
runId: run.id,
type: "approval.decided",
payload: eventPayload,
});
await audit(
c,
{
action: "run.approval.decide",
resourceType: "approval",
resourceId: approval.id,
projectId: run.projectId,
payload: { decision: input.decision },
},
tx,
);
return { updated, event };
});
await audit(c, {
action: "run.approval.decide",
resourceType: "approval",
resourceId: approval.id,
projectId: run.projectId,
payload: { decision: input.decision },
if (!result.updated) {
// already decided, OR a prior attempt decided then failed to enqueue: in
// both cases the durable state is correct, so re-enqueue to unstick the
// run (the sweeper also backstops this) and report the conflict
await c.var.queue?.enqueueRun(run.id);
throw AppError.conflict("already_decided", "Approval is already decided");
}
// publish + enqueue only after the decision durably committed
await c.var.bus?.publish({
runId: run.id,
seq: result.event.seq,
type: "approval.decided",
payload: eventPayload,
createdAt: result.event.createdAt.toISOString(),
});
await c.var.queue?.enqueueRun(run.id); // resume at the gated phase
return c.json(updated);
return c.json(result.updated);
})

// ── Artifacts ───────────────────────────────────────────────────────────────
Expand Down Expand Up @@ -419,9 +429,6 @@ export const executionRoutes = new Hono<AppEnv>()
return rows;
};

// 1) replay history (gap-free by construction)
await replay();

const isTerminal = async (): Promise<boolean> => {
const [row] = await db
.select({ status: runs.status })
Expand All @@ -430,40 +437,38 @@ export const executionRoutes = new Hono<AppEnv>()
return row ? isTerminalRunStatus(row.status) : true;
};

if (await isTerminal()) return;

// 2) live: bridge the bus when present, else poll the DB
// Live: bridge the bus when present, else poll the DB. The bus is only a
// WAKE-UP — every event is delivered by an ordered `replay()` from
// Postgres, so the cursor advances contiguously. Sending bus events
// directly would advance a high-water cursor past a dropped seq, and that
// gap would then be skipped forever (even on Last-Event-ID reconnect).
if (bus) {
const queue: Array<() => Promise<void>> = [];
let notify: (() => void) | null = null;
const unsubscribe = bus.subscribe(run.id, (event) => {
queue.push(async () => {
if (event.seq > cursor) {
await sendRow({ ...event, createdAt: event.createdAt });
}
});
notify?.();
});
const subscription = bus.subscribe(run.id, () => notify?.());
try {
// wait until the subscription is actually live, THEN replay history,
// so nothing published in between is dropped (ADR-0007)
await subscription.ready;
await replay();
while (!closed) {
while (queue.length > 0) {
const job = queue.shift();
if (job) await job();
}
if (await isTerminal()) {
await replay(); // drain anything raced between bus and DB
await replay(); // final ordered drain
break;
}
// sleep until woken by the bus (or a 2 s safety tick), then replay
await new Promise<void>((resolve) => {
notify = resolve;
setTimeout(resolve, 2000);
});
notify = null;
await replay();
}
} finally {
unsubscribe();
subscription.unsubscribe();
}
} else {
await replay();
if (await isTerminal()) return;
while (!closed) {
await new Promise((resolve) => setTimeout(resolve, 1000));
await replay();
Expand Down
Loading
Loading