Coverage for agentos/workflows/engine.py: 63%

52 statements  

« prev     ^ index     » next       coverage.py v7.14.3, created at 2026-07-08 21:26 +0800

1""" 

2AgentOS v0.20 预设工作流模板。 

3开箱即用的 Agent 协作模式。 

4""" 

5 

6from __future__ import annotations 

7 

8from collections.abc import Callable 

9from dataclasses import dataclass, field 

10from enum import StrEnum 

11from typing import Any 

12 

13 

14class WorkflowType(StrEnum): 

15 """工作流类型枚举。""" 

16 

17 CODE_REVIEW = "code_review" 

18 RESEARCH = "research" 

19 DEBATE = "debate" 

20 QA = "qa" 

21 CUSTOM = "custom" 

22 

23 

24@dataclass 

25class WorkflowStep: 

26 """工作流步骤定义。""" 

27 

28 agent_role: str 

29 instruction: str 

30 input_from: int | None = None # 上一步的index,None=原始输入 

31 parallel: bool = False 

32 

33 

34@dataclass 

35class Workflow: 

36 """预设工作流定义。""" 

37 

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) 

44 

45 

46# ── 内置工作流 ────────────────────────────────── 

47 

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) 

59 

60 

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) 

71 

72 

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) 

83 

84 

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) 

95 

96 

97BUILTIN_WORKFLOWS: dict[WorkflowType, Workflow] = { 

98 WorkflowType.CODE_REVIEW: CODE_REVIEW, 

99 WorkflowType.RESEARCH: RESEARCH, 

100 WorkflowType.DEBATE: DEBATE, 

101 WorkflowType.QA: QA, 

102} 

103 

104 

105class WorkflowEngine: 

106 """工作流引擎 — 按预设步骤调度多个Agent协作。""" 

107 

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] = {} 

112 

113 async def execute(self, input_text: str, context: dict | None = None) -> str: 

114 """执行工作流。""" 

115 

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 

123 

124 agent = self.agent_factory(step.agent_role) 

125 full_prompt = f"""你是一名{step.agent_role}。 

126 

127输入内容: 

128{feed} 

129 

130任务: 

131{step.instruction} 

132 

133请直接给出你的分析和结论。""" 

134 

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] 

138 

139 if self.workflow.auto_merge: 

140 break # 单轮工作流 

141 

142 if self.workflow.auto_merge: 

143 return self._results.get(len(self.workflow.steps) - 1, last_output) 

144 

145 return last_output