首页
/ DeerFlow IM Channels 消息网关架构与实战指南:从聊天接入到事件驱动 Agent

DeerFlow IM Channels 消息网关架构与实战指南:从聊天接入到事件驱动 Agent

2026-09-06 18:31:53作者:裴麒琰

DeerFlow 是一个面向长时程(minutes-to-hours)任务的开放源码 SuperAgent 框架,而 backend/app/channels/AGENTS.md 正是其 IM 通道子系统的权威开发文档,完整描述了飞书(Feishu/Lark)、Slack、Telegram、Discord、钉钉(DingTalk)、GitHub、Buzz(Nostr relay)等外部消息平台如何桥接到 DeerFlow Agent。本文以该文档为主体,结合仓库源码与示例配置深入展开,读完你将掌握:IM 通道的整体架构与消息流转链路、平台通道的差异实现与「fire-and-forget / 流式增量回复」等运行策略、channel_user_id 与 owner-scoped 文件存储等身份与安全设计,以及 GitHub 事件驱动 Agent 与 Buzz 通道的重连与信任模型。

架构总览:通道如何接入 DeerFlow Agent

app/channels/ 目录的作用一句话概括:把外部消息平台(Feishu、Slack、Telegram、Discord、DingTalk、GitHub 等)桥接到 DeerFlow Agent —— 桥接的对象是 Gateway 的 LangGraph 兼容 HTTP API。

关键架构约束是「通道与前端使用同一种客户端」:Channels 通过 langgraph-sdk HTTP 客户端与 Gateway 通信(与前端完全相同),从而保证线程(thread)在服务端创建并被服务端统一管理。通道工作进程要调用 Gateway 上会改变状态的线程/运行请求(thread/run),但通道进程内并没有浏览器会话 cookie,因此内部 SDK 客户端会注入两层凭据:

  • 进程内的内部认证(process-local internal auth);
  • 一对匹配的 CSRF cookie/header。

这样 Gateway 才会接受来自通道 worker 的写请求。

从源码结构看,backend/app/channels/ 下每个文件的职责划分清晰:

