返回博客

门面模式:统一 API 入口

coomia-dip 内部有 8 个 Layer,每个 Layer 有自己的 gRPC 服务、数据模型和通信协议。如果客户端直接与各 Layer 交互,复杂度将不可控:

Coomia发布于 2025年12月26日9 分钟阅读
分享本文Twitter / X

门面模式:统一 API 入口

系列:S10 设计模式 · 第 11 篇 | 难度:高级 | 阅读时间:18 分钟

#TL;DR

  • 门面模式(Facade Pattern)为 coomia-dip 的多 Layer 架构提供统一的 API 入口。外部客户端只需与一个 API Gateway 交互,无需了解内部 8 个 Layer 的复杂拓扑。
  • coomia-dip 的门面层实现了协议转换(REST → gRPC)、请求路由、认证鉴权、限流熔断、请求聚合和响应转换等横切关注点。
  • SDK 是门面模式的客户端体现——将复杂的多步 API 调用封装为简洁的方法调用。

#引言:复杂性隐藏

coomia-dip 内部有 8 个 Layer,每个 Layer 有自己的 gRPC 服务、数据模型和通信协议。如果客户端直接与各 Layer 交互,复杂度将不可控:

Code
问题 1:客户端需要知道每个 Layer 的地址和端口
问题 2:不同 Layer 使用不同的 Protobuf 定义
问题 3:一个业务操作可能需要调用 3-4 个 Layer
问题 4:认证、限流、日志等横切关注点在每个 Layer 都要实现一遍
问题 5:客户端升级时需要同步升级对多个 Layer 的调用

门面模式的核心思想:用一个简洁的接口隐藏内部系统的复杂性

#一、API Gateway 门面

#1.1 统一入口架构

Python
from dataclasses import dataclass, field
from typing import Any
from enum import Enum

class PlaneTarget(Enum):
    CONTROL = "control-Layer"
    DATA = "data-Layer"
    INTELLIGENCE = "intelligence-Layer"
    AGENT = "agent-runtime"
    SDK = "sdk-Layer"

@dataclass
class RouteConfig:
    """API route configuration."""
    path: str
    method: str
    target_plane: PlaneTarget
    grpc_service: str
    grpc_method: str
    requires_auth: bool = True
    rate_limit: int = 100          # requests/minute
    timeout_seconds: int = 30
    cache_ttl: int = 0             # 0 = no cache
    aggregate_from: list[str] = field(default_factory=list)  # 需要聚合的 Layer

class APIGatewayFacade:
    """Unified API Gateway facade for coomia-dip."""

    def __init__(
        self,
        routes: list[RouteConfig],
        auth_service: "AuthService",
        rate_limiter: "RateLimiter",
        circuit_breaker: "CircuitBreaker",
    ):
        self._routes = {(r.path, r.method): r for r in routes}
        self._auth = auth_service
        self._rate_limiter = rate_limiter
        self._circuit_breaker = circuit_breaker

    async def handle_request(self, request: "HTTPRequest") -> "HTTPResponse":
        """Handle an incoming HTTP request through the facade."""
        route = self._routes.get((request.path, request.method))
        if not route:
            return HTTPResponse(status=404, body={"error": "Not found"})

        # 1. 认证
        if route.requires_auth:
            auth_result = await self._auth.authenticate(request)
            if not auth_result.authenticated:
                return HTTPResponse(status=401, body={"error": "Unauthorized"})

        # 2. 限流
        if not await self._rate_limiter.allow(
            request.client_id, route.rate_limit
        ):
            return HTTPResponse(status=429, body={"error": "Rate limit exceeded"})

        # 3. 熔断检查
        if self._circuit_breaker.is_open(route.target_plane.value):
            return HTTPResponse(
                status=503,
                body={"error": f"Service {route.target_plane.value} unavailable"},
            )

        try:
            # 4. 路由到目标 Layer
            if route.aggregate_from:
                result = await self._aggregate_request(request, route)
            else:
                result = await self._route_request(request, route)

            return HTTPResponse(status=200, body=result)

        except TimeoutError:
            return HTTPResponse(status=504, body={"error": "Gateway timeout"})
        except Exception as e:
            self._circuit_breaker.record_failure(route.target_plane.value)
            return HTTPResponse(status=500, body={"error": str(e)})

    async def _route_request(
        self, request: "HTTPRequest", route: RouteConfig
    ) -> dict:
        """Route request to target Layer via gRPC."""
        endpoint = self._service_discovery.resolve(route.target_plane.value)

        async with grpc_channel(endpoint) as channel:
            # REST → gRPC 协议转换
            grpc_request = self._rest_to_grpc(
                request.body, route.grpc_service, route.grpc_method
            )
            stub = self._get_stub(channel, route.grpc_service)
            method = getattr(stub, route.grpc_method)
            response = await asyncio.wait_for(
                method(grpc_request),
                timeout=route.timeout_seconds,
            )
            return self._grpc_to_rest(response)

    async def _aggregate_request(
        self, request: "HTTPRequest", route: RouteConfig
    ) -> dict:
        """Aggregate responses from multiple Layers."""
        tasks = []
        for plane_name in route.aggregate_from:
            endpoint = self._service_discovery.resolve(plane_name)
            tasks.append(self._call_plane(endpoint, request, route))

        results = await asyncio.gather(*tasks, return_exceptions=True)

        aggregated = {}
        for plane_name, result in zip(route.aggregate_from, results):
            if isinstance(result, Exception):
                aggregated[plane_name] = {"error": str(result)}
            else:
                aggregated[plane_name] = result

        return aggregated

