feat(coding-agent): add swarm extension for multi-agent pipeline orchestration
Adds a standalone extension that orchestrates multi-agent workflows defined in YAML. Supports pipeline (iterative), parallel (fan-in/out), sequential, and arbitrary DAG execution patterns. Each agent runs as a full oh-my-pi subagent with complete tool access. Agents communicate through the shared workspace filesystem. The orchestrator handles lifecycle, dependency ordering, and state tracking. Key components: - YAML schema parser with validation and cycle detection - DAG builder with topological sort into execution waves - Pipeline controller with iteration loop and wave execution - Filesystem state persistence for monitoring and resumability - Standalone runner (run-pipeline.ts) for long-running unattended work - TUI integration via /swarm command Designed for any task type: research, coding, data processing, content creation, analysis workflows, or any multi-step objective that benefits from specialized agents working in coordination.
This commit is contained in:
@@ -0,0 +1 @@
|
||||
node_modules/
|
||||
@@ -0,0 +1,457 @@
|
||||
# Swarm Extension
|
||||
|
||||
Multi-agent orchestration for oh-my-pi. Define agent workflows in YAML — pipelines, parallel fan-outs, sequential chains, or any DAG — and run them unattended until completion.
|
||||
|
||||
Each agent is a full oh-my-pi subagent with access to every tool: bash, python, read, write, edit, grep, find, fetch, web_search, browser. The orchestrator manages lifecycle and ordering; agents communicate through the shared workspace filesystem.
|
||||
|
||||
Use it for anything: research pipelines, code generation, data processing, content creation, analysis workflows, CI-like automation — any multi-step task that benefits from specialized agents working in coordination.
|
||||
|
||||
## Setup
|
||||
|
||||
```bash
|
||||
cd packages/coding-agent/extensions/swarm
|
||||
bun install
|
||||
```
|
||||
|
||||
## Running
|
||||
|
||||
### Standalone (recommended for long-running work)
|
||||
|
||||
```bash
|
||||
# Foreground — runs until complete, no timeout:
|
||||
bun packages/coding-agent/extensions/swarm/run-pipeline.ts path/to/swarm.yaml
|
||||
|
||||
# Background — survives terminal close:
|
||||
nohup bun packages/coding-agent/extensions/swarm/run-pipeline.ts path/to/swarm.yaml \
|
||||
> pipeline.log 2>&1 & disown
|
||||
```
|
||||
|
||||
The standalone runner has no timeout. It runs iteration after iteration until the pipeline finishes or you kill it.
|
||||
|
||||
### Inside oh-my-pi (TUI)
|
||||
|
||||
Register the extension in your config (`~/.omp/config.json` or `.omp/config.json`):
|
||||
|
||||
```json
|
||||
{
|
||||
"extensions": ["packages/coding-agent/extensions/swarm"]
|
||||
}
|
||||
```
|
||||
|
||||
Then:
|
||||
|
||||
```
|
||||
/swarm run path/to/swarm.yaml
|
||||
/swarm status <name>
|
||||
/swarm help
|
||||
```
|
||||
|
||||
## Monitoring
|
||||
|
||||
State persists to `<workspace>/.swarm_<name>/` while the pipeline runs:
|
||||
|
||||
```
|
||||
.swarm_<name>/
|
||||
state/pipeline.json # Live pipeline + per-agent status
|
||||
logs/orchestrator.log # Wave transitions, iteration progress
|
||||
logs/<agent>.log # Per-agent timestamps and errors
|
||||
context/ # Agent session artifacts
|
||||
```
|
||||
|
||||
Check on a running pipeline:
|
||||
|
||||
```bash
|
||||
# Quick status
|
||||
cat workspace/.swarm_mypipeline/state/pipeline.json | python -m json.tool
|
||||
|
||||
# Watch the orchestrator log
|
||||
tail -f workspace/.swarm_mypipeline/logs/orchestrator.log
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## YAML Reference
|
||||
|
||||
Every swarm is a single YAML file with a top-level `swarm` key:
|
||||
|
||||
```yaml
|
||||
swarm:
|
||||
name: my-pipeline # Identifier (state stored in .swarm_<name>/)
|
||||
workspace: ./workspace # Working directory (relative to YAML file location)
|
||||
mode: pipeline # pipeline | parallel | sequential
|
||||
target_count: 10 # Iterations (pipeline mode only, default: 1)
|
||||
model: claude-opus-4-6 # Model for all agents (optional)
|
||||
|
||||
agents:
|
||||
first_agent:
|
||||
role: short-role-name
|
||||
task: |
|
||||
Full instructions for this agent.
|
||||
extra_context: |
|
||||
Optional additional system prompt text.
|
||||
reports_to:
|
||||
- downstream_agent
|
||||
waits_for:
|
||||
- upstream_agent
|
||||
```
|
||||
|
||||
### Top-Level Fields
|
||||
|
||||
| Field | Required | Default | Description |
|
||||
|-------|----------|---------|-------------|
|
||||
| `name` | yes | — | Pipeline identifier. State directory is `.swarm_<name>/` |
|
||||
| `workspace` | yes | — | Shared working directory. Relative paths resolve from YAML file location |
|
||||
| `mode` | no | `sequential` | Execution mode (see below) |
|
||||
| `target_count` | no | `1` | How many times to repeat the full pipeline. Only meaningful in `pipeline` mode |
|
||||
| `model` | no | session default | Model ID for all agents. Any omp-configured model works |
|
||||
|
||||
### Agent Fields
|
||||
|
||||
| Field | Required | Description |
|
||||
|-------|----------|-------------|
|
||||
| `role` | yes | Short role identifier — becomes the agent's system prompt |
|
||||
| `task` | yes | Complete instructions sent as user prompt. Use YAML `\|` for multi-line |
|
||||
| `extra_context` | no | Additional text appended to system prompt |
|
||||
| `reports_to` | no | List of agent names that depend on this agent |
|
||||
| `waits_for` | no | List of agent names this agent depends on |
|
||||
|
||||
### Execution Modes
|
||||
|
||||
**`pipeline`** — Repeat the full agent graph `target_count` times. Each iteration runs all waves in order. Use for accumulative work: "find 50 things, one per iteration."
|
||||
|
||||
**`sequential`** — Run agents once, chained by declaration order (unless explicit dependencies override). The default mode.
|
||||
|
||||
**`parallel`** — Run all agents simultaneously (unless explicit dependencies impose ordering).
|
||||
|
||||
### Dependency Resolution
|
||||
|
||||
The orchestrator builds a DAG from `waits_for` and `reports_to`, then groups agents into **waves** using topological sort. Agents in the same wave run in parallel; waves execute in sequence.
|
||||
|
||||
- `waits_for: [a, b]` — this agent won't start until both `a` and `b` finish
|
||||
- `reports_to: [x]` — equivalent to `x` having `waits_for: [this_agent]`
|
||||
- No explicit deps + `pipeline`/`sequential` mode — agents chain by YAML declaration order
|
||||
- No explicit deps + `parallel` mode — all agents run in one wave
|
||||
- Cycles are detected and rejected before execution
|
||||
|
||||
---
|
||||
|
||||
## Patterns
|
||||
|
||||
### Pipeline: Iterative Accumulation
|
||||
|
||||
Run the same agent chain N times. Each iteration builds on the previous one's output. Good for: research collection, data gathering, batch processing, iterative refinement.
|
||||
|
||||
```yaml
|
||||
swarm:
|
||||
name: research-collector
|
||||
workspace: ./workspace
|
||||
mode: pipeline
|
||||
target_count: 25
|
||||
model: claude-opus-4-6
|
||||
|
||||
agents:
|
||||
finder:
|
||||
role: researcher
|
||||
task: |
|
||||
Find ONE new source on the topic defined in workspace/topic.md.
|
||||
|
||||
1. Read processed.txt to see what's already been found
|
||||
2. Use web_search to find a new, high-quality source
|
||||
3. Append the URL to processed.txt
|
||||
4. Write the URL to signals/finder_out.txt: FOUND:<url>
|
||||
|
||||
analyzer:
|
||||
role: analyst
|
||||
task: |
|
||||
Read signals/finder_out.txt for the URL.
|
||||
Fetch the page and extract key findings.
|
||||
Read tracking/count.txt, increment it, write back.
|
||||
Write analysis to analyzed/item_<N>.md
|
||||
Write to signals/analyzer_out.txt: DONE:<N>
|
||||
|
||||
compiler:
|
||||
role: technical-writer
|
||||
task: |
|
||||
Read signals/analyzer_out.txt for the item number.
|
||||
Read analyzed/item_<N>.md.
|
||||
Append a summary to output/report.md under a new section.
|
||||
```
|
||||
|
||||
After 25 iterations: 25 sources found, analyzed, and compiled into a single report.
|
||||
|
||||
### Fan-In: Parallel Specialists
|
||||
|
||||
Multiple agents work independently, one synthesizer combines results. Good for: multi-perspective analysis, parallel code review, comprehensive audits.
|
||||
|
||||
```yaml
|
||||
swarm:
|
||||
name: codebase-audit
|
||||
workspace: ./workspace
|
||||
|
||||
agents:
|
||||
security:
|
||||
role: security-auditor
|
||||
task: |
|
||||
Audit all code in src/ for security vulnerabilities.
|
||||
Write findings to reports/security.md with severity ratings.
|
||||
reports_to:
|
||||
- lead
|
||||
|
||||
performance:
|
||||
role: performance-analyst
|
||||
task: |
|
||||
Profile and analyze src/ for performance bottlenecks.
|
||||
Write findings to reports/performance.md with benchmarks.
|
||||
reports_to:
|
||||
- lead
|
||||
|
||||
architecture:
|
||||
role: architecture-reviewer
|
||||
task: |
|
||||
Review src/ for architectural issues, coupling, and tech debt.
|
||||
Write findings to reports/architecture.md with refactoring suggestions.
|
||||
reports_to:
|
||||
- lead
|
||||
|
||||
lead:
|
||||
role: engineering-lead
|
||||
task: |
|
||||
Read all reports in reports/.
|
||||
Create a prioritized action plan in output/action_plan.md.
|
||||
Rank issues by impact and effort.
|
||||
waits_for:
|
||||
- security
|
||||
- performance
|
||||
- architecture
|
||||
```
|
||||
|
||||
Execution: security + performance + architecture run in parallel (wave 1), lead starts after all three complete (wave 2).
|
||||
|
||||
### Sequential Chain: Staged Handoff
|
||||
|
||||
Linear progression through distinct phases. Good for: content pipelines, multi-stage processing, review chains.
|
||||
|
||||
```yaml
|
||||
swarm:
|
||||
name: blog-post
|
||||
workspace: ./workspace
|
||||
mode: sequential
|
||||
|
||||
agents:
|
||||
researcher:
|
||||
role: researcher
|
||||
task: |
|
||||
Research the topic in topic.md using web_search.
|
||||
Write raw findings and source links to research/notes.md
|
||||
|
||||
writer:
|
||||
role: technical-writer
|
||||
task: |
|
||||
Read research/notes.md.
|
||||
Write a complete blog post draft to drafts/post.md.
|
||||
Include code examples where relevant.
|
||||
|
||||
editor:
|
||||
role: editor
|
||||
task: |
|
||||
Read drafts/post.md.
|
||||
Fix grammar, improve flow, tighten prose.
|
||||
Rewrite to drafts/post.md.
|
||||
|
||||
reviewer:
|
||||
role: senior-reviewer
|
||||
task: |
|
||||
Read drafts/post.md.
|
||||
Check technical accuracy against research/notes.md.
|
||||
Add an editorial note at top if issues found, otherwise
|
||||
copy to output/final.md.
|
||||
```
|
||||
|
||||
Execution: researcher -> writer -> editor -> reviewer, one after another.
|
||||
|
||||
### Diamond: Fan-Out Then Fan-In
|
||||
|
||||
One planner, parallel workers, one integrator. Good for: divide-and-conquer, modular code generation, multi-file refactors.
|
||||
|
||||
```yaml
|
||||
swarm:
|
||||
name: feature-implementation
|
||||
workspace: ./workspace
|
||||
|
||||
agents:
|
||||
planner:
|
||||
role: architect
|
||||
task: |
|
||||
Read the feature spec in spec.md.
|
||||
Break it into independent implementation tasks.
|
||||
Write the plan to plan.md with file assignments.
|
||||
reports_to:
|
||||
- api
|
||||
- ui
|
||||
- tests
|
||||
|
||||
api:
|
||||
role: backend-developer
|
||||
task: |
|
||||
Read plan.md for your assigned files.
|
||||
Implement the API layer. Write to src/api/.
|
||||
reports_to:
|
||||
- integrator
|
||||
|
||||
ui:
|
||||
role: frontend-developer
|
||||
task: |
|
||||
Read plan.md for your assigned files.
|
||||
Implement the UI components. Write to src/ui/.
|
||||
reports_to:
|
||||
- integrator
|
||||
|
||||
tests:
|
||||
role: test-engineer
|
||||
task: |
|
||||
Read plan.md for the full feature scope.
|
||||
Write integration tests to tests/.
|
||||
reports_to:
|
||||
- integrator
|
||||
|
||||
integrator:
|
||||
role: tech-lead
|
||||
task: |
|
||||
Read plan.md and review all code in src/ and tests/.
|
||||
Wire everything together. Fix any integration issues.
|
||||
Run the tests and fix failures.
|
||||
Write status to output/done.md.
|
||||
```
|
||||
|
||||
Execution: planner (wave 1) -> api + ui + tests in parallel (wave 2) -> integrator (wave 3).
|
||||
|
||||
### Hybrid: Mixed Dependencies
|
||||
|
||||
Any DAG is valid. Combine patterns freely.
|
||||
|
||||
```yaml
|
||||
swarm:
|
||||
name: data-pipeline
|
||||
workspace: ./workspace
|
||||
mode: pipeline
|
||||
target_count: 10
|
||||
|
||||
agents:
|
||||
scraper_a:
|
||||
role: web-scraper
|
||||
task: |
|
||||
Scrape data source A. Write to raw/source_a.json
|
||||
reports_to:
|
||||
- transformer
|
||||
|
||||
scraper_b:
|
||||
role: web-scraper
|
||||
task: |
|
||||
Scrape data source B. Write to raw/source_b.json
|
||||
reports_to:
|
||||
- transformer
|
||||
|
||||
transformer:
|
||||
role: data-engineer
|
||||
task: |
|
||||
Read raw/source_a.json and raw/source_b.json.
|
||||
Clean, normalize, merge. Write to processed/merged.json
|
||||
reports_to:
|
||||
- loader
|
||||
- validator
|
||||
|
||||
validator:
|
||||
role: qa-analyst
|
||||
task: |
|
||||
Read processed/merged.json.
|
||||
Validate schema, check for anomalies.
|
||||
Write report to qa/validation.md
|
||||
|
||||
loader:
|
||||
role: data-engineer
|
||||
task: |
|
||||
Read processed/merged.json.
|
||||
Append to output/dataset.jsonl
|
||||
```
|
||||
|
||||
Execution per iteration: scraper_a + scraper_b (wave 1) -> transformer (wave 2) -> loader + validator (wave 3).
|
||||
|
||||
---
|
||||
|
||||
## Writing Agent Tasks
|
||||
|
||||
### What Agents Can Do
|
||||
|
||||
Each agent is a full oh-my-pi session. It can:
|
||||
|
||||
- **bash/python**: Run commands, scripts, install packages, process data
|
||||
- **read/write/edit**: Create and modify files in the workspace
|
||||
- **grep/find**: Search the workspace (or anywhere on disk)
|
||||
- **web_search**: Search the internet (via configured provider)
|
||||
- **fetch**: Download web pages, APIs, documents
|
||||
- **browser**: Navigate websites, scrape dynamic content, take screenshots
|
||||
|
||||
### Inter-Agent Communication
|
||||
|
||||
The orchestrator starts and stops agents in the right order. It does **not** pass data between them. Agents communicate through files in the shared workspace.
|
||||
|
||||
Design your own protocol. Common patterns:
|
||||
|
||||
**Signal files** — lightweight status flags an agent writes when done:
|
||||
```
|
||||
signals/finder_out.txt -> "FOUND:https://example.com"
|
||||
signals/analyzer_out.txt -> "DONE:42"
|
||||
signals/reviewer_out.txt -> "APPROVED" or "REJECTED:reason"
|
||||
```
|
||||
|
||||
**Structured output** — detailed results other agents read:
|
||||
```
|
||||
analyzed/item_1.md -> Full analysis document
|
||||
results/report.json -> Machine-readable data
|
||||
output/final.docx -> Accumulated deliverable
|
||||
```
|
||||
|
||||
**Tracking files** — prevent duplicate work across pipeline iterations:
|
||||
```
|
||||
processed.txt -> Items already handled (one per line)
|
||||
tracking/count.txt -> Current item counter
|
||||
tracking/status.json -> Cumulative state
|
||||
```
|
||||
|
||||
### Tips for Reliable Agents
|
||||
|
||||
- **Be explicit about paths.** Agents start fresh each iteration — they don't remember previous runs. Tell them exactly where to read input and write output.
|
||||
- **Check existing state.** In pipeline mode, tell agents to read tracking files before doing work: "Read processed.txt to avoid duplicates."
|
||||
- **Use numbered outputs.** `item_1.md`, `item_2.md` etc. so iterations don't clobber each other.
|
||||
- **Handle failure.** Tell agents what to do when things go wrong: "If the source lacks depth, write SKIP to signals/out.txt and explain why."
|
||||
- **Keep signal files simple.** One line, parseable format. Complex data goes in structured output files.
|
||||
- **Scope the task tightly.** An agent that tries to do five things will do zero well. One clear objective per agent.
|
||||
|
||||
---
|
||||
|
||||
## Models
|
||||
|
||||
Any model configured in omp works. Set it in the YAML:
|
||||
|
||||
```yaml
|
||||
swarm:
|
||||
model: claude-opus-4-6
|
||||
```
|
||||
|
||||
Or omit `model` to use your session's default. Check `packages/ai/src/models.json` for available model IDs.
|
||||
|
||||
---
|
||||
|
||||
## Architecture
|
||||
|
||||
```
|
||||
extension.ts TUI entry point (registers /swarm command)
|
||||
run-pipeline.ts Standalone runner (no TUI, no timeout)
|
||||
swarm/
|
||||
schema.ts YAML parsing + validation
|
||||
dag.ts Dependency graph, cycle detection, topological sort
|
||||
executor.ts Spawns agents via oh-my-pi's runSubprocess
|
||||
pipeline.ts Iteration loop + wave controller
|
||||
state.ts Filesystem state persistence
|
||||
render.ts Progress display formatting
|
||||
```
|
||||
@@ -0,0 +1,15 @@
|
||||
{
|
||||
"lockfileVersion": 1,
|
||||
"configVersion": 1,
|
||||
"workspaces": {
|
||||
"": {
|
||||
"name": "omp-extension-swarm",
|
||||
"dependencies": {
|
||||
"yaml": "^2.7.0",
|
||||
},
|
||||
},
|
||||
},
|
||||
"packages": {
|
||||
"yaml": ["yaml@2.8.2", "", { "bin": { "yaml": "bin.mjs" } }, "sha512-mplynKqc1C2hTVYxd0PU2xQAc22TI1vShAYGksCCfxbn/dFwnHTNi1bvYsBTkhdUNtGIf5xNOg938rrSSYvS9A=="],
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,276 @@
|
||||
/**
|
||||
* Swarm Extension — Multi-agent pipeline orchestration from YAML definitions.
|
||||
*
|
||||
* Registers:
|
||||
* - /swarm run <file.yaml> — Execute a swarm pipeline
|
||||
* - /swarm status — Show current pipeline status
|
||||
*
|
||||
* Usage: Add this extension's directory to your extensions config,
|
||||
* then use /swarm in any oh-my-pi session.
|
||||
*/
|
||||
import * as path from "node:path";
|
||||
import * as fs from "node:fs/promises";
|
||||
import type { ExtensionAPI, ExtensionCommandContext } from "@oh-my-pi/pi-coding-agent";
|
||||
import { parseSwarmYaml, validateSwarmDefinition, type SwarmDefinition } from "./swarm/schema";
|
||||
import { buildDependencyGraph, detectCycles, buildExecutionWaves } from "./swarm/dag";
|
||||
import { StateTracker } from "./swarm/state";
|
||||
import { PipelineController } from "./swarm/pipeline";
|
||||
import { renderSwarmProgress } from "./swarm/render";
|
||||
|
||||
export default function swarmExtension(pi: ExtensionAPI): void {
|
||||
pi.setLabel("Swarm Orchestrator");
|
||||
|
||||
pi.registerCommand("swarm", {
|
||||
description: "Run a multi-agent swarm pipeline from YAML",
|
||||
getArgumentCompletions: (prefix) => {
|
||||
const subcommands = ["run", "status", "help"];
|
||||
if (!prefix) return subcommands.map((s) => ({ label: s, value: s }));
|
||||
return subcommands
|
||||
.filter((s) => s.startsWith(prefix))
|
||||
.map((s) => ({ label: s, value: s }));
|
||||
},
|
||||
handler: async (args: string, ctx: ExtensionCommandContext) => {
|
||||
const parts = args.trim().split(/\s+/);
|
||||
const subcommand = parts[0] ?? "help";
|
||||
|
||||
switch (subcommand) {
|
||||
case "run": {
|
||||
const yamlPath = parts[1];
|
||||
if (!yamlPath) {
|
||||
ctx.ui.notify("Usage: /swarm run <path/to/pipeline.yaml>", "error");
|
||||
return;
|
||||
}
|
||||
await handleRun(yamlPath, ctx, pi);
|
||||
return;
|
||||
}
|
||||
case "status": {
|
||||
await handleStatus(parts[1], ctx);
|
||||
return;
|
||||
}
|
||||
case "help":
|
||||
default:
|
||||
ctx.ui.notify(
|
||||
[
|
||||
"Swarm — multi-agent pipeline orchestrator",
|
||||
"",
|
||||
" /swarm run <file.yaml> Run a pipeline",
|
||||
" /swarm status [name] Show pipeline status",
|
||||
" /swarm help Show this help",
|
||||
].join("\n"),
|
||||
"info",
|
||||
);
|
||||
return;
|
||||
}
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// /swarm run
|
||||
// ============================================================================
|
||||
|
||||
async function handleRun(yamlPath: string, ctx: ExtensionCommandContext, pi: ExtensionAPI): Promise<void> {
|
||||
// 1. Resolve and read YAML
|
||||
const resolvedPath = path.isAbsolute(yamlPath) ? yamlPath : path.resolve(ctx.cwd, yamlPath);
|
||||
|
||||
let content: string;
|
||||
try {
|
||||
content = await Bun.file(resolvedPath).text();
|
||||
} catch {
|
||||
ctx.ui.notify(`Cannot read file: ${resolvedPath}`, "error");
|
||||
return;
|
||||
}
|
||||
|
||||
// 2. Parse YAML
|
||||
let def: SwarmDefinition;
|
||||
try {
|
||||
def = parseSwarmYaml(content);
|
||||
} catch (err) {
|
||||
ctx.ui.notify(`YAML error: ${err instanceof Error ? err.message : String(err)}`, "error");
|
||||
return;
|
||||
}
|
||||
|
||||
// 3. Validate
|
||||
const validationErrors = validateSwarmDefinition(def);
|
||||
if (validationErrors.length > 0) {
|
||||
ctx.ui.notify(`Validation errors:\n${validationErrors.map((e) => ` - ${e}`).join("\n")}`, "error");
|
||||
return;
|
||||
}
|
||||
|
||||
// 4. Build DAG
|
||||
const deps = buildDependencyGraph(def);
|
||||
const cycleNodes = detectCycles(deps);
|
||||
if (cycleNodes) {
|
||||
ctx.ui.notify(`Cycle detected in agent dependencies: [${cycleNodes.join(", ")}]`, "error");
|
||||
return;
|
||||
}
|
||||
const waves = buildExecutionWaves(deps);
|
||||
|
||||
// 5. Resolve workspace (relative to YAML file location)
|
||||
const workspace = path.isAbsolute(def.workspace)
|
||||
? def.workspace
|
||||
: path.resolve(path.dirname(resolvedPath), def.workspace);
|
||||
|
||||
// Ensure workspace exists
|
||||
await fs.mkdir(workspace, { recursive: true });
|
||||
|
||||
// 6. Initialize state tracker
|
||||
const stateTracker = new StateTracker(workspace, def.name);
|
||||
await stateTracker.init([...def.agents.keys()], def.targetCount, def.mode);
|
||||
|
||||
// 7. Log start
|
||||
const agentList = [...def.agents.keys()].join(", ");
|
||||
const waveDesc = waves.map((w, i) => `wave ${i + 1}: [${w.join(", ")}]`).join("; ");
|
||||
pi.logger.debug("Swarm starting", {
|
||||
name: def.name,
|
||||
mode: def.mode,
|
||||
agents: agentList,
|
||||
waves: waveDesc,
|
||||
workspace,
|
||||
});
|
||||
|
||||
ctx.ui.notify(
|
||||
`Starting swarm '${def.name}': ${def.agents.size} agents, ${waves.length} waves, ${def.targetCount} iteration(s)`,
|
||||
"info",
|
||||
);
|
||||
|
||||
// 8. Set up progress widget
|
||||
const widgetKey = `swarm-${def.name}`;
|
||||
const updateWidget = () => {
|
||||
const lines = renderSwarmProgress(stateTracker.state);
|
||||
ctx.ui.setWidget(widgetKey, lines);
|
||||
};
|
||||
updateWidget();
|
||||
|
||||
// 9. Resolve infrastructure for agent execution
|
||||
let authStorage: Awaited<ReturnType<typeof pi.pi.discoverAuthStorage>> | undefined;
|
||||
try {
|
||||
authStorage = await pi.pi.discoverAuthStorage();
|
||||
} catch {
|
||||
// Let runSubprocess discover auth per-agent as fallback
|
||||
}
|
||||
|
||||
// 10. Run pipeline
|
||||
const controller = new PipelineController(def, waves, stateTracker);
|
||||
|
||||
const result = await controller.run({
|
||||
workspace,
|
||||
onProgress: () => updateWidget(),
|
||||
authStorage,
|
||||
modelRegistry: ctx.modelRegistry,
|
||||
settings: pi.pi.settings,
|
||||
});
|
||||
|
||||
// 11. Clear widget and show summary
|
||||
ctx.ui.setWidget(widgetKey, undefined);
|
||||
|
||||
const elapsed = stateTracker.state.completedAt
|
||||
? formatDuration(stateTracker.state.completedAt - stateTracker.state.startedAt)
|
||||
: "unknown";
|
||||
|
||||
const summaryParts = [
|
||||
`Swarm '${def.name}' ${result.status}`,
|
||||
`${result.iterations}/${def.targetCount} iterations`,
|
||||
`elapsed: ${elapsed}`,
|
||||
];
|
||||
|
||||
if (result.errors.length > 0) {
|
||||
summaryParts.push(`${result.errors.length} error(s)`);
|
||||
}
|
||||
|
||||
const summaryType = result.status === "completed" ? "info" : "error";
|
||||
ctx.ui.notify(summaryParts.join(" | "), summaryType);
|
||||
|
||||
// Log errors
|
||||
if (result.errors.length > 0) {
|
||||
pi.logger.warn("Swarm completed with errors", { errors: result.errors });
|
||||
}
|
||||
|
||||
// 12. Send summary to the conversation so the LLM knows what happened
|
||||
const summaryMessage = buildSummaryMessage(def, result, stateTracker, workspace);
|
||||
pi.sendMessage(
|
||||
{
|
||||
customType: "swarm-result",
|
||||
content: [{ type: "text", text: summaryMessage }],
|
||||
display: true,
|
||||
details: {
|
||||
swarmName: def.name,
|
||||
status: result.status,
|
||||
iterations: result.iterations,
|
||||
errorCount: result.errors.length,
|
||||
},
|
||||
},
|
||||
{ triggerTurn: false },
|
||||
);
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// /swarm status
|
||||
// ============================================================================
|
||||
|
||||
async function handleStatus(name: string | undefined, ctx: ExtensionCommandContext): Promise<void> {
|
||||
if (!name) {
|
||||
ctx.ui.notify("Usage: /swarm status <name> (reads .swarm_<name>/state/pipeline.json from cwd)", "info");
|
||||
return;
|
||||
}
|
||||
|
||||
const stateTracker = new StateTracker(ctx.cwd, name);
|
||||
const state = await stateTracker.load();
|
||||
if (!state) {
|
||||
ctx.ui.notify(`No state found for swarm '${name}' in ${ctx.cwd}`, "error");
|
||||
return;
|
||||
}
|
||||
|
||||
const lines = renderSwarmProgress(state);
|
||||
ctx.ui.notify(lines.join("\n"), "info");
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Helpers
|
||||
// ============================================================================
|
||||
|
||||
function buildSummaryMessage(
|
||||
def: SwarmDefinition,
|
||||
result: { status: string; iterations: number; errors: string[] },
|
||||
stateTracker: StateTracker,
|
||||
workspace: string,
|
||||
): string {
|
||||
const lines: string[] = [];
|
||||
lines.push(`## Swarm Pipeline: ${def.name}`);
|
||||
lines.push("");
|
||||
lines.push(`- **Status**: ${result.status}`);
|
||||
lines.push(`- **Mode**: ${def.mode}`);
|
||||
lines.push(`- **Iterations**: ${result.iterations}/${def.targetCount}`);
|
||||
lines.push(`- **Workspace**: ${workspace}`);
|
||||
lines.push(`- **State dir**: ${stateTracker.swarmDir}`);
|
||||
lines.push("");
|
||||
|
||||
lines.push("### Agent Results");
|
||||
lines.push("");
|
||||
for (const [name, agent] of Object.entries(stateTracker.state.agents)) {
|
||||
const duration =
|
||||
agent.startedAt && agent.completedAt
|
||||
? formatDuration(agent.completedAt - agent.startedAt)
|
||||
: "n/a";
|
||||
lines.push(`- **${name}**: ${agent.status} (${duration})${agent.error ? ` — ${agent.error}` : ""}`);
|
||||
}
|
||||
|
||||
if (result.errors.length > 0) {
|
||||
lines.push("");
|
||||
lines.push("### Errors");
|
||||
lines.push("");
|
||||
for (const error of result.errors) {
|
||||
lines.push(`- ${error}`);
|
||||
}
|
||||
}
|
||||
|
||||
return lines.join("\n");
|
||||
}
|
||||
|
||||
function formatDuration(ms: number): string {
|
||||
if (ms < 1000) return `${ms}ms`;
|
||||
if (ms < 60_000) return `${(ms / 1000).toFixed(1)}s`;
|
||||
const mins = Math.floor(ms / 60_000);
|
||||
const secs = Math.floor((ms % 60_000) / 1000);
|
||||
return `${mins}m${secs}s`;
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
{
|
||||
"name": "omp-extension-swarm",
|
||||
"version": "1.0.0",
|
||||
"type": "module",
|
||||
"pi": {
|
||||
"extensions": ["./extension.ts"]
|
||||
},
|
||||
"dependencies": {
|
||||
"yaml": "^2.7.0"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,105 @@
|
||||
/**
|
||||
* Direct pipeline runner — executes a swarm pipeline outside of the TUI.
|
||||
*
|
||||
* Usage: bun run-pipeline.ts <path-to-yaml>
|
||||
*/
|
||||
import * as path from "node:path";
|
||||
import * as fs from "node:fs/promises";
|
||||
import { parseSwarmYaml, validateSwarmDefinition } from "./swarm/schema";
|
||||
import { buildDependencyGraph, detectCycles, buildExecutionWaves } from "./swarm/dag";
|
||||
import { StateTracker } from "./swarm/state";
|
||||
import { PipelineController } from "./swarm/pipeline";
|
||||
import { renderSwarmProgress } from "./swarm/render";
|
||||
import { discoverAuthStorage } from "@oh-my-pi/pi-coding-agent";
|
||||
import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry";
|
||||
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
|
||||
|
||||
const yamlPath = process.argv[2];
|
||||
if (!yamlPath) {
|
||||
console.error("Usage: bun run-pipeline.ts <path-to-yaml>");
|
||||
process.exit(1);
|
||||
}
|
||||
|
||||
const resolvedPath = path.resolve(yamlPath);
|
||||
console.log(`Reading: ${resolvedPath}`);
|
||||
|
||||
const content = await Bun.file(resolvedPath).text();
|
||||
const def = parseSwarmYaml(content);
|
||||
|
||||
console.log(`Swarm: ${def.name}`);
|
||||
console.log(`Mode: ${def.mode}`);
|
||||
console.log(`Target count: ${def.targetCount}`);
|
||||
console.log(`Agents: ${[...def.agents.keys()].join(", ")}`);
|
||||
|
||||
// Validate
|
||||
const errors = validateSwarmDefinition(def);
|
||||
if (errors.length > 0) {
|
||||
console.error("Validation errors:", errors);
|
||||
process.exit(1);
|
||||
}
|
||||
|
||||
// Build DAG
|
||||
const deps = buildDependencyGraph(def);
|
||||
const cycles = detectCycles(deps);
|
||||
if (cycles) {
|
||||
console.error("Cycle detected:", cycles);
|
||||
process.exit(1);
|
||||
}
|
||||
const waves = buildExecutionWaves(deps);
|
||||
console.log(`Waves: ${waves.map((w, i) => `W${i + 1}:[${w.join(",")}]`).join(" -> ")}`);
|
||||
|
||||
// Resolve workspace
|
||||
const workspace = path.isAbsolute(def.workspace)
|
||||
? def.workspace
|
||||
: path.resolve(path.dirname(resolvedPath), def.workspace);
|
||||
|
||||
await fs.mkdir(workspace, { recursive: true });
|
||||
console.log(`Workspace: ${workspace}`);
|
||||
|
||||
// Initialize
|
||||
const stateTracker = new StateTracker(workspace, def.name);
|
||||
await stateTracker.init([...def.agents.keys()], def.targetCount, def.mode);
|
||||
|
||||
// Auth + settings
|
||||
const authStorage = await discoverAuthStorage();
|
||||
const modelRegistry = new ModelRegistry(authStorage);
|
||||
const settings = Settings.isolated();
|
||||
|
||||
// Progress display
|
||||
let lastProgressDump = 0;
|
||||
const PROGRESS_INTERVAL_MS = 5000;
|
||||
|
||||
// Run
|
||||
console.log("\n--- Pipeline starting ---\n");
|
||||
|
||||
const controller = new PipelineController(def, waves, stateTracker);
|
||||
const result = await controller.run({
|
||||
workspace,
|
||||
onProgress: () => {
|
||||
const now = Date.now();
|
||||
if (now - lastProgressDump > PROGRESS_INTERVAL_MS) {
|
||||
lastProgressDump = now;
|
||||
const lines = renderSwarmProgress(stateTracker.state);
|
||||
console.log(lines.join("\n"));
|
||||
console.log();
|
||||
}
|
||||
},
|
||||
authStorage,
|
||||
modelRegistry,
|
||||
settings,
|
||||
});
|
||||
|
||||
console.log("\n--- Pipeline finished ---\n");
|
||||
console.log(`Status: ${result.status}`);
|
||||
console.log(`Iterations completed: ${result.iterations}/${def.targetCount}`);
|
||||
if (result.errors.length > 0) {
|
||||
console.log(`Errors (${result.errors.length}):`);
|
||||
for (const err of result.errors) {
|
||||
console.log(` - ${err}`);
|
||||
}
|
||||
}
|
||||
console.log(`\nState saved to: ${stateTracker.swarmDir}`);
|
||||
|
||||
// Final state dump
|
||||
const lines = renderSwarmProgress(stateTracker.state);
|
||||
console.log(lines.join("\n"));
|
||||
@@ -0,0 +1,146 @@
|
||||
/**
|
||||
* Directed Acyclic Graph operations for swarm agent dependencies.
|
||||
*
|
||||
* Builds a dependency graph from waits_for / reports_to relationships,
|
||||
* detects cycles, and produces execution waves via topological sort.
|
||||
*/
|
||||
import type { SwarmDefinition } from "./schema";
|
||||
|
||||
/**
|
||||
* Build a dependency map: agent name → set of agents it depends on.
|
||||
*
|
||||
* Dependencies come from:
|
||||
* 1. Explicit `waits_for` declarations
|
||||
* 2. Implicit from `reports_to` (if A reports_to B, then B depends on A)
|
||||
* 3. For pipeline/sequential mode with no explicit deps: chain by YAML declaration order
|
||||
*/
|
||||
export function buildDependencyGraph(def: SwarmDefinition): Map<string, Set<string>> {
|
||||
const deps = new Map<string, Set<string>>();
|
||||
|
||||
for (const name of def.agents.keys()) {
|
||||
deps.set(name, new Set());
|
||||
}
|
||||
|
||||
// Explicit waits_for
|
||||
for (const [name, agent] of def.agents) {
|
||||
for (const dep of agent.waitsFor) {
|
||||
if (deps.has(dep)) {
|
||||
deps.get(name)!.add(dep);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// reports_to implies the target waits for the reporter
|
||||
for (const [name, agent] of def.agents) {
|
||||
for (const target of agent.reportsTo) {
|
||||
if (deps.has(target)) {
|
||||
deps.get(target)!.add(name);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// For pipeline/sequential with no explicit deps, chain by declaration order
|
||||
if ((def.mode === "pipeline" || def.mode === "sequential") && !hasExplicitDeps(deps)) {
|
||||
for (let i = 1; i < def.agentOrder.length; i++) {
|
||||
deps.get(def.agentOrder[i])!.add(def.agentOrder[i - 1]);
|
||||
}
|
||||
}
|
||||
|
||||
return deps;
|
||||
}
|
||||
|
||||
function hasExplicitDeps(deps: Map<string, Set<string>>): boolean {
|
||||
for (const s of deps.values()) {
|
||||
if (s.size > 0) return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/**
|
||||
* Detect cycles in the dependency graph.
|
||||
* Returns the names of agents involved in cycles, or null if acyclic.
|
||||
*/
|
||||
export function detectCycles(deps: Map<string, Set<string>>): string[] | null {
|
||||
// Kahn's algorithm: if topological sort doesn't include all nodes, cycles exist
|
||||
const inDegree = new Map<string, number>();
|
||||
const forward = new Map<string, string[]>(); // dependency → its dependents
|
||||
|
||||
for (const [node, nodeDeps] of deps) {
|
||||
inDegree.set(node, nodeDeps.size);
|
||||
for (const dep of nodeDeps) {
|
||||
const list = forward.get(dep) ?? [];
|
||||
list.push(node);
|
||||
forward.set(dep, list);
|
||||
}
|
||||
}
|
||||
|
||||
const queue: string[] = [];
|
||||
for (const [node, degree] of inDegree) {
|
||||
if (degree === 0) queue.push(node);
|
||||
}
|
||||
|
||||
const sorted: string[] = [];
|
||||
while (queue.length > 0) {
|
||||
const node = queue.shift()!;
|
||||
sorted.push(node);
|
||||
for (const dependent of forward.get(node) ?? []) {
|
||||
const newDegree = inDegree.get(dependent)! - 1;
|
||||
inDegree.set(dependent, newDegree);
|
||||
if (newDegree === 0) queue.push(dependent);
|
||||
}
|
||||
}
|
||||
|
||||
if (sorted.length < deps.size) {
|
||||
return [...deps.keys()].filter(k => !sorted.includes(k));
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Build execution waves from dependency graph via topological sort.
|
||||
*
|
||||
* Each wave contains agents whose dependencies are all in earlier waves.
|
||||
* Agents within a wave can execute in parallel.
|
||||
*/
|
||||
export function buildExecutionWaves(deps: Map<string, Set<string>>): string[][] {
|
||||
const waves: string[][] = [];
|
||||
const completed = new Set<string>();
|
||||
const remaining = new Set(deps.keys());
|
||||
|
||||
while (remaining.size > 0) {
|
||||
const wave: string[] = [];
|
||||
|
||||
for (const node of remaining) {
|
||||
const nodeDeps = deps.get(node)!;
|
||||
let ready = true;
|
||||
for (const dep of nodeDeps) {
|
||||
if (!completed.has(dep)) {
|
||||
ready = false;
|
||||
break;
|
||||
}
|
||||
}
|
||||
if (ready) {
|
||||
wave.push(node);
|
||||
}
|
||||
}
|
||||
|
||||
if (wave.length === 0) {
|
||||
throw new Error(
|
||||
`Deadlock: agents [${[...remaining].join(", ")}] cannot make progress. This indicates a bug in cycle detection.`,
|
||||
);
|
||||
}
|
||||
|
||||
// Sort for deterministic execution order
|
||||
wave.sort();
|
||||
|
||||
for (const node of wave) {
|
||||
remaining.delete(node);
|
||||
completed.add(node);
|
||||
}
|
||||
|
||||
waves.push(wave);
|
||||
}
|
||||
|
||||
return waves;
|
||||
}
|
||||
@@ -0,0 +1,108 @@
|
||||
/**
|
||||
* Swarm agent execution via oh-my-pi's subagent infrastructure.
|
||||
*
|
||||
* Wraps `runSubprocess` to spawn individual swarm agents with full tool access.
|
||||
* Each agent runs in the swarm workspace with its task instructions as the user prompt.
|
||||
*/
|
||||
import * as path from "node:path";
|
||||
import type { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage";
|
||||
import type { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry";
|
||||
import type { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
|
||||
import { runSubprocess } from "@oh-my-pi/pi-coding-agent/task/executor";
|
||||
import type { AgentDefinition, AgentProgress, AgentSource, SingleResult } from "@oh-my-pi/pi-coding-agent/task/types";
|
||||
import type { SwarmAgent } from "./schema";
|
||||
import type { StateTracker } from "./state";
|
||||
|
||||
export interface SwarmExecutorOptions {
|
||||
workspace: string;
|
||||
swarmName: string;
|
||||
iteration: number;
|
||||
modelOverride?: string;
|
||||
signal?: AbortSignal;
|
||||
onProgress?: (agentName: string, progress: AgentProgress) => void;
|
||||
authStorage?: AuthStorage;
|
||||
modelRegistry?: ModelRegistry;
|
||||
settings?: Settings;
|
||||
stateTracker: StateTracker;
|
||||
}
|
||||
|
||||
/**
|
||||
* Execute a single swarm agent as an oh-my-pi subagent.
|
||||
*
|
||||
* The agent receives:
|
||||
* - System prompt: built from role + extra_context
|
||||
* - User prompt (task): the full task instructions from the YAML
|
||||
* - Working directory: the swarm workspace
|
||||
* - Full tool access (bash, python, read, write, edit, grep, find, fetch, web_search, browser)
|
||||
*/
|
||||
export async function executeSwarmAgent(
|
||||
agent: SwarmAgent,
|
||||
index: number,
|
||||
options: SwarmExecutorOptions,
|
||||
): Promise<SingleResult> {
|
||||
const { workspace, swarmName, iteration, modelOverride, signal, onProgress, authStorage, modelRegistry, settings, stateTracker } = options;
|
||||
|
||||
const agentId = `swarm-${swarmName}-${agent.name}-${iteration}`;
|
||||
|
||||
const agentDef: AgentDefinition = {
|
||||
name: agent.name,
|
||||
description: `Swarm agent: ${agent.role}`,
|
||||
systemPrompt: buildSystemPrompt(agent),
|
||||
source: "project" as AgentSource,
|
||||
};
|
||||
|
||||
await stateTracker.updateAgent(agent.name, {
|
||||
status: "running",
|
||||
iteration,
|
||||
startedAt: Date.now(),
|
||||
});
|
||||
await stateTracker.appendLog(agent.name, `Starting iteration ${iteration}`);
|
||||
|
||||
try {
|
||||
const result = await runSubprocess({
|
||||
cwd: workspace,
|
||||
agent: agentDef,
|
||||
task: agent.task,
|
||||
index,
|
||||
id: agentId,
|
||||
modelOverride,
|
||||
signal,
|
||||
onProgress: (progress) => onProgress?.(agent.name, progress),
|
||||
authStorage,
|
||||
modelRegistry,
|
||||
settings,
|
||||
enableLsp: false,
|
||||
artifactsDir: path.join(stateTracker.swarmDir, "context"),
|
||||
});
|
||||
|
||||
const status = result.exitCode === 0 ? "completed" as const : "failed" as const;
|
||||
await stateTracker.updateAgent(agent.name, {
|
||||
status,
|
||||
completedAt: Date.now(),
|
||||
error: result.error,
|
||||
});
|
||||
await stateTracker.appendLog(
|
||||
agent.name,
|
||||
`Iteration ${iteration} ${status}${result.error ? `: ${result.error}` : ""}`,
|
||||
);
|
||||
|
||||
return result;
|
||||
} catch (err) {
|
||||
const error = err instanceof Error ? err.message : String(err);
|
||||
await stateTracker.updateAgent(agent.name, {
|
||||
status: "failed",
|
||||
completedAt: Date.now(),
|
||||
error,
|
||||
});
|
||||
await stateTracker.appendLog(agent.name, `Iteration ${iteration} error: ${error}`);
|
||||
throw err;
|
||||
}
|
||||
}
|
||||
|
||||
function buildSystemPrompt(agent: SwarmAgent): string {
|
||||
const parts = [`You are a ${agent.role}.`];
|
||||
if (agent.extraContext) {
|
||||
parts.push(agent.extraContext);
|
||||
}
|
||||
return parts.join("\n\n");
|
||||
}
|
||||
@@ -0,0 +1,219 @@
|
||||
/**
|
||||
* Pipeline controller for swarm execution.
|
||||
*
|
||||
* Orchestrates execution waves within each iteration:
|
||||
* - Agents in the same wave execute in parallel
|
||||
* - Waves execute sequentially (wave N+1 starts after wave N completes)
|
||||
* - For pipeline mode, iterations repeat the full DAG execution
|
||||
*/
|
||||
import type { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage";
|
||||
import type { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry";
|
||||
import type { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
|
||||
import type { AgentProgress, AgentSource, SingleResult } from "@oh-my-pi/pi-coding-agent/task/types";
|
||||
import type { SwarmDefinition } from "./schema";
|
||||
import type { StateTracker } from "./state";
|
||||
import { executeSwarmAgent } from "./executor";
|
||||
|
||||
// ============================================================================
|
||||
// Types
|
||||
// ============================================================================
|
||||
|
||||
export interface PipelineOptions {
|
||||
workspace: string;
|
||||
signal?: AbortSignal;
|
||||
onProgress?: (state: PipelineProgress) => void;
|
||||
authStorage?: AuthStorage;
|
||||
modelRegistry?: ModelRegistry;
|
||||
settings?: Settings;
|
||||
}
|
||||
|
||||
export interface PipelineProgress {
|
||||
iteration: number;
|
||||
targetCount: number;
|
||||
currentWave: number;
|
||||
totalWaves: number;
|
||||
agents: Record<string, { status: string; iteration: number }>;
|
||||
}
|
||||
|
||||
export interface PipelineResult {
|
||||
status: "completed" | "failed" | "aborted";
|
||||
iterations: number;
|
||||
agentResults: Map<string, SingleResult[]>;
|
||||
errors: string[];
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Controller
|
||||
// ============================================================================
|
||||
|
||||
export class PipelineController {
|
||||
#def: SwarmDefinition;
|
||||
#waves: string[][];
|
||||
#stateTracker: StateTracker;
|
||||
|
||||
constructor(def: SwarmDefinition, waves: string[][], stateTracker: StateTracker) {
|
||||
this.#def = def;
|
||||
this.#waves = waves;
|
||||
this.#stateTracker = stateTracker;
|
||||
}
|
||||
|
||||
async run(options: PipelineOptions): Promise<PipelineResult> {
|
||||
const { workspace, signal, onProgress, authStorage, modelRegistry, settings } = options;
|
||||
const allResults = new Map<string, SingleResult[]>();
|
||||
const errors: string[] = [];
|
||||
|
||||
for (const name of this.#def.agents.keys()) {
|
||||
allResults.set(name, []);
|
||||
}
|
||||
|
||||
const targetCount = this.#def.targetCount;
|
||||
|
||||
await this.#stateTracker.appendOrchestratorLog(
|
||||
`Pipeline '${this.#def.name}' starting: mode=${this.#def.mode} iterations=${targetCount} waves=${this.#waves.length} agents=${this.#def.agents.size}`,
|
||||
);
|
||||
|
||||
try {
|
||||
for (let iteration = 0; iteration < targetCount; iteration++) {
|
||||
if (signal?.aborted) {
|
||||
await this.#stateTracker.updatePipeline({ status: "aborted" });
|
||||
return { status: "aborted", iterations: iteration, agentResults: allResults, errors };
|
||||
}
|
||||
|
||||
await this.#stateTracker.updatePipeline({ iteration });
|
||||
await this.#stateTracker.appendOrchestratorLog(
|
||||
`--- Iteration ${iteration + 1}/${targetCount} ---`,
|
||||
);
|
||||
|
||||
const emitProgress = (currentWave: number) => {
|
||||
onProgress?.({
|
||||
iteration,
|
||||
targetCount,
|
||||
currentWave,
|
||||
totalWaves: this.#waves.length,
|
||||
agents: this.#buildProgressSnapshot(),
|
||||
});
|
||||
};
|
||||
|
||||
const iterationResults = await this.#runIteration(iteration, {
|
||||
workspace,
|
||||
signal,
|
||||
emitProgress,
|
||||
authStorage,
|
||||
modelRegistry,
|
||||
settings,
|
||||
});
|
||||
|
||||
for (const [agentName, result] of iterationResults) {
|
||||
allResults.get(agentName)!.push(result);
|
||||
if (result.exitCode !== 0) {
|
||||
errors.push(`${agentName} (iteration ${iteration + 1}): ${result.error || "exit code " + result.exitCode}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const status = errors.length > 0 ? "failed" as const : "completed" as const;
|
||||
await this.#stateTracker.updatePipeline({ status, completedAt: Date.now() });
|
||||
await this.#stateTracker.appendOrchestratorLog(`Pipeline ${status} (${errors.length} errors)`);
|
||||
return { status, iterations: targetCount, agentResults: allResults, errors };
|
||||
} catch (err) {
|
||||
const error = err instanceof Error ? err.message : String(err);
|
||||
await this.#stateTracker.updatePipeline({ status: "failed", completedAt: Date.now() });
|
||||
await this.#stateTracker.appendOrchestratorLog(`Pipeline fatal error: ${error}`);
|
||||
errors.push(error);
|
||||
return { status: "failed", iterations: 0, agentResults: allResults, errors };
|
||||
}
|
||||
}
|
||||
|
||||
async #runIteration(
|
||||
iteration: number,
|
||||
options: {
|
||||
workspace: string;
|
||||
signal?: AbortSignal;
|
||||
emitProgress: (currentWave: number) => void;
|
||||
authStorage?: AuthStorage;
|
||||
modelRegistry?: ModelRegistry;
|
||||
settings?: Settings;
|
||||
},
|
||||
): Promise<Map<string, SingleResult>> {
|
||||
const results = new Map<string, SingleResult>();
|
||||
let agentIndex = 0;
|
||||
|
||||
for (let waveIdx = 0; waveIdx < this.#waves.length; waveIdx++) {
|
||||
const wave = this.#waves[waveIdx];
|
||||
|
||||
if (options.signal?.aborted) break;
|
||||
|
||||
await this.#stateTracker.appendOrchestratorLog(
|
||||
`Wave ${waveIdx + 1}/${this.#waves.length}: [${wave.join(", ")}]`,
|
||||
);
|
||||
|
||||
// Mark agents in this wave as waiting
|
||||
for (const agentName of wave) {
|
||||
await this.#stateTracker.updateAgent(agentName, {
|
||||
status: "waiting",
|
||||
iteration,
|
||||
wave: waveIdx,
|
||||
});
|
||||
}
|
||||
options.emitProgress(waveIdx);
|
||||
|
||||
// Execute all agents in wave in parallel, catching per-agent errors
|
||||
const waveResults = await Promise.all(
|
||||
wave.map(async (agentName) => {
|
||||
const agent = this.#def.agents.get(agentName)!;
|
||||
const currentIndex = agentIndex++;
|
||||
try {
|
||||
const result = await executeSwarmAgent(agent, currentIndex, {
|
||||
workspace: options.workspace,
|
||||
swarmName: this.#def.name,
|
||||
iteration,
|
||||
modelOverride: this.#def.model,
|
||||
signal: options.signal,
|
||||
onProgress: (_name, _progress) => {
|
||||
options.emitProgress(waveIdx);
|
||||
},
|
||||
authStorage: options.authStorage,
|
||||
modelRegistry: options.modelRegistry,
|
||||
settings: options.settings,
|
||||
stateTracker: this.#stateTracker,
|
||||
});
|
||||
return { agentName, result };
|
||||
} catch (err) {
|
||||
const error = err instanceof Error ? err.message : String(err);
|
||||
const failResult: SingleResult = {
|
||||
index: currentIndex,
|
||||
id: `swarm-${this.#def.name}-${agentName}-${iteration}`,
|
||||
agent: agentName,
|
||||
agentSource: "project" as AgentSource,
|
||||
task: agent.task,
|
||||
exitCode: 1,
|
||||
output: "",
|
||||
stderr: error,
|
||||
truncated: false,
|
||||
durationMs: 0,
|
||||
tokens: 0,
|
||||
error,
|
||||
};
|
||||
return { agentName, result: failResult };
|
||||
}
|
||||
}),
|
||||
);
|
||||
|
||||
for (const { agentName, result } of waveResults) {
|
||||
results.set(agentName, result);
|
||||
}
|
||||
|
||||
options.emitProgress(waveIdx);
|
||||
}
|
||||
|
||||
return results;
|
||||
}
|
||||
|
||||
#buildProgressSnapshot(): Record<string, { status: string; iteration: number }> {
|
||||
const snapshot: Record<string, { status: string; iteration: number }> = {};
|
||||
for (const [name, agent] of Object.entries(this.#stateTracker.state.agents)) {
|
||||
snapshot[name] = { status: agent.status, iteration: agent.iteration };
|
||||
}
|
||||
return snapshot;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,75 @@
|
||||
/**
|
||||
* TUI progress rendering for swarm pipeline status.
|
||||
*/
|
||||
import type { SwarmState } from "./state";
|
||||
|
||||
const STATUS_LABELS: Record<string, string> = {
|
||||
completed: "[done]",
|
||||
running: "[....]",
|
||||
failed: "[FAIL]",
|
||||
pending: "[ ]",
|
||||
waiting: "[wait]",
|
||||
idle: "[idle]",
|
||||
aborted: "[stop]",
|
||||
};
|
||||
|
||||
export function renderSwarmProgress(state: SwarmState): string[] {
|
||||
const lines: string[] = [];
|
||||
|
||||
const statusLabel = state.status.toUpperCase();
|
||||
lines.push(`Swarm: ${state.name} [${statusLabel}]`);
|
||||
lines.push(`Mode: ${state.mode} | Iteration: ${state.iteration + 1}/${state.targetCount}`);
|
||||
lines.push("");
|
||||
|
||||
const agents = Object.values(state.agents);
|
||||
if (agents.length === 0) {
|
||||
lines.push(" (no agents)");
|
||||
return lines;
|
||||
}
|
||||
|
||||
for (const agent of agents) {
|
||||
const icon = STATUS_LABELS[agent.status] ?? "[????]";
|
||||
const duration = formatAgentDuration(agent);
|
||||
const errorSuffix = agent.error ? ` - ${truncate(agent.error, 60)}` : "";
|
||||
lines.push(` ${icon} ${agent.name}: ${agent.status}${duration}${errorSuffix}`);
|
||||
}
|
||||
|
||||
// Summary line
|
||||
const completed = agents.filter((a) => a.status === "completed").length;
|
||||
const failed = agents.filter((a) => a.status === "failed").length;
|
||||
const running = agents.filter((a) => a.status === "running").length;
|
||||
|
||||
lines.push("");
|
||||
const parts = [`${completed}/${agents.length} done`];
|
||||
if (running > 0) parts.push(`${running} running`);
|
||||
if (failed > 0) parts.push(`${failed} failed`);
|
||||
if (state.startedAt) {
|
||||
parts.push(`elapsed: ${formatDuration(Date.now() - state.startedAt)}`);
|
||||
}
|
||||
lines.push(` ${parts.join(" | ")}`);
|
||||
|
||||
return lines;
|
||||
}
|
||||
|
||||
function formatAgentDuration(agent: { startedAt?: number; completedAt?: number; status: string }): string {
|
||||
if (agent.startedAt && agent.completedAt) {
|
||||
return ` (${formatDuration(agent.completedAt - agent.startedAt)})`;
|
||||
}
|
||||
if (agent.startedAt && (agent.status === "running" || agent.status === "waiting")) {
|
||||
return ` (${formatDuration(Date.now() - agent.startedAt)}...)`;
|
||||
}
|
||||
return "";
|
||||
}
|
||||
|
||||
function formatDuration(ms: number): string {
|
||||
if (ms < 1000) return `${ms}ms`;
|
||||
if (ms < 60_000) return `${(ms / 1000).toFixed(1)}s`;
|
||||
const mins = Math.floor(ms / 60_000);
|
||||
const secs = Math.floor((ms % 60_000) / 1000);
|
||||
return `${mins}m${secs}s`;
|
||||
}
|
||||
|
||||
function truncate(str: string, maxLen: number): string {
|
||||
if (str.length <= maxLen) return str;
|
||||
return `${str.slice(0, maxLen - 1)}…`;
|
||||
}
|
||||
@@ -0,0 +1,146 @@
|
||||
/**
|
||||
* YAML schema parsing, validation, and normalized types for swarm definitions.
|
||||
*/
|
||||
import YAML from "yaml";
|
||||
|
||||
// ============================================================================
|
||||
// Raw YAML shape (snake_case, optional fields)
|
||||
// ============================================================================
|
||||
|
||||
interface RawSwarmAgentConfig {
|
||||
role: string;
|
||||
task: string;
|
||||
extra_context?: string;
|
||||
reports_to?: string[];
|
||||
waits_for?: string[];
|
||||
}
|
||||
|
||||
interface RawSwarmConfig {
|
||||
name: string;
|
||||
workspace: string;
|
||||
mode?: string;
|
||||
target_count?: number;
|
||||
model?: string;
|
||||
agents: Record<string, RawSwarmAgentConfig>;
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Normalized types (camelCase, defaults applied)
|
||||
// ============================================================================
|
||||
|
||||
export type SwarmMode = "pipeline" | "parallel" | "sequential";
|
||||
|
||||
export interface SwarmAgent {
|
||||
name: string;
|
||||
role: string;
|
||||
task: string;
|
||||
extraContext?: string;
|
||||
reportsTo: string[];
|
||||
waitsFor: string[];
|
||||
}
|
||||
|
||||
export interface SwarmDefinition {
|
||||
name: string;
|
||||
workspace: string;
|
||||
mode: SwarmMode;
|
||||
targetCount: number;
|
||||
model?: string;
|
||||
agents: Map<string, SwarmAgent>;
|
||||
/** Preserves YAML declaration order for implicit pipeline sequencing. */
|
||||
agentOrder: string[];
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Parsing
|
||||
// ============================================================================
|
||||
|
||||
const VALID_MODES = new Set<string>(["pipeline", "parallel", "sequential"]);
|
||||
|
||||
export function parseSwarmYaml(content: string): SwarmDefinition {
|
||||
const raw = YAML.parse(content) as { swarm?: RawSwarmConfig } | null;
|
||||
if (!raw?.swarm) {
|
||||
throw new Error("YAML must have a top-level 'swarm' key");
|
||||
}
|
||||
const swarm = raw.swarm;
|
||||
|
||||
if (!swarm.name || typeof swarm.name !== "string") {
|
||||
throw new Error("swarm.name is required and must be a string");
|
||||
}
|
||||
if (!swarm.workspace || typeof swarm.workspace !== "string") {
|
||||
throw new Error("swarm.workspace is required and must be a string");
|
||||
}
|
||||
if (!swarm.agents || typeof swarm.agents !== "object" || Object.keys(swarm.agents).length === 0) {
|
||||
throw new Error("swarm.agents must contain at least one agent");
|
||||
}
|
||||
|
||||
const mode = swarm.mode ?? "sequential";
|
||||
if (!VALID_MODES.has(mode)) {
|
||||
throw new Error(`Invalid mode '${mode}'. Must be one of: ${[...VALID_MODES].join(", ")}`);
|
||||
}
|
||||
|
||||
const agentOrder: string[] = [];
|
||||
const agents = new Map<string, SwarmAgent>();
|
||||
|
||||
for (const [name, config] of Object.entries(swarm.agents)) {
|
||||
if (!config.role || typeof config.role !== "string") {
|
||||
throw new Error(`Agent '${name}': 'role' is required`);
|
||||
}
|
||||
if (!config.task || typeof config.task !== "string") {
|
||||
throw new Error(`Agent '${name}': 'task' is required`);
|
||||
}
|
||||
|
||||
agentOrder.push(name);
|
||||
agents.set(name, {
|
||||
name,
|
||||
role: config.role,
|
||||
task: config.task.trim(),
|
||||
extraContext: config.extra_context?.trim(),
|
||||
reportsTo: Array.isArray(config.reports_to) ? config.reports_to : [],
|
||||
waitsFor: Array.isArray(config.waits_for) ? config.waits_for : [],
|
||||
});
|
||||
}
|
||||
|
||||
return {
|
||||
name: swarm.name,
|
||||
workspace: swarm.workspace,
|
||||
mode: mode as SwarmMode,
|
||||
targetCount: swarm.target_count ?? 1,
|
||||
model: swarm.model,
|
||||
agents,
|
||||
agentOrder,
|
||||
};
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Validation (semantic — references, constraints)
|
||||
// ============================================================================
|
||||
|
||||
export function validateSwarmDefinition(def: SwarmDefinition): string[] {
|
||||
const errors: string[] = [];
|
||||
const agentNames = new Set(def.agents.keys());
|
||||
|
||||
for (const [name, agent] of def.agents) {
|
||||
for (const dep of agent.waitsFor) {
|
||||
if (!agentNames.has(dep)) {
|
||||
errors.push(`Agent '${name}' waits_for unknown agent '${dep}'`);
|
||||
}
|
||||
if (dep === name) {
|
||||
errors.push(`Agent '${name}' cannot wait for itself`);
|
||||
}
|
||||
}
|
||||
for (const target of agent.reportsTo) {
|
||||
if (!agentNames.has(target)) {
|
||||
errors.push(`Agent '${name}' reports_to unknown agent '${target}'`);
|
||||
}
|
||||
if (target === name) {
|
||||
errors.push(`Agent '${name}' cannot report to itself`);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (def.targetCount < 1) {
|
||||
errors.push("target_count must be at least 1");
|
||||
}
|
||||
|
||||
return errors;
|
||||
}
|
||||
@@ -0,0 +1,127 @@
|
||||
/**
|
||||
* Filesystem state tracker for swarm pipeline execution.
|
||||
*
|
||||
* Persists pipeline and per-agent state to `.swarm_<name>/` in the workspace.
|
||||
* Supports resumability by loading state from disk.
|
||||
*/
|
||||
import * as fs from "node:fs/promises";
|
||||
import * as path from "node:path";
|
||||
|
||||
// ============================================================================
|
||||
// State types
|
||||
// ============================================================================
|
||||
|
||||
export type PipelineStatus = "idle" | "running" | "completed" | "failed" | "aborted";
|
||||
export type AgentStatus = "pending" | "waiting" | "running" | "completed" | "failed";
|
||||
|
||||
export interface AgentState {
|
||||
name: string;
|
||||
status: AgentStatus;
|
||||
iteration: number;
|
||||
wave: number;
|
||||
startedAt?: number;
|
||||
completedAt?: number;
|
||||
error?: string;
|
||||
}
|
||||
|
||||
export interface SwarmState {
|
||||
name: string;
|
||||
status: PipelineStatus;
|
||||
mode: string;
|
||||
iteration: number;
|
||||
targetCount: number;
|
||||
agents: Record<string, AgentState>;
|
||||
startedAt: number;
|
||||
completedAt?: number;
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// State tracker
|
||||
// ============================================================================
|
||||
|
||||
export class StateTracker {
|
||||
#swarmDir: string;
|
||||
#state: SwarmState;
|
||||
|
||||
constructor(workspaceDir: string, name: string) {
|
||||
this.#swarmDir = path.join(workspaceDir, `.swarm_${name}`);
|
||||
this.#state = {
|
||||
name,
|
||||
status: "idle",
|
||||
mode: "sequential",
|
||||
iteration: 0,
|
||||
targetCount: 1,
|
||||
agents: {},
|
||||
startedAt: Date.now(),
|
||||
};
|
||||
}
|
||||
|
||||
get swarmDir(): string {
|
||||
return this.#swarmDir;
|
||||
}
|
||||
|
||||
get state(): Readonly<SwarmState> {
|
||||
return this.#state;
|
||||
}
|
||||
|
||||
async init(agentNames: string[], targetCount: number, mode: string): Promise<void> {
|
||||
await fs.mkdir(path.join(this.#swarmDir, "state"), { recursive: true });
|
||||
await fs.mkdir(path.join(this.#swarmDir, "logs"), { recursive: true });
|
||||
await fs.mkdir(path.join(this.#swarmDir, "context"), { recursive: true });
|
||||
|
||||
this.#state.targetCount = targetCount;
|
||||
this.#state.mode = mode;
|
||||
this.#state.status = "running";
|
||||
this.#state.startedAt = Date.now();
|
||||
|
||||
for (const name of agentNames) {
|
||||
this.#state.agents[name] = {
|
||||
name,
|
||||
status: "pending",
|
||||
iteration: 0,
|
||||
wave: 0,
|
||||
};
|
||||
}
|
||||
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
async updateAgent(name: string, update: Partial<AgentState>): Promise<void> {
|
||||
const agent = this.#state.agents[name];
|
||||
if (!agent) return;
|
||||
Object.assign(agent, update);
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
async updatePipeline(update: Partial<SwarmState>): Promise<void> {
|
||||
Object.assign(this.#state, update);
|
||||
await this.#persist();
|
||||
}
|
||||
|
||||
async appendLog(agentName: string, message: string): Promise<void> {
|
||||
const logPath = path.join(this.#swarmDir, "logs", `${agentName}.log`);
|
||||
const timestamp = new Date().toISOString();
|
||||
await fs.appendFile(logPath, `[${timestamp}] ${message}\n`);
|
||||
}
|
||||
|
||||
async appendOrchestratorLog(message: string): Promise<void> {
|
||||
const logPath = path.join(this.#swarmDir, "logs", "orchestrator.log");
|
||||
const timestamp = new Date().toISOString();
|
||||
await fs.appendFile(logPath, `[${timestamp}] ${message}\n`);
|
||||
}
|
||||
|
||||
async load(): Promise<SwarmState | null> {
|
||||
const statePath = path.join(this.#swarmDir, "state", "pipeline.json");
|
||||
try {
|
||||
const content = await Bun.file(statePath).text();
|
||||
this.#state = JSON.parse(content) as SwarmState;
|
||||
return this.#state;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
async #persist(): Promise<void> {
|
||||
await Bun.write(path.join(this.#swarmDir, "state", "pipeline.json"), JSON.stringify(this.#state, null, 2));
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
{
|
||||
"compilerOptions": {
|
||||
"target": "ESNext",
|
||||
"module": "ESNext",
|
||||
"moduleResolution": "bundler",
|
||||
"strict": true,
|
||||
"noEmit": true,
|
||||
"skipLibCheck": true,
|
||||
"types": ["bun-types"],
|
||||
"paths": {
|
||||
"@oh-my-pi/pi-coding-agent": ["../../src/index.ts"],
|
||||
"@oh-my-pi/pi-coding-agent/*": ["../../src/*.ts"]
|
||||
}
|
||||
},
|
||||
"include": ["extension.ts", "swarm/**/*.ts"]
|
||||
}
|
||||
Reference in New Issue
Block a user