适用版本: OpenClaw v2026.3 阅读时间: 约 25 分钟 前置知识: 阅读第一篇:Gateway——系统的心脏 源码位置:
src/channels/,src/sessions/,src/gateway/server-chat.ts
开篇:追踪一条消息
上一篇文章,我们了解了 Gateway 作为控制平面如何协调各子系统。今天,让我们追踪一条真实的消息:
场景: 你在 WhatsApp 发送"今天北京天气怎么样?" 结果: OpenClaw 回复"北京今天晴,气温 15-25°C"
这条消息经历了什么?让我们画出它的旅程地图:
flowchart LR
A[📱 WhatsApp<br/>发送消息] --> B[🔌 Channel Plugin<br/>接收原始消息]
B --> C[📝 Message Normalizer<br/>标准化格式]
C --> D[🔐 Auth & Pairing<br/>权限验证]
D --> E[🗺️ Session Router<br/>路由到会话]
E --> F[🤖 Agent<br/>AI 处理]
F --> G[📤 Response Delivery<br/>回复下发]
G --> H[📱 WhatsApp<br/>收到回复]
style A fill:#e1f5fe
style H fill:#e8f5e9
本文将回答核心问题:
一条消息从发送到收到回复,经历了哪些站点?每个站点解决了什么问题?
第一站:Channel Plugin——消息的入口
问题:如何抹平平台差异?
WhatsApp、Telegram、Slack、Discord… 每个平台的消息格式都不同:
| 平道 | 消息 ID 格式 | 发送者标识 | 媒体处理 |
|---|---|---|---|
[email protected] |
[email protected] |
需下载后上传 | |
| Telegram | 数字 ID | user_id |
直接 URL |
| Slack | ts 时间戳 |
UXXXXXX |
URL 直接访问 |
| Discord | UUID | user_id |
CDN URL |
解决方案:插件化适配
flowchart TB
subgraph External["外部消息平台"]
WA[WhatsApp]
TG[Telegram]
SL[Slack]
DC[Discord]
end
subgraph Plugins["Channel Plugins"]
WP[whatsapp.ts<br/>@whiskeysockets/baileys]
TP[telegram.ts<br/>grammy]
SP[slack.ts<br/>@slack/bolt]
DP[discord.ts<br/>discord.js]
end
subgraph Normalizer["Message Normalizer"]
MN[标准化消息格式]
end
WA --> WP
TG --> TP
SL --> SP
DC --> DP
WP --> MN
TP --> MN
SP --> MN
DP --> MN
Channel Plugin 接口
// src/channels/types.ts
interface ChannelPlugin {
// 插件标识
id: ChannelId;
// 生命周期
start(context: ChannelContext): Promise<void>;
stop(): Promise<void>;
// 消息发送
sendMessage(params: SendMessageParams): Promise<void>;
// 消息接收回调
onMessage(callback: MessageCallback): void;
// 媒体处理
downloadMedia?(messageId: string): Promise<Buffer>;
uploadMedia?(data: Buffer, filename: string): Promise<string>;
}
// 标准化消息格式
interface NormalizedMessage {
id: string; // 消息唯一 ID
channelId: ChannelId; // 来源渠道
accountId: string; // 账号 ID
senderId: string; // 发送者 ID
senderName?: string; // 发送者名称
// 消息内容
text?: string;
attachments?: Attachment[];
// 元数据
timestamp: number;
replyTo?: string; // 回复的消息 ID
isGroup?: boolean;
groupId?: string;
}
WhatsApp 插件实现示例
// src/whatsapp/channel.ts
export class WhatsAppChannel implements ChannelPlugin {
id = "whatsapp" as ChannelId;
private socket: WASocket;
private messageCallbacks: MessageCallback[] = [];
async start(context: ChannelContext): Promise<void> {
// 1. 连接 WhatsApp Web
const { state, saveCreds } = await useMultiFileAuthState(
context.authPath
);
this.socket = makeWASocket({
auth: state,
printQRInTerminal: true,
});
// 2. 监听消息
this.socket.ev.on("messages.upsert", async (update) => {
for (const msg of update.messages) {
if (msg.key.fromMe) continue; // 跳过自己发的消息
// 3. 标准化消息
const normalized = this.normalizeMessage(msg);
// 4. 触发回调
for (const cb of this.messageCallbacks) {
await cb(normalized);
}
}
});
// 5. 保存认证凭据
this.socket.ev.on("creds.update", saveCreds);
}
private normalizeMessage(msg: WAMessage): NormalizedMessage {
return {
id: msg.key.id!,
channelId: "whatsapp",
accountId: msg.key.remoteJid!.split("@")[0],
senderId: msg.key.participant || msg.key.remoteJid!,
senderName: msg.pushName,
text: msg.message?.conversation ||
msg.message?.extendedTextMessage?.text,
timestamp: msg.messageTimestamp as number * 1000,
isGroup: msg.key.remoteJid!.endsWith("@g.us"),
groupId: msg.key.remoteJid!.endsWith("@g.us")
? msg.key.remoteJid!
: undefined,
};
}
async sendMessage(params: SendMessageParams): Promise<void> {
await this.socket.sendMessage(params.to, {
text: params.text,
});
}
onMessage(callback: MessageCallback): void {
this.messageCallbacks.push(callback);
}
async stop(): Promise<void> {
this.socket?.end();
}
}
设计启示: 插件模式将平台差异隔离在适配层,核心逻辑不需要关心消息来自哪个平台。
第二站:Message Normalizer——统一消息格式
问题:如何让后续处理不关心平台差异?
Channel Plugin 接收到的消息格式各异,但 Agent 只需要理解一种格式。
解决方案:标准化为 NormalizedMessage
flowchart TB
subgraph Raw["原始消息(平台特定)"]
WA_MSG["WhatsApp: {key: {id, remoteJid}, message: {conversation}, ...}"]
TG_MSG["Telegram: {message_id, from: {id}, text, ...}"]
SL_MSG["Slack: {ts, user, text, ...}"]
end
subgraph Normalized["标准化消息"]
NORM["NormalizedMessage {<br/>id: string,<br/>channelId: ChannelId,<br/>senderId: string,<br/>text: string,<br/>timestamp: number,<br/>...<br/>}"]
end
WA_MSG --> |whatsapp.ts| NORM
TG_MSG --> |telegram.ts| NORM
SL_MSG --> |slack.ts| NORM
NORM --> |后续处理| HANDLER["统一的处理逻辑"]
消息增强:添加元数据
// src/gateway/server-chat.ts
async function enrichMessage(
msg: NormalizedMessage,
config: OpenClawConfig
): Promise<EnrichedMessage> {
return {
...msg,
// 解析发送者信息
sender: {
id: msg.senderId,
name: msg.senderName,
isPaired: await checkPairing(msg.senderId, config),
isInAllowList: checkAllowList(msg.senderId, config),
},
// 解析群组信息(如果是群聊)
group: msg.isGroup ? {
id: msg.groupId!,
name: await getGroupName(msg.groupId!),
} : undefined,
// 添加处理元数据
receivedAt: Date.now(),
traceId: generateTraceId(),
};
}
第三站:Auth & Pairing——权限验证
问题:谁可以和我对话?
OpenClaw 支持 DM(私信)配对机制,防止陌生人随意发送消息。
DM 配对流程
sequenceDiagram
autonumber
participant User as 用户
participant Channel as Channel Plugin
participant Gateway as Gateway
participant Config as Config
User->>Channel: 发送消息
Channel->>Gateway: NormalizedMessage
Note over Gateway: 检查 DM 配对状态
Gateway->>Config: getPairingStatus(senderId)
alt 未配对
Config-->>Gateway: NOT_PAIRED
alt 配对模式 = OPEN
Gateway->>Gateway: 自动配对
Gateway->>Config: savePairing(senderId)
Config-->>Gateway: PAIRED
else 配对模式 = PAIRING
Gateway-->>Channel: 发送配对提示
Channel-->>User: "请先发送 !pair 命令配对"
else 配对模式 = CLOSED
Gateway-->>Channel: 拒绝消息
end
else 已配对
Config-->>Gateway: PAIRED
end
Gateway->>Gateway: 继续处理消息
配对模式配置
# ~/.openclaw/openclaw.yaml
channels:
whatsapp:
enabled: true
dmPolicy: pairing # open | pairing | closed
pairingCode: "123456" # 可选的配对码
telegram:
enabled: true
dmPolicy: open # 允许所有人
slack:
enabled: true
dmPolicy: closed # 只响应配置的用户
allowList:
- "U123456"
- "U789012"
群组权限检查
// src/channels/permissions.ts
async function checkGroupPermission(
msg: EnrichedMessage,
config: OpenClawConfig
): Promise<PermissionResult> {
if (!msg.isGroup) {
return { allowed: true }; // 私聊由 DM 配对处理
}
const groupConfig = config.channels?.[msg.channelId]?.groups?.[msg.groupId!];
if (!groupConfig) {
// 群组未配置,使用默认策略
const defaultPolicy = config.channels?.[msg.channelId]?.groupPolicy || "ignore";
return {
allowed: defaultPolicy === "accept",
reason: defaultPolicy === "ignore" ? "Group not configured" : undefined,
};
}
// 检查是否需要在群组中被 @
if (groupConfig.requireMention) {
const botId = await getBotId(msg.channelId);
const isMentioned = msg.text?.includes(`@${botId}`);
if (!isMentioned) {
return { allowed: false, reason: "Bot not mentioned" };
}
}
return { allowed: true };
}
第四站:Session Router——消息路由
问题:消息该发给哪个会话?
一个用户可能同时有多个会话:主会话、项目专用会话、特定 Agent 会话。如何决定?
Session 类型
// src/sessions/types.ts
interface SessionEntry {
key: string; // 会话唯一标识
agentId?: string; // 绑定的 Agent
channelId?: string; // 来源渠道
senderId?: string; // 发送者
groupId?: string; // 群组 ID(如果是群聊)
// 会话数据
transcript?: TranscriptEntry[];
metadata?: Record<string, unknown>;
// 生命周期
createdAt: number;
updatedAt: number;
lastMessageAt?: number;
}
会话路由规则
flowchart TB
MSG["收到消息"] --> TYPE{消息类型?}
TYPE --> |群聊消息| GROUP["群组会话<br/>key: {channel}:{groupId}"]
TYPE --> |私信消息| PRIVATE
PRIVATE --> AGENT_CHECK{是否有 Agent 别名?}
AGENT_CHECK --> |是| AGENT_SESSION["Agent 会话<br/>key: {channel}:{sender}:{agentAlias}"]
AGENT_CHECK --> |否| CHANNEL_CHECK{渠道配置?}
CHANNEL_CHECK --> |per-sender| SENDER["发送者会话<br/>key: {channel}:{sender}"]
CHANNEL_CHECK --> |global| GLOBAL["全局会话<br/>key: main"]
GROUP --> FINAL["Session Router<br/>返回 sessionKey"]
AGENT_SESSION --> FINAL
SENDER --> FINAL
GLOBAL --> FINAL
FINAL --> LOAD["加载会话上下文"]
路由实现
// src/sessions/router.ts
export function resolveSessionKey(
msg: EnrichedMessage,
config: OpenClawConfig
): string {
const channelId = msg.channelId;
// 1. 群组消息:按群组隔离
if (msg.isGroup) {
return `${channelId}:${msg.groupId}`;
}
// 2. 检查消息是否指定了 Agent
const agentAlias = extractAgentAlias(msg.text);
if (agentAlias) {
return `${channelId}:${msg.senderId}:${agentAlias}`;
}
// 3. 根据渠道配置决定会话策略
const sessionScope = config.channels?.[channelId]?.sessionScope ||
config.session?.scope ||
"per-sender";
if (sessionScope === "per-sender") {
return `${channelId}:${msg.senderId}`;
} else {
return "main";
}
}
// 从消息中提取 Agent 别名
// 例如: "@coder 帮我写个函数" -> "coder"
function extractAgentAlias(text?: string): string | undefined {
if (!text) return undefined;
const match = text.match(/^@(\w+)\s/);
return match?.[1];
}
会话隔离示例
# 不同场景的会话隔离
# 场景1: 用户在 WhatsApp 私聊
# sessionKey = "whatsapp:+1234567890"
# 该用户的所有 WhatsApp 私聊共享一个会话
# 场景2: 用户在 WhatsApp 群组
# sessionKey = "whatsapp:[email protected]"
# 该群组的所有成员共享一个会话
# 场景3: 用户指定 Agent
# 消息: "@coder 帮我写代码"
# sessionKey = "whatsapp:+1234567890:coder"
# 该会话绑定到 coder Agent,独立于主会话
第五站:Agent Processing——AI 处理
问题:如何让 AI 理解并回复消息?
Session Router 确定了会话后,消息进入 Agent 处理流程。
Turn 执行流程
sequenceDiagram
autonumber
participant Router as Session Router
participant Agent as Agent Runtime
participant LLM as LLM Provider
participant Tools as Tools
Router->>Agent: runTurn(sessionKey, message)
Note over Agent: 1. 构建上下文
Agent->>Agent: loadTranscript(sessionKey)
Agent->>Agent: buildPrompt(message, transcript)
Note over Agent: 2. 调用 LLM
Agent->>LLM: streamCompletion(prompt)
loop 流式响应
LLM-->>Agent: text_delta
Agent-->>Router: broadcast("text_delta")
end
Note over Agent: 3. 工具调用(如果需要)
LLM-->>Agent: tool_call
Agent->>Tools: executeTool(name, params)
Tools-->>Agent: result
Agent->>LLM: continue(tool_result)
Note over Agent: 4. 完成
LLM-->>Agent: done
Agent-->>Router: broadcast("done")
Agent-->>Router: return
Note over Agent: 5. 保存会话
Agent->>Agent: saveTranscript(sessionKey)
流式响应处理
// src/gateway/server-chat.ts
export async function* handleChatMessage(params: {
sessionKey: string;
message: string;
broadcast: BroadcastFn;
}): AsyncIterable<AcpRuntimeEvent> {
const { sessionKey, message, broadcast } = params;
// 1. 获取 Agent 运行时
const runtime = getAgentRuntime(sessionKey);
const handle = await runtime.ensureSession({ sessionKey });
// 2. 执行 Turn
const events = runtime.runTurn({
handle,
text: message,
mode: "prompt",
requestId: generateRequestId(),
});
// 3. 流式处理事件
for await (const event of events) {
// 广播给客户端
broadcast("agent", {
sessionId: sessionKey,
type: event.type,
...event,
});
yield event;
// 如果是完成事件,保存会话
if (event.type === "done") {
await saveSessionTranscript(sessionKey);
}
}
}
工具调用示例
sequenceDiagram
autonumber
participant User as 用户
participant Agent as Agent
participant Weather as weather.get
User->>Agent: "今天北京天气怎么样?"
Note over Agent: AI 决定调用天气工具
Agent-->>User: event: status "正在查询天气..."
Agent->>Weather: tool_call: weather.get({city: "北京"})
Note over Weather: 调用天气 API
Weather-->>Agent: {temp: "15-25°C", condition: "晴"}
Agent-->>User: event: tool_result
Agent-->>User: event: text_delta "北京今天"
Agent-->>User: event: text_delta "天气晴朗"
Agent-->>User: event: text_delta ",气温 15-25°C。"
Agent-->>User: event: done
第六站:Response Delivery——回复下发
问题:如何把回复发送回原渠道?
Agent 生成的回复需要通过原渠道发送回去。
回复下发流程
// src/gateway/server-chat.ts
async function deliverResponse(params: {
sessionKey: string;
response: string;
originalMessage: EnrichedMessage;
}): Promise<void> {
const { sessionKey, response, originalMessage } = params;
// 1. 解析会话键,获取渠道信息
const sessionMeta = parseSessionKey(sessionKey);
const channelId = sessionMeta.channelId;
// 2. 获取对应的 Channel Plugin
const channel = getChannelPlugin(channelId);
if (!channel) {
throw new Error(`Channel ${channelId} not found`);
}
// 3. 确定回复目标
const to = originalMessage.isGroup
? originalMessage.groupId
: originalMessage.senderId;
// 4. 发送回复
await channel.sendMessage({
to,
text: response,
replyTo: originalMessage.id, // 引用原消息
});
}
流式回复的魔法
对于支持流式回复的渠道(如 Web UI),可以实现"打字机效果":
// 流式回复处理
async function handleStreamingResponse(params: {
sessionKey: string;
events: AsyncIterable<AcpRuntimeEvent>;
channel: ChannelPlugin;
to: string;
}): Promise<void> {
let buffer = "";
for await (const event of params.events) {
if (event.type === "text_delta") {
buffer += event.text;
// 每 50ms 或每 10 个字符发送一次更新
if (shouldFlush(buffer)) {
await params.channel.sendTyping?.(params.to);
// 或者直接发送部分文本(如果渠道支持)
}
} else if (event.type === "done") {
// 最终发送完整消息
await params.channel.sendMessage({
to: params.to,
text: buffer,
});
}
}
}
实战演练:追踪一条消息
让我们用日志追踪一条真实的消息流程。
启用调试日志
LOG_LEVEL=debug LOG_INCLUDE="gateway,channels,sessions,agent" openclaw gateway start
发送测试消息
通过 WhatsApp 发送 “Hello”,观察日志输出:
[DEBUG] channels:whatsapp: Received message {
id: "3EB0ABC123",
from: "+1234567890",
text: "Hello",
timestamp: 1736789123456
}
[DEBUG] gateway: Normalized message {
id: "3EB0ABC123",
channelId: "whatsapp",
senderId: "+1234567890",
text: "Hello"
}
[DEBUG] gateway: Checking DM pairing for +1234567890
[DEBUG] gateway: DM status: PAIRED
[DEBUG] sessions: Resolving session key {
channelId: "whatsapp",
senderId: "+1234567890",
scope: "per-sender"
}
[DEBUG] sessions: Session key resolved: "whatsapp:+1234567890"
[DEBUG] agent: Starting turn {
sessionKey: "whatsapp:+1234567890",
message: "Hello"
}
[DEBUG] agent: Loading transcript, found 5 entries
[DEBUG] agent: Calling LLM provider: openai/gpt-4
[DEBUG] agent: text_delta "Hello"
[DEBUG] agent: text_delta "!"
[DEBUG] agent: text_delta " How"
[DEBUG] agent: text_delta " can"
[DEBUG] agent: text_delta " I"
[DEBUG] agent: text_delta " help"
[DEBUG] agent: text_delta " you"
[DEBUG] agent: text_delta "?"
[DEBUG] agent: done {stopReason: "end_turn"}
[DEBUG] sessions: Saving transcript for whatsapp:+1234567890
[DEBUG] channels:whatsapp: Sending message to +1234567890
[DEBUG] channels:whatsapp: Message sent, id: "3EB0DEF456"
使用 wscat 追踪
wscat -c "ws://127.0.0.1:18789" -H "Authorization: Bearer $TOKEN"
> {"kind":"request","id":"trace-1","method":"message.send","params":{"sessionKey":"main","text":"Hello"}}
< {"kind":"event","event":"agent","payload":{"sessionId":"main","seq":1,"type":"status","text":"Thinking..."}}
< {"kind":"event","event":"agent","payload":{"sessionId":"main","seq":2,"type":"text_delta","text":"Hello"}}
< {"kind":"event","event":"agent","payload":{"sessionId":"main","seq":3,"type":"text_delta","text":"!"}}
< {"kind":"event","event":"agent","payload":{"sessionId":"main","seq":4,"type":"text_delta","text":" How"}}
< {"kind":"event","event":"agent","payload":{"sessionId":"main","seq":5,"type":"text_delta","text":" can"}}
< {"kind":"event","event":"agent","payload":{"sessionId":"main","seq":6,"type":"text_delta","text":" I"}}
< {"kind":"event","event":"agent","payload":{"sessionId":"main","seq":7,"type":"text_delta","text":" help"}}
< {"kind":"event","event":"agent","payload":{"sessionId":"main","seq":8,"type":"text_delta","text":" you"}}
< {"kind":"event","event":"agent","payload":{"sessionId":"main","seq":9,"type":"text_delta","text":"?"}}
< {"kind":"event","event":"agent","payload":{"sessionId":"main","seq":10,"type":"done","stopReason":"end_turn"}}
< {"kind":"response","id":"trace-1","ok":true}
常见陷阱
消息处理流程看似简单,但实际开发中有很多容易踩的坑。
陷阱 1:Session 键解析错误
症状:
# 私聊消息错误地路由到群组会话
[DEBUG] sessions: Session key resolved: "whatsapp:[email protected]"
# 但消息实际来自私聊,senderId = "+1234567890"
原因: Session 键的构造逻辑错误,导致消息被路由到错误的会话。
错误代码:
// 错误:直接使用 remoteJid,没有判断是否是群组
function resolveSessionKey(msg: NormalizedMessage): string {
// remoteJid 在私聊时是用户 ID,群聊时是群组 ID
const sessionKey = msg.key.remoteJid; // ❌ 逻辑错误
return `whatsapp:${sessionKey}`;
}
正确代码:
// 正确:根据消息类型决定 sessionKey
function resolveSessionKey(msg: NormalizedMessage): string {
const channelId = msg.channelId;
// 1. 群组消息:使用 groupId
if (msg.isGroup) {
return `${channelId}:${msg.groupId}`;
}
// 2. 私聊消息:使用 senderId
return `${channelId}:${msg.senderId}`;
}
调试方法:
# 启用 Session 路由调试日志
LOG_LEVEL=debug LOG_INCLUDE="session:router" openclaw gateway start
# 查看路由决策日志
[DEBUG] session:router: Incoming message {
isGroup: false,
senderId: "+1234567890",
groupId: undefined
}
[DEBUG] session:router: Resolved sessionKey: "whatsapp:+1234567890"
源码位置: src/sessions/router.ts:434-461
陷阱 2:消息重复处理
症状:
# 同一条消息被处理了两次
[DEBUG] gateway: Processing message 3EB0ABC123 from +1234567890
[DEBUG] gateway: Processing message 3EB0ABC123 from +1234567890 # 重复!
原因: 渠道的消息回调可能被触发多次(网络重试、多设备同步等),导致同一消息被重复处理。
错误理解:
// 错误:假设消息只会到达一次
socket.ev.on("messages.upsert", async (update) => {
for (const msg of update.messages) {
// 直接处理,没有去重
await processMessage(msg); // ❌ 可能重复处理
}
});
正确代码:
// 正确:使用消息 ID 去重
const processedMessages = new Set<string>();
const MESSAGE_TTL = 5 * 60 * 1000; // 5 分钟过期
async function handleWithDedupe(msg: NormalizedMessage): Promise<void> {
const key = `${msg.channelId}:${msg.id}`;
// 检查是否已处理
if (processedMessages.has(key)) {
log.debug(`Message ${msg.id} already processed, skipping`);
return;
}
// 标记为已处理
processedMessages.add(key);
// 设置过期清理
setTimeout(() => processedMessages.delete(key), MESSAGE_TTL);
// 处理消息
await processMessage(msg);
}
调试方法:
# 启用去重日志
LOG_LEVEL=debug LOG_INCLUDE="gateway:dedupe" openclaw gateway start
# 观察去重效果
[DEBUG] gateway:dedupe: New message 3EB0ABC123, adding to cache
[DEBUG] gateway:dedupe: Duplicate message 3EB0ABC123, skipping
源码位置: src/gateway/server-runtime-state.ts (DEDUPE Map)
陷阱 3:群组 @ 提及检测失败
症状:
# 群组消息没有 @ 机器人,但机器人却回复了
# 或:群组消息 @ 了机器人,但机器人没有回复
原因: 群组中的 @ 提及检测逻辑不正确,可能是:
- 提及格式与平台不符
- 没有考虑不同的提及语法
错误代码:
// 错误:只检查文本中是否包含 @bot
const isMentioned = msg.text?.includes("@bot"); // ❌ 不准确
正确代码:
// 正确:检查实际的提及实体
async function checkMention(msg: EnrichedMessage, channelId: string): Promise<boolean> {
// WhatsApp: 检查 message.mentionedJid
if (channelId === "whatsapp") {
const mentionedJids = msg.message?.extendedTextMessage?.contextInfo?.mentionedJid || [];
const botId = await getBotId(channelId);
return mentionedJids.includes(botId);
}
// Telegram: 检查 entities
if (channelId === "telegram") {
const entities = msg.entities || [];
return entities.some(e => e.type === "mention" && e.user?.id === botId);
}
// Slack: 检查 <@U123456> 格式
if (channelId === "slack") {
const botUserId = await getBotUserId(channelId);
return msg.text?.includes(`<@${botUserId}>`);
}
return false;
}
调试方法:
# 启用群组消息调试
LOG_LEVEL=debug LOG_INCLUDE="channels:group" openclaw gateway start
# 查看提及检测日志
[DEBUG] channels:group: Checking mention for message in group [email protected]
[DEBUG] channels:group: Mentioned JIDs: ["[email protected]"]
[DEBUG] channels:group: Bot ID: [email protected]
[DEBUG] channels:group: Is mentioned: false
源码位置: src/channels/permissions.ts:344-376
陷阱 4:消息顺序错乱
症状:
# 用户发送消息 A,然后发送消息 B
# 但 Agent 先回复了 B,再回复 A
原因: 异步处理导致消息顺序无法保证,特别是当:
- 消息 A 触发了工具调用(耗时较长)
- 消息 B 只需要简单回复
错误代码:
// 错误:每个消息独立处理,没有顺序保证
socket.ev.on("messages.upsert", async (update) => {
for (const msg of update.messages) {
// 并行处理,顺序无法保证
processMessage(msg); // ❌ 没有 await
}
});
正确代码:
// 正确:使用队列保证顺序
class MessageQueue {
private queue: NormalizedMessage[] = [];
private processing = false;
async enqueue(msg: NormalizedMessage): Promise<void> {
this.queue.push(msg);
if (!this.processing) {
this.processQueue();
}
}
private async processQueue(): Promise<void> {
this.processing = true;
while (this.queue.length > 0) {
const msg = this.queue.shift()!;
await processMessage(msg); // ✅ 顺序处理
}
this.processing = false;
}
}
// 使用队列
const messageQueue = new MessageQueue();
socket.ev.on("messages.upsert", async (update) => {
for (const msg of update.messages) {
const normalized = normalizeMessage(msg);
await messageQueue.enqueue(normalized);
}
});
调试方法:
# 启用消息队列日志
LOG_LEVEL=debug LOG_INCLUDE="gateway:queue" openclaw gateway start
# 观察消息处理顺序
[DEBUG] gateway:queue: Message 3EB0ABC123 enqueued, queue length: 1
[DEBUG] gateway:queue: Message 3EB0DEF456 enqueued, queue length: 2
[DEBUG] gateway:queue: Processing message 3EB0ABC123
[DEBUG] gateway:queue: Processing message 3EB0DEF456
预防原则:
同一会话的消息必须顺序处理,使用队列或序列号机制保证顺序。
设计启示:消息处理的优化空间
当前设计的优势
- 模块化——每个站点职责清晰,易于独立测试和修改
- 可扩展——新增渠道只需实现 Channel Plugin 接口
- 灵活性——会话路由策略可配置
潜在问题
- 延迟累积——每个站点都会增加延迟
- 错误传播——一个站点失败可能影响整条链路
- 监控困难——需要跨多个模块追踪问题
优化方向
flowchart TB
subgraph Current["当前架构"]
C1["Channel"] --> C2["Normalizer"] --> C3["Auth"] --> C4["Router"] --> C5["Agent"] --> C6["Delivery"]
end
subgraph Optimized["优化架构"]
O1["Channel"] --> O2["Pipeline<br/>(并行处理)"]
O2 --> O3["Agent<br/>(异步执行)"]
O3 --> O4["Delivery<br/>(批量发送)"]
end
Current --> |优化| Optimized
| 优化点 | 方案 | 预期收益 |
|---|---|---|
| 并行处理 | Auth 和 Session 解析并行 | 减少 50-100ms |
| 异步日志 | 日志写入异步化 | 减少阻塞 |
| 连接池 | 复用 Channel 连接 | 减少连接开销 |
| 缓存 | Session 元数据缓存 | 减少磁盘 I/O |
小结
一条消息的生命周期:
| 站点 | 职责 | 关键问题 |
|---|---|---|
| Channel Plugin | 消息入口 | 如何抹平平台差异? |
| Message Normalizer | 格式统一 | 如何让后续处理不关心平台? |
| Auth & Pairing | 权限验证 | 谁可以和我对话? |
| Session Router | 消息路由 | 消息该发给哪个会话? |
| Agent Processing | AI 处理 | 如何理解并回复? |
| Response Delivery | 回复下发 | 如何发送回原渠道? |
理解这个流程的关键是每个站点只解决一个特定问题,通过管道式处理完成完整的消息生命周期。
关键源码文件
| 文件 | 职责 |
|---|---|
src/channels/types.ts |
Channel Plugin 接口定义 |
src/whatsapp/channel.ts |
WhatsApp 插件实现 |
src/gateway/server-chat.ts |
消息处理主流程 |
src/sessions/router.ts |
会话路由逻辑 |
src/sessions/store.ts |
会话存储 |
在下一篇文章中,我们将深入 Agent 的内部世界——ACP 协议、Turn 执行、工具调用,揭示 AI 如何理解和生成回复。
系列索引: OpenClaw 源码解析:目录索引
上一篇: OpenClaw 源码解析(一):Gateway——系统的心脏 下一篇: OpenClaw 源码解析(三):Agent——AI 的躯壳与灵魂