Skip to content
Merged
Show file tree
Hide file tree
Changes from 4 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
2,661 changes: 2,661 additions & 0 deletions conductor/Cargo.lock

Large diffs are not rendered by default.

37 changes: 37 additions & 0 deletions conductor/Cargo.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
[package]
name = "iii-conductor"
version = "0.1.0"
edition = "2021"
description = "Multi-agent fan-out + verifier-gated merge worker for iii-engine"
license = "Apache-2.0"
authors = ["Rohit Ghumare <ghumare64@gmail.com>"]
repository = "https://github.com/iii-hq/workers"
homepage = "https://github.com/iii-hq/workers"
rust-version = "1.85"
keywords = ["iii-engine", "agents", "orchestration", "ai", "worker"]
categories = ["command-line-utilities"]
publish = false

[[bin]]
name = "iii-conductor"
path = "src/main.rs"

[lib]
name = "iii_conductor"
path = "src/lib.rs"

[dependencies]
iii-sdk = "=0.11.3"
tokio = { version = "1", features = ["macros", "rt-multi-thread", "io-util", "sync", "time", "process", "fs", "signal"] }
serde = { version = "1", features = ["derive"] }
serde_json = "1"
clap = { version = "4", features = ["derive"] }
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
anyhow = "1"
uuid = { version = "1", features = ["v4"] }
async-trait = "0.1"
futures = "0.3"

[dev-dependencies]
tokio = { version = "1", features = ["macros", "rt-multi-thread", "test-util"] }
137 changes: 137 additions & 0 deletions conductor/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,137 @@
# conductor

Multi-agent fan-out + verifier-gated merge for the iii engine. Runs a task
across N agent CLIs in parallel, each in its own git worktree, runs a
configurable list of verifier gates against each result, and picks a winning
diff.

Designed to be transport-agnostic: every function the worker registers shows
up over MCP (via `iii-mcp`) and A2A (via `iii-a2a`) without any extra wiring,
subject to the RBAC policy on `iii-worker-manager`.

## Functions

| Function | Input | Output |
|---|---|---|
| `conductor::dispatch` | `{ task, agents[], gates?[], cwd, timeout_ms? }` | `{ ok, run_id, agents, gates }` |
| `conductor::status` | `{ run_id }` | `RunState \| null` |
| `conductor::list` | `{}` | `RunState[]` |
| `conductor::merge` | `{ run_id }` | `MergeResult` |

### `AgentSpec`

```jsonc
{
"kind": "claude", // claude | codex | gemini | aider | cursor | amp | opencode | qwen | remote
"bin": "claude", // optional, override the default CLI binary
"args": ["--print", "..."], // optional, override the default arg vector
"function_id": "a2a.foo::write_code", // required when kind=remote
"prompt": "Add /healthz", // optional, defaults to the dispatch task
"worktree": false // pass --worktree to the CLI when supported
}
```

For `kind: "remote"`, the conductor still creates a worktree and passes
that worktree path as `cwd` in the trigger payload. Remote handlers that
write to `cwd` produce a real diff and can win the merge. This is how
external A2A agents (registered via `iii-a2a-client`) and remote MCP tool
servers (via `iii-mcp-client`) participate in a fan-out on equal footing
with local CLI agents.

### `GateSpec`

```jsonc
{ "function_id": "verify::tests", "description": "unit tests pass" }
```

Each gate is just an iii function the conductor invokes against the agent's
worktree. The gate returns `{ ok: boolean, reason?: string }`. Recommended
gates: `verify::tests`, `verify::lint`, `verify::types`, `verify::build`,
`verify::diff_clean`. The `eval`, `guardrails`, and `proof` workers in this
repo register suitable gate functions.

Gate results are stored as an ordered `Vec<GateRunResult>`, not a map. The
same `function_id` can appear more than once with different descriptions
(e.g. `verify::tests` for unit and again for integration), and the original
order is preserved for `conductor::status`.

## How a run flows

1. `dispatch` records a seed `RunState` under `state::set` scope
`conductor`, key `runs::<run_id>`.
2. For each agent (local or remote), conductor creates a git worktree
(`conductor/<run_id>/<i>-<kind>`) off the current branch under
`~/.iii/conductor/worktrees/`.
3. Local agents are spawned via `tokio::process::Command` inside their
worktree. Remote agents are reached via
`iii.trigger(spec.function_id, { task, cwd: <worktree path> })`.
4. As each agent completes, gates run in series against its worktree and
the run is written back to `state::set`. Mid-run crashes preserve the
transitions of every agent that already finished.
5. `merge` picks the eligible agent with the **smallest `finished_at`**
(true "first finished agent wins" semantics). An agent is eligible
when `status == Finished`, `diff` is non-empty, and every gate passed.
Losers' worktrees are removed; the winner's worktree and branch survive
for review.

