health_analysis.py 3.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687
  1. """
  2. 健康分析工作流服务
  3. 负责串联多个 Agent,完成一次完整的健康报告分析
  4. """
  5. import asyncio
  6. from typing import Dict, Any
  7. from uuid import uuid4
  8. from agents.planner import PlannerAgent
  9. from agents.health_indicator import HealthIndicatorAgent
  10. from agents.risk_assess import RiskAssessmentAgent
  11. from agents.advice import AdviceAgent
  12. from agents.report import ReportAgent
  13. from agents.base import create_task, update_agent_state, complete_task
  14. class HealthAnalysisService:
  15. def __init__(self, task_id: str = None):
  16. self.task_id = task_id or str(uuid4())
  17. # 任务初始化
  18. create_task(self.task_id)
  19. self.planner = PlannerAgent(task_id=self.task_id)
  20. self.indicator_agent = HealthIndicatorAgent(task_id=self.task_id)
  21. self.risk_agent = RiskAssessmentAgent(task_id=self.task_id)
  22. self.advice_agent = AdviceAgent(task_id=self.task_id)
  23. self.report_agent = ReportAgent(task_id=self.task_id)
  24. async def run(self, report_text: str) -> Dict[str, Any]:
  25. """
  26. 执行完整的健康分析流程
  27. """
  28. # 1.任务规划
  29. update_agent_state(self.task_id, "PlannerAgent", "running")
  30. plan_result = await self.planner.run({"goal": f"分析以下体检报告并制定执行计划:\n{report_text}"})
  31. update_agent_state(self.task_id, "PlannerAgent", "completed")
  32. # 2.健康指标分析
  33. update_agent_state(self.task_id, "HealthIndicatorAgent", "running")
  34. indicator_result = await self.indicator_agent.run({
  35. "report_text": report_text,
  36. "plan": plan_result
  37. })
  38. update_agent_state(self.task_id, "HealthIndicatorAgent", "completed", partial_report={"indicator_results": indicator_result})
  39. # 3. 风险评估
  40. update_agent_state(self.task_id, "RiskAssessmentAgent", "running")
  41. risk_result = await self.risk_agent.run({
  42. "indicator_results": indicator_result
  43. })
  44. update_agent_state(self.task_id, "RiskAssessmentAgent", "completed", partial_report={"risk_assessment": risk_result})
  45. # 4. 健康建议生成
  46. update_agent_state(self.task_id, "AdviceAgent", "running")
  47. advice_result = await self.advice_agent.run({
  48. "risk_assessment": risk_result
  49. })
  50. update_agent_state(self.task_id, "AdviceAgent", "completed", partial_report={"advice": advice_result})
  51. # 5. 报告汇总
  52. update_agent_state(self.task_id, "ReportAgent", "running")
  53. final_report = await self.report_agent.run({
  54. "indicators": indicator_result,
  55. "risk_assessment": risk_result,
  56. "advice": advice_result
  57. })
  58. update_agent_state(self.task_id, "ReportAgent", "completed")
  59. complete_task(self.task_id, final_report)
  60. return self.task_id
  61. # ---------- 临时本地验证入口 ----------
  62. async def _demo():
  63. demo_text = """
  64. 男性,28岁,BMI 27.3,血压 145/95 mmHg,
  65. 总胆固醇 6.2 mmol/L,空腹血糖 6.1 mmol/L。
  66. """
  67. workflow = HealthAnalysisService()
  68. result = await workflow.run(demo_text)
  69. print(result)
  70. if __name__ == "__main__":
  71. asyncio.run(_demo())