首页
/ DeerFlow IM Channel Connections 实战指南:用户自有渠道的连接绑定、消息分发与文件管线

DeerFlow IM Channel Connections 实战指南:用户自有渠道的连接绑定、消息分发与文件管线

2026-09-06 18:12:56作者:董斯意

本文是 DeerFlow 中 IM Channel Connections(用户自有 IM 渠道连接) 体系的技术实战指南。该特性允许在 Telegram、Slack、Discord、飞书/Lark、钉钉、微信、企业微信 WeCom 与 Buzz 上,把某个平台账号/工作区绑定到某一个 DeerFlow 用户,从而让经渠道发起的每一次运行都落进该用户的专属"桶"(记忆、上传文件、产物与自定义 Agent)。读完本文,你将掌握:Connect 一次性绑定码的完整时序、数据库层面实现的"单一活跃归属者"不变量、绑定后消息从入站到出站的完整分发链路、同步/流式两类渠道的差异化运行策略,以及 config.yaml 中每个开关(含 Buzz 的订阅模型与信任模型)的真实含义与运维后果。

原始设计文档见 IM_CHANNEL_CONNECTIONS.md,高层索引(配置旋钮、消息流、组件清单)见 AGENTS.md 中 "IM Channels System" 一节,以及渠道子系统自己的 AGENTS.md


一、特性概览:绑定层在 bot 凭证之上补上了什么

DeerFlow 的 IM 渠道连接在语义上是一个**"按 DeerFlow 用户"的绑定层**,它叠在既有的 channels.* 出站/入站 bot 凭证之上,而不是重新发明一套出站传输。因此本地部署与私有化部署可以复用 DeerFlow 已经支持的全部出站通道,且整套实现不需要公网 IP、OAuth 回调 URL,也不需要向渠道平台注册 webhook——所有传输都建立在各平台已有的长连接/轮询/Socket Mode 机制之上。

为什么有了 bot 凭证还不够?绑定层在既有凭证能力之外,新增了三样东西:

  1. 归属者身份(Owner identity) —— 每一个 (provider, external account, workspace) 三元组都唯一映射到一个 DeerFlow 账号(owner_user_id)。由此连接发起的每个 run,都运行在该归属者的桶里(记忆、上传、产物、自定义 Agent 均按此隔离)。
  2. 一次性绑定码(One-time bind codes) —— 浏览器 Connect 流程铸造一个短期随机码(secrets.token_urlsafe(16),600 秒 TTL、单次使用),且只在发起用户的浏览器中展示。各平台 worker 在应用 allowed_users 过滤之前消费 /connect <code>(Telegram 使用 deep link 方式的 /start <code>),因此尚未进入白名单的用户也能完成首次绑定。
  3. 严格的归属转移(Strict ownership transfer) —— 最近一次成功的绑定"后到者胜";upsert_connection 会撤销其他归属者针对同一外部身份的有效行。数据库层的部分唯一索引 uq_channel_connection_active_identityWHERE status != 'revoked')让该不变量在并发写入下也是无竞态(race-free)的。

一句话概括安全边界:Connect 码是"绑定时刻"的防线,不是"聊天时刻"的防线。绑定完成后,普通消息仍由既有的 allowed_users 按原样把关。

二、Connect 绑定码流程:浏览器发起、worker 消费、manager 不碰码

关键设计约束是:浏览器发起;provider worker 消费码;调度中心(ChannelManager)永远看不到码本身。

sequenceDiagram
    autonumber
    participant Browser as Browser (Settings)
    participant Gateway as Gateway<br/>/api/channels/...
    participant Store as SQL store<br/>channel_oauth_states
    participant Worker as Provider worker<br/>(Telegram/Slack/...)
    participant Repo as ChannelConnection repo<br/>(upsert_connection)

    Browser->>Gateway: POST /api/channels/{provider}/connect
    Gateway->>Store: insert code (token_urlsafe(16), TTL=600s, single-use)
    Gateway-->>Browser: code + (Telegram: deep-link URL)

    Note over Browser,Worker: User sends /connect <code> (or /start <code>) to the provider bot

    Worker->>Store: consume_oauth_state(code)
    alt valid + unexpired
        Store-->>Worker: ok (state consumed once)
        Worker->>Repo: upsert_connection(provider, external_account_id, workspace_id, owner_user_id)
        Repo-->>Worker: connection row (active)
        Worker-->>Browser: success reply (via channel callback)
    else invalid / expired / used
        Store-->>Worker: reject
        Worker-->>Browser: rejected (no reply in chat)
    end

    Note over Repo: Partial unique index uq_channel_connection_active_identity<br/>revokes prior owner's active row for the same identity