#1.2 路由配置

Python
# coomia-dip API 路由配置
ROUTES = [
    # Ontology CRUD
    RouteConfig("/api/v1/ontology/types", "GET", PlaneTarget.CONTROL,
                "OntologyService", "ListObjectTypes"),
    RouteConfig("/api/v1/ontology/types", "POST", PlaneTarget.CONTROL,
                "OntologyService", "CreateObjectType"),
    RouteConfig("/api/v1/ontology/objects/{type}", "GET", PlaneTarget.DATA,
                "ObjectQueryService", "QueryObjects"),

    # Actions
    RouteConfig("/api/v1/actions/{action_id}/execute", "POST", PlaneTarget.CONTROL,
                "ActionService", "ExecuteAction", timeout_seconds=60),

    # Reasoning
    RouteConfig("/api/v1/reasoning/evaluate", "POST", PlaneTarget.INTELLIGENCE,
                "ReasoningService", "Evaluate", timeout_seconds=120),

    # Dashboard (聚合多个 Layer)
    RouteConfig("/api/v1/dashboard/overview", "GET", PlaneTarget.CONTROL,
                "DashboardService", "GetOverview",
                aggregate_from=["control-Layer", "data-Layer", "intelligence-Layer"]),

    # Metrics
    RouteConfig("/api/v1/metrics", "GET", PlaneTarget.CONTROL,
                "MetricsService", "GetMetrics", cache_ttl=60),
]

#二、协议转换层

#2.1 REST 到 gRPC 转换

Python
class ProtocolTranslator:
    """Translate between REST/JSON and gRPC/Protobuf."""

    def __init__(self, proto_registry: "ProtoRegistry"):
        self._proto_registry = proto_registry

    def rest_to_grpc(
        self, json_body: dict, service_name: str, method_name: str
    ) -> "Message":
        """Convert REST JSON body to gRPC Protobuf message."""
        message_type = self._proto_registry.get_request_type(
            service_name, method_name
        )
        return self._json_to_proto(json_body, message_type)

    def grpc_to_rest(self, proto_message: "Message") -> dict:
        """Convert gRPC Protobuf response to REST JSON."""
        return self._proto_to_json(proto_message)

    def _json_to_proto(self, data: dict, message_type: type) -> "Message":
        """Convert JSON dict to Protobuf message."""
        from google.protobuf.json_format import ParseDict
        return ParseDict(data, message_type())

    def _proto_to_json(self, message: "Message") -> dict:
        """Convert Protobuf message to JSON dict."""
        from google.protobuf.json_format import MessageToDict
        return MessageToDict(message, preserving_proto_field_name=True)

