跳到主要内容

接入层:多渠道 drivers、统一队列、MCP

30 秒导读: QwenPaw 的 agent 内核只认一种输入(AgentRequest)和一种输出(Event 流)。本章讲清两件事:左边,十几种 IM(飞书、钉钉、Discord、Telegram……)怎么把各自的原生消息标准化成这一种输入,再经一套统一优先级队列喂给内核;右边,外部工具怎么通过 driver / MCP 接进来,受统一策略引擎 + 加密凭据 + OAuth/白名单管控后,变成 agent 能调的 tool。


1. 这是什么(零基础也能懂)

一句话定义: 接入层是 agent 内核的两副插座——一副朝用户(渠道 channels),一副朝工具(drivers)。

先记住一个前提:agent 内核(见 第 01 章)不关心消息从飞书来还是从 Telegram 来,也不关心工具是本地函数还是远程 MCP 服务。它只吃 AgentRequest,只吐 Event 流。接入层的全部职责,就是把五花八门的外部世界,翻译成内核认得的这一种语言。

这一层解决的两个真实痛点:

痛点场景化白话接入层怎么答
一个 agent 想同时上十几个 IM你写好一个客服 bot,老板要它同时在飞书、钉钉、企业微信、Discord 上线渠道抽象:每个 IM 只实现"怎么收、怎么发",标准化逻辑复用父类
想让 agent 用别人家的工具你想让 agent 能查 GitHub、操作数据库,但这些能力在别的进程/别的服务里driver + MCP:把外部工具注册成受管控的 capability,再暴露成 tool

它能做什么(功能清单):

  • 把 18 种内置渠道的原生消息,统一标准化成 AgentRequest(飞书群消息、钉钉卡片、语音……)。
  • 按"渠道 × 会话 × 优先级"三元组分队列,不同会话并发、同会话严格串行,/stop 这类控制命令走高优先级插队。
  • 渠道级访问控制:黑白名单 + 待审批(pending),陌生人第一条消息自动挂起等放行。
  • 把外部 MCP 服务器注册成 driver,其工具经策略引擎(allow/ask/deny)、加密凭据、OAuth、工具白名单层层过滤后,变成 agent 可调的 tool。

一句话直觉/类比: 把接入层想成跨国机场的两道翻译闸机。入境闸(渠道)不管你说哪国话,统一翻成"标准语"送进城;出境闸(driver)不管城里要用哪国的服务,先查签证(策略)、验证件(凭据),放行了才让接触。城内(agent 内核)永远只讲标准语。


2. 顶层全景(它大概怎么转)

本节讲"大盘":一条消息进来、一次工具调用出去,分别流经哪些部件。

2.1 两条主线一张图

先看入站主线(用户消息 → 内核)和工具主线(内核要调工具 → 外部服务)。图从左到右是数据流向, 是部件, 是流向。

入站主线(渠道 channels)
┌────────┐ 原生消息 ┌──────────────┐ AgentRequest ┌──────────────┐
│ 飞书/钉钉│ ───────────▶ │ 具体 Channel │ ─────────────▶ │ UnifiedQueue │
│ /TG/... │ webhook/ws │ 标准化+ACL闸 │ 入队(带优先级)│ 三元组分队列 │
└────────┘ └──────────────┘ └──────┬───────┘
串行消费 │

┌────────────┐
│ agent 内核 │
│ (第01章) │
└─────┬──────┘
工具主线(drivers) │ Event 流
┌────────┐ call_tool ┌──────────────┐ invoke ▼(回发经 Channel.send)
│ 外部MCP │ ◀─────────── │ MCPDriverHandler│ ◀────────── ┌────────────┐
│ 服务器 │ ───────────▶ │ 策略+凭据闸 │ ───────────▶ │ DriverManager│
└────────┘ 结果 └──────────────┘ capability │ 注册/分发 │
└────────────┘

怎么读:上半是"用户 → agent",下半是"agent → 工具"。两侧各有一道"闸"——入站的 ACL 闸、工具侧的 策略+凭据闸。城中央(agent 内核)只跟标准契约打交道。

2.2 部件一句话职责

渠道侧(app/channels/):

