Fluent builder API → DAG engine → auto-serializable config — swarm/builder/
class Swarm:
def __init__(self) -> None:
self._dag = DAG()
self._last_node: str | None = None
# ── pattern methods (each adds a DAG node) ──────────────────────────
def sequential(self, name: str, agents: list[Agent], *, after: str | list[str] | None = None) -> "Swarm":
return self._add_node(name, SequentialPattern(), agents, after)
def parallel(self, name: str, agents: list[Agent], *, after: str | list[str] | None = None, merge: str = "concatenate") -> "Swarm":
return self._add_node(name, ParallelPattern(merge_strategy=merge), agents, after)
def hierarchical(self, name: str, agents: list[Agent], *, after: str | list[str] | None = None) -> "Swarm":
return self._add_node(name, HierarchicalPattern(), agents, after)
def decentralized(self, name: str, agents: list[Agent], *, after: str | list[str] | None = None) -> "Swarm":
return self._add_node(name, DecentralizedPattern(), agents, after)
def adaptive(self, name: str, agents: list[Agent], *, threshold: float = 0.8, after: str | list[str] | None = None) -> "Swarm":
return self._add_node(name, AdaptivePattern(threshold), agents, after)
def mesh(self, name: str, agents: list[Agent], *, after: str | list[str] | None = None) -> "Swarm":
return self._add_node(name, MeshPattern(), agents, after)
# ── execution ────────────────────────────────────────────────────────
async def run(self, task: str, ctx: SwarmContext | None = None) -> SwarmResult:
ctx = ctx or SwarmContext()
return await self._dag.execute(task, ctx)
def run_sync(self, task: str) -> SwarmResult:
return asyncio.run(self.run(task))
# ── config serialization ─────────────────────────────────────────────
def to_config(self) -> dict:
return self._dag.to_config()
@classmethod
def from_config(cls, config: dict, agents: dict[str, Agent]) -> "Swarm":
swarm = cls()
swarm._dag = DAG.from_config(config, agents)
return swarm
def _add_node(self, name, pattern, agents, after) -> "Swarm":
deps = [after] if isinstance(after, str) else (after or ([self._last_node] if self._last_node else []))
self._dag.add_node(name, pattern, agents, deps)
self._last_node = name
return self
class DAG:
"""Directed acyclic graph of pattern nodes."""
def __init__(self) -> None:
self._nodes: dict[str, DAGNode] = {}
async def execute(self, task: str, ctx: SwarmContext) -> SwarmResult:
order = self._topological_sort()
last_result = None
for node_name in order:
node = self._nodes[node_name]
# pass prior output as task if deps exist
node_task = ctx.state.get(f"{node_name}.input", task)
result = await node.pattern.execute(node.agents, node_task, ctx)
ctx.state[f"{node_name}.output"] = result.final_output
last_result = result
return last_result
# Linear (implicit chaining — no after= needed)
result = await (Swarm()
.hierarchical("plan", [coordinator, *workers])
.parallel("research", [r1, r2, r3])
.sequential("write", [writer, editor])
.adaptive("review", [reviewer, specialist])
.run("Write a comprehensive market analysis"))
# Branching (explicit after=)
s = Swarm()
s.hierarchical("plan", [coordinator, *workers])
s.parallel("research", [r1, r2, r3], after="plan")
s.parallel("data", [d1, d2], after="plan")
s.sequential("synthesize", [synth], after=["research", "data"])
s.adaptive("review", [reviewer], after="synthesize")
result = await s.run("Build a startup business plan")
# From config (CLI usage)
swarm = Swarm.from_config(yaml.safe_load(config_file), agents=agent_registry)
result = await swarm.run(task)