#2.2 GraphQL 门面

除了 REST,coomia-dip 还提供 GraphQL 接口作为另一种门面:

Python
class GraphQLFacade:
    """GraphQL facade for flexible querying."""

    async def execute_query(
        self, query: str, variables: dict | None = None, context: dict | None = None
    ) -> dict:
        """Execute a GraphQL query against the Ontology."""
        parsed = self._parser.parse(query)

        # 将 GraphQL 查询分解为对各 Layer 的调用
        execution_plan = self._planner.plan(parsed, variables)

        results = {}
        for step in execution_plan.steps:
            result = await self._execute_step(step, context)
            results[step.field_name] = result

        return self._assemble_response(results, parsed)

#三、请求聚合

#3.1 BFF(Backend For Frontend)模式

不同前端页面需要的数据来自不同 Layer。BFF 层将多个 Layer 的数据聚合为一个响应:

Python
class DashboardBFF:
    """Backend-For-Frontend for the dashboard page."""

    async def get_dashboard_data(
        self, tenant_id: str, world_id: str
    ) -> dict:
        """Aggregate dashboard data from multiple Layers."""
        # 并行调用各 Layer
        ontology_stats, data_stats, reasoning_stats, agent_stats = (
            await asyncio.gather(
                self._control_plane.get_ontology_stats(tenant_id),
                self._data_plane.get_storage_stats(tenant_id),
                self._intelligence_plane.get_reasoning_stats(tenant_id),
                self._agent_runtime.get_agent_stats(tenant_id),
            )
        )

        return {
            "ontology": {
                "total_types": ontology_stats.total_types,
                "total_objects": ontology_stats.total_objects,
                "total_links": ontology_stats.total_links,
            },
            "data": {
                "storage_used_gb": data_stats.storage_used_gb,
                "total_tables": data_stats.total_tables,
                "query_count_today": data_stats.query_count_today,
            },
            "reasoning": {
                "active_rules": reasoning_stats.active_rules,
                "decisions_today": reasoning_stats.decisions_today,
                "avg_latency_ms": reasoning_stats.avg_latency_ms,
            },
            "agents": {
                "active_agents": agent_stats.active_agents,
                "tasks_completed_today": agent_stats.tasks_completed_today,
            },
        }

#3.2 批量请求

客户端可以在一个 HTTP 请求中发送多个 API 调用:

Python
class BatchRequestHandler:
    """Handle batch API requests."""

    async def handle_batch(
        self, batch: list[dict], context: dict
    ) -> list[dict]:
        """Execute multiple API calls in a single request."""
        results = []
        # 分析依赖关系,独立请求并行执行
        independent, dependent = self._analyze_dependencies(batch)

        # 并行执行独立请求
        independent_results = await asyncio.gather(
            *[self._execute_single(req, context) for req in independent]
        )

        # 按序执行有依赖的请求
        for req in dependent:
            # 替换引用(如 $ref:0.id)
            resolved = self._resolve_references(req, results)
            result = await self._execute_single(resolved, context)
            results.append(result)

        return list(independent_results) + results

#四、SDK 作为客户端门面

#4.1 SDK 封装

SDK 是门面模式在客户端的体现。用户不需要知道背后调用了几个 API:

Python
class OntoPlatformSDK:
    """Python SDK as client-side facade."""

    def __init__(self, endpoint: str, api_key: str):
        self._client = HTTPClient(endpoint, api_key)

    async def create_risk_rule(
        self,
        name: str,
        conditions: list[dict],
        actions: list[dict],
    ) -> "RiskRule":
        """Create a risk rule — hides 4 API calls behind one method."""
        # 1. 创建 ObjectType(如果不存在)
        obj_type = await self._ensure_object_type("RiskRule")

        # 2. 注册 Action
        action = await self._client.post("/api/v1/actions", {
            "name": f"evaluate_{name}",
            "object_type": "RiskRule",
            "logic": conditions,
        })

        # 3. 创建规则对象
        rule = await self._client.post("/api/v1/ontology/objects/RiskRule", {
            "name": name,
            "conditions": conditions,
            "actions": actions,
            "action_id": action["id"],
            "_state": "draft",
        })

        # 4. 初始化推理模板
        await self._client.post("/api/v1/reasoning/templates", {
            "rule_id": rule["id"],
            "conditions": conditions,
        })

        return RiskRule.from_dict(rule)

