news 2026/9/10 13:13:01

FastGPT 微信个人号 ClawBot 渠道设计:OutLink 轮询链路、双队列与幂等消息消费

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
FastGPT 微信个人号 ClawBot 渠道设计:OutLink 轮询链路、双队列与幂等消息消费

FastGPT 微信个人号 ClawBot 渠道设计:OutLink 轮询链路、双队列与幂等消息消费

【免费下载链接】FastGPTFastGPT is a knowledge-based platform built on the LLMs, offers a comprehensive suite of out-of-the-box capabilities such as data processing, RAG retrieval, and visual AI workflow orchestration, letting you easily develop and deploy complex question-answering systems without the need for extensive setup or configuration.项目地址: https://gitcode.com/GitHub_Trending/fa/FastGPT

导读

本文是 FastGPT 仓库中微信个人号发布渠道(ClawBot)的唯一设计文档讲解,完整覆盖 iLink 长轮询消息摄入、wechatPoll/wechatReply双队列编排、syncBuf消费游标推进、连续失败熔断与多实例恢复等核心机制。读完本文,你将理解 FastGPT 如何在不依赖长期双写的情况下,把“拉取消息”与“生成回复”彻底解耦,并掌握 at-least-once 摄入、stalled retry 防重复回复、渠道状态机等可直接复用的工程模式。

1. 目标与边界:为什么微信渠道需要独立的轮询与回复设计

微信个人号(ClawBot)渠道的链路是:通过 iLink 长轮询(long polling)接收用户消息,消息进入 FastGPT OutLink 工作流生成回复,再调用 iLink 将回复发送给用户。整条链路以 WechatMQService 提供的 BullMQ/Redis 队列为骨架,领域 processor 位于 mq.ts。

该设计需要同时满足六条硬性目标:

  • 每个shareId同一时刻只有一条 poll 链(单实例轮询,避免消息重复摄入);
  • 拉取与回复解耦,慢回复不阻塞后续消息摄入;
  • enqueue、stalled retry、多实例恢复都不会产生重复回复;
  • 推进syncBuf之前,消息必须已经进入 reply queue;
  • 渠道下线、登出或连续错误后能够停止续链(不再无限轮询)。

在职责边界上,DAL(Data Access Layer)只拥有 Redis Cache 与 BullMQ 数据合同;iLink client、Mongo 状态、消息解析、工作流调用和渠道编排全部保留在 service/project 层。这一边界与>Wechat Publish UI | v QR Login API ---- WechatQrLoginCache | v MongoOutLink.app = { token, baseUrl, syncBuf, status } | v wechatPoll Queue ---- getUpdates ---- groupMessagesByUser | v wechatReply Queue | v provider adapter -> runOutlinkRuntime -> sendMessage

链路要点:

  • 用户在前端发布页触发扫码登录,登录二维码与登录态通过WechatQrLoginCache(QR JSON,480 秒 TTL)暂存;
  • 登录成功后,渠道配置写入MongoOutLink.app(token、baseUrl、syncBuf、status 等字段的事实来源);
  • 后台轮询 Worker 消费wechatPoll队列,调用 iLinkgetUpdates长轮询;拉到消息后经groupMessagesByUser按用户聚合;
  • 聚合结果投递到wechatReply队列,由回复 Worker 经 provider adapter →runOutlinkRuntime生成回复 →sendMessage发回微信。

队列合同定义在 wechat.ts,其中WechatMQService暴露getPollQueue/getReplyQueue/addPollJob/addReplyJob/removePollJob等操作,并允许注入BullMQBinding便于单元测试,不在模块加载时连接 Redis。

3. 渠道状态:MongoOutLink.app 的字段与状态机

WechatAppType的稳定字段由 Zod Schema 定义在 packages/global/support/outLink/type.ts:

