DeerFlow IM Channels:从飞书 Slack Telegram 接入 Agent
网页 Gateway 解决了浏览器侧的登录与 API。实际办公里更常见的入口是飞书、Slack、Telegram、钉钉:在聊天窗口 @ 机器人一句,就把任务丢给 agent。各平台消息格式、认证方式和限流策略都不一样,按平台各写一套调度逻辑会迅速失控。DeerFlow 的 IM Channels 把外部消息翻译成统一入站结构,经 MessageBus 交给 agent,再把回复翻译回原平台。渠道只负责收发,调度与线程映射集中在 Manager;新增一个平台,主要是实现一个 Channel 子类,不必复制半套运行时。
四件套与数据流
Channel(app/channels/base.py)抽象三个动作:start 开始监听,stop 优雅停止,send 把 OutboundMessage 发回平台。仓库里已有 Telegram、Slack、飞书、钉钉、企业微信、Discord 等实现。统一采用 WebSocket 或 long-polling,不依赖公网 IP 与入站 Webhook,本机或内网部署也能挂机器人。对本地开发和防火墙后部署更友好。
MessageBus(message_bus.py):入站用 asyncio.Queue,出站是 listener 列表。渠道 publish_inbound,Manager get_inbound;处理完 publish_outbound,各渠道在回调里按 channel_name 过滤后调用自己的 send。入站 FIFO 保证顺序,出站广播式分发;渠道与 agent 不直接耦合,一边慢了不会把另一边的协议细节扯进来。
InboundMessage / OutboundMessage:标准字段包括 channel_name、chat_id、user_id、text、msg_type、topic_id、附件列表等。topic_id 决定映射到哪个 DeerFlow thread。例如 Telegram 私聊常设 None,同一私聊共用一个 thread;群聊常用 message id,每条新消息可开新 thread。飞书群聊类似,用消息 id 做 topic。搞清楚 chat_id(回哪)和 topic_id(对应哪个对话)之后,IM 与 Web thread 模型就能对上。
ChannelManager(manager.py):_dispatch_loop 循环取入站消息 → store.resolve_thread_id(channel, chat_id, topic_id)(没有则创建)→ asyncio.create_task(_run_agent(...))。_run_agent 经 Gateway API 或 LangGraph SDK 对流式 chunk 累积,再发 outbound;非 final 可先推流式片段,is_final=True 表示收束。用户侧大致就是:@ 一下机器人,过几秒开始出字,最后一条完整回复落在同一聊天窗口。
平台差异缩在 Channel 子类
以飞书为例:WebSocket 推送到达后 _make_inbound,设置 topic_id,再 publish_inbound。Telegram 的 start 用 python-telegram-bot 注册文本 handler,并把 long-polling 放在独立 daemon 线程,避免同步 SDK 堵住 asyncio 主循环。私聊 / 群聊分支决定 topic_id;send 对非 final 编辑同一条消息做打字机效果,final 再发完整条,并对群聊限流做节流。supports_streaming = True 标明该渠道支持边生成边改消息。
1 | outbound = OutboundMessage( |
出站回调里每个渠道都会收到广播,但只有名字匹配的渠道真正调用平台 API。读某个 Channel 实现时,盯住「入站翻译」和「出站 send」两头即可。
绑定码:IM 身份对齐 DeerFlow 用户
IM 侧有自己的 user id(如飞书 open_id),DeerFlow 有自己的用户体系,二者靠绑定关联。网页端 POST /{provider}/connect(channel_connections.py)生成一次性 code(默认约 10 分钟有效),提示用户在机器人聊天里发送 /connect <code>。渠道解析到 code 后 consume_oauth_state,再 upsert_connection(owner_user_id, provider, external_account_id)。之后每条入站消息经 attach_connection_identity 查表,填上 connection_id 与 owner_user_id,后续权限与数据隔离才能落到正确用户。绑定码一次性且过期,降低泄露后被复用的风险。
ChannelService 与长连接取舍
service.py 用懒加载注册表把名字映射到类路径(feishu → app.channels.feishu:FeishuChannel 等),只实例化配置里 enabled: true 的渠道,避免无用依赖拖进进程。start 先启 Manager,再逐个 _start_channel。restart_channel 可停旧启新,并用 asyncio.to_thread 重载磁盘配置,网页上改完飞书参数不必重启整个进程。
相对 Webhook,长连接要保活和断线重连,但换来本地友好:不用公网 IP、不用开端口映射、不必先配域名和 HTTPS。读代码时先跟通「平台事件 → InboundMessage → Bus → Manager → Agent → OutboundMessage → Channel.send」,再看绑定与 topic_id 映射;平台差异基本都缩在各个 Channel 子类里,Gateway 与 Runtime 仍走原来的鉴权与流式路径。
DeerFlow IM Channels:从飞书 Slack Telegram 接入 Agent