## Example

```bash
iii trigger conductor::dispatch \
--payload '{
"task": "Add a /healthz endpoint to the public API",
"cwd": "/abs/path/to/repo",
"agents": [
{ "kind": "claude" },
{ "kind": "codex" },
{ "kind": "remote", "function_id": "a2a.codex_web::write_code" }
],
"gates": [
{ "function_id": "verify::tests" },
{ "function_id": "verify::types" },
{ "function_id": "verify::build" }
],
"timeout_ms": 600000
}'

iii trigger conductor::merge --payload '{ "run_id": "<id from dispatch>" }'
```

## RBAC

This worker registers its functions with `metadata.public = true`. To expose
them over MCP or A2A, list them in `iii-worker-manager`'s `expose_functions`:

```yaml
workers:
- name: iii-worker-manager
config:
rbac:
auth_function_id: myproject::auth
expose_functions:
- match("conductor::*")
- metadata:
public: true
- name: iii-mcp
- name: iii-a2a
- name: conductor
```

## CLI flags

```text
--engine-url <URL> WebSocket URL of the iii engine (default ws://localhost:49134)
--debug Verbose logging
```

## Dependencies

- `git` on PATH (worktree creation, diffs).
- The agent CLIs you list in `agents[]` must be installed on PATH for local
kinds, or registered with the engine for `kind: "remote"`.
- `state::set` / `state::get` / `state::list` / `state::delete` must be
registered (the engine ships these by default).

## Layout

Worktrees land under `~/.iii/conductor/worktrees/`.
7 changes: 7 additions & 0 deletions conductor/iii.worker.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
iii: v1
name: conductor
language: rust
deploy: binary
manifest: Cargo.toml
bin: iii-conductor
description: Multi-agent fan-out + verifier-gated merge worker. Dispatches a task across N agent CLIs in parallel, runs verifier gates per result, picks a winning diff.
89 changes: 89 additions & 0 deletions conductor/src/agents.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
use std::path::Path;

use crate::git::{run_cmd, CmdResult};
use crate::types::{AgentKind, AgentSpec};

fn default_bin(kind: AgentKind) -> Option<&'static str> {
match kind {
AgentKind::Claude => Some("claude"),
AgentKind::Codex => Some("codex"),
AgentKind::Gemini => Some("gemini"),
AgentKind::Aider => Some("aider"),
AgentKind::Cursor => Some("cursor-agent"),
AgentKind::Amp => Some("amp"),
AgentKind::Opencode => Some("opencode"),
AgentKind::Qwen => Some("qwen"),
AgentKind::Remote => None,
}
}

fn build_args(spec: &AgentSpec) -> Vec<String> {
if let Some(args) = &spec.args {
if !args.is_empty() {
return args.clone();
}
}
let prompt = spec.prompt.clone().unwrap_or_default();
match spec.kind {
AgentKind::Claude => {
let mut a = vec!["--print".to_string()];
if spec.worktree {
a.push("--worktree".to_string());
}
if !prompt.is_empty() {
a.push(prompt);
}
a
}
AgentKind::Codex => {
if prompt.is_empty() {
vec!["exec".to_string()]
} else {
vec!["exec".to_string(), prompt]
}
}
AgentKind::Gemini => {
if prompt.is_empty() {
Vec::new()
} else {
vec!["--prompt".to_string(), prompt]
}
}
_ => {
if prompt.is_empty() {
Vec::new()
} else {
vec![prompt]
}
}
}
}

pub async fn run_local_agent(spec: &AgentSpec, cwd: &Path, timeout_ms: Option<u64>) -> CmdResult {
if spec.kind == AgentKind::Remote {
return CmdResult {
ok: false,
code: None,
stdout: String::new(),
stderr: "remote agent must be invoked through iii.trigger".to_string(),
};
}
let bin = match spec
.bin
.clone()
.or_else(|| default_bin(spec.kind).map(String::from))
{
Some(b) => b,
None => {
return CmdResult {
ok: false,
code: None,
stdout: String::new(),
stderr: format!("no binary configured for agent kind {:?}", spec.kind),
};
}
};
let args = build_args(spec);
let arg_refs: Vec<&str> = args.iter().map(String::as_str).collect();
run_cmd(cwd, &bin, &arg_refs, timeout_ms).await
}
Loading
Loading