在存储实现上(channel_connections/sql.py),consume_oauth_state 使用一次条件 UPDATE:只有能把 consumed_at 从 NULL 翻转为当前时刻的那个 writer 才算赢,两个 worker 并发消费同一码时只有一个能成功——这就是"单次使用"的无竞态落地点。码在数据库中不以明文落盘,而是先经 SHA-256 哈希(hash_state)再以 state_hash 作为主键存储(见 model.pyChannelOAuthStateRow)。

三、单一活跃归属者:数据库索引才是"真话"

整个归属转移体系的核心论点是:应用代码永远不需要显式地"撤销上一任归属者"。原因是,当一次重新使用同一身份的 upsert 触发唯一约束冲突而失败后,失败方会"输掉"这场竞争,并在下一次重试时对着现在已可见的 revoked 状态行动。

graph LR
    classDef prior fill:#E5D2C4,stroke:#806A5B,color:#30251E
    classDef new fill:#C9D7D2,stroke:#5D706A,color:#21302C
    classDef db fill:#D7D3E8,stroke:#6B6680,color:#29263A

    Prior["Prior owner<br/>connection_id=A<br/>status=connected"]:::prior
    New["New owner<br/>connection_id=B"]:::new
    Upsert["upsert_connection()<br/>(owner_user_id=B)"]:::new
    Idx["Partial unique index<br/>uq_channel_connection_active_identity<br/>WHERE status != 'revoked'"]:::db

    Prior -->|"loser: revoke"| Idx
    New -->|"winner: insert"| Idx
    Upsert -->|"trigger"| Idx
    Idx -->|"returns"| Prior
    Idx -.->|"retry against new state"| Upsert

尘埃落定之后的状态是:新连接 connection_id=Bowner=Bstatus=connected;旧连接 connection_id=Aowner=Astatus=revoked

源码层面可以看到两层配合(channel_connections/model.py):

  • 普通唯一约束 uq_channel_connection_owner_provider_identityowner_user_id + provider + external_account_id + workspace_id)保证同一归属者不会为同一外部身份重复建行;
  • 部分唯一索引 uq_channel_connection_active_identityprovider + external_account_id + workspace_idWHERE status != 'revoked')保证任何时刻、全库范围内一个外部身份至多只有一条未撤销行,SQLite 与 PostgreSQL 均支持这种 partial unique index。

sql.py 中的 upsert_connection 把事务拆成了"先撤销其他归属者的 active 行(连带删除其 channel_credentials),再提交自己的 connected 行",冲突时回滚并最多重试 3 次(_UPSERT_MAX_ATTEMPTS = 3),每次重试都重读当前可见状态,因此并发归属转移在真实竞争下也能收敛。

同一个不变量还保护着 find_connection_by_external_identity 查找——它是 ChannelManager._get_bound_identity_rejection 判断"这条消息是否来自已绑定身份"的底层依据。因为未撤销的行在任何时刻只会解析到唯一归属者,所以这条查找天然不会出现一身份映射多归属者的歧义。

四、绑定之后:Provider 消息如何走完全程

一旦连接绑定成功,每条入站消息都会经由 ChannelManager 走同一条路径。Slack/Discord(无流式)与飞书/Telegram(有流式)只在 run 边界处分叉:

sequenceDiagram
    autonumber
    participant Platform as Provider<br/>(Slack/Telegram/...)
    participant Worker as Provider worker
    participant Bus as MessageBus<br/>bounded admission + queue
    participant Mgr as ChannelManager<br/>fixed worker pool
    participant Client as langgraph_sdk<br/>async client
    participant Gateway as Gateway<br/>/api/* routers

    Platform->>Worker: inbound chat message<br/>(resolved to connection_id + owner_user_id)
    Worker->>Bus: reserve by Gateway handoff, then commit InboundMessage
    Bus->>Mgr: fixed worker gets msg and awaits handler inline
    Mgr->>Mgr: _channel_storage_user_id(msg)<br/>→ owner-bound user_id
    Mgr->>Mgr: _get_bound_identity_rejection()<br/>(re-check identity by provider+ext+ws)
    Mgr->>Client: _get_or_create_thread(thread_id or new)
    Client->>Gateway: threads.create(metadata={channel_source})
    Gateway-->>Client: thread_id
    Mgr->>Mgr: receive_file(msg, thread_id, user_id=...)<br/>(owner-bound bucket)
    Mgr->>Mgr: _ingest_inbound_files(thread_id, user_id=...)

    alt channel supports streaming
        Mgr->>Client: runs.stream(messages-tuple + values)
        loop each chunk
            Client-->>Mgr: delta / values snapshot
            Mgr->>Bus: publish_outbound(is_final=False)
        end
    else Slack/Discord (no streaming)
        Mgr->>Client: runs.wait()
        Client-->>Mgr: final state
    end

    Mgr->>Bus: publish_outbound(is_final=True)
    Bus->>Worker: outbound callback
    Worker->>Platform: post reply (Telegram editMessageText,<br/>Feishu patch card, etc.)

4.1 入站容量与过载行为

三个顶层 channels 设置控制着 MessageBus/manager 的生命周期(默认值、语义与启动校验分别实现在 service.pymessage_bus.pymanager.py 中):

配置项 默认值 含义
channels.inbound_queue_maxsize 1000 排队消息 + provider 侧仍在做最终身份/ack 准备的预留之和
channels.max_concurrency 5 常驻 ChannelManager worker 的精确数量
channels.shutdown_grace_period_seconds 3 主动 handler 被取消前允许的排空(graceful draining)时间上限

active handler 在这些 worker 内联执行(inline),所以突发流量不会造成"每条消息一个 task"。由此可得 manager 实际持有的最大在途 intake = 待处理容量 + 固定 worker 数,是有上界的。

一个值得注意的实现细节是:准入(admission)永不等待队列空间。因为让生产者协程"等待队列腾位"只会把无界积压搬到队列外面,并不会缓解过载。因此在容量打满时:

  • Slack、Discord、飞书/Lark、钉钉、Telegram、微信、WeCom:在 DeerFlow 发出工作确认(working acknowledgment)之前就丢弃新消息,MessageBus 以限速方式输出一条带累计拒绝计数的告警日志;
  • Buzz:保持该渠道的按渠道回放水位(per-channel replay watermark)不变并重连,让 relay 的历史回放机制把该事件重新投递进来;
  • GitHub webhook 扇出:返回 503,GitHub 侧把投递记为失败,运维人员或恢复任务可通过 GitHub 的 Recent Deliveries 界面或 REST redelivery API 手动重试(GitHub 不会自动重试失败的投递)。

4.2 优雅停机:先关准入、后断传输

关停顺序决定了会不会出现"handler 用已关闭的传输发消息"。实现上:先关闭准入(admission)并取消后续 watcher,但保持 provider 传输存活,让 worker 在 shutdown_grace_period_seconds 内把已接受消息排空;宽限期耗尽后,取消 active handler、丢弃从未开始的队列条目,并 await 每一个 manager 持有的 worker 与 watcher。从 SDK 线程提交的 provider 协程同样被保留、取消并 await,之后该渠道才会拆除 SDK 资源——因此一次成功的停机会保证没有任何 owned handler 还能使用已关闭的传输。Gateway 外层停机超时仍是进程级的总预算:如果它取消了清理动作,服务会保留其传输与单例,而不是虚报"成功停机"或掩盖未完成的归属清理。

五、同步 vs 流式渠道:一切分歧源于一个能力位

两种路径在 ChannelRunPolicy.supports_streaming 处分叉。该能力位来自 ChannelManager 里的渠道能力登记表 CHANNEL_CAPABILITIESmanager.py,也见渠道实现各自声明的 supports_streaming(),例如 telegram.pyfeishu.pywecom.pybuzz.py):

graph TB
    classDef sync fill:#E5D2C4,stroke:#806A5B,color:#30251E
    classDef stream fill:#C9D7D2,stroke:#5D706A,color:#21302C

    Msg["InboundMessage<br/>(channel, chat_id, text, files)"]:::sync
    Sync1["Slack"]:::sync
    Sync2["Discord"]:::sync
    Sync3["DingTalk"]:::sync
    Wait["runs.wait()<br/>→ extract final AI text"]:::sync
    Out1["publish_outbound(is_final=True)"]:::sync

    Stream1["Feishu"]:::stream
    Stream2["Telegram"]:::stream
    Stream3["WeCom (AI card)"]:::stream
    Stream["runs.stream(messages-tuple + values)"]:::stream
    Mid1["publish_outbound(is_final=False)<br/>throttled"]:::stream
    Mid2["Telegram: edit placeholder message<br/>Feishu: patch running card<br/>WeCom: PUT /v1.0/card/streaming"]:::stream
    Final["publish_outbound(is_final=True)"]:::stream

    Msg --> Sync1 --> Wait --> Out1
    Msg --> Sync2 --> Wait --> Out1
    Msg --> Sync3 --> Wait --> Out1

    Msg --> Stream1 --> Stream --> Mid1 --> Mid2 --> Final
    Msg --> Stream2 --> Stream --> Mid1 --> Mid2 --> Final
    Msg --> Stream3 --> Stream --> Mid1 --> Mid2 --> Final
  • 同步渠道(supports_streaming=false):Slack、Discord、DingTalk —— manager 调用 runs.wait() 阻塞等待,从最终 state 中抽取 AI 文本,然后一次性 publish_outbound(is_final=True)
  • 流式渠道(supports_streaming=true):Feishu/Lark、Telegram、WeCom —— manager 调用 runs.stream(messages-tuple + values),逐 chunk 以节流方式 publish_outbound(is_final=False),各平台用各自的"进行中"载体呈现中间结果(Telegram 编辑占位消息、飞书 patch 运行中的卡片、WeCom 走 PUT /v1.0/card/streaming),最终再发布 is_final=True

GitHub 是一个特殊分支:它的渠道策略是 fire_and_forget=True(见 run_policy.py 中的 ChannelRunPolicy 描述符,以及 github/run_policy.py 的登记),manager 只需 runs.create() 并在 run 进入 pending 后即返回,不做任何出站回复——因为 GitHub agent 是在自己的沙箱里通过 gh CLI 直接回帖的。完整 GitHub 流程见 GITHUB_AGENTS.md

5.1 ChannelRunPolicy:一份策略数据类承载全部渠道差异

从源码看,CHANNEL_RUN_POLICY 是一个"渠道名 → 策略"注册表,新增一个 webhook/特殊渠道只需注册一行策略而不是改动 manager 的多处方法(run_policy.py)。值得了解的字段包括:is_interactive(False 时关闭澄清提问,适配"无人同步在场"的自主长任务)、default_recursion_limit(把 recursion_limit 抬到 max(现有值, 配置值),GitHub 这类长跑任务需要更高上限,普通聊天回合则保持全局默认)、credentials_provider(向 run_context 注入平台专用凭证的异步钩子,异常会被捕获降级为只读运行)、requires_bound_identity(False 时跳过按发送者的绑定身份闸门——webhook 渠道用 HMAC 做真实性认证,没有 /connect 握手)、fire_and_forgetserialize_thread_runs 等。

六、归属者作用域的文件存储(Owner-scoped File Storage)

ChannelManager_handle_chat 顶部只解析一次存储归属者:_channel_storage_user_id(msg) 会做安全化处理(sanitized owner id),并在无绑定身份的渠道上回退到 safe(msg.user_id)(与 _resolve_run_params 的 run 身份保持一致);只有当任何身份都不可用时才返回 None。这个值被贯穿到整个文件管线:

  • 同时用作 run_context["user_id"](run 身份);
  • 同时用作 memory、uploads、outputs 的存储桶。

于是 agent 读写的桶,永远是渠道文件被暂存的那个桶——二者不会错位

flowchart TB
    classDef owner fill:#D8CFC4,stroke:#6E6259,color:#2F2A26
    classDef resolve fill:#C9D7D2,stroke:#5D706A,color:#21302C
    classDef bucket fill:#D7D3E8,stroke:#6B6680,color:#29263A
    classDef agent fill:#E5D2C4,stroke:#806A5B,color:#30251E

    Inbound["InboundMessage<br/>connection_id, owner_user_id, workspace_id"]:::owner
    Resolve["_channel_storage_user_id(msg)<br/>sanitized + fall back to safe(msg.user_id)"]:::resolve
    UserID["user_id = OWNER"]:::resolve

    RunID["run_context['user_id']<br/>(run identity)"]:::agent
    RunUploads["ensure_uploads_dir(thread_id, user_id=OWNER)"]:::bucket
    Ingest["_ingest_inbound_files(user_id=OWNER)"]:::bucket
    Receive["Channel.receive_file(msg, thread_id, user_id=OWNER)"]:::bucket
    Resolved["_resolve_attachments(user_id=OWNER)"]:::bucket
    Artifact["_prepare_artifact_delivery(user_id=OWNER)"]:::bucket
    Memory["_resolve_memory_user_id<br/>(make_safe_user_id match)"]:::bucket

    Bucket["backend/.deer-flow/users/OWNER/.../user-data/{uploads,outputs}"]:::bucket

    Inbound --> Resolve --> UserID
    UserID --> RunID
    UserID --> Receive
    UserID --> Ingest
    UserID --> RunUploads
    UserID --> Resolved
    UserID --> Artifact
    UserID --> Memory
    RunUploads --> Bucket
    Ingest --> Bucket
    Receive --> Bucket
    Resolved --> Bucket
    Artifact --> Bucket

缓存下来的归属者值在阻塞路径(runs.wait)与流式路径(_handle_streaming_chat)之间是复用的——即使未来某个 Channel.receive_file 返回了被改写过的 InboundMessage,uploads 与产物投递依然命中同一个桶。

七、IM 附件管线:下载、落桶、注入与发现

入站文件(图片、文档)先经各渠道自己的 Channel.receive_fileprovider 特定物化。物化后存在两条路:

  • 走共享元数据路径的附件_ingest_inbound_files 暂存,元数据写入 HumanMessage.additional_kwargs.files,然后 UploadsMiddleware 为当前消息注入一段 <current_uploads> 上下文块;
  • 部分 provider 在下载时直接消费自己的描述符,把占位符或消息文本改写为最终的虚拟路径(或一条失败提示)。

历史上传不会在后续轮次被自动注入;agent 需要借助 list_uploaded_files 工具主动发现历史文件。飞书/Lark 的入站资源流在被持久化或同步进非本地沙箱之前有一个 20,000,000 字节上限;超限资源与个别文件的路径失败会以失败占位符形式呈现在消息文本中,而不会中断同一条入站消息里后续附件的处理

sequenceDiagram
    autonumber
    participant IM as Provider message<br/>(file attachment)
    participant Worker as Provider worker
    participant Mgr as ChannelManager
    participant Ch as Channel impl<br/>.receive_file
    participant FS as Uploads directory<br/>users/OWNER/.../uploads/
    participant MW as UploadsMiddleware
    participant Agent as Agent run

    IM->>Worker: message with file URL/bytes
    Worker->>Mgr: InboundMessage(files=[...], connection_id, owner_user_id)
    Mgr->>Mgr: storage_user_id = _channel_storage_user_id(msg)
    Mgr->>Ch: receive_file(msg, thread_id, user_id=storage_user_id)
    Note over Ch: provider-specific download/decrypt/read;<br/>may persist or hand bytes to manager
    Ch-->>Mgr: materialized message<br/>(provider may rewrite placeholders/text)
    alt attachment continues through shared metadata path
        Mgr->>FS: _ingest_inbound_files(<br/>thread_id, msg, user_id=storage_user_id)
        FS-->>Mgr: uploaded file metadata
        Mgr->>MW: HumanMessage with<br/>additional_kwargs.files
        MW->>Agent: prepend <current_uploads><br/>(paths under /mnt/user-data/uploads/)
    else provider supplies virtual path in message text
        Ch->>FS: persist and/or sync attachment
        Mgr->>Agent: HumanMessage with rewritten path<br/>or failure notice
    end
    Agent->>FS: read_file / view_image (sandbox)

八、配置详解:两步启用 + 一个必须知晓的升级语义

8.1 第一步:在既有的 channels 块中配置各平台 bot

实际的 IM bot 凭证仍然位于既有的 channels 块(环境变量 $VAR 形式由 AppConfig.resolve_env_variables 在加载时递归展开,见 app_config.py;完整示例可对照 config.example.yaml):

channels:
  inbound_queue_maxsize: 1000
  max_concurrency: 5
  shutdown_grace_period_seconds: 3

  telegram:
    enabled: true
    bot_token: $TELEGRAM_BOT_TOKEN

  slack:
    enabled: true
    bot_token: $SLACK_BOT_TOKEN
    app_token: $SLACK_APP_TOKEN

  discord:
    enabled: true
    bot_token: $DISCORD_BOT_TOKEN

  feishu:
    enabled: true
    app_id: $FEISHU_APP_ID
    app_secret: $FEISHU_APP_SECRET

  dingtalk:
    enabled: true
    client_id: $DINGTALK_CLIENT_ID
    client_secret: $DINGTALK_CLIENT_SECRET

  wechat:
    enabled: true
    bot_token: $WECHAT_BOT_TOKEN

  wecom:
    enabled: true
    bot_id: $WECOM_BOT_ID
    bot_secret: $WECOM_BOT_SECRET

  buzz:
    enabled: true
    relay_url: wss://buzz.example.com
    private_key: $BUZZ_PRIVATE_KEY   # hex or nsec1…

8.2 第二步:在 channel_connections 中开启用户绑定

channel_connections:
  enabled: true
  # Auth-enabled deployments require ordinary IM messages to come from a
  # connected DeerFlow user by default. Set this to false only for legacy
  # operator-owned/open-bot deployments that intentionally route unbound
  # platform users to platform-ID user buckets.
  require_bound_identity: true

  telegram:
    enabled: true
    bot_username: $TELEGRAM_BOT_USERNAME

  slack:
    enabled: true

  discord:
    enabled: true

  feishu:
    enabled: true

  dingtalk:
    enabled: true

  wechat:
    enabled: true

  wecom:
    enabled: true

  buzz:
    enabled: true

几个关键语义:

  • channel_connections 不复制 provider 密钥。它只控制浏览器侧的 Connect UI,并存储按用户的绑定记录。Telegram 需要 bot_username 的唯一理由,是让前端能拼出 deep link。
  • 配置模型见 channel_connections_config.py:顶层的 enabledrequire_bound_identity 之下,每个渠道是独立的子配置;其中 telegram.configured 取决于 bot_username 是否非空,其余渠道只要 enabled 即为 configured。
  • 强制绑定身份语义:当 channel_connections.enabledrequire_bound_identity 同时为 true,启用认证的部署会在创建 DeerFlow thread/run 之前拒绝普通的未绑定 IM 消息,用户必须先在 DeerFlow Settings 里完成渠道连接。而关闭认证的本地模式仍会把渠道消息路由给默认用户;需要恢复旧式 open-bot 行为时,显式设置 require_bound_identity: false 即可。

8.3 升级注意事项(重要)

require_bound_identity 的默认值是 true。这意味着:升级前已开启 channel_connections.enabled: true 的认证部署,在本字段引入后,会开始拒绝普通的未绑定 IM 消息。对于刻意允许未绑定平台用户创建 DeerFlow run 的旧式 operator-owned/open-bot 部署,请在升级之前显式设置 require_bound_identity: false 并重启服务。

8.4 各平台的 Connect 操作流程

  • Telegram:前端生成一次性短码 → Connect 按钮打开 https://t.me/<bot_username>?start=<code> → 既有 Telegram 长轮询 worker 收到 /start <code>,把该 Telegram chat/user 绑定到当前 DeerFlow 用户。
  • Slack:前端生成一次性短码 → UI 提示 Send /connect <code> to the DeerFlow Slack bot. → 既有 Slack Socket Mode worker 收到消息后绑定该 Slack user/team。
  • Discord:前端生成一次性短码 → UI 提示 Send /connect <code> to the DeerFlow Discord bot. → 既有 Discord Gateway worker 收到消息后绑定该 Discord user/guild。
  • 飞书/Lark、钉钉、微信、WeCom:前端生成一次性短码 → UI 提示 Send /connect <code> to the DeerFlow <Provider> bot. → 已运行的 long-connection 或轮询 worker 收到消息后绑定平台 user/workspace 身份。
  • Buzz:见下一节单独展开。

对带 allowed_users 白名单的 provider(Telegram、Slack、钉钉、微信等),有效的 /connect <code>(或 Telegram 的 /start <code>会在白名单检查之前被消费。这是有意为之:尚未进入白名单、bot 因此从未见过其平台身份的用户,仍然可以完成第一次由浏览器发起的绑定。绑定完成后,allowed_users 继续如常把关普通(非绑定)消息。

九、Buzz:无开发者控制台的特殊成员身份

Buzz 与上述 bot/app 凭证渠道不同——它没有独立的开发者控制台:DeerFlow 以普通成员身份加入 relay。需要为这个身份生成一对 Nostr 密钥,以下所有描述都指其 hex 公钥

9.1 两步式 onboarding(缺一不可)

Relay 成员身份与渠道成员身份是两回事,只做第一步会产生一个"能连上、能通过认证、但收不到任何东西"的 connector:

  1. 把 pubkey 注册为 relay 成员buzz-admin add-member --pubkey <hex>。这一步让该身份能够认证(NIP-42)并能发布消息。
  2. 把它加入应参与的每个渠道buzz channels add-member --channel <uuid> --pubkey <hex> --role bot。聊天事件只会投递给渠道成员,且 relay 会拒绝任何 p-mentions 非成员的消息(报错 mentioned pubkeys are not channel members)——缺这一步,connector 既听不到 mention 也答不了 mention。

配置上只认两个键:channels.buzz.relay_urlchannels.buzz.private_key(hex 或 nsec1…),再开启 channel_connections.buzz。Connect 时前端生成一次性短码,UI 提示 Send /connect <code> to the DeerFlow Buzz bot.,Buzz relay-loop worker 收到以 DM 或双方共同所在渠道 @mention 发送的消息后,把发送者的 Nostr pubkey 绑定到当前 DeerFlow 用户。

两个重要事实:

  • 渠道是自动发现的,不需要写进 config.yaml 每次连接时 connector 会向 relay 查询该身份属于哪些渠道并逐一下订;之后新加入的渠道实时生效(relay 下发 membership 通知即开始监听,移除即停止),无需重启或重连。日志里出现 channel discovery returned no channels 说明第 2 步没做。
  • 依赖 buzz extra(uv sync --extra buzz,为 coincurve 库)。detect_uv_extras.py(以及通过 backend/Dockerfile 的 Docker/生产构建)会在 config.yamlchannels.buzz.enabled: true 时自动探测并保留该 extra——这和 browser extra 为 browser_navigate 被自动探测的方式一致。

9.2 Buzz 订阅模型

Buzz relay 只向渠道作用域的订阅投递聊天事件,因此 connector 的订阅形态是刚性的:全局 REQ {"kinds":[9]} 会被接受并应答 EOSE,但不会扇出任何聊天事件;单条订阅也无法覆盖多个渠道(多值 #h 匹配不到任何东西)。所以在每次连接、NIP-42 认证完成后:

订阅 过滤器 用途
buzz-discovery {"kinds":[39000]} 历史查询:精确列出该身份所属的全部渠道(每个渠道一条已存事件,随后 EOSE)。提供各渠道的名称与类型——DM 的 mention 豁免也读取这里。不要#p 收窄它——那会什么都匹配不到。
buzz-membership {"kinds":[44100,44101], "#p":["<our pubkey>"], "since": …} 实时成员资格通知。44100(added)立刻订阅新渠道;44101(removed)关闭对应渠道订阅。这就是"新增渠道无需重启"的原因。
buzz-chat-<uuid> {"kinds":[9], "#h":["<uuid>"], "since": …} 每个已发现渠道一条——唯一真正能收到消息的形态

值得运维人员记住的五个推论:

  • 回放按渠道分别跟踪。 每个渠道有自己的 since 水位,只有 DeerFlow 实际处理过的事件才会推进它。若共用一条水位,繁忙渠道会把游标拖过安静渠道的未读消息,导致重连后跳过它们;而按渠道游标的最坏代价只是重复投递(manager 的入站去重会吸收),永远不会漏消息。
  • 成员资格事件只存在于实时流。 relay 会存储 44100/44101,若不设 since,每次连接都会把整段成员资格历史当作"刚刚发生"重放:对每条已存储的 add 重跑一次渠道发现(一次连接会看到多条 channel discovery complete,渠道还会被记为 <unnamed>),重订已被移除的渠道,并短暂退订仍在的渠道。因此订阅锚定在 socket 打开的瞬间并往前留 60 秒余量,让连接/认证握手期间发生的成员变化或 relay 时钟偏差仍能被捕捉。
  • 渠道订阅数有上限(256)。 渠道列表来自网络,与其他远端输入一样有界。触顶时新渠道会被拒绝,并在 per-channel subscription limit reached 告警中点名,而不会驱逐任何在工作的订阅。
  • relay 关闭的订阅会被重新打开,每次连接最多重试 3 次。 relay 丢弃订阅时 socket 上是静默失败的:聊天订阅哑掉一个渠道、buzz-membership 让 DeerFlow 再也学不到入退群、buzz-discovery 杀死完整性清扫。所以 CLOSED 帧不只是被记录而是被恢复,静默的订阅总会以 WARNING 级被点名。当 relay 明示原因说明"该订阅已不属于你"时(NIP-01/NIP-42 的 auth-required:/restricted:/blocked:/invalid: 前缀,或 buzz-relay 自己的吊销措辞),恢复会被跳过——因为重发同一条 REQ 只会和 relay 对抗;其余任何原因(包括不带任何原因的 CLOSED)都被当作小故障重试,以 3 次尝试为兜底,之后保持静默直到下次重连从零重建。
  • 已知边界:单渠道在断连期间积压超过 2000 条未读,会丢掉最旧的一批。 relay 把单订阅历史投递上限设为 2000 条且最新优先(即使带 since 也如此)。DeerFlow 只处理收到的部分,渠道水位随之越过其余消息——那些更旧的消息永远不被投递也不被重试。设计中其余所有缺口都偏向"重复投递"(由入站去重吸收),这是唯一会"跳过"的遗留情形,且需要"断连 + 单渠道 >2000 条积压"两个条件同时成立。

9.3 Buzz 信任模型

在团队自营的 relay 上,relay 运营者未必等于 DeerFlow 运营者,所以必须分清什么被密码学验证、什么只是被信任:

已验证(每次入站事件都在密码学上成立):DeerFlow 从收到的载荷重算每个事件的 NIP-01 id,并用 BIP-340 Schnorr 签名对照声明中的 pubkey 做验证,验证通过前事件不能影响任何东西。因此 relay 无法改写成员消息、无法把一位作者的签名搬到另一载荷上、也无法冒称一个它不掌握密钥的已白名单作者。这一点同样适用于 /connect 绑定——relay 无法把别人的 pubkey 绑到攻击者的 DeerFlow 账号。验证失败的事件会被丢弃并告警。

被信任(未验证):kind-39000 渠道元数据的作者身份。Buzz 用 relay 自己的密钥对发布渠道发现事件,但现有配置没有标识这把密钥(relay_url 是网络地址而非签名密钥),DeerFlow 只能证明该事件被某个成员签名。由于渠道发现与订阅恰好由这些事件驱动,一条伪造的 kind-39000 会有两个(而非一个)后果:

  1. 它可以把渠道标记为 type: "dm",从而放宽该渠道的 require_mention 要求;
  2. 它可以让 DeerFlow 为一个伪造者指定的渠道打开聊天订阅——因为 DeerFlow 监听的渠道集合就是它持有元数据的渠道集合。

两者都无法让任何东西被执行allowed_users 白名单与逐事件签名验证是两条独立的闸门:无论渠道类型如何、订阅如何被打开,非白名单作者一律被丢弃。(2) 的爆炸半径是"relay 把自己流量读回给一个忽略它的订阅者",且受 256 订阅上限约束(只拒绝新订阅、不驱逐在工作的订阅,诱导出来的订阅顶不掉真实渠道)。伪造的 kind-44100 成员通知同理,只是其 p 标签会在本地被重新核对,所以它至少必须点名本身份。如果你需要在一个不能全信成员的 relay 上让 mention 要求不可伪造,就把这些渠道排除在 mention_free_channels 之外,并把 DM 检测当作便利功能而非安全边界。

默认拒绝的白名单:与其他 provider(空 allowed_users = 允许所有人)不同,channels.buzz.allowed_users 刻意是默认拒绝——空列表意味着没人能触发 run,DeerFlow 会就此打一条启动告警。应把每个允许触达 agent 的成员 pubkey(hex 或 npub1…)加进去。个别丢弃只在 DEBUG 级记录。

绑定身份:pubkey 一旦完成 /connect,其入站消息就解析到该连接并在绑定 DeerFlow 用户身份下运行(记忆、文件、产物都落在该用户桶里)。绑定按 relay 主机隔离——同一个 pubkey 在不同 relay 上是不同身份,必须分别绑定。

绑定码使用 128 位随机性,10 分钟后过期,且单次使用。

十、运行时模型:四张 SQL 表与身份字段

连接记录落在 deerflow.persistence.channel_connections 下的 SQL 表中(模型定义见 model.py,SQL 仓库见 sql.py):

职责
channel_connections 归属用户、provider 身份、workspace/guild/team、状态、元数据
channel_oauth_states 一次性 Connect 码与 Telegram deep-link 状态
channel_conversations 连接作用域的 IM 会话 → DeerFlow thread 映射
channel_credentials 为未来的 provider-token 流程预留;本地/私有绑定流程不使用

解析到某条连接的入站消息携带 connection_idowner_user_idworkspace_id 三个字段。ChannelManagerowner_user_id 用作 DeerFlow run 用户 id,同时保留原始平台用户 id 为 channel_user_id

运行时 provider 凭证是部署级 bot 秘密,不是用户所有的连接凭证。它们既可以来自 config.yamlchannels.*,也可以来自浏览器运行时配置流程——后者通过 ChannelRuntimeConfigStore 持久化,让本地/私有部署无需编辑 YAML 即可配置 bot。运行时存储的本地回退是明文 JSON 文件,owner-only 文件权限(0600;仅当 DeerFlow 数据目录已被视为秘密存储时才使用它。微信扫码登录的认证状态遵循同一本地运行时模型,并可能把二维码导出的 bot token 持久化在渠道状态目录里。

十一、安全笔记汇总

  • 浏览器 API 保持已认证与 CSRF 防护;
  • Connect 码 128 位随机、短时效、单次使用;
  • 运行时 provider bot token 是共享部署秘密。运行时设置响应会脱敏密码字段,变更运行时/渠道 worker 的 API 需要 admin 用户;
  • 存储的按连接凭证使用 channel_credentials 加密路径。若存量的凭证材料无法解密,DeerFlow 将其视为不可用,而不是使用损坏的秘密(ChannelCredentialCipher 基于 Fernet,见 sql.py);
  • 本地明文运行时凭证回退(0600)已在上一节说明;在部署专用 secret backend 之前,非本地部署应优先使用由部署管理的环境/配置秘密;
  • allowed_users 不是绑定时刻防线。由于 Connect 码在白名单之前被处理(见 8.4 节),任何持有有效码的人都能消费它——不只是白名单用户。绑定安全因此完全系于码的机密性:128 位随机、10 分钟过期、单次使用、只显示在发起用户浏览器中(从不回显到聊天)。请把 Connect 码当作一次性密码,不要转发。
  • 一个外部身份——(provider, external account, workspace/team/guild)——至多一个活跃归属者。最近一次成功绑定胜出:绑定一个已被另一 DeerFlow 用户持有的身份会转移归属权并撤销前任的绑定(连同其存储的凭证)。这在数据库层强制,两个用户竞速绑定同一身份不可能双双处于已连接状态;
  • provider bot token 仍留在 channels.*永不返回给浏览器
  • 本实现不新增任何公共 provider 回调或 webhook 路由

十二、测试与源码索引:如何继续深入

上述机制的许多不变量都可以在测试中找到可验证的落点,例如:

  • 归属转移/凭证加解密/一次性码消费等仓库行为:backend/tests/test_channel_connections_repository.py
  • 浏览器 Connect API 的路由与权限:backend/tests/test_channel_connections_router.pybackend/tests/test_channel_connections_config.py
  • 各 provider 的绑定与消息处理:backend/tests/test_telegram_channel_connections.pybackend/tests/test_slack_channel_connections.pybackend/tests/test_discord_channel_connections.pybackend/tests/test_additional_channel_connections.pybackend/tests/test_buzz_channel.py
  • 过载、去重与去重后消息处理:backend/tests/test_channels.py

核心代码入口:调度与绑定闸门在 manager.py_channel_storage_user_id_handle_chat_handle_streaming_chat_get_bound_identity_rejectionCHANNEL_CAPABILITIES);有界队列在 message_bus.py;渠道策略描述符在 run_policy.py;持久化在 channel_connections/sql.py。总体而言,"浏览器铸造一次性码 → provider worker 消费并绑定 → 数据库保证单一活跃归属者 → ChannelManager 按归属者身份隔离地跑 run、存文件、回消息"这条主线清晰、独立、无回调依赖,正是它让 IM 渠道在本地与私有部署中也能拥有完整的多用户归属语义。

登录后查看全文
热门项目推荐
相关项目推荐