部件干什么文件
BaseChannel所有渠道的父类:定义"入站标准化 + 出站渲染 + ACL + 流式"的骨架app/channels/base.py:80
渠道注册表懒加载 18 个内置渠道 + 发现工作目录里的自定义渠道app/channels/registry.py:193
ChannelManager拥有队列与消费循环;对外提供线程安全的 enqueueapp/channels/manager.py:68
UnifiedQueueManager按三元组 QueueKey 动态建队列、动态起消费者、闲置回收app/channels/unified_queue_manager.py:60
CommandRegistry/stop/status 等命令映射到优先级 0/10/20/30app/channels/command_registry.py:23
MessageRenderer把内核 Message 渲染成可发送的 content parts(markdown/emoji 可配)app/channels/renderer.py:78
AccessControlStore每渠道的黑白名单 + 待审批,持久化到 JSONapp/channels/access_control.py:157

工具侧(drivers/app/mcp/):

部件干什么文件
DriverManager拥有外部能力的存储、生命周期、分发;协议中立drivers/manager.py:47
DriverHandler模板方法基类:策略评估 + 凭据解析 + 审批,子类填协议细节drivers/handler.py:42
MCPDriverHandler具体协议实现:连 MCP 服务器,把其 tools 暴露成 capabilitydrivers/handlers/mcp.py:51
evaluate_policy策略引擎:对一次工具调用返回 allow / ask / denydrivers/policy.py:77
AsyncCredentialStore每工作区的 YAML 凭据库,secret 全程加密落盘drivers/credentials/store.py:22
DriverCapabilityTool适配器:把一个 capability 包成 AgentScope 的 ToolBasedrivers/adapters/agentscope_tool.py:135
MCPConfigServiceConsole 管理 MCP 客户端:增删改、工具白名单、访问策略app/mcp/config_service.py:74

2.3 主线走一遍(高层,不进代码)

入站: 飞书推来一条群消息 → FeishuChannel 解析出 content_parts 和会话标识 → ChannelManager.enqueue 按命令算优先级、按会话算 session_id,丢进 UnifiedQueueManager → 对应队列的消费者取出、标准化成 AgentRequest、过 ACL 闸 → 交给内核 _process → 内核吐 Event 流 → MessageRenderer 渲染 → Channel.send 回发。

工具: 内核准备工具时,build_driver_agent_toolsDriverManager 要所有 active driver 的 capability,每个包成 ToolBase → 模型决定调用某工具 → DriverCapabilityTool.__call__DriverManager.invoke_capability 按 capability_id 路由到 MCPDriverHandler → 策略引擎判 allow/ask/deny(ask 走审批闸,见 第 05 章)→ 解析加密凭据 → 调远程 MCP tool → 结果转回内核。


3. 核心原理(逐个机制,由浅入深)

3.1 渠道抽象:一条入站消息如何被标准化

它要解决的小问题: 飞书的消息是一坨 JSON,钉钉又是另一坨,内核只认 AgentRequest。谁来翻译,翻译逻辑怎么不重复写 18 遍?

思路: 把"每个渠道都一样"的部分(建 AgentRequest、去抖、ACL、渲染、流式)全部沉到父类 BaseChannel;把"每个渠道不同"的部分(解析原生 payload、send 怎么发)留成抽象方法给子类填。这是经典的模板方法

父类定的关键契约:

方法谁实现职责
build_agent_request_from_native子类必填原生 payload → AgentRequest(base.py:1153 默认抛 NotImplementedError)
resolve_session_id子类可覆写发送者 + 会话元数据 → 会话键(默认 f"{channel}:{sender_id}",base.py:1105)
send子类必填把一段文本(+附件)发到 to_handle(base.py:1887)
build_agent_request_from_user_content父类提供content_parts 包成标准 Message/AgentRequest(base.py:1117)

一条飞书消息的标准化(真实实现):

FeishuChannel.build_agent_request_from_native 把飞书原生 dict 拆成三样——content_parts(内容)、session_id(会话)、user_id(发送目标),再交给父类的通用组装器:

# app/channels/feishu/channel.py:402 build_agent_request_from_native(节选,已简化)
content_parts = payload.get("content_parts") or []
session_id = payload.get("session_id") or self.resolve_session_id(sender_id, meta)
user_id = meta.get("feishu_sender_id") or payload.get("user_id") or sender_id
request = self.build_agent_request_from_user_content(
channel_id=channel_id, sender_id=user_id,
session_id=session_id, content_parts=content_parts, channel_meta=meta,
)

