更新智能体

This commit is contained in:
zhanggangyong
2026-02-21 03:00:16 +00:00
parent 3a9773febd
commit 387c80b901
18 changed files with 2337 additions and 0 deletions
+22
View File
@@ -0,0 +1,22 @@
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"]
+168
View File
@@ -0,0 +1,168 @@
# Agent 职责越界检测 Agent 🛡️
**职责边界保护 Agent** - 专门负责检测 Agent 是否超出职责范围的 AI Agent。
## 功能特点
只做一件事:**判断一个 agent 是否开始"多管闲事"**
- 🎯 检测 Agent 是否超出职责范围
- 🔍 分析 Agent 的职责定义与实际行为
- ❌ 识别 Agent 是否在处理不属于自己的任务
- 🚫 判断 Agent 是否在"多管闲事"或越界操作
- 🛡️ 为多 agent 系统提供职责边界保护
## 快速开始
### 本地测试
```bash
cd agent_boundary_check_agent
pip install -r requirements.txt
python run_api_server.py
```
### 构建镜像
```bash
docker build -t agent-boundary-check-agent:latest .
```
## API 端点
### REST API
```bash
# 职责越界检测 API
curl -X POST http://localhost:8000/api/v1/check \
-H "Content-Type: application/json" \
-H "api-key: your-api-key" \
-d '{
"agent_name": "数据分析 Agent",
"agent_defined_duty": "负责数据分析和报表生成",
"agent_actual_behavior": "我正在帮用户写代码,因为数据分析需要一些代码支持",
"current_task": "生成数据分析报告"
}'
```
响应示例:
```json
{
"success": true,
"is_violation": true,
"analysis": "该 Agent 明显超出了职责范围。数据分析 Agent 的职责是数据分析和报表生成,但实际行为是在写代码,这属于代码开发 Agent 的职责范围。虽然数据分析可能需要代码支持,但直接帮用户写代码已经超出了边界..."
}
```
### MCP 端点
| 端点 | 方法 | 说明 |
|------|------|------|
| `/mcp` | POST | MCP HTTP 端点 |
| `/mcp/sse` | GET | MCP SSE 连接 |
| `/mcp/sse` | POST | MCP SSE 请求 |
### MCP 工具
#### `check_boundary_violation` - 职责越界检测工具
检测 Agent 是否超出职责范围(是否"多管闲事")。
**参数:**
| 参数 | 类型 | 必需 | 说明 |
|------|------|------|------|
| `agent_name` | string | 是 | Agent 的名称或标识 |
| `agent_defined_duty` | string | 是 | Agent 定义的职责范围(应该做什么) |
| `agent_actual_behavior` | string | 是 | Agent 的实际行为或输出(实际做了什么) |
| `current_task` | string | 否 | 当前正在处理的任务描述 |
**MCP 调用示例:**
```json
{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "check_boundary_violation",
"arguments": {
"agent_name": "文档整理 Agent",
"agent_defined_duty": "负责文档的整理、分类和归档",
"agent_actual_behavior": "我正在帮用户修改代码中的 bug,因为文档中提到了一些代码问题",
"current_task": "整理项目文档"
}
}
}
```
## MCP 配置示例
在你的 MCP 客户端配置中添加:
```json
{
"mcpServers": {
"agent-boundary-check": {
"url": "http://localhost:8000/mcp/sse",
"transport": "sse"
}
}
}
```
## 项目结构
```
agent_boundary_check_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
"agent_boundary_check": 8000,
# image_map
"agent_boundary_check": "agnettaiji.azurecr.io/ai-agents/agent-boundary-check-agent:latest",
```
在 `app.py` 的 `valid_templates` 中添加 `"agent_boundary_check"`。
## 使用场景
- 🏢 **多 Agent 系统**:确保每个 Agent 专注于自己的职责
- 🔄 **工作流管理**:在 Agent 协作流程中检测职责混乱
- 🎯 **职责审计**:定期检查 Agent 是否偏离原始职责
- 🛡️ **系统保护**:防止 Agent 越界操作导致系统混乱
- 📊 **职责监控**:实时监控 Agent 行为是否符合职责定义
## 检测维度
- **职责匹配度**:Agent 的行为是否与定义的职责一致
- **任务相关性**:当前任务是否属于该 Agent 的职责范围
- **越界行为**:是否在处理其他 Agent 应该处理的任务
- **职责边界**:是否清晰地区分了"该做"和"不该做"的事情
+189
View File
@@ -0,0 +1,189 @@
# Agent Boundary(模糊输入补齐智能体)
简述:将模糊、不完整或口语化的人类输入补齐为可执行、结构化的意图与参数表示,适合作为任何上游 Agent 的预处理层,提升下游任务解析与执行的准确性与鲁棒性。
## 功能概览
- 补齐与规范化:将不完整/口语化输入拓展为具备明确意图与槽位的可执行对象。
- 实体消歧与补全:基于上下文与知识库消除歧义并补全缺失槽位。
- 语义验证与约束应用:基于业务 schema/类型约束验证并格式化输出。
- 候选生成与置信度评估:为关键补全项提供候选集与置信度,支持人机确认或自动选择策略。
- 异步与可追踪:支持异步长时任务并提供查询接口。
---
## API: complete_input — 输入补齐与规范化
功能说明:对原始文本进行语义解析、上下文回溯与推断,输出结构化意图对象(含槽位、类型、来源与置信度)。
REST API 调用:
```
POST /api/v1/complete
Content-Type: application/json
```
请求示例:
```
{
"raw_text": "帮我订明天从北上到上海的票",
"context": {"user_id":"u123","history":[...]},
"schema": {"intent":"book_ticket","slots":["from","to","date"]},
"strategy": "conservative"
}
```
MCP 调用示例:
```
{
"jsonrpc":"2.0",
"id":1,
"method":"tools/call",
"params":{
"name":"complete_input",
"arguments":{...}
}
}
```
参数说明:
| 参数 | 类型 | 必需 | 默认值 | 说明 |
| --- | --- | ---: | --- | --- |
| raw_text | string | ✅ | - | 用户原始文本 |
| context | object | ❌ | null | 对话历史/外部上下文,用于消歧 |
| schema | object | ❌ | null | 业务意图与槽位约束(含类型/正则) |
| strategy | string | ❌ | "balanced" | 补齐策略(conservative/balanced/aggressive) |
| max_candidates | integer | ❌ | 3 | 返回候选数量 |
返回结果(示例):
```
{
"success": true,
"completed": {
"text": "为您预订 2026-02-22 从 北京 到 上海 的火车票",
"intent": "book_ticket",
"slots": {
"from": {"value":"北京","source":"inference","confidence":0.87},
"to": {"value":"上海","source":"explicit","confidence":0.99},
"date": {"value":"2026-02-22","source":"resolve_date","confidence":0.95}
},
"metadata": {"strategy":"conservative","explain":"推断‘北上’为北京"}
},
"candidates": [ ... ]
}
```
---
## API: resolve_entities — 实体消歧与标准化
功能说明:对识别的实体执行消歧、标准化(地名/时间/单位)与外部查证,返回标准标识与来源证明。
REST API 调用:
```
POST /api/v1/resolve_entities
Content-Type: application/json
```
请求示例:
```
{
"entities": [{"text":"北上","type":"location"}],
"context": {...},
"kb_providers": ["geo","user_profile"]
}
```
参数说明:
| 参数 | 类型 | 必需 | 默认值 | 说明 |
| --- | --- | ---: | --- | --- |
| entities | array | ✅ | - | 待消歧实体列表(text + type) |
| context | object | ❌ | null | 上下文限制候选范围 |
| kb_providers | array | ❌ | [] | 外部知识源优先级列表 |
| prefer_user_profile | bool | ❌ | true | 是否优先使用用户偏好 |
返回示例:
```
{
"success": true,
"resolved": [
{
"original":"北上",
"id":"city:beijing",
"canonical":"北京",
"confidence":0.87,
"sources":["geo","user_profile"]
}
]
}
```
---
## API: validate_and_format — 验证与结构化输出
功能说明:基于传入 `schema` 或目标接口规范,对意图对象执行类型验证、格式化(如日期/时区/单位归一化)与必要修正或注释。
REST API 调用:
```
POST /api/v1/validate
Content-Type: application/json
```
请求示例:
```
{
"intent_object": { "intent":"book_ticket", "slots":{...} },
"schema": { "slots": { "date": {"type":"date","format":"YYYY-MM-DD"} } },
"strict": true
}
```
参数说明:
| 参数 | 类型 | 必需 | 默认值 | 说明 |
| --- | --- | ---: | --- | --- |
| intent_object | object | ✅ | - | 待验证的意图与槽位结构 |
| schema | object | ❌ | null | 验证/格式化规则 |
| strict | boolean | ❌ | false | 严格模式:验证失败返回错误 |
| allow_coerce | boolean | ❌ | true | 允许类型强制转换 |
返回示例:
```
{
"success": true,
"validated": {
"intent":"book_ticket",
"slots":{
"date":{"value":"2026-02-22","type":"date","format":"YYYY-MM-DD","confidence":0.95},
"passengers":{"value":2,"type":"integer"}
}
},
"issues": []
}
```
---
## 统一错误格式
成功:
```
{
"success": true,
"data": {}
}
```
失败:
```
{
"success": false,
"error": "错误描述",
"code": "INVALID_SCHEMA"
}
```
---
如需进一步导出为 OpenAPI/MCP Schema、或生成工具注册说明与时序图,可基于本手册的接口示例扩展。
@@ -0,0 +1,15 @@
# 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
@@ -0,0 +1,19 @@
#!/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 职责越界检测 Agent API: http://{host}:{port}")
uvicorn.run(app, host=host, port=port, log_level="info")
@@ -0,0 +1,3 @@
"""Agent 职责越界检测 Agent 源代码包"""
@@ -0,0 +1,3 @@
"""Agent 服务器模块"""
@@ -0,0 +1,263 @@
"""
HTTP API 服务器 - Agent 职责越界检测 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 = "Agent 职责越界检测 Agent API"
# ==================== FastAPI 应用 ====================
@asynccontextmanager
async def lifespan(app: FastAPI):
print(f"🛡️ {SERVER_NAME} 启动")
yield
print(f"🛑 {SERVER_NAME} 关闭")
app = FastAPI(
title=SERVER_NAME,
description="专门负责检测 Agent 是否超出职责范围的 AI Agent - 判断一个 agent 是否开始'多管闲事',用于多 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": "专门负责检测 Agent 是否超出职责范围的 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 BoundaryCheckRequest(BaseModel):
"""职责越界检测请求模型"""
agent_name: str = Field(..., description="Agent 的名称或标识")
agent_defined_duty: str = Field(..., description="Agent 定义的职责范围(应该做什么)")
agent_actual_behavior: str = Field(..., description="Agent 的实际行为或输出(实际做了什么)")
current_task: Optional[str] = Field(None, description="可选,当前正在处理的任务描述")
class BoundaryCheckResponse(BaseModel):
"""职责越界检测响应模型"""
success: bool
is_violation: Optional[bool] = None
analysis: Optional[str] = None
error: Optional[str] = None
@app.post("/api/v1/check", response_model=BoundaryCheckResponse)
async def api_check_boundary(request: BoundaryCheckRequest, api_key: str = Depends(verify_api_key)):
"""
职责越界检测 API - 判断 Agent 是否超出职责范围
检测 agent 是否开始"多管闲事",是否超出了自己的职责范围。
用于多 agent 系统必备的职责边界保护。
"""
try:
# 设置 API Key
old_key = os.environ.get('OPENAI_API_KEY')
os.environ['OPENAI_API_KEY'] = api_key
try:
result_json = await TOOL_MAP['check_boundary_violation'](
agent_name=request.agent_name,
agent_defined_duty=request.agent_defined_duty,
agent_actual_behavior=request.agent_actual_behavior,
current_task=request.current_task
)
result = json.loads(result_json)
if result.get("success"):
return BoundaryCheckResponse(
success=True,
is_violation=result.get("is_violation"),
analysis=result.get("analysis")
)
else:
return BoundaryCheckResponse(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,175 @@
"""
MCP 服务器 - Agent 职责越界检测 Agent
专门负责检测 Agent 是否超出职责范围的 AI Agent。
判断一个 agent 是否开始"多管闲事",用于多 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("Agent 职责越界检测 Agent")
# 系统提示词 - 专注于检测职责越界
SYSTEM_PROMPT = '''你是一个专业的 Agent 职责边界检测专家,你的唯一职责是:判断一个 agent 是否开始"多管闲事"。
你的核心任务:
1. 分析 agent 的原始职责定义和实际行为
2. 判断 agent 是否超出了自己的职责范围
3. 识别 agent 是否在处理不属于自己的任务
4. 检测 agent 是否在"多管闲事"或越界操作
5. 为多 agent 系统提供职责边界保护
你的分析维度:
- 职责匹配度:agent 的行为是否与定义的职责一致
- 任务相关性:当前任务是否属于该 agent 的职责范围
- 越界行为:是否在处理其他 agent 应该处理的任务
- 职责边界:是否清晰地区分了"该做"和"不该做"的事情
你的判断标准:
- 严格但合理:既要防止越界,也要允许必要的协作
- 明确职责边界:清楚区分核心职责和辅助功能
- 识别越界信号:当 agent 开始处理明显不属于自己的任务时
- 提供改进建议:指出越界行为,并建议如何回归职责范围
输出格式要求:
- 明确给出是否越界的判断(是/否)
- 越界程度评分(0-10,0=无越界,10=严重越界)
- 具体的越界行为描述
- 职责边界分析
- 改进建议(如果越界)
记住:你的存在价值就是保护多 agent 系统的职责边界,防止 agent 之间职责混乱。要严格、客观、专业!'''
def get_agent() -> Agent:
"""创建 Agent 实例(每次调用使用最新的 API Key)"""
return Agent(MODEL_NAME, system_prompt=SYSTEM_PROMPT)
# ==================== MCP 工具定义 ====================
@server.tool()
async def check_boundary_violation(
agent_name: str,
agent_defined_duty: str,
agent_actual_behavior: str,
current_task: Optional[str] = None
) -> str:
"""
检测 Agent 是否超出职责范围(是否"多管闲事")
Args:
agent_name: Agent 的名称或标识
agent_defined_duty: Agent 定义的职责范围(应该做什么)
agent_actual_behavior: Agent 的实际行为或输出(实际做了什么)
current_task: 可选,当前正在处理的任务描述
Returns:
职责越界检测结果(JSON 格式)
"""
try:
# 构建分析提示
prompt = f"""请检测以下 Agent 是否超出职责范围:
Agent 名称:{agent_name}
定义的职责:
{agent_defined_duty}
实际行为:
{agent_actual_behavior}"""
if current_task:
prompt += f"\n\n当前任务:{current_task}"
prompt += "\n\n请严格分析该 Agent 是否开始'多管闲事',是否超出了自己的职责范围。"
# 调用 AI Agent 进行检测
result = await get_agent().run(prompt)
# 尝试解析结果,提取关键信息
output_text = result.output
# 判断是否越界(简单关键词检测,作为辅助)
violation_keywords = ['越界', '超出', '多管闲事', '不属于', '不应该', '职责范围外']
is_violation = any(keyword in output_text for keyword in violation_keywords)
return json.dumps({
"success": True,
"agent_name": agent_name,
"is_violation": is_violation,
"analysis": output_text,
"defined_duty": agent_defined_duty[:200] + "..." if len(agent_defined_duty) > 200 else agent_defined_duty,
"actual_behavior": agent_actual_behavior[:200] + "..." if len(agent_actual_behavior) > 200 else agent_actual_behavior
}, ensure_ascii=False, indent=2)
except Exception as e:
return json.dumps({
"success": False,
"error": str(e)
}, ensure_ascii=False)
# ==================== 工具映射(供 API 使用)====================
TOOL_MAP = {
'check_boundary_violation': check_boundary_violation,
}
TOOL_LIST = [
{
"name": "check_boundary_violation",
"description": "检测 Agent 是否超出职责范围(是否'多管闲事')。用于多 agent 系统必备的职责边界保护。",
"inputSchema": {
"type": "object",
"properties": {
"agent_name": {
"type": "string",
"description": "Agent 的名称或标识"
},
"agent_defined_duty": {
"type": "string",
"description": "Agent 定义的职责范围(应该做什么)"
},
"agent_actual_behavior": {
"type": "string",
"description": "Agent 的实际行为或输出(实际做了什么)"
},
"current_task": {
"type": "string",
"description": "可选,当前正在处理的任务描述"
}
},
"required": ["agent_name", "agent_defined_duty", "agent_actual_behavior"]
}
}
]
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"]
+114
View File
@@ -0,0 +1,114 @@
# 意图澄清 Agent (Intent Clarify Agent)
把模糊、不完整的人类输入补齐到可执行状态。
## 核心功能
**只做一件事**:将模糊、不完整、含糊的人类输入转化为清晰、完整、可执行的指令。
极适合放在任何 Agent 前面,让模型更准确理解用户要求。
### 提供能力
| 工具 | 功能 |
|------|------|
| `clarify_intent` | 澄清用户意图,补全缺失信息 |
| `expand_instruction` | 将简短指令扩展为详细步骤 |
| `infer_parameters` | 从自然语言推断参数值 |
| `batch_clarify` | 批量澄清多条输入 |
## 快速开始
### 本地运行
```bash
cd intent_clarify_agent
pip install -r requirements.txt
python run_api_server.py
```
### Docker 运行
```bash
docker build -t intent-clarify-agent:latest .
docker run -p 8000:8000 -e OPENAI_API_KEY=your-key intent-clarify-agent:latest
```
## 使用示例
### 输入
```
"帮我写个脚本"
```
### 输出
```json
{
"original_input": "帮我写个脚本",
"analysis": {
"core_intent": "请求编写脚本",
"missing_elements": ["脚本语言", "脚本功能", "运行环境"],
"ambiguous_parts": ["'脚本'指什么类型的脚本"],
"assumptions": ["假设是 Python 脚本", "假设用于自动化任务"]
},
"clarified_intent": {
"full_instruction": "请使用 Python 编写一个自动化脚本,实现指定功能,输出为可直接运行的 .py 文件",
"action": "编写",
"target": "Python 脚本",
"constraints": ["使用 Python 3.x", "代码需有注释"],
"expected_output": "可运行的 Python 脚本文件"
},
"confidence": "low",
"clarification_needed": ["具体需要脚本实现什么功能?", "有特定的编程语言偏好吗?"]
}
```
## 环境变量
| 变量 | 必需 | 说明 |
|------|------|------|
| OPENAI_API_KEY | 是 | OpenAI 或 LiteLLM API Key |
| OPENAI_BASE_URL | 否 | API Base URL |
| MODEL_NAME | 否 | 模型名称,默认 taiji/gpt-4o-mini |
| API_PORT | 否 | 服务端口,默认 8000 |
## 项目结构
```
intent_clarify_agent/
├── Dockerfile
├── README.md
├── USAGE.md # 详细使用文档
├── requirements.txt
├── run_api_server.py
└── src/
├── __init__.py
└── server/
├── __init__.py
├── api_server.py # FastAPI + MCP HTTP
└── mcp_server.py # MCP 工具定义
```
## 服务端点
| 端点 | 方法 | 说明 |
|------|------|------|
| `/` | GET | 服务状态 |
| `/health` | GET | 健康检查 |
| `/mcp` | POST | MCP JSON-RPC |
| `/mcp/sse` | GET/POST | MCP SSE 流式 |
| `/api/v1/clarify` | POST | 意图澄清 |
| `/api/v1/expand` | POST | 指令扩展 |
| `/api/v1/infer-params` | POST | 参数推断 |
| `/api/v1/batch-clarify` | POST | 批量澄清 |
## 部署信息
| 配置项 | 值 |
|--------|-----|
| 镜像地址 | agnettaiji.azurecr.io/ai-agents/intent-clarify-agent:latest |
| 服务端口 | 8000 |
| 健康检查 | /health |
+485
View File
@@ -0,0 +1,485 @@
# 意图澄清 Agent
---
本 Agent 提供 **意图澄清与输入补全** 能力,通过 **HTTP API** 与 **MCP(Model Context Protocol)** 对外提供服务。
核心能力:
- **意图澄清**:将模糊输入转化为清晰可执行的指令
- **指令扩展**:将简短指令展开为详细执行步骤
- **参数推断**:从自然语言中提取结构化参数
- **批量处理**:一次性处理多条模糊输入
---
## 功能概览
提供人类输入的 **语义理解、缺失补全、意图推断、参数提取** 能力,返回结构化的可执行指令。
支持能力:
- 识别输入中的缺失信息
- 基于上下文推断合理默认值
- 补全为完整可执行指令
- 标注不确定性与假设
---
## 1⃣ clarify_intent — 意图澄清
### 功能说明
将模糊、不完整的用户输入澄清为 **清晰、完整、可执行** 的指令,识别缺失要素并补全。
---
### REST API 调用
```
POST /api/v1/clarify
Content-Type: application/json
api-key: {your-api-key}
```
```json
{
"user_input": "帮我写个接口",
"context": "正在开发一个用户管理系统",
"domain": "后端开发"
}
```
---
### MCP 调用
```json
{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "clarify_intent",
"arguments": {
"user_input": "帮我写个接口",
"context": "正在开发一个用户管理系统",
"domain": "后端开发"
}
}
}
```
---
### 参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
| --- | --- | --- | --- | --- |
| user_input | string | ✅ | - | 用户的原始输入(可能模糊、不完整) |
| context | string | ❌ | null | 上下文信息(如之前的对话、当前任务背景) |
| domain | string | ❌ | null | 领域说明(如"编程"、"写作"、"数据分析") |
---
### 返回结果
```json
{
"success": true,
"original_input": "帮我写个接口",
"analysis": {
"core_intent": "请求编写 API 接口",
"missing_elements": ["接口功能", "接口路径", "请求方法", "返回格式"],
"ambiguous_parts": ["'接口'具体指 REST API 还是其他类型"],
"assumptions": ["假设是 REST API", "假设使用 JSON 格式", "基于上下文推断是用户相关接口"]
},
"clarified_intent": {
"full_instruction": "请使用 Python FastAPI 框架编写一个用户管理 REST API 接口,包含用户 CRUD 操作,使用 JSON 格式进行数据交换,需要包含请求参数验证和错误处理",
"action": "编写",
"target": "用户管理 REST API 接口",
"constraints": ["使用 FastAPI 框架", "JSON 数据格式", "包含参数验证"],
"expected_output": "可运行的 Python API 代码"
},
"confidence": "medium",
"clarification_needed": ["需要支持哪些具体的用户操作(增删改查)?", "是否需要身份验证?"]
}
```
---
### 返回字段说明
| 字段 | 类型 | 说明 |
| --- | --- | --- |
| original_input | string | 用户原始输入 |
| analysis.core_intent | string | 用户核心意图(一句话) |
| analysis.missing_elements | array | 缺失的关键要素 |
| analysis.ambiguous_parts | array | 模糊不清的部分 |
| analysis.assumptions | array | 为补全做出的假设 |
| clarified_intent.full_instruction | string | 完整可执行指令 |
| clarified_intent.action | string | 核心动作 |
| clarified_intent.target | string | 操作对象 |
| clarified_intent.constraints | array | 约束条件 |
| clarified_intent.expected_output | string | 期望输出 |
| confidence | string | 置信度:high/medium/low |
| clarification_needed | array | 仍需用户确认的问题 |
---
## 2️⃣ expand_instruction — 指令扩展
### 功能说明
将简短的指令扩展为 **详细的、可执行的步骤**,包含前置条件和交付物说明。
---
### REST API 调用
```
POST /api/v1/expand
Content-Type: application/json
api-key: {your-api-key}
```
```json
{
"brief_input": "部署应用到生产环境",
"task_type": "运维部署"
}
```
---
### MCP 调用
```json
{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "expand_instruction",
"arguments": {
"brief_input": "部署应用到生产环境",
"task_type": "运维部署"
}
}
}
```
---
### 参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
| --- | --- | --- | --- | --- |
| brief_input | string | ✅ | - | 简短的用户指令 |
| task_type | string | ❌ | null | 任务类型(如"代码编写"、"文档撰写") |
---
### 返回结果
```json
{
"success": true,
"brief_input": "部署应用到生产环境",
"expanded": {
"goal": "将应用程序安全、稳定地部署到生产环境",
"steps": [
{
"step": 1,
"action": "代码审查与测试",
"details": "确保所有代码通过 review,单元测试和集成测试全部通过",
"expected_result": "测试报告显示 100% 通过"
},
{
"step": 2,
"action": "构建生产镜像",
"details": "使用 CI/CD 流水线构建生产环境的 Docker 镜像",
"expected_result": "生成带版本标签的 Docker 镜像"
},
{
"step": 3,
"action": "备份当前生产环境",
"details": "备份数据库和配置文件,记录当前版本",
"expected_result": "备份文件和回滚方案就绪"
},
{
"step": 4,
"action": "部署新版本",
"details": "使用蓝绿部署或滚动更新策略发布新版本",
"expected_result": "新版本容器运行正常"
},
{
"step": 5,
"action": "验证与监控",
"details": "执行冒烟测试,检查日志和监控指标",
"expected_result": "服务健康,无异常告警"
}
],
"prerequisites": ["代码已合并到 main 分支", "获得部署审批", "通知相关团队"],
"deliverables": ["部署完成通知", "版本更新记录", "监控仪表盘链接"]
},
"full_instruction": "执行生产环境部署:首先确保代码审查和测试通过,然后构建生产镜像,备份当前环境后使用蓝绿部署策略发布新版本,最后执行验证并持续监控"
}
```
---
## 3️⃣ infer_parameters — 参数推断
### 功能说明
从用户的自然语言输入中 **提取/推断参数值**,用于填充工具或 API 调用的参数。
---
### REST API 调用
```
POST /api/v1/infer-params
Content-Type: application/json
api-key: {your-api-key}
```
```json
{
"user_input": "查询上海最近一周的天气,要详细的",
"available_params": ["city", "days", "detail_level", "format"],
"param_descriptions": "{\"city\": \"城市名称\", \"days\": \"查询天数\", \"detail_level\": \"详细程度(simple/detailed)\", \"format\": \"输出格式(json/text)\"}"
}
```
---
### MCP 调用
```json
{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "infer_parameters",
"arguments": {
"user_input": "查询上海最近一周的天气,要详细的",
"available_params": ["city", "days", "detail_level", "format"],
"param_descriptions": "{\"city\": \"城市名称\", \"days\": \"查询天数\"}"
}
}
}
```
---
### 参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
| --- | --- | --- | --- | --- |
| user_input | string | ✅ | - | 用户的自然语言输入 |
| available_params | array\<string\> | ✅ | - | 可用的参数名列表 |
| param_descriptions | string | ❌ | null | 参数描述(JSON 格式) |
---
### 返回结果
```json
{
"success": true,
"extracted_params": {
"city": "上海",
"days": "7"
},
"inferred_params": {
"detail_level": {
"value": "detailed",
"reason": "用户明确要求'要详细的'"
},
"format": {
"value": "json",
"reason": "未指定格式,使用默认值 json"
}
},
"missing_params": [],
"final_params": {
"city": "上海",
"days": "7",
"detail_level": "detailed",
"format": "json"
}
}
```
---
## 4️⃣ batch_clarify — 批量澄清
### 功能说明
一次性处理多条模糊输入,批量返回澄清结果。
---
### REST API 调用
```
POST /api/v1/batch-clarify
Content-Type: application/json
api-key: {your-api-key}
```
```json
{
"inputs": [
"帮我看看这个bug",
"优化一下性能",
"加个功能"
],
"shared_context": "正在维护一个 Python Web 应用"
}
```
---
### MCP 调用
```json
{
"jsonrpc": "2.0",
"id": 1,
"method": "tools/call",
"params": {
"name": "batch_clarify",
"arguments": {
"inputs": [
"帮我看看这个bug",
"优化一下性能",
"加个功能"
],
"shared_context": "正在维护一个 Python Web 应用"
}
}
}
```
---
### 参数说明
| 参数 | 类型 | 必需 | 默认值 | 说明 |
| --- | --- | --- | --- | --- |
| inputs | array\<string\> | ✅ | - | 待澄清的输入列表 |
| shared_context | string | ❌ | null | 共享的上下文信息 |
---
### 返回结果
```json
{
"success": true,
"results": [
{
"index": 1,
"original": "帮我看看这个bug",
"clarified": "请分析并定位当前 Python Web 应用中的 bug,提供错误原因分析和修复建议,需要提供具体的错误信息或日志",
"key_assumptions": ["需要提供错误信息", "假设是运行时错误"]
},
{
"index": 2,
"original": "优化一下性能",
"clarified": "请对 Python Web 应用进行性能优化,包括响应时间、内存使用和并发处理能力的改进,提供优化方案和预期效果",
"key_assumptions": ["全面性能优化", "假设关注响应速度"]
},
{
"index": 3,
"original": "加个功能",
"clarified": "请为 Python Web 应用添加新功能,需要提供功能需求描述、预期行为和接口设计",
"key_assumptions": ["需要具体功能需求", "假设是业务功能"]
}
],
"total": 3
}
```
---
## 统一错误格式
成功:
```json
{
"success": true,
"data": {}
}
```
失败:
```json
{
"success": false,
"error": "错误描述"
}
```
---
## 服务端点
| 端点 | 方法 | 说明 |
| --- | --- | --- |
| / | GET | 服务状态 |
| /health | GET | 健康检查 |
| /mcp | POST | MCP JSON-RPC |
| /mcp/sse | GET/POST | MCP SSE 流式 |
| /api/v1/clarify | POST | 意图澄清 |
| /api/v1/expand | POST | 指令扩展 |
| /api/v1/infer-params | POST | 参数推断 |
| /api/v1/batch-clarify | POST | 批量澄清 |
---
## 部署信息
| 配置项 | 值 |
| --- | --- |
| 镜像地址 | agnettaiji.azurecr.io/ai-agents/intent-clarify-agent:latest |
| 服务端口 | 8000 |
| 健康检查 | /health |
---
## 典型使用场景
### 场景 1:作为前置处理器
在调用其他 Agent 之前,先用本 Agent 澄清用户意图:
```
用户输入 → Intent Clarify Agent → 澄清后的指令 → 业务 Agent
```
### 场景 2:参数提取
从用户自然语言中提取 API 调用参数:
```
"搜索北京到上海明天的机票"
→ {from: "北京", to: "上海", date: "明天", type: "机票"}
```
### 场景 3:需求分析
将模糊的需求描述转化为具体的任务列表:
```
"做个登录功能"
→ 详细的登录功能需求,包含表单字段、验证规则、安全要求等
```
+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"🚀 启动 Intent Clarify Agent API: http://{host}:{port}")
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,348 @@
"""
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 = "Intent Clarify Agent API"
# ==================== 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",
description="意图澄清 Agent - 把模糊、不完整的人类输入补齐到可执行状态",
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": "意图澄清 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 ClarifyRequest(BaseModel):
"""意图澄清请求"""
user_input: str = Field(..., description="用户的原始输入(可能模糊、不完整)")
context: Optional[str] = Field(None, description="可选的上下文信息")
domain: Optional[str] = Field(None, description="可选的领域说明")
class ExpandRequest(BaseModel):
"""指令扩展请求"""
brief_input: str = Field(..., description="简短的用户指令")
task_type: Optional[str] = Field(None, description="任务类型")
class InferParamsRequest(BaseModel):
"""参数推断请求"""
user_input: str = Field(..., description="用户的自然语言输入")
available_params: List[str] = Field(..., description="可用的参数名列表")
param_descriptions: Optional[str] = Field(None, description="可选的参数描述")
class BatchClarifyRequest(BaseModel):
"""批量澄清请求"""
inputs: List[str] = Field(..., description="待澄清的输入列表")
shared_context: Optional[str] = Field(None, description="共享的上下文信息")
class ApiResponse(BaseModel):
"""API 响应"""
success: bool
data: Optional[Dict[str, Any]] = None
error: Optional[str] = None
@app.post("/api/v1/clarify", response_model=ApiResponse)
async def api_clarify(request: ClarifyRequest, api_key: str = Depends(verify_api_key)):
"""
澄清用户意图 - 将模糊输入转化为清晰可执行的指令
"""
try:
old_key = os.environ.get('OPENAI_API_KEY')
os.environ['OPENAI_API_KEY'] = api_key
try:
result = await TOOL_MAP['clarify_intent'](
user_input=request.user_input,
context=request.context,
domain=request.domain
)
parsed = json.loads(result)
if parsed.get("success"):
return ApiResponse(success=True, data=parsed)
else:
return ApiResponse(success=False, error=parsed.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))
@app.post("/api/v1/expand", response_model=ApiResponse)
async def api_expand(request: ExpandRequest, api_key: str = Depends(verify_api_key)):
"""
扩展指令 - 将简短指令扩展为详细步骤
"""
try:
old_key = os.environ.get('OPENAI_API_KEY')
os.environ['OPENAI_API_KEY'] = api_key
try:
result = await TOOL_MAP['expand_instruction'](
brief_input=request.brief_input,
task_type=request.task_type
)
parsed = json.loads(result)
if parsed.get("success"):
return ApiResponse(success=True, data=parsed)
else:
return ApiResponse(success=False, error=parsed.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))
@app.post("/api/v1/infer-params", response_model=ApiResponse)
async def api_infer_params(request: InferParamsRequest, api_key: str = Depends(verify_api_key)):
"""
参数推断 - 从自然语言输入中提取参数值
"""
try:
old_key = os.environ.get('OPENAI_API_KEY')
os.environ['OPENAI_API_KEY'] = api_key
try:
result = await TOOL_MAP['infer_parameters'](
user_input=request.user_input,
available_params=request.available_params,
param_descriptions=request.param_descriptions
)
parsed = json.loads(result)
if parsed.get("success"):
return ApiResponse(success=True, data=parsed)
else:
return ApiResponse(success=False, error=parsed.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))
@app.post("/api/v1/batch-clarify", response_model=ApiResponse)
async def api_batch_clarify(request: BatchClarifyRequest, api_key: str = Depends(verify_api_key)):
"""
批量澄清 - 一次性处理多条模糊输入
"""
try:
old_key = os.environ.get('OPENAI_API_KEY')
os.environ['OPENAI_API_KEY'] = api_key
try:
result = await TOOL_MAP['batch_clarify'](
inputs=request.inputs,
shared_context=request.shared_context
)
parsed = json.loads(result)
if parsed.get("success"):
return ApiResponse(success=True, data=parsed)
else:
return ApiResponse(success=False, error=parsed.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,476 @@
"""
意图澄清 MCP 服务器
核心功能:把模糊、不完整的人类输入补齐到可执行状态。
极适合放在任何 agent 前面,让模型更准确理解用户要求。
"""
import json
import os
from typing import Optional, List
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('Intent Clarify Agent')
# 系统提示词 - 专注于意图澄清和输入补全
SYSTEM_PROMPT = '''你是一个专业的意图澄清助手。你的核心任务是:
将模糊、不完整、含糊的人类输入转化为清晰、完整、可执行的指令。
你需要:
1. **识别缺失信息**:找出用户输入中缺少的关键要素(对象、范围、条件、格式等)
2. **推断合理默认值**:基于上下文和常识推断缺失部分的合理取值
3. **补全为可执行指令**:将模糊输入转化为精确、无歧义的指令
4. **标注不确定性**:对于无法确定的部分,给出合理假设并说明
处理原则:
- 保持用户原始意图不变,只补全缺失部分
- 优先使用上下文推断,而非随意假设
- 如有多种可能解释,选择最常见/最合理的
- 对于专业领域输入,推断专业默认值
- 输出必须是有效的 JSON 格式'''
def get_agent() -> Agent:
"""创建 Agent 实例(每次调用使用最新的 API Key)"""
return Agent(MODEL_NAME, system_prompt=SYSTEM_PROMPT)
# ==================== MCP 工具定义 ====================
@server.tool()
async def clarify_intent(
user_input: str,
context: Optional[str] = None,
domain: Optional[str] = None
) -> str:
"""
澄清并补全用户的模糊输入,转化为清晰可执行的指令。
Args:
user_input: 用户的原始输入(可能模糊、不完整)
context: 可选的上下文信息(如之前的对话、当前任务背景)
domain: 可选的领域说明(如"编程"、"写作"、"数据分析"),帮助推断专业默认值
Returns:
澄清结果(JSON 格式),包含补全后的指令和分析
"""
if not user_input or not user_input.strip():
return json.dumps({
"success": False,
"error": "用户输入为空"
}, ensure_ascii=False)
try:
# 构建提示
context_hint = f"\n\n上下文背景:{context}" if context else ""
domain_hint = f"\n\n领域:{domain}" if domain else ""
prompt = f'''请分析并澄清以下用户输入:
{context_hint}{domain_hint}
用户原始输入:
「{user_input}」
请按以下 JSON 格式输出(不要添加其他内容):
{{
"original_input": "用户原始输入",
"analysis": {{
"core_intent": "用户的核心意图(一句话概括)",
"missing_elements": ["缺失的关键要素1", "缺失的关键要素2"],
"ambiguous_parts": ["模糊不清的部分1", "模糊不清的部分2"],
"assumptions": ["为补全做出的假设1", "为补全做出的假设2"]
}},
"clarified_intent": {{
"full_instruction": "完整、清晰、可执行的指令",
"action": "核心动作(动词)",
"target": "操作对象",
"constraints": ["约束条件1", "约束条件2"],
"expected_output": "期望的输出格式/结果描述"
}},
"confidence": "high/medium/low",
"clarification_needed": ["如果仍需用户确认的问题1", "问题2"]
}}
要求:
1. "full_instruction" 必须是一个完整的、可直接执行的指令
2. 即使原输入很模糊,也要给出合理的补全结果
3. "assumptions" 说明你做了哪些假设来补全
4. 如果有关键信息必须由用户确认,放在 "clarification_needed"'''
# 调用 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)
return json.dumps({
"success": True,
**parsed
}, 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 expand_instruction(
brief_input: str,
task_type: Optional[str] = None
) -> str:
"""
将简短的指令扩展为详细的、可执行的步骤。
Args:
brief_input: 简短的用户指令
task_type: 任务类型(如"代码编写"、"文档撰写"、"数据处理")
Returns:
扩展后的详细指令(JSON 格式)
"""
if not brief_input or not brief_input.strip():
return json.dumps({
"success": False,
"error": "输入为空"
}, ensure_ascii=False)
try:
task_hint = f"\n任务类型:{task_type}" if task_type else ""
prompt = f'''将以下简短指令扩展为详细的、可执行的步骤:
{task_hint}
简短指令:
「{brief_input}」
请按以下 JSON 格式输出:
{{
"brief_input": "原始简短指令",
"expanded": {{
"goal": "最终目标",
"steps": [
{{
"step": 1,
"action": "具体动作",
"details": "详细说明",
"expected_result": "该步骤的预期结果"
}}
],
"prerequisites": ["前置条件1", "前置条件2"],
"deliverables": ["最终交付物1", "交付物2"]
}},
"full_instruction": "整合后的完整指令(一段话)"
}}'''
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,
**parsed
}, ensure_ascii=False, indent=2)
except Exception as e:
return json.dumps({
"success": False,
"error": str(e)
}, ensure_ascii=False)
@server.tool()
async def infer_parameters(
user_input: str,
available_params: List[str],
param_descriptions: Optional[str] = None
) -> str:
"""
从用户输入中推断参数值,用于填充工具/API 调用的参数。
Args:
user_input: 用户的自然语言输入
available_params: 可用的参数名列表
param_descriptions: 可选的参数描述(JSON 格式的参数说明)
Returns:
推断出的参数值(JSON 格式)
"""
if not user_input or not available_params:
return json.dumps({
"success": False,
"error": "用户输入或参数列表为空"
}, ensure_ascii=False)
try:
desc_hint = f"\n\n参数说明:\n{param_descriptions}" if param_descriptions else ""
params_list = "\n".join([f"- {p}" for p in available_params])
prompt = f'''从用户输入中提取/推断参数值。
用户输入:
「{user_input}」
可用参数:
{params_list}
{desc_hint}
请按以下 JSON 格式输出:
{{
"extracted_params": {{
"参数名1": "从输入中提取的值",
"参数名2": "从输入中提取的值"
}},
"inferred_params": {{
"参数名3": {{
"value": "推断的值",
"reason": "推断依据"
}}
}},
"missing_params": ["无法确定的参数1", "无法确定的参数2"],
"final_params": {{
"参数名1": "最终值",
"参数名2": "最终值"
}}
}}
要求:
1. "extracted_params" 是从用户输入中直接提取的
2. "inferred_params" 是基于上下文推断的,要说明推断依据
3. "missing_params" 是无法确定的参数
4. "final_params" 是所有参数的最终取值(包括默认值)'''
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,
**parsed
}, ensure_ascii=False, indent=2)
except Exception as e:
return json.dumps({
"success": False,
"error": str(e)
}, ensure_ascii=False)
@server.tool()
async def batch_clarify(
inputs: List[str],
shared_context: Optional[str] = None
) -> str:
"""
批量澄清多条用户输入。
Args:
inputs: 待澄清的输入列表
shared_context: 共享的上下文信息
Returns:
批量澄清结果(JSON 格式)
"""
if not inputs:
return json.dumps({
"success": True,
"results": [],
"total": 0
}, ensure_ascii=False)
try:
inputs_text = "\n".join([f"{i+1}. 「{inp}」" for i, inp in enumerate(inputs)])
context_hint = f"\n\n共享上下文:{shared_context}" if shared_context else ""
prompt = f'''批量澄清以下用户输入:
{context_hint}
输入列表:
{inputs_text}
请按以下 JSON 格式输出:
{{
"results": [
{{
"index": 1,
"original": "原始输入",
"clarified": "澄清后的完整指令",
"key_assumptions": ["主要假设"]
}}
]
}}
要求:
1. 每条输入都要澄清为完整、可执行的指令
2. 保持输入顺序
3. "key_assumptions" 只列出关键假设'''
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)
results = parsed.get('results', [])
return json.dumps({
"success": True,
"results": results,
"total": len(results)
}, ensure_ascii=False, indent=2)
except Exception as e:
return json.dumps({
"success": False,
"error": str(e)
}, ensure_ascii=False)
# ==================== 工具映射(供 API 使用)====================
TOOL_MAP = {
'clarify_intent': clarify_intent,
'expand_instruction': expand_instruction,
'infer_parameters': infer_parameters,
'batch_clarify': batch_clarify,
}
TOOL_LIST = [
{
"name": "clarify_intent",
"description": "澄清并补全用户的模糊输入,转化为清晰可执行的指令",
"inputSchema": {
"type": "object",
"properties": {
"user_input": {
"type": "string",
"description": "用户的原始输入(可能模糊、不完整)"
},
"context": {
"type": "string",
"description": "可选的上下文信息"
},
"domain": {
"type": "string",
"description": "可选的领域说明(如'编程'、'写作')"
}
},
"required": ["user_input"]
}
},
{
"name": "expand_instruction",
"description": "将简短的指令扩展为详细的、可执行的步骤",
"inputSchema": {
"type": "object",
"properties": {
"brief_input": {
"type": "string",
"description": "简短的用户指令"
},
"task_type": {
"type": "string",
"description": "任务类型(如'代码编写'、'文档撰写')"
}
},
"required": ["brief_input"]
}
},
{
"name": "infer_parameters",
"description": "从用户输入中推断参数值,用于填充工具/API 调用的参数",
"inputSchema": {
"type": "object",
"properties": {
"user_input": {
"type": "string",
"description": "用户的自然语言输入"
},
"available_params": {
"type": "array",
"items": {"type": "string"},
"description": "可用的参数名列表"
},
"param_descriptions": {
"type": "string",
"description": "可选的参数描述"
}
},
"required": ["user_input", "available_params"]
}
},
{
"name": "batch_clarify",
"description": "批量澄清多条用户输入",
"inputSchema": {
"type": "object",
"properties": {
"inputs": {
"type": "array",
"items": {"type": "string"},
"description": "待澄清的输入列表"
},
"shared_context": {
"type": "string",
"description": "共享的上下文信息"
}
},
"required": ["inputs"]
}
}
]
if __name__ == '__main__':
server.run()