Skip to content

Latest commit

 

History

History
328 lines (290 loc) · 23.6 KB

File metadata and controls

328 lines (290 loc) · 23.6 KB

Architecture

How Operator is put together. This is the public companion to CLAUDE.md (the in-repo codebase map agents read); if the two ever disagree, trust the code.

Three processes, one origin

  • server.js — custom Next.js server (plain Node, Turbopack in dev). Fronts Next on one port, proxies /pty WebSocket upgrades to the sidecar, and forwards dev HMR upgrades to Next. Because everything rides one origin, a Cloudflare Tunnel (or any reverse proxy) exposing a single https hostname carries both the app and the terminal — no second port, and wss:// is used automatically over https.
  • pty-server.js — the node-pty terminal sidecar, bound to 127.0.0.1 only; never exposed directly. The browser reaches it through the app origin at /pty.
  • The Next app — UI in app/, REST under app/api/, server logic in lib/.

The turn lifecycle

lib/runner.ts is the detached turn runner: POST /api/tasks/[id]/messages launches a turn and returns immediately — the turn runs server-side, owned by the process (not by any HTTP request), persisting every event to SQLite and publishing it on lib/events.ts (in-process pub/sub keyed by task id, plus a wildcard channel that sees every task's events). Stopping is only ever explicit, via lib/abort.ts. If a turn is already running, the message parks in the pending_messages queue to run next.

GET on the same route is the SSE watch stream: a snapshot of the persisted transcript, then a live tail — reconnect-safe, any number of viewers, zero viewers fine.

GET /api/events is the global lifecycle stream: one always-open SSE connection per client tab broadcasting coarse turn boundaries (turn started / awaiting input / answered / suggestion created / turn ended) for every task across every project. Each event carries a fresh snapshot of the task row plus its project's awaiting count — that's how spinners, project badges, and the "N need you" pill update instantly for tasks whose transcript stream isn't open. There is no task-list polling.

A task is a lineage of sessions. Generation N ends at /clear; its transcript is condensed to a summary, and generation N+1 starts with a clean context window seeded by all prior summaries. The task persists — only the context window resets.

/clear resets tasks.started to 0, so the next send is an opening turn — but not a first one. buildOpeningPrompt() (lib/agents/shared.ts) splits on the generation: generation 1 gets the task text (buildInitialPrompt), generation > 1 gets the short continue-from-here buildResumePrompt, carrying whatever the user typed on that send. Both launchers — POST /api/tasks/[id]/messages and lib/autoStart.ts — go through that one helper. Only a never-started task takes the route's inline first-turn branch; a resumed generation goes through startResumeTurn() like any other follow-up (its session_id is null, so the driver still opens a fresh session).

The agent-driver seam (lib/agents/)

The app talks to coding agents only through the AgentDriver interface.

  • types.ts defines the interface: a normalized StreamEvent turn contract, one-shot summarize/draft/recap helpers, a capability descriptor, and the login/verify auth surface.
  • registry.ts resolves a driver by id — getDriver(task.agent), persisted per task, defaulted per project via projects.default_agent.
  • shared.ts holds the agent-agnostic pieces every driver reuses: project-context and conflict prompts, tool-call → title/peek/diff normalizers, the event queue.
  • GET /api/agents serves each driver's capability descriptor plus its persisted connection state to the client, which renders every run-control picker (model / reasoning / permission), the per-task agent picker, agent badges, and the cost/ask feature gates from that data — no hardcoded per-agent lists in the UI.
  • Connecting an agent is driver-driven and route-generic: /api/agents/[id]/{login,login/code,verify,api-key,status} resolve getDriverStrict(id) and call its auth surface, so a new agent costs zero new routes; connections.ts records which agents are connected. Each agent's credentials live under $HOME (Claude ~/.claude, Codex ~/.codex); the optional per-token API-key paths persist to a 0600 file (lib/anthropic-key.ts / lib/openai-key.ts).

The Claude driver (lib/agents/claude/driver.ts)

runTurn() via the Claude Agent SDK (resume or fresh session, project context appended to the Claude Code system prompt), the suggest_task + expose_service MCP tools, summarizeTranscript() for /clear, and draftProjectContext() (a read-only agent loop that explores the repo to refresh a project's saved context). Auth delegates to lib/claude-auth.ts. Sessions run permissionMode: "bypassPermissions".

The Codex driver (lib/agents/codex/driver.ts)

Driven by the user's ChatGPT-plan codex login (no API key). Built on @openai/codex-sdk (it spawns the codex CLI and speaks JSONL over stdio, same architecture as the Claude driver): startThread() / resumeThread(session_id), with the codex thread id emitted as the session event so lineage/resume works unchanged. events.ts normalizes codex's ThreadItem stream (agent_message → assistant; command_execution / file_change / mcp_tool_call / web_search / todo_list / reasoning → tool + tool_result; turn.completed usage → tokens plus an estimated cost_usd) into the StreamEvent contract.

Run controls map our permission modes to codex's sandbox/approval policy (bypassPermissions → workspace-write + approvals-never; plan → read-only); reasoning levels are the CLI's own effort names (per model, from model/list), passed straight to model_reasoning_effort (see lib/agents/reasoning.ts). Capabilities declare supportsMcpTools: true (the orchestrator's tools reach codex through the portable stdio MCP bridge below, registered per turn with a ~1-day tool_timeout_sec so a parked ask survives), supportsAsks: true (codex has no native interactive-ask hook, but the bridge's ask_user tool surfaces the same question card and blocks until the user answers) and reportsCostUsd: false + costIsEstimated: true — ChatGPT-plan auth reports token counts only, so pricing.ts estimates the dollar cost per turn (tokens × published API prices for the resolved model) and the UI renders those figures with a ~. The one upstream limitation not papered over: the non-interactive CLI cannot pause a turn for command approval, so on-request approval modes aren't offered — permission modes are Auto-run (workspace-write, approvals never) and Plan (read-only). Auth (auth.ts) drives codex login --device-auth + codex login status. The one-shot helpers run as codex exec one-shots in a read-only sandbox (no writes, no approvals, no network), bounded by an item cap — the codex analog of the Claude helpers' maxTurns — so a runaway helper turn is cut off rather than looping unbounded. Binary via CODEX_CLI_PATH (else the SDK auto-resolves its bundled binary / PATH).

Internal one-shots (lib/agents/oneshots.ts)

Routing for the internal jobs that run a turn outside the main chat: /clear handoff summaries, project recaps, and "Refresh with AI" context drafts. Two policies: task-scoped one-shots (/clear transcript summarization) follow the task's own agent, so a Codex task's handoff note is written by Codex and counted against the Codex login; project-scoped one-shots (recap, context draft) aren't tied to any one task, so they run on the utility agent, resolved connected-first: the utility_agent app setting when that agent is actually connected → the app default agent → the built-in default → any connected agent at all — so a Codex-only instance gets working recaps and context drafts with zero configuration, and when NO agent is connected the job fails fast with an actionable "connect an agent in Settings → Agents" error instead of driving a dead CLI. Either way, if the chosen driver doesn't implement a given helper, the utility agent backstops it — so a new driver can ship runTurn() alone and still get working summaries/recaps/drafts. AI conflict-resolution turns need no special routing: buildConflictPrompt() (lib/agents/shared.ts) produces the prompt and the client sends it as an ordinary message, so it flows through startTurn() → the task's driver like any turn.

Unattended one-shots are server-gated by the background_jobs setting (default on). Project recaps add a second recap_mode gate: automatic (default), on_open, or off. The five-minute sweep requires automatic; opening a project accepts automatic or on_open. Explicit /clear, Refresh with AI, and manual recap refreshes bypass the unattended gate. Settings reads a single 30-day aggregation from internal_usage so the controls show their run count and API-price-equivalent cost without polling.

The agent-tool bridge (scripts/orch-mcp.mjs + lib/agentTools.ts)

complete_step (auto-advance chains, lib/chains.ts) rides the same seams but is mounted only for a chain step: the Claude driver adds it to its in-process server when isAutoAdvanceTask(task), and the Codex driver sets ORCH_COMPLETE_STEP=1 in the bridge's env so the bridge registers it (→ /api/internal/agent-tools/complete-step). Both call recordStepComplete(), which only stamps the task; the runner's turn-end chainAdvanceBlocker() decides whether to advance, and finishChainStep() (lib/autoStart.ts) commits, moves the step to in_review, and launches the next step on a worktree stacked on this step's branch. The blocker rule shared by the server and the client lives in lib/chainRules.ts.

The end-of-chain review lives in lib/chainMerge.ts (routes under app/api/chains/[id]/). Because steps are stacked, landing a chain is one mergeTask of the target step's branch (the last step, or through) into chains.base_branch, run under the task lock of every step (taken in chain order) so the "nothing running" check is atomic with the git work. Before merging it verifies the stack (each earlier branch is an ancestor of the next, no uncommitted edits below the target) and reads each step's own line stats (base_sha → stepOwnTip, which peels base-branch merges off the tip so a conflict resolution doesn't inflate them); after a successful merge each step gets merged_at, status done, its base_sha advanced to what landed, and its own task_merges row. The conflict path reuses prepareWorktreeMerge on the target step's worktree; a retry sees the staged merge and finishes it with completeWorktreeMerge. Merging never rebases later steps. The tasks-column card is derived client-side (app/orchestrator/chains.ts) from rows the global event stream keeps live. Those events carry each task's step_pause, so a paused step's reason is live too.

The review actions live in lib/chainActions.ts, all under the same every-step locks. Send back records a chain_fixups row (feedback, the step it's about, the last step's HEAD as start_sha), claims the last step's turn slot, and launches the fix-up through startResumeTurn. While a fix-up is open, recordStepComplete writes the summary to the fix-up instead of step_summary. finishChainStep commits it as Fix-up: … and closes it. So does a manual status change to In review. Discard from step k runs the task DELETE route's teardown for k..n, last first. Rebase stack runs rebaseWorktreeOnto for each unmerged step and rolls everything back with reset --hard on any conflict. It refuses up front when a step's base_sha..HEAD contains a merge commit (hasMergeCommits), because rebase --onto would drop it along with any edits made inside it. relinkAfterStepDelete (called by the task DELETE route) adds a dependency from the next unstarted step to the nearest earlier one. previousChainStep skips gaps, so the next step stacks there. orphanedStackProblems (checked by the view and the merge) catches a step that stacked on a step that was later deleted: after a gap, its base_sha is no longer reachable from the step before it. The runner saves pauseReason(blocker) into tasks.step_pause when a chain step's turn pauses, and clears it when the next turn starts. The "needs you" count adds CHAIN_AWAITS_REVIEW (lib/store.ts, mirrored by chainAwaitsReview on the client) to the task predicate, once per finished chain.

suggest_task / expose_service / ask_user are the same orchestrator tools every driver exposes. The Claude driver mounts the first two as an in-process SDK MCP server (createSdkMcpServer) and gets asks natively via its AskUserQuestion hook; the portable equivalent is scripts/orch-mcp.mjs, a plain-Node stdio MCP server (@modelcontextprotocol/sdk) the non-Claude drivers spawn per turn. It's a thin proxy: it reads ORCH_TASK_ID / ORCH_PROJECT_ID / ORCH_BASE_URL / SERVICE_TOKEN from env (injected by the driver) and POSTs each tool call to the app's internal endpoints (app/api/internal/agent-tools/{suggest-task,expose-service,ask-user}, gated by the strict per-instance SERVICE_TOKEN in middleware.ts). ask_user is the asynchronous one: the endpoint persists + publishes the same interactive question card the Claude hook produces, parks a detached waiter on the user's answer (lib/asks.ts, tied to the turn's abort signal), and the bridge polls the sibling ask-user/wait endpoint for the settled outcome — no long-held HTTP request, and the ask survives page reloads because the card lives in the transcript. Both the in-process server and the endpoints call the SAME shared logic in lib/agentTools.ts, and both build their tool defs from the SAME constants in lib/agentToolDefs.mjs, so the two paths can't drift.

Whether an agent may file a task the user never asked for is the suggestion_policy setting (ask_first default, auto; lib/suggestionPolicy.ts). buildProjectContext() reads it and swaps the suggest_task paragraph: ask_first tells the agent to list proposed follow-ups in chat and confirm — via AskUserQuestion on Claude, the bridge's ask_user elsewhere, picked from task.agent — before calling the tool, while an explicit "plan/break down/scope/roadmap" request stays free to file tasks directly. auto restores the old always-proactive wording. The tool DESCRIPTION in lib/agentToolDefs.mjs states the ask-first rule too, so the tool itself says it; those strings are static (the stdio bridge has no DB), so the auto branch of the prompt says outright that it overrides them. Enforcement stays prompt-based on purpose: the confirmation is prose in the transcript, not a tool call, so no server-side check can tell an approved suggestion from an unrequested one. Instead createSuggestedTask() emits a suggestion_created analytics event carrying the policy in force, which is what makes drift visible — under ask_first a rising rate means the prompt has stopped landing.

Every suggest_task call is stamped with its proposer: createSuggestedTask() takes an optional source: { taskId, generation } and writes it to tasks.suggested_by_task_id / tasks.suggested_by_generation, which is what the tray groups by. Both callers already know the running task — the Claude driver's MCP server closes over it, and the bridge posts the ORCH_TASK_ID it uses for ask_user, from which the endpoint re-reads the live generation. The id is a self-referential FOREIGN KEY with ON DELETE SET NULL, so deleting the proposing task orphans its suggestions rather than taking them with it; a source whose task has vanished mid-turn is dropped at insert time for the same reason the project is re-read there. Provenance is write-once — updateTask() deliberately doesn't carry those columns.

suggested_by_generation has no FK, so it SURVIVES that SET NULL — and the tray leans on the asymmetry: a null id next to a non-null generation is durable evidence the proposer was hard-deleted ("From a deleted task", rendered stale), while both-null is a row that never recorded a proposer at all ("Other", which says nothing about staleness). tests/ suggestionSource.test.ts pins the invariant against the real delete path.

Everything the tray derives from those columns lives in app/orchestrator/suggestions.ts (pure: grouping, newest-first ordering, the 7-day staleness threshold, the new-since counts) and is rendered by the one SuggestionGroup component, which both tray surfaces — the list column and the board's Suggested column — mount. The "new since you last looked" mark is a per-project timestamp in the localStorage prefs blob (usePrefs), advanced by an IntersectionObserver on the tray container and frozen for the life of the mounted tray so the pills don't vanish out from under the eye that's reading them.

The transcript side of the same link: a suggest_task call's persisted tool message carries a suggestion: { title, taskId } (ToolSuggestion in lib/types.ts). The title comes from the tool input via describeToolUse(); the id can't — the MCP handler never sees the tool_use id it's answering — so it travels in the tool result text (formatSuggestedTaskText / parseSuggestedTaskId in lib/agents/shared.ts, the one place both ends of the format live) and each driver lifts it onto its tool_result event as taskId: the Claude driver for tool_use ids it saw go to suggest_task, the Codex normalizer for orchestrator-server MCP items (which also render with Claude's title instead of the generic ⚙ server: tool line), the e2e mock driver directly. The runner merges it into the row's ToolData and publishes suggested right then rather than at turn end — so the task list refreshes while the turn is live, and a driver that never emits a trailing suggested event (Codex) still updates the tray.

SuggestionChip (app/orchestrator/Transcript.tsx) renders that card against the project's task list, which reaches it through SuggestionContext rather than props: MessageView is memoized on the message alone, and only the chips should re-render when tasks change. The context value (useSuggestionActions, provided once in Orchestrator.tsx) is a by-id index of the selected project's tasks plus the tray's handlers behind ref-backed stable functions. A chip resolves its state from the index — in the tray, accepted (with status), gone — with two guards against reading "missing" as "dismissed": ready is false until the project's tasks have loaded, and a card filed within the last half-minute by a still-running turn reads "adding…" (the suggested event's reload is in flight) rather than dismissed. Inline rename is a PATCH /api/tasks/[id] of title via renameTask() in useOrchestrator, optimistic so chip and tray move together. suggestionBatches() (suggestions.ts, pinned by tests/suggestionChips.test.ts) keys the "Suggested this session" block by the last message of any turn — user message to user message, /clear breaks and queued bubbles included — that filed two or more suggestions.

Adding a third agent (e.g. Gemini, Cursor)

Implement the AgentDriver interface in lib/agents/<id>/driver.ts (runTurn() is the only required method — the one-shot helpers are optional and fall back to the utility agent), register it in lib/agents/registry.ts, and ship its CLI in the Dockerfile (installed on PATH next to claude / codex). No edits to the runner, routes, recap/refresh jobs, or UI data flow — the capability descriptor drives the pickers, the /api/agents/[id]/* routes are generic, and getDriver(task.agent) resolves it everywhere. The driver contract test (tests/agentDriver.test.ts) and the event-mapping test (tests/codexEvents.test.ts) are the templates for pinning a new driver to the same StreamEvent contract.

Everything else, by module

  • lib/db.ts — SQLite schema, migrations, seed. lib/store.ts — typed queries for projects / tasks / messages / summaries / sessions.
  • lib/git.ts — per-task worktrees/branches, diffs, and merging (mergeTask(), plus prepareWorktreeMerge() / completeWorktreeMerge() / abortWorktreeMerge() for AI/manual conflict resolution; worktreeSyncStatus() / fastForwardWorktree() to catch a stale branch up to base).
  • lib/services.ts — the managed-services supervisor: starts/stops/restarts a project's configured dev/setup/test commands as detached process-group children owned by the server (not a turn or a tab), captures their stdout/stderr into a per-service ring buffer, and publishes status/log events over SSE. State lives on globalThis (survives HMR), like lib/events.ts. Each project gets a stable PORT (projects.port, deterministic from ORCH_SERVICE_PORT_BASE) injected into every service's env and the PTY shell. On by default (ORCH_FEATURE_SERVICES=0 disables): the registry is persisted (services table) and server.js restores + auto-restarts managed services on boot — first reaping any process group a crashed server left orphaned (the spawn pid is persisted per row; the reaper verifies the group still runs the service's command before SIGKILLing it, so a recycled pid is never killed by mistake), and probing the port first so a conflict with an unmanaged process surfaces as a readable error on the service instead of an EADDRINUSE crash loop. A clean process exit SIGKILLs every managed group on the way out. Running services don't block idle-stop (/api/instance/idle reports runningServices informationally — sleeping is safe because boot restore relaunches them). Public hostnames are a separate opt-in (ORCH_SERVICE_HOSTS): each service then gets a stable <slug>--<appHost> hostname with per-service visibility (private / shared-link / public), dispatched through the reverse-proxy router in lib/service-router.mjs (WebSocket/HMR passthrough included), with the pure hostname/token helpers in lib/service-host.mjs.
  • lib/contextRefresh.ts — "Refresh with AI" as a detached background job (a multi-minute draft never holds an HTTP request open across a tunnel): startRefreshJob() seeds the utility agent with recent git activity, runs draftProjectContext() in the repo (read-only), and persists the result for the client to poll via GET /api/projects/[id]/refresh-context. The draft is for the user to review — never auto-saved. In-flight work marks the instance busy (lib/idle.ts) so an idle daemon won't stop the container mid-refresh.
  • lib/recap.ts — "where you left off" staleness/activity logic + background sweep.
  • app/Orchestrator.tsx — the dark mission-control client UI (projects rail · task list · live session, the session split into transcript + SessionRail DIFF/PREVIEW/CONTEXT tabs); one EventSource per selected task renders from server events, so a reload, sleep, or task switch mid-turn just catches up.
  • app/Terminal.tsx + pty-server.js — xterm.js ↔ same-origin /pty WebSocket (proxied by server.js) ↔ node-pty sidecar bound to 127.0.0.1.

Where data lives

What Where
Projects, tasks, transcripts, summaries, session index orchestrator.db (SQLite) in ORCH_DB_DIR, default ~/.zen-orchestrator
Per-task git worktrees ORCH_WORKTREES_DIR, default ~/.agent-orchestrator/worktrees — deliberately outside every repo
Cloned project repos ORCH_PROJECTS_DIR, default ~/projects
Your apps' actual code each project's working directory — never inside Operator's own tree
Claude Code's raw session logs ~/.claude/projects/... (managed by Claude Code)

Stack: Next.js (App Router) + TypeScript · React 19 · better-sqlite3 · @anthropic-ai/claude-agent-sdk · xterm.js + node-pty sidecar · streaming over SSE.