Token导航 LogoToken导航TokenDH.com

手把手构建企业级 Agent 框架(二):Gateway 网关与多渠道接入

更新时间 2026-05-18来源 Aike正文 1.1万字阅读约 35分钟1 张图片

在上篇文章手把手构建企业级 Agent 框架:从 OpenClaw 架构到自主实现中,我们剖析了 OpenClaw 的架构骨架,并搭建了一个包含 Gateway、Agent、Skill 的最小原型。今天,我们将深入框架的“咽喉要道”——Gateway 网关。如果说 Agent 是大脑,那么 Gateway 就是整个系统的耳朵和嘴巴,它负责倾听来自不同渠道的声音,并精准地将回复送达用户。一个设计良好的 Gateway,能够让你的 Agent 无缝出现在飞书、企微、WebChat、甚至自定义 IoT 设备中,而核心逻辑零改动。

一、为什么 Gateway 是企业 Agent 的第一道关卡?

OpenClaw 的设计哲学中,Gateway 被赋予了极高的地位:它不只是反向代理,更是安全边界、协议适配器和流量调度中心。在企业环境中,这种思想更加重要:

  • 多租户隔离:
    不同部门、不同客户的消息必须严格路由到对应的 Agent 实例,数据不得混淆。
  • 安全校验:
    认证、鉴权、API Key 验证必须在入口处完成,绝不能将裸请求直接抛给 Agent。
  • 协议统一:
    飞书的 Webhook 格式、微信的 XML 消息、WebSocket 的 JSON 帧……Gateway 负责将它们翻译成内部标准事件。
  • 会话连续性:
    同一个用户在不同渠道发言,Agent 应当能识别出这是同一个人,并延续上下文。

因此,我们的企业级 Gateway 将承担六大核心职能,直接参考 OpenClaw 的设计并强化:

Gateway 六大职责:

  1. 会话管理:
    创建、恢复、过期清理用户会话,绑定渠道身份。
  2. 智能路由:
    根据用户标签、关键词、时间段等规则,将消息转发给不同的 Agent 或技能组。
  3. 协议转换:
    将异构平台消息标准化为统一的 AgentRequest 内部格式。
  4. 安全校验:
    JWT 验证、签名校验、租户白名单、速率限制。
  5. 设备配对:
    支持同一用户在多个设备上同时连接,消息实时同步。
  6. 访问控制:
    基于角色的工具/功能权限,在 Gateway 层即可拦截未授权操作。

二、整体架构:Channel 抽象与内部事件流

为了实现多渠道无缝接入,我们引入 Channel Adapter(渠道适配器) 的概念。每个外部平台对应一个 Adapter,它负责两件事:

  • 将平台特定的消息解析为内部 AgentRequest
  • 将内部 AgentEvent 流转换回平台所需的回复格式(文本、卡片、Markdown 等)

下图展示了 Gateway 内部组件以及数据流向:

图片

关键设计决策:Gateway 通过消息总线(Redis Pub/Sub)与 Agent 运行时解耦。Gateway 不直接调用 Agent,而是发布标准化请求到特定队列,并监听响应流。这样做的好处是: - Agent 可以独立扩缩容,Gateway 无感知。 - 支持请求持久化与重试,避免 Agent 崩溃导致消息丢失。 - 方便加入监控、限流等中间件。

三、接口与数据流设计

让我们用 Python 协议类定义核心接口。这保证了未来任何新渠道只需遵循契约即可接入。

3.1 渠道适配器协议

from typing import Protocol, AsyncIterator
from dataclasses import dataclass

classChannelAdapter(Protocol):
"""渠道适配器必须实现的接口"""
    channel_name:str

asyncdefparse_incoming(self, raw_data:bytes, metadata:dict)->'AgentRequest':
"""将原始请求解析为内部标准格式"""
...

asyncdefformat_outgoing(self, event:'AgentEvent')->bytes:
"""将内部事件转换为平台特定回复格式"""
...

