gmf_forge_ai_orchestration.multi_agent
Multi-agent orchestration — Supervisor, Pipeline, Debate, and Swarm.
1"""Multi-agent orchestration — Supervisor, Pipeline, Debate, and Swarm.""" 2 3from gmf_forge_ai_orchestration.multi_agent.base import BaseOrchestrator, OrchestratorResult 4from gmf_forge_ai_orchestration.multi_agent.supervisor import SupervisorOrchestrator 5from gmf_forge_ai_orchestration.multi_agent.pipeline import PipelineOrchestrator 6from gmf_forge_ai_orchestration.multi_agent.debate import DebateOrchestrator 7from gmf_forge_ai_orchestration.multi_agent.swarm import SwarmOrchestrator 8 9__all__ = [ 10 "BaseOrchestrator", 11 "OrchestratorResult", 12 "SupervisorOrchestrator", 13 "PipelineOrchestrator", 14 "DebateOrchestrator", 15 "SwarmOrchestrator", 16]
60class BaseOrchestrator(ABC): 61 """ 62 Abstract base class for multi-agent orchestrators. 63 64 Args: 65 logger: Optional :class:`BasicLogger`. 66 metrics: Optional :class:`BasicMetricsCollector`. 67 tracer: Optional :class:`TracingProvider`. Falls back to ``get_tracer()``. 68 """ 69 70 def __init__( 71 self, 72 router: Optional["BaseRouter"] = None, 73 logger: Optional[BasicLogger] = None, 74 metrics: Optional[BasicMetricsCollector] = None, 75 tracer: Optional[TracingProvider] = None, 76 registry: Optional["OrchestratorRegistry"] = None, 77 ) -> None: 78 self._logger = logger or BasicLogger( 79 f"gmf_forge_ai.orchestrator.{self.__class__.__name__}" 80 ) 81 self._metrics = metrics 82 self._tracer = tracer or get_tracer() 83 self.router: Optional["BaseRouter"] = router 84 self._orchestrator_id = str(uuid.uuid4()) 85 self._registry = registry 86 self._run_history: Dict[str, RunRecord] = {} 87 self._current_run_id: Optional[str] = None 88 self._status: Literal["idle", "running", "completed", "error"] = "idle" 89 self._pending_agent_operations: List[PendingAgentOperation] = [] 90 91 async def register(self) -> None: 92 """Register with the configured registry, if any.""" 93 if self._registry: 94 await self._registry.register(self) 95 96 async def close(self) -> None: 97 """Unregister from the registry during application shutdown.""" 98 if self._registry: 99 await self._registry.unregister(self._orchestrator_id) 100 101 @abstractmethod 102 async def run(self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult: 103 """Execute the multi-agent orchestration and return the aggregated result.""" 104 105 @abstractmethod 106 def get_orchestrator_type(self) -> str: 107 """Return a stable orchestrator type literal.""" 108 109 @abstractmethod 110 def list_configured_agents(self) -> List[Dict[str, Any]]: 111 """Return configured orchestrator agents for management and discovery.""" 112 113 @abstractmethod 114 def _add_agent_now(self, payload: Dict[str, Any]) -> None: 115 """Apply an add-agent operation immediately.""" 116 117 @abstractmethod 118 def _remove_agent_now(self, agent_id: str) -> None: 119 """Apply a remove-agent operation immediately.""" 120 121 def get_orchestrator_id(self) -> str: 122 return self._orchestrator_id 123 124 def get_current_status(self) -> str: 125 return self._status 126 127 def list_runs(self, limit: int = 20, offset: int = 0) -> List[RunRecord]: 128 runs = sorted( 129 self._run_history.values(), 130 key=lambda run: run.start_time, 131 reverse=True, 132 ) 133 return runs[offset : offset + limit] 134 135 def get_run(self, run_id: str) -> Optional[RunRecord]: 136 return self._run_history.get(run_id) 137 138 def list_run_agents(self, run_id: str) -> List[str]: 139 run = self.get_run(run_id) 140 if not run: 141 return [] 142 return list(run.agent_outputs.keys()) 143 144 def get_run_agent(self, run_id: str, agent_id: str) -> Optional[AgentResult]: 145 run = self.get_run(run_id) 146 if not run: 147 return None 148 return run.agent_outputs.get(agent_id) 149 150 def queue_add_agent(self, payload: Dict[str, Any]) -> None: 151 """Queue add-agent mutation to be applied before the next run.""" 152 self._pending_agent_operations.append(PendingAgentOperation(op="add", payload=payload)) 153 154 def queue_remove_agent(self, agent_id: str) -> None: 155 """Queue remove-agent mutation to be applied before the next run.""" 156 self._pending_agent_operations.append( 157 PendingAgentOperation(op="remove", payload={"agent_id": agent_id}) 158 ) 159 160 def apply_pending_agent_operations(self) -> None: 161 """Apply pending membership changes before starting a new run.""" 162 if not self._pending_agent_operations: 163 return 164 165 for operation in self._pending_agent_operations: 166 if operation.op == "add": 167 self._add_agent_now(operation.payload) 168 else: 169 agent_id = str(operation.payload.get("agent_id", "")) 170 if agent_id: 171 self._remove_agent_now(agent_id) 172 self._pending_agent_operations.clear() 173 174 def begin_run(self) -> str: 175 """Create a run record and mark orchestrator as running.""" 176 self.apply_pending_agent_operations() 177 run_id = str(uuid.uuid4()) 178 self._current_run_id = run_id 179 self._status = "running" 180 self._run_history[run_id] = RunRecord( 181 run_id=run_id, 182 status="running", 183 start_time=datetime.now(timezone.utc), 184 ) 185 return run_id 186 187 def update_run_type_specific(self, run_id: str, data: Dict[str, Any]) -> None: 188 run = self._run_history.get(run_id) 189 if not run: 190 return 191 run.type_specific = {**run.type_specific, **data} 192 193 def complete_run( 194 self, 195 run_id: str, 196 final_output: str, 197 agent_outputs: Dict[str, AgentResult], 198 rounds: int, 199 ) -> None: 200 run = self._run_history.get(run_id) 201 if not run: 202 return 203 run.status = "completed" 204 run.end_time = datetime.now(timezone.utc) 205 run.final_output = final_output 206 run.agent_outputs = agent_outputs 207 run.rounds = rounds 208 run.success = True 209 run.error = None 210 self._status = "idle" 211 self._current_run_id = None 212 self._persist_run(run) 213 214 def fail_run( 215 self, 216 run_id: str, 217 error: str, 218 agent_outputs: Optional[Dict[str, AgentResult]] = None, 219 rounds: int = 0, 220 final_output: str = "", 221 ) -> None: 222 run = self._run_history.get(run_id) 223 if not run: 224 return 225 run.status = "error" 226 run.end_time = datetime.now(timezone.utc) 227 run.final_output = final_output 228 run.agent_outputs = agent_outputs or {} 229 run.rounds = rounds 230 run.success = False 231 run.error = error 232 self._status = "idle" 233 self._current_run_id = None 234 self._persist_run(run) 235 236 def _persist_run(self, run: "RunRecord") -> None: 237 """Fire-and-forget run persistence via the registry (backend-agnostic).""" 238 if not self._registry or not hasattr(self._registry, "save_run"): 239 return 240 try: 241 loop = asyncio.get_event_loop() 242 if loop.is_running(): 243 loop.create_task(self._async_persist_run(self._orchestrator_id, run)) 244 except RuntimeError: 245 pass 246 247 async def _async_persist_run(self, orchestrator_id: str, run: "RunRecord") -> None: 248 """Serialize and persist a RunRecord through the registry's state store.""" 249 payload = { 250 "run_id": run.run_id, 251 "status": run.status, 252 "start_time": run.start_time.isoformat() if run.start_time else None, 253 "end_time": run.end_time.isoformat() if run.end_time else None, 254 "final_output": run.final_output, 255 "rounds": run.rounds, 256 "success": run.success, 257 "error": run.error, 258 "type_specific": run.type_specific, 259 "agent_outputs": { 260 agent_id: { 261 "output": r.output, 262 "success": r.success, 263 "error": r.error, 264 "steps": [ 265 { 266 "step_index": i, 267 "thought": s.thought, 268 "action": s.action, 269 "observation": s.observation, 270 "action_input": s.action_input if hasattr(s, "action_input") else {}, 271 } 272 for i, s in enumerate(r.steps) 273 ], 274 } 275 for agent_id, r in run.agent_outputs.items() 276 }, 277 } 278 await self._registry.save_run(orchestrator_id, payload) 279 280 def _log_start(self, task: str) -> None: 281 self._logger.info("Orchestrator started", orchestrator=self.__class__.__name__, task=task) 282 if self._metrics: 283 self._metrics.increment("orchestrator.runs", orchestrator=self.__class__.__name__) 284 285 def _log_agent_dispatch(self, agent_id: str, task: str) -> None: 286 self._logger.info("Dispatching to agent", agent_id=agent_id, task=task[:100]) 287 288 def _log_finished(self, rounds: int, success: bool) -> None: 289 self._logger.info( 290 "Orchestrator finished", 291 orchestrator=self.__class__.__name__, 292 rounds=rounds, 293 success=success, 294 ) 295 if self._metrics: 296 self._metrics.increment("orchestrator.rounds", count=rounds)
Abstract base class for multi-agent orchestrators.
Args:
logger: Optional BasicLogger.
metrics: Optional BasicMetricsCollector.
tracer: Optional TracingProvider. Falls back to get_tracer().
91 async def register(self) -> None: 92 """Register with the configured registry, if any.""" 93 if self._registry: 94 await self._registry.register(self)
Register with the configured registry, if any.
96 async def close(self) -> None: 97 """Unregister from the registry during application shutdown.""" 98 if self._registry: 99 await self._registry.unregister(self._orchestrator_id)
Unregister from the registry during application shutdown.
101 @abstractmethod 102 async def run(self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult: 103 """Execute the multi-agent orchestration and return the aggregated result."""
Execute the multi-agent orchestration and return the aggregated result.
105 @abstractmethod 106 def get_orchestrator_type(self) -> str: 107 """Return a stable orchestrator type literal."""
Return a stable orchestrator type literal.
109 @abstractmethod 110 def list_configured_agents(self) -> List[Dict[str, Any]]: 111 """Return configured orchestrator agents for management and discovery."""
Return configured orchestrator agents for management and discovery.
150 def queue_add_agent(self, payload: Dict[str, Any]) -> None: 151 """Queue add-agent mutation to be applied before the next run.""" 152 self._pending_agent_operations.append(PendingAgentOperation(op="add", payload=payload))
Queue add-agent mutation to be applied before the next run.
154 def queue_remove_agent(self, agent_id: str) -> None: 155 """Queue remove-agent mutation to be applied before the next run.""" 156 self._pending_agent_operations.append( 157 PendingAgentOperation(op="remove", payload={"agent_id": agent_id}) 158 )
Queue remove-agent mutation to be applied before the next run.
160 def apply_pending_agent_operations(self) -> None: 161 """Apply pending membership changes before starting a new run.""" 162 if not self._pending_agent_operations: 163 return 164 165 for operation in self._pending_agent_operations: 166 if operation.op == "add": 167 self._add_agent_now(operation.payload) 168 else: 169 agent_id = str(operation.payload.get("agent_id", "")) 170 if agent_id: 171 self._remove_agent_now(agent_id) 172 self._pending_agent_operations.clear()
Apply pending membership changes before starting a new run.
174 def begin_run(self) -> str: 175 """Create a run record and mark orchestrator as running.""" 176 self.apply_pending_agent_operations() 177 run_id = str(uuid.uuid4()) 178 self._current_run_id = run_id 179 self._status = "running" 180 self._run_history[run_id] = RunRecord( 181 run_id=run_id, 182 status="running", 183 start_time=datetime.now(timezone.utc), 184 ) 185 return run_id
Create a run record and mark orchestrator as running.
193 def complete_run( 194 self, 195 run_id: str, 196 final_output: str, 197 agent_outputs: Dict[str, AgentResult], 198 rounds: int, 199 ) -> None: 200 run = self._run_history.get(run_id) 201 if not run: 202 return 203 run.status = "completed" 204 run.end_time = datetime.now(timezone.utc) 205 run.final_output = final_output 206 run.agent_outputs = agent_outputs 207 run.rounds = rounds 208 run.success = True 209 run.error = None 210 self._status = "idle" 211 self._current_run_id = None 212 self._persist_run(run)
214 def fail_run( 215 self, 216 run_id: str, 217 error: str, 218 agent_outputs: Optional[Dict[str, AgentResult]] = None, 219 rounds: int = 0, 220 final_output: str = "", 221 ) -> None: 222 run = self._run_history.get(run_id) 223 if not run: 224 return 225 run.status = "error" 226 run.end_time = datetime.now(timezone.utc) 227 run.final_output = final_output 228 run.agent_outputs = agent_outputs or {} 229 run.rounds = rounds 230 run.success = False 231 run.error = error 232 self._status = "idle" 233 self._current_run_id = None 234 self._persist_run(run)
22@dataclass 23class OrchestratorResult: 24 """The aggregated result of a multi-agent orchestration run.""" 25 26 final_output: str 27 agent_outputs: Dict[str, AgentResult] = field(default_factory=dict) 28 subtask_outputs: List[Tuple[str, str, AgentResult]] = field(default_factory=list) 29 rounds: int = 0 30 success: bool = True 31 error: Optional[str] = None 32 metadata: Dict[str, Any] = field(default_factory=dict) 33 run_id: Optional[str] = None
The aggregated result of a multi-agent orchestration run.
55class SupervisorOrchestrator(BaseOrchestrator): 56 """ 57 LLM-based supervisor that plans subtask assignment and synthesizes results. 58 59 Phase 1 — Plan: The LLM supervisor decomposes the task and assigns each 60 subtask to the most appropriate worker agent. 61 Phase 2 — Execute: All worker agents run (concurrently if not interdependent). 62 Phase 3 — Synthesize: The supervisor LLM merges all results into a final answer. 63 64 Args: 65 supervisor_gateway: The LLM gateway used by the supervisor. 66 agents: Mapping of agent name → :class:`BaseAgent`. 67 agent_descriptions: Optional human-readable descriptions for the supervisor prompt. 68 supervisor_model: LLM model for supervisor calls (optional). 69 logger, metrics, tracer: Observability (optional). 70 """ 71 72 def __init__( 73 self, 74 supervisor_gateway: Any, 75 agents: Optional[Dict[str, BaseAgent]] = None, 76 agent_descriptions: Optional[Dict[str, str]] = None, 77 supervisor_model: Optional[str] = None, 78 router: Optional[BaseRouter] = None, 79 logger: Optional[BasicLogger] = None, 80 metrics: Optional[BasicMetricsCollector] = None, 81 tracer: Optional[TracingProvider] = None, 82 registry: Optional[Any] = None, 83 ) -> None: 84 super().__init__( 85 router=router, 86 logger=logger, 87 metrics=metrics, 88 tracer=tracer, 89 registry=registry, 90 ) 91 self.supervisor_gateway = supervisor_gateway 92 self.agents: Dict[str, BaseAgent] = agents or {} 93 self.agent_descriptions: Dict[str, str] = agent_descriptions or {} 94 self.supervisor_model = supervisor_model 95 96 def get_orchestrator_type(self) -> str: 97 return "supervisor" 98 99 def list_configured_agents(self) -> List[Dict[str, Any]]: 100 return [ 101 { 102 "agent_id": agent_id, 103 "name": agent_id, 104 "role": self.agent_descriptions.get(agent_id), 105 "model": None, 106 "metadata": {}, 107 } 108 for agent_id in self.agents.keys() 109 ] 110 111 def _add_agent_now(self, payload: Dict[str, Any]) -> None: 112 agent_id = str(payload["agent_id"]) 113 agent_obj = payload.get("agent") 114 if agent_obj is None: 115 raise ValueError("Supervisor add agent requires payload['agent']") 116 self.agents[agent_id] = agent_obj 117 role = payload.get("role") 118 if role: 119 self.agent_descriptions[agent_id] = str(role) 120 121 def _remove_agent_now(self, agent_id: str) -> None: 122 if agent_id not in self.agents: 123 raise ValueError(f"Agent not found: {agent_id}") 124 if len(self.agents) <= 1: 125 raise ValueError("Cannot remove the last agent from supervisor orchestrator") 126 self.agents.pop(agent_id, None) 127 self.agent_descriptions.pop(agent_id, None) 128 129 async def _decompose(self, task: str) -> List[str]: 130 """Use LLM to decompose a task into subtask strings (no agent assignment).""" 131 prompt = _SUPERVISOR_DECOMPOSE_PROMPT.format(task=task) 132 response = await self.supervisor_gateway.complete( 133 prompt, model=self.supervisor_model, temperature=0.0 134 ) 135 raw = response.content.strip() 136 match = re.search(r"\[.*\]", raw, re.DOTALL) 137 if match: 138 try: 139 subtasks = json.loads(match.group()) 140 if isinstance(subtasks, list) and all(isinstance(s, str) for s in subtasks): 141 return subtasks 142 except (json.JSONDecodeError, ValueError): 143 pass 144 return [task] # fallback: treat whole task as one subtask 145 146 async def _plan(self, task: str) -> List[Dict[str, str]]: 147 """Decompose task and assign agents. 148 149 When a ``router`` is injected: the LLM only decomposes the task into 150 subtasks; the router selects the best agent for each subtask. 151 152 When no router is provided: the LLM both decomposes and assigns agents 153 (original behaviour). 154 """ 155 if self.router: 156 subtasks = await self._decompose(task) 157 available = list(self.agents.keys()) 158 plan: List[Dict[str, str]] = [] 159 for subtask in subtasks: 160 decision = await self.router.route( 161 RoutingRequest(input=subtask, available_agents=available) 162 ) 163 agent_name = decision.target if decision.target in self.agents else available[0] 164 plan.append({"agent": agent_name, "subtask": subtask}) 165 self._logger.info( 166 "Router assigned subtask", 167 agent=agent_name, 168 confidence=decision.confidence, 169 subtask=subtask[:80], 170 ) 171 return plan 172 173 # --- Original LLM-based plan (no router) --- 174 desc_block = "\n".join( 175 f"- {name}: {self.agent_descriptions.get(name, 'General agent.')}" 176 for name in self.agents 177 ) 178 prompt = _SUPERVISOR_PLAN_PROMPT.format( 179 agent_descriptions=desc_block, task=task 180 ) 181 response = await self.supervisor_gateway.complete( 182 prompt, model=self.supervisor_model, temperature=0.0 183 ) 184 raw = response.content.strip() 185 match = re.search(r"\[.*\]", raw, re.DOTALL) 186 if match: 187 try: 188 llm_plan = json.loads(match.group()) 189 if isinstance(llm_plan, list): 190 return [p for p in llm_plan if "agent" in p and "subtask" in p] 191 except (json.JSONDecodeError, ValueError): 192 pass 193 # Fallback: assign whole task to first agent 194 first_agent = next(iter(self.agents), "") 195 return [{"agent": first_agent, "subtask": task}] 196 197 async def _synthesize(self, task: str, subtask_results: List[Tuple[str, str, AgentResult]]) -> str: 198 results_block = "\n\n".join( 199 f"[{agent}]: {result.output}" for agent, _subtask, result in subtask_results 200 ) 201 prompt = _SUPERVISOR_SYNTHESIS_PROMPT.format( 202 task=task, agent_results=results_block 203 ) 204 response = await self.supervisor_gateway.complete( 205 prompt, model=self.supervisor_model, temperature=0.1 206 ) 207 return response.content.strip() 208 209 async def run(self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult: 210 self._log_start(task) 211 try: 212 run_id = self.begin_run() 213 except Exception as exc: 214 self._logger.error("SupervisorOrchestrator failed before run start", error=str(exc)) 215 return OrchestratorResult( 216 final_output="", 217 rounds=0, 218 success=False, 219 error=str(exc), 220 ) 221 222 with self._tracer.trace( 223 "supervisor.run", input=task, metadata={"agent_count": len(self.agents)} 224 ) as trace: 225 try: 226 # Phase 1: Plan 227 with trace.span("supervisor.plan", input=task) as plan_span: 228 plan = await self._plan(task) 229 plan_span.set_output({"plan": plan}) 230 231 # Phase 2: Execute subtasks 232 agent_outputs: Dict[str, AgentResult] = {} 233 subtask_list: List[Tuple[str, str, AgentResult]] = [] 234 for assignment in plan: 235 agent_name = assignment["agent"] 236 subtask = assignment["subtask"] 237 agent = self.agents.get(agent_name) 238 if agent is None: 239 self._logger.warning("Unknown agent in plan", agent=agent_name) 240 continue 241 242 self._log_agent_dispatch(agent_name, subtask) 243 with trace.span(f"agent.{agent_name}", input=subtask) as aspan: 244 result = await agent.execute(subtask, context=context) 245 agent_outputs[agent_name] = result 246 subtask_list.append((agent_name, subtask, result)) 247 aspan.set_output(result.output) 248 249 self.update_run_type_specific( 250 run_id, 251 { 252 "subtask_outputs": [ 253 { 254 "agent": agent_name, 255 "subtask": subtask_text, 256 "output": result.output, 257 "success": result.success, 258 "steps": [s.__dict__ for s in result.steps], 259 } 260 for agent_name, subtask_text, result in subtask_list 261 ] 262 }, 263 ) 264 265 # Phase 3: Synthesize 266 with trace.span("supervisor.synthesize") as syn_span: 267 final_output = await self._synthesize(task, subtask_list) 268 syn_span.set_output(final_output) 269 270 self._log_finished(rounds=1, success=True) 271 trace.set_output(final_output) 272 self.complete_run( 273 run_id=run_id, 274 final_output=final_output, 275 agent_outputs=agent_outputs, 276 rounds=1, 277 ) 278 return OrchestratorResult( 279 final_output=final_output, 280 agent_outputs=agent_outputs, 281 subtask_outputs=subtask_list, 282 rounds=1, 283 success=True, 284 ) 285 286 except Exception as exc: 287 self._logger.error("SupervisorOrchestrator failed", error=str(exc)) 288 trace.set_error(exc) 289 self.fail_run( 290 run_id=run_id, 291 error=str(exc), 292 rounds=1, 293 ) 294 return OrchestratorResult( 295 final_output="", 296 rounds=1, 297 success=False, 298 error=str(exc), 299 )
LLM-based supervisor that plans subtask assignment and synthesizes results.
Phase 1 — Plan: The LLM supervisor decomposes the task and assigns each subtask to the most appropriate worker agent. Phase 2 — Execute: All worker agents run (concurrently if not interdependent). Phase 3 — Synthesize: The supervisor LLM merges all results into a final answer.
Args:
supervisor_gateway: The LLM gateway used by the supervisor.
agents: Mapping of agent name → BaseAgent.
agent_descriptions: Optional human-readable descriptions for the supervisor prompt.
supervisor_model: LLM model for supervisor calls (optional).
logger, metrics, tracer: Observability (optional).
72 def __init__( 73 self, 74 supervisor_gateway: Any, 75 agents: Optional[Dict[str, BaseAgent]] = None, 76 agent_descriptions: Optional[Dict[str, str]] = None, 77 supervisor_model: Optional[str] = None, 78 router: Optional[BaseRouter] = None, 79 logger: Optional[BasicLogger] = None, 80 metrics: Optional[BasicMetricsCollector] = None, 81 tracer: Optional[TracingProvider] = None, 82 registry: Optional[Any] = None, 83 ) -> None: 84 super().__init__( 85 router=router, 86 logger=logger, 87 metrics=metrics, 88 tracer=tracer, 89 registry=registry, 90 ) 91 self.supervisor_gateway = supervisor_gateway 92 self.agents: Dict[str, BaseAgent] = agents or {} 93 self.agent_descriptions: Dict[str, str] = agent_descriptions or {} 94 self.supervisor_model = supervisor_model
99 def list_configured_agents(self) -> List[Dict[str, Any]]: 100 return [ 101 { 102 "agent_id": agent_id, 103 "name": agent_id, 104 "role": self.agent_descriptions.get(agent_id), 105 "model": None, 106 "metadata": {}, 107 } 108 for agent_id in self.agents.keys() 109 ]
Return configured orchestrator agents for management and discovery.
209 async def run(self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult: 210 self._log_start(task) 211 try: 212 run_id = self.begin_run() 213 except Exception as exc: 214 self._logger.error("SupervisorOrchestrator failed before run start", error=str(exc)) 215 return OrchestratorResult( 216 final_output="", 217 rounds=0, 218 success=False, 219 error=str(exc), 220 ) 221 222 with self._tracer.trace( 223 "supervisor.run", input=task, metadata={"agent_count": len(self.agents)} 224 ) as trace: 225 try: 226 # Phase 1: Plan 227 with trace.span("supervisor.plan", input=task) as plan_span: 228 plan = await self._plan(task) 229 plan_span.set_output({"plan": plan}) 230 231 # Phase 2: Execute subtasks 232 agent_outputs: Dict[str, AgentResult] = {} 233 subtask_list: List[Tuple[str, str, AgentResult]] = [] 234 for assignment in plan: 235 agent_name = assignment["agent"] 236 subtask = assignment["subtask"] 237 agent = self.agents.get(agent_name) 238 if agent is None: 239 self._logger.warning("Unknown agent in plan", agent=agent_name) 240 continue 241 242 self._log_agent_dispatch(agent_name, subtask) 243 with trace.span(f"agent.{agent_name}", input=subtask) as aspan: 244 result = await agent.execute(subtask, context=context) 245 agent_outputs[agent_name] = result 246 subtask_list.append((agent_name, subtask, result)) 247 aspan.set_output(result.output) 248 249 self.update_run_type_specific( 250 run_id, 251 { 252 "subtask_outputs": [ 253 { 254 "agent": agent_name, 255 "subtask": subtask_text, 256 "output": result.output, 257 "success": result.success, 258 "steps": [s.__dict__ for s in result.steps], 259 } 260 for agent_name, subtask_text, result in subtask_list 261 ] 262 }, 263 ) 264 265 # Phase 3: Synthesize 266 with trace.span("supervisor.synthesize") as syn_span: 267 final_output = await self._synthesize(task, subtask_list) 268 syn_span.set_output(final_output) 269 270 self._log_finished(rounds=1, success=True) 271 trace.set_output(final_output) 272 self.complete_run( 273 run_id=run_id, 274 final_output=final_output, 275 agent_outputs=agent_outputs, 276 rounds=1, 277 ) 278 return OrchestratorResult( 279 final_output=final_output, 280 agent_outputs=agent_outputs, 281 subtask_outputs=subtask_list, 282 rounds=1, 283 success=True, 284 ) 285 286 except Exception as exc: 287 self._logger.error("SupervisorOrchestrator failed", error=str(exc)) 288 trace.set_error(exc) 289 self.fail_run( 290 run_id=run_id, 291 error=str(exc), 292 rounds=1, 293 ) 294 return OrchestratorResult( 295 final_output="", 296 rounds=1, 297 success=False, 298 error=str(exc), 299 )
Execute the multi-agent orchestration and return the aggregated result.
14class PipelineOrchestrator(BaseOrchestrator): 15 """ 16 Runs a sequence of agents where the output of each becomes the input to the next. 17 18 The first agent in the pipeline receives the original ``task``. Each 19 subsequent agent receives the previous agent's ``output`` as its task. 20 21 Args: 22 agents: Ordered list of :class:`BaseAgent` instances. 23 pass_context: If True, the full context dict (including all prior outputs) 24 is also forwarded to each agent (default: True). 25 relay_output: If True (default) each agent receives the previous agent's 26 output as its task (classic pipeline chaining). If False, every agent 27 receives the original ``task`` instead, so each one acts on the user's 28 request directly (prior outputs remain available via context). 29 logger, metrics, tracer: Observability (optional). 30 31 Example:: 32 33 pipeline = PipelineOrchestrator( 34 agents=[search_agent, summarise_agent, translate_agent] 35 ) 36 result = await pipeline.run("Summarise the latest AI news in Spanish") 37 """ 38 39 def __init__( 40 self, 41 agents: Optional[List[BaseAgent]] = None, 42 pass_context: bool = True, 43 relay_output: bool = True, 44 router: Optional[BaseRouter] = None, 45 logger: Optional[BasicLogger] = None, 46 metrics: Optional[BasicMetricsCollector] = None, 47 tracer: Optional[TracingProvider] = None, 48 registry: Optional[Any] = None, 49 ) -> None: 50 super().__init__( 51 router=router, 52 logger=logger, 53 metrics=metrics, 54 tracer=tracer, 55 registry=registry, 56 ) 57 self.agents: List[BaseAgent] = agents or [] 58 self.pass_context = pass_context 59 self.relay_output = relay_output 60 61 def get_orchestrator_type(self) -> str: 62 return "pipeline" 63 64 def list_configured_agents(self) -> List[Dict[str, Any]]: 65 return [ 66 { 67 "agent_id": a.agent_id, 68 "name": a.agent_id, 69 "role": None, 70 "model": None, 71 "metadata": {"endpoint_url": getattr(a, "endpoint_url", None)}, 72 } 73 for a in self.agents 74 ] 75 76 def _add_agent_now(self, payload: Dict[str, Any]) -> None: 77 agent_obj = payload.get("agent") 78 if agent_obj is None: 79 raise ValueError("Pipeline add agent requires payload['agent']") 80 self.agents.append(agent_obj) 81 82 def _remove_agent_now(self, agent_id: str) -> None: 83 if len(self.agents) <= 1: 84 raise ValueError("Cannot remove the last agent from pipeline orchestrator") 85 original_len = len(self.agents) 86 self.agents = [a for a in self.agents if a.agent_id != agent_id] 87 if len(self.agents) == original_len: 88 raise ValueError(f"Agent not found: {agent_id}") 89 90 async def run(self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult: 91 self._log_start(task) 92 try: 93 run_id = self.begin_run() 94 except Exception as exc: 95 self._logger.error("PipelineOrchestrator failed before run start", error=str(exc)) 96 return OrchestratorResult( 97 final_output=task, 98 rounds=0, 99 success=False, 100 error=str(exc), 101 ) 102 103 agent_outputs: Dict[str, AgentResult] = {} 104 step_records: List[Dict[str, Any]] = [] 105 current_task = task 106 final_output = task # last agent's output; used as the pipeline result 107 # Merge caller-supplied context (user_assertion, etc.) so it flows to agents; 108 # step outputs are accumulated on top under step_N_output keys. 109 pipeline_context: Dict[str, Any] = {"original_task": task, **(context or {})} 110 111 with self._tracer.trace( 112 "pipeline.run", 113 input=task, 114 metadata={"pipeline_length": len(self.agents)}, 115 ) as trace: 116 try: 117 for i, agent in enumerate(self.agents): 118 # If a router is provided, let it pick the best agent from 119 # the pool for the current task rather than using fixed order. 120 if self.router: 121 available = [a.agent_id for a in self.agents] 122 decision = await self.router.route( 123 RoutingRequest(input=current_task, available_agents=available) 124 ) 125 agent = next( 126 (a for a in self.agents if a.agent_id == decision.target), 127 self.agents[i], # fallback to positional agent 128 ) 129 self._logger.info( 130 "Router selected pipeline agent", 131 step=i, 132 agent=agent.agent_id, 133 confidence=decision.confidence, 134 ) 135 agent_id = agent.agent_id or f"agent_{i}" 136 self._log_agent_dispatch(agent_id, current_task) 137 138 with trace.span(f"pipeline.step_{i}.{agent_id}", input=current_task) as span: 139 # pass_context=True: full accumulated pipeline context (step outputs included) 140 # pass_context=False: only caller's original context (no step outputs) 141 exec_context = pipeline_context if self.pass_context else (context or None) 142 result = await agent.execute(current_task, context=exec_context) 143 output_key = f"step_{i}_{agent_id}" 144 agent_outputs[output_key] = result 145 pipeline_context[f"step_{i}_output"] = result.output 146 span.set_output(result.output) 147 step_records.append( 148 { 149 "step_index": i, 150 "agent_id": agent_id, 151 "input": current_task, 152 "output": result.output, 153 "success": result.success, 154 "steps": [s.__dict__ for s in result.steps], 155 } 156 ) 157 158 if not result.success: 159 self._logger.warning( 160 "Pipeline step failed", 161 step=i, 162 agent=agent_id, 163 error=result.error, 164 ) 165 166 # In classic chaining, this step's output becomes the next 167 # task. With relay_output=False, every agent instead acts on 168 # the original task (prior outputs stay available via context). 169 final_output = result.output 170 if self.relay_output: 171 current_task = result.output 172 173 self._log_finished(rounds=len(self.agents), success=True) 174 trace.set_output(final_output) 175 self.update_run_type_specific( 176 run_id, 177 { 178 "step_sequence": step_records, 179 "total_steps": len(self.agents), 180 "completed_steps": len(agent_outputs), 181 }, 182 ) 183 self.complete_run( 184 run_id=run_id, 185 final_output=final_output, 186 agent_outputs=agent_outputs, 187 rounds=len(self.agents), 188 ) 189 return OrchestratorResult( 190 final_output=final_output, 191 agent_outputs=agent_outputs, 192 rounds=len(self.agents), 193 success=True, 194 run_id=run_id, 195 ) 196 197 except Exception as exc: 198 self._logger.error("PipelineOrchestrator failed", error=str(exc)) 199 trace.set_error(exc) 200 self.fail_run( 201 run_id=run_id, 202 error=str(exc), 203 agent_outputs=agent_outputs, 204 rounds=len(agent_outputs), 205 final_output=current_task, 206 ) 207 return OrchestratorResult( 208 final_output=current_task, 209 agent_outputs=agent_outputs, 210 rounds=len(agent_outputs), 211 success=False, 212 error=str(exc), 213 )
Runs a sequence of agents where the output of each becomes the input to the next.
The first agent in the pipeline receives the original task. Each
subsequent agent receives the previous agent's output as its task.
Args:
agents: Ordered list of BaseAgent instances.
pass_context: If True, the full context dict (including all prior outputs)
is also forwarded to each agent (default: True).
relay_output: If True (default) each agent receives the previous agent's
output as its task (classic pipeline chaining). If False, every agent
receives the original task instead, so each one acts on the user's
request directly (prior outputs remain available via context).
logger, metrics, tracer: Observability (optional).
Example::
pipeline = PipelineOrchestrator(
agents=[search_agent, summarise_agent, translate_agent]
)
result = await pipeline.run("Summarise the latest AI news in Spanish")
39 def __init__( 40 self, 41 agents: Optional[List[BaseAgent]] = None, 42 pass_context: bool = True, 43 relay_output: bool = True, 44 router: Optional[BaseRouter] = None, 45 logger: Optional[BasicLogger] = None, 46 metrics: Optional[BasicMetricsCollector] = None, 47 tracer: Optional[TracingProvider] = None, 48 registry: Optional[Any] = None, 49 ) -> None: 50 super().__init__( 51 router=router, 52 logger=logger, 53 metrics=metrics, 54 tracer=tracer, 55 registry=registry, 56 ) 57 self.agents: List[BaseAgent] = agents or [] 58 self.pass_context = pass_context 59 self.relay_output = relay_output
64 def list_configured_agents(self) -> List[Dict[str, Any]]: 65 return [ 66 { 67 "agent_id": a.agent_id, 68 "name": a.agent_id, 69 "role": None, 70 "model": None, 71 "metadata": {"endpoint_url": getattr(a, "endpoint_url", None)}, 72 } 73 for a in self.agents 74 ]
Return configured orchestrator agents for management and discovery.
90 async def run(self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult: 91 self._log_start(task) 92 try: 93 run_id = self.begin_run() 94 except Exception as exc: 95 self._logger.error("PipelineOrchestrator failed before run start", error=str(exc)) 96 return OrchestratorResult( 97 final_output=task, 98 rounds=0, 99 success=False, 100 error=str(exc), 101 ) 102 103 agent_outputs: Dict[str, AgentResult] = {} 104 step_records: List[Dict[str, Any]] = [] 105 current_task = task 106 final_output = task # last agent's output; used as the pipeline result 107 # Merge caller-supplied context (user_assertion, etc.) so it flows to agents; 108 # step outputs are accumulated on top under step_N_output keys. 109 pipeline_context: Dict[str, Any] = {"original_task": task, **(context or {})} 110 111 with self._tracer.trace( 112 "pipeline.run", 113 input=task, 114 metadata={"pipeline_length": len(self.agents)}, 115 ) as trace: 116 try: 117 for i, agent in enumerate(self.agents): 118 # If a router is provided, let it pick the best agent from 119 # the pool for the current task rather than using fixed order. 120 if self.router: 121 available = [a.agent_id for a in self.agents] 122 decision = await self.router.route( 123 RoutingRequest(input=current_task, available_agents=available) 124 ) 125 agent = next( 126 (a for a in self.agents if a.agent_id == decision.target), 127 self.agents[i], # fallback to positional agent 128 ) 129 self._logger.info( 130 "Router selected pipeline agent", 131 step=i, 132 agent=agent.agent_id, 133 confidence=decision.confidence, 134 ) 135 agent_id = agent.agent_id or f"agent_{i}" 136 self._log_agent_dispatch(agent_id, current_task) 137 138 with trace.span(f"pipeline.step_{i}.{agent_id}", input=current_task) as span: 139 # pass_context=True: full accumulated pipeline context (step outputs included) 140 # pass_context=False: only caller's original context (no step outputs) 141 exec_context = pipeline_context if self.pass_context else (context or None) 142 result = await agent.execute(current_task, context=exec_context) 143 output_key = f"step_{i}_{agent_id}" 144 agent_outputs[output_key] = result 145 pipeline_context[f"step_{i}_output"] = result.output 146 span.set_output(result.output) 147 step_records.append( 148 { 149 "step_index": i, 150 "agent_id": agent_id, 151 "input": current_task, 152 "output": result.output, 153 "success": result.success, 154 "steps": [s.__dict__ for s in result.steps], 155 } 156 ) 157 158 if not result.success: 159 self._logger.warning( 160 "Pipeline step failed", 161 step=i, 162 agent=agent_id, 163 error=result.error, 164 ) 165 166 # In classic chaining, this step's output becomes the next 167 # task. With relay_output=False, every agent instead acts on 168 # the original task (prior outputs stay available via context). 169 final_output = result.output 170 if self.relay_output: 171 current_task = result.output 172 173 self._log_finished(rounds=len(self.agents), success=True) 174 trace.set_output(final_output) 175 self.update_run_type_specific( 176 run_id, 177 { 178 "step_sequence": step_records, 179 "total_steps": len(self.agents), 180 "completed_steps": len(agent_outputs), 181 }, 182 ) 183 self.complete_run( 184 run_id=run_id, 185 final_output=final_output, 186 agent_outputs=agent_outputs, 187 rounds=len(self.agents), 188 ) 189 return OrchestratorResult( 190 final_output=final_output, 191 agent_outputs=agent_outputs, 192 rounds=len(self.agents), 193 success=True, 194 run_id=run_id, 195 ) 196 197 except Exception as exc: 198 self._logger.error("PipelineOrchestrator failed", error=str(exc)) 199 trace.set_error(exc) 200 self.fail_run( 201 run_id=run_id, 202 error=str(exc), 203 agent_outputs=agent_outputs, 204 rounds=len(agent_outputs), 205 final_output=current_task, 206 ) 207 return OrchestratorResult( 208 final_output=current_task, 209 agent_outputs=agent_outputs, 210 rounds=len(agent_outputs), 211 success=False, 212 error=str(exc), 213 )
Execute the multi-agent orchestration and return the aggregated result.
42class DebateOrchestrator(BaseOrchestrator): 43 """ 44 Runs structured multi-agent debate to improve answer quality. 45 46 Each agent generates an initial position. Then for ``debate_rounds`` 47 rounds, each agent critiques other agents' positions and refines its own. 48 A synthesis LLM call (or the first agent's gateway) produces the final answer. 49 50 Args: 51 agents: List of debating :class:`BaseAgent` instances. 52 debate_rounds: Number of critique/refine cycles (default: 1). 53 synthesis_gateway: Optional separate LLM gateway for the final synthesis. 54 If None, uses the first agent's gateway. 55 synthesis_model: Model for synthesis (optional). 56 logger, metrics, tracer: Observability (optional). 57 58 Example:: 59 60 debate = DebateOrchestrator( 61 agents=[agent_a, agent_b, agent_c], 62 debate_rounds=2, 63 ) 64 result = await debate.run("What is the best AI architecture for RAG?") 65 """ 66 67 def __init__( 68 self, 69 agents: Optional[List[BaseAgent]] = None, 70 debate_rounds: int = 1, 71 synthesis_gateway: Optional[Any] = None, 72 synthesis_model: Optional[str] = None, 73 router: Optional[BaseRouter] = None, 74 logger: Optional[BasicLogger] = None, 75 metrics: Optional[BasicMetricsCollector] = None, 76 tracer: Optional[TracingProvider] = None, 77 registry: Optional[Any] = None, 78 ) -> None: 79 super().__init__( 80 router=router, 81 logger=logger, 82 metrics=metrics, 83 tracer=tracer, 84 registry=registry, 85 ) 86 self.agents: List[BaseAgent] = agents or [] 87 self.debate_rounds = debate_rounds 88 self._synthesis_gateway = synthesis_gateway 89 self.synthesis_model = synthesis_model 90 91 def get_orchestrator_type(self) -> str: 92 return "debate" 93 94 def list_configured_agents(self) -> List[Dict[str, Any]]: 95 return [ 96 { 97 "agent_id": a.agent_id, 98 "name": a.agent_id, 99 "role": "debater", 100 "model": None, 101 "metadata": {}, 102 } 103 for a in self.agents 104 ] 105 106 def _add_agent_now(self, payload: Dict[str, Any]) -> None: 107 agent_obj = payload.get("agent") 108 if agent_obj is None: 109 raise ValueError("Debate add agent requires payload['agent']") 110 self.agents.append(agent_obj) 111 112 def _remove_agent_now(self, agent_id: str) -> None: 113 if len(self.agents) <= 1: 114 raise ValueError("Cannot remove the last agent from debate orchestrator") 115 original_len = len(self.agents) 116 self.agents = [a for a in self.agents if a.agent_id != agent_id] 117 if len(self.agents) == original_len: 118 raise ValueError(f"Agent not found: {agent_id}") 119 120 def _synthesis_gw(self) -> Any: 121 if self._synthesis_gateway: 122 return self._synthesis_gateway 123 if self.agents: 124 return self.agents[0].llm_gateway 125 raise ValueError("DebateOrchestrator: no agents or synthesis_gateway provided.") 126 127 async def run(self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult: 128 self._log_start(task) 129 try: 130 run_id = self.begin_run() 131 except Exception as exc: 132 self._logger.error("DebateOrchestrator failed before run start", error=str(exc)) 133 return OrchestratorResult( 134 final_output="", 135 rounds=0, 136 success=False, 137 error=str(exc), 138 ) 139 140 with self._tracer.trace( 141 "debate.run", 142 input=task, 143 metadata={"agents": len(self.agents), "rounds": self.debate_rounds}, 144 ) as trace: 145 try: 146 # Phase 1: Initial positions 147 positions: Dict[str, str] = {} 148 agent_outputs: Dict[str, AgentResult] = {} 149 debate_rounds_data: List[Dict[str, Any]] = [] 150 151 with trace.span("debate.initial_positions") as init_span: 152 for agent in self.agents: 153 prompt = _DEBATE_INITIAL_PROMPT.format(task=task) 154 result = await agent.execute(prompt, context=context) 155 positions[agent.agent_id] = result.output 156 agent_outputs[agent.agent_id] = result 157 158 debate_rounds_data.append( 159 { 160 "round_num": 0, 161 "positions": dict(positions), 162 } 163 ) 164 165 init_span.set_output({"agents": list(positions.keys())}) 166 167 # Phase 2: Critique rounds 168 for round_num in range(self.debate_rounds): 169 self._logger.info("Debate round", round=round_num + 1) 170 if self._metrics: 171 self._metrics.increment("orchestrator.debate_rounds") 172 173 with trace.span(f"debate.round_{round_num + 1}") as round_span: 174 new_positions: Dict[str, str] = {} 175 for agent in self.agents: 176 others = "\n\n".join( 177 f"[{aid}]: {pos}" 178 for aid, pos in positions.items() 179 if aid != agent.agent_id 180 ) 181 critique_prompt = _DEBATE_CRITIQUE_PROMPT.format( 182 task=task, 183 positions=others, 184 own_position=positions.get(agent.agent_id, ""), 185 ) 186 result = await agent.execute(critique_prompt, context=context) 187 new_positions[agent.agent_id] = result.output 188 agent_outputs[agent.agent_id] = result 189 190 positions = new_positions 191 debate_rounds_data.append( 192 { 193 "round_num": round_num + 1, 194 "positions": dict(positions), 195 } 196 ) 197 round_span.set_output({"positions_updated": len(positions)}) 198 199 # Phase 3: Synthesis 200 with trace.span("debate.synthesis") as syn_span: 201 all_positions = "\n\n".join( 202 f"[{aid}] (round {self.debate_rounds}):\n{pos}" 203 for aid, pos in positions.items() 204 ) 205 synthesis_prompt = _DEBATE_SYNTHESIS_PROMPT.format( 206 task=task, 207 rounds=self.debate_rounds, 208 all_positions=all_positions, 209 ) 210 synthesis_response = await self._synthesis_gw().complete( 211 synthesis_prompt, model=self.synthesis_model, temperature=0.1 212 ) 213 final_output = synthesis_response.content.strip() 214 syn_span.set_output(final_output) 215 216 self._log_finished(rounds=self.debate_rounds, success=True) 217 trace.set_output(final_output) 218 self.update_run_type_specific( 219 run_id, 220 { 221 "debate_rounds": debate_rounds_data, 222 "total_rounds": self.debate_rounds, 223 "synthesis": final_output, 224 }, 225 ) 226 self.complete_run( 227 run_id=run_id, 228 final_output=final_output, 229 agent_outputs=agent_outputs, 230 rounds=self.debate_rounds, 231 ) 232 return OrchestratorResult( 233 final_output=final_output, 234 agent_outputs=agent_outputs, 235 rounds=self.debate_rounds, 236 success=True, 237 ) 238 239 except Exception as exc: 240 self._logger.error("DebateOrchestrator failed", error=str(exc)) 241 trace.set_error(exc) 242 self.fail_run( 243 run_id=run_id, 244 error=str(exc), 245 rounds=self.debate_rounds, 246 ) 247 return OrchestratorResult( 248 final_output="", 249 rounds=self.debate_rounds, 250 success=False, 251 error=str(exc), 252 )
Runs structured multi-agent debate to improve answer quality.
Each agent generates an initial position. Then for debate_rounds
rounds, each agent critiques other agents' positions and refines its own.
A synthesis LLM call (or the first agent's gateway) produces the final answer.
Args:
agents: List of debating BaseAgent instances.
debate_rounds: Number of critique/refine cycles (default: 1).
synthesis_gateway: Optional separate LLM gateway for the final synthesis.
If None, uses the first agent's gateway.
synthesis_model: Model for synthesis (optional).
logger, metrics, tracer: Observability (optional).
Example::
debate = DebateOrchestrator(
agents=[agent_a, agent_b, agent_c],
debate_rounds=2,
)
result = await debate.run("What is the best AI architecture for RAG?")
67 def __init__( 68 self, 69 agents: Optional[List[BaseAgent]] = None, 70 debate_rounds: int = 1, 71 synthesis_gateway: Optional[Any] = None, 72 synthesis_model: Optional[str] = None, 73 router: Optional[BaseRouter] = None, 74 logger: Optional[BasicLogger] = None, 75 metrics: Optional[BasicMetricsCollector] = None, 76 tracer: Optional[TracingProvider] = None, 77 registry: Optional[Any] = None, 78 ) -> None: 79 super().__init__( 80 router=router, 81 logger=logger, 82 metrics=metrics, 83 tracer=tracer, 84 registry=registry, 85 ) 86 self.agents: List[BaseAgent] = agents or [] 87 self.debate_rounds = debate_rounds 88 self._synthesis_gateway = synthesis_gateway 89 self.synthesis_model = synthesis_model
94 def list_configured_agents(self) -> List[Dict[str, Any]]: 95 return [ 96 { 97 "agent_id": a.agent_id, 98 "name": a.agent_id, 99 "role": "debater", 100 "model": None, 101 "metadata": {}, 102 } 103 for a in self.agents 104 ]
Return configured orchestrator agents for management and discovery.
127 async def run(self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult: 128 self._log_start(task) 129 try: 130 run_id = self.begin_run() 131 except Exception as exc: 132 self._logger.error("DebateOrchestrator failed before run start", error=str(exc)) 133 return OrchestratorResult( 134 final_output="", 135 rounds=0, 136 success=False, 137 error=str(exc), 138 ) 139 140 with self._tracer.trace( 141 "debate.run", 142 input=task, 143 metadata={"agents": len(self.agents), "rounds": self.debate_rounds}, 144 ) as trace: 145 try: 146 # Phase 1: Initial positions 147 positions: Dict[str, str] = {} 148 agent_outputs: Dict[str, AgentResult] = {} 149 debate_rounds_data: List[Dict[str, Any]] = [] 150 151 with trace.span("debate.initial_positions") as init_span: 152 for agent in self.agents: 153 prompt = _DEBATE_INITIAL_PROMPT.format(task=task) 154 result = await agent.execute(prompt, context=context) 155 positions[agent.agent_id] = result.output 156 agent_outputs[agent.agent_id] = result 157 158 debate_rounds_data.append( 159 { 160 "round_num": 0, 161 "positions": dict(positions), 162 } 163 ) 164 165 init_span.set_output({"agents": list(positions.keys())}) 166 167 # Phase 2: Critique rounds 168 for round_num in range(self.debate_rounds): 169 self._logger.info("Debate round", round=round_num + 1) 170 if self._metrics: 171 self._metrics.increment("orchestrator.debate_rounds") 172 173 with trace.span(f"debate.round_{round_num + 1}") as round_span: 174 new_positions: Dict[str, str] = {} 175 for agent in self.agents: 176 others = "\n\n".join( 177 f"[{aid}]: {pos}" 178 for aid, pos in positions.items() 179 if aid != agent.agent_id 180 ) 181 critique_prompt = _DEBATE_CRITIQUE_PROMPT.format( 182 task=task, 183 positions=others, 184 own_position=positions.get(agent.agent_id, ""), 185 ) 186 result = await agent.execute(critique_prompt, context=context) 187 new_positions[agent.agent_id] = result.output 188 agent_outputs[agent.agent_id] = result 189 190 positions = new_positions 191 debate_rounds_data.append( 192 { 193 "round_num": round_num + 1, 194 "positions": dict(positions), 195 } 196 ) 197 round_span.set_output({"positions_updated": len(positions)}) 198 199 # Phase 3: Synthesis 200 with trace.span("debate.synthesis") as syn_span: 201 all_positions = "\n\n".join( 202 f"[{aid}] (round {self.debate_rounds}):\n{pos}" 203 for aid, pos in positions.items() 204 ) 205 synthesis_prompt = _DEBATE_SYNTHESIS_PROMPT.format( 206 task=task, 207 rounds=self.debate_rounds, 208 all_positions=all_positions, 209 ) 210 synthesis_response = await self._synthesis_gw().complete( 211 synthesis_prompt, model=self.synthesis_model, temperature=0.1 212 ) 213 final_output = synthesis_response.content.strip() 214 syn_span.set_output(final_output) 215 216 self._log_finished(rounds=self.debate_rounds, success=True) 217 trace.set_output(final_output) 218 self.update_run_type_specific( 219 run_id, 220 { 221 "debate_rounds": debate_rounds_data, 222 "total_rounds": self.debate_rounds, 223 "synthesis": final_output, 224 }, 225 ) 226 self.complete_run( 227 run_id=run_id, 228 final_output=final_output, 229 agent_outputs=agent_outputs, 230 rounds=self.debate_rounds, 231 ) 232 return OrchestratorResult( 233 final_output=final_output, 234 agent_outputs=agent_outputs, 235 rounds=self.debate_rounds, 236 success=True, 237 ) 238 239 except Exception as exc: 240 self._logger.error("DebateOrchestrator failed", error=str(exc)) 241 trace.set_error(exc) 242 self.fail_run( 243 run_id=run_id, 244 error=str(exc), 245 rounds=self.debate_rounds, 246 ) 247 return OrchestratorResult( 248 final_output="", 249 rounds=self.debate_rounds, 250 success=False, 251 error=str(exc), 252 )
Execute the multi-agent orchestration and return the aggregated result.
29class SwarmOrchestrator(BaseOrchestrator): 30 """ 31 Dynamically routes work to the best available agent each round. 32 33 Each round: 34 1. A coordinator LLM call decides whether the task is done or describes the 35 next sub-task. 36 2. The :class:`BaseRouter` selects the agent best suited for that sub-task. 37 3. The selected agent executes the sub-task. 38 4. Results accumulate until ``DONE`` is signalled or ``max_rounds`` is reached. 39 40 Args: 41 coordinator_gateway: LLM gateway used by the coordinator. 42 agents: Mapping of agent name → :class:`BaseAgent`. 43 router: :class:`BaseRouter` that selects agents each round. 44 max_rounds: Maximum dispatch rounds (default: 10). 45 coordinator_model: Model for coordinator calls (optional). 46 logger, metrics, tracer: Observability (optional). 47 """ 48 49 def __init__( 50 self, 51 coordinator_gateway: Any, 52 agents: Optional[Dict[str, BaseAgent]] = None, 53 router: Optional[BaseRouter] = None, 54 max_rounds: int = 10, 55 coordinator_model: Optional[str] = None, 56 logger: Optional[BasicLogger] = None, 57 metrics: Optional[BasicMetricsCollector] = None, 58 tracer: Optional[TracingProvider] = None, 59 registry: Optional[Any] = None, 60 ) -> None: 61 super().__init__( 62 router=router, 63 logger=logger, 64 metrics=metrics, 65 tracer=tracer, 66 registry=registry, 67 ) 68 self.coordinator_gateway = coordinator_gateway 69 self.agents: Dict[str, BaseAgent] = agents or {} 70 self.max_rounds = max_rounds 71 self.coordinator_model = coordinator_model 72 73 def get_orchestrator_type(self) -> str: 74 return "swarm" 75 76 def list_configured_agents(self) -> List[Dict[str, Any]]: 77 return [ 78 { 79 "agent_id": agent_id, 80 "name": agent_id, 81 "role": None, 82 "model": None, 83 "metadata": {}, 84 } 85 for agent_id in self.agents.keys() 86 ] 87 88 def _add_agent_now(self, payload: Dict[str, Any]) -> None: 89 agent_id = str(payload["agent_id"]) 90 agent_obj = payload.get("agent") 91 if agent_obj is None: 92 raise ValueError("Swarm add agent requires payload['agent']") 93 self.agents[agent_id] = agent_obj 94 95 def _remove_agent_now(self, agent_id: str) -> None: 96 if agent_id not in self.agents: 97 raise ValueError(f"Agent not found: {agent_id}") 98 if len(self.agents) <= 1: 99 raise ValueError("Cannot remove the last agent from swarm orchestrator") 100 self.agents.pop(agent_id, None) 101 102 async def _next_step(self, task: str, history: List[str]) -> Optional[str]: 103 """Returns None to signal DONE, otherwise the next sub-task description. 104 105 When a router is present the coordinator is asked to name an explicit agent 106 using the structured ``AGENT: / TASK:`` format. The router then matches on 107 the agent name directly rather than performing keyword matching on free text. 108 """ 109 history_text = "\n".join( 110 f"Round {i+1}: {h}" for i, h in enumerate(history) 111 ) or "Nothing yet." 112 agents_list = ", ".join(self.agents.keys()) if self.agents else "(none)" 113 prompt = _SWARM_CONTINUE_PROMPT.format( 114 task=task, history=history_text, agents=agents_list 115 ) 116 response = await self.coordinator_gateway.complete( 117 prompt, model=self.coordinator_model, temperature=0.0 118 ) 119 text = response.content.strip() 120 if text.upper().startswith("DONE"): 121 return None # Task complete 122 return text 123 124 async def run(self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult: 125 self._log_start(task) 126 try: 127 run_id = self.begin_run() 128 except Exception as exc: 129 self._logger.error("SwarmOrchestrator failed before run start", error=str(exc)) 130 return OrchestratorResult( 131 final_output="", 132 rounds=0, 133 success=False, 134 error=str(exc), 135 ) 136 137 agent_outputs: Dict[str, AgentResult] = {} 138 history: List[str] = [] 139 rounds_data: List[Dict[str, Any]] = [] 140 final_output = "" 141 available = list(self.agents.keys()) 142 143 with self._tracer.trace( 144 "swarm.run", input=task, metadata={"max_rounds": self.max_rounds} 145 ) as trace: 146 try: 147 for round_num in range(self.max_rounds): 148 with trace.span(f"swarm.round_{round_num + 1}") as round_span: 149 # Coordinator decides next step 150 next_task = await self._next_step(task, history) 151 if next_task is None: 152 # Extract final answer from DONE line 153 if history: 154 done_response = await self.coordinator_gateway.complete( 155 f"Summarise the work done:\n" + "\n".join(history), 156 model=self.coordinator_model, 157 temperature=0.1, 158 ) 159 final_output = done_response.content.strip() 160 round_span.set_output("DONE") 161 break 162 163 # Route to best agent 164 routing_req = RoutingRequest( 165 input=next_task, available_agents=available 166 ) 167 if self.router: 168 decision = await self.router.route(routing_req) 169 agent_name = decision.target 170 else: 171 agent_name = available[round_num % len(available)] 172 173 agent = self.agents.get(agent_name) 174 if agent is None: 175 self._logger.warning("Swarm: agent not found", agent=agent_name) 176 continue 177 178 # Extract just the task description from structured AGENT:/TASK: format 179 agent_task = next_task 180 for line in next_task.splitlines(): 181 stripped = line.strip() 182 if stripped.upper().startswith("TASK:"): 183 agent_task = stripped[5:].strip() 184 break 185 186 self._log_agent_dispatch(agent_name, agent_task) 187 result = await agent.execute(agent_task, context=context) 188 history.append(f"[{agent_name}]: {result.output[:300]}") 189 agent_outputs[f"round_{round_num + 1}_{agent_name}"] = result 190 final_output = result.output 191 rounds_data.append( 192 { 193 "round_num": round_num + 1, 194 "coordinator_decision": next_task, 195 "agent_dispatched": agent_name, 196 "agent_task": agent_task, 197 "output": result.output, 198 "success": result.success, 199 } 200 ) 201 round_span.set_output(result.output[:200]) 202 203 if self._metrics: 204 self._metrics.increment( 205 "orchestrator.swarm_dispatches", agent=agent_name 206 ) 207 208 self._log_finished(rounds=len(history), success=True) 209 trace.set_output(final_output) 210 self.update_run_type_specific( 211 run_id, 212 { 213 "rounds": rounds_data, 214 "max_rounds": self.max_rounds, 215 }, 216 ) 217 self.complete_run( 218 run_id=run_id, 219 final_output=final_output, 220 agent_outputs=agent_outputs, 221 rounds=len(history), 222 ) 223 return OrchestratorResult( 224 final_output=final_output, 225 agent_outputs=agent_outputs, 226 rounds=len(history), 227 success=True, 228 ) 229 230 except Exception as exc: 231 self._logger.error("SwarmOrchestrator failed", error=str(exc)) 232 trace.set_error(exc) 233 self.fail_run( 234 run_id=run_id, 235 error=str(exc), 236 agent_outputs=agent_outputs, 237 rounds=len(history), 238 final_output=final_output, 239 ) 240 return OrchestratorResult( 241 final_output=final_output, 242 agent_outputs=agent_outputs, 243 rounds=len(history), 244 success=False, 245 error=str(exc), 246 )
Dynamically routes work to the best available agent each round.
Each round:
- A coordinator LLM call decides whether the task is done or describes the next sub-task.
- The
BaseRouterselects the agent best suited for that sub-task. - The selected agent executes the sub-task.
- Results accumulate until
DONEis signalled ormax_roundsis reached.
Args:
coordinator_gateway: LLM gateway used by the coordinator.
agents: Mapping of agent name → BaseAgent.
router: BaseRouter that selects agents each round.
max_rounds: Maximum dispatch rounds (default: 10).
coordinator_model: Model for coordinator calls (optional).
logger, metrics, tracer: Observability (optional).
49 def __init__( 50 self, 51 coordinator_gateway: Any, 52 agents: Optional[Dict[str, BaseAgent]] = None, 53 router: Optional[BaseRouter] = None, 54 max_rounds: int = 10, 55 coordinator_model: Optional[str] = None, 56 logger: Optional[BasicLogger] = None, 57 metrics: Optional[BasicMetricsCollector] = None, 58 tracer: Optional[TracingProvider] = None, 59 registry: Optional[Any] = None, 60 ) -> None: 61 super().__init__( 62 router=router, 63 logger=logger, 64 metrics=metrics, 65 tracer=tracer, 66 registry=registry, 67 ) 68 self.coordinator_gateway = coordinator_gateway 69 self.agents: Dict[str, BaseAgent] = agents or {} 70 self.max_rounds = max_rounds 71 self.coordinator_model = coordinator_model
76 def list_configured_agents(self) -> List[Dict[str, Any]]: 77 return [ 78 { 79 "agent_id": agent_id, 80 "name": agent_id, 81 "role": None, 82 "model": None, 83 "metadata": {}, 84 } 85 for agent_id in self.agents.keys() 86 ]
Return configured orchestrator agents for management and discovery.
124 async def run(self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult: 125 self._log_start(task) 126 try: 127 run_id = self.begin_run() 128 except Exception as exc: 129 self._logger.error("SwarmOrchestrator failed before run start", error=str(exc)) 130 return OrchestratorResult( 131 final_output="", 132 rounds=0, 133 success=False, 134 error=str(exc), 135 ) 136 137 agent_outputs: Dict[str, AgentResult] = {} 138 history: List[str] = [] 139 rounds_data: List[Dict[str, Any]] = [] 140 final_output = "" 141 available = list(self.agents.keys()) 142 143 with self._tracer.trace( 144 "swarm.run", input=task, metadata={"max_rounds": self.max_rounds} 145 ) as trace: 146 try: 147 for round_num in range(self.max_rounds): 148 with trace.span(f"swarm.round_{round_num + 1}") as round_span: 149 # Coordinator decides next step 150 next_task = await self._next_step(task, history) 151 if next_task is None: 152 # Extract final answer from DONE line 153 if history: 154 done_response = await self.coordinator_gateway.complete( 155 f"Summarise the work done:\n" + "\n".join(history), 156 model=self.coordinator_model, 157 temperature=0.1, 158 ) 159 final_output = done_response.content.strip() 160 round_span.set_output("DONE") 161 break 162 163 # Route to best agent 164 routing_req = RoutingRequest( 165 input=next_task, available_agents=available 166 ) 167 if self.router: 168 decision = await self.router.route(routing_req) 169 agent_name = decision.target 170 else: 171 agent_name = available[round_num % len(available)] 172 173 agent = self.agents.get(agent_name) 174 if agent is None: 175 self._logger.warning("Swarm: agent not found", agent=agent_name) 176 continue 177 178 # Extract just the task description from structured AGENT:/TASK: format 179 agent_task = next_task 180 for line in next_task.splitlines(): 181 stripped = line.strip() 182 if stripped.upper().startswith("TASK:"): 183 agent_task = stripped[5:].strip() 184 break 185 186 self._log_agent_dispatch(agent_name, agent_task) 187 result = await agent.execute(agent_task, context=context) 188 history.append(f"[{agent_name}]: {result.output[:300]}") 189 agent_outputs[f"round_{round_num + 1}_{agent_name}"] = result 190 final_output = result.output 191 rounds_data.append( 192 { 193 "round_num": round_num + 1, 194 "coordinator_decision": next_task, 195 "agent_dispatched": agent_name, 196 "agent_task": agent_task, 197 "output": result.output, 198 "success": result.success, 199 } 200 ) 201 round_span.set_output(result.output[:200]) 202 203 if self._metrics: 204 self._metrics.increment( 205 "orchestrator.swarm_dispatches", agent=agent_name 206 ) 207 208 self._log_finished(rounds=len(history), success=True) 209 trace.set_output(final_output) 210 self.update_run_type_specific( 211 run_id, 212 { 213 "rounds": rounds_data, 214 "max_rounds": self.max_rounds, 215 }, 216 ) 217 self.complete_run( 218 run_id=run_id, 219 final_output=final_output, 220 agent_outputs=agent_outputs, 221 rounds=len(history), 222 ) 223 return OrchestratorResult( 224 final_output=final_output, 225 agent_outputs=agent_outputs, 226 rounds=len(history), 227 success=True, 228 ) 229 230 except Exception as exc: 231 self._logger.error("SwarmOrchestrator failed", error=str(exc)) 232 trace.set_error(exc) 233 self.fail_run( 234 run_id=run_id, 235 error=str(exc), 236 agent_outputs=agent_outputs, 237 rounds=len(history), 238 final_output=final_output, 239 ) 240 return OrchestratorResult( 241 final_output=final_output, 242 agent_outputs=agent_outputs, 243 rounds=len(history), 244 success=False, 245 error=str(exc), 246 )
Execute the multi-agent orchestration and return the aggregated result.