字段含义
tokeniLink 登录 token(默认空字符串)
baseUrliLink API 地址(默认https://ilinkai.weixin.qq.com
accountIduserId登录身份
syncBuf下一次getUpdates的消费游标(默认空字符串)
statusonlineofflineerror(默认offline
loginTime最近登录时间
lastError停止轮询的最近错误

状态转换:

offline --扫码确认--> online --主动登出/停用--> offline | +--连续失败达到阈值--> error offline/error --重新扫码--> 清空 syncBuf --> online

Worker 每次执行前都会重新读取MongoOutLink:记录不存在、渠道非online或 token 缺失时直接停止处理;completed/failed listener 只有确认渠道仍可用时才续链。具体的检查逻辑见 mq.ts 的pollImpl,它依次校验 outLink 是否存在、app.status === 'online'app.token非空,任一不满足即 throw,从而让 failed listener 走停链分支。

4. Queue 合同:poll 与 reply 两张队列的参数细节

4.1 Poll Queue(消息摄入)

属性合同
QueuewechatPoll
Job namewechatPublishPoll
Job data{ shareId }
Job IDwechat-poll:${shareId}
ConcurrencyWECHAT_CHANNEL_CONCURRENCY
Lock120 秒
Hard timeout120 秒
Stalled interval30 秒
Terminal retentioncompleted/failed 立即删除

Poll job 主要阻塞在约 35 秒的 iLink 长轮询 I/O 上,不执行工作流,拉到消息后只负责解析、分组和投递 reply job。Hard timeout 的语义值得注意:源码在processWechatPollJob中用Promise.race([pollImpl(job), timeout])实现兜底(mq.ts),因为固定 jobId 一旦被一个 hang 住的 processor 长期占用,整个渠道的轮询链就会永久阻塞——硬超时保证该 job 最多 120 秒内让出。

续链节奏:

  • 有消息时 completed listener 立即续链,清空积压;
  • 空响应延迟 10 秒(EMPTY_POLL_DELAY_MS = 10_000),避免上游秒回空包时 completed → 立即续链退化成热循环;
  • failed listener 在渠道仍online时延迟 10 秒重试(FAILURE_BACKOFF_MS = 10_000)。

Worker 的具体配置在initWechatPollWorker(mq.ts):lockDuration: 120_000防止长轮询期间 job 被误判为 stalled;stalledInterval: 30_000每 30 秒检查活跃度;removeOnComplete/removeOnFail均设为立即清理。

4.2 Reply Queue(回复生成)

属性合同
QueuewechatReply
Job namewechatPublishReply
Job datashareIduserIditemscontextTokenlastMsgId
Job IDwechat-reply:${shareId}:${lastMsgId}
ConcurrencyWECHAT_CHANNEL_CONCURRENCY
Lock30 分钟
Stalled interval60 秒
Failed retention500 条或 7 天

Reply processor(processWechatReplyJob,mq.ts)调用 provider adapter →runOutlinkRuntime生成回复,并以稳定的messageId=lastMsgId由聊天写入层保证业务幂等。这里存在两道幂等防线

  1. 队列 jobId 去重wechat-reply:${shareId}:${lastMsgId}是确定 jobId,同一消息重复入队时 BullMQ 自动去重;
  2. 业务 messageId 幂等:即使 job 被 stalled 重试或 processor 中途失败后重放,写入层仍以lastMsgId判断是否已消费,避免副作用重复。

WECHAT_CHANNEL_CONCURRENCY在 packages/service/env.ts 中定义为最小 10、默认 1000 的整数环境变量,poll 与 reply Worker 共用同一并发上限;并发请求数按渠道实际部署规模调低,例如测试用例如 mq.test.ts 中将其设为 1。

5. Poll 处理顺序:固定七步,失败即回退

一次成功 poll 的执行顺序是固定的:

1. 校验渠道状态和 token 2. 使用当前 syncBuf 调用 getUpdates 3. 判断 API ret/errcode 4. 按 userId 聚合消息 5. 并行投递 reply jobs 6. 全部投递成功后更新 Mongo syncBuf 7. completed listener 调度下一条 poll job

在源码pollImpl中的实现要点:

  • 步骤 3ret !== 0 || errcode !== 0视为 API 错误,进入失败计数流程(详见第 7 节);
  • 步骤 5Promise.all(groups.map(...))并行投递 reply jobs,jobId 为replyJobId(shareId, lastMsgId)
  • 步骤 6:只有全部 reply job 入队成功后,才updateOne推进app.syncBuf = resp.get_updates_buf(mq.ts)。

第五步失败时不得推进syncBuf——下一次 poll 会重新拉取同一批消息,由replyJobId去重,从而形成 at-least-once 摄入 + 幂等消费。这条顺序保证是"推进游标前先入队"的核心约束,直接对应设计目标中的第四点。

另一个设计约定是:poll processor 本身不续链。续链统一由 Worker 的completed/failedlistener 负责(scheduleNextPoll),避免 return、throw 和 timeout 三个分支各自维护调度逻辑导致分叉。completed 事件中hadMessages为 false 时追加EMPTY_POLL_DELAY_MS延迟。

6. 消息合并语义:同用户同周期只回一条

groupMessagesByUser(messageParser.ts)在单个 poll 响应内按userId聚合:

  • 文本使用换行拼接(同一用户的多条文本依次 push 进同一组的items);
  • contextTokenlastMsgId取该用户最后一条消息的值;
  • 同一用户在同一 poll 周期只生成一个 reply job 和一次合并回复;
  • 跨 poll 周期生成独立的 reply job,但共享chatId = wechat_${shareId}_${userId}(见 adapter.ts),上下文连续;
  • 多个用户生成多个 reply job,并行处理。

isSupportedMessageItem(messageParser.ts)只保留能转换为 runtime query 的消息项:文本要求有text_item.text,语音要求有voice_item.text(当前 iLink 通过该字段提供上游转写结果),图片要求 CDN 具备下载地址,文件/视频要求同时具备aes_key与下载地址。过滤空项避免空 job 进入工作流

引用消息先转换成带引用前缀的 query item 再与当前消息合并;图片、文件、语音分别通过现有 OutLink 文件/文本处理流程进入工作流(adapter.normalizeMessageresolveQuery负责把媒体资源下载、解密、上传 S3 后合成 query)。设计文档同时注明:微信当前不提供获取被引用消息的接口(截至 2026.7.31),引用解析暂以 @todo 形式保留在 adapter.ts。

7. Redis 与持久化合同:固定窗口失败计数

数据Physical key / 存储TTL/语义
QR Loginfastgpt:cache:publish:wechat:qrcode:${outLinkId}:${tmbId}QR JSON,480 秒
Poll failurefastgpt:cache:wechat:publish:failures:${shareId}integer,300 秒
渠道配置MongoOutLink.apptoken、syncBuf、status 的事实来源

失败计数使用INCRBY + EXPIRE NX:从第一次失败起固定 300 秒,不在后续失败时刷新 TTL;成功 poll 将值重置为带 TTL 的字符串0。这一"固定窗口"与旧版"每次失败刷新 TTL"的滑动窗口语义不同,属于明确业务变化,在 contenteditable="false">【免费下载链接】FastGPTFastGPT is a knowledge-based platform built on the LLMs, offers a comprehensive suite of out-of-the-box capabilities such as data processing, RAG retrieval, and visual AI workflow orchestration, letting you easily develop and deploy complex question-answering systems without the need for extensive setup or configuration.项目地址: https://gitcode.com/GitHub_Trending/fa/FastGPT

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/10 13:12:43

STM32三闭环PID控制:位置-速度-电流实时实现指南

简介:本资源是一套基于STM32F4平台实现直流有刷电机位置-速度-电流三闭环PID控制的完整嵌入式工程,面向自动化、机器人及运动控制方向的嵌入式开发者与高校高年级学生,解决电机多层级动态响应精度低、参数整定难、软硬件协同调试复杂等实际问…

作者头像 李华
网站建设 2026/9/10 13:10:46

基于LSTM的锂电池剩余寿命预测Matlab实现与调优

简介:面向锂电池健康管理、电池管理系统相关研究人员与工程师,提供一套基于长短期记忆(LSTM)神经网络的锂电池剩余寿命预测Matlab完整实现。资源直击剩余使用寿命(RUL)预测中数据准备与网络搭建两大难点&am…

作者头像 李华
网站建设 2026/9/10 13:05:47

Java+SpringBoot教学平台开发实战与优化策略

1. 项目概述:JavaSpringBoot课程教学平台的设计初衷作为一名经历过多次毕业设计指导的老手,我见过太多学生在这个环节踩坑。这个基于JavaSpringBoot的课程教学管理平台,本质上是要解决传统教学中的三个痛点:课程资源分散、师生互动…

作者头像 李华