refactor(whatsapp): split inbox lifecycle coverage (#117910)

* test(whatsapp): split inbox lifecycle coverage

* test(whatsapp): import metadata cache type
This commit is contained in:
Peter Steinberger
2026-08-02 02:09:47 -07:00
committed by GitHub
parent a0d4aa3466
commit d6aafb3dfc
9 changed files with 2503 additions and 2604 deletions
-2
View File
@@ -288,11 +288,9 @@ extensions/voice-call/src/webhook.ts
extensions/voice-call/src/webhook/realtime-handler.test.ts
extensions/voice-call/src/webhook/realtime-handler.ts
extensions/webhooks/src/http.ts
extensions/whatsapp/src/auto-reply.web-auto-reply.connection-and-logging.e2e.test.ts
extensions/whatsapp/src/auto-reply/monitor/inbound-dispatch.test.ts
extensions/whatsapp/src/auto-reply/monitor/inbound-dispatch.ts
extensions/whatsapp/src/connection-controller.ts
extensions/whatsapp/src/monitor-inbox.streams-inbound-messages.test-support.ts
extensions/workboard/src/sqlite-store.ts
extensions/workboard/src/store-card-helpers.ts
extensions/workboard/src/store-core.ts
@@ -65,54 +65,6 @@ function requireOnMessage(
return value as Parameters<typeof sendWebDirectInboundMessage>[0]["onMessage"];
}
async function startWatchdogScenario(params: {
monitorWebChannel: typeof import("./auto-reply/monitor.js").monitorWebChannel;
statusSink?: Parameters<typeof startWebAutoReplyMonitor>[0]["statusSink"];
}) {
const sleep = vi.fn(async () => {});
const scripted = createScriptedWebListenerFactory();
const started = startWebAutoReplyMonitor({
monitorWebChannelFn: params.monitorWebChannel as never,
listenerFactory: scripted.listenerFactory,
sleep,
heartbeatSeconds: 60,
messageTimeoutMs: 30,
watchdogCheckMs: 5,
statusSink: params.statusSink,
});
await vi.waitFor(
() => {
expect(scripted.getListenerCount()).toBe(1);
},
{ timeout: 250, interval: 2 },
);
await vi.waitFor(
() => {
expect(scripted.getOnMessage()).toBeTypeOf("function");
},
{ timeout: 250, interval: 2 },
);
await requireOnMessage(scripted.getOnMessage())(
createTestWebInboundMessage({
event: { id: "m1" },
payload: { body: "ignored" },
admission: {
conversation: { kind: "direct", id: "+1" },
ingress: {
admission: "drop",
decision: "block",
decisiveGateId: "sender",
reasonCode: "no_policy_match",
},
},
}),
);
return { scripted, sleep, ...started };
}
function expectErrorContaining(errorFn: unknown, text: string): void {
const messages = ((errorFn as { mock?: { calls?: unknown[][] } }).mock?.calls ?? []).map(
(call) =>
@@ -136,6 +88,20 @@ function mockCallArg(mocked: unknown, callIndex: number, argIndex: number): unkn
return call[argIndex];
}
async function waitForScriptedListeners(
scripted: ReturnType<typeof createScriptedWebListenerFactory>,
count: number,
atLeast = false,
) {
await vi.waitFor(
() => {
const actual = scripted.getListenerCount();
expect(atLeast ? Math.min(actual, count) : actual).toBe(count);
},
{ timeout: 250, interval: 2 },
);
}
describe("web auto-reply connection", () => {
installWebAutoReplyUnitTestHooks();
@@ -661,318 +627,218 @@ describe("web auto-reply connection", () => {
expectErrorContaining(runtime.error, "Stopping web monitoring");
});
it("forces reconnect when watchdog closes without onClose", async () => {
vi.useFakeTimers();
try {
const statuses: Array<Record<string, unknown>> = [];
const { scripted, controller, run, runtime } = await startWatchdogScenario({
monitorWebChannel,
statusSink: (status) => statuses.push({ ...status }),
});
type WatchdogCaseContext = {
scripted: ReturnType<typeof createScriptedWebListenerFactory>;
statuses: Array<Record<string, unknown>>;
runtime: ReturnType<typeof startWebAutoReplyMonitor>["runtime"];
};
type WatchdogCase = {
name: string;
seedInbound?: boolean;
captureStatus?: boolean;
options?: {
heartbeatSeconds?: number;
transportTimeoutMs?: number;
messageTimeoutMs?: number;
watchdogCheckMs?: number;
};
cleanupWithoutError?: boolean;
exercise: (context: WatchdogCaseContext) => Promise<void>;
assertAfterRun?: (context: WatchdogCaseContext) => void;
};
await vi.advanceTimersByTimeAsync(200);
await Promise.resolve();
await vi.waitFor(
() => {
expect(scripted.getListenerCount()).toBeGreaterThanOrEqual(2);
},
{ timeout: 250, interval: 2 },
);
controller.abort();
scripted.resolveClose(1, { status: 499, isLoggedOut: false });
await Promise.resolve();
await run;
expect(mockStringMessages(runtime.log).join("\n")).toContain(
"WhatsApp Web watchdog is recovering a stale connection",
);
expect(mockStringMessages(runtime.error).join("\n")).not.toContain("status 499");
expect(
statuses.filter(
(status) =>
status.healthState === "reconnecting" &&
status.reconnectAttempts === 1 &&
(status.lastDisconnect as { status?: number } | null)?.status === 499,
),
).not.toEqual([]);
expect(
statuses.filter(
(status) =>
status.lastDisconnect &&
typeof status.lastDisconnect === "object" &&
"expected" in status.lastDisconnect,
),
).toEqual([]);
expect(
statuses.filter(
(status) =>
status.connected === true &&
status.healthState === "healthy" &&
status.reconnectAttempts === 0 &&
status.lastDisconnect === null,
),
).not.toEqual([]);
} finally {
vi.useRealTimers();
}
});
it("keeps quiet linked-device sessions open when transport frames keep arriving", async () => {
vi.useFakeTimers();
try {
const sleep = vi.fn(async () => {});
const scripted = createScriptedWebListenerFactory();
const { controller, run } = startWebAutoReplyMonitor({
monitorWebChannelFn: monitorWebChannel as never,
listenerFactory: scripted.listenerFactory,
sleep,
heartbeatSeconds: 60,
messageTimeoutMs: 30,
watchdogCheckMs: 5,
});
await vi.waitFor(
() => {
expect(scripted.getListenerCount()).toBe(1);
},
{ timeout: 250, interval: 2 },
);
const socket = getLastWebAutoReplySessionSocket();
await vi.advanceTimersByTimeAsync(20);
socket.ws.emit("frame");
await vi.advanceTimersByTimeAsync(20);
socket.ws.emit("frame");
await vi.advanceTimersByTimeAsync(20);
expect(scripted.getListenerCount()).toBe(1);
controller.abort();
scripted.resolveClose(0, { status: 499, isLoggedOut: false });
await Promise.resolve();
await run;
} finally {
vi.useRealTimers();
}
});
it("does not let transport frames mask application silence forever", async () => {
vi.useFakeTimers();
try {
const sleep = vi.fn(async () => {});
const scripted = createScriptedWebListenerFactory();
const { controller, run } = startWebAutoReplyMonitor({
monitorWebChannelFn: monitorWebChannel as never,
listenerFactory: scripted.listenerFactory,
sleep,
heartbeatSeconds: 60,
messageTimeoutMs: 30,
watchdogCheckMs: 5,
});
await vi.waitFor(
() => {
expect(scripted.getListenerCount()).toBe(1);
},
{ timeout: 250, interval: 2 },
);
const socket = getLastWebAutoReplySessionSocket();
for (let elapsedMs = 0; elapsedMs < 140; elapsedMs += 20) {
const watchdogCases = [
{
name: "forces reconnect when watchdog closes without onClose",
seedInbound: true,
captureStatus: true,
cleanupWithoutError: true,
exercise: async ({ scripted }) => {
await vi.advanceTimersByTimeAsync(200);
await Promise.resolve();
await waitForScriptedListeners(scripted, 2, true);
},
assertAfterRun: ({ runtime, statuses }) => {
expect(mockStringMessages(runtime.log).join("\n")).toContain(
"WhatsApp Web watchdog is recovering a stale connection",
);
expect(mockStringMessages(runtime.error).join("\n")).not.toContain("status 499");
expect(
statuses.filter(
(status) =>
status.healthState === "reconnecting" &&
status.reconnectAttempts === 1 &&
(status.lastDisconnect as { status?: number } | null)?.status === 499,
),
).not.toEqual([]);
expect(
statuses.filter(
(status) =>
status.lastDisconnect &&
typeof status.lastDisconnect === "object" &&
"expected" in status.lastDisconnect,
),
).toEqual([]);
expect(
statuses.filter(
(status) =>
status.connected === true &&
status.healthState === "healthy" &&
status.reconnectAttempts === 0 &&
status.lastDisconnect === null,
),
).not.toEqual([]);
},
},
{
name: "keeps quiet linked-device sessions open when transport frames keep arriving",
cleanupWithoutError: true,
exercise: async ({ scripted }) => {
const socket = getLastWebAutoReplySessionSocket();
await vi.advanceTimersByTimeAsync(20);
socket.ws.emit("frame");
await vi.advanceTimersByTimeAsync(20);
}
await vi.waitFor(
() => {
expect(scripted.getListenerCount()).toBeGreaterThanOrEqual(2);
},
{ timeout: 250, interval: 2 },
);
controller.abort();
scripted.resolveClose(scripted.getListenerCount() - 1, {
status: 499,
isLoggedOut: false,
error: "aborted",
});
await Promise.resolve();
await run;
} finally {
vi.useRealTimers();
}
});
it("publishes frame-driven transport activity for quiet sessions", async () => {
vi.useFakeTimers();
try {
const sleep = vi.fn(async () => {});
const statuses: Array<Record<string, unknown>> = [];
const scripted = createScriptedWebListenerFactory();
const { controller, run } = startWebAutoReplyMonitor({
monitorWebChannelFn: monitorWebChannel as never,
listenerFactory: scripted.listenerFactory,
sleep,
socket.ws.emit("frame");
await vi.advanceTimersByTimeAsync(20);
expect(scripted.getListenerCount()).toBe(1);
},
},
{
name: "does not let transport frames mask application silence forever",
exercise: async ({ scripted }) => {
const socket = getLastWebAutoReplySessionSocket();
for (let elapsedMs = 0; elapsedMs < 140; elapsedMs += 20) {
socket.ws.emit("frame");
await vi.advanceTimersByTimeAsync(20);
}
await waitForScriptedListeners(scripted, 2, true);
},
},
{
name: "publishes frame-driven transport activity for quiet sessions",
captureStatus: true,
options: {
heartbeatSeconds: 1,
transportTimeoutMs: 60_000,
messageTimeoutMs: 60_000,
watchdogCheckMs: 5,
statusSink: (next) => statuses.push({ ...next }),
});
await vi.waitFor(
() => {
expect(scripted.getListenerCount()).toBe(1);
},
{ timeout: 250, interval: 2 },
);
const initialTransportAt = Number(statuses.at(-1)?.lastTransportActivityAt ?? 0);
const socket = getLastWebAutoReplySessionSocket();
await vi.advanceTimersByTimeAsync(250);
socket.ws.emit("frame");
await vi.advanceTimersByTimeAsync(1_000);
const lastTransportAt = Number(statuses.at(-1)?.lastTransportActivityAt ?? 0);
expect(lastTransportAt).toBeGreaterThan(initialTransportAt);
controller.abort();
scripted.resolveClose(0, { status: 499, isLoggedOut: false, error: "aborted" });
await Promise.resolve();
await run;
} finally {
vi.useRealTimers();
}
});
it("reconnects on transport stall before the long app-silence window", async () => {
vi.useFakeTimers();
try {
const sleep = vi.fn(async () => {});
const scripted = createScriptedWebListenerFactory();
const { controller, run } = startWebAutoReplyMonitor({
monitorWebChannelFn: monitorWebChannel as never,
listenerFactory: scripted.listenerFactory,
sleep,
},
exercise: async ({ statuses }) => {
const initialTransportAt = Number(statuses.at(-1)?.lastTransportActivityAt ?? 0);
const socket = getLastWebAutoReplySessionSocket();
await vi.advanceTimersByTimeAsync(250);
socket.ws.emit("frame");
await vi.advanceTimersByTimeAsync(1_000);
const lastTransportAt = Number(statuses.at(-1)?.lastTransportActivityAt ?? 0);
expect(lastTransportAt).toBeGreaterThan(initialTransportAt);
},
},
{
name: "reconnects on transport stall before the long app-silence window",
options: {
heartbeatSeconds: 1,
transportTimeoutMs: 30,
messageTimeoutMs: 3_000,
watchdogCheckMs: 5,
});
},
exercise: async ({ scripted }) => {
await vi.advanceTimersByTimeAsync(36);
await Promise.resolve();
await waitForScriptedListeners(scripted, 2, true);
},
},
{
name: "recovers a post-408 listener when transport frames continue but app delivery stays silent",
seedInbound: true,
exercise: async ({ scripted }) => {
scripted.resolveClose(0, {
status: 408,
isLoggedOut: false,
error: "status=408 Request Time-out",
});
await waitForScriptedListeners(scripted, 2);
const reconnectedSocket = getLastWebAutoReplySessionSocket();
for (let elapsedMs = 0; elapsedMs < 45; elapsedMs += 5) {
reconnectedSocket.ws.emit("frame");
await vi.advanceTimersByTimeAsync(5);
}
await waitForScriptedListeners(scripted, 3, true);
},
},
{
name: "gives a reconnected listener a fresh watchdog window",
seedInbound: true,
exercise: async ({ scripted }) => {
scripted.resolveClose(0, { status: 499, isLoggedOut: false, error: "first-close" });
await waitForScriptedListeners(scripted, 2);
await vi.advanceTimersByTimeAsync(20);
await Promise.resolve();
expect(scripted.getListenerCount()).toBe(2);
await vi.advanceTimersByTimeAsync(20);
await Promise.resolve();
await waitForScriptedListeners(scripted, 3, true);
},
},
] satisfies WatchdogCase[];
await vi.waitFor(
() => {
expect(scripted.getListenerCount()).toBe(1);
},
{ timeout: 250, interval: 2 },
);
await vi.advanceTimersByTimeAsync(36);
await Promise.resolve();
await vi.waitFor(
() => {
expect(scripted.getListenerCount()).toBeGreaterThanOrEqual(2);
},
{ timeout: 250, interval: 2 },
);
controller.abort();
scripted.resolveClose(scripted.getListenerCount() - 1, {
status: 499,
isLoggedOut: false,
error: "aborted",
});
await Promise.resolve();
await run;
} finally {
vi.useRealTimers();
}
});
it("recovers a post-408 listener when transport frames continue but app delivery stays silent", async () => {
it.each(watchdogCases)("$name", async (scenario) => {
vi.useFakeTimers();
const sleep = vi.fn(async () => {});
const statuses: Array<Record<string, unknown>> = [];
const scripted = createScriptedWebListenerFactory();
const { controller, run, runtime } = startWebAutoReplyMonitor({
monitorWebChannelFn: monitorWebChannel as never,
listenerFactory: scripted.listenerFactory,
sleep,
heartbeatSeconds: 60,
messageTimeoutMs: 30,
watchdogCheckMs: 5,
...scenario.options,
statusSink: scenario.captureStatus ? (status) => statuses.push({ ...status }) : undefined,
});
const context = { scripted, statuses, runtime };
try {
const { scripted, controller, run } = await startWatchdogScenario({
monitorWebChannel,
});
scripted.resolveClose(0, {
status: 408,
isLoggedOut: false,
error: "status=408 Request Time-out",
});
await vi.waitFor(
() => {
expect(scripted.getListenerCount()).toBe(2);
},
{ timeout: 250, interval: 2 },
);
const reconnectedSocket = getLastWebAutoReplySessionSocket();
for (let elapsedMs = 0; elapsedMs < 45; elapsedMs += 5) {
reconnectedSocket.ws.emit("frame");
await vi.advanceTimersByTimeAsync(5);
await waitForScriptedListeners(scripted, 1);
if (scenario.seedInbound) {
await vi.waitFor(
() => {
expect(scripted.getOnMessage()).toBeTypeOf("function");
},
{ timeout: 250, interval: 2 },
);
await requireOnMessage(scripted.getOnMessage())(
createTestWebInboundMessage({
event: { id: "m1" },
payload: { body: "ignored" },
admission: {
conversation: { kind: "direct", id: "+1" },
ingress: {
admission: "drop",
decision: "block",
decisiveGateId: "sender",
reasonCode: "no_policy_match",
},
},
}),
);
}
await vi.waitFor(
() => {
expect(scripted.getListenerCount()).toBeGreaterThanOrEqual(3);
},
{ timeout: 250, interval: 2 },
);
await scenario.exercise(context);
} finally {
controller.abort();
scripted.resolveClose(scripted.getListenerCount() - 1, {
status: 499,
isLoggedOut: false,
error: "aborted",
...(scenario.cleanupWithoutError ? {} : { error: "aborted" }),
});
await Promise.resolve();
await run;
} finally {
vi.useRealTimers();
try {
await run;
} finally {
vi.useRealTimers();
}
}
});
it("gives a reconnected listener a fresh watchdog window", async () => {
vi.useFakeTimers();
try {
const { scripted, controller, run } = await startWatchdogScenario({
monitorWebChannel,
});
scripted.resolveClose(0, { status: 499, isLoggedOut: false, error: "first-close" });
await vi.waitFor(
() => {
expect(scripted.getListenerCount()).toBe(2);
},
{ timeout: 250, interval: 2 },
);
await vi.advanceTimersByTimeAsync(20);
await Promise.resolve();
expect(scripted.getListenerCount()).toBe(2);
await vi.advanceTimersByTimeAsync(20);
await Promise.resolve();
await vi.waitFor(
() => {
expect(scripted.getListenerCount()).toBeGreaterThanOrEqual(3);
},
{ timeout: 250, interval: 2 },
);
controller.abort();
scripted.resolveClose(scripted.getListenerCount() - 1, {
status: 499,
isLoggedOut: false,
error: "aborted",
});
await Promise.resolve();
await run;
} finally {
vi.useRealTimers();
}
scenario.assertAfterRun?.(context);
});
it("passes the global inbound debounce into the live listener", async () => {
@@ -1229,4 +1095,3 @@ function toLintErrorObject(value: unknown, fallbackMessage: string): Error {
}
return error;
}
/* oxlint-disable max-lines -- TODO: split this grandfathered oversized file. */
@@ -653,6 +653,53 @@ async function runWhatsAppReplyPlan(
return plan.finalize(dispatchResult);
}
const DEFERRED_TOOL_MEDIA = {
text: "tool image",
mediaUrls: ["/tmp/generated.jpg"],
};
const CAPTIONED_MEDIA_REPLACEMENT = {
text: "captioned replacement",
mediaUrls: ["/tmp/generated.jpg"],
};
function requireCapturedDeliver(params: CapturedDispatchParams) {
const deliver = params.dispatcherOptions?.deliver;
if (!deliver) {
throw new Error("expected captured deliver callback");
}
return deliver;
}
async function dispatchDeferredMediaReplacement(overrides: BufferedReplyOverrides): Promise<{
finalization: Promise<unknown>;
replacement: PromiseSettledResult<unknown>;
}> {
let finalization: Promise<unknown> | undefined;
let replacement: PromiseSettledResult<unknown> | undefined;
dispatchReplyWithBufferedBlockDispatcherMock.mockImplementationOnce(
async (params: CapturedDispatchParams) => {
capturedDispatchParams = params;
const deliver = requireCapturedDeliver(params);
const deferred = requireRecord(
await deliver(DEFERRED_TOOL_MEDIA, { kind: "tool" }),
"deferred media result",
);
finalization = deferred.finalization as Promise<unknown>;
[replacement] = await Promise.allSettled([
deliver(CAPTIONED_MEDIA_REPLACEMENT, { kind: "block" }),
]);
await params.dispatcherOptions?.onSettled?.();
return { queuedFinal: false, counts: { tool: 1, block: 1, final: 0 } };
},
);
await dispatchBufferedReply(overrides);
if (!finalization || !replacement) {
throw new Error("expected deferred media replacement lifecycle");
}
return { finalization, replacement };
}
describe("whatsapp inbound dispatch", () => {
beforeEach(() => {
capturedDispatchParams = undefined;
@@ -1055,42 +1102,19 @@ describe("whatsapp inbound dispatch", () => {
it("retains approved deferred media when its captioned replacement is cancelled", async () => {
const deliverReply = vi.fn(async () => acceptedDeliveryResult());
let deferredFinalization: Promise<unknown> | undefined;
dispatchReplyWithBufferedBlockDispatcherMock.mockImplementationOnce(
async (params: CapturedDispatchParams) => {
capturedDispatchParams = params;
const deliver = params.dispatcherOptions?.deliver;
if (!deliver) {
throw new Error("expected captured deliver callback");
}
const deferred = requireRecord(
await deliver(
{ text: "tool image", mediaUrls: ["/tmp/generated.jpg"] },
{ kind: "tool" },
),
"deferred media result",
);
deferredFinalization = deferred.finalization as Promise<unknown>;
await expect(
deliver(
{ text: "captioned replacement", mediaUrls: ["/tmp/generated.jpg"] },
{ kind: "block" },
),
).resolves.toMatchObject({
visibleReplySent: false,
suppression: { reason: "cancelled_by_message_sending_hook" },
});
await params.dispatcherOptions?.onSettled?.();
return { queuedFinal: false, counts: { tool: 1, block: 1, final: 0 } };
},
);
await dispatchBufferedReply({
const { finalization, replacement } = await dispatchDeferredMediaReplacement({
deliverReply,
cancelAfterPrepare: (payload) => payload.text === "captioned replacement",
});
await expect(deferredFinalization).resolves.toMatchObject({ visibleReplySent: true });
expect(replacement).toMatchObject({
status: "fulfilled",
value: {
visibleReplySent: false,
suppression: { reason: "cancelled_by_message_sending_hook" },
},
});
await expect(finalization).resolves.toMatchObject({ visibleReplySent: true });
expect(deliverReply).toHaveBeenCalledTimes(1);
expectReplyResultFields(deliverReply, {
mediaUrls: ["/tmp/generated.jpg"],
@@ -1104,39 +1128,10 @@ describe("whatsapp inbound dispatch", () => {
visibleReplySent: true,
});
const deliverReply = vi.fn().mockRejectedValueOnce(error);
let deferredFinalization: Promise<unknown> | undefined;
let replacementFailure: unknown;
dispatchReplyWithBufferedBlockDispatcherMock.mockImplementationOnce(
async (params: CapturedDispatchParams) => {
capturedDispatchParams = params;
const deliver = params.dispatcherOptions?.deliver;
if (!deliver) {
throw new Error("expected captured deliver callback");
}
const deferred = requireRecord(
await deliver(
{ text: "tool image", mediaUrls: ["/tmp/generated.jpg"] },
{ kind: "tool" },
),
"deferred media result",
);
deferredFinalization = deferred.finalization as Promise<unknown>;
try {
await deliver(
{ text: "captioned replacement", mediaUrls: ["/tmp/generated.jpg"] },
{ kind: "block" },
);
} catch (deliveryError: unknown) {
replacementFailure = deliveryError;
}
await params.dispatcherOptions?.onSettled?.();
return { queuedFinal: false, counts: { tool: 1, block: 1, final: 0 } };
},
);
const { finalization, replacement } = await dispatchDeferredMediaReplacement({ deliverReply });
const replacementFailure = replacement.status === "rejected" ? replacement.reason : undefined;
await dispatchBufferedReply({ deliverReply });
await expect(deferredFinalization).resolves.toEqual({ visibleReplySent: false });
await expect(finalization).resolves.toEqual({ visibleReplySent: false });
expect(deliverReply).toHaveBeenCalledTimes(1);
expect(replacementFailure).toMatchObject({
code: "CHANNEL_PARTIAL_DELIVERY",
@@ -1156,39 +1151,13 @@ describe("whatsapp inbound dispatch", () => {
const rememberSentText = vi.fn(() => {
throw error;
});
let deferredFinalization: Promise<unknown> | undefined;
let replacementFailure: unknown;
dispatchReplyWithBufferedBlockDispatcherMock.mockImplementationOnce(
async (params: CapturedDispatchParams) => {
capturedDispatchParams = params;
const deliver = params.dispatcherOptions?.deliver;
if (!deliver) {
throw new Error("expected captured deliver callback");
}
const deferred = requireRecord(
await deliver(
{ text: "tool image", mediaUrls: ["/tmp/generated.jpg"] },
{ kind: "tool" },
),
"deferred media result",
);
deferredFinalization = deferred.finalization as Promise<unknown>;
try {
await deliver(
{ text: "captioned replacement", mediaUrls: ["/tmp/generated.jpg"] },
{ kind: "block" },
);
} catch (deliveryError: unknown) {
replacementFailure = deliveryError;
}
await params.dispatcherOptions?.onSettled?.();
return { queuedFinal: false, counts: { tool: 1, block: 1, final: 0 } };
},
);
const { finalization, replacement } = await dispatchDeferredMediaReplacement({
deliverReply,
rememberSentText,
});
const replacementFailure = replacement.status === "rejected" ? replacement.reason : undefined;
await dispatchBufferedReply({ deliverReply, rememberSentText });
await expect(deferredFinalization).resolves.toEqual({ visibleReplySent: false });
await expect(finalization).resolves.toEqual({ visibleReplySent: false });
expect(deliverReply).toHaveBeenCalledTimes(1);
expect(replacementFailure).toMatchObject({
code: "CHANNEL_PARTIAL_DELIVERY",
@@ -0,0 +1,702 @@
// WhatsApp monitor inbox behavior split by ownership.
import { describe, expect, it, vi } from "vitest";
import { createWhatsAppDurableInboundQueue } from "./inbound/durable-receive.js";
import { resolveWhatsAppIngressLifecycle } from "./inbound/ingress-lifecycle.js";
import type { WebInboundMessage } from "./inbound/types.js";
import {
nextMessageId,
inboundMessage,
expectDeprecatedAdmissionAliases,
installStreamsInboundMessageHooks,
} from "./monitor-inbox.streams-inbound-messages.test-support.js";
import {
buildNotifyMessageUpsert,
DEFAULT_ACCOUNT_ID,
resetWebInboundDedupeForTests,
settleInboundWork,
startInboxMonitor,
waitForMessageCalls,
type InboxOnMessage,
} from "./monitor-inbox.test-harness.js";
describe("web monitor inbox delivery and dedupe", () => {
installStreamsInboundMessageHooks();
it("delivery coordinator streams inbound messages", async () => {
const onMessage = vi.fn(async (msg) => {
await msg.sendComposing();
await msg.reply("flat reply works");
await msg.sendMedia({ text: "flat media works" });
});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
expect(sock.sendPresenceUpdate).toHaveBeenNthCalledWith(1, "available");
const messageId = nextMessageId("stream");
const upsert = buildNotifyMessageUpsert({
id: messageId,
remoteJid: "999@s.whatsapp.net",
text: "ping",
timestamp: 1_700_000_000,
pushName: "Tester",
});
sock.ev.emit("messages.upsert", upsert);
await waitForMessageCalls(onMessage, 1);
const inbound = inboundMessage(onMessage);
expect(inbound.payload.body).toBe("ping");
expect(inbound.platform.recipientJid).toBe("+123");
expect(inbound.admission).toMatchObject({
accountId: DEFAULT_ACCOUNT_ID,
conversation: {
kind: "direct",
id: "+999",
},
sender: {
id: "+999",
},
senderAccess: {
allowed: true,
decision: "allow",
},
});
expectDeprecatedAdmissionAliases(inbound);
expect(sock.readMessages).toHaveBeenCalledWith([
{
remoteJid: "999@s.whatsapp.net",
id: messageId,
participant: undefined,
fromMe: false,
},
]);
expect(sock.sendPresenceUpdate).toHaveBeenCalledWith("available");
expect(sock.sendPresenceUpdate).toHaveBeenCalledWith("composing", "999@s.whatsapp.net");
expect(sock.sendMessage).toHaveBeenNthCalledWith(1, "999@s.whatsapp.net", {
text: "flat reply works",
});
expect(sock.sendMessage).toHaveBeenNthCalledWith(2, "999@s.whatsapp.net", {
text: "flat media works",
});
await listener.close();
});
it("delivery coordinator delays read receipts until inbound handlers complete", async () => {
let finishMessage: (() => void) | undefined;
const handlerGate = new Promise<void>((resolve) => {
finishMessage = resolve;
});
const onMessage = vi.fn(async () => {
await handlerGate;
});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
const messageId = nextMessageId("delayed-read");
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: messageId,
remoteJid: "999@s.whatsapp.net",
text: "ping",
timestamp: 1_700_000_000,
pushName: "Tester",
}),
);
await waitForMessageCalls(onMessage, 1);
expect(sock.readMessages).not.toHaveBeenCalled();
finishMessage?.();
await vi.waitFor(() => {
expect(sock.readMessages).toHaveBeenCalledWith([
{
remoteJid: "999@s.whatsapp.net",
id: messageId,
participant: undefined,
fromMe: false,
},
]);
});
await listener.close();
});
it("delivery coordinator keeps the first durable delivery when a duplicate arrives", async () => {
const onMessage = vi.fn(async () => undefined);
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
const messageId = nextMessageId("dup-prepared");
const upsert = buildNotifyMessageUpsert({
id: messageId,
remoteJid: "999@s.whatsapp.net",
text: "first",
timestamp: 1_700_000_000,
pushName: "Tester",
});
sock.ev.emit("messages.upsert", upsert);
// Duplicate delivery of the same message id stays pending behind the first claim.
sock.ev.emit("messages.upsert", upsert);
await waitForMessageCalls(onMessage, 1);
await settleInboundWork();
expect(onMessage).toHaveBeenCalledTimes(1);
expect(inboundMessage(onMessage).payload.body).toBe("first");
await listener.close();
});
it("delivery coordinator retries a transient persistence failure through the drain", async () => {
const onMessage = vi.fn(async () => undefined);
const queue = createWhatsAppDurableInboundQueue(DEFAULT_ACCOUNT_ID);
// One transient rejection absorbs into the bounded retry; the message then
// flows durably. The retired live-dispatch fallback is gone: it bypassed
// drain dedupe and lane serialization once the replay guard was deleted.
const durableInboundQueue = {
...queue,
// First attempt rejects; the bounded retry's second attempt must reach
// the real queue so the message flows durably.
enqueue: vi.fn(queue.enqueue.bind(queue)).mockRejectedValueOnce(new Error("SQLITE_FULL")),
};
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage, {
durableInboundQueue,
});
const messageId = nextMessageId("durable-fallback");
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: messageId,
remoteJid: "999@s.whatsapp.net",
text: "ping",
timestamp: 1_700_000_000,
pushName: "Tester",
}),
);
await waitForMessageCalls(onMessage, 1);
expect(inboundMessage(onMessage).payload.body).toBe("ping");
await vi.waitFor(() => {
expect(sock.readMessages).toHaveBeenCalledWith([
{
remoteJid: "999@s.whatsapp.net",
id: messageId,
participant: undefined,
fromMe: false,
},
]);
});
await listener.close();
});
it("delivery coordinator does not dispatch duplicates with pending durable delivery", async () => {
let finishMessage: (() => void) | undefined;
const handlerGate = new Promise<void>((resolve) => {
finishMessage = resolve;
});
const onMessage = vi.fn(async () => {
await handlerGate;
});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
const messageId = nextMessageId("durable-pending");
const upsert = buildNotifyMessageUpsert({
id: messageId,
remoteJid: "999@s.whatsapp.net",
text: "ping",
timestamp: 1_700_000_000,
pushName: "Tester",
});
sock.ev.emit("messages.upsert", upsert);
await waitForMessageCalls(onMessage, 1);
resetWebInboundDedupeForTests();
sock.ev.emit("messages.upsert", upsert);
await settleInboundWork();
expect(onMessage).toHaveBeenCalledTimes(1);
finishMessage?.();
await listener.close();
});
it("delivery coordinator does not redispatch a completed transport-key duplicate", async () => {
const onMessage = vi.fn(async () => undefined);
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
const upsert = buildNotifyMessageUpsert({
id: nextMessageId("durable-completed"),
remoteJid: "999@s.whatsapp.net",
text: "ping",
timestamp: 1_700_000_000,
pushName: "Tester",
});
sock.ev.emit("messages.upsert", upsert);
await waitForMessageCalls(onMessage, 1);
await vi.waitFor(() => expect(sock.readMessages).toHaveBeenCalledTimes(1));
sock.readMessages.mockClear();
sock.ev.emit("messages.upsert", upsert);
await vi.waitFor(() => expect(sock.readMessages).toHaveBeenCalledTimes(1));
expect(onMessage).toHaveBeenCalledTimes(1);
await listener.close();
});
it("delivery coordinator lets a later same-key flush steer during an active turn", async () => {
let releaseFirst: (() => void) | undefined;
const firstTurn = new Promise<void>((resolve) => {
releaseFirst = resolve;
});
const onMessage = vi.fn(async (message: WebInboundMessage) => {
const lifecycle = resolveWhatsAppIngressLifecycle(message);
if (!lifecycle) {
throw new Error("expected durable ingress lifecycle");
}
await lifecycle.onAdopted();
if (onMessage.mock.calls.length === 1) {
await firstTurn;
}
});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage, {
debounceMs: 20,
});
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: nextMessageId("debounce-steer-1"),
remoteJid: "999@s.whatsapp.net",
text: "first",
timestamp: 1_700_000_000,
pushName: "Tester",
}),
);
await waitForMessageCalls(onMessage, 1);
expect(inboundMessage(onMessage).payload.body).toBe("first");
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: nextMessageId("debounce-steer-2"),
remoteJid: "999@s.whatsapp.net",
text: "steer",
timestamp: 1_700_000_001,
pushName: "Tester",
}),
);
await waitForMessageCalls(onMessage, 2);
expect(inboundMessage(onMessage, 1).payload.body).toBe("steer");
releaseFirst?.();
await listener.close();
});
it("delivery coordinator drains admitted same-lane turns before close completes", async () => {
let releaseFirst: (() => void) | undefined;
const firstTurn = new Promise<void>((resolve) => {
releaseFirst = resolve;
});
const onMessage = vi.fn(async (message: WebInboundMessage) => {
const lifecycle = resolveWhatsAppIngressLifecycle(message);
if (!lifecycle) {
throw new Error("expected durable ingress lifecycle");
}
await lifecycle.onAdopted();
if (onMessage.mock.calls.length === 1) {
await firstTurn;
}
});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage, {
debounceMs: 20,
});
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: nextMessageId("debounce-close-1"),
remoteJid: "999@s.whatsapp.net",
text: "first",
timestamp: 1_700_000_000,
pushName: "Tester",
}),
);
await waitForMessageCalls(onMessage, 1);
const second = buildNotifyMessageUpsert({
id: nextMessageId("debounce-close-2"),
remoteJid: "999@s.whatsapp.net",
text: "second",
timestamp: 1_700_000_001,
pushName: "Tester",
});
const third = buildNotifyMessageUpsert({
id: nextMessageId("debounce-close-3"),
remoteJid: "999@s.whatsapp.net",
text: "third",
timestamp: 1_700_000_002,
pushName: "Tester",
});
sock.ev.emit("messages.upsert", {
type: "notify",
messages: [...second.messages, ...third.messages],
});
let closed = false;
const closePromise = listener.close().then(() => {
closed = true;
});
await waitForMessageCalls(onMessage, 3);
expect(closed).toBe(false);
expect(inboundMessage(onMessage, 1).payload.body).toBe("second");
expect(inboundMessage(onMessage, 2).payload.body).toBe("third");
releaseFirst?.();
await closePromise;
expect(closed).toBe(true);
});
it("delivery coordinator keeps a reused debounce key pending across turns", async () => {
let releaseFirst!: () => void;
let finishFirst!: () => void;
const firstTurn = new Promise<void>((resolve) => {
releaseFirst = resolve;
});
const firstFinished = new Promise<void>((resolve) => {
finishFirst = resolve;
});
const onMessage = vi.fn(async (message: WebInboundMessage) => {
const lifecycle = resolveWhatsAppIngressLifecycle(message);
if (!lifecycle) {
throw new Error("expected durable ingress lifecycle");
}
await lifecycle.onAdopted();
if (onMessage.mock.calls.length === 1) {
await firstTurn;
finishFirst();
}
});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage, {
debounceMs: 60_000,
shouldDebounce: (message) => message.payload.body !== "first",
});
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: nextMessageId("debounce-reused-key-1"),
remoteJid: "999@s.whatsapp.net",
text: "first",
timestamp: 1_700_000_000,
pushName: "Tester",
}),
);
await waitForMessageCalls(onMessage, 1);
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: nextMessageId("debounce-reused-key-2"),
remoteJid: "999@s.whatsapp.net",
text: "second",
timestamp: 1_700_000_001,
pushName: "Tester",
}),
);
await settleInboundWork();
expect(onMessage).toHaveBeenCalledTimes(1);
releaseFirst();
await firstFinished;
await settleInboundWork();
const closeStarted = Date.now();
await listener.close();
expect(Date.now() - closeStarted).toBeLessThan(5_000);
expect(onMessage).toHaveBeenCalledTimes(2);
expect(inboundMessage(onMessage, 1).payload.body).toBe("second");
});
it("delivery coordinator force-flushes long durable debounce during shutdown", async () => {
// Durable pump tasks await claim flush waiters; close must force-flush
// debounced batches before waiting on those pumps (socket-close timeout).
const onMessage = vi.fn(async () => undefined);
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage, {
debounceMs: 60_000,
});
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: nextMessageId("debounce-shutdown-durable"),
remoteJid: "999@s.whatsapp.net",
text: "held in debounce",
timestamp: 1_700_000_100,
pushName: "Tester",
}),
);
// Let accept+pump reach the flush waiter without exhausting the debounce window.
await settleInboundWork();
const closeStarted = Date.now();
await listener.close();
const closeElapsedMs = Date.now() - closeStarted;
expect(closeElapsedMs).toBeLessThan(5_000);
expect(onMessage).toHaveBeenCalledTimes(1);
expect(inboundMessage(onMessage).payload.body).toBe("held in debounce");
expect(sock.end).toHaveBeenCalledTimes(1);
});
it("delivery coordinator drains serialized same-lane replies before socket close", async () => {
vi.useFakeTimers();
try {
let releaseFirst: (() => void) | undefined;
const firstTurn = new Promise<void>((resolve) => {
releaseFirst = resolve;
});
const onMessage = vi.fn(async (msg) => {
await msg.platform.reply("pong");
await msg.platform.sendMedia({ text: "media" });
if (onMessage.mock.calls.length === 1) {
await firstTurn;
}
});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage, {
debounceMs: 50,
});
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: nextMessageId("debounce-close-reply-1"),
remoteJid: "999@s.whatsapp.net",
text: "first",
timestamp: 1_700_000_000,
pushName: "Tester",
}),
);
await vi.advanceTimersByTimeAsync(50);
await waitForMessageCalls(onMessage, 1);
expect(inboundMessage(onMessage).payload.body).toBe("first");
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: nextMessageId("debounce-close-reply-2"),
remoteJid: "999@s.whatsapp.net",
text: "second",
timestamp: 1_700_000_001,
pushName: "Tester",
}),
);
const closePromise = listener.close();
expect(onMessage).toHaveBeenCalledTimes(1);
releaseFirst?.();
await closePromise;
expect(onMessage).toHaveBeenCalledTimes(2);
expect(inboundMessage(onMessage, 1).payload.body).toBe("second");
expect(sock.sendMessage).toHaveBeenNthCalledWith(1, "999@s.whatsapp.net", {
text: "pong",
});
expect(sock.sendMessage).toHaveBeenNthCalledWith(2, "999@s.whatsapp.net", {
text: "media",
});
expect(sock.sendMessage).toHaveBeenNthCalledWith(3, "999@s.whatsapp.net", {
text: "pong",
});
expect(sock.sendMessage).toHaveBeenNthCalledWith(4, "999@s.whatsapp.net", {
text: "media",
});
expect(sock.end).toHaveBeenCalledTimes(1);
expect(sock.sendMessage.mock.invocationCallOrder.at(-1)).toBeLessThan(
sock.end.mock.invocationCallOrder.at(0),
);
} finally {
vi.useRealTimers();
}
});
it("delivery coordinator waits for in-flight handlers before close drain", async () => {
let releaseHandler: (() => void) | undefined;
const handlerGate = new Promise<void>((resolve) => {
releaseHandler = resolve;
});
let markHandlerStarted: (() => void) | undefined;
const handlerStarted = new Promise<void>((resolve) => {
markHandlerStarted = resolve;
});
const onMessage = vi.fn(async (msg) => {
await msg.platform.reply("pong");
markHandlerStarted?.();
await handlerGate;
});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: nextMessageId("close-inflight"),
remoteJid: "999@s.whatsapp.net",
text: "first",
timestamp: 1_700_000_000,
pushName: "Tester",
}),
);
await handlerStarted;
const closePromise = listener.close();
await Promise.resolve();
expect(sock.end).not.toHaveBeenCalled();
if (!releaseHandler) {
throw new Error("Expected handler release callback to be initialized");
}
releaseHandler();
await closePromise;
expect(onMessage).toHaveBeenCalledTimes(1);
expect(sock.sendMessage).toHaveBeenCalledWith("999@s.whatsapp.net", {
text: "pong",
});
expect(sock.end).toHaveBeenCalledTimes(1);
expect(sock.sendMessage.mock.invocationCallOrder.at(0)).toBeLessThan(
sock.end.mock.invocationCallOrder.at(0),
);
});
it("delivery coordinator deduplicates redelivered messages by id", async () => {
const onMessage = vi.fn(async () => {});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
const upsert = buildNotifyMessageUpsert({
id: nextMessageId("dedupe"),
remoteJid: "999@s.whatsapp.net",
text: "ping",
timestamp: 1_700_000_000,
pushName: "Tester",
});
sock.ev.emit("messages.upsert", upsert);
sock.ev.emit("messages.upsert", upsert);
await waitForMessageCalls(onMessage, 1);
expect(onMessage).toHaveBeenCalledTimes(1);
await listener.close();
});
it("delivery coordinator retries redelivery after an explicit retryable failure", async () => {
let attempts = 0;
const onMessage = vi.fn(async () => {
attempts += 1;
if (attempts === 1) {
// Any non-permanent error is retryable to the drain classifier.
throw new Error("retry me");
}
});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
const upsert = buildNotifyMessageUpsert({
id: nextMessageId("retryable-dedupe"),
remoteJid: "999@s.whatsapp.net",
text: "ping",
timestamp: 1_700_000_000,
pushName: "Tester",
});
sock.ev.emit("messages.upsert", upsert);
await waitForMessageCalls(onMessage, 1);
expect(sock.readMessages).not.toHaveBeenCalled();
sock.ev.emit("messages.upsert", upsert);
await waitForMessageCalls(onMessage, 2);
await vi.waitFor(() => {
expect(sock.readMessages).toHaveBeenCalledTimes(1);
});
await listener.close();
});
it("delivery coordinator retries redelivery after reply session conflicts", async () => {
let attempts = 0;
const onMessage = vi.fn(async () => {
attempts += 1;
if (attempts === 1) {
throw new Error(
"reply session initialization conflicted for agent:main:whatsapp:direct:+15551234567",
);
}
});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
const upsert = buildNotifyMessageUpsert({
id: nextMessageId("session-init-conflict"),
remoteJid: "999@s.whatsapp.net",
text: "ping",
timestamp: 1_700_000_000,
pushName: "Tester",
});
sock.ev.emit("messages.upsert", upsert);
await waitForMessageCalls(onMessage, 1);
expect(sock.readMessages).not.toHaveBeenCalled();
sock.ev.emit("messages.upsert", upsert);
await waitForMessageCalls(onMessage, 2);
await vi.waitFor(() => {
expect(sock.readMessages).toHaveBeenCalledTimes(1);
});
await listener.close();
});
it("delivery coordinator keeps same-lane follow-up pending until turn adoption", async () => {
let adoptFirst: (() => void | Promise<void>) | undefined;
const onMessage = vi.fn(async (message: WebInboundMessage) => {
if (!adoptFirst) {
const lifecycle = resolveWhatsAppIngressLifecycle(message);
if (!lifecycle) {
throw new Error("expected durable ingress lifecycle");
}
lifecycle.onDeferred();
adoptFirst = lifecycle.onAdopted;
}
});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
sock.ev.emit("messages.upsert", {
type: "notify",
messages: [
{
key: { id: "abc1", fromMe: false, remoteJid: "999@s.whatsapp.net" },
message: { conversation: "ping" },
messageTimestamp: 1_700_000_000,
},
],
});
await waitForMessageCalls(onMessage, 1);
sock.ev.emit("messages.upsert", {
type: "notify",
messages: [
{
key: { id: "abc2", fromMe: false, remoteJid: "999@s.whatsapp.net" },
message: { conversation: "pong" },
messageTimestamp: 1_700_000_001,
},
],
});
await settleInboundWork();
expect(onMessage).toHaveBeenCalledTimes(1);
expect(inboundMessage(onMessage).payload.body).toBe("ping");
if (!adoptFirst) {
throw new Error("expected first adoption callback");
}
await adoptFirst();
await waitForMessageCalls(onMessage, 2);
expect(inboundMessage(onMessage, 1).payload.body).toBe("pong");
await listener.close();
});
});
@@ -1,2 +0,0 @@
// WhatsApp monitor inbox delivery and lifecycle behavior.
import "./monitor-inbox.streams-inbound-messages.test-support.js";
@@ -0,0 +1,518 @@
// WhatsApp monitor inbox behavior split by ownership.
import type { GroupMetadata } from "baileys";
import { describe, expect, it, vi } from "vitest";
import {
EXPECTED_WHATSAPP_GROUP_METADATA_CACHE_MAX_ENTRIES,
nextMessageId,
inboundMessage,
groupMetadata,
createBaileysCacheSupport,
startInboxMonitorWithBaileysCache,
expectCachedGroupMetadata,
installStreamsInboundMessageHooks,
} from "./monitor-inbox.streams-inbound-messages.test-support.js";
import {
buildNotifyMessageUpsert,
getSock,
settleInboundWork,
startInboxMonitor,
waitForMessageCalls,
type InboxMonitorOptions,
type InboxOnMessage,
} from "./monitor-inbox.test-harness.js";
describe("web monitor inbox metadata cache", () => {
installStreamsInboundMessageHooks();
it("group metadata cache hydrates participating groups once after connect", async () => {
const { listener, sock } = await startInboxMonitor(vi.fn(async () => {}) as InboxOnMessage);
expect(sock.groupFetchAllParticipating).toHaveBeenCalledTimes(1);
await listener.close();
});
it("group metadata cache keeps delivery alive when hydration fails", async () => {
const sock = getSock();
sock.groupFetchAllParticipating.mockRejectedValueOnce(new Error("no groups"));
const { listener } = await startInboxMonitor(vi.fn(async () => {}) as InboxOnMessage);
expect(sock.groupFetchAllParticipating).toHaveBeenCalledTimes(1);
expect(sock.sendPresenceUpdate).toHaveBeenNthCalledWith(1, "available");
await listener.close();
});
it("group metadata cache omits group context when no group facts exist", async () => {
const sock = getSock();
sock.groupFetchAllParticipating.mockRejectedValueOnce(new Error("no groups"));
const onMessage = vi.fn(async () => {});
const { listener } = await startInboxMonitor(onMessage as InboxOnMessage);
sock.groupMetadata.mockRejectedValueOnce(new Error("group metadata unavailable"));
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: nextMessageId("group-no-facts"),
remoteJid: "123@g.us",
participant: "444@s.whatsapp.net",
text: "ping",
timestamp: 1_700_000_000,
}),
);
await waitForMessageCalls(onMessage, 1);
const inbound = inboundMessage(onMessage);
expect(inbound.admission?.conversation.kind).toBe("group");
expect(inbound.group).toBeUndefined();
await listener.close();
});
it("group metadata cache serves reconnect metadata after live fetch failures", async () => {
const groupMetadataCache: NonNullable<InboxMonitorOptions["groupMetadataCache"]> = new Map();
const onMessage = vi.fn(async (_msg: Parameters<InboxOnMessage>[0]) => {});
const firstSock = getSock();
firstSock.groupFetchAllParticipating.mockResolvedValueOnce({
"123@g.us": {
id: "123@g.us",
subject: "Recovered Group",
owner: undefined,
participants: [{ id: "444@s.whatsapp.net" }],
},
});
const first = await startInboxMonitor(onMessage as InboxOnMessage, {
groupMetadataCache,
});
await vi.waitFor(() => {
expect(groupMetadataCache.get("123@g.us")?.subject).toBe("Recovered Group");
});
expect(
(groupMetadataCache.get("123@g.us") as Record<string, unknown>)?.participants,
).toBeUndefined();
await first.listener.close();
const second = await startInboxMonitor(onMessage as InboxOnMessage, {
groupMetadataCache,
});
second.sock.groupMetadata.mockRejectedValueOnce(new Error("408 timed out"));
second.sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: nextMessageId("group-reconnect-cache"),
remoteJid: "123@g.us",
participant: "444@s.whatsapp.net",
text: "ping",
timestamp: 1_700_000_000,
}),
);
await waitForMessageCalls(onMessage, 1);
const inbound = inboundMessage(onMessage);
expect(inbound.payload.body).toBe("ping");
expect(inbound.admission?.conversation.id).toBe("123@g.us");
expect(inbound.group?.subject).toBe("Recovered Group");
expect(inbound.platform.senderE164).toBe("+444");
expect(inbound.admission?.conversation.kind).toBe("group");
expect(inbound.group?.participants).toBeUndefined();
await second.listener.close();
});
it("group metadata cache keeps full participating metadata available to Baileys", async () => {
const sock = getSock();
sock.groupFetchAllParticipating.mockResolvedValueOnce({
"123@g.us": groupMetadata({
subject: "Recovered Group",
participants: ["444@s.whatsapp.net"],
}),
});
const { listener, baileysCache } = await startInboxMonitorWithBaileysCache();
await vi.waitFor(async () => {
await expectCachedGroupMetadata(baileysCache, {
id: "123@g.us",
subject: "Recovered Group",
participants: [{ id: "444@s.whatsapp.net" }],
});
});
await listener.close();
});
it("group metadata cache reuses hydrated participant identities without querying WhatsApp again", async () => {
const participantLid = "277038292303944@lid";
const participantPhone = "15551234567@s.whatsapp.net";
const sock = getSock();
sock.groupFetchAllParticipating.mockResolvedValueOnce({
"123@g.us": {
id: "123@g.us",
subject: "Hydrated Group",
owner: undefined,
participants: [{ id: participantLid, phoneNumber: participantPhone }],
},
});
sock.signalRepository.lidMapping.getPNForLID.mockResolvedValue(participantPhone);
const { listener, baileysCache } = await startInboxMonitorWithBaileysCache();
try {
await vi.waitFor(() => {
expect(baileysCache.baileysGroupMetaCache.has("123@g.us")).toBe(true);
});
sock.groupMetadata.mockRejectedValue(new Error("408 timed out"));
await listener.sendMessage("123@g.us", "ping @+15551234567");
expect(sock.groupMetadata).not.toHaveBeenCalled();
expect(sock.sendMessage).toHaveBeenCalledWith("123@g.us", {
text: "ping @277038292303944",
mentions: [participantLid],
});
} finally {
await listener.close();
}
});
it("group metadata cache refreshes provider snapshots that omit participant identities", async () => {
const sock = getSock();
sock.groupFetchAllParticipating.mockResolvedValueOnce({
"123@g.us": groupMetadata({ subject: "Incomplete Group", participants: [] }),
});
sock.groupMetadata.mockResolvedValueOnce(
groupMetadata({
subject: "Complete Group",
participants: ["15551234567@s.whatsapp.net"],
}),
);
const { listener, baileysCache } = await startInboxMonitorWithBaileysCache();
try {
await vi.waitFor(() => {
expect(baileysCache.baileysGroupMetaCache.get("123@g.us")?.value.participants).toEqual([]);
});
await listener.sendMessage("123@g.us", "recovered @15551234567");
expect(sock.groupMetadata).toHaveBeenCalledOnce();
expect(sock.sendMessage).toHaveBeenCalledWith("123@g.us", {
text: "recovered @15551234567",
mentions: ["15551234567@s.whatsapp.net"],
});
await expectCachedGroupMetadata(baileysCache, {
id: "123@g.us",
subject: "Complete Group",
participants: [{ id: "15551234567@s.whatsapp.net" }],
});
} finally {
await listener.close();
}
});
it("group metadata cache never extends hydrated participant identities beyond their provider expiry", async () => {
const dateNow = vi.spyOn(Date, "now").mockReturnValue(1_700_000_000_000);
const sock = getSock();
sock.groupFetchAllParticipating.mockResolvedValueOnce({
"123@g.us": groupMetadata({
subject: "Expiring Group",
participants: ["15551234567@s.whatsapp.net"],
}),
});
sock.groupMetadata.mockRejectedValue(new Error("408 timed out"));
const { listener, baileysCache } = await startInboxMonitorWithBaileysCache();
try {
await vi.waitFor(() => {
expect(baileysCache.baileysGroupMetaCache.get("123@g.us")?.expiresAt).toBe(
1_700_000_300_000,
);
});
dateNow.mockReturnValue(1_700_000_299_000);
await listener.sendMessage("123@g.us", "fresh @15551234567");
dateNow.mockReturnValue(1_700_000_300_001);
await listener.sendMessage("123@g.us", "expired @15551234567");
expect(sock.sendMessage).toHaveBeenNthCalledWith(1, "123@g.us", {
text: "fresh @15551234567",
mentions: ["15551234567@s.whatsapp.net"],
});
expect(sock.sendMessage).toHaveBeenNthCalledWith(2, "123@g.us", {
text: "expired @15551234567",
});
expect(sock.groupMetadata).toHaveBeenCalledTimes(1);
expect(baileysCache.baileysGroupMetaCache.has("123@g.us")).toBe(false);
} finally {
dateNow.mockRestore();
await listener.close();
}
});
it("group metadata cache drops hydrated local participants when membership changes", async () => {
const sock = getSock();
sock.groupFetchAllParticipating.mockResolvedValueOnce({
"123@g.us": groupMetadata({
subject: "Changing Group",
participants: ["15551234567@s.whatsapp.net"],
}),
});
const { listener, baileysCache } = await startInboxMonitorWithBaileysCache();
try {
await vi.waitFor(() => {
expect(baileysCache.baileysGroupMetaCache.has("123@g.us")).toBe(true);
});
await listener.sendMessage("123@g.us", "before @15551234567");
sock.groupMetadata.mockRejectedValue(new Error("408 timed out"));
sock.ev.emit("group-participants.update", { id: "123@g.us" });
await listener.sendMessage("123@g.us", "after @15551234567");
expect(sock.sendMessage).toHaveBeenNthCalledWith(1, "123@g.us", {
text: "before @15551234567",
mentions: ["15551234567@s.whatsapp.net"],
});
expect(sock.sendMessage).toHaveBeenNthCalledWith(2, "123@g.us", {
text: "after @15551234567",
});
expect(sock.groupMetadata).toHaveBeenCalledTimes(1);
expect(baileysCache.baileysGroupMetaCache.has("123@g.us")).toBe(false);
} finally {
await listener.close();
}
});
it("group metadata cache invalidates partial group and participant updates", async () => {
const groupMetadataCache: NonNullable<InboxMonitorOptions["groupMetadataCache"]> = new Map();
const { listener, sock, baileysCache } = await startInboxMonitorWithBaileysCache({
groupMetadataCache,
});
sock.ev.emit("groups.update", [
groupMetadata({
subject: "Fresh Group",
}),
]);
await expectCachedGroupMetadata(baileysCache, {
id: "123@g.us",
subject: "Fresh Group",
participants: [{ id: "555@s.whatsapp.net" }],
});
expect(groupMetadataCache.has("123@g.us")).toBe(true);
sock.ev.emit("groups.update", [{ id: "123@g.us" }]);
expect(groupMetadataCache.has("123@g.us")).toBe(false);
await expect(
baileysCache.socketOptions.cachedGroupMetadata("123@g.us"),
).resolves.toBeUndefined();
sock.ev.emit("groups.update", [
groupMetadata({
subject: "Fresh Again",
}),
]);
expect(groupMetadataCache.has("123@g.us")).toBe(true);
sock.ev.emit("group-participants.update", { id: "123@g.us" });
expect(groupMetadataCache.has("123@g.us")).toBe(false);
await expect(
baileysCache.socketOptions.cachedGroupMetadata("123@g.us"),
).resolves.toBeUndefined();
await listener.close();
});
it("group metadata cache expires Baileys retry and metadata entries", async () => {
const now = vi.spyOn(Date, "now").mockReturnValue(1_700_000_000_000);
const baileysCache = createBaileysCacheSupport();
const onMessage = vi.fn(async (_msg: Parameters<InboxOnMessage>[0]) => {});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage, {
recentMessageKeys: baileysCache.recentMessageKeys,
baileysGroupMetaCache: baileysCache.baileysGroupMetaCache,
});
const messageId = nextMessageId("baileys-expiry");
try {
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: messageId,
remoteJid: "999@s.whatsapp.net",
text: "retry me",
timestamp: 1_700_000_000,
pushName: "Tester",
}),
);
sock.ev.emit("groups.update", [
groupMetadata({
subject: "Expiring Group",
}),
]);
await waitForMessageCalls(onMessage, 1);
await expect(
baileysCache.socketOptions.getMessage({
id: messageId,
remoteJid: "999@s.whatsapp.net",
}),
).resolves.toEqual({ conversation: "retry me" });
await expectCachedGroupMetadata(baileysCache, {
id: "123@g.us",
subject: "Expiring Group",
participants: [{ id: "555@s.whatsapp.net" }],
});
now.mockReturnValue(1_700_000_000_000 + 5 * 60 * 1000 + 1);
await expect(
baileysCache.socketOptions.cachedGroupMetadata("123@g.us"),
).resolves.toBeUndefined();
now.mockReturnValue(1_700_000_000_000 + 10 * 60 * 1000 + 1);
await expect(
baileysCache.socketOptions.getMessage({
id: messageId,
remoteJid: "999@s.whatsapp.net",
}),
).resolves.toBeUndefined();
} finally {
now.mockRestore();
await listener.close();
}
});
it("group metadata cache does not republish invalidated pending hydration", async () => {
const groupMetadataCache: NonNullable<InboxMonitorOptions["groupMetadataCache"]> = new Map();
const baileysCache = createBaileysCacheSupport();
const sock = getSock();
let resolveHydration!: (groups: Record<string, GroupMetadata>) => void;
sock.groupFetchAllParticipating.mockImplementationOnce(
async () =>
await new Promise<Record<string, GroupMetadata>>((resolve) => {
resolveHydration = resolve;
}),
);
const { listener } = await startInboxMonitor(vi.fn(async () => {}) as InboxOnMessage, {
groupMetadataCache,
recentMessageKeys: baileysCache.recentMessageKeys,
baileysGroupMetaCache: baileysCache.baileysGroupMetaCache,
});
sock.ev.emit("groups.update", [{ id: "123@g.us" }]);
resolveHydration({
"123@g.us": groupMetadata({
subject: "Stale Hydration Group",
}),
});
await settleInboundWork();
expect(groupMetadataCache.has("123@g.us")).toBe(false);
await expect(
baileysCache.socketOptions.cachedGroupMetadata("123@g.us"),
).resolves.toBeUndefined();
await listener.close();
});
it("group metadata cache detaches Baileys listeners on close", async () => {
const baileysCache = createBaileysCacheSupport();
const { listener, sock } = await startInboxMonitor(vi.fn(async () => {}) as InboxOnMessage, {
recentMessageKeys: baileysCache.recentMessageKeys,
baileysGroupMetaCache: baileysCache.baileysGroupMetaCache,
});
expect(sock.ev.listenerCount("groups.upsert")).toBe(1);
expect(sock.ev.listenerCount("groups.update")).toBe(1);
expect(sock.ev.listenerCount("group-participants.update")).toBe(1);
await listener.close();
expect(sock.ev.listenerCount("groups.upsert")).toBe(0);
expect(sock.ev.listenerCount("groups.update")).toBe(0);
expect(sock.ev.listenerCount("group-participants.update")).toBe(0);
});
it("group metadata cache bounds reconnect entries", async () => {
const groupMetadataCache: NonNullable<InboxMonitorOptions["groupMetadataCache"]> = new Map();
const groups = Object.fromEntries(
Array.from({ length: EXPECTED_WHATSAPP_GROUP_METADATA_CACHE_MAX_ENTRIES + 2 }, (_, index) => [
`${index}@g.us`,
{
id: `${index}@g.us`,
subject: `Group ${index}`,
owner: undefined,
participants: [],
},
]),
);
const sock = getSock();
sock.groupFetchAllParticipating.mockResolvedValueOnce(groups);
const { listener } = await startInboxMonitor(vi.fn(async () => {}) as InboxOnMessage, {
groupMetadataCache,
});
await vi.waitFor(() => {
expect(groupMetadataCache.size).toBe(EXPECTED_WHATSAPP_GROUP_METADATA_CACHE_MAX_ENTRIES);
});
expect(groupMetadataCache.has("0@g.us")).toBe(false);
expect(
groupMetadataCache.has(`${EXPECTED_WHATSAPP_GROUP_METADATA_CACHE_MAX_ENTRIES + 1}@g.us`),
).toBe(true);
await listener.close();
});
it("group metadata cache rejects reconnect expiry beyond a valid Date", async () => {
const groupMetadataCache: NonNullable<InboxMonitorOptions["groupMetadataCache"]> = new Map();
const dateNow = vi.spyOn(Date, "now").mockReturnValue(8_640_000_000_000_000);
try {
const sock = getSock();
sock.groupFetchAllParticipating.mockResolvedValueOnce({
"123@g.us": {
id: "123@g.us",
subject: "Boundary Group",
owner: undefined,
participants: [],
},
});
const { listener } = await startInboxMonitor(vi.fn(async () => {}) as InboxOnMessage, {
groupMetadataCache,
});
await vi.waitFor(() => {
expect(sock.groupFetchAllParticipating).toHaveBeenCalledTimes(1);
});
expect(groupMetadataCache.has("123@g.us")).toBe(false);
await listener.close();
} finally {
dateNow.mockRestore();
}
});
it("group metadata cache does not block inbound listeners during hydration", async () => {
let resolveHydration!: () => void;
const sock = getSock();
const pendingHydration = new Promise<Record<string, never>>((resolve) => {
resolveHydration = () => resolve({});
});
sock.groupFetchAllParticipating.mockImplementationOnce(() => pendingHydration);
const onMessage = vi.fn(async () => {});
const { listener } = await startInboxMonitor(onMessage as InboxOnMessage);
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: nextMessageId("pending-hydration"),
remoteJid: "999@s.whatsapp.net",
text: "ping",
timestamp: 1_700_000_000,
pushName: "Tester",
}),
);
await waitForMessageCalls(onMessage, 1);
resolveHydration();
await listener.close();
});
});
@@ -0,0 +1,275 @@
// WhatsApp monitor inbox behavior split by ownership.
import fsSync from "node:fs";
import path from "node:path";
import { describe, expect, it, vi } from "vitest";
import {
imageOps,
nextMessageId,
inboundMessage,
installStreamsInboundMessageHooks,
} from "./monitor-inbox.streams-inbound-messages.test-support.js";
import {
buildNotifyMessageUpsert,
getAuthDir,
startInboxMonitor,
waitForMessageCalls,
type InboxOnMessage,
} from "./monitor-inbox.test-harness.js";
describe("web monitor inbox reply context", () => {
installStreamsInboundMessageHooks();
async function expectQuotedReplyContext(quotedMessage: unknown) {
const onMessage = vi.fn(async (msg) => {
await msg.platform.reply("pong");
});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
const upsert = {
type: "notify",
messages: [
{
key: {
id: nextMessageId("quoted"),
fromMe: false,
remoteJid: "999@s.whatsapp.net",
},
message: {
extendedTextMessage: {
text: "reply",
contextInfo: {
stanzaId: "q1",
participant: "111@s.whatsapp.net",
quotedMessage,
},
},
},
messageTimestamp: 1_700_000_000,
pushName: "Tester",
},
],
};
sock.ev.emit("messages.upsert", upsert);
await waitForMessageCalls(onMessage, 1);
const inbound = inboundMessage(onMessage);
expect(inbound.quote?.id).toBe("q1");
expect(inbound.quote?.body).toBe("original");
expect(inbound.quote?.sender?.displayName).toBe("+111");
const sender = inbound.platform.sender as { e164?: string; name?: string };
expect(sender.e164).toBe("+999");
expect(sender.name).toBe("Tester");
const replyTo = inbound.quote?.context as {
body?: string;
id?: string;
sender?: { e164?: string; jid?: string; label?: string };
};
expect(replyTo.id).toBe("q1");
expect(replyTo.body).toBe("original");
expect(replyTo.sender?.jid).toBe("111@s.whatsapp.net");
expect(replyTo.sender?.e164).toBe("+111");
expect(replyTo.sender?.label).toBe("+111");
const self = inbound.platform.self as { e164?: string; jid?: string };
expect(self.jid).toBe("123@s.whatsapp.net");
expect(self.e164).toBe("+123");
expect(sock.sendMessage).toHaveBeenCalledWith("999@s.whatsapp.net", {
text: "pong",
});
await listener.close();
}
it("prepopulates image previews for inbound sendMedia replies", async () => {
const onMessage = vi.fn(async () => undefined);
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: nextMessageId("image-preview"),
remoteJid: "999@s.whatsapp.net",
text: "ping",
timestamp: 1_700_000_000,
pushName: "Tester",
}),
);
await waitForMessageCalls(onMessage, 1);
const inbound = inboundMessage(onMessage);
const image = Buffer.from("img");
const thumbnail = Buffer.from("thumb");
imageOps.getImageMetadata.mockResolvedValueOnce({ width: 640, height: 480 });
imageOps.resizeToJpeg.mockResolvedValueOnce(thumbnail);
await inbound.platform.sendMedia({ image, caption: "cap", mimetype: "image/png" });
expect(imageOps.resizeToJpeg).toHaveBeenCalledWith({
buffer: image,
maxSide: 32,
quality: 50,
withoutEnlargement: true,
});
expect(sock.sendMessage).toHaveBeenCalledWith("999@s.whatsapp.net", {
image,
caption: "cap",
mimetype: "image/png",
width: 640,
height: 480,
jpegThumbnail: thumbnail.toString("base64"),
});
await listener.close();
});
it("resolves LID JIDs using Baileys LID mapping store", async () => {
const onMessage = vi.fn(async () => {});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
const getPNForLID = vi.spyOn(sock.signalRepository.lidMapping, "getPNForLID");
sock.signalRepository.lidMapping.getPNForLID.mockResolvedValueOnce("999:0@s.whatsapp.net");
const upsert = buildNotifyMessageUpsert({
id: nextMessageId("lid-store"),
remoteJid: "999@lid",
text: "ping",
timestamp: 1_700_000_000,
pushName: "Tester",
});
sock.ev.emit("messages.upsert", upsert);
await waitForMessageCalls(onMessage, 1);
expect(getPNForLID).toHaveBeenCalledWith("999@lid");
expect(getPNForLID).toHaveBeenCalledTimes(1);
const inbound = inboundMessage(onMessage);
expect(inbound.payload.body).toBe("ping");
expect(inbound.admission?.conversation.id).toBe("+999");
expect(inbound.platform.recipientJid).toBe("+123");
await listener.close();
});
it("resolves LID JIDs via authDir mapping files", async () => {
const onMessage = vi.fn(async () => {});
fsSync.writeFileSync(
path.join(getAuthDir(), "lid-mapping-555_reverse.json"),
JSON.stringify("1555"),
);
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
const getPNForLID = vi.spyOn(sock.signalRepository.lidMapping, "getPNForLID");
const upsert = buildNotifyMessageUpsert({
id: nextMessageId("lid-authdir"),
remoteJid: "555@lid",
text: "ping",
timestamp: 1_700_000_000,
pushName: "Tester",
});
sock.ev.emit("messages.upsert", upsert);
await waitForMessageCalls(onMessage, 1);
const inbound = inboundMessage(onMessage);
expect(inbound.payload.body).toBe("ping");
expect(inbound.admission?.conversation.id).toBe("+1555");
expect(inbound.platform.recipientJid).toBe("+123");
expect(getPNForLID).not.toHaveBeenCalled();
await listener.close();
});
it("resolves group participant LID JIDs via Baileys mapping", async () => {
const onMessage = vi.fn(async () => {});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
const getPNForLID = vi.spyOn(sock.signalRepository.lidMapping, "getPNForLID");
sock.signalRepository.lidMapping.getPNForLID.mockResolvedValueOnce("444:0@s.whatsapp.net");
const upsert = buildNotifyMessageUpsert({
id: nextMessageId("group-lid"),
remoteJid: "123@g.us",
participant: "444@lid",
text: "ping",
timestamp: 1_700_000_000,
});
sock.ev.emit("messages.upsert", upsert);
await waitForMessageCalls(onMessage, 1);
expect(getPNForLID).toHaveBeenCalledWith("444@lid");
const inbound = inboundMessage(onMessage);
expect(inbound.payload.body).toBe("ping");
expect(inbound.admission?.conversation.id).toBe("123@g.us");
expect(inbound.platform.senderE164).toBe("+444");
expect(inbound.admission?.conversation.kind).toBe("group");
await listener.close();
});
it("captures reply context from quoted messages", async () => {
await expectQuotedReplyContext({ conversation: "original" });
});
it("preserves native reply context when WhatsApp omits the quoted message", async () => {
const onMessage = vi.fn(async () => {});
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
sock.ev.emit("messages.upsert", {
type: "notify",
messages: [
{
key: {
id: nextMessageId("quoted-unavailable"),
fromMe: false,
remoteJid: "999@s.whatsapp.net",
},
message: {
extendedTextMessage: {
text: "yes",
contextInfo: {
stanzaId: "original-message",
participant: "111@s.whatsapp.net",
},
},
},
messageTimestamp: 1_700_000_000,
pushName: "Tester",
},
],
});
await waitForMessageCalls(onMessage, 1);
const inbound = inboundMessage(onMessage);
expect(inbound.payload.body).toBe("yes");
expect(inbound.quote).toMatchObject({
id: "original-message",
body: "[quoted message unavailable]",
sender: { displayName: "+111", jid: "111@s.whatsapp.net", e164: "+111" },
});
await listener.close();
});
it("captures reply context from wrapped quoted messages", async () => {
await expectQuotedReplyContext({
viewOnceMessageV2Extension: {
message: { conversation: "original" },
},
});
});
it("captures reply context from botInvokeMessage wrapped quoted messages", async () => {
await expectQuotedReplyContext({
botInvokeMessage: {
message: { conversation: "original" },
},
});
});
it("captures reply context from groupMentionedMessage wrapped quoted messages", async () => {
await expectQuotedReplyContext({
groupMentionedMessage: {
message: { conversation: "original" },
},
});
});
});
@@ -0,0 +1,722 @@
// WhatsApp monitor inbox behavior split by ownership.
import { defaultRuntime } from "openclaw/plugin-sdk/runtime-env";
import { describe, expect, it, vi } from "vitest";
import {
controllerContexts,
sleepWithAbortMock,
nextMessageId,
createSocketRef,
fastReconnectPolicy,
inboundMessage,
expectSocketOperationTimeout,
createBaileysCacheSupport,
primeInboundReplyHandle,
installStreamsInboundMessageHooks,
} from "./monitor-inbox.streams-inbound-messages.test-support.js";
import {
buildNotifyMessageUpsert,
DEFAULT_ACCOUNT_ID,
settleInboundWork,
startInboxMonitor,
waitForMessageCalls,
type InboxMonitorOptions,
type InboxOnMessage,
} from "./monitor-inbox.test-harness.js";
import { lookupInboundMessageMeta } from "./quoted-message.js";
import { DEFAULT_WHATSAPP_SOCKET_TIMING } from "./socket-timing.js";
describe("web monitor inbox socket lifecycle", () => {
installStreamsInboundMessageHooks();
it("socket session stays unavailable on connect in self-chat mode", async () => {
const { listener, sock } = await startInboxMonitor(vi.fn(async () => {}) as InboxOnMessage, {
selfChatMode: true,
});
expect(sock.sendPresenceUpdate).toHaveBeenNthCalledWith(1, "unavailable");
await listener.close();
});
it("socket session uses a replacement socket for replies created before reconnect", async () => {
const onMessage = vi.fn(async () => undefined);
const socketRef: NonNullable<InboxMonitorOptions["socketRef"]> = { current: null };
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage, { socketRef });
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: nextMessageId("replacement-socket"),
remoteJid: "999@s.whatsapp.net",
text: "ping",
timestamp: 1_700_000_000,
pushName: "Tester",
}),
);
await waitForMessageCalls(onMessage, 1);
const inbound = inboundMessage(onMessage);
const replacementSock = {
sendMessage: vi.fn(async () => undefined),
sendPresenceUpdate: vi.fn(async () => undefined),
};
socketRef.current = replacementSock as unknown as NonNullable<
InboxMonitorOptions["socketRef"]
>["current"];
await inbound.platform.reply("pong");
await inbound.platform.sendMedia({ text: "after-reconnect" });
await inbound.platform.sendComposing();
expect(replacementSock.sendMessage).toHaveBeenNthCalledWith(1, "999@s.whatsapp.net", {
text: "pong",
});
expect(replacementSock.sendMessage).toHaveBeenNthCalledWith(2, "999@s.whatsapp.net", {
text: "after-reconnect",
});
expect(replacementSock.sendPresenceUpdate).toHaveBeenCalledWith(
"composing",
"999@s.whatsapp.net",
);
expect(sock.sendMessage).not.toHaveBeenCalled();
await listener.close();
});
it("socket session waits for a replacement socket before sending replies", async () => {
const onMessage = vi.fn(async () => undefined);
const socketRef = createSocketRef();
const { listener, sock, inbound } = await primeInboundReplyHandle({
onMessage,
socketRef,
upsertId: "reconnect-gap",
retryPolicy: {
initialMs: 10,
maxMs: 10,
factor: 1,
jitter: 0,
maxAttempts: 2,
},
});
const replacementSock = {
sendMessage: vi.fn(async () => undefined),
sendPresenceUpdate: vi.fn(async () => undefined),
};
socketRef.current = null;
sleepWithAbortMock.mockImplementationOnce(async () => {
socketRef.current = replacementSock as unknown as NonNullable<
InboxMonitorOptions["socketRef"]
>["current"];
});
await inbound?.platform.reply("pong");
expect(sleepWithAbortMock).toHaveBeenCalledWith(10, undefined);
expect(replacementSock.sendMessage).toHaveBeenCalledWith("999@s.whatsapp.net", {
text: "pong",
});
expect(sock.sendMessage).not.toHaveBeenCalled();
await listener.close();
});
it("socket session retries timed-out sends without clearing the socket ref", async () => {
const onMessage = vi.fn(async () => undefined);
const socketRef = createSocketRef();
const { listener, sock, inbound } = await primeInboundReplyHandle({
onMessage,
socketRef,
upsertId: "timeout-retry",
retryPolicy: fastReconnectPolicy(2),
});
sock.sendMessage
.mockRejectedValueOnce(new Error("operation timed out"))
.mockResolvedValueOnce({ key: { id: "after-timeout" } });
await inbound?.platform.reply("pong");
expect(sock.sendMessage).toHaveBeenNthCalledWith(1, "999@s.whatsapp.net", {
text: "pong",
});
expect(sock.sendMessage).toHaveBeenNthCalledWith(2, "999@s.whatsapp.net", {
text: "pong",
});
expect(socketRef.current).toBe(sock);
expect(sleepWithAbortMock).toHaveBeenCalledTimes(1);
await listener.close();
});
type ReachoutTimelockCase = {
name: string;
fetchedStates?: Array<"active" | "inactive">;
updatesBefore?: boolean[];
preflight?: boolean;
updateAfterPreflight?: boolean;
target?: string;
sends: Array<{ text: string; result: "resolve" | "reject" }>;
expectedFetches: number;
expectedSends: number;
};
const reachoutTimelockCases = [
{
name: "rejects direct sends while reachout timelock is active",
fetchedStates: ["active"],
sends: [{ text: "hello", result: "reject" }],
expectedFetches: 1,
expectedSends: 0,
},
{
name: "uses connection updates before direct sends",
updatesBefore: [true],
sends: [{ text: "hello", result: "reject" }],
expectedFetches: 0,
expectedSends: 0,
},
{
name: "allows direct sends after reachout timelock clears",
updatesBefore: [true, false],
sends: [{ text: "hello", result: "resolve" }],
expectedFetches: 1,
expectedSends: 1,
},
{
name: "refreshes inactive reachout state before later direct sends",
fetchedStates: ["inactive", "active"],
sends: [
{ text: "first", result: "resolve" },
{ text: "second", result: "reject" },
],
expectedFetches: 2,
expectedSends: 1,
},
{
name: "reuses readiness preflight for the immediate direct send",
fetchedStates: ["inactive"],
preflight: true,
sends: [{ text: "hello", result: "resolve" }],
expectedFetches: 1,
expectedSends: 1,
},
{
name: "invalidates readiness after an active timelock update",
fetchedStates: ["inactive"],
preflight: true,
updateAfterPreflight: true,
sends: [{ text: "hello", result: "reject" }],
expectedFetches: 1,
expectedSends: 0,
},
{
name: "does not apply account reachout timelock to group sends",
updatesBefore: [true],
target: "120363401234567890@g.us",
sends: [{ text: "hello", result: "resolve" }],
expectedFetches: 0,
expectedSends: 1,
},
] satisfies ReachoutTimelockCase[];
it.each(reachoutTimelockCases)("socket session $name", async (scenario) => {
const onMessage = vi.fn(async () => undefined);
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
for (const state of scenario.fetchedStates ?? []) {
sock.fetchAccountReachoutTimelock.mockResolvedValueOnce(
state === "active"
? {
isActive: true,
enforcementType: "WEB_COMPANION_ONLY",
timeEnforcementEnds: new Date(Date.now() + 60_000),
}
: { isActive: false },
);
}
for (const isActive of scenario.updatesBefore ?? []) {
sock.ev.emit("connection.update", {
reachoutTimeLock: {
isActive,
...(isActive ? { enforcementType: "WEB_COMPANION_ONLY" } : {}),
},
});
}
try {
const target = scenario.target ?? "+1555";
if (scenario.preflight) {
await listener.assertSendReady?.(target);
}
if (scenario.updateAfterPreflight) {
sock.ev.emit("connection.update", {
reachoutTimeLock: {
isActive: true,
enforcementType: "WEB_COMPANION_ONLY",
},
});
}
for (const send of scenario.sends) {
const sendPromise = listener.sendMessage(target, send.text);
if (send.result === "reject") {
await expect(sendPromise).rejects.toThrow("WhatsApp reachout timelock is active");
} else {
await expect(sendPromise).resolves.toBeDefined();
if (target.endsWith("@g.us")) {
expect(sock.sendMessage).toHaveBeenCalledWith(target, { text: send.text });
}
}
}
expect(sock.fetchAccountReachoutTimelock).toHaveBeenCalledTimes(scenario.expectedFetches);
expect(sock.sendMessage).toHaveBeenCalledTimes(scenario.expectedSends);
} finally {
await listener.close();
}
});
it("socket session blocks direct composing presence during reachout timelock", async () => {
const onMessage = vi.fn(async () => undefined);
const socketRef = createSocketRef();
const { listener, sock, inbound } = await primeInboundReplyHandle({
onMessage,
socketRef,
upsertId: "reachout-composing",
retryPolicy: fastReconnectPolicy(2),
});
sock.ev.emit("connection.update", {
reachoutTimeLock: {
isActive: true,
enforcementType: "WEB_COMPANION_ONLY",
},
});
sock.sendPresenceUpdate.mockClear();
try {
await inbound.platform.sendComposing();
expect(sock.fetchAccountReachoutTimelock).not.toHaveBeenCalled();
expect(sock.sendPresenceUpdate).not.toHaveBeenCalled();
} finally {
await listener.close();
}
});
it("socket session times out stalled sends at the Baileys query timeout", async () => {
const onMessage = vi.fn(async () => undefined);
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
vi.useFakeTimers();
try {
sock.sendMessage.mockImplementationOnce(() => new Promise(() => {}));
const sendPromise = listener.sendMessage("+1555", "hello");
await expectSocketOperationTimeout("sendMessage", sendPromise);
expect(vi.getTimerCount()).toBe(0);
expect(sock.sendMessage).toHaveBeenCalledTimes(1);
} finally {
vi.useRealTimers();
await listener.close();
}
});
it("socket session preserves the socket after a local send timeout", async () => {
const onMessage = vi.fn(async () => undefined);
const socketRef = createSocketRef();
const { listener, sock, inbound } = await primeInboundReplyHandle({
onMessage,
socketRef,
upsertId: "local-timeout-terminal",
retryPolicy: fastReconnectPolicy(2),
});
vi.useFakeTimers();
try {
sock.sendMessage.mockImplementationOnce(() => new Promise(() => {}));
const replyPromise = inbound.platform.reply("pong");
await expectSocketOperationTimeout("sendMessage", replyPromise);
expect(sock.sendMessage).toHaveBeenCalledTimes(1);
expect(socketRef.current).toBe(sock);
expect(sleepWithAbortMock).not.toHaveBeenCalled();
expect(vi.getTimerCount()).toBe(0);
} finally {
vi.useRealTimers();
await listener.close();
}
});
it("socket session records outbound replies for Baileys retry lookup", async () => {
const onMessage = vi.fn(async () => undefined);
const socketRef = createSocketRef();
const baileysCache = createBaileysCacheSupport();
const message = { conversation: "pong" };
const { listener, sock, inbound } = await primeInboundReplyHandle({
onMessage,
socketRef,
baileysCache,
upsertId: "outbound-retry-cache",
retryPolicy: fastReconnectPolicy(2),
});
sock.sendMessage.mockResolvedValueOnce({
key: { id: "outbound-cached" },
message,
});
await inbound.platform.reply("pong");
await expect(
baileysCache.socketOptions.getMessage({
id: "outbound-cached",
remoteJid: "999@s.whatsapp.net",
}),
).resolves.toBe(message);
expect(
lookupInboundMessageMeta(DEFAULT_ACCOUNT_ID, "999@s.whatsapp.net", "outbound-cached"),
).toMatchObject({ fromMe: true, body: "pong" });
await listener.close();
});
it("socket session suppresses self-echo after a late accepted send", async () => {
const onMessage = vi.fn(async () => undefined);
const socketRef = createSocketRef();
const baileysCache = createBaileysCacheSupport();
const { listener, sock, inbound } = await primeInboundReplyHandle({
onMessage,
socketRef,
baileysCache,
upsertId: "late-accept",
retryPolicy: fastReconnectPolicy(2),
});
const message = { conversation: "pong" };
let acceptLateSend:
| ((value: { key: { id: string }; message: { conversation: string } }) => void)
| undefined;
vi.useFakeTimers();
try {
sock.sendMessage.mockImplementationOnce(
async () =>
await new Promise((resolve) => {
acceptLateSend = resolve;
}),
);
const replyPromise = inbound.platform.reply("pong");
await expectSocketOperationTimeout("sendMessage", replyPromise);
} finally {
vi.useRealTimers();
}
acceptLateSend?.({ key: { id: "late-accepted" }, message });
await settleInboundWork();
await expect(
baileysCache.socketOptions.getMessage({
id: "late-accepted",
remoteJid: "999@s.whatsapp.net",
}),
).resolves.toBe(message);
sock.ev.emit("messages.upsert", {
type: "notify",
messages: [
{
key: {
id: "late-accepted",
fromMe: true,
remoteJid: "999@s.whatsapp.net",
},
message: { conversation: "pong" },
messageTimestamp: 1_700_000_001,
pushName: "Tester",
},
],
});
await settleInboundWork();
expect(onMessage).toHaveBeenCalledTimes(1);
expect(socketRef.current).toBe(sock);
expect(sleepWithAbortMock).not.toHaveBeenCalled();
await listener.close();
});
it("socket session times out stalled send-api presence updates", async () => {
const onMessage = vi.fn(async () => undefined);
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage);
vi.useFakeTimers();
try {
sock.sendPresenceUpdate.mockClear();
sock.sendPresenceUpdate.mockImplementationOnce(() => new Promise(() => {}));
const presencePromise = listener.sendComposingTo("+1555");
await expectSocketOperationTimeout("sendPresenceUpdate", presencePromise);
expect(sock.sendPresenceUpdate).toHaveBeenCalledTimes(1);
expect(vi.getTimerCount()).toBe(0);
} finally {
vi.useRealTimers();
await listener.close();
}
});
it("socket session bounds stalled read-receipt operations", async () => {
const onMessage = vi.fn(async () => undefined);
const logSpy = vi.spyOn(defaultRuntime, "log").mockImplementation(() => undefined);
const { listener, sock } = await startInboxMonitor(onMessage as InboxOnMessage, {
verbose: true,
});
vi.useFakeTimers();
try {
const messageId = nextMessageId("read-receipt-timeout");
// A WhatsApp socket whose read-receipt acknowledgement never resolves
// (e.g. a stalled Baileys privacy IQ query) would otherwise hang the
// inbound delivery pipeline forever.
sock.readMessages.mockImplementationOnce(() => new Promise(() => {}));
sock.ev.emit(
"messages.upsert",
buildNotifyMessageUpsert({
id: messageId,
remoteJid: "999@s.whatsapp.net",
text: "ping",
timestamp: 1_700_000_000,
pushName: "Tester",
}),
);
await waitForMessageCalls(onMessage, 1);
// The read receipt is attempted on the stalled socket...
await vi.waitFor(() => {
expect(sock.readMessages).toHaveBeenCalledWith([
{ remoteJid: "999@s.whatsapp.net", id: messageId, participant: undefined, fromMe: false },
]);
});
// ...and is bounded by the socket operation timeout rather than hanging.
await vi.advanceTimersByTimeAsync(DEFAULT_WHATSAPP_SOCKET_TIMING.defaultQueryTimeoutMs);
await vi.waitFor(() => {
const loggedTimeoutFailure = logSpy.mock.calls.some(
([message]) =>
typeof message === "string" &&
message.includes(`Failed to mark message ${messageId} read`) &&
message.includes("readMessages timed out"),
);
expect(loggedTimeoutFailure).toBe(true);
});
expect(vi.getTimerCount()).toBe(0);
} finally {
vi.useRealTimers();
await listener.close();
logSpy.mockRestore();
}
});
it("socket session bounds reconnect-gap retries when attempts are unlimited", async () => {
const onMessage = vi.fn(async () => undefined);
const socketRef = createSocketRef();
const { listener, inbound } = await primeInboundReplyHandle({
onMessage,
socketRef,
upsertId: "unlimited-reconnect-send-bound",
retryPolicy: fastReconnectPolicy(0),
useCurrentSock: true,
});
socketRef.current = null;
await expect(inbound?.platform.reply("pong")).rejects.toThrow(
"no active socket - reconnection in progress",
);
expect(sleepWithAbortMock).toHaveBeenCalledTimes(11);
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("socket session routes captured replies through a matching successor", 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: fastReconnectPolicy(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);
// === Simulate health-monitor-driven shutdown of controller A ===
socketRefA.current = null;
aShouldRetryDisconnect = false;
// The mock harness socket exposes user.id = "123@s.whatsapp.net"; the
// successor must report an overlapping identity for the handoff to succeed.
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;
controllerContexts.set(DEFAULT_ACCOUNT_ID, handleB);
try {
await inbound.reply("pong");
await inbound.sendMedia({ text: "media after restart" });
// Captured A reply routed through B via the runtime context.
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 {
controllerContexts.delete(DEFAULT_ACCOUNT_ID);
await listenerA.close();
}
});
// VE-513 PN/LID normalization: Baileys sockets may expose self identity in
// different forms across reconnects (PN JID `<n>@s.whatsapp.net` vs LID
// `<x>@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("socket session accepts a successor with equivalent LID and phone identity", 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: fastReconnectPolicy(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);
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;
controllerContexts.set(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 {
controllerContexts.delete(DEFAULT_ACCOUNT_ID);
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("socket session refuses a successor with a mismatched self identity", 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: fastReconnectPolicy(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);
// 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;
controllerContexts.set(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 {
controllerContexts.delete(DEFAULT_ACCOUNT_ID);
await listenerA.close();
}
});
});
File diff suppressed because it is too large Load Diff