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]
class BaseOrchestrator(abc.ABC):
 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().

async def register(self) -> None:
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.

async def close(self) -> None:
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.

@abstractmethod
async def run( self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult:
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.

@abstractmethod
def get_orchestrator_type(self) -> str:
105    @abstractmethod
106    def get_orchestrator_type(self) -> str:
107        """Return a stable orchestrator type literal."""

Return a stable orchestrator type literal.

@abstractmethod
def list_configured_agents(self) -> List[Dict[str, Any]]:
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.

def get_orchestrator_id(self) -> str:
121    def get_orchestrator_id(self) -> str:
122        return self._orchestrator_id
def get_current_status(self) -> str:
124    def get_current_status(self) -> str:
125        return self._status
def list_runs( self, limit: int = 20, offset: int = 0) -> List[gmf_forge_ai_orchestration.multi_agent.base.RunRecord]:
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]
def get_run( self, run_id: str) -> Optional[gmf_forge_ai_orchestration.multi_agent.base.RunRecord]:
135    def get_run(self, run_id: str) -> Optional[RunRecord]:
136        return self._run_history.get(run_id)
def list_run_agents(self, run_id: str) -> List[str]:
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())
def get_run_agent( self, run_id: str, agent_id: str) -> Optional[gmf_forge_ai_orchestration.AgentResult]:
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)
def queue_add_agent(self, payload: Dict[str, Any]) -> None:
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.

def queue_remove_agent(self, agent_id: str) -> None:
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.

def apply_pending_agent_operations(self) -> None:
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.

def begin_run(self) -> str:
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.

def update_run_type_specific(self, run_id: str, data: Dict[str, Any]) -> None:
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}
def complete_run( self, run_id: str, final_output: str, agent_outputs: Dict[str, gmf_forge_ai_orchestration.AgentResult], rounds: int) -> None:
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)
def fail_run( self, run_id: str, error: str, agent_outputs: Optional[Dict[str, gmf_forge_ai_orchestration.AgentResult]] = None, rounds: int = 0, final_output: str = '') -> None:
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)
@dataclass
class OrchestratorResult:
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.

OrchestratorResult( final_output: str, agent_outputs: Dict[str, gmf_forge_ai_orchestration.AgentResult] = <factory>, subtask_outputs: List[Tuple[str, str, gmf_forge_ai_orchestration.AgentResult]] = <factory>, rounds: int = 0, success: bool = True, error: Optional[str] = None, metadata: Dict[str, Any] = <factory>, run_id: Optional[str] = None)
final_output: str
agent_outputs: Dict[str, gmf_forge_ai_orchestration.AgentResult]
subtask_outputs: List[Tuple[str, str, gmf_forge_ai_orchestration.AgentResult]]
rounds: int = 0
success: bool = True
error: Optional[str] = None
metadata: Dict[str, Any]
run_id: Optional[str] = None
class SupervisorOrchestrator(gmf_forge_ai_orchestration.multi_agent.BaseOrchestrator):
 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).

SupervisorOrchestrator( supervisor_gateway: Any, agents: Optional[Dict[str, gmf_forge_ai_orchestration.BaseAgent]] = None, agent_descriptions: Optional[Dict[str, str]] = None, supervisor_model: Optional[str] = None, router: Optional[gmf_forge_ai_orchestration.BaseRouter] = None, logger: Optional[gmf_forge_ai_shared_core.observability.BasicLogger] = None, metrics: Optional[gmf_forge_ai_shared_core.observability.BasicMetricsCollector] = None, tracer: Optional[gmf_forge_ai_shared_core.observability.TracingProvider] = None, registry: Optional[Any] = None)
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
supervisor_gateway
agent_descriptions: Dict[str, str]
supervisor_model
def get_orchestrator_type(self) -> str:
96    def get_orchestrator_type(self) -> str:
97        return "supervisor"

Return a stable orchestrator type literal.

def list_configured_agents(self) -> List[Dict[str, Any]]:
 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.

async def run( self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult:
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.

class PipelineOrchestrator(gmf_forge_ai_orchestration.multi_agent.BaseOrchestrator):
 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")
PipelineOrchestrator( agents: Optional[List[gmf_forge_ai_orchestration.BaseAgent]] = None, pass_context: bool = True, relay_output: bool = True, router: Optional[gmf_forge_ai_orchestration.BaseRouter] = None, logger: Optional[gmf_forge_ai_shared_core.observability.BasicLogger] = None, metrics: Optional[gmf_forge_ai_shared_core.observability.BasicMetricsCollector] = None, tracer: Optional[gmf_forge_ai_shared_core.observability.TracingProvider] = None, registry: Optional[Any] = None)
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
pass_context
relay_output
def get_orchestrator_type(self) -> str:
61    def get_orchestrator_type(self) -> str:
62        return "pipeline"

Return a stable orchestrator type literal.

def list_configured_agents(self) -> List[Dict[str, Any]]:
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.

async def run( self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult:
 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.

class DebateOrchestrator(gmf_forge_ai_orchestration.multi_agent.BaseOrchestrator):
 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?")
DebateOrchestrator( agents: Optional[List[gmf_forge_ai_orchestration.BaseAgent]] = None, debate_rounds: int = 1, synthesis_gateway: Optional[Any] = None, synthesis_model: Optional[str] = None, router: Optional[gmf_forge_ai_orchestration.BaseRouter] = None, logger: Optional[gmf_forge_ai_shared_core.observability.BasicLogger] = None, metrics: Optional[gmf_forge_ai_shared_core.observability.BasicMetricsCollector] = None, tracer: Optional[gmf_forge_ai_shared_core.observability.TracingProvider] = None, registry: Optional[Any] = None)
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
debate_rounds
synthesis_model
def get_orchestrator_type(self) -> str:
91    def get_orchestrator_type(self) -> str:
92        return "debate"

Return a stable orchestrator type literal.

def list_configured_agents(self) -> List[Dict[str, Any]]:
 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.

async def run( self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult:
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.

class SwarmOrchestrator(gmf_forge_ai_orchestration.multi_agent.BaseOrchestrator):
 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:

  1. A coordinator LLM call decides whether the task is done or describes the next sub-task.
  2. The BaseRouter selects the agent best suited for that sub-task.
  3. The selected agent executes the sub-task.
  4. Results accumulate until DONE is signalled or max_rounds is 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).

SwarmOrchestrator( coordinator_gateway: Any, agents: Optional[Dict[str, gmf_forge_ai_orchestration.BaseAgent]] = None, router: Optional[gmf_forge_ai_orchestration.BaseRouter] = None, max_rounds: int = 10, coordinator_model: Optional[str] = None, logger: Optional[gmf_forge_ai_shared_core.observability.BasicLogger] = None, metrics: Optional[gmf_forge_ai_shared_core.observability.BasicMetricsCollector] = None, tracer: Optional[gmf_forge_ai_shared_core.observability.TracingProvider] = None, registry: Optional[Any] = None)
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
coordinator_gateway
max_rounds
coordinator_model
def get_orchestrator_type(self) -> str:
73    def get_orchestrator_type(self) -> str:
74        return "swarm"

Return a stable orchestrator type literal.

def list_configured_agents(self) -> List[Dict[str, Any]]:
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.

async def run( self, task: str, context: Optional[Dict[str, Any]] = None) -> OrchestratorResult:
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.