Coverage for agentos/workflows/templates.py: 58%

125 statements  

« prev     ^ index     » next       coverage.py v7.14.3, created at 2026-07-06 19:15 +0800

1""" # noqa: E501 

2Workflow Templates — Declarative, reusable multi-step agent workflows. 

3 

4Define workflows as YAML/JSON templates with conditional branching, 

5parallel execution, retry policies, and human-in-the-loop checkpoints. 

6""" 

7 

8from __future__ import annotations 

9 

10import json 

11from dataclasses import dataclass, field 

12from enum import Enum 

13from typing import Any 

14 

15import yaml 

16 

17 

18class StepType(Enum): 

19 """步骤类型枚举。""" 

20 

21 AGENT = "agent" 

22 TOOL = "tool" 

23 CONDITION = "condition" 

24 PARALLEL = "parallel" 

25 HUMAN_REVIEW = "human_review" 

26 TRANSFORM = "transform" 

27 WAIT = "wait" 

28 

29 

30class RetryPolicy(Enum): 

31 """重试策略类。""" 

32 

33 NONE = "none" 

34 LINEAR = "linear" 

35 EXPONENTIAL = "exponential" 

36 

37 

38@dataclass 

39class WorkflowStep: 

40 """Single step in a workflow template.""" 

41 

42 name: str 

43 step_type: StepType = StepType.AGENT 

44 description: str = "" 

45 

46 # Agent/Tool step 

47 agent_type: str = "default" 

48 task_template: str = "" 

49 tool_name: str = "" 

50 

51 # Condition step 

52 condition: str = "" 

53 """Python expression evaluated with step outputs as variables.""" 

54 

55 then_steps: list[WorkflowStep] = field(default_factory=list) 

56 else_steps: list[WorkflowStep] = field(default_factory=list) 

57 

58 # Parallel step 

59 sub_steps: list[WorkflowStep] = field(default_factory=list) 

60 max_concurrency: int = 5 

61 

62 # Human review 

63 review_prompt: str = "" 

64 timeout_minutes: int = 30 

65 

66 # Retry 

67 retry_policy: RetryPolicy = RetryPolicy.NONE 

68 max_retries: int = 3 

69 retry_delay_seconds: float = 1.0 

70 

71 # Transform 

72 transform_expr: str = "" 

73 """Python expression to transform output.""" 

74 

75 # Input/output 

76 depends_on: list[str] = field(default_factory=list) 

77 output_key: str = "" 

78 """Store output under this key for downstream steps.""" 

79 

80 

81@dataclass 

82class WorkflowTemplate: 

83 """ 

84 Declarative workflow template. 

85 

86 Example (YAML):: 

87 

88 name: research_report 

89 description: Research a topic and generate a report 

90 steps: 

91 - name: research 

92 step_type: agent 

93 agent_type: researcher 

94 task_template: "Research: {{input.topic}}" 

95 output_key: research_result 

96 - name: review 

97 step_type: human_review 

98 review_prompt: "Review the research: {{research_result}}" 

99 depends_on: [research] 

100 - name: write_report 

101 step_type: agent 

102 agent_type: writer 

103 task_template: "Write report based on: {{research_result}}" 

104 depends_on: [review] 

105 output_key: final_report 

106 """ 

107 

108 name: str 

109 description: str = "" 

110 version: str = "1.0" 

111 steps: list[WorkflowStep] = field(default_factory=list) 

112 metadata: dict[str, Any] = field(default_factory=dict) 

113 

114 def to_dict(self) -> dict[str, Any]: 

115 """Serialize to dict.""" 

116 return { 

117 "name": self.name, 

118 "description": self.description, 

119 "version": self.version, 

120 "steps": [self._step_to_dict(s) for s in self.steps], 

121 "metadata": self.metadata, 

122 } 

123 

124 def _step_to_dict(self, step: WorkflowStep) -> dict[str, Any]: 

125 d: dict[str, Any] = { 

126 "name": step.name, 

127 "step_type": step.step_type.value, 

128 "description": step.description, 

129 } 

130 if step.agent_type != "default": 

131 d["agent_type"] = step.agent_type 

132 if step.task_template: 

133 d["task_template"] = step.task_template 

134 if step.tool_name: 

135 d["tool_name"] = step.tool_name 

136 if step.output_key: 

137 d["output_key"] = step.output_key 

138 if step.condition: 

139 d["condition"] = step.condition 

140 if step.depends_on: 

141 d["depends_on"] = step.depends_on 

142 if step.then_steps: 

143 d["then_steps"] = [self._step_to_dict(s) for s in step.then_steps] 

144 if step.else_steps: 

145 d["else_steps"] = [self._step_to_dict(s) for s in step.else_steps] 

146 if step.sub_steps: 

147 d["sub_steps"] = [self._step_to_dict(s) for s in step.sub_steps] 

148 d["max_concurrency"] = step.max_concurrency 

149 if step.retry_policy != RetryPolicy.NONE: 

150 d["retry_policy"] = step.retry_policy.value 

151 d["max_retries"] = step.max_retries 

152 return d 

153 

154 @classmethod 

155 def from_dict(cls, data: dict[str, Any]) -> WorkflowTemplate: 

156 """Deserialize from dict.""" 

157 return cls( 

158 name=data["name"], 

159 description=data.get("description", ""), 

160 version=data.get("version", "1.0"), 

161 steps=[cls._step_from_dict(s) for s in data.get("steps", [])], 

162 metadata=data.get("metadata", {}), 

163 ) 

164 

165 @classmethod 

166 def _step_from_dict(cls, data: dict[str, Any]) -> WorkflowStep: 

