十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

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渠道下线、登出或连续错误后能够停止续链不再无限轮询。在职责边界上DALData 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链路要点用户在前端发布页触发扫码登录登录二维码与登录态通过WechatQrLoginCacheQR JSON480 秒 TTL暂存登录成功后渠道配置写入MongoOutLink.apptoken、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.comaccountId、userId登录身份syncBuf下一次getUpdates的消费游标默认空字符串statusonline、offline、error默认offlineloginTime最近登录时间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消息摄入属性合同QueuewechatPollJob namewechatPublishPollJob data{ shareId }Job IDwechat-poll:${shareId}ConcurrencyWECHAT_CHANNEL_CONCURRENCYLock120 秒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 的具体配置在initWechatPollWorkermq.tslockDuration: 120_000防止长轮询期间 job 被误判为 stalledstalledInterval: 30_000每 30 秒检查活跃度removeOnComplete/removeOnFail均设为立即清理。4.2 Reply Queue回复生成属性合同QueuewechatReplyJob namewechatPublishReplyJob datashareId、userId、items、contextToken、lastMsgIdJob IDwechat-reply:${shareId}:${lastMsgId}ConcurrencyWECHAT_CHANNEL_CONCURRENCYLock30 分钟Stalled interval60 秒Failed retention500 条或 7 天Reply processorprocessWechatReplyJobmq.ts调用 provider adapter →runOutlinkRuntime生成回复并以稳定的messageIdlastMsgId由聊天写入层保证业务幂等。这里存在两道幂等防线队列 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中的实现要点步骤 3ret ! 0 || errcode ! 0视为 API 错误进入失败计数流程详见第 7 节步骤 5Promise.all(groups.map(...))并行投递 reply jobsjobId 为replyJobId(shareId, lastMsgId)步骤 6只有全部 reply job 入队成功后才updateOne推进app.syncBuf resp.get_updates_bufmq.ts。第五步失败时不得推进syncBuf——下一次 poll 会重新拉取同一批消息由replyJobId去重从而形成 at-least-once 摄入 幂等消费。这条顺序保证是推进游标前先入队的核心约束直接对应设计目标中的第四点。另一个设计约定是poll processor 本身不续链。续链统一由 Worker 的completed/failedlistener 负责scheduleNextPoll避免 return、throw 和 timeout 三个分支各自维护调度逻辑导致分叉。completed 事件中hadMessages为 false 时追加EMPTY_POLL_DELAY_MS延迟。6. 消息合并语义同用户同周期只回一条groupMessagesByUsermessageParser.ts在单个 poll 响应内按userId聚合文本使用换行拼接同一用户的多条文本依次 push 进同一组的itemscontextToken和lastMsgId取该用户最后一条消息的值同一用户在同一 poll 周期只生成一个 reply job 和一次合并回复跨 poll 周期生成独立的 reply job但共享chatId wechat_${shareId}_${userId}见 adapter.ts上下文连续多个用户生成多个 reply job并行处理。isSupportedMessageItemmessageParser.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 Loginfastgpt:cache:publish:wechat:qrcode:${outLinkId}:${tmbId}QR JSON480 秒Poll failurefastgpt:cache:wechat:publish:failures:${shareId}integer300 秒渠道配置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),仅供参考
返回列表