文件 职责
message_bus.py 异步 pub/sub 中枢:InboundMessage 入队 → dispatcher 消费;OutboundMessage → 回调 → 各通道发送
store.py JSON 文件持久化,维护 IM 会话 → DeerFlow thread_id 映射
manager.py 核心分发器:建线程、路由命令、按通道策略执行 runs.wait / runs.stream / runs.create
base.py 抽象 Channel 基类,提供 start/stop/send 生命周期与线程安全协程提交设施
service.py 依据 config.yaml 管理全部已配置通道的生命周期
slack.py / feishu.py / telegram.py / discord.py / dingtalk.py / wechat.py / wecom.py / github.py / buzz.py 各平台专属实现
run_policy.py 每通道运行策略注册表(ChannelRunPolicy

核心组件逐个拆解

message_bus.py —— 通道与分发器的解耦枢纽

message_bus.py 实现了一个异步 pub/sub 总线,两条路径方向相反:

  • 入站(inbound)InboundMessage → 队列 → dispatcher(消费者即 ChannelManager._dispatch_loop());
  • 出站(outbound)OutboundMessage → 回调 → 平台通道回发。

InboundMessage 数据类(见源码 docstring)承载了通道与 DeerFlow 之间的全部上下文字段:

  • channel_namechat_iduser_id —— 平台侧标识;
  • topic_id —— 会话内主题,同 chat_id 下相同 topic_id 的消息复用同一个 DeerFlow 线程;为 None 时每次消息新建线程(一次性问答);
  • connection_id / owner_user_id / workspace_id —— 用户自有连接(user-owned connection)场景下携带;owner_user_id 成为 DeerFlow run 的 user_id,平台原生用户 id 保留在 user_id
  • files —— 附件(平台专有字典);
  • 模块还定义了 INBOUND_FILE_CONTENT_KEY = "_content"(源码第 23 行):通道下载到的原始字节只能通过该临时字段穿越 adapter/manager 边界,manager 消费该字段后就会在持久化安全上传元数据之前移除它,避免令牌或大字节数据残留。

store.py —— 会话到线程的持久映射

store.py 用单个 JSON 文件把 IM 会话映射到 DeerFlow thread_id。映射键规则为:

  • 根会话:channel_name:chat_id
  • 子主题会话:channel_name:chat_id:topic_id

磁盘数据结构(源码第 19–34 行的 docstring 明示)形如:

{
  "<channel_name>:<chat_id>": {
    "thread_id": "<uuid>",
    "user_id": "<platform_user>",
    "created_at": 1700000000.0,
    "updated_at": 1700000000.0
  }
}

该文件有两个值得注意的实现细节:写入采用临时文件 + Path.replace() 的原子重写(崩溃不会留下半截 JSON);每次对 _data 的访问都必须由 _lock 保护。文档特别强调:list_entries() 在锁内对 keys 与条目做快照拷贝,释放锁后再格式化结果——否则并发通道线程会在遍历期间改变字典尺寸,导致 RuntimeError: dictionary changed size during iteration

base.py —— Channel 抽象基类与线程安全提交

base.py 定义了抽象 Channel 的生命周期契约(start / stop / send / send_file / receive_file)。对于由 SDK 回调线程提交协程的场景,文档要求使用 _submit_threadsafe_coroutine(),其背后有一个常被忽视的 asyncio 陷阱:

它会在 owner loop 上创建并保留真实的 asyncio.Task,而不是把 run_coroutine_threadsafe() 返回的代理 Future 当作完成信号。

同时,提交通道与关闭原子化绑定,stop() 在拆除 SDK 资源之前必须调用 _close_and_drain_threadsafe_futures()。通道子类由此获得线程安全的入站保留(reservation)与提交原语。

service.py —— 通道生命周期编排

service.py 负责从 config.yaml 加载并管理所有已配置通道的生命周期。关闭顺序上刻意做了设计:shutdown 先关闭 manager 的准入(admission),但保持传输层存活,直到每个 manager worker 与 follow-up watcher 都退出。这样一次成功的 manager stop 后不再持有任何活跃 handler;如果 Gateway 的外层超时把 shutdown 取消了,service 会保留通道对象与全局单例,使未完成资源不被脱离、清理可被重试。

run_policy.py —— 每通道运行策略

run_policy.py 定义了全局 CHANNEL_RUN_POLICY 注册表(channel_name → ChannelRunPolicy)。该 dataclass 一次性声明了影响通道 run 行为的全部开关,它的设计动机在源码注释里说得很清楚:让「新增一个 webhook 通道」变成一行注册,而不是去改 manager 的多个方法。各字段含义:

字段 默认值 语义
is_interactive True False 时 manager 设 run_context["disable_clarification"]=TrueClarificationMiddleware 直接返回 "proceed with best judgment" 而不是中断(无人同步在场的自主长任务)
default_recursion_limit None 为数值时把 run_config["recursion_limit"] 提到 max(现有值, limit)None 保留全局默认 100
credentials_provider None 异步钩子,_resolve_run_params 之后向 run_context 注入平台凭据;异常被捕获并记录,凭据失败降级为只读运行
requires_bound_identity True False 时跳过按发送者绑定的身份门禁(webhook 通道由 HMAC 保证真实性)
fire_and_forget False True 时用 runs.create(返回 pending 即结束)替代 runs.wait,规避 SDK 300 秒 httpx.ReadTimeout
serialize_thread_runs False True 时同线程入站消息排队串行(用于飞书话题这类快速追问场景),而不是触发运行时 busy 错误
buffer_followups_on_busy False Truefire_and_forget 路径遇到 ConflictError 会把触发消息缓冲起来,run 结束后自动合并重放

平台专属实现与「通道策略」

不同平台的能力差异很大,DeerFlow 的选择是把「运行策略」做成可配置的组合,而不是给每个通道写死一套:

  • Feishu/Telegram(聊天):走 client.runs.stream(["messages-tuple", "values"]),把 AI 文本增量累计后发布多条 is_final=False 的出站更新,最后发布 is_final=True 的终态回复。
  • Slack/Discord(聊天):走 client.runs.wait(),run 结束后提取最终响应一次发布。
  • Feishu 话题的快速追问:当通道的 ChannelRunPolicy.serialize_thread_runs=True 时,同线程的飞书轮次在 manager 内串行化,紧跟的追问会排队而不是撞上 busy 回复。
  • GitHub(事件驱动)fire_and_forget=Trueruns.create() 在 run 进入 pending 即返回;manager 不等待终态也不主动发布出站消息(GitHub 通道出站默认仅打日志),因为 Agent 会在沙箱里用 gh CLI 自主发布评论或建 PR。

各通道平台实现还带各自的能力细节:

  • feishu.py:内存中跟踪正在运行的消息卡片 message_id,原地对同一卡片反复 patch(卡片 JSON 设 config.update_multi=true 以满足飞书 patch API);已有飞书话题内的消息会在卡片上附带一条紧凑的源消息预览。
  • telegram.py:接受文本/图片/文档,保留媒体 caption;把无令牌的附件字节交给共享上传管线;用 editMessageText 原地编辑 "Working on it..." 流式目标;可选的 channels.telegram.rich_messages 让最终 Markdown 回复以 Bot API 10.1 Rich Message 形式发送。
  • discord.py:入站处理让出前注册 typing 循环;_start_typing()_running 为假时拒绝工作。因 stop() 跑在主 loop 而 typing 任务属于 _discord_loop,跨线程的取消、await 与 map 清理都在目标 loop 上调度并做有界等待;_run_client()finally 中先清空任务再让异常/断线破坏 loop。
  • dingtalk.py:配置 card_template_id 时使用 AI Card 流式模式(创建卡片 → 经 PUT /v1.0/card/streaming 流式更新 → is_final=True 时终结,失败回退 sampleMarkdown);并重写 receive_file,按 downloadCode 把入站图片(picture/richText)与文档(file)下载进线程上传桶,镜像 feishu 的做法。

运行策略在源码中的落地

策略由 manager 在 _resolve_run_params 之后经 _apply_channel_policy 应用。以 GitHub 为例,backend/app/gateway/github/run_policy.py 从 Gateway 引导时调用 register_policy() 注册自己的条目,并在 inject_github_credentials 中把短时(1h)GitHub App installation token 以字符串形式写入 run_context["github_token"] —— 源码注释解释得很透彻:run_context 会被 langgraph_sdk HTTP 客户端 JSON 编码后发给 Gateway,Python callable 过不了 TypeError: Type is not JSON serializable,只有字符串能穿过 SDK 传输层往返。

消息流转全链路

文档给出了完整消息流编号,整理如下:

  1. 外部平台 → Channel 实现 → MessageBus.publish_inbound()
    • GitHub 例外:webhook 路由校验投递后调用 fanout_event(bus, ...),每个匹配的 agent 绑定各发布一个 InboundMessage,而不是一个长轮询通道 worker;
    • Telegram 的图片/文档更新取最大图片尺寸或文档元数据,保留 message.caption,下载前后都强制 Bot API 20000000 字节上限,绝不暴露带令牌的 Bot API 文件 URL;
    • 飞书/Lark 入站图片/文件下载最多读取 20000001 字节,超过 20000000 字节在持久化或沙箱同步前即拒绝;超限、不安全路径、路径解析失败只把该附件占位符改写为 Failed to obtain the [type],同消息的后续附件仍可加载。
  2. ChannelManager._dispatch_loop() 从队列消费;
  3. 用户自有连接(user-owned)的入站消息携带 connection_idowner_user_idworkspace_id(细节见下文安全章节);
  4. 聊天消息:经 Gateway 的 LangGraph 兼容 API 查找/创建线程;
  5. 飞书/Telegram:runs.stream() → 累积 AI 文本 → 多条 is_final=False 更新 → 终态 is_final=True
  6. Slack/Discord:runs.wait() → 提取最终响应 → 发布出站; 6b. GitHub(fire_and_forget=True):runs.create() 一旦 pending 即返回;忙线程上的 ConflictError 仍走标准 THREAD_BUSY_MESSAGE 路径,当策略还设置 buffer_followups_on_busy=True(GitHub 默认)时触发消息额外被捕获进 per-thread 缓冲,而不是只打日志后静默丢弃;
  7. 飞书先发送一张 running 回复卡片,随后每次出站更新原地 patch 同一张卡;话题内的已发送消息携带紧凑源消息预览,排队的同线程追问会从 queued → running → final 原地 patch,不回退到通用 busy 回复;
  8. Telegram 流式:「Working on it...」占位消息注册为流目标;非终态更新 editMessageText 原地编辑(通道侧限流:私聊 1s、群组 3s,应对 Telegram 群组 20 条/分钟上限;4096 字符截断;被限流的更新直接丢弃);终态更新做最后一次编辑并把超过 4096 字符的文本拆分后追加发送;
  9. 钉钉 AI Card 模式:runs.stream() → 建卡 → 流式更新 → 终结;卡片创建/流式失败回退 sampleMarkdown
  10. 命令类(/new/status/models/memory/goal/help):本地处理或查询 Gateway API。/goal 设置目标会先经 Gateway 持久化,然后把目标作为一条聊天轮次路由给 Agent;
  11. 出站 → 通道回调 → 平台回复;GitHub 例外(打日志,不自动发帖,见下文)。

流式输出的白名单机制:为什么只允许 assistant 消息

文档浓墨重彩地强调一个安全问题,对应 manager.py_accumulate_stream_text / _stream_payload_type / _is_assistant_stream_type

流中可以发布的内容是一份白名单(allowlist),而不是黑名单(denylist)。

原因要追溯到 DeerFlow 对运行时上下文的隐藏注入:

  • DynamicContextMiddleware 把召回的记忆注入为一条隐藏的 HumanMessagetype == "human"),并把用户本轮发言改写为新的 HumanMessage
  • DurableContextMiddleware 注入隐藏的 <durable_context_data> HumanMessage
  • LangGraph 会把这些状态写入在 messages-tuple 流上扇出。

