Memoh 群聊/社交网络实现实录
从玩具 Demo 到生产级群聊机器人的完整工程路径
写在前面
Memoh 是一个多成员、结构化长记忆的容器化 AI agent 平台。用户可以通过 Telegram、Discord、飞书、Matrix 等渠道创建 Bot 并与之对话。在传统的 Chat 模式下,用户通过 @ 机器人 + 输入文字触发对话,Bot 完成指令后返回结果——这在一对一场景很自然,但在群聊中显然不够"群友"。
我们希望 Bot 能像真正的群成员一样:听得见所有对话、自己判断要不要开口、开口就开口,闭嘴就闭嘴。 这就是 Discuss 模式的起源。
本文记录了我们从头把 Discuss 模式从概念做到生产级系统的全过程——包括架构设计的演进、与上游的对比,以及踩过的坑。
一、上游给了什么:一个概念验证级 Demo
Discuss 模式的概念来自上游 memohai/Memoh 和其研究性项目 Cahciua(Menci 的开源贡献)。上游提供了两个核心基础:
1.1 DCP 确定性的上下文管线
DCP(Deterministic Context Pipeline)是 Discuss 模式的上下文基石,三个纯函数层:
- Adaptation(适配层):将不同平台(Telegram/Discord/飞书)的原始消息标准化为
CanonicalEvent - Projection(投影层):
IC' = Reduce(IC, Event),事件驱动的状态机,消息追加/修改/删除/撤回全部映射 - Rendering(渲染层):
RC = Render(IC, RenderParams),将中间上下文序列化为 LLM 可读的 XML
这条管线是纯函数的,输入确定则输出确定,可测试、可重放、可审计。
但管线只是"看得见群聊",看到之后要不要说话、怎么说话、什么时候说——这些全部不在管线里。
1.2 上游 Driver:一个 449 行的演示级实现
上游在 internal/pipeline/driver.go 里提供了一个 DiscussDriver,结构如下:
// 上游的 handleReply —— 就一行
func (d *DiscussDriver) handleReply(ctx context.Context, sess *discussSession,
rc RenderedContext, log *slog.Logger) {
d.handleReplyWithAgent(ctx, sess, rc, log, d.deps.Agent) // 无条件直接调 LLM
}
功能清单:
- ✅ 有 session goroutine(基础生命周期管理)
- ✅ 有 RC drain(批量合并更新)
- ✅ 有 idle timeout(10 分钟自动退出)
- ✅ 有 cold-start anchor(从 TR 锚定游标)
- ✅ 有 Late Binding Prompt
- ❌ 没有智能计时(每条消息无条件触发 LLM)
- ❌ 没有渠道消息投递(只有 WebUI 广播,消息根本到不了 Telegram)
- ❌ 没有 dispatcher 队列(LLM 运行期间新消息会丢)
- ❌ 没有中断重试
- ❌ 没有被动记忆提取
- ❌ 没有 replyer 改写
- ❌ 没有表达式/行话学习
更关键的是架构问题:上游 Driver 直接调用 agent.Stream(),和 Chat 模式走两条完全独立的路径。这意味着:
- 压缩(compaction)需要单独实现
- 重试逻辑需要单独实现
- 消息持久化需要单独实现
- 修一个 bug,两边各修一遍
二、我们的重构:「注入式架构」
2.1 核心洞察
Chat 模式和 Discuss 模式的差异,本质上是 触发策略 和 输出的含义 不同,而不是 LLM 调用本身有什么本质区别。
| Chat 模式 | Discuss 模式 | |
|---|---|---|
| 触发时机 | 用户发消息即触发 | 智能计时决定 |
| 输入来源 | 数据库历史 + 当前消息 | DCP 管线 (RC+TR) |
| 文本输出 | 就是回复内容 | 内心独白,不可见 |
| 是否回复 | 必须回复 | Bot 自主决定(send 工具) |
LLM 调用本身完全一样:组装上下文 → 调 StreamChat → 消费事件流 → 持久化结果。
2.2 注入式架构设计
我们的核心设计决策:不写第二条 LLM 调用路径。改为通过 ChatRequest 的三个注入字段,让 Discuss 模式复用 Chat 模式的完整 Resolver.StreamChat() 管线:
┌───────────────────────┐
│ DiscussTrigger │
│ (讨论触发策略层) │
│ │
│ • 智能计时决策 │
│ • 会话生命周期 │
│ • 调度器/队列管理 │
└───────────┬───────────┘
│ ChatRequest { SessionType="discuss",
│ InjectCh=...,
│ DiscussLateBindingPrompt=...
│ }
▼
┌───────────────────────┐
│ Resolver.StreamChat()│
│ (共享的 LLM 执行管线) │
│ │
│ • 模型选择/凭证解析 │
│ • 上下文组装 (RC+TR) │
│ • 压缩/重试/持久化 │
│ • Token 预算管理 │
│ • 工具调用/流式响应 │
└───────────┬───────────┘
│
▼
┌───────────────────────┐
│ agent.Stream() │
│ (LLM 模型调用) │
└───────────────────────┘
2.3 三个注入点
| 注入点 | 类型 | 作用 |
|---|---|---|
SessionType = "discuss" |
字符串 | resolve() 中跳过 query 校验、走 Pipeline 消息源 |
DiscussLateBindingPrompt |
字符串 | 作为最后一条 user 消息注入,告诉 LLM "你的输出是内心独白" |
InjectCh |
channel | agent 运行期间接收外部注入消息,支持中断式交互 |
这三个字段就是我们和上游最根本的架构差异。上游需要一个独立的 DiscussDriver 做所有事情,我们只需要构造一个正确的 ChatRequest。
2.4 内心独白机制
在 consumeStream() 中,Discuss 模式的关键逻辑:
// In discuss mode, text deltas are internal monologue —
// only the send/reply tool delivers visible messages.
if event.Type == channel.StreamEventDelta && event.Phase == channel.StreamPhaseText {
continue // 跳过文本阶段,不投递到渠道
}
finalizeOutboundStream() 同理——纯文本助手输出不投递,只有 send/reply 工具调用才产生可见消息。
三、智能计时子系统(internal/chattiming/)
这是上游完全没有的东西,也是"群聊机器人不被踢出群"的核心。
3.1 Debounce(消息去抖)
群聊中经常出现连续多条消息同时到达的场景。不能每条都触发一次 LLM。Debounce 采用两阶段策略:
- 安静期:收到消息后等 2s,期间无新消息才开始处理
- 最大等待:15s 强制超时,防止无限等待
if sess.debounce != nil {
sess.debounce.Reset()
if err := sess.debounce.Wait(ctx); err != nil {
continue
}
}
// Drain any additional RCs that arrived during the debounce window
3.2 TalkValue(多言值)
Bot 有一个可配置的"多言值"(0.0–1.0),控制它插话的频率。通过 JSONB 存储在 bot 设置中,支持按时间分时调节:
type TalkValueConfig struct {
Value float64 `json:"value,omitempty"` // 基础值
Rules []TimeBasedValueRule `json:"rules,omitempty"` // 按时间调节
}
func (c *TalkValueConfig) TriggerThreshold(now time.Time) int {
// 多言值越高,需要的消息数越少,Bot 越容易说话
}
3.3 IdleCompensation(空闲补偿)
长时间没人说话时,积累"信用值",Bot 更容易插话;群聊热火朝天时,信用值不足则不打扰:
if newMsgCount < threshold && sess.idleCompensate != nil {
idleDuration := time.Since(time.UnixMilli(lastMsgMs))
credit := sess.idleCompensate.ComputeCredit(idleDuration,
chattiming.ComputeCreditRateFromIntervals(sess.msgIntervals))
newMsgCount += credit
}
3.4 TimingGate(探针门)
在正式调用大模型之前,先跑一次轻量级 LLM 判断,返回三种结果:
continue— 继续执行完整 LLMwait— 等 N 秒后重新评估no_reply— 本轮不说话,做被动记忆提取
if sess.timingGate != nil && sess.chatTimingCfg.TimingGate && !isMentioned {
result := sess.timingGate.Evaluate(ctx, params, probeCfg)
if result.Decision == chattiming.TimingNoReply {
d.extractPassiveMemory(ctx, sess, rc, log)
return // 跳过本轮 LLM
}
}
目前 TimingGate 和主 LLM 共用同一个模型。discuss_probe_model_id 字段已在数据库中预留(migration 0057),计划后续独立为一个更小的模型以节省成本。
3.5 Interrupt(中断重试)
这是处理"Bot 正在回复时新消息又来了"的关键机制。上游完全不存在,我们的实现包含 4 层防护:
- 中断请求:新消息到达时调用
interrupt.RequestInterrupt()向运行中的 agent 发信号 - 去抖再等:中断后等一个安静期再重试
- 轮次限制:最多 7 轮(
maxRounds = 7),防止死循环 - 连续中断限制:最多连续 2 次中断,防止无限重启
for round := 0; round < maxRounds; round++ {
agentCtx, agentCancel = sess.interrupt.Bind(ctx)
chunkCh, errCh := d.deps.ChatRunner.StreamChat(agentCtx, chatReq)
// ... consume stream ...
if wasInterrupted && sess.interrupt.CanRetry() {
if hadOutput { break } // 已经产生输出了就跳过重试
_ = sess.debounce.Wait(ctx)
continue // 重试
}
break
}
四、被动记忆与表达式学习
4.1 被动记忆提取
Bot 决定不说话的时候,不是真的在发呆。extractPassiveMemory() 收集本轮新消息,异步提取记忆:
func (d *DiscussTrigger) extractPassiveMemory(ctx context.Context, sess *discussSession,
rc RenderedContext, log *slog.Logger) {
// 收集本轮新消息
// → MemoryFormation.OnAfterChat() 提取关键信息
// → ExpressionAccumulator() 积累表达方式
}
4.2 Replyer 回复改写
reply 工具是 send 的增强版:LLM 只输出推理逻辑,replyer(一个更小的 LLM)将推理改写成自然的群聊语言:
LLM 推理: "用户问今晚吃什么,我可以推荐附近餐厅"
↓ (replyer LLM 改写)
群聊消息: "今晚不如去吃楼下的火锅?上周试过还不错"
expression.Selector 还会从 Expression 仓库中检索 Bot 之前学到的表达风格,注入 replyer 提示中,保持输出风格一致。
五、Dispatcher 队列
在 Bot 调用 LLM(可能长达数分钟)期间,群聊不会停下来。DiscussDispatcherAdapter 确保这期间的新消息不会丢失:
新消息到达 → RouteDispatcher.IsActive(routeID)?
├── false → 直接 NotifyRC,启动新 LLM
└── true → EnqueueNotification() 排队
↓
Bot LLM 调用完成
↓
MarkDone() → 重放队列中的通知
六、与上游的量化对比
| 维度 | 上游 DiscussDriver |
我们的 DiscussTrigger |
|---|---|---|
| Driver 代码行数 | 449 行 | 671 行 |
handleReply 策略 |
1 行(无条件 LLM 调用) | 48 行(7 层决策) |
| 智能计时 | 无 | 去抖/中断/TimingGate/空闲补偿/TalkValue |
| 渠道消息投递 | 仅 WebUI 广播 | 完整 ChannelSender + 流式投递 + TextPhase 跳过 |
| 内心独白 | 靠 prompt 约定(无代码保障) | consumeStream 硬拦截 |
| Dispatcher 队列 | 无 | MarkActive/MarkDone/队列重放 |
| 中断重试 | 无 | 4 层防护 + 7 轮上限 |
| LLM 执行路径 | 独立 agent.Stream() |
注入 Resolver.StreamChat() |
| 与 Chat 模式关系 | 并行独立路径 | 共享同一管线 |
| 被动记忆 | 无 | 异步 MemoryFormation |
| Replyer 改写 | 无 | reply 工具 + replyer LLM |
| 表达式学习 | 无 | internal/expression/ 完整子系统 |
| 单元测试 | 299 行(基础测试) | 544 行 |
| Timeline 支持 | 无 | 全链路 IsTimeline |
| 硬超时保护 | 无 | 15 分钟 Deadline |
| 冷却时间 | 无 | 15 秒最小间隔 |
七、踩过的坑
7.1 冷启动游标锚定
上游的 anchorFromTRs() 是一个巧妙的设计:服务器重启或 10 分钟空闲超时后重连时,从 TR 流推测上次处理到哪里。我们因为架构不同(统一走 Resolver.StreamChat),消息去重由 Resolver 的 dedupePersistedCurrentUserMessage() 处理,所以没有移植这个功能。
7.2 send 调用频次控制
早期版本让 LLM 在每轮工具调用时都附带一个 send,结果是群聊被刷屏。后来在 system_discuss.md 中明确约束 at most ONCE per turn,并加了异常规则:
- 图片生成请求 → 直接调
generate_image,不通过send - 多步任务 → 第一轮可以
send("在查了"),但不能每步都说
7.3 中断重试的边界条件
中断机制最棘手的情况是:LLM 已经调用了 send 工具(消息已通过 SendDirect 投递到渠道),但后续又收到了中断信号。我们的处理策略是:如果已经有可见输出(hadOutput == true),跳过重试——重复发消息比不说话更糟糕。
八、未来方向
- Probe Model 独立化:目前 TimingGate 和主 LLM 共用模型,
discuss_probe_model_id字段已预留,接线上一个小模型可以大幅降低探针阶段成本 - 跨会话感知:让 Bot 感知在不同群/不同频道中的状态,实现跨会话上下文关联
- Timeline 社区学习:公开时间线上的社区风格学习已通过
system_discuss.md中的 "Timeline intelligence" 段落指导 LLM 行为,后续可接expression.Learner自动化
尾声
做 Discuss 模式的本质问题是:怎么让一个调用外部 API 的软件,"像人一样"在群聊里说话?
答案是分两层:策略层决定什么时候开口(智能计时 + 中断重试),执行层复用了 Chat 模式的完整 LLM 管线。通过三个注入字段(SessionType / InjectCh / DiscussLateBindingPrompt),我们把设计和执行拆开了——每层只做一件事,加起来就是 bot 的"人格"。
上游给了一个很好的起点(DCP 管线),我们在上面搭出了真正能跑在群里的,bot。
鸣谢 Menci 和 Cahciua 项目提出的 DCP 架构概念,以及上游 memohai/Memoh 提供的基础 Pipeline 实现。本文所述 Discuss Driver、智能计时、Dispatcher 队列、Replyer 改写及表达式学习系统为本文作者独立完成的工程实现。