Coverage for agentos/workflows/engine.py: 63%
52 statements
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-10 01:20 +0800
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-10 01:20 +0800
1"""
2AgentOS v0.20 预设工作流模板。
3开箱即用的 Agent 协作模式。
4"""
6from __future__ import annotations
8from collections.abc import Callable
9from dataclasses import dataclass, field
10from enum import StrEnum
11from typing import Any
14class WorkflowType(StrEnum):
15 """工作流类型枚举。"""
17 CODE_REVIEW = "code_review"
18 RESEARCH = "research"
19 DEBATE = "debate"
20 QA = "qa"
21 CUSTOM = "custom"
24@dataclass
25class WorkflowStep:
26 """工作流步骤定义。"""
28 agent_role: str
29 instruction: str
30 input_from: int | None = None # 上一步的index,None=原始输入
31 parallel: bool = False
34@dataclass
35class Workflow:
36 """预设工作流定义。"""
38 name: str
39 workflow_type: WorkflowType
40 steps: list[WorkflowStep]
41 max_rounds: int = 5
42 auto_merge: bool = True
43 metadata: dict[str, Any] = field(default_factory=dict)
46# ── 内置工作流 ──────────────────────────────────
48CODE_REVIEW = Workflow(
49 name="代码审查",
50 workflow_type=WorkflowType.CODE_REVIEW,
51 steps=[
52 WorkflowStep("architect", "审查代码架构和设计模式"),
53 WorkflowStep("security_expert", "审查安全漏洞和注入风险"),
54 WorkflowStep("performance_expert", "审查性能瓶颈和资源消耗"),
55 WorkflowStep("reviewer", "综合以上意见,输出最终审查报告"),
56 ],
57 max_rounds=1,
58)
61RESEARCH = Workflow(
62 name="深度调研",
63 workflow_type=WorkflowType.RESEARCH,
64 steps=[
65 WorkflowStep("researcher", "搜索并收集相关资料", parallel=False),
66 WorkflowStep("analyst", "分析数据并提取关键insights", input_from=0),
67 WorkflowStep("synthesizer", "综合所有发现,撰写调研报告", input_from=1),
68 ],
69 max_rounds=1,
70)
73DEBATE = Workflow(
74 name="辩证讨论",
75 workflow_type=WorkflowType.DEBATE,
76 steps=[
77 WorkflowStep("proponent", "提出论点并给出论据"),
78 WorkflowStep("opponent", "反驳对方论点,指出逻辑漏洞", input_from=0),
79 WorkflowStep("judge", "综合双方观点,给出平衡结论", input_from=1),
80 ],
81 max_rounds=3,
82)
85QA = Workflow(
86 name="智能问答",
87 workflow_type=WorkflowType.QA,
88 steps=[
89 WorkflowStep("retriever", "从知识库检索相关信息"),
90 WorkflowStep("reasoner", "基于检索结果进行推理回答", input_from=0),
91 WorkflowStep("verifier", "验证答案准确性并修正", input_from=1),
92 ],
93 max_rounds=2,
94)
97BUILTIN_WORKFLOWS: dict[WorkflowType, Workflow] = {
98 WorkflowType.CODE_REVIEW: CODE_REVIEW,
99 WorkflowType.RESEARCH: RESEARCH,
100 WorkflowType.DEBATE: DEBATE,
101 WorkflowType.QA: QA,
102}
105class WorkflowEngine:
106 """工作流引擎 — 按预设步骤调度多个Agent协作。"""
108 def __init__(self, workflow: Workflow, agent_factory: Callable[[str], Any]):
109 self.workflow = workflow
110 self.agent_factory = agent_factory
111 self._results: dict[int, Any] = {}
113 async def execute(self, input_text: str, context: dict | None = None) -> str:
114 """执行工作流。"""
116 last_output = input_text
117 for round_idx in range(self.workflow.max_rounds):
118 for step_idx, step in enumerate(self.workflow.steps):
119 if step.input_from is not None:
120 feed = self._results.get(step.input_from, input_text)
121 else:
122 feed = input_text if round_idx == 0 else last_output
124 agent = self.agent_factory(step.agent_role)
125 full_prompt = f"""你是一名{step.agent_role}。
127输入内容:
128{feed}
130任务:
131{step.instruction}
133请直接给出你的分析和结论。"""
135 result = await agent.run(full_prompt, context=context or {})
136 self._results[step_idx] = result.get("output", str(result))
137 last_output = self._results[step_idx]
139 if self.workflow.auto_merge:
140 break # 单轮工作流
142 if self.workflow.auto_merge:
143 return self._results.get(len(self.workflow.steps) - 1, last_output)
145 return last_output