Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

- **Scale-out, phase A: the foundation for API pods that share one database and runs that outlive them** (`docs/plans/scale-out.md`). Every API process has an id (`OPTIO_INSTANCE_ID`, the pod's name from the chart, plus a suffix new at every boot; `GET /api/health` says which instance answered, to a signed-in request only). One migration adds every table and column the plan's phases build on: the attachment columns on runs (`exec_state`, `exec_pid`, `consumed_bytes`, `attached_by`, `attach_lease_until`, on tasks, Job runs, PR-review runs and agent turns), `leases`, `ws_upgrade_tokens`, `inbound_webhook_deliveries`, `glance_state`, `installed_skill_files`, `ticket_sync_claims`, and one active PR review per PR URL. `services/lease-service.ts` is the one way an instance claims something for a while. The **run protocol** (`packages/container-runtime/src/run-protocol.ts`) is on every container runtime — `startRun`, `attachRun`, `deliverStdin`, `killRun`: a start script launches the agent under a supervisor in its own session in the run's home, with stdin fed from a file through a FIFO until the `__OPTIO_STDIN_EOF__` line, its output appended to `output.ndjson` and its exit code to `exit`; the supervisor writes its own pid file (`<pid> <starttime>`) once its TERM trap is in place, and attach, kill and the pooled start guard trust a pid only when the live process's start time matches it; stdin bytes travel over the exec's stdin (`head -c <n>`), never in a script or its environment; an attach is `tail --pid` from a byte offset; the exec that started a run can go away and any instance can attach later. `tini` is the pods' pid 1 (`images/base.Dockerfile`, every init script), so finished runs are reaped — **the agent images must be rebuilt for the protocol**. The repo-cleanup worker removes the run homes of finished runs from every pod (`run-home-sweep-service.ts`), the backstop for the attached worker that removes them once a run is terminal and consumed (phase B). `buildTaskStartScript` / `buildPooledStartScript` build the start scripts beside the exec scripts the workers still use (they switch in phase B). The fake runtime keeps its pods on disk (`OPTIO_FAKE_RUNTIME_DIR`) and plays a protocol run as a detached process (`fake-agent.mjs`, every event with a `seq`), and `test-utils/e2e/api-cluster.ts` boots N real API servers on one database for the e2e tier: two boot together, one is SIGKILLed, a third joins.

- **Scale-out, phase C1: coordination that holds across API instances** (`docs/plans/scale-out.md` §3). A run's concurrency claim — fewer than the limit running, so take this one — runs under a Postgres advisory lock on the global limit (`claim:tasks`, `claim:workflows`; a repo's or a Job's own limit is counted under it) in one READ COMMITTED transaction (`services/claim-lock.ts`), in place of a mutex one process kept to itself; the count and the compare-and-swap claim run on the lock's transaction (`claimTransitionIn`, `claimWorkflowRunIn`), and what follows the claim (events, webhooks, the reconciler) is registered with `afterCommit` and runs once it commits. A claimer waits for the lock at most 10 s (`lock_timeout`) and otherwise re-queues as when a limit is full. Every periodic sweep (ticket sync, the schedule checker, external PR review, pod cleanup, the PR watcher, the reconcile resync, skill sync, token validation, the config directory) runs on one instance at a time under a lease (`services/poller-lease.ts`), and each is safe to overlap anyway: ticket sync claims each (source, ticket, repo) in `ticket_sync_claims` before creating a task (a claim whose task was deleted, or whose sweep died before making one, is taken over); the schedule checker advances a trigger's `next_fire_at` by compare-and-swap before firing it, so a tick fires once; an external PR review sweep launches only while the review is still the one it saw (none, for a new PR) and treats a lost race or the one-active-review-per-PR conflict as someone else's launch. Repeat jobs are BullMQ job schedulers with stable ids (`services/repeat-jobs.ts`), so instances booting in any order leave exactly one schedule per tick; boot no longer wipes every repeat job, only the hash-keyed ones an older build registered. Inbound webhook deliveries (GitHub, GitLab, Slack, Linear, Jira, PagerDuty, Sentry) are claimed in `inbound_webhook_deliveries` (`services/inbound-delivery-service.ts`) instead of a per-process set and a Redis key; WebSocket upgrade tokens live in `ws_upgrade_tokens`, so a token minted on one instance is accepted by whichever one the browser's upgrade lands on, once; the Optio assistant's one-conversation-per-user rule (30 s lease renewed every 10 s), a user's GitHub token refresh and the config directory's apply are leases (`withLease` takes a `wait` option for the latter two). A GitHub user token is only ever refreshed under its lease — GitHub rotates the refresh token, so two instances refreshing at once would strand one: a waiter re-reads the token the holder stored, the call to GitHub times out after 15 s, a wait that times out uses the stored token while it works and otherwise the PAT, and a refused refresh deletes the stored tokens only while the refresh token is still the one it sent. A skill sync that finds the skill's lease held retries with backoff, and records its result only while the skill's ref and subpath are still the ones it resolved (a compare-and-swap in the same transaction): a sync of the old ref that finishes after a PATCH moved the ref writes nothing, and the skill stays due for the next pass. "Sync now" answers 409 while another instance is applying the directory. A webhook delivery's claim gets one second (`statement_timeout`) before the delivery is accepted unclaimed; an upgrade token's expiry is the database's clock. A housekeeping tick (`workers/sweep-worker.ts`, every 30 s under a lease, `registerSweep` for more) removes deliveries older than a day, expired upgrade tokens and long-expired leases. `docs/production-eks.md` says what stays per pod. Coverage: integration tests drive every claim with two racing callers; `apps/api/e2e/scale-out-coordination.e2e.test.ts` runs two real API servers with auth on — an upgrade token minted on A accepted by B once, one signed webhook delivery posted to both firing once, a due schedule fired once by two sweepers, and one scheduler per repeat job after both boot.
- **Android: every trigger type in the automations and agent trigger sheets.** A Local automation's Add trigger sheet and a persistent agent's New trigger sheet take all fourteen trigger types — GitLab, Jira, Pylon, PagerDuty, Sentry, Alertmanager and Datadog events included, with their event kinds, identity and filters — through the same rows as the New work form (one trigger editor in `core:ui`, `triggers/`). A Pylon / Alertmanager / Datadog trigger's own URL and shared secret are shown once, with copy buttons, right after it is created, in the sheets and in the New work form.
- **Every agent runtime gets the work's MCP servers.** Connections' tools and MCP servers reached only Claude Code (`.mcp.json`) and Codex (`config.toml`); a Job on Gemini, a persistent agent on Copilot, or a Task on OpenCode or Cursor had the connections' credentials and notes but no tools. Each runtime now gets the same servers in the file it reads, established from the versions the agent image installs: Gemini CLI in the user `settings.json` of a `GEMINI_CLI_HOME` of the run's own (the adapter's settings merged in, each server `trust: true`), OpenCode in an `OPENCODE_CONFIG` file merged before the project's own, GitHub Copilot CLI through `--additional-mcp-config @<file>` (`type: "local"`, every tool enabled), and Cursor in the project's `.cursor/mcp.json` with `--approve-mcps`. See `docs/connections.md` → "Which runtime reads which file".
- The per-run files live in one **run home**, `/home/agent/optio/runs/<run id>` (`OPTIO_RUN_HOME`), named after the task, Job run, or agent turn — so a retry lands in the same place — and removed when the run's script exits and again when the task's worktree is cleaned up. Codex's per-run `CODEX_HOME` moves there too; it used to be a random directory that was never removed from the repo pod's home volume.
Expand Down
2 changes: 1 addition & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -281,7 +281,7 @@ After changing backend code, ALWAYS rebuild + redeploy the local cluster and ver
- **Drizzle ORM**: schema in `apps/api/src/db/schema.ts`, run `drizzle-kit generate` after changes. **New migrations use unix-timestamp prefixes** (`migrations.prefix: "unix"` in `drizzle.config.ts`). Existing `00xx_*` files are frozen — never rename them. Migrations are hand-written SQL (the drizzle-kit snapshots stopped at `0012`). **Run columns**: add them to `tasks` as always (`ALTER TABLE "tasks" ...`) and to `runColumns()` in `schema.ts` (the table `workRuns` and the view `tasks` = `repo_tasks` both use it); `migrate-safe.ts` drops the run views before every migration and remakes them after (`db/run-views.ts`), so `repo_tasks` picks the column up by itself. A column Job runs need goes in `workflow_runs`'s list in `run-views.ts` too, and in `workflowRuns` in the schema
- **Zustand**: use `useStore.getState()` in callbacks/effects, not hook selectors (avoids infinite re-renders)
- **Next.js webpack**: `extensionAlias` in `next.config.ts` resolves `.js` → `.ts` for workspace packages
- **State transitions**: always go through `taskService.transitionTask()` — validates, updates DB, records event, publishes WebSocket
- **State transitions**: always go through `taskService.transitionTask()` — validates, updates DB, records event, publishes WebSocket. Its in-transaction form is `claimTransitionIn(claim, …)` inside `withClaimLock` (`services/claim-lock.ts`): the write runs on the claim's transaction and the announcement is registered with `claim.afterCommit`, so it runs only once the transaction commits
- **Secrets**: never log or return secret values. Encrypted at rest with AES-256-GCM
- **Cost tracking**: stored as string (`costUsd`) to avoid float precision issues
- **K8s RBAC**: namespace-scoped Role (pods, exec, secrets, PVCs) + ClusterRole (nodes, namespaces, metrics)
Expand Down
255 changes: 255 additions & 0 deletions apps/api/e2e/scale-out-coordination.e2e.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,255 @@
/**
* E2E: phase C1 of docs/plans/scale-out.md through two REAL API servers on
* one database (the api-cluster harness), with auth ENABLED so the
* WebSocket upgrade path is the real one:
*
* 6. a WebSocket upgrade token minted on A is accepted by B, once;
* 7. the same signed webhook delivery posted to A and to B fires once;
* 8. a schedule trigger due now, with both servers sweeping, fires once;
* + two servers booting leave exactly one scheduler per repeat job.
*/
import { createHash, randomBytes, randomUUID } from "node:crypto";
import { Queue } from "bullmq";
import postgres from "postgres";
import { afterAll, beforeAll, describe, expect, it } from "vitest";
import { startApiCluster, type ApiCluster } from "../src/test-utils/e2e/api-cluster.js";
import { waitFor } from "../src/test-utils/e2e/api-server.js";
import { REPEAT_SCHEDULERS } from "../src/services/repeat-jobs.js";

const GITLAB_WEBHOOK_SECRET = "e2e-scale-out-gitlab-token";

let cluster: ApiCluster;
let adminToken = "";

/** An admin with a session, seeded straight into the DB (like webhook-ingress-auth.e2e). */
async function seedAdmin(): Promise<void> {
const sql = postgres(process.env.DATABASE_URL!, { max: 1 });
try {
const wsId = randomUUID();
const userId = randomUUID();
adminToken = `e2e-admin-${randomBytes(16).toString("hex")}`;
const tokenHash = createHash("sha256").update(adminToken).digest("hex");
await sql`INSERT INTO workspaces (id, name, slug) VALUES (${wsId}, 'Scale-out e2e', ${`scale-out-e2e-${wsId.slice(0, 8)}`})`;
await sql`
INSERT INTO users (id, provider, external_id, email, display_name, default_workspace_id)
VALUES (${userId}, 'github', 'scale-out-e2e-admin', 'admin@scale-out.e2e', 'Scale-out admin', ${wsId})`;
await sql`INSERT INTO workspace_members (workspace_id, user_id, role) VALUES (${wsId}, ${userId}, 'admin')`;
await sql`
INSERT INTO sessions (user_id, token_hash, expires_at)
VALUES (${userId}, ${tokenHash}, NOW() + INTERVAL '1 day')`;
} finally {
await sql.end();
}
}

beforeAll(async () => {
await seedAdmin();
cluster = await startApiCluster({
size: 2,
env: {
OPTIO_AUTH_DISABLED: "false",
GITLAB_WEBHOOK_SECRET,
OPTIO_WORKFLOW_TRIGGER_INTERVAL: "1000",
},
});
}, 240_000);

afterAll(async () => {
await cluster?.stop();
});

function server(i: number): string {
return cluster.servers[i].handle.baseUrl;
}

async function api<T = any>(
i: number,
method: string,
path: string,
body?: unknown,
): Promise<{ status: number; body: T }> {
const res = await fetch(`${server(i)}${path}`, {
method,
headers: {
authorization: `Bearer ${adminToken}`,
...(body !== undefined ? { "content-type": "application/json" } : {}),
},
...(body !== undefined ? { body: JSON.stringify(body) } : {}),
});
return { status: res.status, body: (await res.json().catch(() => null)) as T };
}

/**
* Opens a WebSocket with the upgrade token in the subprotocol. The upgrade
* itself always completes; the server authenticates after it and closes
* with 4401 when the token is no good. So: accepted when the socket is
* still open a moment later (then we close it), else the server's code.
*/
function upgrade(i: number, token: string): Promise<{ accepted: boolean; closeCode?: number }> {
return new Promise((resolve) => {
const url = `${server(i).replace(/^http/, "ws")}/ws/events`;
const ws = new WebSocket(url, ["optio-ws-v1", `optio-auth-${token}`]);
let settled = false;
let timer: NodeJS.Timeout | undefined;
ws.addEventListener("open", () => {
timer = setTimeout(() => {
settled = true;
ws.close(1000);
resolve({ accepted: true });
}, 1500);
});
ws.addEventListener("close", (ev) => {
clearTimeout(timer);
if (!settled) resolve({ accepted: false, closeCode: ev.code });
});
ws.addEventListener("error", () => {
/* close follows */
});
});
}

describe("scale-out coordination across two API servers", () => {
it("a WebSocket upgrade token minted on A is accepted by B, and only once", async () => {
const minted = await api<{ token: string }>(0, "GET", "/api/auth/ws-token");
expect(minted.status).toBe(200);
expect(minted.body.token).not.toBe("auth-disabled");

const onB = await upgrade(1, minted.body.token);
expect(onB).toEqual({ accepted: true });

// Consumed on B: A rejects the same token.
const again = await upgrade(0, minted.body.token);
expect(again).toEqual({ accepted: false, closeCode: 4401 });
}, 30_000);

it("the same signed webhook delivery posted to A and to B fires once", async () => {
const project = `acme/scale-out-${Date.now()}`;
const created = await api<{ workflow: { id: string } }>(0, "POST", "/api/jobs", {
name: `push summary ${Date.now()}`,
promptTemplate: "Summarize {{commits}} on {{sourceBranch}}",
agentRuntime: "claude-code",
});
expect(created.status, JSON.stringify(created.body)).toBe(201);
const jobId = created.body.workflow.id;
const trigger = await api(0, "POST", `/api/jobs/${jobId}/triggers`, {
type: "gitlab",
config: { events: ["push"], projects: [project], branches: ["main"] },
});
expect(trigger.status, JSON.stringify(trigger.body)).toBe(201);

const raw = JSON.stringify({
object_kind: "push",
ref: "refs/heads/main",
before: "aaaa1111",
after: "bbbb2222",
user_username: "alice",
project: { path_with_namespace: project, web_url: `https://gitlab.com/${project}` },
commits: [{ id: "bbbb2222", title: "fix: thing", message: "fix: thing" }],
});
const deliveryId = randomUUID();
const deliver = async (i: number) => {
const res = await fetch(`${server(i)}/api/webhooks/gitlab`, {
method: "POST",
headers: {
"content-type": "application/json",
"x-gitlab-event": "Push Hook",
"x-gitlab-event-uuid": deliveryId,
"x-gitlab-token": GITLAB_WEBHOOK_SECRET,
},
body: raw,
});
return { status: res.status, body: await res.json() };
};

// The provider's retry lands on the other instance.
const [onA, onB] = await Promise.all([deliver(0), deliver(1)]);
expect(onA.status).toBe(200);
expect(onB.status).toBe(200);
expect([onA.body, onB.body]).toEqual(
expect.arrayContaining([{ ok: true }, { ok: true, duplicate: true }]),
);
const runs = async () =>
(await api<{ runs: unknown[] }>(1, "GET", `/api/jobs/${jobId}/runs`)).body.runs;
await waitFor(async () => ((await runs()).length >= 1 ? true : null), {
timeoutMs: 30_000,
label: "a run of the GitLab-triggered Job",
});
await new Promise((r) => setTimeout(r, 1500));
expect(await runs()).toHaveLength(1);
}, 60_000);

it("a schedule trigger due now fires once with both servers sweeping", async () => {
const created = await api<{ workflow: { id: string } }>(1, "POST", "/api/jobs", {
name: `nightly ${Date.now()}`,
promptTemplate: "Nightly sweep",
agentRuntime: "claude-code",
});
expect(created.status).toBe(201);
const jobId = created.body.workflow.id;
// A yearly cron: once nextFireAt is pulled into the past it is due
// exactly once, and the advance puts it a year away.
const trigger = await api<{ trigger: { id: string } }>(
1,
"POST",
`/api/jobs/${jobId}/triggers`,
{
type: "schedule",
config: { cronExpression: "0 0 1 1 *" },
},
);
expect(trigger.status, JSON.stringify(trigger.body)).toBe(201);
const sql = postgres(process.env.DATABASE_URL!, { max: 1 });
try {
await sql`UPDATE workflow_triggers SET next_fire_at = NOW() - INTERVAL '1 minute' WHERE id = ${trigger.body.trigger.id}`;
} finally {
await sql.end();
}

const runs = async () =>
(await api<{ runs: unknown[] }>(0, "GET", `/api/jobs/${jobId}/runs`)).body.runs;
await waitFor(async () => ((await runs()).length >= 1 ? true : null), {
timeoutMs: 30_000,
label: "the schedule to fire",
});
// Both servers sweep every second; give them several more ticks.
await new Promise((r) => setTimeout(r, 4000));
expect(await runs()).toHaveLength(1);
const listed = await api<{ triggers: Array<{ id: string; nextFireAt: string | null }> }>(
0,
"GET",
`/api/jobs/${jobId}/triggers`,
);
const after = listed.body.triggers.find((t) => t.id === trigger.body.trigger.id);
expect(new Date(after!.nextFireAt!).getTime()).toBeGreaterThan(Date.now());
}, 60_000);

it("two servers booting leave exactly one scheduler per repeat job", async () => {
const { getBullMQOptions } = await import("../src/services/redis-config.js");
const opts = getBullMQOptions();
const expectations: Array<[string, number]> = [
["pr-watcher", 1],
["external-pr-review", 1],
["repo-cleanup", 2],
["ticket-sync", 1],
["workflow-trigger-checker", 1],
["token-validation", 1],
["reconcile-resync", 1],
["skill-sync", 1],
["housekeeping", 1],
];
for (const [name, expected] of expectations) {
const queue = new Queue(name, { ...opts });
try {
const schedulers = await queue.getJobSchedulers();
expect(schedulers.map((s) => s.key).sort(), `${name} schedulers`).toEqual(
[...REPEAT_SCHEDULERS[name as keyof typeof REPEAT_SCHEDULERS]].sort(),
);
expect(schedulers).toHaveLength(expected);
// And nothing hash-keyed beside them.
expect(await queue.getRepeatableJobs()).toHaveLength(expected);
} finally {
await queue.close();
}
}
}, 30_000);
});
Loading
Loading