重点看:飞书覆写了 resolve_session_id(channel.py:381)——群聊用 chat_id 的短后缀 + app_id 后缀(区分同群多 bot),私聊用 open_id。这样"同一个会话"的定义才对得上,后面队列串行化和 cron 回发才不串台。

共性抽出来的价值: 一个新渠道基本只需实现 build_agent_request_from_native + send + start/stop,其余(去抖合并、ACL、优先级、流式卡片)全白拿。18 个内置渠道就是这么共享一套骨架的。

3.2 渠道注册:内置懒加载 + 自定义热发现

它要解决的小问题: 18 个渠道各自依赖不同的第三方 SDK(飞书 SDK、discord.py……),用户大概率只装了其中两三个。一个渠道 import 失败,不能拖垮整个 CLI 启动。

思路: 注册表用懒加载 + 容错。逐个 import,失败的记 debug 日志跳过,只有 console 这个必需渠道失败才真的抛错。

# app/channels/registry.py:48 _load_builtin_channels(节选)
for key, (module_name, class_name) in _BUILTIN_SPECS.items():
try:
mod = importlib.import_module(module_name, package=__package__)
cls = getattr(mod, class_name)
...
except Exception:
if key in _REQUIRED_CHANNEL_KEYS: # 只有 {"console"}
raise
continue # 其它渠道缺依赖就跳过

内置清单在 _BUILTIN_SPECS(registry.py:20),覆盖:imessage / discord / dingtalk / feishu / qq / telegram / mattermost / mqtt / console / matrix / slack / voice / sip / wecom / xiaoyi / yuanbao / wechat / onebot(源码里还有 SIP/语音等电话类渠道)。结果缓存一次(_get_cached_builtin_channels,registry.py:84)。

自定义渠道从工作目录的 CUSTOM_CHANNELS_DIR 动态发现(_discover_custom_channels,registry.py:100):任何 BaseChannel 子类,只要类上有 channel 属性(渠道 key),就自动注册。自定义渠道还能通过模块级 register_app_routes(app) 钩子挂自己的 HTTP 路由(如二维码登录页),但路由必须挂在 /api/ 前缀下,否则会被 SPA 的 catch-all 吞掉(register_custom_channel_routes,registry.py:138)。

ChannelType 干脆就是 str(schema.py:48)——内置是固定几个,插件渠道用任意字符串 key,这样类型系统不挡插件。

3.3 统一优先级队列:三元组 QueueKey

它要解决的小问题: 三个矛盾要同时满足——(1) 不同用户的消息要并发处理,不能互相排队;(2) 同一个会话的消息必须严格串行,不能乱序;(3) /stop 这种控制命令要能插队,不能排在一长串正常消息后面。

思路: 把队列的粒度切到 三元组 QueueKey = (channel_id, session_id, priority_level)(unified_queue_manager.py:31)。每个三元组一条独立 asyncio.Queue + 一个独立消费者协程。于是:

  • 不同 session_id → 不同队列 → 天然并发。
  • session_id 同优先级 → 同一队列 → 天然串行。
  • /stop(优先级 0)和正常消息(优先级 20)→ 不同队列 → 高优先级不被堵。

优先级从命令前缀查表得来(CommandRegistry,command_registry.py:23):

级别名字典型命令
0critical/stop
10high/status /restart /approve /deny
20normal普通对话(默认)
30low批处理(预留)

动态起消费者、闲置回收: 队列和消费者都是按需创建的(_get_or_create_queue,unified_queue_manager.py:165)——第一条消息到达某三元组时才建。后台 _cleanup_idle_queues(unified_queue_manager.py:376)每分钟扫一次,空且闲置超 10 分钟的队列连带消费者一起回收。没有固定 worker 池,省内存。

入队全流程(高层):

enqueue(channel_id, payload) [线程安全,可从 ws/轮询线程调]
└─ call_soon_threadsafe → _enqueue_one manager.py:260
├─ query = 取文本 → 算 priority_level (命令查表)
├─ session_id = 归一化会话键 (等价 debounce_key)
└─ UnifiedQueueManager.enqueue(...) 带 30s 超时保护
└─ _get_or_create_queue → queue.put
└─ 消费者 _consume_queue manager.py:367
├─ 排空同键队列 → batch (合并连发的图片等)
└─ _process_batch → 合并 → 标准化 → 内核

enqueue 本身线程安全(manager.py:354):它用 loop.call_soon_threadsafe 把真正的入队动作弹回事件循环,所以 WebSocket 回调线程、轮询线程都能直接调。