@dataclass
classAgentRequest:
    session_id:str
    user_id:str
    tenant_id:str
    content:str
    channel:str
    metadata:dict# 包含原始请求头、签名等

3.2 标准化消息格式

所有进入 Gateway 的消息,无论来源,最终都会被转换为以下 JSON 结构(内部流转):

{
"session_id":"uuid-xxxx",
"user_id":"user_123",
"tenant_id":"tenant_abc",
"content":"查询订单 O12345",
"channel":"feishu",
"metadata":{
"timestamp":1715000000,
"raw_signature":"...",
"device_id":"mobile_01"
}
}

这种标准化让后续的路由、Agent 处理、日志审计变得极其简单——所有组件都只关心这一种格式。

3.3 路由规则引擎

路由规则采用插件式设计,每条规则实现一个简单的判定函数。在实际项目里,可以集成表达式引擎(如 expr 库)实现更灵活的配置。

from typing import Callable, Optional

classRouteRule:
def__init__(self, predicate: Callable[[AgentRequest],bool], target:str):
        self.predicate = predicate
        self.target = target  # 目标 Agent 名称或队列名

# 示例规则:将飞书渠道且包含“工单”关键词的消息路由到工单 Agent
rule1 = RouteRule(
    predicate=lambda req: req.channel =="feishu"and"工单"in req.content,
    target="agent-workorder"
)
# 默认路由
default_rule = RouteRule(lambda_:True, target="agent-default")

路由器会按顺序匹配规则,命中即停止。同时支持基于用户标签(VIP、地域)的精确路由。

四、代码实战:构建可运行的多渠道 Gateway

我们将实现一个支持 WebSocket(模拟 WebChat) 和 HTTP Webhook(模拟飞书) 双渠道的 Gateway。它具备会话管理、JWT 认证、Redis 消息总线、路由分发等完整能力。你可以直接运行以下代码,并使用 WebSocket 客户端和 curl 分别测试两个渠道。

4.1 项目结构

eclaw-gateway/
├── main.py              # 启动入口
├── gateway/
│   ├── server.py        # FastAPI 应用与 WebSocket/HTTP 端点
│   ├── adapters/
│   │   ├── base.py      # ChannelAdapter 协议
│   │   ├── websocket_adapter.py
│   │   └── feishu_webhook_adapter.py
│   ├── auth.py          # JWT 验证与租户提取
│   ├── session.py       # Redis 会话管理
│   ├── router.py        # 规则引擎
│   └── bus.py           # Redis 消息总线封装
├── agent_stub.py        # 模拟下游 Agent 服务
└── requirements.txt

4.2 核心依赖

pip install fastapi uvicorn redis pyjwt websockets

4.3 认证模块 (auth.py)

我们实现一个极简的 JWT 验证器,实际企业环境可对接 LDAP/OAuth2。

import jwt
from fastapi import HTTPException, WebSocketException

SECRET_KEY ="eclaw-secret-2026"
ALGORITHM ="HS256"

defverify_token(token:str)->dict:
try:
        payload = jwt.decode(token, SECRET_KEY, algorithms=[ALGORITHM])
return payload  # 应包含 user_id, tenant_id
except jwt.PyJWTError:
raise HTTPException(status_code=401, detail="Invalid token")

defverify_ws_token(token:str)->dict:
try:
return jwt.decode(token, SECRET_KEY, algorithms=[ALGORITHM])
except jwt.PyJWTError:
raise WebSocketException(code=4001, reason="Invalid token")

4.4 渠道适配器示例

WebSocket 适配器 (websocket_adapter.py)

import json
from.base import ChannelAdapter
from dataclasses import dataclass, field

@dataclass
classAgentRequest:
    session_id:str
    user_id:str
    tenant_id:str
    content:str
    channel:str
    metadata:dict= field(default_factory=dict)