167 return WorkflowStep( 

168 name=data["name"], 

169 step_type=StepType(data.get("step_type", "agent")), 

170 description=data.get("description", ""), 

171 agent_type=data.get("agent_type", "default"), 

172 task_template=data.get("task_template", ""), 

173 tool_name=data.get("tool_name", ""), 

174 output_key=data.get("output_key", ""), 

175 condition=data.get("condition", ""), 

176 depends_on=data.get("depends_on", []), 

177 then_steps=[cls._step_from_dict(s) for s in data.get("then_steps", [])], 

178 else_steps=[cls._step_from_dict(s) for s in data.get("else_steps", [])], 

179 sub_steps=[cls._step_from_dict(s) for s in data.get("sub_steps", [])], 

180 max_concurrency=data.get("max_concurrency", 5), 

181 retry_policy=RetryPolicy(data.get("retry_policy", "none")), 

182 max_retries=data.get("max_retries", 3), 

183 retry_delay_seconds=data.get("retry_delay_seconds", 1.0), 

184 ) 

185 

186 def to_yaml(self) -> str: 

187 """Export workflow as YAML string.""" 

188 return yaml.dump(self.to_dict(), default_flow_style=False, sort_keys=False) 

189 

190 def to_json(self, indent: int = 2) -> str: 

191 """Export workflow as JSON string.""" 

192 return json.dumps(self.to_dict(), indent=indent, ensure_ascii=False) 

193 

194 @classmethod 

195 def from_yaml(cls, yaml_str: str) -> WorkflowTemplate: 

196 """Load workflow from YAML string.""" 

197 data = yaml.safe_load(yaml_str) 

198 return cls.from_dict(data) 

199 

200 @classmethod 

201 def from_json(cls, json_str: str) -> WorkflowTemplate: 

202 """Load workflow from JSON string.""" 

203 data = json.loads(json_str) 

204 return cls.from_dict(data) 

205 

206 def get_step(self, name: str) -> WorkflowStep | None: 

207 """Find a step by name (searches recursively).""" 

208 for step in self.steps: 

209 result = self._find_step(step, name) 

210 if result: 

211 return result 

212 return None 

213 

214 def _find_step(self, step: WorkflowStep, name: str) -> WorkflowStep | None: 

215 if step.name == name: 

216 return step 

217 for sub in step.then_steps + step.else_steps + step.sub_steps: 

218 result = self._find_step(sub, name) 

219 if result: 

220 return result 

221 return None 

222 

223 def flatten_steps(self) -> list[WorkflowStep]: 

224 """Return all steps in a flat list.""" 

225 result: list[WorkflowStep] = [] 

226 for step in self.steps: 

227 self._flatten(step, result) 

228 return result 

229 

230 def _flatten(self, step: WorkflowStep, result: list[WorkflowStep]) -> None: 

231 result.append(step) 

232 for sub in step.then_steps + step.else_steps + step.sub_steps: 

233 self._flatten(sub, result) 

234 

235 @property 

236 def step_count(self) -> int: 

237 return len(self.flatten_steps()) 

238 

239 

240# ---- Built-in Workflow Templates ---- 

241 

242BUILTIN_TEMPLATES: dict[str, WorkflowTemplate] = {} 

243 

244 

245def _init_builtins() -> None: 

246 """Initialize built-in workflow templates.""" 

247 # Research → Summarize → Report 

248 BUILTIN_TEMPLATES["research_report"] = WorkflowTemplate( 

249 name="research_report", 

250 description="Research a topic, summarize findings, generate report", 

251 steps=[ 

252 WorkflowStep( 

253 name="research", 

254 step_type=StepType.AGENT, 

255 agent_type="researcher", 

256 task_template="Deep research on: {{input.topic}}", 

257 output_key="research", 

258 ), 

259 WorkflowStep( 

260 name="summarize", 

261 step_type=StepType.AGENT, 

262 agent_type="summarizer", 

263 task_template="Summarize key findings from: {{research}}", 

264 depends_on=["research"], 

265 output_key="summary", 

266 ), 

267 WorkflowStep( 

268 name="report", 

269 step_type=StepType.AGENT, 

270 agent_type="writer", 

271 task_template="Write a comprehensive report based on: {{research}}\\nSummary: {{summary}}", 

272 depends_on=["research", "summarize"], 

273 output_key="report", 

274 ), 

275 ], 

276 ) 

277 

278 # Code Review → Fix → Test 

279 BUILTIN_TEMPLATES["code_review"] = WorkflowTemplate( 

280 name="code_review", 

281 description="Review code, apply fixes, run tests", 

282 steps=[ 

283 WorkflowStep( 

284 name="review", 

285 step_type=StepType.AGENT, 

286 agent_type="code_reviewer", 

287 task_template="Review this code for bugs and improvements:\\n```\\n{{input.code}}\\n```", 

288 output_key="review_feedback", 

289 ), 

290 WorkflowStep( 

291 name="human_approval", 

292 step_type=StepType.HUMAN_REVIEW, 

293 review_prompt="Approve fixes based on: {{review_feedback}}", 

294 depends_on=["review"], 

295 ), 

296 WorkflowStep( 

297 name="apply_fixes", 

298 step_type=StepType.AGENT, 

299 agent_type="coder", 

300 task_template="Apply fixes based on review:\\n{{review_feedback}}\\n\\nOriginal code:\\n```\\n{{input.code}}\\n```", # noqa: E501 

301 depends_on=["human_approval"], 

302 output_key="fixed_code", 

303 retry_policy=RetryPolicy.EXPONENTIAL, 

304 max_retries=3, 

305 ), 

306 WorkflowStep( 

307 name="test", 

308 step_type=StepType.TOOL, 

309 tool_name="run_tests", 

310 depends_on=["apply_fixes"], 

311 ), 

312 ], 

313 ) 

314 

315 

316_init_builtins()