批量合并(batch merge): 消费者取出第一条后,会把同队列里已到的消息一次性排空成一个 batch(manager.py:423 附近),再 _process_batch(manager.py:39)。多条连发的原生消息(比如用户连发三张图)会被 merge_native_items 合并成一次请求——这是"发一半没发完"的去抖思路的队列版。

3.4 非阻塞流式推送:边生成边刷新气泡

它要解决的小问题: 模型是流式吐字的。想让用户看到"逐字出现"的打字机效果,又不能让每个 delta 都阻塞住事件循环(网络慢的渠道一卡,整条流就断)。

思路: 父类在 _stream_with_tracker(base.py:859)里对每个 Event 分两条路走——先照常走"消息完成"的老路,若 streaming_enabled 再额外派发流式钩子。流式刷新用发射后不管 + 在途守卫:

# app/channels/base.py:776 附近 _on_stream_content_delta(节选,已简化)
# 上一次刷新还在飞 → 本次直接跳过,避免堆积
task = flush_meta.get("task")
if task and not task.done():
...
return True
# 否则 fire-and-forget 起一个刷新任务,不 await
flush_meta["task"] = asyncio.create_task(
self._safe_streaming_delta(request, to_handle, event, send_meta,
stream_type, streaming_buffers[stream_type]),
)

三个钩子留给子类覆写:on_streaming_start / on_streaming_delta / on_streaming_end(base.py:1540 起)。像企业微信这种"用全量文本覆盖同一条气泡"的渠道,就在 on_streaming_delta 里拿 accumulated_text 去改消息。渠道靠 streaming_enabled 开关(类属性或 from_config 里配)决定要不要走这条路(base.py:119)。

主动推送 / 非对话回发: 除了应答,内核还需要在没有用户消息时主动发东西(cron 定时任务、审批通知)。ChannelManager 提供 send_text(manager.py:812)和 push_approval_notification(manager.py:865),把 (user_id, session_id)to_handle_from_target 翻成渠道目标再发。控制台这种无长连接的渠道,则把推送暂存到内存的 console_push_store(app/console_push_store.py:22append / :41take),前端轮询取走——有界(最多 500 条、超 60 秒丢弃),按 id 去重。

3.5 渠道访问控制:黑白名单 + 待审批

它要解决的小问题: bot 上线后是公开的,任何人加了就能发消息。得能"只允许白名单"、"拉黑骚扰者"、"陌生人先挂起等我批"。

思路:_consume_one_request 里、真正跑内核之前,插一道统一 ACL 闸 _access_control_gate(base.py:374)。返回 True(拦截)就直接 return,不进内核。判定顺序:

收到消息 sender_id

├─ ACL 未开(dm/group 都没开) → 放行
├─ 在白名单 → 放行
├─ 在黑名单 → 回一句"已被禁止",拦截
└─ 都不在(陌生人) → 记入 pending + 回"需审批,ID:xxx",拦截

拒绝话术按 agent 语言本地化(_ACL_I18N,base.py:336,含 zh/en/ja/ru/pt-BR/id)。判定用真实发送者 acl_sender_id(不受共享会话影响)。

持久化: AccessControlStore(access_control.py:157)按渠道存 whitelist / blacklist / pending 三份到 JSON,线程安全,并且检测文件 mtime 变化自动重载(_reload_if_stale,access_control.py:182)——因为工作区热重载可能重建 store,而旧引用还在写。审批动作齐全:add_pending(:361)挂起、approve_pending(:408)转白名单、deny_pending(:442)转黑名单,批准时自动带过 username。

二维码登录: 部分渠道(微信/企微/钉钉/飞书)的 bot 授权走扫码。qrcode_auth_handler.py 抽象了 QRCodeAuthHandler(:55,两个方法 fetch_qrcode / poll_status),各渠道实现自己的设备码流程(如 DingtalkQRCodeAuthHandler:277FeishuQRCodeAuthHandler:404 用 OAuth Device Flow),回调 token 用 AES-256-GCM 加密成无状态串(:602 起)。

技能语义与工具的权限审批引擎不在本章: 技能怎么变成 tool 见 第 03 章;工具调用的审批引擎(ApprovalGate 背后的人机交互)见 第 05 章。本章的 driver 侧只讲"接入 + 策略判定",判定为 ask 之后怎么弹给人、怎么恢复,是第 05 章的事。


