返回

OpenClaw 源码解析(二):消息的一生

追踪一条消息从发送到收到回复的完整旅程,深入理解 Channel 插件、消息标准化、Session 路由、Agent 处理等核心流程。

适用版本: 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 格式 发送者标识 媒体处理
WhatsApp [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

预防原则:

同一会话的消息必须顺序处理,使用队列或序列号机制保证顺序。


设计启示:消息处理的优化空间

当前设计的优势

  1. 模块化——每个站点职责清晰,易于独立测试和修改
  2. 可扩展——新增渠道只需实现 Channel Plugin 接口
  3. 灵活性——会话路由策略可配置

潜在问题

  1. 延迟累积——每个站点都会增加延迟
  2. 错误传播——一个站点失败可能影响整条链路
  3. 监控困难——需要跨多个模块追踪问题

优化方向

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 的躯壳与灵魂