编程 MCP 协议 0.28 重大升级:无状态核心、能力治理体系与 Agent 生产级落地的完整指南 [range 24212-36318]

2026-07-26 07:51:08 +0800 CST views 8
"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
        }

class EvidenceGenerator:
    def __init__(self, server_instance: str, redis_url: str):
        self.server_instance = server_instance
        self.redis = None  # 懒加载
    
    @staticmethod
    def hash_response(data: dict) -> str:
        """对响应数据生成哈希,用于审计"""
        import json
        normalized = json.dumps(data, sort_keys=True, ensure_ascii=False)
        return hashlib.sha256(normalized.encode()).hexdigest()[:16]
    
    async def generate(
        self,
        trace_id: str,
        tool_name: str,
        data: dict,
        data_source: str,
        data_timepoint: str,
        confidence: float,
        limitations: List[str],
        fields_used: List[str],
        cache_hit: bool = False
    ) -> Evidence:
        """生成证据元信息"""
        return Evidence(
            evidence_id=f"ev-{uuid.uuid4().hex[:12]}",
            trace_id=trace_id,
            data_source=data_source,
            query_timestamp=datetime.now(timezone.utc).isoformat(),
            data_timepoint=data_timepoint,
            confidence=confidence,
            limitations=limitations,
            fields_used=fields_used,
            raw_response_hash=self.hash_response(data),
            server_instance=self.server_instance,
            cache_hit=cache_hit
        )

A.5 主服务器入口

# src/server.py
import asyncio
import uuid
from contextlib import asynccontextmanager
from fastapi import FastAPI, HTTPException, Request, Response
from fastapi.middleware.cors import CORSMiddleware
from pydantic import BaseModel
import structlog
from opentelemetry import trace

from mcp.server import Server
from mcp.types import Tool, TextContent
from mcp.server.stdio import stdio_server

from src.governance.rate_limiter import TieredRateLimiter
from src.utils.evidence import EvidenceGenerator

logger = structlog.get_logger()
tracer = trace.get_tracer(__name__)

app = FastAPI(title="Enterprise MCP Server v0.28")
app.add_middleware(
    CORSMiddleware,
    allow_origins=["*"],
    allow_methods=["*"],
    allow_headers=["*"],
)

# 全局组件(通过 lifespan 管理生命周期)
rate_limiter: TieredRateLimiter = None
evidence_gen: EvidenceGenerator = None

@asynccontextmanager
async def lifespan(app: FastAPI):
    global rate_limiter, evidence_gen
    import aioredis
    redis = await aioredis.from_url("redis://localhost:6379")
    rate_limiter = TieredRateLimiter("redis://localhost:6379")
    evidence_gen = EvidenceGenerator(
        server_instance=f"server-{uuid.uuid4().hex[:8]}",
        redis_url="redis://localhost:6379"
    )
    logger.info("server.started", instance=evidence_gen.server_instance)
    yield
    await redis.close()
    logger.info("server.stopped")

app.router.lifespan_context = lifespan

# ========== MCP v0.28 HTTP 传输 ==========

class ToolCallRequest(BaseModel):
    name: str
    arguments: dict
    trace_id: Optional[str] = None
    tenant_id: Optional[str] = None
    meta: Optional[dict] = None

@app.post("/tools/call")
async def call_tool(req: ToolCallRequest, request: Request) -> Response:
    trace_id = req.trace_id or uuid.uuid4().hex
    tenant_id = req.tenant_id or "anonymous"
    
    with tracer.start_as_current_span(f"tool.{req.name}") as span:
        span.set_attribute("trace_id", trace_id)
        span.set_attribute("tenant_id", tenant_id)
        span.set_attribute("tool_name", req.name)
        
        # Step 1: 限流检查
        allowed, limit_meta = await rate_limiter.check(
            tenant_id, req.name, trace_id
        )
        if not allowed:
            headers = rate_limiter.get_headers(limit_meta)
            return Response(
                status_code=429,
                content='{"error": "rate_limit_exceeded"}',
                headers={**headers, "Content-Type": "application/json"}
            )
        
        # Step 2: 执行工具
        try:
            result = await execute_tool(req.name, req.arguments)
            
            # Step 3: 生成证据元信息
            evidence = await evidence_gen.generate(
                trace_id=trace_id,
                tool_name=req.name,
                data=result,
                data_source="business_database",
                data_timepoint="current",
                confidence=0.95,
                limitations=["仅覆盖中国大陆企业数据"],
                fields_used=list(result.keys()),
                cache_hit=False
            )
            
            response_body = {
                "result": result,
                "_meta": evidence.to_meta()
            }
            
            span.set_attribute("success", True)
            return Response(
                content=json.dumps(response_body),
                headers={
                    "Content-Type": "application/json",
                    "X-Trace-Id": trace_id,
                    "X-Evidence-Id": evidence.evidence_id
                }
            )
            
        except Exception as e:
            span.record_exception(e)
            logger.error(
                "tool.execution_failed",
                tool=req.name,
                trace_id=trace_id,
                error=str(e)
            )
            raise HTTPException(status_code=500, detail=str(e))

@app.get("/tools/list")
async def list_tools():
    """返回带完整 v0.28 annotations 的工具清单"""
    return {
        "tools": [
            {
                "name": "company_search",
                "description": "模糊搜索企业名称,返回匹配的企业列表",
                "inputSchema": CompanySearchInput.model_json_schema(),
                "outputSchema": CompanySearchOutput.model_json_schema(),
                "annotations": {
                    "dataScope": "current",
                    "dataFreshness": "realtime",
                    "capabilityType": "search",
                    "requiresAnchor": False,
                    "prerequisites": [],
                    "mutuallyExclusiveWith": [],
                    "cacheTtlMs": 300000,  # 5分钟缓存
                    "costUnits": 2,
                    "rateLimitTier": "default"
                }
            },
            {
                "name": "risk_sc
复制全文 生成海报 MCP AI Agent

推荐文章

Elasticsearch 的索引操作
2024-11-19 03:41:41 +0800 CST
如何在Rust中使用UUID?
2024-11-19 06:10:59 +0800 CST
程序员茄子在线接单