4. 深入实现:driver 层与 MCP 接入

本节讲右半边——外部工具怎么接进来。先建立心智模型,再逐个部件走读。

4.1 心智模型:DriverCard 是"一张外部服务的身份证"

一个外部 MCP 服务器,在 QwenPaw 里被建模成一张 DriverCard(drivers/contracts.py:103)。它是声明式的、不含明文密钥的配置:

字段装什么
namedriver 唯一名(也是存储文件名,校验不含路径分隔符,contracts.py:48)
protocol协议,目前恒为 "mcp"(constants.py:27)
endpoint连接方式:stdio 的 command/args/env,或 http 的 url/headers/transport
credentials别名 → CredentialRef(只是指针,指向凭据库里的条目,不含密钥本体)
config展示名、描述、工具白名单 tools
policy访问策略(rules + default_effect)
enabled是否启用

关键设计:卡里没有密钥endpoint.env / endpoint.headers 里的敏感值以 {"source": "credential", "credential": "static", "field": "api_key"} 这种绑定引用存在(校验见 _validate_endpoint_bindings,contracts.py:227),真正的密文另存于加密凭据库。卡可以随便看、随便备份,不泄密。

4.2 DriverManager:生命周期与分发

DriverManager(drivers/manager.py:47)是外部能力的总管,协议中立——它不懂 MCP,只懂"有一种 handler 类型、按协议注册"。启动时扫描 cards 目录,把每张 enabled 的卡构建成一个 handler(build_drivers,manager.py:89)。

构建 handler 的关键一步(_build_handler,manager.py:311):按卡里的 credentials 别名,为每个别名建一个凭据 provider,再把卡、主 provider、全部 provider、审批闸一起塞给 handler 构造器。build-before-swap:注册/重载时先把新 handler 初始化成功,再换掉旧的(register_driver,manager.py:137;reload_driver,manager.py:153),失败则保留旧的、不中断服务。

一次工具调用怎么路由到对的 handler: capability 有个稳定 id,形如 driver://mcp/<driver>/tools/<tool>#invoke(format_capability_id,capabilities.py:89)。invoke_capability(manager.py:253)从 id 里 parse 出 driver 名,查到 handler,转发过去。id 用 URL 编码,行号漂移也不影响路由——路由只认这个字符串。

4.3 DriverHandler:策略 + 凭据 + 审批的模板方法

DriverHandler(drivers/handler.py:42)是 handler 基类,用模板方法把"每种协议都一样"的三件事固化:

invoke 一次能力

├─ _authorize_invocation handler.py:102
│ └─ evaluate_policy → allow / ask / deny
│ ├─ deny → 抛 DriverPermissionDeniedError
│ └─ ask → _request_approval → ApprovalGate(第05章)

├─ 解析凭据(provider.resolve,可能解密/刷新)

└─ _execute(credential, context, **kwargs) [子类填协议细节]

ApprovalGate 只是个 Protocol(drivers/approval.py:11)——driver 核心只负责判出"要审批",至于怎么问人,由应用外壳注入实现。这是核心与产品 UI 的干净解耦。

4.4 MCPDriverHandler:把 MCP tools 变成 capability

MCPDriverHandler(drivers/handlers/mcp.py:51)是目前唯一的具体协议实现。

连接(_setup,mcp.py:60):endpoint.transport 分两种客户端——stdioStdIOStatefulClient(拉起子进程),streamable_http/sseHttpStatefulClient。连接前先解析凭据,把绑定的 env/header 填进去(resolve_binding),HTTP 还会从凭据推断 Authorization 头(implicit_auth_headers,见 4.6)。

发现(list_capabilities,mcp.py:130): 拉远程 tools 列表,每个转成一个 DriverCapability,带 10 秒缓存。转换里有两处工程细节值得记(_mcp_tool_to_capability,mcp.py:324):

  • 工具名消毒:OpenAI 要求工具名匹配 ^[a-zA-Z0-9_-]+$,而 MCP 工具名可能带别的字符。_sanitize_tool_name(mcp.py:399)把非法字符替成 _,但 capability_id 里保留原始名(URL 编码),这样 invoke 路由回服务端时用的还是真名。
  • 命名空间:暴露给模型的工具名是 {display_namespace}__{sanitized},namespace 从 display_name 推导。

