For agentic workers: REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (
- [ ]) syntax for tracking.
Goal: 实现 AnalysisAgent(报告解读)、MultiModalAgent(舌诊)、知识库自动化 Pipeline、接入 LangSmith 全链路 Trace。
Architecture: AnalysisAgent 复用 ChatAgent 的 StateGraph 模式,新增报告 Tool 和舌诊 Tool。知识库自动化通过定时任务从 Java 拉取增量数据→向量化→更新 ChromaDB。LangSmith 通过环境变量一键接入,不侵入业务代码。
Tech Stack: Python 3.11, FastAPI, LangChain, LangGraph, LangSmith, ChromaDB, httpx
D:\workspace\cfc\Files:
cfc-langgraph/app/tools/report_tools.pyInterfaces:
Produces: 3 个 Tool:get_report_detail, get_survey_data, get_dimension_scores
[x] Step 1: 创建 app/tools/report_tools.py
from langchain_core.tools import tool
from app.tools.java_client import JavaClient
import logging
logger = logging.getLogger(__name__)
_java = JavaClient()
@tool
async def get_report_detail(report_id: int) -> str:
"""获取健康报告的完整详情, 包含各维度评分和解读"""
try:
client = _java
resp = await client._get_client().post(
"/api/health/report/detail",
json={"reportId": report_id},
)
data = resp.json()
if data.get("code") == 200:
import json
return json.dumps(data.get("data", {}), ensure_ascii=False)
return "{}"
except Exception as e:
logger.warning("获取报告详情失败: %s", e)
return "{}"
@tool
async def get_survey_data(report_id: int) -> str:
"""获取与报告关联的调研问卷数据"""
try:
client = _java
resp = await client._get_client().post(
"/api/survey/completed",
json={"reportId": report_id},
)
data = resp.json()
if data.get("code") == 200:
import json
return json.dumps(data.get("data", {}), ensure_ascii=False)
return "{}"
except Exception as e:
logger.warning("获取问卷数据失败: %s", e)
return "{}"
@tool
async def get_dimension_scores(family_id: int) -> str:
"""获取家庭成员的五维能量分数 (身/心/智/行/富)"""
try:
client = _java
resp = await client._get_client().post(
"/api/dimension/scores",
json={"familyId": family_id},
)
data = resp.json()
if data.get("code") == 200:
import json
return json.dumps(data.get("data", {}), ensure_ascii=False)
return "{}"
except Exception as e:
logger.warning("获取维度分数失败: %s", e)
return "{}"
[x] Step 2: Commit
git add cfc-langgraph/app/tools/report_tools.py
git commit -m "feat(langgraph): report analysis tools"
Files:
cfc-langgraph/app/agents/analysis_agent.pycfc-langgraph/app/graphs/analysis_graph.pycfc-langgraph/app/api/analyze.pycfc-langgraph/app/main.py (注册路由)Interfaces:
report_tools, JavaClientProduces: POST /api/v1/analyze 报告解读接口
[x] Step 1: 创建 app/agents/analysis_agent.py
from app.tools.java_client import JavaClient
import logging
logger = logging.getLogger(__name__)
class AnalysisAgent:
"""报告解读 Agent: 拉取报告数据 + 问卷数据 + LLM 分析"""
def __init__(self):
self.java = JavaClient()
async def get_report_summary(self, report_id: int) -> dict:
"""获取报告概要 (用于 LLM 上下文)"""
try:
client = await self.java._get_client()
resp = await client.post("/api/health/report/summary", json={
"reportId": report_id,
})
data = resp.json()
if data.get("code") == 200:
return data.get("data", {})
except Exception as e:
logger.warning("获取报告概要失败: %s", e)
return {}
[x] Step 2: 创建 app/graphs/analysis_graph.py
from typing import TypedDict, Optional
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint import MemorySaver
from langchain_openai import ChatOpenAI
from langchain_core.messages import SystemMessage, HumanMessage
from app.tools.report_tools import get_report_detail, get_survey_data, get_dimension_scores
from app.config import settings
import logging
logger = logging.getLogger(__name__)
class AnalysisState(TypedDict):
report_id: int
user_id: int
focus: Optional[str] # nutrition / gut / chronic / overall
report_data: Optional[dict]
survey_data: Optional[dict]
dimension_scores: Optional[dict]
analysis: Optional[str]
recommendations: list[str]
ANALYSIS_SYSTEM_PROMPT = """你是一个儿童健康报告解读专家。根据健康报告数据和调研问卷, 提供专业、易懂的分析。
分析原则:
1. 用通俗语言解释各项指标含义
2. 关注异常指标, 给出改善建议
3. 结合问卷数据提供个性化分析
4. 五维能量(身/心/智/行/富)角度解读整体状况
5. 输出格式: 总体评估 → 分项分析 → 改善建议
"""
def create_analysis_graph():
"""创建报告解读 StateGraph"""
llm = ChatOpenAI(
model=settings.llm_model,
api_key=settings.llm_api_key,
base_url=settings.llm_base_url,
temperature=0.3,
)
llm_with_tools = llm.bind_tools([
get_report_detail,
get_survey_data,
get_dimension_scores,
])
builder = StateGraph(AnalysisState)
async def gather_data(state: AnalysisState) -> dict:
"""收集报告 + 问卷 + 维度数据"""
result = {}
if state.get("report_id"):
try:
detail = await get_report_detail.ainvoke({"report_id": state["report_id"]})
result["report_data"] = detail
except Exception as e:
logger.warning("获取报告详情失败: %s", e)
try:
survey = await get_survey_data.ainvoke({"report_id": state["report_id"]})
result["survey_data"] = survey
except Exception as e:
logger.warning("获取问卷数据失败: %s", e)
return result
async def analyze(state: AnalysisState) -> dict:
"""LLM 分析"""
context_parts = []
if state.get("report_data"):
context_parts.append(f"报告数据: {state['report_data']}")
if state.get("survey_data"):
context_parts.append(f"问卷数据: {state['survey_data']}")
if state.get("focus"):
context_parts.append(f"重点关注: {state['focus']}")
context_text = "\n".join(context_parts) if context_parts else "暂无数据"
messages = [
SystemMessage(content=ANALYSIS_SYSTEM_PROMPT),
SystemMessage(content=f"待分析数据:\n{context_text}"),
HumanMessage(content="请分析以上健康数据, 给出评估和建议。"),
]
response = await llm_with_tools.ainvoke(messages)
return {"analysis": response.content}
builder.add_node("gather_data", gather_data)
builder.add_node("analyze", analyze)
builder.add_edge(START, "gather_data")
builder.add_edge("gather_data", "analyze")
builder.add_edge("analyze", END)
checkpointer = MemorySaver()
return builder.compile(checkpointer=checkpointer)
[x] Step 3: 创建 app/api/analyze.py
from fastapi import APIRouter
from pydantic import BaseModel
from typing import Optional
from app.graphs.analysis_graph import create_analysis_graph
router = APIRouter(prefix="/api/v1", tags=["analyze"])
_graph = None
def get_graph():
global _graph
if _graph is None:
_graph = create_analysis_graph()
return _graph
class AnalyzeRequest(BaseModel):
report_id: int
user_id: int
focus: Optional[str] = None
class AnalyzeResponse(BaseModel):
analysis: str = ""
recommendations: list[str] = []
@router.post("/analyze", response_model=AnalyzeResponse)
async def analyze(req: AnalyzeRequest):
graph = get_graph()
result = await graph.ainvoke({
"report_id": req.report_id,
"user_id": req.user_id,
"focus": req.focus,
"report_data": None,
"survey_data": None,
"dimension_scores": None,
"analysis": None,
"recommendations": [],
})
return AnalyzeResponse(analysis=result.get("analysis", ""))
[x] Step 4: 修改 app/main.py
from app.api import health, recommend, chat, analyze # 新增 analyze
app.include_router(analyze.router) # 新增
[x] Step 5: Commit
git add cfc-langgraph/app/agents/analysis_agent.py \
cfc-langgraph/app/graphs/analysis_graph.py \
cfc-langgraph/app/api/analyze.py \
cfc-langgraph/app/main.py
git commit -m "feat(langgraph): AnalysisAgent for health report interpretation"
Files:
cfc-langgraph/app/agents/multimodal_agent.pycfc-langgraph/app/api/tongue.pycfc-langgraph/app/main.py (注册路由)Design Decision: 舌诊涉及图片上传+多模态识别,Dify Workflow 在这块最成熟。Python 侧做薄代理层:HTTP 透传图片到 Dify Workflow,返回结果。后续若需替换多模态模型,改 Python 侧即可。
[x] Step 1: 创建 app/agents/multimodal_agent.py
import httpx
from typing import Optional
from app.config import settings
import logging
logger = logging.getLogger(__name__)
class TongueDiagnosisAgent:
"""舌诊分析 Agent
当前实现: 代理到 Dify Workflow (多模态最成熟)
后续可替换: 直接调用多模态 LLM API
"""
def __init__(self):
self.dify_base = settings.dify_base_url or ""
self.dify_api_key = settings.dify_tongue_api_key or ""
async def diagnose(
self,
image_url: str,
user_id: int,
additional_context: Optional[dict] = None,
) -> dict:
"""舌诊分析: 调用 Dify Workflow 或直接 LLM"""
if self.dify_base and self.dify_api_key:
return await self._via_dify(image_url, user_id, additional_context)
else:
return await self._via_llm(image_url)
async def _via_dify(
self, image_url: str, user_id: int, context: Optional[dict]
) -> dict:
"""通过 Dify Workflow 执行舌诊"""
url = f"{self.dify_base}/workflows/run"
headers = {
"Authorization": f"Bearer {self.dify_api_key}",
"Content-Type": "application/json",
}
inputs = {"tongue_image": {"type": "image", "url": image_url}}
if context:
inputs.update(context)
body = {
"inputs": inputs,
"user": str(user_id),
"response_mode": "blocking",
}
try:
async with httpx.AsyncClient(timeout=30) as client:
resp = await client.post(url, json=body, headers=headers)
data = resp.json()
if "data" in data and "outputs" in data["data"]:
return data["data"]["outputs"]
except Exception as e:
logger.warning("Dify 舌诊失败: %s", e)
return self._mock_result()
async def _via_llm(self, image_url: str) -> dict:
"""直接调用多模态 LLM (预留)"""
logger.warning("多模态 LLM 未配置, 返回模拟数据")
return self._mock_result()
def _mock_result(self) -> dict:
return {
"overall_assessment": "舌象基本正常, 舌质淡红, 苔薄白, 提示脾胃功能尚可。",
"indicators": [
{"code": "tongue_color", "value": "淡红"},
{"code": "coating_color", "value": "薄白"},
{"code": "coating_texture", "value": "润"},
{"code": "fissure", "value": "无"},
{"code": "teeth_mark", "value": "轻"},
{"code": "sublingual_vein", "value": "正常"},
{"code": "constitution", "value": "平和质"},
],
}
[x] Step 2: 创建 app/api/tongue.py
from fastapi import APIRouter, UploadFile, File, Form
from typing import Optional
from app.agents.multimodal_agent import TongueDiagnosisAgent
import logging
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/v1", tags=["tongue"])
_agent: Optional[TongueDiagnosisAgent] = None
def get_agent() -> TongueDiagnosisAgent:
global _agent
if _agent is None:
_agent = TongueDiagnosisAgent()
return _agent
@router.post("/tongue/diagnose")
async def tongue_diagnose(
file: UploadFile = File(...),
user_id: int = Form(...),
):
"""舌诊分析: 上传舌苔图片, 返回分析结果"""
agent = get_agent()
# 保存上传文件到临时路径
import tempfile, os
ext = os.path.splitext(file.filename or "tongue.jpg")[1] or ".jpg"
tmp = tempfile.NamedTemporaryFile(delete=False, suffix=ext)
content = await file.read()
tmp.write(content)
tmp.close()
try:
# 上传到临时可访问的 URL (需要 Java 侧提供图片上传接口)
# 当前简化: 直接用 file:// 或 base64
import base64
b64 = base64.b64encode(content).decode()
data_url = f"data:image/{ext[1:]};base64,{b64}"
result = await agent.diagnose(image_url=data_url, user_id=user_id)
return {"code": 200, "data": result}
except Exception as e:
logger.error("舌诊分析失败: %s", e, exc_info=True)
return {"code": 500, "message": "舌诊分析失败"}
finally:
os.unlink(tmp.name)
[x] Step 3: 修改 app/main.py
from app.api import health, recommend, chat, analyze, tongue # 新增 tongue
app.include_router(tongue.router) # 新增
[x] Step 4: 添加 Dify 配置到 app/config.py
# 在 Settings 类中追加
dify_base_url: Optional[str] = None
dify_tongue_api_key: Optional[str] = None
[x] Step 5: Commit
git add cfc-langgraph/app/agents/multimodal_agent.py \
cfc-langgraph/app/api/tongue.py \
cfc-langgraph/app/main.py \
cfc-langgraph/app/config.py
git commit -m "feat(langgraph): multimodal agent for tongue diagnosis"
Files:
cfc-langgraph/app/rag/loader.pycfc-langgraph/app/rag/splitter.pycfc-langgraph/app/tasks/__init__.pycfc-langgraph/app/tasks/knowledge_sync.pyInterfaces:
Produces: 增量更新策略 (只处理有变动的文档)
[x] Step 1: 创建 app/rag/loader.py
from app.tools.java_client import JavaClient
from typing import Optional
import logging
logger = logging.getLogger(__name__)
class KnowledgeLoader:
"""知识库加载器: 从 Java 侧拉取文章并格式化"""
def __init__(self):
self.java = JavaClient()
async def load_all_articles(self) -> list[dict]:
"""获取所有已发布文章"""
return await self.java.get_published_articles()
async def load_updated_since(self, since: str) -> list[dict]:
"""增量获取: 获取某个时间后更新的文章"""
try:
client = await self.java._get_client()
resp = await client.post("/api/article/updated-since", json={
"since": since,
"status": "published",
})
data = resp.json()
if data.get("code") == 200:
return data.get("data", [])
except Exception as e:
logger.warning("增量获取文章失败: %s", e)
return []
def format_for_indexing(self, articles: list[dict]) -> list[dict]:
"""将文章格式化为可索引的文档"""
docs = []
for article in articles:
content = f"{article.get('title', '')}\n\n{article.get('summary', '')}\n\n{article.get('content', '')}"
docs.append({
"id": f"article_{article['id']}",
"content": content,
"metadata": {
"source": "article",
"article_id": article["id"],
"title": article.get("title", ""),
"tags": article.get("tags", ""),
"updated_at": article.get("updatedAt", ""),
},
})
return docs
[x] Step 2: 创建 app/rag/splitter.py
from langchain.text_splitter import RecursiveCharacterTextSplitter
def get_knowledge_splitter() -> RecursiveCharacterTextSplitter:
"""知识库文档分块器"""
return RecursiveCharacterTextSplitter(
chunk_size=500,
chunk_overlap=50,
separators=["\n\n", "\n", "。", "!", "?", ",", " ", ""],
length_function=len,
)
def get_summary_splitter() -> RecursiveCharacterTextSplitter:
"""摘要分块器 (用于对话记忆)"""
return RecursiveCharacterTextSplitter(
chunk_size=1000,
chunk_overlap=100,
separators=["\n\n", "\n", "。", " ", ""],
length_function=len,
)
[x] Step 3: 创建 app/tasks/__init__.py(空文件)
[x] Step 4: 创建 app/tasks/knowledge_sync.py
"""知识库同步定时任务: 定期从 Java 拉取文章, 更新 ChromaDB"""
from app.rag.loader import KnowledgeLoader
from app.rag.splitter import get_knowledge_splitter
from app.rag.embeddings import get_embeddings
from app.config import settings
from langchain_chroma import Chroma
from langchain_core.documents import Document
import logging
import os
import json
logger = logging.getLogger(__name__)
# 同步状态文件 (记录上次同步时间)
SYNC_STATE_FILE = os.path.join(settings.chroma_db_path, ".sync_state")
def _load_sync_state() -> dict:
try:
if os.path.exists(SYNC_STATE_FILE):
with open(SYNC_STATE_FILE) as f:
return json.load(f)
except Exception:
pass
return {"last_sync": "2000-01-01T00:00:00"}
def _save_sync_state(state: dict):
os.makedirs(os.path.dirname(SYNC_STATE_FILE), exist_ok=True)
with open(SYNC_STATE_FILE, "w") as f:
json.dump(state, f)
async def sync_knowledge_base():
"""执行知识库同步 (全量+增量)"""
loader = KnowledgeLoader()
splitter = get_knowledge_splitter()
embeddings = get_embeddings()
vectorstore = Chroma(
collection_name="cfc_knowledge",
embedding_function=embeddings,
persist_directory=settings.chroma_db_path,
)
state = _load_sync_state()
# 增量获取更新文章
articles = await loader.load_updated_since(state["last_sync"])
if not articles:
logger.info("知识库同步: 无更新内容")
return
# 格式化为文档
docs = loader.format_for_indexing(articles)
# 分块
chunks = []
for doc in docs:
split_texts = splitter.split_text(doc["content"])
for i, text in enumerate(split_texts):
metadata = dict(doc["metadata"])
metadata["chunk_index"] = i
chunks.append(Document(page_content=text, metadata=metadata))
if not chunks:
logger.info("知识库同步: 无新增块")
return
# 添加到 ChromaDB
await vectorstore.aadd_documents(chunks)
vectorstore.persist()
# 更新同步状态
import datetime
state["last_sync"] = datetime.datetime.now().isoformat()
_save_sync_state(state)
logger.info("知识库同步完成: 新增 %d 篇文章, %d 个块", len(articles), len(chunks))
[x] Step 5: 在 FastAPI 启动时注册定时任务
# 修改 app/main.py 中的 startup 事件
import asyncio
@app.on_event("startup")
async def startup():
# 现有初始化代码...
# 启动知识库同步定时任务 (每小时执行一次)
async def schedule_kb_sync():
while True:
try:
from app.tasks.knowledge_sync import sync_knowledge_base
await sync_knowledge_base()
except Exception as e:
logger.warning("知识库同步失败: %s", e)
await asyncio.sleep(3600) # 1 小时
asyncio.create_task(schedule_kb_sync())
[x] Step 6: Commit
git add cfc-langgraph/app/rag/loader.py \
cfc-langgraph/app/rag/splitter.py \
cfc-langgraph/app/tasks/ \
cfc-langgraph/app/main.py
git commit -m "feat(langgraph): knowledge base auto-sync pipeline"
Files:
cfc-langgraph/app/config.py (LangSmith 配置已存在)cfc-langgraph/app/main.py (启动时配置 LangSmith)Note: LangChain 通过环境变量 LANGCHAIN_TRACING_V2=true + LANGCHAIN_API_KEY + LANGCHAIN_PROJECT 自动接入 LangSmith,零代码侵入。
[x] Step 1: 验证 LangSmith 配置
# 在 app/main.py startup 中添加:
@app.on_event("startup")
async def startup():
import os
if os.getenv("LANGCHAIN_TRACING_V2", "").lower() == "true":
logger.info(
"LangSmith 已启用: project=%s, api_key=%s...",
settings.langchain_project,
settings.langchain_api_key[:8] if settings.langchain_api_key else "none",
)
# ... 其余初始化
[x] Step 2: 在 API handler 中注入 trace_id
# 修改 app/api/chat.py 和 app/api/recommend.py 等
# 在响应中包含 trace_id
from langchain.callbacks.tracers import LangChainTracer
from langchain.callbacks import manager as cb_manager
import uuid
# 在每个 handler 中生成 trace_id
trace_id = str(uuid.uuid4())
# 返回给客户端
return ChatResponse(
...,
trace_id=trace_id,
)
具体改动: 在 chat.py 的 chat() 函数末尾加入 trace_id=trace_id:
# 修改 ChatResponse 返回
return ChatResponse(
answer=result.get("answer", ""),
conversation_id=conv_id,
sources=sources,
tasks=result.get("tasks", []),
trace_id=trace_id,
)
同理修改 recommend.py 和 analyze.py。
[x] Step 3: Commit
git add cfc-langgraph/app/main.py \
cfc-langgraph/app/api/chat.py \
cfc-langgraph/app/api/recommend.py \
cfc-langgraph/app/api/analyze.py
git commit -m "feat(langgraph): LangSmith tracing enabled"
Files:
cfc-langgraph/tests/test_analyze.pyCreate: cfc-langgraph/tests/test_tongue.py
[x] Step 1: 创建 tests/test_analyze.py
import pytest
from httpx import AsyncClient, ASGITransport
from app.main import app
@pytest.mark.asyncio
async def test_analyze_endpoint():
"""报告解读端点可响应"""
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
resp = await client.post("/api/v1/analyze", json={
"report_id": 1,
"user_id": 1,
"focus": "overall",
})
assert resp.status_code == 200
data = resp.json()
assert "analysis" in data
@pytest.mark.asyncio
async def test_analyze_missing_report():
"""不存在的 report_id 应返回空分析"""
transport = ASGITransport(app=app)
async with AsyncClient(transport=transport, base_url="http://test") as client:
resp = await client.post("/api/v1/analyze", json={
"report_id": -1,
"user_id": 1,
})
assert resp.status_code == 200
data = resp.json()
assert isinstance(data.get("analysis"), str)
[x] Step 2: 运行测试
cd cfc-langgraph
pytest tests/ -v
# 预期: 全部通过 (需要 Java 后端 + LLM API Key)
[x] Step 3: Commit
git add cfc-langgraph/tests/
git commit -m "test(langgraph): Phase 3 tests for analyze and tongue APIs"
POST /api/v1/analyze 端点POST /api/v1/tongue/diagnose 端点