classWebSocketAdapter(ChannelAdapter):
    channel_name ="webchat"

asyncdefparse_incoming(self, raw_data:bytes, metadata:dict)-> AgentRequest:
        msg = json.loads(raw_data)
return AgentRequest(
            session_id=metadata.get("session_id"),
            user_id=metadata.get("user_id"),
            tenant_id=metadata.get("tenant_id"),
            content=msg["content"],
            channel=self.channel_name,
            metadata=metadata
)

asyncdefformat_outgoing(self, event:dict)->str:
return json.dumps(event)

飞书 Webhook 适配器 (feishu_webhook_adapter.py)

import json
from.base import ChannelAdapter
from.websocket_adapter import AgentRequest

classFeishuWebhookAdapter(ChannelAdapter):
    channel_name ="feishu"

asyncdefparse_incoming(self, raw_data:bytes, metadata:dict)-> AgentRequest:
# 模拟飞书的回调格式,真实场景需处理签名验证
        payload = json.loads(raw_data)
return AgentRequest(
            session_id=payload.get("session_id"),
            user_id=payload.get("open_id"),
            tenant_id=metadata.get("tenant_id","default"),
            content=payload.get("text",{}).get("content",""),
            channel=self.channel_name,
            metadata={"raw": payload}
)

asyncdefformat_outgoing(self, event:dict)->str:
# 飞书需要特定的卡片格式,这里简化返回文本
if event.get("type")=="final":
return json.dumps({"msg_type":"text","content":{"text": event["data"]}})
return""

4.5 会话管理 (session.py)

使用 Redis 存储会话状态,支持过期和跨渠道关联。

import redis.asyncio as aioredis
import uuid
from datetime import timedelta

classSessionManager:
def__init__(self, redis_url="redis://localhost:6379"):
        self.redis = aioredis.from_url(redis_url, decode_responses=True)

asyncdefcreate_session(self, user_id:str, tenant_id:str, channel:str)->str:
        session_id =str(uuid.uuid4())
        key =f"session:{tenant_id}:{session_id}"
await self.redis.hset(key, mapping={
"user_id": user_id,
"channel": channel,
"active":"1"
})
await self.redis.expire(key, timedelta(hours=24))
return session_id

asyncdefget_session(self, tenant_id:str, session_id:str)->dict|None:
        key =f"session:{tenant_id}:{session_id}"
returnawait self.redis.hgetall(key)

4.6 消息总线 (bus.py)

基于 Redis Pub/Sub 实现请求发布与响应订阅。

import redis.asyncio as aioredis
import json

classMessageBus:
def__init__(self, redis_url="redis://localhost:6379"):
        self.redis = aioredis.from_url(redis_url, decode_responses=True)

asyncdefpublish_request(self, queue_name:str, request_dict:dict):
await self.redis.publish(f"agent:req:{queue_name}", json.dumps(request_dict))

asyncdefsubscribe_responses(self, session_id:str):
"""返回一个异步生成器,监听属于此 session 的响应"""
        pubsub = self.redis.pubsub()
await pubsub.subscribe(f"agent:resp:{session_id}")
try:
asyncfor message in pubsub.listen():
if message["type"]=="message":
yield json.loads(message["data"])
finally:
await pubsub.unsubscribe(f"agent:resp:{session_id}")

4.7 路由器 (router.py)

from typing import List
from.adapters.websocket_adapter import AgentRequest

classRuleEngine:
def__init__(self):
        self.rules =[]

defadd_rule(self, predicate, target):
        self.rules.append((predicate, target))

defroute(self, request: AgentRequest)->str:
for predicate, target in self.rules:
if predicate(request):
return target
return"agent-default"

4.8 主服务 (server.py)

这是最关键的整合代码,展示了 Gateway 如何串联所有组件。

