"uvicorn>=0.27.0",
"redis>=5.0.0",
"aioredis>=2.0.0",
"httpx>=0.27.0",
"pydantic>=2.6.0",
"structlog>=24.0.0",
"opentelemetry-api>=1.22.0",
"opentelemetry-sdk>=1.22.0",
"opentelemetry-instrumentation-fastapi>=0.43b0",
]
### A.2 工具定义:完整的 Schema 与 Evidence
这是 v0.28 规范下工具定义的完整示例:
```python
# src/schemas/business.py
from pydantic import BaseModel, Field
from typing import Optional, List, Literal
from datetime import datetime
class CompanySearchInput(BaseModel):
keyword: str = Field(
description="企业名称关键词(支持模糊匹配)",
min_length=2,
max_length=100
)
data_scope: Literal["current", "history", "all"] = Field(
default="current",
description="数据范围:current=当前在营,history=历史记录,all=全部"
)
province: Optional[str] = Field(
default=None,
description="省份筛选(行政区划代码,如 110000)"
)
limit: int = Field(
default=10,
ge=1,
le=100,
description="返回结果数量上限"
)
class CompanySearchOutput(BaseModel):
companies: List[dict] = Field(
description="匹配的企业列表"
)
total: int = Field(description="符合条件的总企业数")
search_id: str = Field(description="本次查询的唯一标识,用于关联后续查询")
class CompanySearchEvidence(BaseModel):
"""证据元信息:每个工具返回结果必须附带"""
evidence_id: str
trace_id: str
data_source: str
query_timestamp: datetime
data_timepoint: str # "current" | "history"
confidence: float
limitations: List[str]
cache_hit: bool
server_instance: str
A.3 限流器实现
# src/governance/rate_limiter.py
import time
import hashlib
from typing import Dict, Tuple
from collections import defaultdict
from dataclasses import dataclass
import structlog
logger = structlog.get_logger()
@dataclass
class RateLimitRule:
"""限流规则"""
requests_per_minute: int
requests_per_hour: int
burst_size: int # 允许的突发请求数
class TieredRateLimiter:
"""分层限流器:支持租户级、工具级的多层次限流"""
def __init__(self, redis_url: str):
self.redis_url = redis_url
self._local_burst_cache: Dict[str, Tuple[int, float]] = {}
self._rules: Dict[str, RateLimitRule] = {
"default": RateLimitRule(100, 1000, 20),
"company_search": RateLimitRule(500, 5000, 50),
"risk_scan": RateLimitRule(20, 200, 5),
"legal_case_query": RateLimitRule(100, 1000, 20),
}
async def check(
self,
tenant_id: str,
tool_name: str,
trace_id: str
) -> Tuple[bool, dict]:
"""
检查请求是否允许通过
返回: (是否允许, 限流元信息)
"""
rule = self._rules.get(tool_name, self._rules["default"])
now = time.time()
# 突发限流(本地内存,毫秒级)
burst_key = f"{tenant_id}:{tool_name}"
if burst_key in self._local_burst_cache:
count, window_start = self._local_burst_cache[burst_key]
window_duration = now - window_start
if window_duration < 1.0: # 1秒窗口
if count >= rule.burst_size:
return False, {
"reason": "burst_limit",
"retry_after_ms": int(1000 - window_duration * 1000)
}
self._local_burst_cache[burst_key] = (count + 1, window_start)
else:
self._local_burst_cache[burst_key] = (1, now)
else:
self._local_burst_cache[burst_key] = (1, now)
# 分钟级限流(Redis)
minute_key = f"rl:minute:{tenant_id}:{tool_name}"
minute_count = await self.redis.get(minute_key)
if minute_count and int(minute_count) >= rule.requests_per_minute:
ttl = await self.redis.ttl(minute_key)
return False, {
"reason": "minute_limit",
"retry_after_ms": ttl * 1000
}
# 小时级限流(Redis)
hour_key = f"rl:hour:{tenant_id}:{tool_name}"
hour_count = await self.redis.get(hour_key)
if hour_count and int(hour_count) >= rule.requests_per_hour:
ttl = await self.redis.ttl(hour_key)
return False, {
"reason": "hour_limit",
"retry_after_ms": ttl * 1000
}
# 记录请求
pipe = self.redis.pipeline()
pipe.incr(minute_key)
pipe.expire(minute_key, 60)
pipe.incr(hour_key)
pipe.expire(hour_key, 3600)
await pipe.execute()
return True, {"allowed": True}
def get_headers(self, metadata: dict) -> dict:
"""生成返回给客户端的限流头"""
headers = {}
if "retry_after_ms" in metadata:
headers["Retry-After"] = str(metadata["retry_after_ms"] // 1000)
headers["X-RateLimit-Retry-After-Ms"] = str(metadata["retry_after_ms"])
return headers
A.4 证据元信息生成器
# src/utils/evidence.py
import hashlib
import uuid
from datetime import datetime, timezone
from typing import List, Optional
from dataclasses import dataclass, asdict
@dataclass
class Evidence:
"""MCP v0.28 证据元信息:每个工具调用结果的核心组成部分"""
evidence_id: str
trace_id: str
data_source: str
query_timestamp: str
data_timepoint: str
confidence: float
limitations: List[str]
fields_used: List[str]
raw_response_hash: str
server_instance: str
protocol_version: str = "0.28.0"
cache_hit: bool = False
def to_meta(self) -> dict:
"""转换为返回给客户端的 _meta 字段"""
return {
"evidence_id": self.evidence_id,
"trace_id": self.trace_id,
"data_source": self.data_source,
"query_timestamp": self.query_timestamp,
"data_timepoint": self.data_timepoint,
"confidence": self.confidence,
"limitations": self.limitations,
"fields_used": self.fields_used,
"raw_response_hash": self.raw_response_hash,
"server_instance": self.server_instance,
"cache_hit": self.cache_hit,
"protocol_version": self.protocol_version