若按旧的「type"tool" 才拒绝、其余全发布」的规则,等于把 DeerFlow 的隐藏模型上下文泄露到所有流式 IM 通道。文档记录了实锤:在一次 Buzz relay 上曾发布出 <memory> 事实块;另一次运行把用户消息逐字回显成了 assistant 回复。

白名单的具体形态:

  • 只放行 assistant 类型消息——LangChain 把 AIMessage.type 序列化为 "ai"AIMessageChunk.type 序列化为 "AIMessageChunk",外加 OpenAI 风格的外来运行时拼写 "assistant"
  • 匹配用前缀(ai / assistant)而绝不用子串,因为普通英文单词(chaindomain)也含 "ai"
  • 类型解析由 _stream_payload_type 完成,同时兼容两种形态:DeerFlow 自家 Gateway 发出的 model_dump() 形状,以及 LangChain to_json() 构造器形状(其顶层 type 是字面量 "constructor",类名在 id 路径尾部);
  • str 载荷被彻底拒绝——它不携带类型信息,无法归属到 assistant,且 DeerFlow 内没有任何路径产生它(runtime/serialization.pyserialize_messages_tuple 总是发出 [message_dict, metadata])。

此外,manager 会吞掉一次失败的流式:在释放入站去重键之前发布该流的最终出站,这样 provider 的重投递可以重试而不超车终态回复。

