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:
| 字段 | 含义 |
|---|---|
token | iLink 登录 token(默认空字符串) |
baseUrl | iLink API 地址(默认https://ilinkai.weixin.qq.com) |
accountId、userId | 登录身份 |
syncBuf | 下一次getUpdates的消费游标(默认空字符串) |
status | online、offline、error(默认offline) |
loginTime | 最近登录时间 |
lastError | 停止轮询的最近错误 |
状态转换:
offline --扫码确认--> online --主动登出/停用--> offline | +--连续失败达到阈值--> error offline/error --重新扫码--> 清空 syncBuf --> onlineWorker 每次执行前都会重新读取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(消息摄入)
| 属性 | 合同 |
|---|---|
| Queue | wechatPoll |
| Job name | wechatPublishPoll |
| Job data | { shareId } |
| Job ID | wechat-poll:${shareId} |
| Concurrency | WECHAT_CHANNEL_CONCURRENCY |
| Lock | 120 秒 |
| Hard timeout | 120 秒 |
| Stalled interval | 30 秒 |
| Terminal retention | completed/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(回复生成)
| 属性 | 合同 |
|---|---|
| Queue | wechatReply |
| Job name | wechatPublishReply |
| Job data | shareId、userId、items、contextToken、lastMsgId |
| Job ID | wechat-reply:${shareId}:${lastMsgId} |
| Concurrency | WECHAT_CHANNEL_CONCURRENCY |
| Lock | 30 分钟 |
| Stalled interval | 60 秒 |
| Failed retention | 500 条或 7 天 |
Reply processor(processWechatReplyJob,mq.ts)调用 provider adapter →runOutlinkRuntime生成回复,并以稳定的messageId=lastMsgId由聊天写入层保证业务幂等。这里存在两道幂等防线:
- 队列 jobId 去重:
wechat-reply:${shareId}:${lastMsgId}是确定 jobId,同一消息重复入队时 BullMQ 自动去重; - 业务 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中的实现要点:
- 步骤 3:
ret !== 0 || errcode !== 0视为 API 错误,进入失败计数流程(详见第 7 节); - 步骤 5:
Promise.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); contextToken和lastMsgId取该用户最后一条消息的值;- 同一用户在同一 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.normalizeMessage的resolveQuery负责把媒体资源下载、解密、上传 S3 后合成 query)。设计文档同时注明:微信当前不提供获取被引用消息的接口(截至 2026.7.31),引用解析暂以 @todo 形式保留在 adapter.ts。
7. Redis 与持久化合同:固定窗口失败计数
| 数据 | Physical key / 存储 | TTL/语义 |
|---|---|---|
| QR Login | fastgpt:cache:publish:wechat:qrcode:${outLinkId}:${tmbId} | QR JSON,480 秒 |
| Poll failure | fastgpt:cache:wechat:publish:failures:${shareId} | integer,300 秒 |
| 渠道配置 | MongoOutLink.app | token、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),仅供参考