编程 MCP 协议升级测试[片段5]

2026-07-26 07:52:06 +0800 CST views 7
"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
    
复制全文 生成海报 MCP AI Agent

推荐文章

介绍Vue3的Tree Shaking是什么?
2024-11-18 20:37:41 +0800 CST
一键压缩图片代码
2024-11-19 00:41:25 +0800 CST
程序员茄子在线接单