Run tasks in topological order. Args: inputs: The initial input passed to root tasks (as the same object). Returns: Dict[str, Any]: mapping of terminal task name to its output.
(self, inputs: Any, **kwargs)
| 61 | self.parents[child].append(parent) |
| 62 | |
| 63 | async def run(self, inputs: Any, **kwargs): |
| 64 | """Run tasks in topological order. |
| 65 | |
| 66 | Args: |
| 67 | inputs: The initial input passed to root tasks (as the same object). |
| 68 | |
| 69 | Returns: |
| 70 | Dict[str, Any]: mapping of terminal task name to its output. |
| 71 | """ |
| 72 | outputs: Dict[str, Any] = {} |
| 73 | for task in self.topo_order: |
| 74 | # Prepare input for task |
| 75 | if task in self.roots: |
| 76 | task_input = inputs |
| 77 | else: |
| 78 | parent_outs = [outputs[p] for p in self.parents[task]] |
| 79 | task_input = parent_outs if len( |
| 80 | parent_outs) > 1 else parent_outs[0] |
| 81 | |
| 82 | task_info: DictConfig = getattr(self.config, task) |
| 83 | agent_cfg_path = os.path.join(self.config.local_dir, |
| 84 | task_info.agent_config) |
| 85 | if not hasattr(task_info, 'agent'): |
| 86 | task_info.agent = DictConfig({}) |
| 87 | init_args = getattr(task_info.agent, 'kwargs', {}) |
| 88 | init_args['trust_remote_code'] = self.trust_remote_code |
| 89 | init_args['mcp_server_file'] = self.mcp_server_file |
| 90 | init_args['task'] = task |
| 91 | init_args['load_cache'] = self.load_cache |
| 92 | init_args['config_dir_or_id'] = agent_cfg_path |
| 93 | init_args['env'] = self.env |
| 94 | if 'tag' not in init_args: |
| 95 | init_args['tag'] = task |
| 96 | engine = AgentLoader.build(**init_args) |
| 97 | result = await engine.run(task_input) |
| 98 | outputs[task] = result |
| 99 | |
| 100 | # Return results of terminal nodes (no outgoing edges) |
| 101 | terminals = [ |
| 102 | t for t in self.config.keys() |
| 103 | if t not in self.graph and t in self.nodes |
| 104 | ] |
| 105 | return {t: outputs[t] for t in terminals} |