更新agent

This commit is contained in:
zhanggangyong
2026-02-05 16:01:34 +00:00
commit 3a9773febd
35 changed files with 3623 additions and 0 deletions
+20
View File
@@ -0,0 +1,20 @@
FROM python:3.12-slim
WORKDIR /app
ENV PYTHONUNBUFFERED=1
ENV PYTHONDONTWRITEBYTECODE=1
RUN apt-get update && apt-get install -y gcc curl && rm -rf /var/lib/apt/lists/*
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
EXPOSE 8000
HEALTHCHECK --interval=30s --timeout=10s --start-period=5s --retries=3 \
CMD curl -f http://localhost:8000/health || exit 1
CMD ["python", "run_api_server.py"]
+67
View File
@@ -0,0 +1,67 @@
# Agent 模板
基于 **Pydantic AI** 的轻量级 Agent 模板。
## 快速开始
### 1. 复制模板
```bash
cp -r _template your_agent_name
cd your_agent_name
# 全局替换 "your_agent" 为你的 agent 名称
```
### 2. 修改核心文件
- `src/server/mcp_server.py` - 添加你的 MCP 工具
- `src/server/api_server.py` - 添加你的 API 端点(可选)
### 3. 本地测试
```bash
python run_api_server.py
```
### 4. 构建镜像
```bash
docker build -t your-agent:latest .
```
### 5. 注册到 Agent Manager
在 `k8s_manager.py` 中添加:
```python
# TEMPLATE_PORTS
"your_agent": 8000,
# image_map
"your_agent": "agnettaiji.azurecr.io/ai-agents/your-agent:latest",
```
在 `app.py` 的 `valid_templates` 中添加 `"your_agent"`。
## 项目结构
```
your_agent/
├── Dockerfile
├── requirements.txt
├── run_api_server.py # 启动脚本
└── src/
├── __init__.py
└── server/
├── __init__.py
├── api_server.py # FastAPI + MCP HTTP
└── mcp_server.py # MCP 工具定义
```
## 环境变量
| 变量 | 必需 | 说明 |
|------|------|------|
| LITELLM_GATEWAY_URL | 是 | LiteLLM Gateway URL |
| LITELLM_MODEL | 否 | 模型名称,默认 taiji/gpt-4o-mini |
| API_PORT | 否 | 端口,默认 8000 |
+13
View File
@@ -0,0 +1,13 @@
# Pydantic AI
pydantic-ai>=0.0.14
# MCP
mcp>=0.9.0
fastmcp>=0.1.0
# FastAPI
fastapi>=0.109.0
uvicorn[standard]>=0.27.0
# HTTP Client
aiohttp>=3.9.0
+17
View File
@@ -0,0 +1,17 @@
#!/usr/bin/env python
"""启动 API 服务器"""
import sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).parent))
if __name__ == '__main__':
from src.server.api_server import app
import uvicorn
import os
host = os.getenv('API_HOST', '0.0.0.0')
port = int(os.getenv('API_PORT', '8000'))
print(f"🚀 启动 Agent API: http://{host}:{port}")
uvicorn.run(app, host=host, port=port, log_level="info")
+1
View File
@@ -0,0 +1 @@
"""Agent 源代码包"""
+1
View File
@@ -0,0 +1 @@
"""服务器模块"""
+237
View File
@@ -0,0 +1,237 @@
"""
HTTP API 服务器
提供 REST API 和 MCP HTTP/SSE 端点。
"""
import json
import uuid
import os
from typing import Optional, Dict, Any, AsyncGenerator
from contextlib import asynccontextmanager
from fastapi import FastAPI, HTTPException, Request, Header, Depends
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import StreamingResponse, JSONResponse
from pydantic import BaseModel, Field
from .mcp_server import TOOL_MAP, TOOL_LIST
# ==================== 配置 ====================
SERVER_NAME = "Your Agent API" # 修改为你的 Agent 名称
# ==================== FastAPI 应用 ====================
@asynccontextmanager
async def lifespan(app: FastAPI):
print(f"🚀 {SERVER_NAME} 启动")
yield
print(f"🛑 {SERVER_NAME} 关闭")
app = FastAPI(
title=SERVER_NAME,
version="1.0.0",
lifespan=lifespan
)
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
# ==================== API Key 验证 ====================
async def verify_api_key(
api_key: Optional[str] = Header(None, alias="api-key"),
authorization: Optional[str] = Header(None)
) -> str:
"""验证 API Key"""
if api_key and api_key.strip() and api_key.strip() != "sk":
return api_key.strip()
if authorization:
key = authorization[7:].strip() if authorization.startswith("Bearer ") else authorization.strip()
if key and key != "sk":
return key
raise HTTPException(status_code=401, detail="缺少 API Key")
def get_api_key_from_request(request: Request) -> Optional[str]:
"""从请求头提取 API Key(不验证)"""
api_key = request.headers.get("api-key") or request.headers.get("api_key")
if not api_key:
auth = request.headers.get("Authorization")
if auth:
api_key = auth[7:] if auth.startswith("Bearer ") else auth
return api_key
# ==================== 健康检查 ====================
@app.get("/")
async def root():
return {
"service": SERVER_NAME,
"status": "running",
"tools": list(TOOL_MAP.keys())
}
@app.get("/health")
async def health():
return {"status": "healthy", "service": SERVER_NAME}
# ==================== MCP 端点 ====================
sessions: Dict[str, Dict] = {}
async def handle_mcp_request(data: Dict, session_id: str = None, api_key: str = None) -> Dict:
"""处理 MCP JSON-RPC 请求"""
method = data.get("method")
params = data.get("params", {})
req_id = data.get("id")
# tools/call 需要验证 API Key
if method == "tools/call" and (not api_key or api_key == "sk"):
return {"jsonrpc": "2.0", "id": req_id, "error": {"code": -32001, "message": "缺少 API Key"}}
try:
if method == "initialize":
session_id = session_id or str(uuid.uuid4())
sessions[session_id] = {"initialized": True}
return {
"jsonrpc": "2.0", "id": req_id,
"result": {
"protocolVersion": "2024-11-05",
"capabilities": {"tools": {}},
"serverInfo": {"name": SERVER_NAME, "version": "1.0.0"}
}
}
elif method == "tools/list":
return {"jsonrpc": "2.0", "id": req_id, "result": {"tools": TOOL_LIST}}
elif method == "tools/call":
tool_name = params.get("name")
args = params.get("arguments", {})
if tool_name not in TOOL_MAP:
raise ValueError(f"Unknown tool: {tool_name}")
# 设置 API Key 到环境变量
old_key = os.environ.get('OPENAI_API_KEY')
if api_key:
os.environ['OPENAI_API_KEY'] = api_key
try:
result = await TOOL_MAP[tool_name](**args)
finally:
if old_key:
os.environ['OPENAI_API_KEY'] = old_key
return {
"jsonrpc": "2.0", "id": req_id,
"result": {"content": [{"type": "text", "text": str(result)}]}
}
elif method == "ping":
return {"jsonrpc": "2.0", "id": req_id, "result": {}}
else:
raise ValueError(f"Unknown method: {method}")
except Exception as e:
return {"jsonrpc": "2.0", "id": req_id, "error": {"code": -32603, "message": str(e)}}
@app.post("/mcp")
async def mcp_endpoint(request: Request):
"""MCP HTTP 端点"""
try:
body = await request.json()
session_id = request.headers.get("x-mcp-session-id")
api_key = get_api_key_from_request(request)
response = await handle_mcp_request(body, session_id, api_key)
return JSONResponse(content=response, headers={"x-mcp-session-id": session_id or ""})
except Exception as e:
return JSONResponse(status_code=400, content={"jsonrpc": "2.0", "error": {"code": -32700, "message": str(e)}})
@app.get("/mcp/sse")
async def mcp_sse(request: Request):
"""MCP SSE 端点"""
session_id = request.headers.get("x-mcp-session-id") or str(uuid.uuid4())
async def stream() -> AsyncGenerator[str, None]:
yield f"data: {json.dumps({'type': 'connection', 'sessionId': session_id})}\n\n"
import asyncio
while True:
await asyncio.sleep(30)
yield f"data: {json.dumps({'type': 'ping'})}\n\n"
return StreamingResponse(stream(), media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "x-mcp-session-id": session_id})
@app.post("/mcp/sse")
async def mcp_sse_post(request: Request):
"""MCP SSE POST 端点"""
try:
body = await request.json()
session_id = request.headers.get("x-mcp-session-id") or str(uuid.uuid4())
api_key = get_api_key_from_request(request)
async def stream() -> AsyncGenerator[str, None]:
response = await handle_mcp_request(body, session_id, api_key)
yield f"data: {json.dumps(response)}\n\n"
return StreamingResponse(stream(), media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "x-mcp-session-id": session_id})
except Exception as e:
return JSONResponse(status_code=400, content={"jsonrpc": "2.0", "error": {"code": -32700, "message": str(e)}})
# ==================== 业务 API(可选)====================
class QueryRequest(BaseModel):
"""请求模型"""
query: str = Field(..., description="查询内容")
option: Optional[str] = Field(None, description="可选参数")
class QueryResponse(BaseModel):
"""响应模型"""
success: bool
result: Optional[str] = None
error: Optional[str] = None
@app.post("/api/v1/query", response_model=QueryResponse)
async def api_query(request: QueryRequest, api_key: str = Depends(verify_api_key)):
"""业务 API 端点(示例)"""
try:
# 设置 API Key
old_key = os.environ.get('OPENAI_API_KEY')
os.environ['OPENAI_API_KEY'] = api_key
try:
result = await TOOL_MAP['your_tool'](query=request.query, option=request.option)
return QueryResponse(success=True, result=result)
finally:
if old_key:
os.environ['OPENAI_API_KEY'] = old_key
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
if __name__ == '__main__':
import uvicorn
uvicorn.run(app, host="0.0.0.0", port=8000)
+109
View File
@@ -0,0 +1,109 @@
"""
MCP 服务器 - 定义 Agent 工具
使用 Pydantic AI 和 FastMCP 框架。
在此文件中添加你的 MCP 工具。
"""
import json
import os
from typing import Optional
from mcp.server.fastmcp import FastMCP
from pydantic_ai import Agent
# ==================== 配置 ====================
# LiteLLM Gateway 配置
_BASE_URL = os.getenv('OPENAI_BASE_URL',
os.getenv('LLM_BASE_URL', 'https://litellm.graystone-fb459c5d.southeastasia.azurecontainerapps.io/v1'))
_API_KEY = os.getenv('OPENAI_API_KEY', 'sk')
os.environ.setdefault('OPENAI_API_KEY', _API_KEY)
os.environ.setdefault('OPENAI_BASE_URL', _BASE_URL)
# 模型名称(pydantic_ai 需要 openai: 前缀)
def _get_model_name() -> str:
model = os.getenv('MODEL_NAME', os.getenv('LITELLM_MODEL', 'taiji/gpt-4o-mini'))
return model if ':' in model else f'openai:{model}'
MODEL_NAME = _get_model_name()
# ==================== MCP 服务器 ====================
server = FastMCP('Your Agent') # 修改为你的 Agent 名称
# 系统提示词
SYSTEM_PROMPT = '''你是一个专业的 AI 助手。
请根据用户的需求提供帮助。'''
def get_agent() -> Agent:
"""创建 Agent 实例(每次调用使用最新的 API Key)"""
return Agent(MODEL_NAME, system_prompt=SYSTEM_PROMPT)
# ==================== MCP 工具定义 ====================
@server.tool()
async def your_tool(
query: str,
option: Optional[str] = None
) -> str:
"""
你的工具描述
Args:
query: 查询内容
option: 可选参数
Returns:
处理结果(JSON 格式)
"""
try:
# 1. 调用 AI Agent 处理
result = await get_agent().run(f"请处理: {query}")
# 2. 返回结果
return json.dumps({
"success": True,
"query": query,
"result": result.output
}, ensure_ascii=False, indent=2)
except Exception as e:
return json.dumps({
"success": False,
"error": str(e)
}, ensure_ascii=False)
# 添加更多工具...
# @server.tool()
# async def another_tool(...) -> str:
# pass
# ==================== 工具映射(供 API 使用)====================
TOOL_MAP = {
'your_tool': your_tool,
}
TOOL_LIST = [
{
"name": "your_tool",
"description": "你的工具描述",
"inputSchema": {
"type": "object",
"properties": {
"query": {"type": "string", "description": "查询内容"},
"option": {"type": "string", "description": "可选参数"}
},
"required": ["query"]
}
}
]
if __name__ == '__main__':
server.run()
+21
View File
@@ -0,0 +1,21 @@
FROM python:3.12-slim
WORKDIR /app
ENV PYTHONUNBUFFERED=1
ENV PYTHONDONTWRITEBYTECODE=1
RUN apt-get update && apt-get install -y gcc curl && rm -rf /var/lib/apt/lists/*
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
EXPOSE 8000
HEALTHCHECK --interval=30s --timeout=10s --start-period=5s --retries=3 \
CMD curl -f http://localhost:8000/health || exit 1
CMD ["python", "run_api_server.py"]
+113
View File
@@ -0,0 +1,113 @@
# 信息去重 Agent
基于 **Pydantic AI** 的信息去重 Agent,专注于从一堆信息中识别并合并"本质相同"的内容。
## 核心功能
- **语义去重**:不只是文字相同,而是识别本质含义相同的信息
- **聚类分组**:将相似信息归为一组,保留最具代表性的表述
- **差异分析**:识别看似相同实则有细微差异的信息
## 快速开始
### 1. 本地测试
```bash
cd dedup_agent
python run_api_server.py
```
### 2. 构建镜像
```bash
docker build -t dedup-agent:latest .
```
## API 使用
### MCP 工具调用
```json
{
"method": "tools/call",
"params": {
"name": "deduplicate_info",
"arguments": {
"items": [
"苹果公司发布了新款 iPhone",
"Apple 推出最新 iPhone 系列",
"微软发布 Windows 更新",
"苹果今天发布了 iPhone 新产品"
]
}
}
}
```
### REST API
```bash
curl -X POST http://localhost:8000/api/v1/deduplicate \
-H "Content-Type: application/json" \
-H "api-key: your-api-key" \
-d '{
"items": [
"苹果公司发布了新款 iPhone",
"Apple 推出最新 iPhone 系列",
"微软发布 Windows 更新"
]
}'
```
## 项目结构
```
dedup_agent/
├── Dockerfile
├── requirements.txt
├── run_api_server.py # 启动脚本
├── README.md
└── src/
├── __init__.py
└── server/
├── __init__.py
├── api_server.py # FastAPI + MCP HTTP
└── mcp_server.py # MCP 工具定义(去重逻辑)
```
## 环境变量
| 变量 | 必需 | 说明 |
|------|------|------|
| OPENAI_BASE_URL | 否 | LiteLLM Gateway URL |
| OPENAI_API_KEY | 是 | API Key |
| MODEL_NAME | 否 | 模型名称,默认 taiji/gpt-4o-mini |
| API_PORT | 否 | 端口,默认 8000 |
## 响应示例
```json
{
"success": true,
"groups": [
{
"representative": "苹果公司发布了新款 iPhone",
"members": [
"苹果公司发布了新款 iPhone",
"Apple 推出最新 iPhone 系列",
"苹果今天发布了 iPhone 新产品"
],
"essence": "苹果发布新iPhone"
},
{
"representative": "微软发布 Windows 更新",
"members": ["微软发布 Windows 更新"],
"essence": "微软Windows更新"
}
],
"total_input": 4,
"total_groups": 2,
"dedup_ratio": "50%"
}
```
+351
View File
@@ -0,0 +1,351 @@
# 信息去重 Agent
---
本 Agent 提供基于 **领域理解** 的信息去重能力,通过 **HTTP API** 与 **MCP(Model Context Protocol)** 对外提供服务。
核心能力:
- **语义去重**:基于深度语义分析识别本质相同的信息
- **重复检测**:多粒度阈值的重复项识别
- **本质提取**:信息核心语义的精准提炼
---
## 功能概览
提供信息集合的 **语义理解、领域判断、重复检测、本质提取** 能力,返回结构化分析结果。
支持能力:
- 语义聚类与去重
- 可调阈值的重复检测
- 信息本质与关键词提取
- 去重率统计
---
## 1⃣ deduplicate_info — 语义领域去重
### 功能说明
对输入的信息集合进行深度领域分析,识别 **本质相同** 的内容并归类合并,返回去重后的分组结果。
---
### REST API 调用
```
POST /api/v1/deduplicate
Content-Type: application/json
api-key: {your-api-key}
```
```json
{
"items": [
"苹果公司发布了新款 iPhone",
"Apple 推出最新 iPhone 系列",
"微软发布 Windows 更新",
"苹果今天发布了 iPhone 新产品"
],
"context": "科技新闻"
}
```
---
### MCP 调用
```json
{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "deduplicate_info",
"arguments": {
"items": [
"苹果公司发布了新款 iPhone",
"Apple 推出最新 iPhone 系列",
"微软发布 Windows 更新",
"苹果今天发布了 iPhone 新产品"
],
"context": "科技新闻"
}
}
}
```
---
### 参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
| --- | --- | --- | --- | --- |
| items | array\<string\> | ✅ | - | 待去重的信息列表 |
| context | string | ❌ | null | 上下文说明,辅助语义理解 |
---
### 返回结果
```json
{
"success": true,
"groups": [
{
"representative": "苹果公司发布了新款 iPhone",
"members": [
"苹果公司发布了新款 iPhone",
"Apple 推出最新 iPhone 系列",
"苹果今天发布了 iPhone 新产品"
],
"essence": "苹果发布新iPhone"
},
{
"representative": "微软发布 Windows 更新",
"members": ["微软发布 Windows 更新"],
"essence": "微软Windows更新"
}
],
"total_input": 4,
"total_groups": 2,
"dedup_ratio": "50%"
}
```
---
### 返回字段说明
| 字段 | 类型 | 说明 |
| --- | --- | --- |
| groups | array | 去重后的分组列表 |
| groups[].representative | string | 该组的代表性表述(原文) |
| groups[].members | array | 该组包含的所有原始信息 |
| groups[].essence | string | 该组信息的本质概括 |
| total_input | integer | 输入信息总数 |
| total_groups | integer | 去重后分组数 |
| dedup_ratio | string | 去重率 |
---
## 2️⃣ find_duplicates — 重复检测
### 功能说明
快速识别信息集合中的重复项,支持 **多粒度阈值** 控制检测严格程度。
---
### REST API 调用
```
POST /api/v1/find-duplicates
Content-Type: application/json
api-key: {your-api-key}
```
```json
{
"items": [
"今天天气很好",
"今日阳光明媚",
"明天会下雨",
"天气晴朗适合外出"
],
"threshold": "normal"
}
```
---
### MCP 调用
```json
{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "find_duplicates",
"arguments": {
"items": [
"今天天气很好",
"今日阳光明媚",
"明天会下雨"
],
"threshold": "strict"
}
}
}
```
---
### 参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
| --- | --- | --- | --- | --- |
| items | array\<string\> | ✅ | - | 待检测的信息列表 |
| threshold | string | ❌ | normal | 检测阈值:strict / normal / loose |
---
### 阈值说明
| 阈值 | 检测策略 |
| --- | --- |
| strict | 仅识别近乎完全相同的内容 |
| normal | 识别本质含义相同的内容 |
| loose | 识别主题相关的内容 |
---
### 返回结果
```json
{
"success": true,
"has_duplicates": true,
"duplicate_pairs": [
{
"items": [1, 2, 4],
"reason": "均表达天气状况良好"
}
],
"unique_count": 2
}
```
---
## 3️⃣ extract_essence — 本质提取
### 功能说明
对每条信息进行语义分析,提取其 **核心本质** 与 **关键词**,便于后续比对或索引。
---
### REST API 调用
```
POST /api/v1/extract-essence
Content-Type: application/json
api-key: {your-api-key}
```
```json
{
"items": [
"特斯拉宣布下调全系车型售价",
"OpenAI 发布了 GPT-5 模型",
"中国央行决定降息25个基点"
]
}
```
---
### MCP 调用
```json
{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "extract_essence",
"arguments": {
"items": [
"特斯拉宣布下调全系车型售价",
"OpenAI 发布了 GPT-5 模型"
]
}
}
}
```
---
### 参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
| --- | --- | --- | --- | --- |
| items | array\<string\> | ✅ | - | 待提取本质的信息列表 |
---
### 返回结果
```json
{
"success": true,
"essences": [
{
"original": "特斯拉宣布下调全系车型售价",
"essence": "特斯拉降价",
"keywords": ["特斯拉", "降价", "车型"]
},
{
"original": "OpenAI 发布了 GPT-5 模型",
"essence": "OpenAI发布新模型",
"keywords": ["OpenAI", "GPT-5", "模型"]
},
{
"original": "中国央行决定降息25个基点",
"essence": "央行降息",
"keywords": ["央行", "降息", "利率"]
}
]
}
```
---
## 统一错误格式
成功:
```json
{
"success": true,
"data": {}
}
```
失败:
```json
{
"success": false,
"error": "错误描述"
}
```
---
## 服务端点
| 端点 | 方法 | 说明 |
| --- | --- | --- |
| / | GET | 服务状态 |
| /health | GET | 健康检查 |
| /mcp | POST | MCP JSON-RPC |
| /mcp/sse | GET/POST | MCP SSE 流式 |
| /api/v1/deduplicate | POST | 语义去重 |
| /api/v1/find-duplicates | POST | 重复检测 |
| /api/v1/extract-essence | POST | 本质提取 |
---
## 部署信息
| 配置项 | 值 |
| --- | --- |
| 镜像地址 | agnettaiji.azurecr.io/ai-agents/dedup-agent:latest |
| 服务端口 | 8000 |
| 健康检查 | /health |
+14
View File
@@ -0,0 +1,14 @@
# Pydantic AI
pydantic-ai>=0.0.14
# MCP
mcp>=0.9.0
fastmcp>=0.1.0
# FastAPI
fastapi>=0.109.0
uvicorn[standard]>=0.27.0
# HTTP Client
aiohttp>=3.9.0
+18
View File
@@ -0,0 +1,18 @@
#!/usr/bin/env python
"""启动信息去重 Agent API 服务器"""
import sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).parent))
if __name__ == '__main__':
from src.server.api_server import app
import uvicorn
import os
host = os.getenv('API_HOST', '0.0.0.0')
port = int(os.getenv('API_PORT', '8000'))
print(f"🚀 启动信息去重 Agent API: http://{host}:{port}")
uvicorn.run(app, host=host, port=port, log_level="info")
+2
View File
@@ -0,0 +1,2 @@
"""信息去重 Agent 源代码包"""
+2
View File
@@ -0,0 +1,2 @@
"""服务器模块"""
+289
View File
@@ -0,0 +1,289 @@
"""
信息去重 Agent - HTTP API 服务器
提供 REST API 和 MCP HTTP/SSE 端点。
"""
import json
import uuid
import os
from typing import Optional, Dict, Any, AsyncGenerator, List
from contextlib import asynccontextmanager
from fastapi import FastAPI, HTTPException, Request, Header, Depends
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import StreamingResponse, JSONResponse
from pydantic import BaseModel, Field
from .mcp_server import TOOL_MAP, TOOL_LIST
# ==================== 配置 ====================
SERVER_NAME = "信息去重 Agent API"
# ==================== FastAPI 应用 ====================
@asynccontextmanager
async def lifespan(app: FastAPI):
print(f"🚀 {SERVER_NAME} 启动")
yield
print(f"🛑 {SERVER_NAME} 关闭")
app = FastAPI(
title=SERVER_NAME,
description="从一堆信息中找出本质相同的内容并去重",
version="1.0.0",
lifespan=lifespan
)
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
# ==================== API Key 验证 ====================
async def verify_api_key(
api_key: Optional[str] = Header(None, alias="api-key"),
authorization: Optional[str] = Header(None)
) -> str:
"""验证 API Key"""
if api_key and api_key.strip() and api_key.strip() != "sk":
return api_key.strip()
if authorization:
key = authorization[7:].strip() if authorization.startswith("Bearer ") else authorization.strip()
if key and key != "sk":
return key
raise HTTPException(status_code=401, detail="缺少 API Key")
def get_api_key_from_request(request: Request) -> Optional[str]:
"""从请求头提取 API Key(不验证)"""
api_key = request.headers.get("api-key") or request.headers.get("api_key")
if not api_key:
auth = request.headers.get("Authorization")
if auth:
api_key = auth[7:] if auth.startswith("Bearer ") else auth
return api_key
# ==================== 健康检查 ====================
@app.get("/")
async def root():
return {
"service": SERVER_NAME,
"description": "信息去重 Agent - 找出本质相同的信息",
"status": "running",
"tools": list(TOOL_MAP.keys())
}
@app.get("/health")
async def health():
return {"status": "healthy", "service": SERVER_NAME}
# ==================== MCP 端点 ====================
sessions: Dict[str, Dict] = {}
async def handle_mcp_request(data: Dict, session_id: str = None, api_key: str = None) -> Dict:
"""处理 MCP JSON-RPC 请求"""
method = data.get("method")
params = data.get("params", {})
req_id = data.get("id")
# tools/call 需要验证 API Key
if method == "tools/call" and (not api_key or api_key == "sk"):
return {"jsonrpc": "2.0", "id": req_id, "error": {"code": -32001, "message": "缺少 API Key"}}
try:
if method == "initialize":
session_id = session_id or str(uuid.uuid4())
sessions[session_id] = {"initialized": True}
return {
"jsonrpc": "2.0", "id": req_id,
"result": {
"protocolVersion": "2024-11-05",
"capabilities": {"tools": {}},
"serverInfo": {"name": SERVER_NAME, "version": "1.0.0"}
}
}
elif method == "tools/list":
return {"jsonrpc": "2.0", "id": req_id, "result": {"tools": TOOL_LIST}}
elif method == "tools/call":
tool_name = params.get("name")
args = params.get("arguments", {})
if tool_name not in TOOL_MAP:
raise ValueError(f"Unknown tool: {tool_name}")
# 设置 API Key 到环境变量
old_key = os.environ.get('OPENAI_API_KEY')
if api_key:
os.environ['OPENAI_API_KEY'] = api_key
try:
result = await TOOL_MAP[tool_name](**args)
finally:
if old_key:
os.environ['OPENAI_API_KEY'] = old_key
return {
"jsonrpc": "2.0", "id": req_id,
"result": {"content": [{"type": "text", "text": str(result)}]}
}
elif method == "ping":
return {"jsonrpc": "2.0", "id": req_id, "result": {}}
else:
raise ValueError(f"Unknown method: {method}")
except Exception as e:
return {"jsonrpc": "2.0", "id": req_id, "error": {"code": -32603, "message": str(e)}}
@app.post("/mcp")
async def mcp_endpoint(request: Request):
"""MCP HTTP 端点"""
try:
body = await request.json()
session_id = request.headers.get("x-mcp-session-id")
api_key = get_api_key_from_request(request)
response = await handle_mcp_request(body, session_id, api_key)
return JSONResponse(content=response, headers={"x-mcp-session-id": session_id or ""})
except Exception as e:
return JSONResponse(status_code=400, content={"jsonrpc": "2.0", "error": {"code": -32700, "message": str(e)}})
@app.get("/mcp/sse")
async def mcp_sse(request: Request):
"""MCP SSE 端点"""
session_id = request.headers.get("x-mcp-session-id") or str(uuid.uuid4())
async def stream() -> AsyncGenerator[str, None]:
yield f"data: {json.dumps({'type': 'connection', 'sessionId': session_id})}\n\n"
import asyncio
while True:
await asyncio.sleep(30)
yield f"data: {json.dumps({'type': 'ping'})}\n\n"
return StreamingResponse(stream(), media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "x-mcp-session-id": session_id})
@app.post("/mcp/sse")
async def mcp_sse_post(request: Request):
"""MCP SSE POST 端点"""
try:
body = await request.json()
session_id = request.headers.get("x-mcp-session-id") or str(uuid.uuid4())
api_key = get_api_key_from_request(request)
async def stream() -> AsyncGenerator[str, None]:
response = await handle_mcp_request(body, session_id, api_key)
yield f"data: {json.dumps(response)}\n\n"
return StreamingResponse(stream(), media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "x-mcp-session-id": session_id})
except Exception as e:
return JSONResponse(status_code=400, content={"jsonrpc": "2.0", "error": {"code": -32700, "message": str(e)}})
# ==================== 业务 API ====================
class DeduplicateRequest(BaseModel):
"""去重请求"""
items: List[str] = Field(..., description="待去重的信息列表", min_length=1)
context: Optional[str] = Field(None, description="可选的上下文说明")
class FindDuplicatesRequest(BaseModel):
"""查找重复请求"""
items: List[str] = Field(..., description="待检测的信息列表", min_length=1)
threshold: Optional[str] = Field("normal", description="重复判定严格程度: strict/normal/loose")
class ExtractEssenceRequest(BaseModel):
"""提取本质请求"""
items: List[str] = Field(..., description="信息列表", min_length=1)
@app.post("/api/v1/deduplicate")
async def api_deduplicate(request: DeduplicateRequest, api_key: str = Depends(verify_api_key)):
"""
信息去重 API
对一组信息进行语义去重,找出本质相同的内容并合并分组。
"""
try:
old_key = os.environ.get('OPENAI_API_KEY')
os.environ['OPENAI_API_KEY'] = api_key
try:
result = await TOOL_MAP['deduplicate_info'](items=request.items, context=request.context)
return JSONResponse(content=json.loads(result))
finally:
if old_key:
os.environ['OPENAI_API_KEY'] = old_key
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
@app.post("/api/v1/find-duplicates")
async def api_find_duplicates(request: FindDuplicatesRequest, api_key: str = Depends(verify_api_key)):
"""
查找重复项 API
快速识别列表中哪些信息是重复的。
"""
try:
old_key = os.environ.get('OPENAI_API_KEY')
os.environ['OPENAI_API_KEY'] = api_key
try:
result = await TOOL_MAP['find_duplicates'](items=request.items, threshold=request.threshold)
return JSONResponse(content=json.loads(result))
finally:
if old_key:
os.environ['OPENAI_API_KEY'] = old_key
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
@app.post("/api/v1/extract-essence")
async def api_extract_essence(request: ExtractEssenceRequest, api_key: str = Depends(verify_api_key)):
"""
提取本质 API
提取每条信息的本质/核心含义。
"""
try:
old_key = os.environ.get('OPENAI_API_KEY')
os.environ['OPENAI_API_KEY'] = api_key
try:
result = await TOOL_MAP['extract_essence'](items=request.items)
return JSONResponse(content=json.loads(result))
finally:
if old_key:
os.environ['OPENAI_API_KEY'] = old_key
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
if __name__ == '__main__':
import uvicorn
uvicorn.run(app, host="0.0.0.0", port=8000)
+380
View File
@@ -0,0 +1,380 @@
"""
信息去重 MCP 服务器
核心功能:从一堆信息里,找"本质相同的点"并去重。
使用 AI 进行语义理解,识别本质相同但表述不同的信息。
"""
import json
import os
from typing import List, Optional
from mcp.server.fastmcp import FastMCP
from pydantic_ai import Agent
# ==================== 配置 ====================
# LiteLLM Gateway 配置
_BASE_URL = os.getenv('OPENAI_BASE_URL',
os.getenv('LLM_BASE_URL', 'https://litellm.graystone-fb459c5d.southeastasia.azurecontainerapps.io/v1'))
_API_KEY = os.getenv('OPENAI_API_KEY', 'sk')
os.environ.setdefault('OPENAI_API_KEY', _API_KEY)
os.environ.setdefault('OPENAI_BASE_URL', _BASE_URL)
# 模型名称(pydantic_ai 需要 openai: 前缀)
def _get_model_name() -> str:
model = os.getenv('MODEL_NAME', os.getenv('LITELLM_MODEL', 'taiji/gpt-4o-mini'))
return model if ':' in model else f'openai:{model}'
MODEL_NAME = _get_model_name()
# ==================== MCP 服务器 ====================
server = FastMCP('Dedup Agent')
# 系统提示词 - 专注于信息去重
SYSTEM_PROMPT = '''你是一个专业的信息去重分析师。你的任务是:
1. 分析一组信息,找出"本质相同"的内容
2. "本质相同"意味着:核心含义、关键事实、主要观点相同,即使表述方式不同
3. 将本质相同的信息归为一组
4. 为每组提取出"本质"(用最精简的方式表达核心含义)
5. 选择每组中最具代表性、最清晰的表述作为代表
注意事项:
- 不要仅根据表面文字相似度判断,要理解深层含义
- 考虑同一事物的不同表述方式(如"苹果公司"和"Apple"指同一实体)
- 如果两条信息有细微但重要的差异,应该分开
- 输出必须是有效的 JSON 格式'''
def get_agent() -> Agent:
"""创建 Agent 实例(每次调用使用最新的 API Key)"""
return Agent(MODEL_NAME, system_prompt=SYSTEM_PROMPT)
# ==================== MCP 工具定义 ====================
@server.tool()
async def deduplicate_info(
items: List[str],
context: Optional[str] = None
) -> str:
"""
对一组信息进行语义去重,找出本质相同的内容并合并。
Args:
items: 待去重的信息列表,每个元素是一条信息
context: 可选的上下文说明,帮助理解信息的背景
Returns:
去重结果(JSON 格式),包含分组信息和去重统计
"""
if not items:
return json.dumps({
"success": True,
"groups": [],
"total_input": 0,
"total_groups": 0,
"dedup_ratio": "0%"
}, ensure_ascii=False)
if len(items) == 1:
return json.dumps({
"success": True,
"groups": [{
"representative": items[0],
"members": items,
"essence": items[0]
}],
"total_input": 1,
"total_groups": 1,
"dedup_ratio": "0%"
}, ensure_ascii=False)
try:
# 构建提示
items_text = "\n".join([f"{i+1}. {item}" for i, item in enumerate(items)])
context_hint = ""
if context:
context_hint = f"\n\n背景说明:{context}"
prompt = f'''请分析以下 {len(items)} 条信息,找出本质相同的内容并分组:
{context_hint}
信息列表:
{items_text}
请按以下 JSON 格式输出(不要添加其他内容):
{{
"groups": [
{{
"representative": "最具代表性的原文",
"members": ["原文1", "原文2", ...],
"essence": "用最精简的方式表达这组信息的本质"
}}
]
}}
要求:
1. 每条信息必须且只能属于一个组
2. "members" 必须是原文,不要修改
3. "representative" 必须是 members 中的一条
4. "essence" 要简洁,抓住核心'''
# 调用 AI 分析
result = await get_agent().run(prompt)
output = result.output.strip()
# 提取 JSON(处理可能的 markdown 代码块)
if '```json' in output:
output = output.split('```json')[1].split('```')[0].strip()
elif '```' in output:
output = output.split('```')[1].split('```')[0].strip()
# 解析结果
parsed = json.loads(output)
groups = parsed.get('groups', [])
# 计算去重率
total_input = len(items)
total_groups = len(groups)
dedup_count = total_input - total_groups
dedup_ratio = f"{round(dedup_count / total_input * 100)}%" if total_input > 0 else "0%"
return json.dumps({
"success": True,
"groups": groups,
"total_input": total_input,
"total_groups": total_groups,
"dedup_ratio": dedup_ratio
}, ensure_ascii=False, indent=2)
except json.JSONDecodeError as e:
return json.dumps({
"success": False,
"error": f"AI 返回的结果不是有效的 JSON: {str(e)}",
"raw_output": result.output if 'result' in dir() else None
}, ensure_ascii=False)
except Exception as e:
return json.dumps({
"success": False,
"error": str(e)
}, ensure_ascii=False)
@server.tool()
async def find_duplicates(
items: List[str],
threshold: Optional[str] = "normal"
) -> str:
"""
快速识别列表中哪些信息是重复的(仅返回重复项)。
Args:
items: 待检测的信息列表
threshold: 重复判定严格程度 - "strict"(严格)/"normal"(正常)/"loose"(宽松)
Returns:
重复检测结果(JSON 格式)
"""
if len(items) < 2:
return json.dumps({
"success": True,
"has_duplicates": False,
"duplicate_pairs": [],
"unique_count": len(items)
}, ensure_ascii=False)
try:
threshold_desc = {
"strict": "只有几乎完全相同的才算重复",
"normal": "本质含义相同就算重复",
"loose": "主题相关就算重复"
}.get(threshold, "本质含义相同就算重复")
items_text = "\n".join([f"{i+1}. {item}" for i, item in enumerate(items)])
prompt = f'''分析以下信息,找出重复的内容。
判定标准:{threshold_desc}
信息列表:
{items_text}
请按以下 JSON 格式输出:
{{
"duplicate_pairs": [
{{
"items": [1, 3],
"reason": "简短说明为什么这些是重复的"
}}
]
}}
注意:
- "items" 是信息的序号(从1开始)
- 如果有多条重复,放在同一个 pair 里
- 如果没有重复,duplicate_pairs 为空数组'''
result = await get_agent().run(prompt)
output = result.output.strip()
# 提取 JSON
if '```json' in output:
output = output.split('```json')[1].split('```')[0].strip()
elif '```' in output:
output = output.split('```')[1].split('```')[0].strip()
parsed = json.loads(output)
duplicate_pairs = parsed.get('duplicate_pairs', [])
# 计算唯一项数量
duplicated_indices = set()
for pair in duplicate_pairs:
for idx in pair.get('items', []):
duplicated_indices.add(idx)
return json.dumps({
"success": True,
"has_duplicates": len(duplicate_pairs) > 0,
"duplicate_pairs": duplicate_pairs,
"unique_count": len(items) - len(duplicated_indices) + len(duplicate_pairs)
}, ensure_ascii=False, indent=2)
except Exception as e:
return json.dumps({
"success": False,
"error": str(e)
}, ensure_ascii=False)
@server.tool()
async def extract_essence(
items: List[str]
) -> str:
"""
提取每条信息的本质/核心含义,便于后续比较。
Args:
items: 信息列表
Returns:
每条信息的本质提取结果(JSON 格式)
"""
if not items:
return json.dumps({
"success": True,
"essences": []
}, ensure_ascii=False)
try:
items_text = "\n".join([f"{i+1}. {item}" for i, item in enumerate(items)])
prompt = f'''为以下每条信息提取其"本质"——用最精简的方式表达核心含义。
信息列表:
{items_text}
请按以下 JSON 格式输出:
{{
"essences": [
{{
"original": "原文",
"essence": "本质(10字以内)",
"keywords": ["关键词1", "关键词2"]
}}
]
}}'''
result = await get_agent().run(prompt)
output = result.output.strip()
if '```json' in output:
output = output.split('```json')[1].split('```')[0].strip()
elif '```' in output:
output = output.split('```')[1].split('```')[0].strip()
parsed = json.loads(output)
return json.dumps({
"success": True,
"essences": parsed.get('essences', [])
}, ensure_ascii=False, indent=2)
except Exception as e:
return json.dumps({
"success": False,
"error": str(e)
}, ensure_ascii=False)
# ==================== 工具映射(供 API 使用)====================
TOOL_MAP = {
'deduplicate_info': deduplicate_info,
'find_duplicates': find_duplicates,
'extract_essence': extract_essence,
}
TOOL_LIST = [
{
"name": "deduplicate_info",
"description": "对一组信息进行语义去重,找出本质相同的内容并合并分组",
"inputSchema": {
"type": "object",
"properties": {
"items": {
"type": "array",
"items": {"type": "string"},
"description": "待去重的信息列表"
},
"context": {
"type": "string",
"description": "可选的上下文说明"
}
},
"required": ["items"]
}
},
{
"name": "find_duplicates",
"description": "快速识别列表中哪些信息是重复的",
"inputSchema": {
"type": "object",
"properties": {
"items": {
"type": "array",
"items": {"type": "string"},
"description": "待检测的信息列表"
},
"threshold": {
"type": "string",
"enum": ["strict", "normal", "loose"],
"description": "重复判定严格程度"
}
},
"required": ["items"]
}
},
{
"name": "extract_essence",
"description": "提取每条信息的本质/核心含义",
"inputSchema": {
"type": "object",
"properties": {
"items": {
"type": "array",
"items": {"type": "string"},
"description": "信息列表"
}
},
"required": ["items"]
}
}
]
if __name__ == '__main__':
server.run()
+21
View File
@@ -0,0 +1,21 @@
FROM python:3.12-slim
WORKDIR /app
ENV PYTHONUNBUFFERED=1
ENV PYTHONDONTWRITEBYTECODE=1
RUN apt-get update && apt-get install -y gcc curl && rm -rf /var/lib/apt/lists/*
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
EXPOSE 8000
HEALTHCHECK --interval=30s --timeout=10s --start-period=5s --retries=3 \
CMD curl -f http://localhost:8000/health || exit 1
CMD ["python", "run_api_server.py"]
+145
View File
@@ -0,0 +1,145 @@
# Devil's Advocate Agent 😈
**反对意见 Agent** - 专门负责"挑刺"的 AI Agent。
## 功能特点
只做一件事:**挑刺**
- 🎯 对任何观点、方案、想法提出反对意见和质疑
- 🔍 找出潜在的问题、漏洞、风险和盲点
- ❓ 提出尖锐但有建设性的批判性问题
- 💡 挑战假设,质疑前提
- 🔬 从不同角度发现可能被忽视的缺陷
## 快速开始
### 本地测试
```bash
cd devils_advocate_agent
pip install -r requirements.txt
python run_api_server.py
```
### 构建镜像
```bash
docker build -t devils-advocate-agent:latest .
```
## API 端点
### REST API
```bash
# 挑刺 API
curl -X POST http://localhost:8000/api/v1/challenge \
-H "Content-Type: application/json" \
-H "api-key: your-api-key" \
-d '{
"content": "我们计划使用微服务架构来重构单体应用",
"focus_area": "可行性"
}'
```
### MCP 端点
| 端点 | 方法 | 说明 |
|------|------|------|
| `/mcp` | POST | MCP HTTP 端点 |
| `/mcp/sse` | GET | MCP SSE 连接 |
| `/mcp/sse` | POST | MCP SSE 请求 |
### MCP 工具
#### `challenge` - 挑刺工具
对任何观点、方案、想法进行挑刺和批判性分析。
**参数:**
| 参数 | 类型 | 必需 | 说明 |
|------|------|------|------|
| `content` | string | 是 | 需要被挑刺的内容(观点、方案、想法、代码、文档等) |
| `focus_area` | string | 否 | 指定关注的领域(如:逻辑漏洞、可行性、安全性、成本、用户体验等) |
**MCP 调用示例:**
```json
{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "challenge",
"arguments": {
"content": "我们的新产品将在3个月内开发完成并上线",
"focus_area": "时间规划"
}
}
}
```
## MCP 配置示例
在你的 MCP 客户端配置中添加:
```json
{
"mcpServers": {
"devils-advocate": {
"url": "http://localhost:8000/mcp/sse",
"transport": "sse"
}
}
}
```
## 项目结构
```
devils_advocate_agent/
├── Dockerfile
├── README.md
├── requirements.txt
├── run_api_server.py # 启动脚本
└── src/
├── __init__.py
└── server/
├── __init__.py
├── api_server.py # FastAPI + MCP HTTP
└── mcp_server.py # MCP 工具定义
```
## 环境变量
| 变量 | 必需 | 说明 |
|------|------|------|
| OPENAI_BASE_URL | 否 | LLM API Base URL |
| OPENAI_API_KEY | 否 | API Key(可通过请求头传递) |
| MODEL_NAME | 否 | 模型名称,默认 taiji/gpt-4o-mini |
| API_PORT | 否 | 端口,默认 8000 |
## 注册到 Agent Manager
在 `k8s_manager.py` 中添加:
```python
# TEMPLATE_PORTS
"devils_advocate": 8000,
# image_map
"devils_advocate": "agnettaiji.azurecr.io/ai-agents/devils-advocate-agent:latest",
```
在 `app.py` 的 `valid_templates` 中添加 `"devils_advocate"`。
## 使用场景
- 📋 **方案评审**:对项目方案、技术方案进行批判性审查
- 💻 **代码审查**:发现代码中的潜在问题和风险
- 📝 **文档审核**:检查文档的逻辑漏洞和表述问题
- 🤔 **决策验证**:在做出重要决策前进行"压力测试"
- 🎨 **设计评估**:对产品设计、架构设计提出质疑
+158
View File
@@ -0,0 +1,158 @@
# 反对意见 Agent(Devil's Advocate)
---
本 Agent 提供基于 **批判性思维模型** 的对抗性分析能力,通过 **HTTP API** 与 **MCP(Model Context Protocol)** 对外提供服务。
核心能力:
- **对抗性分析**:系统性挑战假设、质疑前提、暴露逻辑漏洞
- **风险识别**:多维度识别潜在问题、盲点与失败模式
- **批判性反馈**:提供结构化的反驳论点与改进建议
---
## 功能概览
提供观点、方案、想法的 **逻辑验证、假设挑战、风险评估** 能力,返回结构化批判分析结果。
支持能力:
- 逻辑漏洞与认知偏差识别
- 可行性与风险压力测试
- 假设前提的系统性质疑
- 盲点与失败模式挖掘
---
## 1⃣ challenge — 批判性对抗分析
### 功能说明
对输入内容执行 **Devil's Advocate** 分析范式,系统性识别论证缺陷、隐含假设、潜在风险与逻辑盲点,返回结构化批判报告。
---
### REST API 调用
```
POST /api/v1/challenge
Content-Type: application/json
api-key: {your-api-key}
```
```json
{
"content": "我们计划使用微服务架构重构现有单体应用,预计3个月内完成迁移",
"focus_area": "可行性"
}
```
---
### MCP 调用
```json
{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "challenge",
"arguments": {
"content": "我们计划使用微服务架构重构现有单体应用,预计3个月内完成迁移",
"focus_area": "可行性"
}
}
}
```
---
### 参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
| --- | --- | --- | --- | --- |
| content | string | ✅ | - | 待分析的内容(观点、方案、代码、文档等) |
| focus_area | string | ❌ | null | 聚焦领域:逻辑漏洞 / 可行性 / 安全性 / 成本 / 用户体验 等 |
---
### 聚焦领域说明
| 领域 | 分析维度 |
| --- | --- |
| 逻辑漏洞 | 推理谬误、因果混淆、循环论证、过度泛化 |
| 可行性 | 资源约束、时间估算、技术依赖、执行风险 |
| 安全性 | 攻击面、数据泄露、权限漏洞、合规风险 |
| 成本 | 隐性成本、机会成本、维护成本、规模效应 |
| 用户体验 | 认知负荷、操作摩擦、边界场景、可访问性 |
---
### 返回结果
```json
{
"success": true,
"content": "我们计划使用微服务架构重构现有单体应用,预计3个月内完成迁移",
"focus_area": "可行性",
"critique": "## 批判性分析\n\n### 🔴 高风险问题\n1. **时间估算过于乐观**:...\n\n### 🟡 潜在风险\n1. **团队技能缺口**:...\n\n### 🔍 隐含假设\n1. 假设现有代码边界清晰...\n\n### 💡 改进建议\n1. 采用渐进式迁移策略..."
}
```
---
### 返回字段说明
| 字段 | 类型 | 说明 |
| --- | --- | --- |
| success | boolean | 处理状态 |
| content | string | 原始输入内容(截断显示) |
| focus_area | string | 聚焦领域 |
| critique | string | 结构化批判分析报告 |
---
## 统一错误格式
成功:
```json
{
"success": true,
"data": {}
}
```
失败:
```json
{
"success": false,
"error": "错误描述"
}
```
---
## 服务端点
| 端点 | 方法 | 说明 |
| --- | --- | --- |
| / | GET | 服务状态 |
| /health | GET | 健康检查 |
| /mcp | POST | MCP JSON-RPC |
| /mcp/sse | GET/POST | MCP SSE 流式 |
| /api/v1/challenge | POST | 批判性分析 |
---
## 部署信息
| 配置项 | 值 |
| --- | --- |
| 镜像地址 | agnettaiji.azurecr.io/ai-agents/devils-advocate-agent:latest |
| 服务端口 | 8000 |
| 健康检查 | /health |
+14
View File
@@ -0,0 +1,14 @@
# Pydantic AI
pydantic-ai>=0.0.14
# MCP
mcp>=0.9.0
fastmcp>=0.1.0
# FastAPI
fastapi>=0.109.0
uvicorn[standard]>=0.27.0
# HTTP Client
aiohttp>=3.9.0
+18
View File
@@ -0,0 +1,18 @@
#!/usr/bin/env python
"""启动 API 服务器"""
import sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).parent))
if __name__ == '__main__':
from src.server.api_server import app
import uvicorn
import os
host = os.getenv('API_HOST', '0.0.0.0')
port = int(os.getenv('API_PORT', '8000'))
print(f"😈 启动 Devil's Advocate Agent API: http://{host}:{port}")
uvicorn.run(app, host=host, port=port, log_level="info")
+2
View File
@@ -0,0 +1,2 @@
"""Devil's Advocate Agent 源代码包"""
@@ -0,0 +1,2 @@
"""服务器模块"""
@@ -0,0 +1,252 @@
"""
HTTP API 服务器 - Devil's Advocate Agent
提供 REST API 和 MCP HTTP/SSE 端点。
"""
import json
import uuid
import os
from typing import Optional, Dict, Any, AsyncGenerator
from contextlib import asynccontextmanager
from fastapi import FastAPI, HTTPException, Request, Header, Depends
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import StreamingResponse, JSONResponse
from pydantic import BaseModel, Field
from .mcp_server import TOOL_MAP, TOOL_LIST
# ==================== 配置 ====================
SERVER_NAME = "Devil's Advocate Agent API"
# ==================== FastAPI 应用 ====================
@asynccontextmanager
async def lifespan(app: FastAPI):
print(f"😈 {SERVER_NAME} 启动")
yield
print(f"🛑 {SERVER_NAME} 关闭")
app = FastAPI(
title=SERVER_NAME,
description="专门负责挑刺的 AI Agent - 对任何观点、方案、想法提出反对意见和批判性分析",
version="1.0.0",
lifespan=lifespan
)
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
# ==================== API Key 验证 ====================
async def verify_api_key(
api_key: Optional[str] = Header(None, alias="api-key"),
authorization: Optional[str] = Header(None)
) -> str:
"""验证 API Key"""
if api_key and api_key.strip() and api_key.strip() != "sk":
return api_key.strip()
if authorization:
key = authorization[7:].strip() if authorization.startswith("Bearer ") else authorization.strip()
if key and key != "sk":
return key
raise HTTPException(status_code=401, detail="缺少 API Key")
def get_api_key_from_request(request: Request) -> Optional[str]:
"""从请求头提取 API Key(不验证)"""
api_key = request.headers.get("api-key") or request.headers.get("api_key")
if not api_key:
auth = request.headers.get("Authorization")
if auth:
api_key = auth[7:] if auth.startswith("Bearer ") else auth
return api_key
# ==================== 健康检查 ====================
@app.get("/")
async def root():
return {
"service": SERVER_NAME,
"status": "running",
"description": "专门负责挑刺的 AI Agent",
"tools": list(TOOL_MAP.keys())
}
@app.get("/health")
async def health():
return {"status": "healthy", "service": SERVER_NAME}
# ==================== MCP 端点 ====================
sessions: Dict[str, Dict] = {}
async def handle_mcp_request(data: Dict, session_id: str = None, api_key: str = None) -> Dict:
"""处理 MCP JSON-RPC 请求"""
method = data.get("method")
params = data.get("params", {})
req_id = data.get("id")
# tools/call 需要验证 API Key
if method == "tools/call" and (not api_key or api_key == "sk"):
return {"jsonrpc": "2.0", "id": req_id, "error": {"code": -32001, "message": "缺少 API Key"}}
try:
if method == "initialize":
session_id = session_id or str(uuid.uuid4())
sessions[session_id] = {"initialized": True}
return {
"jsonrpc": "2.0", "id": req_id,
"result": {
"protocolVersion": "2024-11-05",
"capabilities": {"tools": {}},
"serverInfo": {"name": SERVER_NAME, "version": "1.0.0"}
}
}
elif method == "tools/list":
return {"jsonrpc": "2.0", "id": req_id, "result": {"tools": TOOL_LIST}}
elif method == "tools/call":
tool_name = params.get("name")
args = params.get("arguments", {})
if tool_name not in TOOL_MAP:
raise ValueError(f"Unknown tool: {tool_name}")
# 设置 API Key 到环境变量
old_key = os.environ.get('OPENAI_API_KEY')
if api_key:
os.environ['OPENAI_API_KEY'] = api_key
try:
result = await TOOL_MAP[tool_name](**args)
finally:
if old_key:
os.environ['OPENAI_API_KEY'] = old_key
return {
"jsonrpc": "2.0", "id": req_id,
"result": {"content": [{"type": "text", "text": str(result)}]}
}
elif method == "ping":
return {"jsonrpc": "2.0", "id": req_id, "result": {}}
else:
raise ValueError(f"Unknown method: {method}")
except Exception as e:
return {"jsonrpc": "2.0", "id": req_id, "error": {"code": -32603, "message": str(e)}}
@app.post("/mcp")
async def mcp_endpoint(request: Request):
"""MCP HTTP 端点"""
try:
body = await request.json()
session_id = request.headers.get("x-mcp-session-id")
api_key = get_api_key_from_request(request)
response = await handle_mcp_request(body, session_id, api_key)
return JSONResponse(content=response, headers={"x-mcp-session-id": session_id or ""})
except Exception as e:
return JSONResponse(status_code=400, content={"jsonrpc": "2.0", "error": {"code": -32700, "message": str(e)}})
@app.get("/mcp/sse")
async def mcp_sse(request: Request):
"""MCP SSE 端点"""
session_id = request.headers.get("x-mcp-session-id") or str(uuid.uuid4())
async def stream() -> AsyncGenerator[str, None]:
yield f"data: {json.dumps({'type': 'connection', 'sessionId': session_id})}\n\n"
import asyncio
while True:
await asyncio.sleep(30)
yield f"data: {json.dumps({'type': 'ping'})}\n\n"
return StreamingResponse(stream(), media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "x-mcp-session-id": session_id})
@app.post("/mcp/sse")
async def mcp_sse_post(request: Request):
"""MCP SSE POST 端点"""
try:
body = await request.json()
session_id = request.headers.get("x-mcp-session-id") or str(uuid.uuid4())
api_key = get_api_key_from_request(request)
async def stream() -> AsyncGenerator[str, None]:
response = await handle_mcp_request(body, session_id, api_key)
yield f"data: {json.dumps(response)}\n\n"
return StreamingResponse(stream(), media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "x-mcp-session-id": session_id})
except Exception as e:
return JSONResponse(status_code=400, content={"jsonrpc": "2.0", "error": {"code": -32700, "message": str(e)}})
# ==================== 业务 API ====================
class ChallengeRequest(BaseModel):
"""挑刺请求模型"""
content: str = Field(..., description="需要被挑刺的内容(观点、方案、想法、代码、文档等)")
focus_area: Optional[str] = Field(None, description="可选,指定关注的领域(如:逻辑漏洞、可行性、安全性、成本、用户体验等)")
class ChallengeResponse(BaseModel):
"""挑刺响应模型"""
success: bool
critique: Optional[str] = None
error: Optional[str] = None
@app.post("/api/v1/challenge", response_model=ChallengeResponse)
async def api_challenge(request: ChallengeRequest, api_key: str = Depends(verify_api_key)):
"""
挑刺 API - 对任何内容进行批判性分析
对观点、方案、想法、代码、文档等进行挑刺,发现问题、漏洞和风险。
"""
try:
# 设置 API Key
old_key = os.environ.get('OPENAI_API_KEY')
os.environ['OPENAI_API_KEY'] = api_key
try:
result_json = await TOOL_MAP['challenge'](
content=request.content,
focus_area=request.focus_area
)
result = json.loads(result_json)
if result.get("success"):
return ChallengeResponse(success=True, critique=result.get("critique"))
else:
return ChallengeResponse(success=False, error=result.get("error"))
finally:
if old_key:
os.environ['OPENAI_API_KEY'] = old_key
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
if __name__ == '__main__':
import uvicorn
uvicorn.run(app, host="0.0.0.0", port=8000)
@@ -0,0 +1,135 @@
"""
MCP 服务器 - Devil's Advocate Agent
专门负责"挑刺"的 AI Agent。
对任何观点、方案、想法提出反对意见、质疑和批判性分析。
"""
import json
import os
from typing import Optional
from mcp.server.fastmcp import FastMCP
from pydantic_ai import Agent
# ==================== 配置 ====================
# LiteLLM Gateway 配置
_BASE_URL = os.getenv('OPENAI_BASE_URL',
os.getenv('LLM_BASE_URL', 'https://litellm.graystone-fb459c5d.southeastasia.azurecontainerapps.io/v1'))
_API_KEY = os.getenv('OPENAI_API_KEY', 'sk')
os.environ.setdefault('OPENAI_API_KEY', _API_KEY)
os.environ.setdefault('OPENAI_BASE_URL', _BASE_URL)
# 模型名称(pydantic_ai 需要 openai: 前缀)
def _get_model_name() -> str:
model = os.getenv('MODEL_NAME', os.getenv('LITELLM_MODEL', 'taiji/gpt-4o-mini'))
return model if ':' in model else f'openai:{model}'
MODEL_NAME = _get_model_name()
# ==================== MCP 服务器 ====================
server = FastMCP("Devil's Advocate Agent")
# 系统提示词 - 专注于挑刺和批判性分析
SYSTEM_PROMPT = '''你是一个专业的"魔鬼代言人"(Devil's Advocate),你的唯一职责是:挑刺。
你的核心任务:
1. 对任何观点、方案、想法提出反对意见和质疑
2. 找出潜在的问题、漏洞、风险和盲点
3. 提出尖锐但有建设性的批判性问题
4. 挑战假设,质疑前提
5. 从不同角度发现可能被忽视的缺陷
你的分析风格:
- 直接、犀利、不留情面
- 但保持专业和理性
- 每个批评都要有依据
- 目的是帮助改进,而非纯粹否定
输出格式要求:
- 使用清晰的结构化格式
- 按严重程度排列问题
- 每个问题都要解释为什么这是个问题
- 最后提供改进建议(可选)
记住:你的存在价值就是帮助发现问题。不要客气,不要敷衍,认真挑刺!'''
def get_agent() -> Agent:
"""创建 Agent 实例(每次调用使用最新的 API Key)"""
return Agent(MODEL_NAME, system_prompt=SYSTEM_PROMPT)
# ==================== MCP 工具定义 ====================
@server.tool()
async def challenge(
content: str,
focus_area: Optional[str] = None
) -> str:
"""
对任何观点、方案、想法进行挑刺和批判性分析
Args:
content: 需要被挑刺的内容(观点、方案、想法、代码、文档等)
focus_area: 可选,指定关注的领域(如:逻辑漏洞、可行性、安全性、成本、用户体验等)
Returns:
批判性分析结果(JSON 格式)
"""
try:
# 构建提示
prompt = f"请对以下内容进行挑刺和批判性分析:\n\n{content}"
if focus_area:
prompt += f"\n\n请特别关注:{focus_area}"
# 调用 AI Agent 进行批判性分析
result = await get_agent().run(prompt)
return json.dumps({
"success": True,
"content": content[:200] + "..." if len(content) > 200 else content,
"focus_area": focus_area,
"critique": result.output
}, ensure_ascii=False, indent=2)
except Exception as e:
return json.dumps({
"success": False,
"error": str(e)
}, ensure_ascii=False)
# ==================== 工具映射(供 API 使用)====================
TOOL_MAP = {
'challenge': challenge,
}
TOOL_LIST = [
{
"name": "challenge",
"description": "对任何观点、方案、想法进行挑刺和批判性分析。专门负责发现问题、漏洞和风险。",
"inputSchema": {
"type": "object",
"properties": {
"content": {
"type": "string",
"description": "需要被挑刺的内容(观点、方案、想法、代码、文档等)"
},
"focus_area": {
"type": "string",
"description": "可选,指定关注的领域(如:逻辑漏洞、可行性、安全性、成本、用户体验等)"
}
},
"required": ["content"]
}
}
]
if __name__ == '__main__':
server.run()
+21
View File
@@ -0,0 +1,21 @@
FROM python:3.12-slim
WORKDIR /app
ENV PYTHONUNBUFFERED=1
ENV PYTHONDONTWRITEBYTECODE=1
RUN apt-get update && apt-get install -y gcc curl && rm -rf /var/lib/apt/lists/*
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
EXPOSE 8000
HEALTHCHECK --interval=30s --timeout=10s --start-period=5s --retries=3 \
CMD curl -f http://localhost:8000/health || exit 1
CMD ["python", "run_api_server.py"]
+213
View File
@@ -0,0 +1,213 @@
# 格式警察 Agent (Format Police Agent)
🚔 **只做一件事**:检查输出是否符合 JSON 格式,如果不符合则补全格式化 JSON 输出。
## 功能特性
- ✅ **JSON 验证**:检查内容是否为有效 JSON
- 🔧 **自动修复**:尝试修复常见的 JSON 格式问题
- 🤖 **AI 辅助**:无法自动修复时使用 AI 转换
- 📦 **格式化输出**:美化 JSON 输出
## 快速开始
### 1. 本地运行
```bash
# 安装依赖
pip install -r requirements.txt
# 启动服务
python run_api_server.py
```
### 2. Docker 运行
```bash
# 构建镜像
docker build -t format-police-agent:latest .
# 运行容器
docker run -p 8000:8000 -e OPENAI_API_KEY=your_key format-police-agent:latest
```
## API 端点
### 健康检查
```bash
curl http://localhost:8000/health
```
### MCP 端点
| 端点 | 方法 | 说明 |
|------|------|------|
| `/mcp` | POST | MCP HTTP 端点 |
| `/mcp/sse` | GET/POST | MCP SSE 端点 |
### REST API
| 端点 | 方法 | 说明 | 需要 API Key |
|------|------|------|-------------|
| `/api/v1/check` | POST | 检查并修复 JSON | 是(使用 AI 时)|
| `/api/v1/validate` | POST | 仅验证 JSON | 否 |
| `/api/v1/format` | POST | 格式化 JSON | 否 |
## MCP 工具
### check_and_fix_json
检查并修复 JSON 格式。
```json
{
"name": "check_and_fix_json",
"arguments": {
"content": "需要检查的内容",
"use_ai": true
}
}
```
**参数:**
- `content` (必需): 需要检查和修复的内容
- `use_ai` (可选): 是否使用 AI 辅助修复,默认 `true`
### validate_json
仅验证 JSON 格式是否有效。
```json
{
"name": "validate_json",
"arguments": {
"content": "需要验证的内容"
}
}
```
### format_json
格式化已有效的 JSON。
```json
{
"name": "format_json",
"arguments": {
"content": "{\"a\":1}",
"indent": 2
}
}
```
## 使用示例
### 1. 检查并修复 JSON
```bash
curl -X POST http://localhost:8000/api/v1/check \
-H "Content-Type: application/json" \
-H "api-key: your_api_key" \
-d '{"content": "{name: \"test\", value: 123}"}'
```
响应:
```json
{
"success": true,
"result": {
"success": true,
"is_original_valid": false,
"message": "已自动修复 JSON 格式问题",
"formatted_json": {
"name": "test",
"value": 123
}
}
}
```
### 2. 验证 JSON(无需 API Key)
```bash
curl -X POST http://localhost:8000/api/v1/validate \
-H "Content-Type: application/json" \
-d '{"content": "{\"valid\": true}"}'
```
### 3. MCP 调用
```bash
curl -X POST http://localhost:8000/mcp \
-H "Content-Type: application/json" \
-H "api-key: your_api_key" \
-d '{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "check_and_fix_json",
"arguments": {
"content": "这不是JSON,但包含信息:名字是张三,年龄25"
}
}
}'
```
## 自动修复能力
格式警察可以自动修复以下问题:
1. **单引号** → 双引号
2. **末尾多余逗号** → 移除
3. **未加引号的值** → 添加引号
4. **Markdown 代码块** → 提取 JSON
5. **BOM 字符** → 移除
无法自动修复时,会使用 AI 尝试转换。
## 环境变量
| 变量 | 必需 | 说明 | 默认值 |
|------|------|------|--------|
| OPENAI_API_KEY | 使用 AI 时 | OpenAI API Key | - |
| OPENAI_BASE_URL | 否 | LLM API 地址 | LiteLLM Gateway |
| MODEL_NAME | 否 | 模型名称 | taiji/gpt-4o-mini |
| API_HOST | 否 | 监听地址 | 0.0.0.0 |
| API_PORT | 否 | 监听端口 | 8000 |
## 注册到 Agent Manager
在 `k8s_manager.py` 中添加:
```python
# TEMPLATE_PORTS
"format_police": 8000,
# image_map
"format_police": "agnettaiji.azurecr.io/ai-agents/format-police-agent:latest",
```
在 `app.py` 的 `valid_templates` 中添加 `"format_police"`。
## 项目结构
```
format_police_agent/
├── Dockerfile
├── requirements.txt
├── run_api_server.py
├── README.md
└── src/
├── __init__.py
└── server/
├── __init__.py
├── api_server.py # FastAPI + MCP HTTP
└── mcp_server.py # MCP 工具定义
```
## License
MIT
+317
View File
@@ -0,0 +1,317 @@
# 格式警察 Agent
---
本 Agent 提供基于 **智能解析** 的 JSON 格式校验与修复能力,通过 **HTTP API** 与 **MCP(Model Context Protocol)** 对外提供服务。
核心能力:
- **格式校验**:符合 RFC 8259 标准的 JSON 语法检测
- **智能修复**:基于规则引擎与 LLM 的多级修复策略
- **格式美化**:可配置缩进的结构化输出
---
## 功能概览
提供 JSON 文本的 **语法验证、智能修复、格式化输出** 能力,返回结构化处理结果。
支持能力:
- JSON Schema 合规性校验
- 常见语法错误自动修复
- LLM 辅助的非结构化文本转换
- 可配置的格式化输出
---
## 1⃣ check_and_fix_json — 智能校验与修复
### 功能说明
对输入内容执行 **多级修复策略**:优先通过规则引擎修复常见语法问题,失败时启用 LLM 进行语义转换,确保输出符合 JSON 规范。
---
### REST API 调用
```
POST /api/v1/check
Content-Type: application/json
api-key: {your-api-key}
```
```json
{
"content": "{name: 'test', value: 123,}",
"use_ai": true
}
```
---
### MCP 调用
```json
{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "check_and_fix_json",
"arguments": {
"content": "{name: 'test', value: 123,}",
"use_ai": true
}
}
}
```
---
### 参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
| --- | --- | --- | --- | --- |
| content | string | ✅ | - | 待校验或修复的文本内容 |
| use_ai | boolean | ❌ | true | 是否启用 LLM 辅助修复 |
---
### 修复策略
| 阶段 | 策略 | 说明 |
| --- | --- | --- |
| L1 | 直接解析 | 验证是否为合法 JSON |
| L2 | 模式提取 | 从代码块 / 嵌套文本中提取 JSON |
| L3 | 规则修复 | 修复引号、逗号等常见语法问题 |
| L4 | LLM 转换 | 调用大模型将非结构化文本转为 JSON |
---
### 返回结果
```json
{
"success": true,
"result": {
"success": true,
"is_original_valid": false,
"message": "已自动修复 JSON 格式问题",
"formatted_json": {
"name": "test",
"value": 123
}
}
}
```
---
### 返回字段说明
| 字段 | 类型 | 说明 |
| --- | --- | --- |
| success | boolean | 处理是否成功 |
| is_original_valid | boolean | 原始输入是否为合法 JSON |
| message | string | 处理结果描述 |
| formatted_json | object/array | 修复后的 JSON 对象 |
---
## 2️⃣ validate_json — 格式校验
### 功能说明
执行 **只读校验**,验证输入是否符合 JSON 语法规范,返回详细的错误定位信息,不进行任何修改。
---
### REST API 调用
```
POST /api/v1/validate
Content-Type: application/json
```
```json
{
"content": "{\"valid\": true, \"count\": 42}"
}
```
---
### MCP 调用
```json
{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "validate_json",
"arguments": {
"content": "{\"valid\": true}"
}
}
}
```
---
### 参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
| --- | --- | --- | --- | --- |
| content | string | ✅ | - | 待校验的文本内容 |
---
### 返回结果(合法)
```json
{
"success": true,
"result": {
"valid": true,
"message": "有效的 JSON 格式",
"json_type": "dict",
"preview": "{'valid': True, 'count': 42}"
}
}
```
### 返回结果(非法)
```json
{
"success": true,
"result": {
"valid": false,
"message": "无效的 JSON 格式",
"error_detail": "位置 1: Expecting property name enclosed in double quotes",
"suggestion": "可以使用 check_and_fix_json 工具尝试修复"
}
}
```
---
## 3️⃣ format_json — 格式美化
### 功能说明
对合法 JSON 执行 **结构化美化输出**,支持自定义缩进层级,便于阅读与调试。
---
### REST API 调用
```
POST /api/v1/format
Content-Type: application/json
```
```json
{
"content": "{\"a\":1,\"b\":{\"c\":2}}",
"indent": 4
}
```
---
### MCP 调用
```json
{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "format_json",
"arguments": {
"content": "{\"a\":1,\"b\":{\"c\":2}}",
"indent": 2
}
}
}
```
---
### 参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
| --- | --- | --- | --- | --- |
| content | string | ✅ | - | 合法的 JSON 字符串 |
| indent | integer | ❌ | 2 | 缩进空格数(1-8) |
---
### 返回结果
```json
{
"success": true,
"result": {
"success": true,
"formatted_json": {
"a": 1,
"b": {
"c": 2
}
},
"formatted_string": "{\n \"a\": 1,\n \"b\": {\n \"c\": 2\n }\n}"
}
}
```
---
## 统一错误格式
成功:
```json
{
"success": true,
"result": {}
}
```
失败:
```json
{
"success": false,
"error": "错误描述"
}
```
---
## 服务端点
| 端点 | 方法 | 说明 |
| --- | --- | --- |
| / | GET | 服务状态 |
| /health | GET | 健康检查 |
| /mcp | POST | MCP JSON-RPC |
| /mcp/sse | GET/POST | MCP SSE 流式 |
| /api/v1/check | POST | 智能校验与修复 |
| /api/v1/validate | POST | 格式校验 |
| /api/v1/format | POST | 格式美化 |
---
## 部署信息
| 配置项 | 值 |
| --- | --- |
| 镜像地址 | agnettaiji.azurecr.io/ai-agents/format-police-agent:latest |
| 服务端口 | 8000 |
| 健康检查 | /health |
+14
View File
@@ -0,0 +1,14 @@
# Pydantic AI
pydantic-ai>=0.0.14
# MCP
mcp>=0.9.0
fastmcp>=0.1.0
# FastAPI
fastapi>=0.109.0
uvicorn[standard]>=0.27.0
# HTTP Client
aiohttp>=3.9.0
+19
View File
@@ -0,0 +1,19 @@
#!/usr/bin/env python
"""启动格式警察 Agent API 服务器"""
import sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).parent))
if __name__ == '__main__':
from src.server.api_server import app
import uvicorn
import os
host = os.getenv('API_HOST', '0.0.0.0')
port = int(os.getenv('API_PORT', '8000'))
print(f"🚔 启动格式警察 Agent API: http://{host}:{port}")
print(f"📋 MCP 端点: http://{host}:{port}/mcp")
uvicorn.run(app, host=host, port=port, log_level="info")
+2
View File
@@ -0,0 +1,2 @@
"""格式警察 Agent 源代码包"""
@@ -0,0 +1,2 @@
"""服务器模块"""
@@ -0,0 +1,282 @@
"""
HTTP API 服务器 - 格式警察 Agent
提供 REST API 和 MCP HTTP/SSE 端点。
"""
import json
import uuid
import os
from typing import Optional, Dict, Any, AsyncGenerator
from contextlib import asynccontextmanager
from fastapi import FastAPI, HTTPException, Request, Header, Depends
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import StreamingResponse, JSONResponse
from pydantic import BaseModel, Field
from .mcp_server import TOOL_MAP, TOOL_LIST
# ==================== 配置 ====================
SERVER_NAME = "Format Police Agent API"
# ==================== FastAPI 应用 ====================
@asynccontextmanager
async def lifespan(app: FastAPI):
print(f"🚀 {SERVER_NAME} 启动")
print(f"📋 可用工具: {list(TOOL_MAP.keys())}")
yield
print(f"🛑 {SERVER_NAME} 关闭")
app = FastAPI(
title=SERVER_NAME,
description="格式警察 Agent - 检查并修复 JSON 格式",
version="1.0.0",
lifespan=lifespan
)
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
# ==================== API Key 验证 ====================
async def verify_api_key(
api_key: Optional[str] = Header(None, alias="api-key"),
authorization: Optional[str] = Header(None)
) -> str:
"""验证 API Key"""
if api_key and api_key.strip() and api_key.strip() != "sk":
return api_key.strip()
if authorization:
key = authorization[7:].strip() if authorization.startswith("Bearer ") else authorization.strip()
if key and key != "sk":
return key
raise HTTPException(status_code=401, detail="缺少 API Key")
def get_api_key_from_request(request: Request) -> Optional[str]:
"""从请求头提取 API Key(不验证)"""
api_key = request.headers.get("api-key") or request.headers.get("api_key")
if not api_key:
auth = request.headers.get("Authorization")
if auth:
api_key = auth[7:] if auth.startswith("Bearer ") else auth
return api_key
# ==================== 健康检查 ====================
@app.get("/")
async def root():
return {
"service": SERVER_NAME,
"description": "格式警察 Agent - 检查并修复 JSON 格式",
"status": "running",
"tools": list(TOOL_MAP.keys())
}
@app.get("/health")
async def health():
return {"status": "healthy", "service": SERVER_NAME}
# ==================== MCP 端点 ====================
sessions: Dict[str, Dict] = {}
async def handle_mcp_request(data: Dict, session_id: str = None, api_key: str = None) -> Dict:
"""处理 MCP JSON-RPC 请求"""
method = data.get("method")
params = data.get("params", {})
req_id = data.get("id")
# tools/call 中使用 AI 修复时需要 API Key
if method == "tools/call":
tool_name = params.get("name")
args = params.get("arguments", {})
# 如果使用 AI 修复且没有 API Key
if tool_name == "check_and_fix_json" and args.get("use_ai", True) and (not api_key or api_key == "sk"):
return {"jsonrpc": "2.0", "id": req_id, "error": {"code": -32001, "message": "使用 AI 修复需要提供 API Key"}}
try:
if method == "initialize":
session_id = session_id or str(uuid.uuid4())
sessions[session_id] = {"initialized": True}
return {
"jsonrpc": "2.0", "id": req_id,
"result": {
"protocolVersion": "2024-11-05",
"capabilities": {"tools": {}},
"serverInfo": {"name": SERVER_NAME, "version": "1.0.0"}
}
}
elif method == "tools/list":
return {"jsonrpc": "2.0", "id": req_id, "result": {"tools": TOOL_LIST}}
elif method == "tools/call":
tool_name = params.get("name")
args = params.get("arguments", {})
if tool_name not in TOOL_MAP:
raise ValueError(f"Unknown tool: {tool_name}")
# 设置 API Key 到环境变量
old_key = os.environ.get('OPENAI_API_KEY')
if api_key:
os.environ['OPENAI_API_KEY'] = api_key
try:
result = await TOOL_MAP[tool_name](**args)
finally:
if old_key:
os.environ['OPENAI_API_KEY'] = old_key
return {
"jsonrpc": "2.0", "id": req_id,
"result": {"content": [{"type": "text", "text": str(result)}]}
}
elif method == "ping":
return {"jsonrpc": "2.0", "id": req_id, "result": {}}
else:
raise ValueError(f"Unknown method: {method}")
except Exception as e:
return {"jsonrpc": "2.0", "id": req_id, "error": {"code": -32603, "message": str(e)}}
@app.post("/mcp")
async def mcp_endpoint(request: Request):
"""MCP HTTP 端点"""
try:
body = await request.json()
session_id = request.headers.get("x-mcp-session-id")
api_key = get_api_key_from_request(request)
response = await handle_mcp_request(body, session_id, api_key)
return JSONResponse(content=response, headers={"x-mcp-session-id": session_id or ""})
except Exception as e:
return JSONResponse(status_code=400, content={"jsonrpc": "2.0", "error": {"code": -32700, "message": str(e)}})
@app.get("/mcp/sse")
async def mcp_sse(request: Request):
"""MCP SSE 端点"""
session_id = request.headers.get("x-mcp-session-id") or str(uuid.uuid4())
async def stream() -> AsyncGenerator[str, None]:
yield f"data: {json.dumps({'type': 'connection', 'sessionId': session_id})}\n\n"
import asyncio
while True:
await asyncio.sleep(30)
yield f"data: {json.dumps({'type': 'ping'})}\n\n"
return StreamingResponse(stream(), media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "x-mcp-session-id": session_id})
@app.post("/mcp/sse")
async def mcp_sse_post(request: Request):
"""MCP SSE POST 端点"""
try:
body = await request.json()
session_id = request.headers.get("x-mcp-session-id") or str(uuid.uuid4())
api_key = get_api_key_from_request(request)
async def stream() -> AsyncGenerator[str, None]:
response = await handle_mcp_request(body, session_id, api_key)
yield f"data: {json.dumps(response)}\n\n"
return StreamingResponse(stream(), media_type="text/event-stream",
headers={"Cache-Control": "no-cache", "x-mcp-session-id": session_id})
except Exception as e:
return JSONResponse(status_code=400, content={"jsonrpc": "2.0", "error": {"code": -32700, "message": str(e)}})
# ==================== 业务 API ====================
class CheckJsonRequest(BaseModel):
"""检查 JSON 请求"""
content: str = Field(..., description="需要检查的内容")
use_ai: Optional[bool] = Field(True, description="是否使用 AI 辅助修复")
class ValidateJsonRequest(BaseModel):
"""验证 JSON 请求"""
content: str = Field(..., description="需要验证的内容")
class FormatJsonRequest(BaseModel):
"""格式化 JSON 请求"""
content: str = Field(..., description="有效的 JSON 字符串")
indent: Optional[int] = Field(2, description="缩进空格数")
class JsonResponse(BaseModel):
"""JSON 响应"""
success: bool
result: Optional[Dict[str, Any]] = None
error: Optional[str] = None
@app.post("/api/v1/check", response_model=JsonResponse)
async def api_check_json(request: CheckJsonRequest, api_key: str = Depends(verify_api_key)):
"""检查并修复 JSON 格式"""
try:
# 设置 API Key
old_key = os.environ.get('OPENAI_API_KEY')
os.environ['OPENAI_API_KEY'] = api_key
try:
result = await TOOL_MAP['check_and_fix_json'](
content=request.content,
use_ai=request.use_ai
)
return JsonResponse(success=True, result=json.loads(result))
finally:
if old_key:
os.environ['OPENAI_API_KEY'] = old_key
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
@app.post("/api/v1/validate", response_model=JsonResponse)
async def api_validate_json(request: ValidateJsonRequest):
"""验证 JSON 格式(不需要 API Key)"""
try:
result = await TOOL_MAP['validate_json'](content=request.content)
return JsonResponse(success=True, result=json.loads(result))
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
@app.post("/api/v1/format", response_model=JsonResponse)
async def api_format_json(request: FormatJsonRequest):
"""格式化 JSON(不需要 API Key)"""
try:
result = await TOOL_MAP['format_json'](
content=request.content,
indent=request.indent
)
return JsonResponse(success=True, result=json.loads(result))
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
if __name__ == '__main__':
import uvicorn
uvicorn.run(app, host="0.0.0.0", port=8000)
@@ -0,0 +1,351 @@
"""
MCP 服务器 - 格式警察 Agent
功能:检查输出是否符合 JSON 格式,如果不符合则补全格式化 JSON 输出。
"""
import json
import os
import re
from typing import Optional
from mcp.server.fastmcp import FastMCP
from pydantic_ai import Agent
# ==================== 配置 ====================
# LiteLLM Gateway 配置
_BASE_URL = os.getenv('OPENAI_BASE_URL',
os.getenv('LLM_BASE_URL', 'https://litellm.graystone-fb459c5d.southeastasia.azurecontainerapps.io/v1'))
_API_KEY = os.getenv('OPENAI_API_KEY', 'sk')
os.environ.setdefault('OPENAI_API_KEY', _API_KEY)
os.environ.setdefault('OPENAI_BASE_URL', _BASE_URL)
# 模型名称(pydantic_ai 需要 openai: 前缀)
def _get_model_name() -> str:
model = os.getenv('MODEL_NAME', os.getenv('LITELLM_MODEL', 'taiji/gpt-4o-mini'))
return model if ':' in model else f'openai:{model}'
MODEL_NAME = _get_model_name()
# ==================== MCP 服务器 ====================
server = FastMCP('Format Police Agent')
# 系统提示词
SYSTEM_PROMPT = '''你是一个 JSON 格式修复专家。你的任务是:
1. 分析输入内容,判断它是否是有效的 JSON
2. 如果是有效 JSON,直接返回格式化后的 JSON
3. 如果不是有效 JSON,尝试从中提取信息并转换为有效的 JSON 格式
输出规则:
- 只输出 JSON,不要添加任何解释或其他文字
- 使用 2 空格缩进
- 确保输出是有效的 JSON 格式
- 如果无法转换,返回包含原始内容的 JSON 对象'''
def get_agent() -> Agent:
"""创建 Agent 实例(每次调用使用最新的 API Key)"""
return Agent(MODEL_NAME, system_prompt=SYSTEM_PROMPT)
# ==================== 辅助函数 ====================
def is_valid_json(text: str) -> tuple[bool, Optional[dict | list]]:
"""检查文本是否是有效的 JSON"""
try:
parsed = json.loads(text.strip())
return True, parsed
except json.JSONDecodeError:
return False, None
def try_extract_json(text: str) -> Optional[dict | list]:
"""尝试从文本中提取 JSON"""
# 尝试提取 {} 或 [] 包裹的内容
patterns = [
r'```json\s*([\s\S]*?)\s*```', # Markdown JSON 代码块
r'```\s*([\s\S]*?)\s*```', # 普通代码块
r'(\{[\s\S]*\})', # JSON 对象
r'(\[[\s\S]*\])', # JSON 数组
]
for pattern in patterns:
matches = re.findall(pattern, text)
for match in matches:
is_valid, parsed = is_valid_json(match)
if is_valid:
return parsed
return None
def fix_common_json_issues(text: str) -> str:
"""修复常见的 JSON 格式问题"""
fixed = text.strip()
# 移除可能的 BOM
if fixed.startswith('\ufeff'):
fixed = fixed[1:]
# 修复单引号为双引号
# 注意:这是简单替换,可能不完美
fixed = re.sub(r"'([^']*)':", r'"\1":', fixed)
fixed = re.sub(r":\s*'([^']*)'", r': "\1"', fixed)
# 修复末尾多余的逗号
fixed = re.sub(r',(\s*[\]}])', r'\1', fixed)
# 修复缺失的引号(简单情况)
fixed = re.sub(r':\s*([a-zA-Z_][a-zA-Z0-9_]*)\s*([,}\]])', r': "\1"\2', fixed)
return fixed
# ==================== MCP 工具定义 ====================
@server.tool()
async def check_and_fix_json(
content: str,
use_ai: Optional[bool] = True
) -> str:
"""
检查并修复 JSON 格式
Args:
content: 需要检查的内容
use_ai: 是否使用 AI 辅助修复(默认 True)
Returns:
格式化的 JSON 结果
"""
try:
# 1. 首先检查是否已经是有效 JSON
is_valid, parsed = is_valid_json(content)
if is_valid:
return json.dumps({
"success": True,
"is_original_valid": True,
"message": "输入已经是有效的 JSON 格式",
"formatted_json": json.loads(json.dumps(parsed, ensure_ascii=False, indent=2))
}, ensure_ascii=False, indent=2)
# 2. 尝试从内容中提取 JSON
extracted = try_extract_json(content)
if extracted:
return json.dumps({
"success": True,
"is_original_valid": False,
"message": "从内容中提取到有效的 JSON",
"formatted_json": extracted
}, ensure_ascii=False, indent=2)
# 3. 尝试修复常见问题
fixed_content = fix_common_json_issues(content)
is_valid, parsed = is_valid_json(fixed_content)
if is_valid:
return json.dumps({
"success": True,
"is_original_valid": False,
"message": "已自动修复 JSON 格式问题",
"formatted_json": parsed
}, ensure_ascii=False, indent=2)
# 4. 使用 AI 辅助修复
if use_ai:
prompt = f"""请将以下内容转换为有效的 JSON 格式。
只输出 JSON,不要任何解释:
{content}"""
result = await get_agent().run(prompt)
ai_output = result.output.strip()
# 尝试解析 AI 输出
extracted_from_ai = try_extract_json(ai_output)
if extracted_from_ai:
return json.dumps({
"success": True,
"is_original_valid": False,
"message": "已使用 AI 转换为 JSON 格式",
"formatted_json": extracted_from_ai
}, ensure_ascii=False, indent=2)
is_valid, parsed = is_valid_json(ai_output)
if is_valid:
return json.dumps({
"success": True,
"is_original_valid": False,
"message": "已使用 AI 转换为 JSON 格式",
"formatted_json": parsed
}, ensure_ascii=False, indent=2)
# 5. 无法修复,返回包装后的原始内容
return json.dumps({
"success": False,
"is_original_valid": False,
"message": "无法转换为有效的 JSON 格式,已包装原始内容",
"formatted_json": {
"raw_content": content,
"type": "unconvertible"
}
}, ensure_ascii=False, indent=2)
except Exception as e:
return json.dumps({
"success": False,
"error": str(e),
"formatted_json": {
"raw_content": content,
"type": "error"
}
}, ensure_ascii=False, indent=2)
@server.tool()
async def validate_json(content: str) -> str:
"""
仅验证 JSON 格式是否有效(不修复)
Args:
content: 需要验证的内容
Returns:
验证结果(JSON 格式)
"""
try:
is_valid, parsed = is_valid_json(content)
if is_valid:
return json.dumps({
"valid": True,
"message": "有效的 JSON 格式",
"json_type": type(parsed).__name__,
"preview": str(parsed)[:200] + "..." if len(str(parsed)) > 200 else str(parsed)
}, ensure_ascii=False, indent=2)
else:
# 尝试获取具体的错误信息
try:
json.loads(content)
except json.JSONDecodeError as e:
error_detail = f"位置 {e.pos}: {e.msg}"
return json.dumps({
"valid": False,
"message": "无效的 JSON 格式",
"error_detail": error_detail,
"suggestion": "可以使用 check_and_fix_json 工具尝试修复"
}, ensure_ascii=False, indent=2)
except Exception as e:
return json.dumps({
"valid": False,
"error": str(e)
}, ensure_ascii=False, indent=2)
@server.tool()
async def format_json(content: str, indent: Optional[int] = 2) -> str:
"""
格式化已有效的 JSON(美化输出)
Args:
content: 有效的 JSON 字符串
indent: 缩进空格数(默认 2)
Returns:
格式化后的 JSON
"""
try:
is_valid, parsed = is_valid_json(content)
if not is_valid:
return json.dumps({
"success": False,
"message": "输入不是有效的 JSON,无法格式化",
"suggestion": "请先使用 check_and_fix_json 工具修复"
}, ensure_ascii=False, indent=2)
formatted = json.dumps(parsed, ensure_ascii=False, indent=indent)
return json.dumps({
"success": True,
"formatted_json": parsed,
"formatted_string": formatted
}, ensure_ascii=False, indent=indent)
except Exception as e:
return json.dumps({
"success": False,
"error": str(e)
}, ensure_ascii=False, indent=2)
# ==================== 工具映射(供 API 使用)====================
TOOL_MAP = {
'check_and_fix_json': check_and_fix_json,
'validate_json': validate_json,
'format_json': format_json,
}
TOOL_LIST = [
{
"name": "check_and_fix_json",
"description": "检查并修复 JSON 格式。如果输入是有效 JSON 则格式化输出,如果无效则尝试修复或使用 AI 转换。",
"inputSchema": {
"type": "object",
"properties": {
"content": {
"type": "string",
"description": "需要检查和修复的内容"
},
"use_ai": {
"type": "boolean",
"description": "是否使用 AI 辅助修复(默认 True)"
}
},
"required": ["content"]
}
},
{
"name": "validate_json",
"description": "仅验证 JSON 格式是否有效,不进行修复",
"inputSchema": {
"type": "object",
"properties": {
"content": {
"type": "string",
"description": "需要验证的内容"
}
},
"required": ["content"]
}
},
{
"name": "format_json",
"description": "格式化已有效的 JSON,美化输出",
"inputSchema": {
"type": "object",
"properties": {
"content": {
"type": "string",
"description": "有效的 JSON 字符串"
},
"indent": {
"type": "integer",
"description": "缩进空格数(默认 2)"
}
},
"required": ["content"]
}
}
]
if __name__ == '__main__':
server.run()