health.py 3.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107
  1. from io import BytesIO
  2. from uuid import uuid4
  3. import asyncio
  4. import pdfplumber
  5. from fastapi import APIRouter, File, Form, HTTPException, UploadFile
  6. from pydantic import BaseModel, Field, field_validator
  7. from memory.store import get_report_run, list_report_runs_for_user
  8. from service.observability_views import build_report_observability
  9. from service.health_analysis import HealthAnalysisService
  10. router = APIRouter()
  11. class HealthRequest(BaseModel):
  12. report_text: str
  13. user_id: str = Field(..., min_length=1, max_length=256)
  14. @field_validator("user_id")
  15. @classmethod
  16. def normalize_user_id(cls, v: str) -> str:
  17. v = v.strip()
  18. if not v:
  19. raise ValueError("user_id 不能为空")
  20. return v
  21. @router.post("/health/analysis")
  22. async def analysis_health(request: HealthRequest):
  23. task_id = str(uuid4())
  24. service = HealthAnalysisService(task_id=task_id, user_id=request.user_id)
  25. asyncio.create_task(service.run(request.report_text, request.user_id))
  26. return {"task_id": task_id, "user_id": request.user_id}
  27. @router.post("/health/analysis/pdf")
  28. async def analysis_health_pdf(
  29. file: UploadFile = File(...),
  30. user_id: str = Form(...),
  31. ):
  32. uid = user_id.strip()
  33. if not uid:
  34. return {"error": "user_id 不能为空"}
  35. contents = await file.read()
  36. text = ""
  37. with pdfplumber.open(BytesIO(contents)) as pdf:
  38. for page in pdf.pages:
  39. page_text = page.extract_text()
  40. if page_text:
  41. text += page_text + "\n"
  42. if not text.strip():
  43. return {"error": "无法从PDF中提取文本"}
  44. task_id = str(uuid4())
  45. service = HealthAnalysisService(task_id=task_id, user_id=uid)
  46. asyncio.create_task(service.run(text, uid))
  47. return {"task_id": task_id, "user_id": uid}
  48. @router.get("/health/task_status/{task_id}")
  49. async def task_status(task_id: str):
  50. from agents.base import get_task_status
  51. status = get_task_status(task_id)
  52. if not status:
  53. return {"error": "task not found"}
  54. return status
  55. @router.get("/health/users/{user_id}/report_history")
  56. async def report_history(user_id: str, limit: int = 50):
  57. uid = user_id.strip()
  58. if not uid:
  59. return {"error": "user_id 无效", "items": []}
  60. items = list_report_runs_for_user(uid, limit=limit)
  61. return {"user_id": uid, "items": items}
  62. @router.get("/health/report_runs/{task_id}")
  63. async def report_run_detail(task_id: str):
  64. row = get_report_run(task_id)
  65. if not row:
  66. return {"error": "未找到该次分析记录(可能尚未落库或 task_id 无效)"}
  67. return row
  68. @router.get("/health/report_runs/{task_id}/observability")
  69. async def report_run_observability(
  70. task_id: str, include_raw_trace: bool = False
  71. ):
  72. """
  73. 阶段 3:体检分析可观测性 — 各 Agent trace 已随 report_runs 持久化(新产生任务)。
  74. `include_raw_trace=true` 时返回完整 trace(体积可能较大)。
  75. """
  76. row = get_report_run(task_id.strip())
  77. if not row:
  78. raise HTTPException(status_code=404, detail="未找到该次分析记录")
  79. return build_report_observability(row, include_raw_trace=include_raw_trace)