中间件与事件驱动
在分布式多 Agent 系统中,除了点对点的 A2A 调用与服务发现外,VeADK 还提供了两类中间件能力:
- A2A 鉴权中间件:在 A2A 服务端透明地处理身份凭证的提取、解析与存储。
- RocketMQ 事件驱动中间件:让 Agent 通过消息队列以事件驱动的方式异步协作。
A2A 鉴权中间件
A2AAuthMiddleware 是一个基于 Starlette/FastAPI 的中间件,用于在 A2A 服务端从入站请求中提取认证令牌,并将凭证写入凭证服务(credential service),供下游处理逻辑使用。它会:
- 从
Authorization请求头或查询参数中提取认证令牌; - 解析 JWT,提取
user_id与委托链(delegation chain); - 根据认证方式构建
AuthConfig,并将凭证存入凭证服务; - 从
X-Ve-TIP-Token请求头提取 TIP 令牌,用于信任传播; - 通过
IdentityClient将 TIP 令牌换取为 workload 访问令牌,并写入request.scope["auth"]。
推荐通过工厂函数 build_a2a_auth_middleware(...) 创建中间件类,再用 FastAPI 的 add_middleware 注册:
from fastapi import FastAPI
from veadk.a2a.ve_middlewares import build_a2a_auth_middleware
from veadk.auth.ve_credential_service import VeCredentialService
app = FastAPI()
credential_service = VeCredentialService()
# Header 方式(Authorization: Bearer <token>)
app.add_middleware(
build_a2a_auth_middleware(
app_name="my_app",
credential_service=credential_service,
auth_method="header",
)
)
# 或 QueryString 方式(?token=<token>)
app.add_middleware(
build_a2a_auth_middleware(
app_name="my_app",
credential_service=credential_service,
auth_method="querystring",
token_param="token",
)
)build_a2a_auth_middleware 的主要参数:
| 参数 | 类型 | 说明 |
|---|---|---|
app_name | str | 用于凭证存储的应用名称。 |
credential_service | VeCredentialService | 存储凭证的凭证服务。 |
auth_method | "header" | "querystring" | 认证方式,默认 header。 |
token_param | str | querystring 方式下令牌所在的查询参数名,默认 token。 |
credential_key | str | 在凭证存储中标识该凭证的键,默认 inbound_auth。 |
identity_client | Optional[IdentityClient] | 用于 TIP 令牌换取的身份客户端;不传时使用全局 IdentityClient。 |
鉴权成功后,中间件会在请求上设置如下属性,供下游处理逻辑读取:
request.scope["user"]:携带从 JWT 解析出的user_id的SimpleUser实例;request.scope["auth"]:包含workload_access_token的WorkloadToken对象。
提示
该中间件与 A2A Agent 中介绍的 RemoteVeAgent 鉴权方式(auth_method 为 header / querystring)相对应:客户端负责携带令牌,服务端中间件负责解析与存储。
RocketMQ 事件驱动中间件
除了同步的请求/响应调用,VeADK 还支持通过 RocketMQ 消息队列以事件驱动的方式连接 Agent。这适用于异步、解耦、可横向扩展的多 Agent 协作场景。
RocketMQClient 封装了消息的发送与消费。建议通过环境变量提供凭证,避免在代码中硬编码:
export ROCKETMQ_ACCESS_KEY="<access_key>"
export ROCKETMQ_ACCESS_SECRET="<access_secret>"import os
from veadk.a2a.hub.rocketmq_middleware import RocketMQClient
client = RocketMQClient(
name="weather_client",
producer_group="weather_group",
name_server_addr="<rocketmq-name-server-addr>",
access_key=os.environ["ROCKETMQ_ACCESS_KEY"],
access_secret=os.environ["ROCKETMQ_ACCESS_SECRET"],
)
# 发送一条单向(one-way)消息到某个 topic
client.send_msg(topic="weather_topic", msg_body="北京天气如何?", key="req-1", tag="query")RocketMQAgentClient 是一个抽象基类,将一个 Agent 绑定到某个订阅 topic 上。实现其抽象方法 recv_msg_callback,即可在收到消息时驱动 Agent 处理,并返回消费状态:
import asyncio
import os
from rocketmq.client import ConsumeStatus, ReceivedMessage
from veadk import Agent, Runner
from veadk.a2a.hub.rocketmq_middleware import RocketMQAgentClient, RocketMQClient
class WeatherAgentClient(RocketMQAgentClient):
def recv_msg_callback(self, msg: ReceivedMessage) -> ConsumeStatus:
# 用消息内容驱动 self.agent,并返回消费状态。
prompt = msg.body.decode("utf-8")
try:
runner = Runner(agent=self.agent)
response = asyncio.run(runner.run(messages=prompt))
print(f"Agent 回复: {response}")
return ConsumeStatus.CONSUME_SUCCESS
except Exception as e:
print(f"处理消息失败,稍后重试: {e}")
return ConsumeStatus.RECONSUME_LATER
rocketmq_client = RocketMQClient(
name="weather_client",
producer_group="weather_group",
name_server_addr="<rocketmq-name-server-addr>",
access_key=os.environ["ROCKETMQ_ACCESS_KEY"],
access_secret=os.environ["ROCKETMQ_ACCESS_SECRET"],
)
agent_client = WeatherAgentClient(
agent=Agent(name="weather_agent"),
rocketmq_client=rocketmq_client,
subscribe_topic="weather_topic",
group="weather_consumer_group",
)
# 开始监听订阅 topic(阻塞)
agent_client.listen()调用 listen() 后,客户端会在指定 topic 上启动消费者,每当有新消息到达,便回调 recv_msg_callback 来驱动 Agent 处理,从而实现 Agent 之间的异步、事件驱动协作。返回 CONSUME_SUCCESS 表示消费成功,返回 RECONSUME_LATER 则会触发消息重投。