调用(invoke_capability,mcp.py:154): 先 parse + 校验 capability_id(协议/driver/kind/action 四要素都得对),再走 _guarded_execute(mcp.py:230)——即 4.3 那套"授权 → 解析凭据 → _executeclient.call_tool"。各类失败(策略拒绝、要审批、执行错)都转成结构化的 DriverInvocationResult,不抛裸异常给上层。

4.5 策略引擎:一次调用判 allow / ask / deny

它要解决的小问题: "允许所有人用这个 MCP,但 delete_repo 这个工具只有我本人能用、且要二次确认,其它人一律拒绝"——这种细粒度授权怎么表达和计算?

思路: 策略是一组 PolicyRule,每条声明 (subject, principal, target, condition) → effectevaluate_policy(drivers/policy.py:77)对一次调用:先筛出所有匹配的 rule,没有匹配就用 default_effect(缺省 deny);有多条匹配则按具体度排序取最具体的一条

具体度排序的维度(从高到低,policy.py:97 附近):目标名精确 > 目标类型精确 > principal 越具体 > subject 越具体 > 效力越严(deny > ask > allow)。这样"针对某工具某用户"的规则一定压过"针对全部"的兜底。

匹配支持三种通配(subject_matches,policy.py:125):精确 user:alice、类型前缀 user:*、全局 *。principal 还能按来源渠道 + 用户/会话作用域做 AND 匹配(principal_matches,policy.py:114),condition 支持时间段(周几 + 时钟区间,_time_range_satisfied,policy.py:283)。

一句话:策略引擎是"最具体规则胜出、默认拒绝"的本地判定器,判完只给一个字——allow / ask / deny。

4.6 凭据:加密落盘 + 运行时绑定

它要解决的小问题: MCP 服务器要 API key / OAuth token,这些密钥不能明文躺在配置里,也不能在日志和 Console 里裸奔。

加密落盘: AsyncCredentialStore(drivers/credentials/store.py:22)是每工作区一个 YAML 库。写入时 secrets 下每个字符串值都过 encrypt(_encrypt_secrets,store.py:196),读出时 decrypt。文件权限尽力收到 0o600(_restrict_file_permissions,store.py:211,Windows 不适用则靠加密本身兜底)。env:VAR 形式的 ref 直接读环境变量、不落盘(_get_sync,store.py:53)。同步 IO 全部 to_thread 隔离,不阻塞事件循环。

运行时绑定: 卡里的 env/header 用绑定描述,连接时 resolve_binding(drivers/credentials/bindings.py:35)把 source: credential 的项从解密后的凭据里取值填回。HTTP 传输还有 implicit_auth_headers(bindings.py:67):若 header 里没写 Authorization,就从凭据推断——有 access_token 就填 Bearer,有 user/pass 就填 Basic

公开值 vs 密钥的自动分类: 用户在 Console 填 env/header 时,classify_mcp_binding(drivers/adapters/mcp_binding.py:58)按启发式判每个值是 public 还是 secret——header 里 authorization/cookie/x-api-key 判密、env 里 key 含 KEY/TOKEN/SECRET/PASSWORD/AUTH 等判密(SECRET_ENV_KEY_PARTS,mcp_binding.py:33),其余保守判密。public 存成卡里的字面量,secret 存进加密库。Console 回显时密钥经 mask_mcp_secret_value(mcp_binding.py:200)打码。

4.7 OAuth 与工具白名单

OAuth(2.1)概念: 远程 MCP 若要求 OAuth,首次连接会碰 HTTP 401。stateful client 的生命周期循环检出 401 后快速失败并置标志,而不是死重连:

# drivers/handlers/mcp_stateful_client.py:190 附近 _run_lifecycle(节选)
if _is_401_error(e):
logger.info(f"MCP client '{self.name}': server requires OAuth (HTTP 401). "
"Authorize via the UI to connect.")
self._oauth_required = True
self._stop_event.set(); self._ready_event.set()
return

之后 connect 看到 _oauth_required 就抛错,提示去 UI 授权(mcp_stateful_client.py:213 起)。OAuth 凭据用专门的别名 oauth / kind oauth2_auth_code(constants.py:10:14),和静态凭据 static 分开存,授权码换来的 token 落进加密库,连接时经 implicit_auth_headers 注入。

工具白名单: 一个 MCP 服务器可能暴露几十个工具,用户往往只想开其中几个。卡的 config["tools"] 存白名单(None = 全开)。MCPConfigService.list_tools(app/mcp/config_service.py:128)把服务端工具列表和白名单对齐,给每个工具标 enabled;update_tool_whitelist(:163)改白名单。只有 enabled 的工具最终会进 agent 的工具集。