from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Request
import asyncio
import json
from.adapters.websocket_adapter import WebSocketAdapter, AgentRequest
from.adapters.feishu_webhook_adapter import FeishuWebhookAdapter
from.auth import verify_ws_token, verify_token
from.session import SessionManager
from.bus import MessageBus
from.router import RuleEngine

app = FastAPI()
bus = MessageBus()
sessions = SessionManager()
router = RuleEngine()
# 添加示例路由规则
router.add_rule(lambda req: req.channel =="feishu","agent-feishu")
router.add_rule(lambda req: req.channel =="webchat","agent-webchat")

@app.websocket("/ws/chat")
asyncdefwebsocket_chat(ws: WebSocket):
await ws.accept()
# 首次消息应为认证包
    auth_msg =await ws.receive_text()
    auth_data = json.loads(auth_msg)
    token = auth_data.get("token")
try:
        claims = verify_ws_token(token)
except:
await ws.close(code=4001)
return
    user_id = claims["user_id"]
    tenant_id = claims["tenant_id"]
# 创建会话
    session_id =await sessions.create_session(user_id, tenant_id,"webchat")
    adapter = WebSocketAdapter()
# 启动后台任务监听 Agent 响应
asyncdefresponse_listener():
asyncfor event in bus.subscribe_responses(session_id):
            formatted =await adapter.format_outgoing(event)
if formatted:
await ws.send_text(formatted)
    listener_task = asyncio.create_task(response_listener())
try:
whileTrue:
            raw =await ws.receive_text()
            metadata ={"session_id": session_id,"user_id": user_id,"tenant_id": tenant_id}
            req =await adapter.parse_incoming(raw.encode(), metadata)
# 路由选择目标
            target = router.route(req)
# 发布到消息总线
await bus.publish_request(target, req.__dict__)
except WebSocketDisconnect:
        listener_task.cancel()

@app.post("/webhook/feishu")
asyncdeffeishu_webhook(request: Request):
# 从请求头获取认证信息(简化演示)
    token = request.headers.get("Authorization","").replace("Bearer ","")
    claims = verify_token(token)
    adapter = FeishuWebhookAdapter()
    body =await request.body()
    metadata ={"tenant_id": claims["tenant_id"]}
# 创建或恢复会话(飞书通常有 open_id,这里简化新建)
    session_id =await sessions.create_session(claims["user_id"], claims["tenant_id"],"feishu")
    metadata["session_id"]= session_id
    metadata["user_id"]= claims["user_id"]
    req =await adapter.parse_incoming(body, metadata)
    target = router.route(req)
await bus.publish_request(target, req.__dict__)
return{"code":0,"msg":"received"}

4.9 模拟 Agent 服务 (agent_stub.py)

为了让整个流程跑起来,我们编写一个简单的 Agent 订阅者,它监听消息总线并返回模拟回复。

import asyncio
import redis.asyncio as aioredis
import json

asyncdefagent_worker(agent_name:str):
    r = aioredis.from_url("redis://localhost:6379", decode_responses=True)
    pubsub = r.pubsub()
await pubsub.subscribe(f"agent:req:{agent_name}")
print(f"[{agent_name}] 启动,等待请求...")
asyncfor msg in pubsub.listen():
if msg["type"]=="message":
            req = json.loads(msg["data"])
print(f"[{agent_name}] 收到请求: {req['content']}")
# 模拟处理并发送响应
            response ={
"type":"final",
"data":f"Agent [{agent_name}] 已处理: {req['content']}",
"session_id": req["session_id"]
}
await r.publish(f"agent:resp:{req['session_id']}", json.dumps(response))

asyncdefmain():
await asyncio.gather(
        agent_worker("agent-webchat"),
        agent_worker("agent-feishu"),
)

if __name__ =="__main__":
    asyncio.run(main())

