# LangGraph Sidecar — Phase 1 实施计划 > **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:** 在 `cfc-langgraph/` 建立 Python LangGraph sidecar 服务,上线 RecommendAgent,Java 端支持灰度切换(5% 流量走 Python,异常自动 fallback)。 **Architecture:** FastAPI 服务 (port 9000) 承载 LangChain Agent,通过 HTTP 与 Java (port 9082) 通信。数据存储用 ChromaDB 文件模式 + Java 侧 MySQL。Agent 调 LLM API(DeepSeek-V3)进行推荐+解释。 **Tech Stack:** Python 3.11, FastAPI, LangChain, LangGraph, ChromaDB, Docker, Java 8 (不变) ## Global Constraints - Python 3.11+ (不要求最新 3.12/3.13) - ChromaDB 文件模式(不引入额外容器或中间件) - Java 端 JDK 8 / Spring Boot 2.7.18 保持不变 - 所有 Tool 调用通过 HTTP 访问 Java 接口,而非直连 MySQL - Dify 作为 Fallback:Python 超时或异常时 Java 自动回退 Dify - 所有项目文件路径相对于 `D:\workspace\cfc\` --- ### Task 1: Python 项目脚手架 **Files:** - Create: `cfc-langgraph/pyproject.toml` - Create: `cfc-langgraph/.env.example` - Create: `cfc-langgraph/.gitignore` - Create: `cfc-langgraph/Dockerfile` - Create: `cfc-langgraph/app/__init__.py` - Create: `cfc-langgraph/app/main.py` **Interfaces:** - Produces: FastAPI app 入口 `app.main:app`,`/health` 返回 200 - [x] **Step 1: 创建 `pyproject.toml`** ```toml [project] name = "cfc-langgraph" version = "0.1.0" description = "cfc AI sidecar service with LangGraph" requires-python = ">=3.11" dependencies = [ "fastapi>=0.115", "uvicorn[standard]>=0.34", "pydantic>=2.10", "pydantic-settings>=2.7", "httpx>=0.28", "langchain>=0.3", "langchain-community>=0.3", "langchain-openai>=0.3", "langgraph>=0.3", "chromadb>=0.6", "langchain-chroma>=0.2", "python-multipart>=0.0.20", ] [project.optional-dependencies] dev = [ "pytest>=8", "pytest-asyncio>=0.25", "ruff>=0.8", "httpx", ] [tool.ruff] target-version = "py311" line-length = 100 select = ["E", "F", "I", "N", "W"] ignore = ["E501"] ``` - [x] **Step 2: 创建 `.env.example`** ```bash # LLM LLM_API_KEY=sk-your-key LLM_BASE_URL=https://api.deepseek.com/v1 LLM_MODEL=deepseek-chat # Embedding EMBEDDING_API_KEY=${LLM_API_KEY} EMBEDDING_BASE_URL=https://api.deepseek.com/v1 EMBEDDING_MODEL=text-embedding-v3 # Java Backend JAVA_BASE_URL=http://localhost:9082 JAVA_CONTEXT_URL=${JAVA_BASE_URL}/api/ai/context # LangSmith (optional) LANGCHAIN_TRACING_V2=false LANGCHAIN_API_KEY= LANGCHAIN_PROJECT=cfc-langgraph # Service SERVICE_HOST=0.0.0.0 SERVICE_PORT=9000 LOG_LEVEL=info # Chroma CHROMA_DB_PATH=./data/chroma_db ``` - [x] **Step 3: 创建 `.gitignore`** ``` __pycache__/ *.py[cod] .env data/ *.egg-info/ dist/ .ruff_cache/ .pytest_cache/ ``` - [x] **Step 4: 创建 `Dockerfile`** ```dockerfile FROM python:3.11-slim WORKDIR /app COPY pyproject.toml . RUN pip install --no-cache-dir -e . COPY app/ app/ COPY data/ data/ 2>/dev/null || true EXPOSE 9000 CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "9000"] ``` - [x] **Step 5: 创建 `app/__init__.py`**(空文件) - [x] **Step 6: 创建 `app/main.py`** ```python from fastapi import FastAPI from app.api import health app = FastAPI(title="cfc-langgraph", version="0.1.0") app.include_router(health.router) @app.on_event("startup") async def startup(): pass # 后续 Phase 在此初始化 RAG / Agent @app.on_event("shutdown") async def shutdown(): pass # 后续 Phase 在此清理资源 ``` - [x] **Step 7: 验证服务可启动** ```bash cd cfc-langgraph pip install -e . uvicorn app.main:app --port 9000 & curl http://localhost:9000/health # 预期: {"status":"ok"} kill %1 ``` - [x] **Step 8: 创建 `app/api/__init__.py`**(空文件) - [x] **Step 9: 创建 `app/api/health.py`** ```python from fastapi import APIRouter router = APIRouter(tags=["health"]) @router.get("/health") async def health(): return {"status": "ok"} ``` - [x] **Step 10: Commit** ```bash git add cfc-langgraph/ git commit -m "feat(langgraph): scaffold Python sidecar project" ``` --- ### Task 2: 配置管理 **Files:** - Create: `cfc-langgraph/app/config.py` **Interfaces:** - Produces: `app.config.settings` — Pydantic Settings 单例,所有模块通过 `from app.config import settings` 读取 - [x] **Step 1: 创建 `app/config.py`** ```python from pydantic_settings import BaseSettings from typing import Optional class Settings(BaseSettings): # LLM llm_api_key: str llm_base_url: str = "https://api.deepseek.com/v1" llm_model: str = "deepseek-chat" # Embedding embedding_api_key: Optional[str] = None embedding_base_url: Optional[str] = None embedding_model: str = "text-embedding-v3" # Java Backend java_base_url: str = "http://localhost:9082" java_context_url: Optional[str] = None # LangSmith langchain_tracing_v2: bool = False langchain_api_key: Optional[str] = None langchain_project: str = "cfc-langgraph" # Service service_host: str = "0.0.0.0" service_port: int = 9000 log_level: str = "info" # Chroma chroma_db_path: str = "./data/chroma_db" model_config = {"env_file": ".env", "env_file_encoding": "utf-8"} @property def effective_embedding_api_key(self) -> str: return self.embedding_api_key or self.llm_api_key @property def effective_embedding_base_url(self) -> str: return self.embedding_base_url or self.llm_base_url @property def effective_java_context_url(self) -> str: return self.java_context_url or f"{self.java_base_url}/api/ai/context" settings = Settings() ``` - [x] **Step 2: 创建 `.env` 文件(从 `.env.example` 复制,填入真实 Key)** ```bash cp cfc-langgraph/.env.example cfc-langgraph/.env # 手动编辑填入 LLM_API_KEY ``` - [x] **Step 3: 验证配置可加载** ```bash cd cfc-langgraph python -c "from app.config import settings; print(settings.llm_model)" # 预期输出: deepseek-chat ``` - [x] **Step 4: Commit** ```bash git add cfc-langgraph/ git commit -m "feat(langgraph): config management with pydantic-settings" ``` --- ### Task 3: Pydantic 数据模型 **Files:** - Create: `cfc-langgraph/app/models/__init__.py` - Create: `cfc-langgraph/app/models/common.py` - Create: `cfc-langgraph/app/models/chat.py` - Create: `cfc-langgraph/app/models/recommend.py` **Interfaces:** - Produces: Phase 1 所需的请求/响应模型 - [x] **Step 1: 创建 `app/models/__init__.py`** ```python from .common import UserContext, SourceInfo from .chat import ChatRequest, ChatResponse from .recommend import RecommendRequest, RecommendResponse, RecommendItem ``` - [x] **Step 2: 创建 `app/models/common.py`** ```python from pydantic import BaseModel from typing import Optional class UserContext(BaseModel): """Java 侧传来的业务上下文""" child_id: Optional[int] = None report_id: Optional[int] = None survey_id: Optional[int] = None family_id: Optional[int] = None mascot_code: Optional[str] = None class SourceInfo(BaseModel): """回答引用来源""" type: str # knowledge / tool title: str score: Optional[float] = None ``` - [x] **Step 3: 创建 `app/models/chat.py`** ```python from pydantic import BaseModel from typing import Optional from .common import UserContext, SourceInfo class ChatRequest(BaseModel): query: str user_id: int conversation_id: str = "" context: Optional[UserContext] = None class ChatResponse(BaseModel): answer: str conversation_id: str sources: list[SourceInfo] = [] tasks: list[dict] = [] trace_id: str = "" ``` - [x] **Step 4: 创建 `app/models/recommend.py`** ```python from pydantic import BaseModel from typing import Optional class RecommendItem(BaseModel): type: str # product / activity / article id: int name: str description: str = "" cover_image: str = "" price: Optional[float] = None url: str = "" score: float = 0.0 reason: str = "" class RecommendRequest(BaseModel): user_id: int query: str = "" tags: list[str] = [] types: Optional[list[str]] = None limit: int = 5 context: Optional[dict] = None class RecommendResponse(BaseModel): items: list[RecommendItem] source: str = "" # sql / vector / hybrid trace_id: str = "" ``` - [x] **Step 5: Commit** ```bash git add cfc-langgraph/ git commit -m "feat(langgraph): pydantic models for API" ``` --- ### Task 4: ChromaDB + Embedding 基础设施 **Files:** - Create: `cfc-langgraph/app/rag/__init__.py` - Create: `cfc-langgraph/app/rag/embeddings.py` - Create: `cfc-langgraph/app/rag/retriever.py` **Interfaces:** - Produces: `embeddings.get_embeddings()` → Embeddings 实例 - Produces: `retriever.RagRetriever` 类,提供 `retrieve(query, filters, k)` 方法 - [x] **Step 1: 创建 `app/rag/__init__.py`** ```python from .embeddings import get_embeddings from .retriever import RagRetriever ``` - [x] **Step 2: 创建 `app/rag/embeddings.py`** ```python from langchain_openai import OpenAIEmbeddings from app.config import settings _embeddings = None def get_embeddings(): global _embeddings if _embeddings is None: _embeddings = OpenAIEmbeddings( model=settings.embedding_model, api_key=settings.effective_embedding_api_key, base_url=settings.effective_embedding_base_url, ) return _embeddings ``` - [x] **Step 3: 创建 `app/rag/retriever.py`** ```python from langchain_chroma import Chroma from langchain.retrievers import EnsembleRetriever from langchain_community.retrievers import BM25Retriever from langchain.retrievers.document_compressors import LLMChainExtractor from langchain.retrievers import ContextualCompressionRetriever from .embeddings import get_embeddings from app.config import settings from app.tools.java_client import JavaClient from typing import Optional import logging logger = logging.getLogger(__name__) class RagRetriever: """混合检索器: ChromaDB 向量 + BM25 关键词 + 可选 LLM 压缩""" def __init__(self, collection_name: str = "cfc_knowledge"): embeddings = get_embeddings() self.vectorstore = Chroma( collection_name=collection_name, embedding_function=embeddings, persist_directory=settings.chroma_db_path, ) self.java_client = JavaClient() self._bm25_retriever = None self._bm25_texts = [] async def initialize(self): """从 Java 侧拉取知识库数据, 构建 BM25 索引""" try: articles = await self.java_client.get_published_articles() self._bm25_texts = [ f"{a['title']} {a['summary']} {a.get('tags', '')}" for a in articles ] if self._bm25_texts: self._bm25_retriever = BM25Retriever.from_texts( self._bm25_texts, metadatas=articles ) logger.info("BM25 索引就绪: %d 条", len(self._bm25_texts)) except Exception as e: logger.warning("BM25 初始化失败(不影响向量检索): %s", e) async def retrieve( self, query: str, filters: Optional[dict] = None, k: int = 5, use_compression: bool = False, ) -> list[dict]: """混合检索, 返回 [{content, metadata, score}]""" retrievers = [] # 向量检索 vector_retriever = self.vectorstore.as_retriever( search_kwargs={"k": k, "filter": filters} ) retrievers.append(vector_retriever) # BM25 检索 if self._bm25_retriever: retrievers.append(self._bm25_retriever) if len(retrievers) == 1: docs = await retrievers[0].ainvoke(query) else: ensemble = EnsembleRetriever( retrievers=retrievers, weights=[0.6, 0.4] ) docs = await ensemble.ainvoke(query) # 可选: LLM 压缩去噪 if use_compression and docs: from langchain_openai import ChatOpenAI llm = ChatOpenAI( model=settings.llm_model, api_key=settings.llm_api_key, base_url=settings.llm_base_url, ) compressor = LLMChainExtractor.from_llm(llm) compression_retriever = ContextualCompressionRetriever( base_compressor=compressor, base_retriever=self.vectorstore.as_retriever(), ) docs = await compression_retriever.ainvoke(query) results = [] for doc in docs: results.append({ "content": doc.page_content, "metadata": doc.metadata, "score": getattr(doc, "metadata", {}).get("score", 0), }) return results[:k] ``` - [x] **Step 4: Commit** ```bash git add cfc-langgraph/ git commit -m "feat(langgraph): ChromaDB + RAG retriever infrastructure" ``` --- ### Task 5: Java HTTP 客户端(Python 侧 Tool 用) **Files:** - Create: `cfc-langgraph/app/tools/__init__.py` - Create: `cfc-langgraph/app/tools/java_client.py` **Interfaces:** - Produces: `JavaClient` 类, 封装对所有 Java 接口的 HTTP 调用 - Produces: `get_published_articles()`, `search_products()`, `search_activities()`, `search_articles()` - [x] **Step 1: 创建 `app/tools/__init__.py`** ```python from .java_client import JavaClient ``` - [x] **Step 2: 创建 `app/tools/java_client.py`** ```python import httpx from typing import Optional from app.config import settings import logging logger = logging.getLogger(__name__) class JavaClient: """Java 后端 HTTP 客户端 (所有 Python→Java 通信的单一入口)""" def __init__(self): self.base_url = settings.java_base_url self._client: Optional[httpx.AsyncClient] = None async def _get_client(self) -> httpx.AsyncClient: if self._client is None: self._client = httpx.AsyncClient( base_url=self.base_url, timeout=httpx.Timeout(10.0, connect=3.0), ) return self._client async def close(self): if self._client: await self._client.aclose() self._client = None async def get_published_articles(self) -> list[dict]: """获取已发布的文章列表 (用于构建知识库)""" client = await self._get_client() resp = await client.post("/api/article/list", json={"status": "published", "limit": 1000}) data = resp.json() if data.get("code") == 200: return data.get("data", []) return [] async def search_products(self, keyword: str, limit: int = 5) -> list[dict]: """按关键词搜索上架商品""" client = await self._get_client() resp = await client.post("/api/product/search", json={ "keyword": keyword, "status": "上架", "limit": limit }) data = resp.json() if data.get("code") == 200: return data.get("data", []) return [] async def search_activities(self, keyword: str, limit: int = 5) -> list[dict]: """按关键词搜索进行中的活动""" client = await self._get_client() resp = await client.post("/api/activity/search", json={ "keyword": keyword, "status": "published", "limit": limit }) data = resp.json() if data.get("code") == 200: return data.get("data", []) return [] async def search_articles(self, keyword: str, limit: int = 5) -> list[dict]: """按关键词搜索已发布文章""" client = await self._get_client() resp = await client.post("/api/article/search", json={ "keyword": keyword, "status": "published", "limit": limit }) data = resp.json() if data.get("code") == 200: return data.get("data", []) return [] async def get_user_context(self, user_id: int, params: Optional[dict] = None) -> dict: """获取用户上下文 (对应 Java AiContextService)""" client = await self._get_client() resp = await client.post(settings.effective_java_context_url, json={ "user_id": str(user_id), "params": params or {}, }) data = resp.json() if data.get("code") == 200: return data.get("data", {}) return {} ``` - [x] **Step 3: Commit** ```bash git add cfc-langgraph/ git commit -m "feat(langgraph): Java HTTP client for tool calls" ``` --- ### Task 6: RecommendAgent **Files:** - Create: `cfc-langgraph/app/tools/product_tools.py` - Create: `cfc-langgraph/app/agents/__init__.py` - Create: `cfc-langgraph/app/agents/recommend_agent.py` - Create: `cfc-langgraph/app/graphs/__init__.py` - Create: `cfc-langgraph/app/graphs/recommend_graph.py` - Create: `cfc-langgraph/app/api/recommend.py` - Modify: `cfc-langgraph/app/main.py` (注册路由) **Interfaces:** - Consumes: `JavaClient`, `RagRetriever`, `ChatOpenAI` - Produces: `POST /api/v1/recommend` 接口 - [x] **Step 1: 创建 `app/tools/product_tools.py`** ```python from langchain_core.tools import tool from app.tools.java_client import JavaClient import logging logger = logging.getLogger(__name__) _java = JavaClient() @tool async def search_product_by_keyword(keyword: str, limit: int = 5) -> str: """按关键词搜索上架商品, 返回 JSON 商品列表""" try: products = await _java.search_products(keyword, limit) if not products: return "[]" return str([{ "id": p["id"], "name": p["name"], "price": p.get("price"), "description": p.get("intro") or p.get("description", ""), } for p in products]) except Exception as e: logger.warning("搜索商品失败: %s", e) return "[]" @tool async def search_activity_by_keyword(keyword: str, limit: int = 5) -> str: """按关键词搜索进行中的活动, 返回 JSON 活动列表""" try: activities = await _java.search_activities(keyword, limit) if not activities: return "[]" return str([{ "id": a["id"], "name": a["title"], "description": a.get("description", ""), } for a in activities]) except Exception as e: logger.warning("搜索活动失败: %s", e) return "[]" @tool async def search_article_by_keyword(keyword: str, limit: int = 5) -> str: """按关键词搜索已发布文章, 返回 JSON 文章列表""" try: articles = await _java.search_articles(keyword, limit) if not articles: return "[]" return str([{ "id": a["id"], "name": a["title"], "summary": a.get("summary", ""), } for a in articles]) except Exception as e: logger.warning("搜索文章失败: %s", e) return "[]" ``` - [x] **Step 2: 创建 `app/agents/__init__.py`**(空) - [x] **Step 3: 创建 `app/graphs/__init__.py`**(空) - [x] **Step 4: 创建 `app/graphs/recommend_graph.py`** Phase 1 用**最简单的单步 Agent**(不涉及 StateGraph 复杂特性),后续 Phase 再升级为图编排: ```python from langchain_openai import ChatOpenAI from langchain_core.messages import SystemMessage, HumanMessage from app.tools.product_tools import search_product_by_keyword, search_activity_by_keyword, search_article_by_keyword from app.config import settings import json import logging logger = logging.getLogger(__name__) SYSTEM_PROMPT = """你是一个儿童成长营养推荐助手。根据用户的需求和营养标签, 推荐合适的商品、活动或文章。 推荐原则: 1. 首先尝试使用搜索工具查找匹配的内容 2. 如果搜索结果为空, 基于你的知识给出建议 3. 每项推荐必须附带推荐理由 4. 以 JSON 格式输出推荐结果 输出格式: { "items": [ { "source": "tool" 或 "knowledge", "type": "product" / "activity" / "article", "id": 数字, "name": "名称", "description": "描述", "reason": "为什么推荐这个" } ] } """ class RecommendAgent: def __init__(self): self.llm = ChatOpenAI( model=settings.llm_model, api_key=settings.llm_api_key, base_url=settings.llm_base_url, temperature=0.3, ) self.tools = [ search_product_by_keyword, search_activity_by_keyword, search_article_by_keyword, ] self.llm_with_tools = self.llm.bind_tools(self.tools) async def run(self, query: str, tags: list[str], limit: int = 5) -> dict: """执行推荐 Agent, 返回推荐结果""" # 如果传入了 tags, 构造搜索关键词 search_query = query or " ".join(tags) messages = [ SystemMessage(content=SYSTEM_PROMPT), HumanMessage(content=f"用户需求: {search_query}\n最大返回数量: {limit}\n请搜索并推荐合适的内容。"), ] # LangChain Tool calling 自动完成: LLM 决定调哪个 Tool → 工具返回结果 → LLM 组织回答 response = await self.llm_with_tools.ainvoke(messages) # 尝试解析 JSON 输出 content = response.content try: # 提取 JSON 块 if "```json" in content: json_str = content.split("```json")[1].split("```")[0].strip() elif "```" in content: json_str = content.split("```")[1].split("```")[0].strip() else: json_str = content.strip() result = json.loads(json_str) return result except (json.JSONDecodeError, IndexError): # 非 JSON 输出, 包装为文本回答 logger.warning("Agent 输出非 JSON, raw: %s", content[:200]) return {"items": [], "text": content} ``` - [x] **Step 5: 创建 `app/api/recommend.py`** ```python from fastapi import APIRouter from app.models.recommend import RecommendRequest, RecommendResponse, RecommendItem from app.graphs.recommend_graph import RecommendAgent from app.rag.retriever import RagRetriever import logging logger = logging.getLogger(__name__) router = APIRouter(prefix="/api/v1", tags=["recommend"]) _agent: RecommendAgent = None _retriever: RagRetriever = None def get_agent() -> RecommendAgent: global _agent if _agent is None: _agent = RecommendAgent() return _agent def get_retriever() -> RagRetriever: global _retriever if _retriever is None: _retriever = RagRetriever() return _retriever @router.post("/recommend", response_model=RecommendResponse) async def recommend(req: RecommendRequest): """营养推荐: Agent 搜索+LLM 解释""" try: agent = get_agent() result = await agent.run(query=req.query, tags=req.tags, limit=req.limit) items_data = result.get("items", []) items = [] for item in items_data: items.append(RecommendItem( type=item.get("type", "product"), id=item.get("id", 0), name=item.get("name", ""), description=item.get("description", ""), reason=item.get("reason", ""), )) return RecommendResponse(items=items, source="agent") except Exception as e: logger.error("RecommendAgent 调用失败: %s", e, exc_info=True) return RecommendResponse(items=[], source="error") ``` - [x] **Step 6: 修改 `app/main.py` 注册路由** ```python from fastapi import FastAPI from app.api import health, recommend app = FastAPI(title="cfc-langgraph", version="0.1.0") app.include_router(health.router) app.include_router(recommend.router) @app.on_event("startup") async def startup(): from app.rag.retriever import RagRetriever retriever = RagRetriever() await retriever.initialize() @app.on_event("shutdown") async def shutdown(): from app.tools.java_client import JavaClient client = JavaClient() await client.close() ``` - [x] **Step 7: 验证 Recommend API** ```bash cd cfc-langgraph # 确保 Java 后端在运行 (localhost:9082) uvicorn app.main:app --port 9000 & curl -X POST http://localhost:9000/api/v1/recommend \ -H "Content-Type: application/json" \ -d '{"user_id": 1, "tags": ["益生菌", "肠胃"], "limit": 3}' # 预期: {"items": [...], "source": "agent"} kill %1 ``` - [x] **Step 8: Commit** ```bash git add cfc-langgraph/ git commit -m "feat(langgraph): RecommendAgent with tool calling" ``` --- ### Task 7: Java 侧 AiGateway(HTTP → Python + Dify Fallback) **Files:** - Create: `cfc-backend/src/main/java/com/etotem/cfc/service/AiGateway.java` - Modify: `cfc-backend/src/main/resources/application.yml`(加 python.* 配置) **Interfaces:** - Produces: `AiGateway.recommend()` — Java 端统一的 Python 调用入口,超时/异常自动 fallback - Consumes: Python `POST /api/v1/recommend` - [x] **Step 1: 在 `application.yml` 中新增配置** ```yaml # application.yml 底部追加 python: enabled: true base-url: http://localhost:9000 timeout-ms: 15000 circuit-breaker: failure-threshold: 3 reset-timeout-ms: 30000 ``` - [x] **Step 2: 创建 `AiGateway.java`** ```java package com.etotem.cfc.service; import com.etotem.cfc.dto.RecommendationResult; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Value; import org.springframework.http.HttpEntity; import org.springframework.http.ResponseEntity; import org.springframework.stereotype.Service; import org.springframework.web.client.RestTemplate; import javax.annotation.PostConstruct; import java.util.*; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; /** * LangGraph Python 服务网关 * 统一管理 Python 服务调用、熔断、Fallback */ @Service public class AiGateway { private static final Logger log = LoggerFactory.getLogger(AiGateway.class); @Value("${python.enabled:false}") private boolean enabled; @Value("${python.base-url:http://localhost:9000}") private String baseUrl; @Value("${python.timeout-ms:15000}") private int timeoutMs; @Value("${python.circuit-breaker.failure-threshold:3}") private int failureThreshold; @Value("${python.circuit-breaker.reset-timeout-ms:30000}") private int resetTimeoutMs; private final RestTemplate restTemplate = new RestTemplate(); private final ObjectMapper objectMapper = new ObjectMapper(); // 熔断器状态 private final AtomicInteger consecutiveFailures = new AtomicInteger(0); private final AtomicLong lastFailureTime = new AtomicLong(0); private volatile boolean circuitOpen = false; @PostConstruct public void init() { if (enabled) { log.info("AiGateway 已启用: baseUrl={}, timeout={}ms", baseUrl, timeoutMs); } else { log.info("AiGateway 已禁用, 所有请求走 Dify"); } } public boolean isEnabled() { return enabled; } private boolean isCircuitOpen() { if (!circuitOpen) return false; // 检查是否到重置时间 if (System.currentTimeMillis() - lastFailureTime.get() > resetTimeoutMs) { circuitOpen = false; consecutiveFailures.set(0); log.info("AiGateway 熔断器已半开"); return false; } return true; } private void recordFailure() { int failures = consecutiveFailures.incrementAndGet(); lastFailureTime.set(System.currentTimeMillis()); if (failures >= failureThreshold) { circuitOpen = true; log.warn("AiGateway 熔断器已打开 (连续{}次失败)", failures); } } /** * 调用 Python 推荐服务, 失败时返回 null (由调用方决定 fallback) */ public List recommend(Long userId, List tags, List types, Integer limit) { if (!enabled || isCircuitOpen()) return null; try { // 使用 ObjectMapper 构造 JSON 替代手动拼接 ObjectNode body = objectMapper.createObjectNode(); body.put("user_id", userId); if (tags != null && !tags.isEmpty()) { ArrayNode tagsArray = body.putArray("tags"); tags.forEach(tagsArray::add); } if (types != null && !types.isEmpty()) { ArrayNode typesArray = body.putArray("types"); types.forEach(typesArray::add); } body.put("limit", limit != null ? limit : 5); HttpEntity entity = new HttpEntity<>(body.toString(), createJsonHeaders()); String url = baseUrl + "/api/v1/recommend"; ResponseEntity response = restTemplate.postForEntity(url, entity, String.class); if (response.getStatusCode().is2xxSuccessful() && response.getBody() != null) { JsonNode root = objectMapper.readTree(response.getBody()); JsonNode items = root.get("items"); if (items != null && items.isArray()) { List results = new ArrayList<>(); for (JsonNode item : items) { RecommendationResult r = new RecommendationResult(); r.setType(item.get("type").asText()); r.setId(item.get("id").asLong()); r.setName(item.get("name").asText()); r.setDescription(item.has("description") ? item.get("description").asText() : ""); r.setReason(item.has("reason") ? item.get("reason").asText() : ""); r.setSource("langgraph"); results.add(r); } consecutiveFailures.set(0); log.debug("AiGateway recommend 成功: {} 条", results.size()); return results; } } return null; } catch (Exception e) { log.warn("AiGateway recommend 调用失败: {}", e.getMessage()); recordFailure(); return null; } } private org.springframework.http.HttpHeaders createJsonHeaders() { org.springframework.http.HttpHeaders headers = new org.springframework.http.HttpHeaders(); headers.setContentType(org.springframework.http.MediaType.APPLICATION_JSON); return headers; } } ``` - [x] **Step 3: 编译验证** ```bash cd cfc-backend mvn clean compile -q # 预期: BUILD SUCCESS ``` - [x] **Step 4: Commit** ```bash git add cfc-backend/src/main/java/com/etotem/cfc/service/AiGateway.java git add cfc-backend/src/main/resources/application.yml git commit -m "feat(backend): AiGateway with circuit breaker for LangGraph" ``` --- ### Task 8: Java 端 RecommendController 灰度切换 **Files:** - Modify: `cfc-backend/src/main/java/com/etotem/cfc/controller/ai/RecommendationController.java` - Modify: `cfc-backend/src/main/java/com/etotem/cfc/service/RecommendationService.java` **Interfaces:** - Consumes: `AiGateway.recommend()` - 灰度策略: 5% 流量走 Python, 异常自动 fallback 现有逻辑 - [x] **Step 1: 修改 `RecommendationController.java`** ```java // 注入 AiGateway @Resource private AiGateway aiGateway; @Operation(summary = "按营养标签搜索推荐内容") @PostMapping("/search") public Result> search( @RequestBody RecommendationQuery query, HttpServletRequest request) { Long userId = (Long) request.getAttribute("userId"); if (userId == null) userId = 0L; // 灰度: 5% 流量走 Python LangGraph if (aiGateway.isEnabled() && Math.random() < 0.05) { List pythonResults = aiGateway.recommend( userId, query.getNutritionTags(), query.getTypes(), query.getLimit()); if (pythonResults != null && !pythonResults.isEmpty()) { log.info("[灰度] Recommend via LangGraph, userId={}", userId); return Result.success(pythonResults); } log.info("[灰度] LangGraph 无结果/异常, fallback 本地搜索"); } // 现有业务逻辑 List results = recommendationService.search(query); return Result.success(results); } ``` 注意: `Math.random() < 0.05` 需要在 Controller 中注入 `HttpServletRequest`,当前方法签名没有这个参数。Java 端统一通过 `@RequestAttribute("userId")` 获取,所以需要修改方法签名或改用 RequestContextHolder。 安全且兼容的改法: ```java @PostMapping("/search") public Result> search( @RequestBody RecommendationQuery query, @RequestAttribute(value = "userId", required = false) Long userId) { // ... } ``` - [x] **Step 2: 编译验证** ```bash cd cfc-backend mvn clean compile -q # 预期: BUILD SUCCESS ``` - [x] **Step 3: Commit** ```bash git add cfc-backend/src/main/java/com/etotem/cfc/controller/ai/RecommendationController.java git commit -m "feat(backend): 5% recommend traffic routed to LangGraph" ``` --- ### Task 9: Docker Compose 本地开发编排 **Files:** - Create: `cfc-langgraph/docker-compose.yml` - [x] **Step 1: 创建 `docker-compose.yml`** ```yaml version: "3.8" services: langgraph-svc: build: context: . dockerfile: Dockerfile ports: - "9000:9000" env_file: - .env environment: - JAVA_BASE_URL=http://host.docker.internal:9082 - CHROMA_DB_PATH=/data/chroma_db volumes: - langgraph-data:/data - ./app:/app/app # 开发模式: 热重载 command: uvicorn app.main:app --host 0.0.0.0 --port 9000 --reload restart: unless-stopped volumes: langgraph-data: ``` - [x] **Step 2: Commit** ```bash git add cfc-langgraph/docker-compose.yml git commit -m "chore(langgraph): Docker Compose for local dev" ``` --- ### Task 10: Phase 1 集成测试 **Files:** - Create: `cfc-langgraph/tests/__init__.py` - Create: `cfc-langgraph/tests/conftest.py` - Create: `cfc-langgraph/tests/test_recommend.py` - [x] **Step 1: 创建 `tests/__init__.py`**(空文件) - [x] **Step 2: 创建 `tests/conftest.py`** ```python import pytest from httpx import AsyncClient, ASGITransport from app.main import app @pytest.fixture async def client(): transport = ASGITransport(app=app) async with AsyncClient(transport=transport, base_url="http://test") as ac: yield ac @pytest.fixture async def health_response(client): resp = await client.get("/health") return resp ``` - [x] **Step 3: 创建 `tests/test_recommend.py`** ```python import pytest from httpx import AsyncClient, ASGITransport from app.main import app @pytest.mark.asyncio async def test_health_check(): transport = ASGITransport(app=app) async with AsyncClient(transport=transport, base_url="http://test") as client: resp = await client.get("/health") assert resp.status_code == 200 assert resp.json() == {"status": "ok"} @pytest.mark.asyncio async def test_recommend_empty_tags(): """空标签应返回空列表""" transport = ASGITransport(app=app) async with AsyncClient(transport=transport, base_url="http://test") as client: resp = await client.post("/api/v1/recommend", json={ "user_id": 1, "tags": [], "query": "推荐一些适合孩子的活动", "limit": 3, }) assert resp.status_code == 200 data = resp.json() assert "items" in data @pytest.mark.asyncio async def test_recommend_with_tags(): """带标签应返回推荐结果""" transport = ASGITransport(app=app) async with AsyncClient(transport=transport, base_url="http://test") as client: resp = await client.post("/api/v1/recommend", json={ "user_id": 1, "tags": ["益生菌", "肠胃"], "limit": 3, }) assert resp.status_code == 200 data = resp.json() assert "items" in data assert "source" in data ``` - [x] **Step 4: 运行测试** ```bash cd cfc-langgraph pip install -e ".[dev]" pytest tests/ -v # 预期: 2 passed (test_health_check, test_recommend_empty_tags 通过) # test_recommend_with_tags 可能因 Java 后端未运行而返回空列表, 但不应该报错 ``` - [x] **Step 5: Commit** ```bash git add cfc-langgraph/tests/ git commit -m "test(langgraph): Phase 1 integration tests" ``` --- ## 自审清单 - [x] 每个 Task 的产出物独立可测试 - [x] 没有 "TBD" / "TODO" 占位符 - [x] Python 侧没有直连 MySQL - [x] Java 侧灰度开关 + 熔断器 + Fallback 三层保护 - [x] Docker Compose 可一键启动开发环境 - [x] Phase 1 产出物: 运行中的 Recommend API + 知识库基础设施