Console 管控入口: MCPConfigService(config_service.py:74)是 Console 增删改 MCP 客户端的应用服务——建客户端时把公开/密钥拆分、密钥写加密库(create_client,:258);改访问策略时把 Console 友好的 MCPAccessPolicy(工具默认 + 按人/按工具覆盖)翻译成底层 DriverPolicy 的 rule 列表(driver_policy_from_mcp_access_update,:452),反向再翻回展示(mcp_access_policy_from_card,:422)。客户端 key 还禁止撞保留前缀(validate_client_key,:356)。

4.8 适配器:capability → AgentScope ToolBase

最后一跳:capability 怎么变成内核能调的 tool。DriverCapabilityTool(drivers/adapters/agentscope_tool.py:135)把一个 DriverCapability 包成 AgentScope 的 ToolBase,__call__ 里构造 DriverInvocation 交给 DriverManager.invoke_capability,再把结果转成 ToolChunk

注意它的 check_permissions(agentscope_tool.py:161)直接返回 ALLOW——因为权限已经在 driver 的策略引擎里判过了,不在 AgentScope 层重复判。build_driver_agent_tools(agentscope_tool.py:182)是内核 stream_query 装配工具时的入口:向 manager 要所有 as_tool 的 capability,批量包成 tool,顺带带回一条"策略可能变化、需复查"的 prompt 提示。


5. 巧妙之处(可借鉴的技术)

每条先说妙在哪,再给锚点。

  • 一种输入、一种输出收敛两侧复杂度。 内核只认 AgentRequest / Event,渠道多样性锁在 BaseChannel 子类里、工具多样性锁在 DriverHandler 子类里。加渠道加工具都不动内核。(base.py:80drivers/handler.py:42)

  • 三元组队列一招解决"并发 + 串行 + 插队"。 把队列粒度切到 (channel, session, priority),三个看似冲突的诉求各自落到"不同键 = 不同队列"上,不需要显式加锁或调度器。(unified_queue_manager.py:31)

  • 按需建队列、闲置回收,零固定 worker 池。 消费者随第一条消息动态生成,10 分钟无活动自动清。海量会话下内存不爆。(unified_queue_manager.py:165:376)

  • DriverCard 无密钥。 密钥以引用形式在卡、密文本体在加密库,卡可自由分享/备份不泄密;public/secret 还能按名字启发式自动分类。(contracts.py:227mcp_binding.py:58)

  • capability_id 是 URL,路由抗行号漂移。 driver://mcp/<driver>/tools/<tool>#invoke 让分发只认字符串;工具名对模型消毒、对服务端保原名,两头都不出错。(capabilities.py:89mcp.py:399)

  • build-before-swap 的热重载。 新 handler / 新渠道先起成功再换旧的,失败保留旧的,重载不中断在线服务。(manager.py:137app/channels/manager.py:724)

  • 策略"最具体胜出、默认拒绝"。 细粒度授权(某工具某人某时段)用可排序的具体度自然压过兜底规则,语义直观。(policy.py:77)


6. 边界与局限(诚实)

  • 协议目前只有 MCP。 DriverManager 号称协议中立,但仓库里唯一注册的具体协议是 MCP(PROTOCOL_MCP,constants.py:27);capability kind 也只有 tool(capabilities.py:15Literal["tool"]),ACP/A2A、resource/prompt 等都是预留、未实现。

  • 审批的人机交互不在 driver 核心。 driver 只判出 ask 并调 ApprovalGate(drivers/approval.py:11);没有注入 gate 时直接抛 ApprovalRequiredError。怎么弹给人、怎么恢复见 第 05 章

  • console 推送是内存态、易失。 console_push_store 有界(500 条 / 60 秒)且进程内(app/console_push_store.py:18),重启即丢,不适合当可靠消息队列。

  • 入队有超时上限。 队列满时 put 等 30 秒仍失败就丢弃并告警(unified_queue_manager.py:147);极端积压下会丢消息而非无限缓冲。

  • ACL 粒度到"渠道 + 发送者 id"。 群里按发送者判,但"同群不同话题"这类更细的作用域不由 ACL 表达。

  • 凭据加密强度取决于 secret_store 落盘加密依赖 security/secret_store 的密钥管理(store.py:17),本章未展开其密钥来源;文件权限在 Windows 上不构成边界,靠加密兜底(store.py:211)。