消息身份与文件归属:owner-scoped 设计

channel_user_id 的注入链路

对于用户自有通道连接,Gateway 只接受内部认证的通道调用方在顶层 body.context 里给出的 channel_user_id,会从两个 free-form 的 body.config 段中清掉它,并且只写入运行时上下文而绝不写入 configurable(后者会被 checkpoint 持久化)。

该 id 通过 bash_tool 以固定环境变量 DEERFLOW_CHANNEL_USER_ID 暴露给沙箱命令,注入方式值得注意:是带 shell 引号的命令串前缀 export VAR=<id>; ,而不是走 execute_command(env=...) 通道——后者保留给 request-scoped secrets,且会切换 AioSandboxbash.exec 路径(image ≥ 1.9.3、每次调用新建会话)。每次调用注入保证了群聊身份正确(一线程一沙箱、多人发送)。id 有效时注入 export VAR=<id>; ,id 为空/非 str/超 256 字符上限时注入 unset VAR; ——AIO 无 env 路径复用持久 shell 会话(这正是类锁存在的原因,#1433),裸命令可能解析到先前发送者 export 的过期 id;unset 封死这个窗口。Windows 本地沙箱(PowerShell/cmd.exe 无 export/unset)不注入。该 id 还随 task 委派传播:task_tool 捕获派发轮次的 id,subagent executor 把它转发进 subagent 的运行时上下文(与 guardrail 归属字段相同)。

需要特别说明的边界:运行时上下文值在 Gateway/guardrail 边界是授权级的,但导出的 shell 变量仅具信息意义——任何 bash 命令都能改写自身环境,skills 不得把该 shell 变量当作已认证身份。相关测试见 test_gateway_services.pytest_channel_user_id_env.py

文件存储与「运行身份」对齐

入站文件、上传与产物都暂存在 DeerFlow owner 的 bucket 下,落在 Agent run 读写的位置:users/{user_id}/threads/{thread_id}/user-data/{uploads,outputs}ChannelManager._handle_chat 通过 _channel_storage_user_id(msg) 只解析一次存储 owner(已净化的 owner id;未绑定但启用了认证的通道回退到 safe(msg.user_id),与 _resolve_run_params 的 run 身份对齐;仅当完全没有可用身份时返回 None),并以 user_id= kwarg 穿透整个文件管线:

  • Channel.receive_file(msg, thread_id, user_id=...)——owner 绑定通道把下载文件持久化到 owner bucket 而非默认 bucket;
  • FeishuChannel._receive_single_file / DingTalkChannel._receive_single_file——规范化 provider 文件名、在同一通道锁下声明无冲突 basename,并经 write_upload_file_no_symlink 写入;返回的 basename 同时驱动 agent 可见虚拟路径与非本地沙箱同步;
  • sandbox_files.py——非挂载的飞书/钉钉同步获取非释放的执行 holder,跨反复取消排空阻塞中的 update_file worker,直到最后一次沙箱操作完成才释放,防止并行 run 在上传中途关闭共享 client;
  • _ingest_inbound_filesensure_uploads_dir / get_uploads_dir_resolve_attachments / _prepare_artifact_delivery 全部经同一 kwarg owner 化。

缓存值同时用于阻塞(runs.wait)与流式(_handle_streaming_chat)两条路径,所以即便 receive_file 返回改写过的 InboundMessage,上传与产物投递也总是落在同一 bucket。且该 bucket id 与 _resolve_memory_user_id 解析的内存 bucket 一致(都经 make_safe_user_id 归一化)。

配置指南

基础配置(config.yaml → channels

来自 config.example.yaml(IM Channels Configuration 段)与文档:

channels:
  # LangGraph-compatible Gateway API 基础 URL,用于线程/消息管理(默认 http://localhost:8001/api)
  langgraph_url: http://localhost:8001/api
  # Gateway API URL,用于 /models、/memory 等辅助查询(默认 http://localhost:8001)
  gateway_url: http://localhost:8001
  # 队列中最大待处理入站消息数,必须为正整数
  inbound_queue_maxsize: 1000
  # 常驻入站处理 worker 的固定数量,必须为正整数
  max_concurrency: 5
  # 排空已接受入站工作的秒数,之后取消活跃 handler;必须为非负有限数值
  shutdown_grace_period_seconds: 3
  # 可选:所有 IM 通道的默认会话设置
  session:
    assistant_id: lead_agent   # 或自定义 agent 名
    config:
      recursion_limit: 100
    context:
      thinking_enabled: true
      is_plan_mode: false
      subagent_enabled: false

Docker Compose 注意事项:IM 通道跑在 gateway 容器内,localhost 指向的就是该容器自身。应把 langgraph_url 设为 http://gateway:8001/apigateway_url 设为 http://gateway:8001,或通过环境变量 DEER_FLOW_CHANNELS_LANGGRAPH_URL / DEER_FLOW_CHANNELS_GATEWAY_URL 覆盖。所有通道均使用出站连接(WebSocket 或轮询),不需要公网 IP

各平台通道参数

逐平台的凭证与可选参数(完整以 config.example.yaml 注释段为准):

  • feishuapp_idapp_secret;可选 domain(中国默认 https://open.feishu.cn,国际 https://open.larksuite.com);
  • slackbot_tokenxoxb-...)、app_tokenxapp-...,Socket Mode)、allowed_users(空 = 放行所有人,也可填单个 Slack user id,但推荐列表形式);
  • telegrambot_tokenallowed_users(空 = 放行所有人)、可选 rich_messages(Bot API 10.1 最终 Markdown 响应的 Rich Messages);
  • wechatbot_tokenilink_bot_id,可选 qrcode_login_enabled(缺 token 时允许首次扫码引导)、polling_timeout / polling_retry_delay / qrcode_poll_interval / qrcode_poll_timeout 等定时参数(只接受正的有限秒数,非法值回退默认值,避免轮询热循环或永眠)、state_dirmax_inbound_image_bytes / max_outbound_image_bytes / max_inbound_file_bytes / max_outbound_file_bytesallowed_file_extensions
  • wecombot_idbot_secret
  • dingtalkclient_idclient_secretallowed_users、可选 card_template_id(AI Card 流式更新模板 id);
  • discordbot_tokenallowed_guilds(空 = 放行所有 guild)、mention_only(true 则仅被 @ 时响应)、allowed_channelsthread_mode
  • github:操作级 kill-switch enabled,外加 mention-required 触发的 default_mention_login
  • buzzrelay_urlprivate_key(hex 或 nsec1...)、allowed_users(pubkey 列表;deny-by-default,空 = 谁也不放行)、require_mentionmention_free_channels。需要 buzz 依赖 extra(uv sync --extra buzz)。

示例配置中默认全部注释关闭;需将对应段反注释并填充真实凭证。

用户自有通道连接(config.yaml → channel_connections

channel_connections 是在既有 channels.* 运行时配置之上的用户绑定层,不是 provider bot 凭证的替代品。要点:

  • 默认关闭;当前实现不需要公网 IP、OAuth 回调 URL 或 provider webhook 路由
  • Telegram 走既有长轮询 worker 上的深链 /start <code> 流程;Slack、Discord、飞书/Lark、钉钉、WeChat、WeCom 走各自出站通道 worker 上的 /connect <code>
  • 前端 API:GET /api/channels/providersGET /api/channels/connectionsPOST /api/channels/{provider}/connectDELETE /api/channels/connections/{connection_id}
  • provider 级 connection_status 反映用户最新连接行;未绑定时为 not_connected,唯一例外是禁用认证的本地模式下,已配置且运行的通道报 connected(因为所有通道消息都已路由给默认用户);
  • Slack 回复默认使用 channels.slack 里配置的 operator bot token,除非存在 per-connection 凭证;不可读或损坏的已存凭证按不可用处理;
  • WeChat 定时参数只接受正的有限秒数,非法值回退默认,保证轮询既不会热循环也不会永眠;WeCom 对每个通道实例串行化 start()/stop()

文档还明确提示两处安全语义:

  1. connect-code 的消费先于 allowed_users 过滤:入站 worker 在应用 allowed_users 过滤之前消费合法的 /connect <code>(或 Telegram /start <code>),这样刚被加入白名单但尚未绑定的用户能通过浏览器流程完成首次绑定。推论是 allowed_users 不是绑定时的防线——任何持有效 code 的发送者都能消费它。绑定安全模型建立在 code 机密性之上:secrets.token_urlsafe(16)、600 秒 TTL、一次性 consume_oauth_state、code 只出现在发起浏览器里(绝不回显到聊天)。allowed_users 仍然拦截普通(非绑定)消息。
  2. 单一活跃 owner 的转移语义:外部身份以 (provider, external_account_id, workspace_id) 为键,最新成功的绑定胜出——upsert_connection 会吊销同一身份的其他 owner 的活跃行(ownership transfer)。该不变式由数据库层部分唯一索引 uq_channel_connection_active_identityWHERE status != 'revoked')强制,所以不同 owner 的并发 connect 不可能双双 connected;落败的写入者会对着现态重试。因此 find_connection_by_external_identity 的解析是确定性的。

GitHub 事件驱动 Agent:webhook → run 的完整闭环

GitHub 通道本质是 webhook 驱动的 IM 通道,相关文档与架构见 backend/docs/GITHUB_AGENTS.md。要点:

  • 自定义 Agent 在自身 config.yaml 里声明 github: 块,绑定仓库与事件触发器;
  • webhook 路由 POST /api/webhooks/github 默认 fail-closed:只有设置了 GITHUB_WEBHOOK_SECRET 才挂载,且豁免认证/CSRF——真实性由 HMAC 保证;
  • 注册表缓存以 agent 存储的不透明签名为键:文件存储用 agent 配置的 mtime,数据库存储则对排序后的 owner/name/config/soul 内容做哈希,保证同时间戳写入也能让 webhook 路由失效;
  • 出站默认只打日志:每个 Agent 在自己的沙箱里用 gh CLI(gh issue commentgh pr commentgh pr create 等)自主发帖,所以多个 agent 在同一事件上 fan-out 时静默是廉价的。正因如此 manager 使用 fire_and_forget=Trueruns.create() 返回 pending 即结束;
  • 线程确定性:preferred_thread_id = UUID5(repo, number, agent_name),同一 issue/PR 的同一 agent 永远收敛到同一线程;另含 mention-handle 优先级链、经 GH_TOKEN/GITHUB_TOKEN 的 per-call extra_env token 生命周期,以及窄化的 ConflictError(HTTP 409)线程创建竞态恢复。

忙时追问缓冲(issue #4121)

由于出站只打日志,此前忙线程上对 ConflictError 回复的 THREAD_BUSY_MESSAGE 对评论者不可见——run 进行中发的评论看起来像被静默忽略。当 ChannelRunPolicy.buffer_followups_on_busy=True(GitHub 默认)时,runs.create()ConflictError 会把触发消息追加进 per-thread 内存缓冲(ChannelManager._followup_buffers)——按 GitHub delivery id 去重,上限 FOLLOWUP_BUFFER_MAX_PER_THREAD(20,溢出丢最旧并打 WARNING)。之后:

  1. 线程上首次成功的 runs.create() 会捕获 run_id 并派生后台 watcher,订阅该 run 的 StreamBridge 流;
  2. watcher 观察到 END_SENTINEL 后,把缓冲内至多 FOLLOWUP_DRAIN_BATCH_SIZE(10)条合并成一条 <followups-while-busy> 包裹的输入,再触发后续 runs.create()
  3. 该后续 run 自身也走同样 watcher,因此超过一批的积压会链式进入下一轮 drain,而不是长成一条无界输入;
  4. 若后续 runs.create() 又撞上 ConflictError(如 Web UI 手动轮次或定时 run 抢到同一线程),批量数据重新入队而非丢失,等下次该线程成功创建并 watch run 时再重试。

管线细节:watcher 需要 Gateway 的 StreamBridge 单例,而 ChannelManager 此前拿不到它。该依赖从 app.py lifespan(langgraph_runtime 已在那里设置 app.state.stream_bridge)经 start_channel_service(get_stream_bridge=...)ChannelService.__init__ChannelManager.__init__ 一路穿入,形态是零参闭包,镜像了同一 lifespan 函数里为 ScheduledTaskService 使用的 launch_run=lambda **kwargs: ... 模式。测试中直接构造的、无 watcher 的 ChannelManager 仍能安全缓冲,只是没有自动 drain。

范围限制:缓冲与 watcher 状态是 per-process 内存态的。GATEWAY_WORKERS>1 或多 pod 下,评论若路由到与忙线程 agent 不同的 worker 进程就看不到该缓冲。这是文档明示、刻意延后的限制(与 issue #4120 描述的跨 pod 缺口同形,需共享缓冲存储或 IM leader 选举才能闭合)——单进程/单 pod(安全默认)没有任何正确性问题。

Buzz(Nostr relay)通道:重连、水位与信任

buzz.py 是 Buzz(Nostr relay)的实现,需 buzz extra。它在一个 NIP-42 认证 websocket 上工作,叠加 pubkey 白名单 + mention/DM/thread-follow 门控,并通过原地 kind-40003 编辑实现流式回复。operator 视角的完整订阅模型与信任模型另见 backend/docs/IM_CHANNEL_CONNECTIONS.md

订阅模型:relay 把 kind-9 聊天事件扇出给 channel-scoped 订阅(文档记录了对真实 relay 的实测:REQ {"kinds":[9]} 被接受并回 EOSE 却永不收事件;REQ {"kinds":[9],"#h":[uuid]} 有效;多值 #h 匹配不到任何东西)。所以每个连接在 NIP-42 auth 之后要重建三类订阅:

  • buzz-discovery{"kinds":[39000]},历史查询,恰好返回本身份所属各通道(每条一个存储事件后 EOSE;给它加 #p 反而返回零,不要「收窄」它);
  • buzz-membership{"kinds":[44100,44101],"#p":[us],"since":<连接时刻 − MEMBERSHIP_LOOKBACK_SECONDS>},relay 签名的成员加入/移除通知;
  • 每个已发现通道一个 buzz-chat-<uuid>

文档特别强调「membership 订阅必须限定 LIVE 事件,这是 load-bearing 的」:buzz-relay 会存储 44100/44101 并按 newest-first 提供历史(默认上限 2000),无界 filter 会在每次连接时重放全部 membership 历史,造成 M+1 次 discovery、反复订阅已被移除的通道、被历史 44101 瞬时退订仍属于的通道。since 锚定在 socket 打开时刻(_session_started_at)减去 MEMBERSHIP_LOOKBACK_SECONDS(60 秒)的松弛,同时覆盖 relay 时钟偏移与握手期间的成员变更。

CLOSED 恢复与瞬态判定:relay 的 CLOSED 帧意味着该 socket 上所有订阅都被静默丢弃,因此 _handle_closed 会重发,受 MAX_RESUBSCRIBE_ATTEMPTS(3)每订阅 id 每连接 且每 auth epoch 的约束。第一次重试立即执行(常见场景是一次性抖动),后续按 1s、2s 退避,且在 read loop 内内联等待而不是起可能比自身 socket 更长寿的后台任务。auth-required: 出现在 NIP-42 握手完成前是预期的引导序列而非拒绝:连接器立即打开控制 REQ(relay 可能提供未认证读),关闭的 relay 回 auth-required:AUTH challenge,认证分支随后重开一切。该情形只打 DEBUG,不消耗永久拒绝分支,也不占重试预算。per-socket 的 _auth_completed 标志是分界:签名 AUTH 事件发出后置位,会话进入/退出及 stop() 清除。_is_transient_close 决定是否值得重试:NIP-01/NIP-42 的 auth-required:/restricted:/blocked:/mute:/invalid:/pow: 前缀和 buzz-relay 自己的移除/吊销措辞是永久的(不与 relay 争论已不属于我们的通道);rate-limited:/error:/无理由/不可识别都是瞬态——默认偏向「继续收听」,因为静默失聪正是该机制要消灭的失败,且尝试预算约束了误判。聊天 CLOSED 只对已入 _chat_subscriptions 的通道恢复——CLOSED 由 relay 提供,对未知通道动手会让 relay 仅凭点名就诱导出一个订阅。

对远端供给状态的三重上限MAX_CACHED_CHANNELS(512)封顶 kind-39000 元数据缓存、水位 map 与重订阅尝试 map;MAX_CHANNEL_SUBSCRIPTIONS(256,远低于 buzz-relay 每连接 1024 上限)封顶活跃聊天订阅——到顶时新通道被拒绝并在 WARNING 中指名,而不是驱逐工作中的订阅。

一个文档明示的已知边界(记录而非修复):relay 即使给了 since,也把每个订阅的历史投递封顶在 2000 条、newest-first;因此断线期间单个通道里超过 2000 条未读会丢最旧的——relay 从不发送它们,水位照常前进越过。这是唯一剩余可跳过的路径,其余全部偏向重放。

信任模型:每个入站 EVENT 在唯一的 handle_relay_frame 咽喉点被认证——从投递载荷重算 NIP-01 id,并对声称的 pubkey 验证 BIP-340 Schnorr 签名(buzz_nostr.pyverify_event,纯且全:畸形输入返回 False 而绝不抛异常)。因此作为白名单与 /connect 绑定授权主体的 ev["pubkey"] 无法被 DeerFlow operator 不运营的 relay 伪造。仍被信任的是 kind-39000 通道元数据的作者——任何成员都能签一份。per-channel 订阅现在由 discovery 驱动,伪造 kind-39000 有两个后果:可把通道标记为 type: "dm"(放松该通道的 require_mention),也能诱导对伪造者选定的通道建立聊天订阅。但两者都不会让东西被执行allowed_users 与逐事件签名校验是独立门禁;被诱导的订阅只是 relay 把自己的流量读回给一个丢弃它的订阅者,且有 MAX_CHANNEL_SUBSCRIPTIONS 兜底(拒绝而非驱逐)。伪造 kind-44100 同理,只是其 p tag 会被本地复核,必须至少点名我们。要闭合它需要一个配置的可信 relay pubkey——而 relay_url 不是。

与其它通道不同,Buzz 的 allowed_usersdeny-by-default(空 = 谁也不放行),因此 start() 在空配置时打 WARNING,每次 drop 打 DEBUG。

重订阅水位按通道独立since 游标是 per-channel 的,只为实际处理过的事件前进,且绝不越过 now + MAX_FUTURE_SKEW_SECONDS——因为 created_at 是 peer 供给的,一条未来时间戳事件否则会让连接器永久失聪。per-channel 而非全局是安全关键的一半:订阅按通道,一个共享游标等于「任意通道看到的最新事件」,繁忙通道会把游标拖过安静通道的未读消息,导致下次重连跳过它们(真实 relay 实测:某身份的三个通道相距约 28 小时)。per-channel 游标只可能带来重复投递(由 manager 的 event_id 去重吸收),被逐出的游标退化为「无 since」即 relay 默认 backlog——两者都偏向重放。

出站隐藏上下文防护:Buzz 的 send() 直接拒绝发布携带隐藏模型上下文包装的文本(<memory><durable_context_data><system-reminder>——即 _HIDDEN_CONTEXT_MARKERS),命中时打 ERROR 并在被拦截的 is_final 上清空流簿记。这是 manager 白名单之后的纵深防御,之所以放在 Buzz 而非其它兄弟连接器,是因为在 Buzz 上泄露是永久的:每次流式更新都是不可变的公开 Nostr 事件,纠错编辑只改变客户端渲染,原始泄露事件仍留在 relay 上。匹配基于字面开标签,因此一条只是「谈论记忆」的回复仍会照常发布。

测试与验证

通道子系统的行为在 backend/tests 中有大量测试覆盖,可作为实现证据的延伸阅读,例如:

小结:给接入者的实践清单

  • 同一种 langgraph-sdk 客户端 + 内部认证 + 匹配 CSRF 对打通通道与 Gateway,线程永远服务端管理;
  • 各平台按自身能力选策略:交互聊天优先 runs.stream 增量回复;出站自主的平台(GitHub)用 fire_and_forget;忙时追问敏感的平台开 serialize_thread_runs/buffer_followups_on_busy
  • 流式出站内容必须是 assistant 类型白名单,并让 Buzz 类不可变平台在连接器层再做隐藏上下文标记拦截;
  • 所有与用户绑定的入站/文件/产物都收敛到 owner bucket 与统一的 run 身份,shell 里的 DEERFLOW_CHANNEL_USER_ID 仅供展示;
  • GitHub 走 webhook + HMAC + agent 自主 gh 回写,busy 时用 per-thread 缓冲合并;Buzz 则围绕 NIP-42、per-channel 订阅、per-channel 水位与「偏向重放」的故障模型构建。

理解上述架构后,从 config.example.yamlchannels 段开始,按你的平台填好凭证即可把 DeerFlow 接上日常聊天工具,让长时程任务在消息里跑起来。

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