diff --git a/extensions/whatsapp/src/connection-controller-registry.test.ts b/extensions/whatsapp/src/connection-controller-registry.test.ts index 9542d4a1b1c3..58fb69d79020 100644 --- a/extensions/whatsapp/src/connection-controller-registry.test.ts +++ b/extensions/whatsapp/src/connection-controller-registry.test.ts @@ -15,6 +15,8 @@ describe("WhatsApp connection controller registry", () => { const second = await importRegistryModule(`second-${Date.now()}`); const controller = { getActiveListener: vi.fn(() => null), + getCurrentSock: vi.fn(() => null), + getSelfIdentity: vi.fn(() => null), }; first.registerWhatsAppConnectionController("work", controller); diff --git a/extensions/whatsapp/src/connection-controller-registry.ts b/extensions/whatsapp/src/connection-controller-registry.ts index 9e65de3c202c..aa3467fad843 100644 --- a/extensions/whatsapp/src/connection-controller-registry.ts +++ b/extensions/whatsapp/src/connection-controller-registry.ts @@ -1,8 +1,20 @@ // Whatsapp plugin module implements connection controller registry behavior. +import type { WASocket } from "baileys"; +import type { WhatsAppSelfIdentity } from "./identity.js"; import type { ActiveWebListener } from "./inbound/types.js"; type WhatsAppConnectionControllerHandle = { getActiveListener(): ActiveWebListener | null; + getCurrentSock(): WASocket | null; + /** + * The self identity (jid + lid) of the controller's currently-authenticated + * socket, or `null` if the socket is not connected or not authenticated yet. + * Used as the session-identity guard for outbound socket fallback so an + * in-place relink to a different phone number is not silently accepted. + * Compared via `identitiesOverlap()` so JID-vs-LID and device-scoped JID + * differences between the two controllers' user records are normalized away. + */ + getSelfIdentity(): WhatsAppSelfIdentity | null; }; type ConnectionRegistryState = { diff --git a/extensions/whatsapp/src/connection-controller.ts b/extensions/whatsapp/src/connection-controller.ts index 91beb693f35d..84cec8bd6ff1 100644 --- a/extensions/whatsapp/src/connection-controller.ts +++ b/extensions/whatsapp/src/connection-controller.ts @@ -6,6 +6,7 @@ import { registerWhatsAppConnectionController, unregisterWhatsAppConnectionController, } from "./connection-controller-registry.js"; +import { resolveComparableIdentity, type WhatsAppSelfIdentity } from "./identity.js"; import type { ActiveWebListener, WebListenerCloseReason } from "./inbound/types.js"; import { computeBackoff, sleepWithAbort, type ReconnectPolicy } from "./reconnect.js"; import { @@ -370,6 +371,30 @@ export class WhatsAppConnectionController { return this.current?.listener ?? null; } + getCurrentSock(): WASocket | null { + return this.socketRef.current; + } + + getSelfIdentity(): WhatsAppSelfIdentity | null { + const user = this.socketRef.current?.user as + | { id?: string | null; lid?: string | null } + | undefined; + if (!user) { + return null; + } + const jid = user.id ?? null; + const lid = user.lid ?? null; + if (!jid && !lid) { + return null; + } + // Pre-resolve via the controller's authDir so e164 is populated from the + // auth-state PN<->LID mapping. That lets `identitiesOverlap()` recognize a + // successor logged into the same account even when the original socket + // exposes only the PN form and the successor exposes only the LID form. + const resolved = resolveComparableIdentity({ jid, lid }, this.authDir); + return { jid: resolved.jid, lid: resolved.lid, e164: resolved.e164 }; + } + getReconnectAttempts(): number { return this.reconnectAttempts; } diff --git a/extensions/whatsapp/src/inbound/monitor.ts b/extensions/whatsapp/src/inbound/monitor.ts index d87938e89345..d4bbafa39923 100644 --- a/extensions/whatsapp/src/inbound/monitor.ts +++ b/extensions/whatsapp/src/inbound/monitor.ts @@ -21,7 +21,8 @@ import { createSubsystemLogger } from "openclaw/plugin-sdk/runtime-env"; import { uniqueStrings } from "openclaw/plugin-sdk/string-coerce-runtime"; import { maybeResolveWhatsAppApprovalReaction } from "../approval-reactions.js"; import { readWebSelfIdentityForDecision, WhatsAppAuthUnstableError } from "../auth-store.js"; -import { getPrimaryIdentityId, resolveComparableIdentity } from "../identity.js"; +import { getRegisteredWhatsAppConnectionController } from "../connection-controller-registry.js"; +import { getPrimaryIdentityId, identitiesOverlap, resolveComparableIdentity } from "../identity.js"; import { addWhatsAppImagePreviewFields } from "../image-preview.js"; import { cacheInboundMessageMeta } from "../quoted-message.js"; import { DEFAULT_RECONNECT_POLICY, computeBackoff, sleepWithAbort } from "../reconnect.js"; @@ -222,7 +223,6 @@ export async function attachWebInboxToSocket( if (options.socketRef) { options.socketRef.current = sock; } - const getCurrentSock = () => (options.socketRef ? options.socketRef.current : sock); const shouldRetryDisconnect = () => options.shouldRetryDisconnect?.() === true; const disconnectRetryPolicy = options.disconnectRetryPolicy ?? DEFAULT_RECONNECT_POLICY; const sendRetryMaxAttempts = @@ -264,6 +264,29 @@ export async function attachWebInboxToSocket( ); } const self = selfIdentity.identity; + // If this monitor's controller is shutdown while a captured reply is still in + // flight, only hand off to a successor controller authenticated as the same + // WhatsApp identity. Missing or mismatched identity fails closed. + const getCurrentSock = (): WASocket | null => { + if (!options.socketRef) { + return sock; + } + if (options.socketRef.current) { + return options.socketRef.current; + } + if (!self.e164 && !self.jid && !self.lid) { + return null; + } + const successor = getRegisteredWhatsAppConnectionController(options.accountId); + if (!successor) { + return null; + } + const successorIdentity = successor.getSelfIdentity(); + if (!successorIdentity || !identitiesOverlap(self, successorIdentity)) { + return null; + } + return successor.getCurrentSock(); + }; type QueuedInboundMessage = WebInboundMessage & { dedupeKey?: string; debounceKey?: string; diff --git a/extensions/whatsapp/src/monitor-inbox.streams-inbound-messages.test-support.ts b/extensions/whatsapp/src/monitor-inbox.streams-inbound-messages.test-support.ts index 656f70da74a0..39e448d2abbd 100644 --- a/extensions/whatsapp/src/monitor-inbox.streams-inbound-messages.test-support.ts +++ b/extensions/whatsapp/src/monitor-inbox.streams-inbound-messages.test-support.ts @@ -3,11 +3,16 @@ import fsSync from "node:fs"; import path from "node:path"; import "./monitor-inbox.test-harness.js"; import { beforeEach, describe, expect, it, vi } from "vitest"; +import { + registerWhatsAppConnectionController, + unregisterWhatsAppConnectionController, +} from "./connection-controller-registry.js"; import { WhatsAppRetryableInboundError } from "./inbound/dedupe.js"; import { WHATSAPP_GROUP_METADATA_CACHE_MAX_ENTRIES } from "./inbound/monitor.js"; import { type InboxMonitorOptions, buildNotifyMessageUpsert, + DEFAULT_ACCOUNT_ID, failNextWhatsAppPluginStateRegisterIfAbsent, getAuthDir, getSock, @@ -809,6 +814,224 @@ describe("web monitor inbox", () => { await listener.close(); }); + // VE-513: when the gateway's channel-health-monitor tears down controller A + // mid-run (shutdown nulls A.socketRef and aborts A.disconnectRetries), then a + // fresh controller B registers for the same accountId, an in-flight inbound's + // captured reply closure must route through B's socket via the registry instead + // of throwing RECONNECT_IN_PROGRESS. Before the registry fallback was added, + // sendTrackedMessage only knew its own controller's socketRef and the captured + // reply was permanently broken. + it("routes the captured reply through a successor controller when the original controller's socket is gone", async () => { + const onMessage = vi.fn(async () => undefined); + const socketRefA = createSocketRef(); + + let aShouldRetryDisconnect = true; + const { listener: listenerA, sock: sockA } = await startInboxMonitor( + onMessage as InboxOnMessage, + { + socketRef: socketRefA, + shouldRetryDisconnect: () => aShouldRetryDisconnect, + disconnectRetryPolicy: { + initialMs: 1, + maxMs: 1, + factor: 1, + jitter: 0, + maxAttempts: 1, + }, + }, + ); + + sockA.ev.emit( + "messages.upsert", + buildNotifyMessageUpsert({ + id: nextMessageId("outbound-successor"), + remoteJid: "999@s.whatsapp.net", + text: "ping", + timestamp: 1_700_000_000, + pushName: "Tester", + }), + ); + await waitForMessageCalls(onMessage, 1); + const inbound = inboundMessage(onMessage) as { + reply: (text: string) => Promise; + sendMedia: (payload: Record) => Promise; + }; + + // The mock harness socket exposes user.id = "123@s.whatsapp.net"; the + // successor handle must report a self identity that overlaps that JID + // so the session-safety guard accepts the fallback. + const handleA = { + getActiveListener: () => null, + getCurrentSock: () => null, + getSelfIdentity: () => null, + } as never; + registerWhatsAppConnectionController(DEFAULT_ACCOUNT_ID, handleA); + + // === Simulate health-monitor-driven shutdown of controller A === + socketRefA.current = null; + aShouldRetryDisconnect = false; + unregisterWhatsAppConnectionController(DEFAULT_ACCOUNT_ID, handleA); + + // === Successor controller B comes up with its OWN socket and registers === + const sockB = { + sendMessage: vi.fn(async () => ({ key: { id: "post-restart-msg-id" } })), + }; + const handleB = { + getActiveListener: () => null, + getCurrentSock: () => sockB as never, + getSelfIdentity: () => ({ jid: "123@s.whatsapp.net", lid: null }), + } as never; + registerWhatsAppConnectionController(DEFAULT_ACCOUNT_ID, handleB); + + try { + await inbound.reply("pong"); + await inbound.sendMedia({ text: "media after restart" }); + + // Captured A reply routed through B via the registry handle. + expect(sockB.sendMessage).toHaveBeenCalledTimes(2); + expect(sockB.sendMessage).toHaveBeenNthCalledWith(1, "999@s.whatsapp.net", { + text: "pong", + }); + expect(sockB.sendMessage).toHaveBeenNthCalledWith(2, "999@s.whatsapp.net", { + text: "media after restart", + }); + } finally { + unregisterWhatsAppConnectionController(DEFAULT_ACCOUNT_ID, handleB); + await listenerA.close(); + } + }); + + // VE-513 PN/LID normalization: Baileys sockets may expose self identity in + // different forms across reconnects (PN JID `@s.whatsapp.net` vs LID + // `@lid`). The fallback must recognize the same account when one side + // reports only PN and the other only LID, because the captured reply + // shouldn't be dropped just because the identity form rotated. + it("accepts successor-controller fallback when only the LID-vs-PN form differs for the same account e164", async () => { + const onMessage = vi.fn(async () => undefined); + const socketRefA = createSocketRef(); + + let aShouldRetryDisconnect = true; + const { listener: listenerA, sock: sockA } = await startInboxMonitor( + onMessage as InboxOnMessage, + { + socketRef: socketRefA, + shouldRetryDisconnect: () => aShouldRetryDisconnect, + disconnectRetryPolicy: { + initialMs: 1, + maxMs: 1, + factor: 1, + jitter: 0, + maxAttempts: 1, + }, + }, + ); + + sockA.ev.emit( + "messages.upsert", + buildNotifyMessageUpsert({ + id: nextMessageId("outbound-successor-lid"), + remoteJid: "999@s.whatsapp.net", + text: "ping", + timestamp: 1_700_000_000, + pushName: "Tester", + }), + ); + await waitForMessageCalls(onMessage, 1); + const inbound = inboundMessage(onMessage) as { reply: (text: string) => Promise }; + + socketRefA.current = null; + aShouldRetryDisconnect = false; + + // The mock harness socket has user.id="123@s.whatsapp.net" (PN JID, no + // lid). The successor reports only the LID form. Both resolve to the same + // synthetic e164 via resolveComparableIdentity, so identitiesOverlap + // accepts the fallback. (resolveComparableIdentity preserves an explicit + // e164 on the input over deriving from the JID via authDir lookup.) + const sharedE164 = "+123"; + const sockB = { + sendMessage: vi.fn(async () => ({ key: { id: "post-restart-lid-msg-id" } })), + }; + const handleB = { + getActiveListener: () => null, + getCurrentSock: () => sockB as never, + getSelfIdentity: () => ({ jid: null, lid: "12300:1@lid", e164: sharedE164 }), + } as never; + registerWhatsAppConnectionController(DEFAULT_ACCOUNT_ID, handleB); + + try { + await inbound.reply("pong"); + expect(sockB.sendMessage).toHaveBeenCalledTimes(1); + expect(sockB.sendMessage).toHaveBeenCalledWith("999@s.whatsapp.net", { text: "pong" }); + } finally { + unregisterWhatsAppConnectionController(DEFAULT_ACCOUNT_ID, handleB); + await listenerA.close(); + } + }); + + // VE-513 session-safety guard: if the registered successor controller has been + // re-linked to a different WhatsApp identity (different self JID), the + // captured reply must fail closed rather than route through the wrong number. + it("refuses successor-controller fallback when the registered controller's self JID does not match", async () => { + const onMessage = vi.fn(async () => undefined); + const socketRefA = createSocketRef(); + + let aShouldRetryDisconnect = true; + const { listener: listenerA, sock: sockA } = await startInboxMonitor( + onMessage as InboxOnMessage, + { + socketRef: socketRefA, + shouldRetryDisconnect: () => aShouldRetryDisconnect, + disconnectRetryPolicy: { + initialMs: 1, + maxMs: 1, + factor: 1, + jitter: 0, + maxAttempts: 1, + }, + }, + ); + + sockA.ev.emit( + "messages.upsert", + buildNotifyMessageUpsert({ + id: nextMessageId("outbound-successor-mismatch"), + remoteJid: "999@s.whatsapp.net", + text: "ping", + timestamp: 1_700_000_000, + pushName: "Tester", + }), + ); + await waitForMessageCalls(onMessage, 1); + const inbound = inboundMessage(onMessage) as { reply: (text: string) => Promise }; + + // A's shutdown sequence. + socketRefA.current = null; + aShouldRetryDisconnect = false; + + // === Successor is registered but with a DIFFERENT self identity (relink/repair) === + const sockB = { + sendMessage: vi.fn(async () => ({ key: { id: "should-not-be-called" } })), + }; + const handleBMismatch = { + getActiveListener: () => null, + getCurrentSock: () => sockB as never, + getSelfIdentity: () => ({ jid: "456@s.whatsapp.net", lid: null }), + } as never; + registerWhatsAppConnectionController(DEFAULT_ACCOUNT_ID, handleBMismatch); + + try { + await expect(inbound.reply("pong")).rejects.toThrow( + "no active socket - reconnection in progress", + ); + + // The mismatched successor's socket was never used. + expect(sockB.sendMessage).not.toHaveBeenCalled(); + } finally { + unregisterWhatsAppConnectionController(DEFAULT_ACCOUNT_ID, handleBMismatch); + await listenerA.close(); + } + }); + it("deduplicates redelivered messages by id", async () => { const onMessage = vi.fn(async () => {});