7. 横向对比(同 shelf 内的取舍)

  • 与本项目其它章:入站消息进内核后的 8 阶段编排见 01-request-lifecycle;工具/技能的语义与 skill→tool 见 03-skills;工具守卫、沙箱、审计与审批引擎见 05-security

  • 取舍特点:QwenPaw 把"多渠道"和"多工具"都用同一种收敛思路(统一契约 + 模板方法 + 声明式配置)处理,渠道侧还额外做了优先级队列渠道级 ACL,这在偏单渠道(如只做 Web/CLI)或把工具直接硬编码进 agent 的同类项目里通常没有。driver 侧把策略、凭据、审批三者拆成可独立替换的边界(Protocol 注入),偏库化设计。


8. 代码地图(导航索引)

用符号名 grep 比行号抗漂移。

渠道抽象与注册:

主题文件路径符号名
渠道父类(标准化/渲染/ACL/流式)app/channels/base.pyBaseChannel
入站标准化契约app/channels/base.pybuild_agent_request_from_native / build_agent_request_from_user_content
会话键解析app/channels/base.pyresolve_session_id
流式派发app/channels/base.py_stream_with_tracker / _on_stream_content_delta
内置懒加载 + 自定义发现app/channels/registry.py_load_builtin_channels / _discover_custom_channels / get_channel_registry
自定义渠道挂路由app/channels/registry.pyregister_custom_channel_routes
路由地址/类型app/channels/schema.pyChannelAddress / ChannelType / BUILTIN_CHANNEL_TYPES
具体渠道示例(会话键/标准化)app/channels/feishu/channel.pyFeishuChannel.resolve_session_id / build_agent_request_from_native

队列、渲染、访问控制:

主题文件路径符号名
队列/消费循环拥有者app/channels/manager.pyChannelManager / enqueue / _consume_queue / _process_batch
主动推送/审批推送app/channels/manager.pysend_text / push_approval_notification
三元组优先级队列app/channels/unified_queue_manager.pyUnifiedQueueManager / QueueKey / _get_or_create_queue / _cleanup_idle_queues
命令→优先级app/channels/command_registry.pyCommandRegistry / get_priority_level / is_control_command
消息渲染app/channels/renderer.pyMessageRenderer / RenderStyle / message_to_parts
黑白名单/待审批app/channels/access_control.pyAccessControlStore / add_pending / approve_pending
控制台推送暂存app/console_push_store.pyappend / take
扫码授权app/channels/qrcode_auth_handler.pyQRCodeAuthHandler / DingtalkQRCodeAuthHandler / FeishuQRCodeAuthHandler

driver 层与 MCP:

主题文件路径符号名
driver 总管/生命周期/分发drivers/manager.pyDriverManager / build_drivers / register_driver / invoke_capability
handler 模板(策略/凭据/审批)drivers/handler.pyDriverHandler / _authorize_invocation / _request_approval
审批边界drivers/approval.pyApprovalGate
capability 契约/iddrivers/capabilities.pyDriverCapability / format_capability_id / parse_capability_id
DriverCard/绑定校验drivers/contracts.pyDriverCard / CredentialRef / _validate_endpoint_bindings
策略引擎drivers/policy.pyevaluate_policy / subject_matches / principal_matches
MCP handlerdrivers/handlers/mcp.pyMCPDriverHandler / _setup / invoke_capability / _mcp_tool_to_capability / _sanitize_tool_name
MCP stateful client + OAuth 401drivers/handlers/mcp_stateful_client.py_run_lifecycle / connect / _is_401_error
加密凭据库drivers/credentials/store.pyAsyncCredentialStore / _encrypt_secrets
凭据运行时绑定drivers/credentials/bindings.pyresolve_binding / implicit_auth_headers
public/secret 分类 + 打码drivers/adapters/mcp_binding.pyclassify_mcp_binding / mask_mcp_secret_value
capability→ToolBase 适配drivers/adapters/agentscope_tool.pyDriverCapabilityTool / build_driver_agent_tools
Console MCP 管控app/mcp/config_service.pyMCPConfigService / list_tools / update_tool_whitelist / update_policy
Console MCP schemaapp/mcp/schemas.pyMCPClientInfo / MCPAccessPolicy / MCPToolInfo