4.10 启动与测试

  1. 确保本地 Redis 运行在 localhost:6379
  2. 启动 Agent 模拟服务:python agent_stub.py
  3. 启动 Gateway:uvicorn server:app --port 8000
  4. 测试 WebSocket 渠道:
    # 安装 websocat 或使用浏览器控制台
    websocat ws://localhost:8000/ws/chat
    # 先发送认证包
    {"token":"eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.eyJ1c2VyX2lkIjoidXNlcl8xMjMiLCJ0ZW5hbnRfaWQiOiJ0ZW5hbnRfYWJjIn0.xxx"}
    # 再发送业务消息(需生成合法 JWT,可使用 jwt.io 临时生成 payload 包含 user_id,tenant_id)
    {"content":"查询我的订单"}
    # 预期收到 Agent 回复
  5. 测试飞书 Webhook 渠道:
    curl-X POST http://localhost:8000/webhook/feishu \
    -H"Authorization: Bearer <JWT>"\
    -H"Content-Type: application/json"\
    -d'{"session_id":"xxx","open_id":"user_456","text":{"content":"你好,飞书"}}'

📌 运行提示: 为简化演示,JWT 验证允许任意合法签名的 token(可通过 jwt.io 在线生成,payload 包含 {"user_id":"user_123","tenant_id":"tenant_abc"})。实际生产环境必须使用正确的密钥对和过期策略。

五、与 OpenClaw 的对标思考

🔍 “我们做了什么” vs “OpenClaw 为什么这样做”?

渠道抽象:OpenClaw 内置了 20+ 平台的 Channel 适配器,使用 TypeScript 实现,并且通过插件模式管理。我们采用 Python 的Protocol定义适配器接口,虽然目前只实现了两个示例,但完全遵循开闭原则——新增一个渠道只需编写一个 Adapter 文件并注册即可。Python 生态在企业后端更普及,但 TypeScript 在处理大量并发 WebSocket 连接时可能更有优势。不过,FastAPI 的异步能力已经足够应对数千并发。

消息总线:OpenClaw 的内部通信更多是基于事件驱动和进程内传递,而我们引入了 Redis Pub/Sub 作为 Gateway 与 Agent 之间的解耦层。这带来了更好的扩展性,但也增加了网络延迟和运维复杂度。对于中小规模部署,也许直接进程内调用更轻量。这个设计选择体现了“企业级”的考量——你需要为未来的水平扩展预留空间。

安全模型:OpenClaw 作为个人助手,默认信任本机用户。而我们在 Gateway 层强制加入了 JWT 验证、租户隔离和路由级别的权限控制。这是企业 Agent 的刚性需求。我们的认证模块虽然简单,但完全可以替换为 OAuth2 或 LDAP 集成。

会话管理:OpenClaw 通过文件系统或 SQLite 存储会话,而我们的设计直接采用 Redis。这表明了从“单机个人助手”到“分布式多租户系统”的转变。Redis 提供了 TTL、发布订阅等天然优势,但也要求运维一个额外的中间件。

总而言之,我们的 Gateway 继承了 OpenClaw 多渠道统一接入的思想,并在安全、解耦、可观测性方面做了大量企业级增强。下一篇文章,我们将深入 Agent 运行时,实现真正的 ReAct 循环与 LLM 集成。

六、总结与下一步

本文我们完成了:

  1. 深入理解了 Gateway 的六大核心职责,以及为什么它是企业 Agent 的第一道防线。
  2. 设计了基于 Channel Adapter + 消息总线 + 规则路由 的标准化架构。
  3. 用 Python 实现了一个支持双渠道、具备认证和会话管理的可运行 Gateway,并通过模拟 Agent 验证了完整链路。

下一篇文章预告:《手把手构建企业级 Agent 框架:Pi Agent 运行时与 ReAct 循环》。敬请期待!

文章标签智能体
资讯来源:由AI资讯编辑整理自互联网公开内容,版权归原作者所有,未经许可,不得转载。

继续浏览更多资讯

返回资讯目录

相关资讯

更多