#4.2 SDK 错误转换

SDK 将底层的 gRPC/HTTP 错误转换为业务友好的异常:

Python
class SDKErrorTranslator:
    """Translate low-level errors to business-friendly exceptions."""

    ERROR_MAP = {
        "PERMISSION_DENIED": PermissionError,
        "NOT_FOUND": ObjectNotFoundError,
        "ALREADY_EXISTS": DuplicateObjectError,
        "INVALID_ARGUMENT": ValidationError,
        "RESOURCE_EXHAUSTED": QuotaExceededError,
        "UNAVAILABLE": ServiceUnavailableError,
    }

    def translate(self, error: Exception) -> Exception:
        """Translate a gRPC or HTTP error to SDK exception."""
        if hasattr(error, "code"):  # gRPC error
            error_class = self.ERROR_MAP.get(error.code().name, SDKError)
            return error_class(str(error.details()))
        elif hasattr(error, "status_code"):  # HTTP error
            if error.status_code == 404:
                return ObjectNotFoundError(error.text)
            elif error.status_code == 403:
                return PermissionError(error.text)
        return SDKError(str(error))

#五、横切关注点

#5.1 请求日志与追踪

Python
class RequestLogger:
    """Log all requests passing through the facade."""

    async def log_request(
        self, request: "HTTPRequest", response: "HTTPResponse", duration_ms: float
    ) -> None:
        await self._log_store.append({
            "timestamp": datetime.utcnow().isoformat(),
            "method": request.method,
            "path": request.path,
            "client_id": request.client_id,
            "status": response.status,
            "duration_ms": duration_ms,
            "trace_id": request.headers.get("X-Trace-ID"),
            "request_size": len(str(request.body)),
            "response_size": len(str(response.body)),
        })

#5.2 API 版本管理

Python
class APIVersionRouter:
    """Route requests to different API versions."""

    def __init__(self):
        self._versions: dict[str, dict] = {}

    def register_version(self, version: str, routes: list[RouteConfig]) -> None:
        self._versions[version] = {(r.path, r.method): r for r in routes}

    def resolve(self, version: str, path: str, method: str) -> RouteConfig | None:
        version_routes = self._versions.get(version, {})
        return version_routes.get((path, method))

#5.3 响应缓存

Python
class ResponseCache:
    """Cache API responses at the facade level."""

    async def get_or_execute(
        self, cache_key: str, ttl: int, executor: callable
    ) -> dict:
        cached = await self._redis.get(f"api:cache:{cache_key}")
        if cached:
            return json.loads(cached)

        result = await executor()
        if ttl > 0:
            await self._redis.set(
                f"api:cache:{cache_key}",
                json.dumps(result),
                ex=ttl,
            )
        return result

#Key Takeaways

  1. 复杂性隐藏:门面模式将 8 个 Layer 的复杂拓扑隐藏在统一 API 之后
  2. 协议转换:REST/JSON ↔ gRPC/Protobuf 的透明转换
  3. 请求聚合:BFF 层将多 Layer 数据聚合为单一响应,减少客户端往返
  4. SDK 封装:SDK 将多步 API 调用封装为单一方法,提供业务友好的接口
  5. 横切关注点:认证、限流、熔断、日志、缓存在门面层统一处理
  6. 版本管理:API 版本路由确保向后兼容性

#Next Article

下一篇我们将探讨契约优先模式(Contract-First)——coomia-dip 如何用 Protobuf 和 OpenAPI 定义先行的方式确保 Layer 间接口的兼容性。

S10-12: 契约优先:接口定义先行

#Tags

#设计模式 #门面模式 #Facade #APIGateway #协议转换 #BFF #SDK #限流 #熔断