diff --git a/docs/.generated/plugin-sdk-api-baseline.sha256 b/docs/.generated/plugin-sdk-api-baseline.sha256 index 9626a9b51d6f..f5db93d0ff71 100644 --- a/docs/.generated/plugin-sdk-api-baseline.sha256 +++ b/docs/.generated/plugin-sdk-api-baseline.sha256 @@ -29,7 +29,7 @@ c0f910ebfa3dbf283145fb1e3b9c016d03e853ef13e70402b09f3b9d9c2f4ab0 module/channel 9a5aaf650f9242523bb57bdc2556c323ab64e55aa11e25e2685b73b23ee12534 module/channel-dm-policy fbf353eb38ae68d8ded3f2a60b432c7bb2c245d2ec7e7c9f53c6da19a0db0938 module/channel-entry-contract 982f29a18e07228e3da82cae67d06ff38249592a29c2fd28f01f0d2016ff80d9 module/channel-feedback -d645d24bcb7a5f68cc46c692ad0d1fbd19be0a99996f31e9479ce9cffce301c1 module/channel-inbound +6cf2da836db49a267aed470ce6f5f55e68d8d86e1efe980cdc71349cfce30bc2 module/channel-inbound 76bb7f531f3702c801e8fe7479e9e499f601fb361a4303afdcb45fc0da440e4b module/channel-inbound-debounce 4a7ada095f0f483525dcbd848fbccab26473749eba87e6a6c6e5074fd04ee1d1 module/channel-ingress-runtime c97dd36cdf8f83c2893c33e9430a93cd131a03d855725783ca5b545de0cf84f8 module/channel-lifecycle diff --git a/docs/plugins/sdk-channel-inbound.md b/docs/plugins/sdk-channel-inbound.md index d94a65fdffde..f0cebe56fcb3 100644 --- a/docs/plugins/sdk-channel-inbound.md +++ b/docs/plugins/sdk-channel-inbound.md @@ -38,6 +38,78 @@ import { - `dispatchChannelInboundReply(...)`: records and dispatches an already assembled inbound reply with a delivery adapter. +## Prepared inbound envelope + +Channels that finish transport admission and normalization before entering the +shared runner can use `PreparedChannelInbound` as the boundary: + +```ts +import { + projectPreparedChannelInbound, + resolveChannelInboundReplyPolicy, + runChannelInboundEvent, + type PreparedChannelInbound, +} from "openclaw/plugin-sdk/channel-inbound"; + +const inbound = { + channel: "demo", + event: { id: nativeEvent.id, timestamp: nativeEvent.timestamp }, + from: conversation.id, + sender: { id: sender.id, name: sender.name }, + conversation: { kind: "group", id: conversation.id }, + route: { agentId, accountId, routeSessionKey }, + reply: { to: conversation.id, replyToId: quote?.id }, + message: { + rawBody, + body: formattedBody, + bodyForAgent, + commandBody, + }, + command: { + kind: "text-slash", + body: commandBody, + authorization: { kind: "authorized" }, + }, + media, +} satisfies PreparedChannelInbound; + +const prepared = projectPreparedChannelInbound({ + inbound, + control: { messageReceivedHooks: "core" }, +}); +const replyPolicy = resolveChannelInboundReplyPolicy({ + cfg, + ctx: prepared.context, + blockStreamingEnabled, +}); + +await runChannelInboundEvent({ + channel: inbound.channel, + accountId, + raw: inbound, + adapter: { + ingest: () => prepared.input, + resolveTurn: () => buildTurnPlan(prepared.context, replyPolicy), + }, +}); +``` + +The envelope contains prepared product facts only. Keep native callbacks, +socket objects, provider message records, and transport lifecycle symbols in a +channel-private context beside it. Represent command authorization explicitly +as `not_checked`, `authorized`, or `denied`; do not make core infer it from a +transport payload. + +Use `resolveChannelInboundReplyPolicy(...)` after building the portable inbound +context. Core resolves product policy for source reply delivery, block +streaming, and typing suppression. The channel supplies prepared context facts, +including mention state, plus native capabilities such as whether block +streaming is enabled. + +Keep quote encoding, reply callbacks, native target lookup, and other +transport mechanics in the channel. Do not duplicate visible-reply policy in a +channel adapter. + For media-only inbound events, keep the message body and command text empty and pass one `ChannelInboundMediaInput` fact per native attachment. When an ambient history line or another text-only carrier must describe those facts, use diff --git a/extensions/whatsapp/src/auto-reply.web-auto-reply.compresses-common-formats-jpeg-cap.test.ts b/extensions/whatsapp/src/auto-reply.web-auto-reply.compresses-common-formats-jpeg-cap.test.ts index da1bc85ffaee..aaeaf3c04a38 100644 --- a/extensions/whatsapp/src/auto-reply.web-auto-reply.compresses-common-formats-jpeg-cap.test.ts +++ b/extensions/whatsapp/src/auto-reply.web-auto-reply.compresses-common-formats-jpeg-cap.test.ts @@ -10,7 +10,7 @@ import { resetLoadConfigMock, setLoadConfigMock, } from "./auto-reply.test-harness.js"; -import type { WebInboundCallbackMessage, WebInboundMessageInput } from "./inbound.js"; +import type { WebInboundCallbackMessage } from "./inbound.js"; import { createTestWebInboundMessage } from "./inbound/test-message.test-helper.js"; installWebAutoReplyTestHomeHooks(); @@ -45,7 +45,7 @@ describe("web auto-reply", () => { const sendMedia = params.sendMedia ?? spies.sendMedia; const resolver = vi.fn().mockResolvedValue(params.resolverValue); - let capturedOnMessage: ((msg: WebInboundMessageInput) => Promise) | undefined; + let capturedOnMessage: Parameters[0]["onMessage"] | undefined; const listenerFactory: ListenerFactory = async ({ onMessage }) => { capturedOnMessage = onMessage; return createMockWebListener(); diff --git a/extensions/whatsapp/src/auto-reply/deliver-reply.filename-media.test.ts b/extensions/whatsapp/src/auto-reply/deliver-reply.filename-media.test.ts index 7361724acd74..a722dc4d6481 100644 --- a/extensions/whatsapp/src/auto-reply/deliver-reply.filename-media.test.ts +++ b/extensions/whatsapp/src/auto-reply/deliver-reply.filename-media.test.ts @@ -4,7 +4,7 @@ import { createAcceptedWhatsAppSendResult } from "../inbound/send-result.test-he import { createTestWebInboundMessage } from "../inbound/test-message.test-helper.js"; import type { AdmittedWebInboundMessage } from "../inbound/types.js"; import { loadWebMedia } from "../media.js"; -import { deliverWebReply } from "./deliver-reply.js"; +import { createWhatsAppReplyTransportContext, deliverWebReply } from "./deliver-reply.js"; const hoisted = vi.hoisted(() => ({ transcodeAudioBufferToOpus: vi.fn(), @@ -161,7 +161,7 @@ describe("WhatsApp filename-only media delivery", () => { await deliverWebReply({ replyResult: { text: "caption", mediaUrl: "https://example.com/download" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger: { info: vi.fn(), warn: vi.fn() }, diff --git a/extensions/whatsapp/src/auto-reply/deliver-reply.test.ts b/extensions/whatsapp/src/auto-reply/deliver-reply.test.ts index 75ad797fd759..bfc721a0b923 100644 --- a/extensions/whatsapp/src/auto-reply/deliver-reply.test.ts +++ b/extensions/whatsapp/src/auto-reply/deliver-reply.test.ts @@ -43,6 +43,7 @@ vi.mock("../media.js", () => ({ })); let deliverWebReply: typeof import("./deliver-reply.js").deliverWebReply; +let createWhatsAppReplyTransportContext: typeof import("./deliver-reply.js").createWhatsAppReplyTransportContext; let whatsappOutbound: typeof import("../outbound-adapter.js").whatsappOutbound; function unacceptedSendResult(kind: "media" | "text") { @@ -217,7 +218,7 @@ async function expectReplySuppressed(replyResult: { text: string; isReasoning?: const msg = makeMsg(); await deliverWebReply({ replyResult, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -229,7 +230,7 @@ async function expectReplySuppressed(replyResult: { text: string; isReasoning?: describe("deliverWebReply", () => { beforeAll(async () => { - ({ deliverWebReply } = await import("./deliver-reply.js")); + ({ createWhatsAppReplyTransportContext, deliverWebReply } = await import("./deliver-reply.js")); ({ whatsappOutbound } = await import("../outbound-adapter.js")); }); @@ -250,7 +251,7 @@ describe("deliverWebReply", () => { await deliverWebReply({ replyResult: { text: "Intro line\nReasoning: appears in content but is not a prefix" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -269,7 +270,7 @@ describe("deliverWebReply", () => { const delivery = await deliverWebReply({ replyResult: { text: "aaaaaa" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 3, replyLogger, @@ -295,7 +296,7 @@ describe("deliverWebReply", () => { const delivery = await deliverWebReply({ replyResult: { text: "hello" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -319,7 +320,7 @@ describe("deliverWebReply", () => { replyResult: { text: 'Before\nx\nAfter', }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 4000, replyLogger, @@ -349,7 +350,7 @@ describe("deliverWebReply", () => { "
After
", ].join("\n"), }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 4000, replyLogger, @@ -371,7 +372,7 @@ describe("deliverWebReply", () => { "After", ].join("\n"), }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 4000, replyLogger, @@ -392,7 +393,7 @@ describe("deliverWebReply", () => { await deliverWebReply({ replyResult: { text: "aaaaaa", replyToId: "reply-1" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 3, replyLogger, @@ -426,7 +427,7 @@ describe("deliverWebReply", () => { await runWithFakeTimers(() => deliverWebReply({ replyResult: { text: "hi" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -446,7 +447,7 @@ describe("deliverWebReply", () => { await runWithFakeTimers(() => deliverWebReply({ replyResult: { text: "hi" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -467,7 +468,7 @@ describe("deliverWebReply", () => { await expect( deliverWebReply({ replyResult: { text: "hi" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -485,7 +486,7 @@ describe("deliverWebReply", () => { await deliverWebReply({ replyResult: { text: "aaaaaa", mediaUrl: "http://example.com/img.jpg" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), mediaLocalRoots, maxMediaBytes: 1024 * 1024, textLimit: 3, @@ -520,7 +521,7 @@ describe("deliverWebReply", () => { await expect( deliverWebReply({ replyResult: { text: "captiontail", mediaUrl: "http://example.com/img.jpg" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 7, replyLogger, @@ -540,7 +541,7 @@ describe("deliverWebReply", () => { await deliverWebReply({ replyResult: { text: "\n \n indented block" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -566,7 +567,7 @@ describe("deliverWebReply", () => { mediaUrl: "http://example.com/img.jpg", replyToId: "reply-2", }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 7, replyLogger, @@ -606,7 +607,7 @@ describe("deliverWebReply", () => { await runWithFakeTimers(() => deliverWebReply({ replyResult: { text: "caption", mediaUrl: "http://example.com/img.jpg" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -624,7 +625,7 @@ describe("deliverWebReply", () => { await deliverWebReply({ replyResult: { text: "caption", mediaUrl: "http://example.com/img.jpg" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 20, replyLogger, @@ -649,7 +650,7 @@ describe("deliverWebReply", () => { await deliverWebReply({ replyResult: { text: "ALPHALINEBRAVOLINE", mediaUrl: "http://example.com/img.jpg" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 9, replyLogger, @@ -696,7 +697,7 @@ describe("deliverWebReply", () => { text: "caption", mediaUrls: ["http://example.com/bad.jpg", "http://example.com/good.pdf"], }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -746,7 +747,7 @@ describe("deliverWebReply", () => { await deliverWebReply({ replyResult: { text: "caption", mediaUrl, mediaUrls }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -792,7 +793,7 @@ describe("deliverWebReply", () => { text: "caption", mediaUrls: ["http://example.com/img1.jpg", "http://example.com/img2.jpg"], }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -856,7 +857,7 @@ describe("deliverWebReply", () => { await deliverWebReply({ replyResult: payload, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -906,7 +907,7 @@ describe("deliverWebReply", () => { await deliverWebReply({ replyResult: { text: "cap", mediaUrl: "http://example.com/a.ogg" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -940,7 +941,7 @@ describe("deliverWebReply", () => { await deliverWebReply({ replyResult: { text: "cap", mediaUrl: "http://example.com/a.mp3" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -981,7 +982,7 @@ describe("deliverWebReply", () => { await deliverWebReply({ replyResult: { text: "cap", mediaUrl: "http://example.com/v.mp4" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -1011,7 +1012,7 @@ describe("deliverWebReply", () => { await deliverWebReply({ replyResult: { text: "cap", mediaUrl: "http://example.com/x.bin" }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, @@ -1044,7 +1045,7 @@ describe("deliverWebReply", () => { text: "cap", mediaUrl: "https://example.com/report.pdf?X-Amz-Signature=secret#frag", }, - msg, + transport: createWhatsAppReplyTransportContext(msg), maxMediaBytes: 1024 * 1024, textLimit: 200, replyLogger, diff --git a/extensions/whatsapp/src/auto-reply/deliver-reply.ts b/extensions/whatsapp/src/auto-reply/deliver-reply.ts index f7f2bff18b17..8fbdf1418ef7 100644 --- a/extensions/whatsapp/src/auto-reply/deliver-reply.ts +++ b/extensions/whatsapp/src/auto-reply/deliver-reply.ts @@ -37,6 +37,35 @@ export type WhatsAppReplyDeliveryResult = { providerAccepted: boolean; }; +export type WhatsAppReplyTransportContext = { + accountId: string; + conversationId: string; + conversationKind: "direct" | "group"; + chatJid: string; + senderJid?: string; + recipientJid: string; + correlationId?: string; + reply: AdmittedWebInboundMessage["platform"]["reply"]; + sendMedia: AdmittedWebInboundMessage["platform"]["sendMedia"]; +}; + +export function createWhatsAppReplyTransportContext( + msg: AdmittedWebInboundMessage, +): WhatsAppReplyTransportContext { + const admission = requireWhatsAppInboundAdmission(msg); + return { + accountId: admission.accountId, + conversationId: admission.conversation.id, + conversationKind: admission.conversation.kind, + chatJid: msg.platform.chatJid, + senderJid: msg.platform.senderJid, + recipientJid: msg.platform.recipientJid, + correlationId: msg.event.id, + reply: msg.platform.reply, + sendMedia: msg.platform.sendMedia, + }; +} + function resolveWhatsAppReceiptKind( results: readonly WhatsAppSendResult[], ): Parameters[0]["kind"] { @@ -87,7 +116,7 @@ function createWhatsAppReplyDeliveryReceipt( export async function deliverWebReply(params: { replyResult: ReplyPayload; normalizedReplyResult?: DeliverableWhatsAppOutboundPayload; - msg: AdmittedWebInboundMessage; + transport: WhatsAppReplyTransportContext; mediaLocalRoots?: readonly string[]; maxMediaBytes: number; textLimit: number; @@ -100,10 +129,10 @@ export async function deliverWebReply(params: { skipLog?: boolean; tableMode?: MarkdownTableMode; }): Promise { - const { replyResult, msg, maxMediaBytes, textLimit, replyLogger, connectionId, skipLog } = params; - const admission = requireWhatsAppInboundAdmission(msg); - const conversationId = admission.conversation.id; - const isGroupConversation = admission.conversation.kind === "group"; + const { replyResult, transport, maxMediaBytes, textLimit, replyLogger, connectionId, skipLog } = + params; + const conversationId = transport.conversationId; + const isGroupConversation = transport.conversationKind === "group"; const replyStarted = Date.now(); const sendResults: WhatsAppSendResult[] = []; const rememberSendResult = (result: WhatsAppSendResult | undefined) => { @@ -146,16 +175,15 @@ export async function deliverWebReply(params: { // per-message target. Look up cached metadata for the specific // message being quoted — msg.payload.body may be a combined batch body. const cached = lookupInboundMessageMeta( - admission.accountId, - msg.platform.chatJid, + transport.accountId, + transport.chatJid, replyResult.replyToId, ); return buildQuotedMessageOptions({ messageId: replyResult.replyToId, - remoteJid: msg.platform.chatJid, + remoteJid: transport.chatJid, fromMe: cached?.fromMe ?? false, - participant: - cached?.participant ?? (isGroupConversation ? msg.platform.senderJid : undefined), + participant: cached?.participant ?? (isGroupConversation ? transport.senderJid : undefined), messageText: cached?.body ?? "", media: cached?.media, }); @@ -185,7 +213,7 @@ export async function deliverWebReply(params: { for (const [index, chunk] of textChunks.entries()) { const chunkStarted = Date.now(); const quote = getQuote(); - rememberSendResult(await sendWithRetry(() => msg.platform.reply(chunk, quote), "text")); + rememberSendResult(await sendWithRetry(() => transport.reply(chunk, quote), "text")); if (!skipLog) { const durationMs = Date.now() - chunkStarted; whatsappOutboundLog.debug( @@ -195,10 +223,10 @@ export async function deliverWebReply(params: { } const delivery = finishDelivery(); const logPayload = { - correlationId: msg.event.id ?? newConnectionId(), + correlationId: transport.correlationId ?? newConnectionId(), connectionId: connectionId ?? null, to: conversationId, - from: msg.platform.recipientJid, + from: transport.recipientJid, text: elide(replyResult.text, 240), mediaUrl: null, mediaSizeBytes: null, @@ -239,7 +267,7 @@ export async function deliverWebReply(params: { rememberSendResult( await sendWithRetry( () => - msg.platform.sendMedia( + transport.sendMedia( { image: media.buffer, caption, @@ -255,7 +283,7 @@ export async function deliverWebReply(params: { rememberSendResult( await sendWithRetry( () => - msg.platform.sendMedia( + transport.sendMedia( { audio: media.buffer, ptt: true, @@ -268,7 +296,7 @@ export async function deliverWebReply(params: { ); if (caption) { rememberSendResult( - await sendWithRetry(() => msg.platform.reply(caption, quote), "media:audio-text"), + await sendWithRetry(() => transport.reply(caption, quote), "media:audio-text"), ); } } else if (media.kind === "video") { @@ -276,7 +304,7 @@ export async function deliverWebReply(params: { rememberSendResult( await sendWithRetry( () => - msg.platform.sendMedia( + transport.sendMedia( { video: media.buffer, caption, @@ -292,7 +320,7 @@ export async function deliverWebReply(params: { rememberSendResult( await sendWithRetry( () => - msg.platform.sendMedia( + transport.sendMedia( { document: media.buffer, fileName: media.fileName, @@ -310,10 +338,10 @@ export async function deliverWebReply(params: { ); replyLogger.info( { - correlationId: msg.event.id ?? newConnectionId(), + correlationId: transport.correlationId ?? newConnectionId(), connectionId: connectionId ?? null, to: conversationId, - from: msg.platform.recipientJid, + from: transport.recipientJid, text: caption ?? null, mediaUrl, mediaSizeBytes: media.buffer.length, @@ -334,7 +362,7 @@ export async function deliverWebReply(params: { whatsappOutboundLog.warn(`Trailing media failed; sent warning to ${conversationId}`); rememberSendResult( await sendWithRetry( - () => msg.platform.reply("⚠️ Media unavailable.", getQuote()), + () => transport.reply("⚠️ Media unavailable.", getQuote()), "media:fallback-unavailable", ), ); @@ -348,19 +376,14 @@ export async function deliverWebReply(params: { } whatsappOutboundLog.warn(`Media skipped; sent text-only to ${conversationId}`); rememberSendResult( - await sendWithRetry( - () => msg.platform.reply(fallbackText, getQuote()), - "media:fallback-text", - ), + await sendWithRetry(() => transport.reply(fallbackText, getQuote()), "media:fallback-text"), ); }, }); // Remaining text chunks after media for (const chunk of remainingText) { - rememberSendResult( - await sendWithRetry(() => msg.platform.reply(chunk, getQuote()), "media:text"), - ); + rememberSendResult(await sendWithRetry(() => transport.reply(chunk, getQuote()), "media:text")); } return finishDelivery(); } diff --git a/extensions/whatsapp/src/auto-reply/monitor.ts b/extensions/whatsapp/src/auto-reply/monitor.ts index ce0dbb606a58..9bf4b095859e 100644 --- a/extensions/whatsapp/src/auto-reply/monitor.ts +++ b/extensions/whatsapp/src/auto-reply/monitor.ts @@ -33,9 +33,8 @@ import { type WhatsAppBaileysMessageCache, } from "../inbound/baileys-cache.js"; import type { WhatsAppGroupMetadataCache } from "../inbound/group-metadata-cache.js"; -import { normalizeWebInboundMessage } from "../inbound/message-aliases.js"; import { attachWebInboxToSocket } from "../inbound/monitor.js"; -import type { WebInboundMessageInput } from "../inbound/types.js"; +import type { AdmittedWebInboundMessage } from "../inbound/types.js"; import { newConnectionId, resolveHeartbeatSeconds, @@ -230,17 +229,12 @@ export async function monitorWebChannel( cfg, channel: "whatsapp", }); - const shouldDebounce = (msg: WebInboundMessageInput) => { - const normalized = normalizeWebInboundMessage(msg); + const shouldDebounce = (msg: AdmittedWebInboundMessage) => { return shouldDebounceTextInbound({ - text: normalized.payload.commandBody ?? normalized.payload.body, + text: msg.payload.commandBody ?? msg.payload.body, cfg, - hasMedia: Boolean(normalized.payload.media?.path || normalized.payload.media?.type), - allowDebounce: !( - normalized.payload.location || - normalized.quote?.id || - normalized.quote?.body - ), + hasMedia: Boolean(msg.payload.media?.path || msg.payload.media?.type), + allowDebounce: !(msg.payload.location || msg.quote?.id || msg.quote?.body), }); }; @@ -299,12 +293,11 @@ export async function monitorWebChannel( groupMetadataCache, recentMessageKeys, baileysGroupMetaCache, - onMessage: async (msg: WebInboundMessageInput) => { - const normalized = normalizeWebInboundMessage(msg); + onMessage: async (msg: AdmittedWebInboundMessage) => { const inboundAt = Date.now(); controller.noteInbound(inboundAt); statusController.noteInbound(inboundAt); - await onMessage(normalized); + await onMessage(msg); }, onPendingWorkChanged: (pendingWorkCount, at) => { statusController.noteBusy(pendingWorkCount > 0, at); diff --git a/extensions/whatsapp/src/auto-reply/monitor/inbound-dispatch.runtime.ts b/extensions/whatsapp/src/auto-reply/monitor/inbound-dispatch.runtime.ts index 6f2920900a6c..b7c82c905353 100644 --- a/extensions/whatsapp/src/auto-reply/monitor/inbound-dispatch.runtime.ts +++ b/extensions/whatsapp/src/auto-reply/monitor/inbound-dispatch.runtime.ts @@ -4,7 +4,6 @@ export { getAgentScopedMediaLocalRoots, jidToE164, logVerbose, - resolveChannelMessageSourceReplyDeliveryMode, resolveChunkMode, resolveIdentityNamePrefix, resolveInboundLastRouteSessionKey, @@ -12,7 +11,6 @@ export { resolveSendableOutboundReplyParts, resolveTextChunkLimit, shouldLogVerbose, - toLocationContext, type getChildLogger, type getReplyFromConfig, type LoadConfigFn, diff --git a/extensions/whatsapp/src/auto-reply/monitor/inbound-dispatch.test.ts b/extensions/whatsapp/src/auto-reply/monitor/inbound-dispatch.test.ts index 3d1df0ff9038..89e2fef8bcd7 100644 --- a/extensions/whatsapp/src/auto-reply/monitor/inbound-dispatch.test.ts +++ b/extensions/whatsapp/src/auto-reply/monitor/inbound-dispatch.test.ts @@ -54,6 +54,53 @@ vi.mock("openclaw/plugin-sdk/channel-outbound", async (importOriginal) => { }; }); +vi.mock("openclaw/plugin-sdk/channel-inbound", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + resolveChannelInboundReplyPolicy: (params: { + cfg: { + messages?: { + visibleReplies?: "automatic" | "message_tool"; + groupChat?: { visibleReplies?: "automatic" | "message_tool" }; + }; + }; + ctx: { + ChatType?: string; + CommandSource?: "native" | "text"; + CommandAuthorized?: boolean; + WasMentioned?: boolean; + }; + blockStreamingEnabled?: boolean; + }) => { + sourceReplyDeliveryModeContexts.push(params.ctx); + const isRoom = params.ctx.ChatType === "group" || params.ctx.ChatType === "channel"; + const sourceReplyDeliveryMode = !isRoom + ? undefined + : params.ctx.CommandSource === "native" || + (params.ctx.CommandSource === "text" && params.ctx.CommandAuthorized === true) + ? "automatic" + : (params.cfg.messages?.groupChat?.visibleReplies ?? + params.cfg.messages?.visibleReplies) === "automatic" + ? "automatic" + : "message_tool_only"; + const sourceRepliesAreToolOnly = sourceReplyDeliveryMode === "message_tool_only"; + return { + sourceReplyDeliveryMode, + disableBlockStreaming: sourceRepliesAreToolOnly + ? true + : typeof params.blockStreamingEnabled === "boolean" + ? !params.blockStreamingEnabled + : undefined, + suppressTyping: + sourceRepliesAreToolOnly && + params.ctx.ChatType === "group" && + params.ctx.WasMentioned !== true, + }; + }, + }; +}); + vi.mock("./runtime-api.js", async () => { return { dispatchReplyWithBufferedBlockDispatcher: dispatchReplyWithBufferedBlockDispatcherMock, @@ -63,32 +110,6 @@ vi.mock("./runtime-api.js", async () => { return phone ? `+${phone}` : null; }, logVerbose: () => {}, - resolveChannelMessageSourceReplyDeliveryMode: ({ - cfg, - ctx, - }: { - cfg: { - messages?: { - visibleReplies?: "automatic" | "message_tool"; - groupChat?: { visibleReplies?: "automatic" | "message_tool" }; - }; - }; - ctx: { ChatType?: string; CommandSource?: "native" | "text"; CommandAuthorized?: boolean }; - }) => { - sourceReplyDeliveryModeContexts.push(ctx); - if ( - ctx.CommandSource === "native" || - (ctx.CommandSource === "text" && ctx.CommandAuthorized === true) - ) { - return "automatic"; - } - if (ctx.ChatType === "group" || ctx.ChatType === "channel") { - const configuredMode = - cfg.messages?.groupChat?.visibleReplies ?? cfg.messages?.visibleReplies; - return configuredMode === "automatic" ? "automatic" : "message_tool_only"; - } - return cfg.messages?.visibleReplies === "message_tool" ? "message_tool_only" : "automatic"; - }, resolveChunkMode: () => "length", resolveIdentityNamePrefix: (cfg: { agents?: { list?: Array<{ id?: string; default?: boolean; identity?: { name?: string } }> }; @@ -115,20 +136,20 @@ vi.mock("./runtime-api.js", async () => { }, resolveTextChunkLimit: () => 4000, shouldLogVerbose: () => false, - toLocationContext: () => ({}), }; }); import { - buildWhatsAppInboundContext, + buildWhatsAppInboundTransportContext, createWhatsAppReplyPlan, + prepareWhatsAppInboundContext, resolveWhatsAppDmRouteTarget, resolveWhatsAppResponsePrefix, updateWhatsAppMainLastRoute, } from "./inbound-dispatch.js"; -type TestRoute = Parameters[0]["route"]; -type TestMsg = Parameters[0]["msg"]; +type TestRoute = Parameters[0]["route"]; +type TestMsg = Parameters[0]["msg"]; type TestMsgOverrides = NonNullable[0]>; type TestAdmissionOverride = NonNullable; @@ -186,6 +207,72 @@ function makeMsg(overrides: TestMsgOverrides = {}): TestMsg { }); } +function collectNonPortablePaths( + value: unknown, + path = "inbound", + seen = new Set(), +): string[] { + if ((typeof value !== "object" || value === null) && typeof value !== "function") { + return []; + } + if (typeof value === "function") { + return [path]; + } + if (seen.has(value)) { + return []; + } + seen.add(value); + const record = value as Record; + const symbolPaths = Object.getOwnPropertySymbols(record).map( + (symbol) => `${path}.${String(symbol)}`, + ); + return [ + ...symbolPaths, + ...Object.entries(record).flatMap(([key, child]) => + collectNonPortablePaths(child, `${path}.${key}`, seen), + ), + ]; +} + +type PrepareWhatsAppInboundParams = Parameters[0]; +type LegacyTestCommand = Omit< + NonNullable, + "authorization" +> & { + authorized?: boolean; + authorization?: NonNullable["authorization"]; +}; + +async function buildWhatsAppInboundContext( + params: Omit & { + command?: LegacyTestCommand; + }, +) { + const { command: legacyCommand, ...preparedParams } = params; + const command = legacyCommand + ? { + ...legacyCommand, + authorization: + legacyCommand.authorization ?? + (legacyCommand.authorized === undefined + ? { kind: "not_checked" as const } + : legacyCommand.authorized + ? { kind: "authorized" as const } + : { kind: "denied" as const }), + } + : undefined; + if (!command) { + return (await prepareWhatsAppInboundContext(preparedParams)).ctxPayload; + } + const { authorized: _legacyAuthorized, ...preparedCommand } = command; + return ( + await prepareWhatsAppInboundContext({ + ...preparedParams, + command: preparedCommand, + }) + ).ctxPayload; +} + function directAdmission(conversationId: string): TestAdmissionOverride { return { conversation: { @@ -210,6 +297,141 @@ function groupAdmission(conversationId: string): TestAdmissionOverride { }; } +describe("prepared WhatsApp inbound boundary", () => { + it("separates portable facts from WhatsApp transport callbacks", async () => { + const msg = makeMsg({ + event: { id: "current-1", timestamp: 1_710_000_000 }, + payload: { + body: "agent body", + commandBody: "/status", + media: { + path: "/tmp/photo.jpg", + type: "image/jpeg", + kind: "image", + }, + }, + admission: groupAdmission("120363000000000000@g.us"), + groupMention: { + wasMentioned: false, + requireMention: false, + }, + group: { + subject: "Boundary Room", + participants: ["15550001111@s.whatsapp.net"], + }, + }); + const prepared = await prepareWhatsAppInboundContext({ + bodyForAgent: "agent body", + combinedBody: "formatted agent body", + command: { + kind: "text-slash", + body: "/status", + authorization: { kind: "denied", reason: "sender_not_allowed" }, + }, + msg, + route: makeRoute({ + sessionKey: "agent:main:whatsapp:group:120363000000000000@g.us", + }), + sender: { + id: "+15550001111", + name: "Alice", + e164: "+15550001111", + }, + transcript: "prepared transcript", + mediaTranscribedIndexes: [0], + visibleReplyTo: { + id: "quoted-1", + body: "quoted body", + sender: { label: "Bob" }, + }, + replyThreading: { implicitCurrentMessage: "allow" }, + suppressMessageReceivedHooks: true, + }); + + expect(prepared.inbound).toMatchObject({ + event: { + id: "current-1", + timestamp: 1_710_000_000, + }, + message: { + body: "formatted agent body", + bodyForAgent: "agent body", + rawBody: "agent body", + commandBody: "/status", + }, + conversation: { + kind: "group", + id: "120363000000000000@g.us", + label: "120363000000000000@g.us", + }, + reply: { + replyToId: "quoted-1", + }, + command: { + kind: "text-slash", + body: "/status", + authorization: { kind: "denied", reason: "sender_not_allowed" }, + }, + media: [ + { + path: "/tmp/photo.jpg", + contentType: "image/jpeg", + kind: "image", + transcribed: true, + }, + ], + context: { + transcript: "prepared transcript", + groupSubject: "Boundary Room", + senderE164: "+15550001111", + replyThreading: { implicitCurrentMessage: "allow" }, + }, + }); + expect(prepared.ctxPayload).toMatchObject({ + ConversationLabel: "120363000000000000@g.us", + GroupSubject: "Boundary Room", + }); + expect(collectNonPortablePaths(prepared.inbound)).toEqual([]); + expect(prepared.inbound).not.toHaveProperty("platform"); + expect(prepared.inbound).not.toHaveProperty("admission"); + expect(prepared.control).toEqual({ messageReceivedHooks: "channel" }); + + const transport = buildWhatsAppInboundTransportContext(msg); + expect(transport).toMatchObject({ + accountId: "default", + conversationId: "120363000000000000@g.us", + conversationKind: "group", + chatJid: "+1000", + recipientJid: "+2000", + correlationId: "current-1", + }); + expect(transport.reply).toBe(msg.platform.reply); + expect(transport.sendMedia).toBe(msg.platform.sendMedia); + expect(transport.sendComposing).toBe(msg.platform.sendComposing); + expect(transport).not.toHaveProperty("wasMentioned"); + }); + + it("assigns unique portable identities without inventing native message IDs", async () => { + const msg = makeMsg({ + event: { id: undefined, timestamp: 1_710_000_000 }, + }); + const params = { + combinedBody: "hi", + msg, + route: makeRoute(), + sender: { id: "+15550001111" }, + }; + + const [first, second] = await Promise.all([ + prepareWhatsAppInboundContext(params), + prepareWhatsAppInboundContext(params), + ]); + + expect(first.inbound.event.id).not.toBe(second.inbound.event.id); + expect(buildWhatsAppInboundTransportContext(msg).correlationId).toBeUndefined(); + }); +}); + function getCapturedDeliver() { return (capturedDispatchParams as CapturedDispatchParams)?.dispatcherOptions?.deliver; } @@ -276,8 +498,9 @@ function expectRememberSentContextFields( } type BufferedReplyParams = Parameters[0]; -type BufferedReplyOverrides = Partial> & { +type BufferedReplyOverrides = Partial> & { context?: Partial; + msg?: TestMsg; cancelAfterPrepare?: (payload: CapturedReplyPayload) => boolean; }; @@ -319,8 +542,37 @@ function unacceptedDeliveryResult() { }; } +function makePreparedInbound(msg: TestMsg): BufferedReplyParams["inbound"] { + const admission = msg.admission; + return { + channel: "whatsapp", + event: { id: msg.event.id ?? "msg1", timestamp: msg.event.timestamp }, + from: admission.conversation.id, + sender: { id: admission.sender.id }, + conversation: { + kind: admission.conversation.kind, + id: admission.conversation.id, + }, + route: { + agentId: "main", + accountId: admission.accountId, + routeSessionKey: makeRoute().sessionKey, + }, + reply: { + to: msg.platform.recipientJid, + originatingTo: admission.conversation.id, + }, + message: { + body: msg.payload.body, + bodyForAgent: msg.payload.body, + rawBody: msg.payload.commandBody ?? msg.payload.body, + commandBody: msg.payload.commandBody ?? msg.payload.body, + }, + }; +} + async function dispatchBufferedReply(overrides: BufferedReplyOverrides = {}) { - const { cancelAfterPrepare, ...paramOverrides } = overrides; + const { cancelAfterPrepare, msg = makeMsg(), ...paramOverrides } = overrides; const params: BufferedReplyParams = { cfg: { channels: { whatsapp: { streaming: { block: { enabled: true } } } } } as never, connectionId: "conn", @@ -329,13 +581,14 @@ async function dispatchBufferedReply(overrides: BufferedReplyOverrides = {}) { groupHistories: new Map(), groupHistoryKey: "+1000", maxMediaBytes: 1, - msg: makeMsg(), + inbound: makePreparedInbound(msg), rememberSentText: () => {}, replyLogger: makeReplyLogger(), replyPipeline: {} as never, replyResolver: (async () => undefined) as never, route: makeRoute(), shouldClearGroupHistory: false, + transport: buildWhatsAppInboundTransportContext(msg), }; return runWhatsAppReplyPlan( @@ -1135,6 +1388,47 @@ describe("whatsapp inbound dispatch", () => { expect(deliverReply).not.toHaveBeenCalled(); }); + it("does not use a synthetic portable event ID as a WhatsApp reply target", async () => { + deliverInboundReplyWithMessageSendContextMock.mockResolvedValueOnce({ + status: "handled_visible", + delivery: { + messageIds: ["wa-1"], + visibleReplySent: true, + }, + }); + const msg = makeMsg({ + event: { id: undefined, timestamp: 1_710_000_000 }, + }); + + await dispatchBufferedReply({ + context: { + Body: "incoming", + ReplyToId: "quoted-bot-message", + }, + msg, + inbound: { + ...makePreparedInbound(msg), + event: { + id: "120363000000000000@g.us:1710000000", + timestamp: 1_710_000_000, + }, + }, + }); + + const deliver = getCapturedDeliver(); + await deliver?.({ text: "final payload" }, { kind: "final" }); + + const durableParams = requireMockArg( + deliverInboundReplyWithMessageSendContextMock, + 0, + 0, + "durable delivery params", + ); + expectRecordFields(durableParams, { + replyToId: null, + }); + }); + it("does not fall back when durable WhatsApp delivery suppresses a send", async () => { deliverInboundReplyWithMessageSendContextMock.mockResolvedValueOnce({ status: "handled_no_send", @@ -1583,7 +1877,7 @@ describe("whatsapp inbound dispatch", () => { it("suppresses typing for message-tool-only group chat without mention", async () => { await dispatchBufferedReply({ - context: { Body: "hi", ChatType: "group" }, + context: { Body: "hi", ChatType: "group", WasMentioned: false }, msg: makeMsg({ admission: groupAdmission("120363000000000000@g.us"), wasMentioned: false, @@ -1595,7 +1889,7 @@ describe("whatsapp inbound dispatch", () => { it("does not suppress typing for group chat when mentioned", async () => { await dispatchBufferedReply({ - context: { Body: "@bot hi", ChatType: "group" }, + context: { Body: "@bot hi", ChatType: "group", WasMentioned: true }, msg: makeMsg({ admission: groupAdmission("120363000000000000@g.us"), wasMentioned: true, @@ -1700,6 +1994,7 @@ describe("whatsapp inbound dispatch", () => { it("returns true for tool-only media turns after delivering media", async () => { const deliverReply = vi.fn(async () => acceptedDeliveryResult()); const rememberSentText = vi.fn(); + const msg = makeMsg(); dispatchReplyWithBufferedBlockDispatcherMock.mockImplementationOnce( async (params: CapturedDispatchParams) => { capturedDispatchParams = params; @@ -1720,8 +2015,8 @@ describe("whatsapp inbound dispatch", () => { deliverReply, groupHistories: new Map(), groupHistoryKey: "+1000", + inbound: makePreparedInbound(msg), maxMediaBytes: 1, - msg: makeMsg(), rememberSentText, replyLogger: { info: () => {}, @@ -1733,6 +2028,7 @@ describe("whatsapp inbound dispatch", () => { replyResolver: (async () => undefined) as never, route: makeRoute(), shouldClearGroupHistory: false, + transport: buildWhatsAppInboundTransportContext(msg), }), ).resolves.toBe(true); diff --git a/extensions/whatsapp/src/auto-reply/monitor/inbound-dispatch.ts b/extensions/whatsapp/src/auto-reply/monitor/inbound-dispatch.ts index d4cde2bf5928..68fbaa52a3fa 100644 --- a/extensions/whatsapp/src/auto-reply/monitor/inbound-dispatch.ts +++ b/extensions/whatsapp/src/auto-reply/monitor/inbound-dispatch.ts @@ -4,8 +4,10 @@ import { buildChannelInboundEventContext, createChannelPartialDeliveryError, isChannelPartialDeliveryError, - type CommandFacts, + projectPreparedChannelInbound, + resolveChannelInboundReplyPolicy, type ChannelInboundTurnPlan, + type PreparedChannelInbound, toInboundMediaFactsWithMetadata, } from "openclaw/plugin-sdk/channel-inbound"; import { hasVisibleInboundReplyDispatch } from "openclaw/plugin-sdk/channel-inbound"; @@ -18,14 +20,18 @@ import { buildInboundHistoryFromEntries } from "openclaw/plugin-sdk/reply-histor import type { FinalizedMsgContext } from "openclaw/plugin-sdk/reply-runtime"; import { normalizeStringEntries } from "openclaw/plugin-sdk/string-coerce-runtime"; import { requireWhatsAppInboundAdmission } from "../../inbound/admission.js"; -import { resolveWhatsAppIngressLifecycle } from "../../inbound/ingress-lifecycle.js"; import type { AdmittedWebInboundMessage } from "../../inbound/types.js"; import { type DeliverableWhatsAppOutboundPayload, normalizeWhatsAppOutboundPayload, normalizeWhatsAppPayloadTextPreservingIndentation, } from "../../outbound-media-contract.js"; -import type { WhatsAppReplyDeliveryResult } from "../deliver-reply.js"; +import { newConnectionId } from "../../reconnect.js"; +import type { + WhatsAppReplyDeliveryResult, + WhatsAppReplyTransportContext, +} from "../deliver-reply.js"; +import { createWhatsAppReplyTransportContext } from "../deliver-reply.js"; import { markWhatsAppVisibleDeliveryError } from "../util.js"; import type { EchoTracker } from "./echo.js"; import { formatGroupMembers } from "./group-members.js"; @@ -35,7 +41,6 @@ import { getAgentScopedMediaLocalRoots, jidToE164, logVerbose, - resolveChannelMessageSourceReplyDeliveryMode, resolveChunkMode, resolveIdentityNamePrefix, resolveInboundLastRouteSessionKey, @@ -43,7 +48,6 @@ import { resolveSendableOutboundReplyParts, resolveTextChunkLimit, shouldLogVerbose, - toLocationContext, type getChildLogger, type getReplyFromConfig, type LoadConfigFn, @@ -78,6 +82,10 @@ type SenderContext = { e164?: string; }; +type WhatsAppInboundTransportContext = WhatsAppReplyTransportContext & { + sendComposing: AdmittedWebInboundMessage["platform"]["sendComposing"]; +}; + type ReplyDeliveryInfo = { kind: ReplyLifecycleKind }; type PendingWhatsAppMediaOnlyPayload = { @@ -162,20 +170,19 @@ function logWhatsAppReplyDeliveryError(params: { err: unknown; info: ReplyDeliveryInfo; connectionId: string; - msg: AdmittedWebInboundMessage; + transport: WhatsAppInboundTransportContext; replyLogger: ReturnType; }) { - const admission = requireWhatsAppInboundAdmission(params.msg); params.replyLogger.error( { err: normalizeErrForLog(params.err), replyKind: params.info.kind, - correlationId: params.msg.event.id ?? null, + correlationId: params.transport.correlationId ?? null, connectionId: params.connectionId, - conversationId: admission.conversation.id, - chatId: params.msg.platform.chatJid ?? null, - to: admission.conversation.id, - from: params.msg.platform.recipientJid ?? null, + conversationId: params.transport.conversationId, + chatId: params.transport.chatJid, + to: params.transport.conversationId, + from: params.transport.recipientJid, }, "auto-reply delivery failed", ); @@ -184,7 +191,7 @@ function logWhatsAppReplyDeliveryError(params: { function resolveWhatsAppDurableReplyToId(params: { context: FinalizedMsgContext; info: ReplyDeliveryInfo; - msg: AdmittedWebInboundMessage; + currentMessageId?: string; payload: DeliverableWhatsAppOutboundPayload; }): string | null { if (params.payload.replyToId === null) { @@ -197,20 +204,13 @@ function resolveWhatsAppDurableReplyToId(params: { const hasVisibleInboundReplyTarget = Boolean(readTrimmedString(params.context.ReplyToId)) || Boolean(readTrimmedString(params.context.ReplyToIdFull)); - const currentInboundMessageId = readTrimmedString(params.msg.event.id); + const currentInboundMessageId = readTrimmedString(params.currentMessageId); if (params.info.kind === "final" && hasVisibleInboundReplyTarget && currentInboundMessageId) { return currentInboundMessageId; } return null; } -function resolveWhatsAppDisableBlockStreaming(cfg: ReturnType): boolean | undefined { - // The monitor snapshot pins the account-resolved streaming object onto the - // root channel entry, so this root-level read is already account-scoped. - const enabled = resolveChannelStreamingBlockEnabled(cfg.channels?.whatsapp); - return typeof enabled === "boolean" ? !enabled : undefined; -} - function resolveWhatsAppDeliverablePayload( payload: ReplyPayload, info: { kind: ReplyLifecycleKind }, @@ -371,10 +371,19 @@ export function resolveWhatsAppResponsePrefix(params: { ); } -export async function buildWhatsAppInboundContext(params: { +export function buildWhatsAppInboundTransportContext( + msg: AdmittedWebInboundMessage, +): WhatsAppInboundTransportContext { + return { + ...createWhatsAppReplyTransportContext(msg), + sendComposing: msg.platform.sendComposing, + }; +} + +export async function prepareWhatsAppInboundContext(params: { bodyForAgent?: string; combinedBody: string; - command?: CommandFacts; + command?: NonNullable; groupHistory?: GroupHistoryEntry[]; groupHistoryLimit?: number; groupMemberRoster?: Map; @@ -388,7 +397,12 @@ export async function buildWhatsAppInboundContext(params: { replyThreading?: ReplyThreadingContext; visibleReplyTo?: VisibleReplyTarget; suppressMessageReceivedHooks?: boolean; -}): Promise { +}): Promise<{ + inbound: PreparedChannelInbound; + control: Parameters[0]["control"]; + turnInput: ReturnType["input"]; + ctxPayload: FinalizedMsgContext; +}> { const admission = requireWhatsAppInboundAdmission(params.msg); const conversationId = admission.conversation.id; const conversationKind = admission.conversation.kind; @@ -420,7 +434,10 @@ export async function buildWhatsAppInboundContext(params: { : undefined, { transcribed: (_entry, index) => params.mediaTranscribedIndexes?.includes(index) === true }, ); - return buildChannelInboundEventContext({ + const control = { + messageReceivedHooks: params.suppressMessageReceivedHooks ? "channel" : "core", + } as const; + const inbound: PreparedChannelInbound = { channel: "whatsapp", supplemental: { quote: params.visibleReplyTo @@ -434,8 +451,10 @@ export async function buildWhatsAppInboundContext(params: { channelStructuredContext: params.msg.payload.channelStructuredContext, }, media, - messageId: params.msg.event.id, - timestamp: params.msg.event.timestamp, + event: { + id: params.msg.event.id ?? `${conversationId}:${newConnectionId()}`, + timestamp: params.msg.event.timestamp, + }, from: conversationId, sender: { id: params.sender.id ?? params.sender.e164, @@ -455,6 +474,7 @@ export async function buildWhatsAppInboundContext(params: { reply: { to: params.msg.platform.recipientJid, originatingTo: conversationId, + replyToId: params.visibleReplyTo?.id, }, message: { body: params.combinedBody, @@ -469,35 +489,35 @@ export async function buildWhatsAppInboundContext(params: { ? (params.groupHistoryLimit ?? params.groupHistory?.length ?? 0) : 0, }, - access: { - ...(wasMentioned !== undefined + mentions: + wasMentioned !== undefined ? { - mentions: { - canDetectMention: conversationKind === "group", - wasMentioned, - requireMention: params.msg.groupMention?.requireMention, - }, + canDetectMention: conversationKind === "group", + wasMentioned, + requireMention: params.msg.groupMention?.requireMention, } - : {}), - commands: { - authorized: params.command?.authorized === true, - }, - }, + : undefined, command: params.command, - extra: { - Transcript: params.transcript, - GroupSubject: params.msg.group?.subject, - GroupMembers: formatGroupMembers({ + context: { + transcript: params.transcript, + groupSubject: params.msg.group?.subject ?? null, + groupMembers: formatGroupMembers({ participants: params.msg.group?.participants, roster: params.groupMemberRoster, fallbackE164: params.sender.e164, }), - SenderE164: params.sender.e164, - ReplyThreading: params.replyThreading, - SuppressMessageReceivedHooks: params.suppressMessageReceivedHooks, - ...(params.msg.payload.location ? toLocationContext(params.msg.payload.location) : {}), + senderE164: params.sender.e164, + replyThreading: params.replyThreading, + location: params.msg.payload.location, }, - }); + }; + const projected = projectPreparedChannelInbound({ inbound, control }); + return { + inbound, + control, + turnInput: projected.input, + ctxPayload: projected.context, + }; } export function resolveWhatsAppDmRouteTarget(params: { @@ -583,7 +603,7 @@ export function createWhatsAppReplyPlan(params: { deliverReply: (params: { replyResult: ReplyPayload; normalizedReplyResult?: DeliverableWhatsAppOutboundPayload; - msg: AdmittedWebInboundMessage; + transport: WhatsAppReplyTransportContext; mediaLocalRoots: readonly string[]; maxMediaBytes: number; textLimit: number; @@ -597,7 +617,7 @@ export function createWhatsAppReplyPlan(params: { groupHistoryKey: string; maxMediaBytes: number; maxMediaTextChunkLimit?: number; - msg: AdmittedWebInboundMessage; + inbound: PreparedChannelInbound; onModelSelected?: ChannelReplyOnModelSelected; rememberSentText: EchoTracker["rememberText"]; replyLogger: ReturnType; @@ -606,11 +626,12 @@ export function createWhatsAppReplyPlan(params: { route: ReturnType; shouldClearGroupHistory: boolean; statusReactionController?: StatusReactionController | null; + transport: WhatsAppInboundTransportContext; + turnAdoptionLifecycle?: NonNullable< + NonNullable["turnAdoptionLifecycle"] + >; }) { - const admission = requireWhatsAppInboundAdmission(params.msg); - const ingressLifecycle = resolveWhatsAppIngressLifecycle(params.msg); - const conversationId = admission.conversation.id; - const conversationKind = admission.conversation.kind; + const conversationId = params.inbound.conversation.id; const statusReactionController = params.statusReactionController ?? null; const textLimit = params.maxMediaTextChunkLimit ?? resolveTextChunkLimit(params.cfg, "whatsapp"); const chunkMode = resolveChunkMode(params.cfg, "whatsapp", params.route.accountId); @@ -620,17 +641,11 @@ export function createWhatsAppReplyPlan(params: { accountId: params.route.accountId, }); const mediaLocalRoots = getAgentScopedMediaLocalRoots(params.cfg, params.route.agentId); - const sourceReplyDeliveryMode = - params.context.ChatType === "group" || params.context.ChatType === "channel" - ? resolveChannelMessageSourceReplyDeliveryMode({ - cfg: params.cfg, - ctx: params.context, - }) - : undefined; - const sourceRepliesAreToolOnly = sourceReplyDeliveryMode === "message_tool_only"; - const disableBlockStreaming = sourceRepliesAreToolOnly - ? true - : resolveWhatsAppDisableBlockStreaming(params.cfg); + const replyPolicy = resolveChannelInboundReplyPolicy({ + cfg: params.cfg, + ctx: params.context, + blockStreamingEnabled: resolveChannelStreamingBlockEnabled(params.cfg.channels?.whatsapp), + }); let didSendReply = false; let didLogHeartbeatStrip = false; @@ -666,7 +681,7 @@ export function createWhatsAppReplyPlan(params: { delivery = await params.deliverReply({ replyResult: normalizedDeliveryPayload, normalizedReplyResult: normalizedDeliveryPayload, - msg: params.msg, + transport: params.transport, mediaLocalRoots, maxMediaBytes: params.maxMediaBytes, textLimit, @@ -692,12 +707,12 @@ export function createWhatsAppReplyPlan(params: { if (!result.visibleReplySent) { params.replyLogger.warn( { - correlationId: params.msg.event.id ?? null, + correlationId: params.transport.correlationId ?? null, connectionId: params.connectionId, conversationId, - chatId: params.msg.platform.chatJid, + chatId: params.transport.chatJid, to: conversationId, - from: params.msg.platform.recipientJid, + from: params.transport.recipientJid, replyKind: info.kind, }, "auto-reply was not accepted by WhatsApp provider", @@ -736,7 +751,7 @@ export function createWhatsAppReplyPlan(params: { logWhatsAppMediaOnlyFlushResult(flushResult); return whatsAppReplyDeliveryVisibility(didSendReply || flushResult.delivered > 0); }, - onReplyStart: params.msg.platform.sendComposing, + onReplyStart: params.transport.sendComposing, }; const delivery: ChannelInboundTurnPlan["delivery"] = { observeMessageSent: true, @@ -775,7 +790,7 @@ export function createWhatsAppReplyPlan(params: { replyToId: resolveWhatsAppDurableReplyToId({ context: params.context, info, - msg: params.msg, + currentMessageId: params.transport.correlationId, payload, }), formatting: { @@ -836,21 +851,20 @@ export function createWhatsAppReplyPlan(params: { err, info: info as ReplyDeliveryInfo, connectionId: params.connectionId, - msg: params.msg, + transport: params.transport, replyLogger: params.replyLogger, }); }, }; const replyOptions = { - ...(ingressLifecycle ? bindIngressLifecycleToReplyOptions(ingressLifecycle) : {}), - // Message-tool-only unmentioned group turns have no automatic visible reply. - // Suppress composing there so silent background runs do not leak presence. - suppressTyping: - sourceRepliesAreToolOnly && - conversationKind === "group" && - !(params.msg.groupMention?.wasMentioned ?? params.msg.wasMentioned), - disableBlockStreaming, - ...(sourceReplyDeliveryMode ? { sourceReplyDeliveryMode } : {}), + ...(params.turnAdoptionLifecycle + ? { turnAdoptionLifecycle: params.turnAdoptionLifecycle } + : {}), + suppressTyping: replyPolicy.suppressTyping, + disableBlockStreaming: replyPolicy.disableBlockStreaming, + ...(replyPolicy.sourceReplyDeliveryMode + ? { sourceReplyDeliveryMode: replyPolicy.sourceReplyDeliveryMode } + : {}), onModelSelected: params.onModelSelected, ...(statusReactionController ? { diff --git a/extensions/whatsapp/src/auto-reply/monitor/on-message.audio-preflight.test.ts b/extensions/whatsapp/src/auto-reply/monitor/on-message.audio-preflight.test.ts index 67235e4efce2..e40393c5b2e3 100644 --- a/extensions/whatsapp/src/auto-reply/monitor/on-message.audio-preflight.test.ts +++ b/extensions/whatsapp/src/auto-reply/monitor/on-message.audio-preflight.test.ts @@ -95,15 +95,21 @@ vi.mock("openclaw/plugin-sdk/routing", () => ({ }), })); +import { requireAdmittedWhatsAppInboundMessage } from "../../inbound/admission.js"; +import { normalizeWebInboundMessage } from "../../inbound/message-aliases.js"; import { createTestLegacyFlatWebInboundMessage, createTestWebAudioInboundMessage, } from "../../inbound/test-message.test-helper.js"; -import type { WebInboundMessage } from "../../inbound/types.js"; +import type { AdmittedWebInboundMessage, WebInboundMessageInput } from "../../inbound/types.js"; import { createWebOnMessageHandler } from "./on-message.js"; -function makeAudioMsg(): WebInboundMessage { - return createTestWebAudioInboundMessage(); +function admitTestMessage(msg: WebInboundMessageInput): AdmittedWebInboundMessage { + return requireAdmittedWhatsAppInboundMessage(normalizeWebInboundMessage(msg)); +} + +function makeAudioMsg(): AdmittedWebInboundMessage { + return admitTestMessage(createTestWebAudioInboundMessage()); } function makeLegacyAudioMsg() { @@ -114,42 +120,46 @@ function makeLegacyAudioMsg() { }); } -function makeGroupAudioMsg(): WebInboundMessage { - return createTestWebAudioInboundMessage({ - platform: { chatJid: "1203630@g.us" }, - admission: { - conversation: { - kind: "group", - id: "1203630@g.us", +function makeGroupAudioMsg(): AdmittedWebInboundMessage { + return admitTestMessage( + createTestWebAudioInboundMessage({ + platform: { chatJid: "1203630@g.us" }, + admission: { + conversation: { + kind: "group", + id: "1203630@g.us", + }, + senderAccess: { + reasonCode: "group_policy_allowed", + }, }, - senderAccess: { - reasonCode: "group_policy_allowed", - }, - }, - wasMentioned: false, - }); + wasMentioned: false, + }), + ); } -function makeBlockedDirectAudioMsg(): WebInboundMessage { - return createTestWebAudioInboundMessage({ - admission: { - ingress: { - admission: "drop", - decision: "block", - reasonCode: "dm_policy_not_allowlisted", +function makeBlockedDirectAudioMsg(): AdmittedWebInboundMessage { + return admitTestMessage( + createTestWebAudioInboundMessage({ + admission: { + ingress: { + admission: "drop", + decision: "block", + reasonCode: "dm_policy_not_allowlisted", + }, + senderAccess: { + allowed: false, + decision: "block", + reasonCode: "dm_policy_not_allowlisted", + }, + activationAccess: { + allowed: false, + shouldSkip: true, + reasonCode: "dm_policy_not_allowlisted", + }, }, - senderAccess: { - allowed: false, - decision: "block", - reasonCode: "dm_policy_not_allowlisted", - }, - activationAccess: { - allowed: false, - shouldSkip: true, - reasonCode: "dm_policy_not_allowlisted", - }, - }, - }); + }), + ); } function makeEchoTracker() { @@ -281,7 +291,7 @@ describe("createWebOnMessageHandler audio preflight", () => { it("skips early DM ack/preflight for legacy audio without explicit access proof", async () => { const handler = makeHandler(); - await handler(makeLegacyAudioMsg()); + await handler(admitTestMessage(makeLegacyAudioMsg())); expect(events).toStrictEqual([]); expect(transcribeFirstAudioMock).not.toHaveBeenCalled(); @@ -313,7 +323,7 @@ describe("createWebOnMessageHandler audio preflight", () => { } as never, }); - await handler(makeLegacyAudioMsg()); + await handler(admitTestMessage(makeLegacyAudioMsg())); expect(events).toStrictEqual(["stt"]); expect(transcribeFirstAudioMock).toHaveBeenCalledTimes(1); diff --git a/extensions/whatsapp/src/auto-reply/monitor/on-message.ts b/extensions/whatsapp/src/auto-reply/monitor/on-message.ts index 1de1a5c51137..c7d181d2ada2 100644 --- a/extensions/whatsapp/src/auto-reply/monitor/on-message.ts +++ b/extensions/whatsapp/src/auto-reply/monitor/on-message.ts @@ -17,11 +17,11 @@ import { requireAdmittedWhatsAppInboundMessage, requireWhatsAppInboundAdmission, } from "../../inbound/admission.js"; -import { - normalizeWebInboundMessage, - withDeprecatedWebInboundMessageFlatAliases, -} from "../../inbound/message-aliases.js"; -import type { AdmittedWebInboundMessage, WebInboundMessageInput } from "../../inbound/types.js"; +import { withDeprecatedWebInboundMessageFlatAliases } from "../../inbound/message-aliases.js"; +import type { + AdmittedWebInboundMessage, + DeprecatedWebInboundAdmissionTopLevelFields, +} from "../../inbound/types.js"; import { normalizeE164 } from "../../text-runtime.js"; import { buildMentionConfig } from "../mentions.js"; import type { MentionConfig } from "../mentions.js"; @@ -38,6 +38,15 @@ import { type StatusReactionController, } from "./status-reaction.js"; +function readDeprecatedAccessControlPassed(msg: AdmittedWebInboundMessage): boolean | undefined { + // The admitted type hides deprecated flat aliases, but normalized legacy + // listener inputs retain this one tri-state proof for preflight safety. + return ( + msg as AdmittedWebInboundMessage & + Pick + ).accessControlPassed; +} + export function createWebOnMessageHandler(params: { cfg: OpenClawConfig; loadConfig?: () => OpenClawConfig; @@ -54,8 +63,12 @@ export function createWebOnMessageHandler(params: { baseMentionConfig: MentionConfig; account: { authDir?: string; accountId?: string; selfChatMode?: boolean }; }) { - const hasExplicitlyPassedInboundAccess = (msg: WebInboundMessageInput): boolean => - msg.admission ? msg.admission.ingress.decision === "allow" : msg.accessControlPassed === true; + const hasExplicitlyPassedInboundAccess = (msg: AdmittedWebInboundMessage): boolean => { + if (msg.admission.ingress.decisiveGateId === "legacy-flat-compat") { + return readDeprecatedAccessControlPassed(msg) === true; + } + return msg.admission.ingress.decision === "allow"; + }; const withDirectSenderPeer = ( msg: AdmittedWebInboundMessage, @@ -140,9 +153,8 @@ export function createWebOnMessageHandler(params: { return processMessage(processParams); }; - return async (rawMsg: WebInboundMessageInput) => { - const canRunDirectEarlyAudioPreflight = hasExplicitlyPassedInboundAccess(rawMsg); - const normalizedMsg = requireAdmittedWhatsAppInboundMessage(normalizeWebInboundMessage(rawMsg)); + return async (normalizedMsg: AdmittedWebInboundMessage) => { + const canRunDirectEarlyAudioPreflight = hasExplicitlyPassedInboundAccess(normalizedMsg); const cfg = params.loadConfig?.() ?? params.cfg; const peerId = resolvePeerId(normalizedMsg); const msg = withDirectSenderPeer(normalizedMsg, peerId); diff --git a/extensions/whatsapp/src/auto-reply/monitor/process-message.audio-preflight.test.ts b/extensions/whatsapp/src/auto-reply/monitor/process-message.audio-preflight.test.ts index 23ff3e8ac0de..d8c56cfaf0ec 100644 --- a/extensions/whatsapp/src/auto-reply/monitor/process-message.audio-preflight.test.ts +++ b/extensions/whatsapp/src/auto-reply/monitor/process-message.audio-preflight.test.ts @@ -38,9 +38,13 @@ vi.mock("../../session.js", () => ({ formatError: (err: unknown) => String(err), })); -vi.mock("../deliver-reply.js", () => ({ - deliverWebReply: vi.fn(async () => {}), -})); +vi.mock("../deliver-reply.js", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + deliverWebReply: vi.fn(async () => {}), + }; +}); vi.mock("../loggers.js", () => ({ whatsappInboundLog: { info: () => {}, debug: () => {} }, @@ -88,38 +92,41 @@ vi.mock("./runtime-api.js", () => ({ type: undefined, })); -vi.mock("./inbound-dispatch.js", () => ({ - buildWhatsAppInboundContext: (params: { - bodyForAgent?: string; - combinedBody: string; - commandAuthorized?: boolean; - commandBody?: string; - msg: WebInboundMsg; - mediaTranscribedIndexes?: number[]; - rawBody?: string; - transcript?: string; - }) => ({ - Body: params.combinedBody, - BodyForAgent: params.bodyForAgent ?? params.msg.payload.body, - CommandAuthorized: params.commandAuthorized, - CommandBody: params.commandBody ?? params.msg.payload.body, - MediaPath: params.msg.payload.media?.path, - MediaType: params.msg.payload.media?.type, - MediaTranscribedIndexes: params.mediaTranscribedIndexes, - RawBody: params.rawBody ?? params.msg.payload.body, - Transcript: params.transcript, - }), - createWhatsAppReplyPlan: vi.fn((params: { replyResolver?: unknown }) => ({ - dispatcherOptions: {}, - delivery: { deliver: async () => {} }, - replyOptions: {}, - replyResolver: params.replyResolver, - finalize: () => true, - })), - resolveWhatsAppDmRouteTarget: () => "+15550000002", - resolveWhatsAppResponsePrefix: () => undefined, - updateWhatsAppMainLastRoute: () => {}, -})); +vi.mock("./inbound-dispatch.js", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + prepareWhatsAppInboundContext: async ( + params: Parameters[0], + ) => { + const prepared = await actual.prepareWhatsAppInboundContext(params); + return { + ...prepared, + ctxPayload: { + Body: params.combinedBody, + BodyForAgent: params.bodyForAgent ?? params.msg.payload.body, + CommandAuthorized: params.command?.authorization.kind === "authorized", + CommandBody: params.command?.body ?? params.msg.payload.body, + MediaPath: params.msg.payload.media?.path, + MediaType: params.msg.payload.media?.type, + MediaTranscribedIndexes: params.mediaTranscribedIndexes, + RawBody: params.rawBody ?? params.msg.payload.body, + Transcript: params.transcript, + }, + }; + }, + createWhatsAppReplyPlan: vi.fn((params: { replyResolver?: unknown }) => ({ + dispatcherOptions: {}, + delivery: { deliver: async () => {} }, + replyOptions: {}, + replyResolver: params.replyResolver, + finalize: () => true, + })), + resolveWhatsAppDmRouteTarget: () => "+15550000002", + resolveWhatsAppResponsePrefix: () => undefined, + updateWhatsAppMainLastRoute: () => {}, + }; +}); import { createWhatsAppReplyPlan } from "./inbound-dispatch.js"; import { processMessage } from "./process-message.js"; diff --git a/extensions/whatsapp/src/auto-reply/monitor/process-message.test.ts b/extensions/whatsapp/src/auto-reply/monitor/process-message.test.ts index 5426b47b8d99..ea4c7d24162d 100644 --- a/extensions/whatsapp/src/auto-reply/monitor/process-message.test.ts +++ b/extensions/whatsapp/src/auto-reply/monitor/process-message.test.ts @@ -9,6 +9,8 @@ const { buildContextMock, isControlCommandMessageMock, dispatchBufferedReplyMock, + replyPlanParamsMock, + runChannelInboundEventParamsMock, runMessageReceivedMock, shouldComputeCommandAuthorizedMock, trackBackgroundTaskMock, @@ -20,11 +22,24 @@ const { queuedFinal: false, counts: { tool: 0, block: 0, final: 0 }, })), + replyPlanParamsMock: vi.fn(), + runChannelInboundEventParamsMock: vi.fn(), runMessageReceivedMock: vi.fn(async () => undefined), shouldComputeCommandAuthorizedMock: vi.fn(() => false), trackBackgroundTaskMock: vi.fn(), })); +vi.mock("openclaw/plugin-sdk/channel-inbound", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + runChannelInboundEvent: async (params: Parameters[0]) => { + runChannelInboundEventParamsMock(params); + return await actual.runChannelInboundEvent(params); + }, + }; +}); + vi.mock("../../inbound-policy.js", async (importOriginal) => { const actual = await importOriginal(); return { @@ -38,9 +53,18 @@ vi.mock("./inbound-dispatch.js", async (importOriginal) => { const actual = await importOriginal(); return { ...actual, - buildWhatsAppInboundContext: buildContextMock, + prepareWhatsAppInboundContext: async ( + params: Parameters[0], + ) => { + const prepared = await actual.prepareWhatsAppInboundContext(params); + return { + ...prepared, + ctxPayload: buildContextMock(params), + }; + }, createWhatsAppReplyPlan: (...args: unknown[]) => { const params = args[0] as { replyResolver?: unknown }; + replyPlanParamsMock(params); void dispatchBufferedReplyMock(params); return { dispatcherOptions: {}, @@ -151,6 +175,7 @@ vi.mock("./runtime-api.js", async (importOriginal) => { }); import { clearInternalHooks, registerInternalHook } from "openclaw/plugin-sdk/hook-runtime"; +import { attachWhatsAppIngressLifecycle } from "../../inbound/ingress-lifecycle.js"; import { processMessage } from "./process-message.js"; // --------------------------------------------------------------------------- @@ -282,6 +307,8 @@ describe("processMessage group system prompt wiring", () => { isControlCommandMessageMock.mockReset(); isControlCommandMessageMock.mockReturnValue(false); resolvePolicyMock.mockReset(); + replyPlanParamsMock.mockClear(); + runChannelInboundEventParamsMock.mockClear(); runMessageReceivedMock.mockClear(); shouldComputeCommandAuthorizedMock.mockReset(); shouldComputeCommandAuthorizedMock.mockReturnValue(false); @@ -329,7 +356,7 @@ describe("processMessage group system prompt wiring", () => { expect(mockCallArg(buildContextMock, "buildWhatsAppInboundContext")).toMatchObject({ command: { kind: "text-slash", - authorized: true, + authorization: { kind: "authorized" }, body: "/status", }, rawBody: "/status", @@ -354,7 +381,7 @@ describe("processMessage group system prompt wiring", () => { bodyForAgent: "/reset\n\n[whatsapp attachment unavailable]", command: { kind: "text-slash", - authorized: true, + authorization: { kind: "authorized" }, body: "/reset", }, rawBody: "/reset", @@ -373,7 +400,7 @@ describe("processMessage group system prompt wiring", () => { expect(mockCallArg(buildContextMock, "buildWhatsAppInboundContext")).toMatchObject({ command: { kind: "normal", - authorized: true, + authorization: { kind: "authorized" }, body: "please inspect `/tmp/foo`", }, rawBody: "please inspect `/tmp/foo`", @@ -557,6 +584,38 @@ describe("processMessage group system prompt wiring", () => { ); }); + it("passes one lifecycle identity through the portable boundary and reply plan", async () => { + resolvePolicyMock.mockReturnValue(makePolicy(makeAccount())); + buildContextMock.mockImplementationOnce(() => ({ + Body: "hi", + RawBody: "hi", + CommandBody: "hi", + SessionKey: baseRoute.sessionKey, + Provider: "whatsapp", + Surface: "whatsapp", + })); + const lifecycle = { + abortSignal: new AbortController().signal, + onAdopted: vi.fn(async () => undefined), + onDeferred: vi.fn(), + onAbandoned: vi.fn(async () => undefined), + }; + const msg = attachWhatsAppIngressLifecycle(makeBaseMsg(), lifecycle as never); + + await callProcessMessage({ msg }); + + const runParams = mockCallArg(runChannelInboundEventParamsMock, "runChannelInboundEvent") as { + raw?: unknown; + turnAdoptionLifecycle?: unknown; + }; + const replyPlanParams = mockCallArg(replyPlanParamsMock, "createWhatsAppReplyPlan") as { + turnAdoptionLifecycle?: unknown; + }; + expect(runParams.turnAdoptionLifecycle).toBe(replyPlanParams.turnAdoptionLifecycle); + expect(runParams.raw).not.toHaveProperty("platform"); + expect(runParams.raw).not.toHaveProperty("admission"); + }); + it("drops blocked admission before session record and reply dispatch", async () => { resolvePolicyMock.mockReturnValue(makePolicy(makeAccount())); buildContextMock.mockImplementationOnce(() => ({ diff --git a/extensions/whatsapp/src/auto-reply/monitor/process-message.ts b/extensions/whatsapp/src/auto-reply/monitor/process-message.ts index 9d4e1a509d61..4d62410edbbe 100644 --- a/extensions/whatsapp/src/auto-reply/monitor/process-message.ts +++ b/extensions/whatsapp/src/auto-reply/monitor/process-message.ts @@ -8,6 +8,7 @@ import { formatMediaPlaceholderText, runChannelInboundEvent, } from "openclaw/plugin-sdk/channel-inbound"; +import { bindIngressLifecycleToReplyOptions } from "openclaw/plugin-sdk/channel-outbound"; import { createInternalHookEvent, deriveInboundMessageHookContext, @@ -25,6 +26,7 @@ import { resolveWhatsAppInboundPolicy, } from "../../inbound-policy.js"; import { requireWhatsAppInboundAdmission } from "../../inbound/admission.js"; +import { resolveWhatsAppIngressLifecycle } from "../../inbound/ingress-lifecycle.js"; import type { AdmittedWebInboundMessage } from "../../inbound/types.js"; import { newConnectionId } from "../../reconnect.js"; import { formatError } from "../../session.js"; @@ -43,8 +45,9 @@ import { type GroupHistoryEntry, } from "./inbound-context.js"; import { - buildWhatsAppInboundContext, + buildWhatsAppInboundTransportContext, createWhatsAppReplyPlan, + prepareWhatsAppInboundContext, resolveWhatsAppDmRouteTarget, resolveWhatsAppResponsePrefix, updateWhatsAppMainLastRoute, @@ -118,7 +121,7 @@ function shouldEmitWhatsAppMessageReceivedHooks(params: { } function emitWhatsAppMessageReceivedHooks(params: { - ctx: Awaited>; + ctx: Awaited>["ctxPayload"]; sessionKey: string; }): void { const canonical = deriveInboundMessageHookContext(params.ctx); @@ -153,7 +156,7 @@ function emitWhatsAppMessageReceivedHooks(params: { function emitWhatsAppMessageReceivedHooksIfEnabled(params: { cfg: ReturnType; - ctx: Awaited>; + ctx: Awaited>["ctxPayload"]; accountId?: string; sessionKey: string; }): void { @@ -470,13 +473,19 @@ export async function processMessage(params: { peerId: dmRouteTarget ?? conversationId, }); - const ctxPayload = await buildWhatsAppInboundContext({ + const commandAuthorization = + commandAuthorized === undefined + ? ({ kind: "not_checked" } as const) + : commandAuthorized + ? ({ kind: "authorized" } as const) + : ({ kind: "denied" } as const); + const prepared = await prepareWhatsAppInboundContext({ bodyForAgent: msgForAgent.payload.body, combinedBody, command: { kind: isTextCommand ? "text-slash" : "normal", body: commandBody, - authorized: commandAuthorized, + authorization: commandAuthorization, }, groupHistory: visibleGroupHistory, groupHistoryLimit: params.groupHistoryLimit, @@ -496,6 +505,12 @@ export async function processMessage(params: { visibleReplyTo: visibleReplyTo ?? undefined, suppressMessageReceivedHooks: true, }); + const { inbound, turnInput, ctxPayload } = prepared; + const transport = buildWhatsAppInboundTransportContext(params.msg); + const ingressLifecycle = resolveWhatsAppIngressLifecycle(params.msg); + const turnAdoptionLifecycle = ingressLifecycle + ? bindIngressLifecycleToReplyOptions(ingressLifecycle).turnAdoptionLifecycle + : undefined; emitWhatsAppMessageReceivedHooksIfEnabled({ cfg: params.cfg, ctx: ctxPayload, @@ -522,16 +537,10 @@ export async function processMessage(params: { const turnResult = await runChannelInboundEvent({ channel: "whatsapp", accountId: params.route.accountId, - raw: params.msg, + raw: inbound, + ...(turnAdoptionLifecycle ? { turnAdoptionLifecycle } : {}), adapter: { - ingest: () => ({ - id: params.msg.event.id ?? `${conversationId}:${Date.now()}`, - timestamp: params.msg.event.timestamp, - rawText: ctxPayload.RawBody ?? "", - textForAgent: ctxPayload.BodyForAgent, - textForCommands: ctxPayload.CommandBody, - raw: params.msg, - }), + ingest: () => turnInput, preflight: () => { const reason = admission.ingress.reasonCode; if (admission.ingress.admission === "dispatch") { @@ -561,7 +570,7 @@ export async function processMessage(params: { groupHistoryKey: params.groupHistoryKey, maxMediaBytes: params.maxMediaBytes, maxMediaTextChunkLimit: params.maxMediaTextChunkLimit, - msg: params.msg, + inbound, onModelSelected, rememberSentText: params.rememberSentText, replyLogger: params.replyLogger, @@ -573,6 +582,8 @@ export async function processMessage(params: { route: params.route, shouldClearGroupHistory, statusReactionController, + transport, + turnAdoptionLifecycle, }); finalizeReply = finalize; return { diff --git a/extensions/whatsapp/src/auto-reply/monitor/runtime-api.ts b/extensions/whatsapp/src/auto-reply/monitor/runtime-api.ts index 40e260fd495e..643e3318291d 100644 --- a/extensions/whatsapp/src/auto-reply/monitor/runtime-api.ts +++ b/extensions/whatsapp/src/auto-reply/monitor/runtime-api.ts @@ -2,11 +2,7 @@ export { resolveIdentityNamePrefix } from "openclaw/plugin-sdk/agent-runtime"; export { formatInboundEnvelope } from "openclaw/plugin-sdk/channel-inbound"; export { resolveInboundSessionEnvelopeContext } from "openclaw/plugin-sdk/channel-inbound"; -export { toLocationContext } from "openclaw/plugin-sdk/channel-inbound"; -export { - createChannelMessageReplyPipeline, - resolveChannelMessageSourceReplyDeliveryMode, -} from "openclaw/plugin-sdk/channel-outbound"; +export { createChannelMessageReplyPipeline } from "openclaw/plugin-sdk/channel-outbound"; export { isControlCommandMessage, shouldComputeCommandAuthorized, diff --git a/extensions/whatsapp/src/inbound/message-aliases.test.ts b/extensions/whatsapp/src/inbound/message-aliases.test.ts index 55f01934eb71..ef001bf447b6 100644 --- a/extensions/whatsapp/src/inbound/message-aliases.test.ts +++ b/extensions/whatsapp/src/inbound/message-aliases.test.ts @@ -6,7 +6,10 @@ import { } from "./message-aliases.js"; import type { monitorWebInbox } from "./monitor.js"; import { createAcceptedWhatsAppSendResult } from "./send-result.test-helper.js"; -import { createTestWhatsAppInboundAdmission } from "./test-message.test-helper.js"; +import { + createTestLegacyFlatWebInboundMessage, + createTestWhatsAppInboundAdmission, +} from "./test-message.test-helper.js"; import type { LegacyFlatWebInboundMessage, WebInboundCallbackMessage } from "./types.js"; type MonitorWebInboxMessage = Parameters[0]["onMessage"]>[0]; @@ -361,7 +364,7 @@ describe("WhatsApp inbound flat aliases", () => { decisiveGateId: "legacy-flat-compat", reasonCode: "dm_policy_allowlisted", }); - expect(normalized.accessControlPassed).toBe(true); + expect(normalized.accessControlPassed).toBeUndefined(); expect(normalized.quote).toMatchObject({ id: "quote-legacy", body: "legacy quoted", @@ -414,4 +417,16 @@ describe("WhatsApp inbound flat aliases", () => { }); expect(normalized.accessControlPassed).toBe(false); }); + + it("preserves explicit legacy access proof through normalization", () => { + const normalized = normalizeWebInboundMessage({ + ...createTestLegacyFlatWebInboundMessage(), + accessControlPassed: true, + }); + + expect(normalized.admission?.ingress.decisiveGateId).toBe("legacy-flat-compat"); + expect(normalized.accessControlPassed).toBe(true); + normalized.accessControlPassed = false; + expect(normalized.accessControlPassed).toBe(true); + }); }); diff --git a/extensions/whatsapp/src/inbound/message-aliases.ts b/extensions/whatsapp/src/inbound/message-aliases.ts index 343ed329c00f..ce4be8803d28 100644 --- a/extensions/whatsapp/src/inbound/message-aliases.ts +++ b/extensions/whatsapp/src/inbound/message-aliases.ts @@ -192,12 +192,22 @@ function defineDeprecatedAdmissionTopLevelAccessors - msg.admission ? msg.admission.ingress.decision === "allow" : fallbackAccessControlPassed, + get: () => { + // Legacy flat inputs used absence to mean access was not explicitly proven. + // Preserve that tri-state after normalization so preflight work cannot run early. + if (msg.admission?.ingress.decisiveGateId === "legacy-flat-compat") { + return fallbackAccessControlPassed; + } + return msg.admission + ? msg.admission.ingress.decision === "allow" + : fallbackAccessControlPassed; + }, set: (value) => { // The legacy boolean is derived from the ingress graph; writes only preserve // no-admission legacy inputs instead of fabricating a partial graph update. - fallbackAccessControlPassed = value as boolean | undefined; + if (!msg.admission) { + fallbackAccessControlPassed = value as boolean | undefined; + } }, }, chatType: { diff --git a/extensions/whatsapp/src/inbound/monitor.ts b/extensions/whatsapp/src/inbound/monitor.ts index c3b78f6b6518..353c4886c9ed 100644 --- a/extensions/whatsapp/src/inbound/monitor.ts +++ b/extensions/whatsapp/src/inbound/monitor.ts @@ -84,13 +84,8 @@ type MonitorWebInboxOptions = { durableInboundQueue?: WhatsAppDurableInboundQueue; }; -type AttachWebInboxToSocketOptions = Omit< - MonitorWebInboxOptions, - "onMessage" | "shouldDebounce" | "socketTiming" -> & { +type AttachWebInboxToSocketOptions = Omit & { socketTiming: Required; - onMessage: (msg: WebInboundMessageInput) => Promise; - shouldDebounce?: (msg: WebInboundMessageInput) => boolean; }; export async function attachWebInboxToSocket( diff --git a/scripts/plugin-sdk-surface-report.mjs b/scripts/plugin-sdk-surface-report.mjs index 0e3a8b1b42f5..c5595ebd9282 100644 --- a/scripts/plugin-sdk-surface-report.mjs +++ b/scripts/plugin-sdk-surface-report.mjs @@ -204,7 +204,8 @@ export function readPluginSdkSurfaceBudgets(env = process.env) { // +2: shared ingress error factory through channel-outbound and channel-message. // +2: shared ingress retention defaults through channel-outbound and channel-message. // +1: collision-safe MCP server-name assignment for native harness catalogs. - 4774, + // +3: prepared channel inbound envelope, projection, and reply-policy contracts. + 4777, env, ), publicFunctionExports: readPluginSdkSurfaceBudgetEnv( @@ -242,7 +243,8 @@ export function readPluginSdkSurfaceBudgets(env = process.env) { // +1: inbound media-fact metadata projection for plugin-owned channel ingestion. // +2: shared ingress error factory through channel-outbound and channel-message. // +1: collision-safe MCP server-name assignment for native harness catalogs. - 2882, + // +2: prepared channel inbound projection and reply-policy resolution. + 2884, env, ), publicDeprecatedExports: readPluginSdkSurfaceBudgetEnv( diff --git a/src/channels/inbound-event/prepared.ts b/src/channels/inbound-event/prepared.ts new file mode 100644 index 000000000000..d3ebf871f28a --- /dev/null +++ b/src/channels/inbound-event/prepared.ts @@ -0,0 +1,125 @@ +import type { ReplyThreadingPolicy } from "../../auto-reply/types.js"; +import type { NormalizedLocation } from "../location.js"; +import { toLocationContext } from "../location.js"; +import type { CommandFacts, NormalizedTurnInput, SupplementalContextFacts } from "../turn/types.js"; +import { + buildChannelInboundEventContext, + type BuildChannelInboundEventContextParams, + type BuiltChannelInboundEventContext, +} from "./context.js"; + +type PreparedChannelInboundCommandAuthorization = + | { kind: "not_checked" } + | { kind: "authorized" } + | { kind: "denied"; reason?: string }; + +type PreparedChannelInboundCommand = Omit & { + authorization: PreparedChannelInboundCommandAuthorization; +}; + +type PreparedChannelInboundContextFacts = { + transcript?: string; + /** Null suppresses the generic fallback from conversation label to group subject. */ + groupSubject?: string | null; + groupMembers?: string; + senderE164?: string; + replyThreading?: ReplyThreadingPolicy; + location?: NormalizedLocation; +}; + +export type PreparedChannelInbound = Pick< + BuildChannelInboundEventContextParams, + | "channel" + | "accountId" + | "provider" + | "surface" + | "from" + | "sender" + | "conversation" + | "route" + | "reply" + | "message" + | "sessionTranscript" + | "media" + | "contextVisibility" +> & { + event: { + id: string; + fullId?: string; + timestamp?: number; + }; + command?: PreparedChannelInboundCommand; + mentions?: NonNullable["mentions"]; + supplemental?: SupplementalContextFacts; + context?: PreparedChannelInboundContextFacts; +}; + +type PreparedChannelInboundControl = { + messageReceivedHooks: "channel" | "core"; +}; + +function resolvePreparedCommandFacts( + command: PreparedChannelInboundCommand | undefined, +): CommandFacts | undefined { + if (!command) { + return undefined; + } + const { authorization, ...facts } = command; + return { + ...facts, + ...(authorization.kind === "not_checked" + ? {} + : { authorized: authorization.kind === "authorized" }), + }; +} + +export function projectPreparedChannelInbound(params: { + inbound: PreparedChannelInbound; + control: PreparedChannelInboundControl; +}): { + input: NormalizedTurnInput; + context: BuiltChannelInboundEventContext; +} { + const { inbound, control } = params; + const command = resolvePreparedCommandFacts(inbound.command); + const commandAuthorization = inbound.command?.authorization; + const commandAccess = + commandAuthorization && commandAuthorization.kind !== "not_checked" + ? { authorized: commandAuthorization.kind === "authorized" } + : undefined; + return { + input: { + id: inbound.event.id, + timestamp: inbound.event.timestamp, + rawText: inbound.message.rawBody, + textForAgent: inbound.message.bodyForAgent, + textForCommands: inbound.message.commandBody, + raw: inbound, + }, + context: buildChannelInboundEventContext({ + ...inbound, + messageId: inbound.event.id, + messageIdFull: inbound.event.fullId, + timestamp: inbound.event.timestamp, + command, + access: + inbound.mentions || commandAccess + ? { + mentions: inbound.mentions, + commands: commandAccess, + } + : undefined, + extra: { + Transcript: inbound.context?.transcript, + ...(inbound.context && "groupSubject" in inbound.context + ? { GroupSubject: inbound.context.groupSubject ?? undefined } + : {}), + GroupMembers: inbound.context?.groupMembers, + SenderE164: inbound.context?.senderE164, + ReplyThreading: inbound.context?.replyThreading, + SuppressMessageReceivedHooks: control.messageReceivedHooks === "channel", + ...(inbound.context?.location ? toLocationContext(inbound.context.location) : {}), + }, + }), + }; +} diff --git a/src/channels/message/index.ts b/src/channels/message/index.ts index 0b329f470693..10d8a06537af 100644 --- a/src/channels/message/index.ts +++ b/src/channels/message/index.ts @@ -43,6 +43,7 @@ export { createReplyPrefixContext, createReplyPrefixOptions, createTypingCallbacks, + resolveChannelInboundReplyPolicy, resolveChannelSourceReplyDeliveryMode, } from "./reply-pipeline.js"; export type { ChannelIngressDrain } from "./ingress-drain.js"; diff --git a/src/channels/message/reply-pipeline.ts b/src/channels/message/reply-pipeline.ts index 1ea204b6756e..359ae358ca89 100644 --- a/src/channels/message/reply-pipeline.ts +++ b/src/channels/message/reply-pipeline.ts @@ -45,6 +45,43 @@ export function resolveChannelSourceReplyDeliveryMode(params: { return resolveSourceReplyDeliveryMode(params); } +type ChannelInboundReplyPolicy = { + sourceReplyDeliveryMode?: SourceReplyDeliveryMode; + disableBlockStreaming?: boolean; + suppressTyping: boolean; +}; + +/** Resolves product reply policy from portable turn facts and channel capabilities. */ +export function resolveChannelInboundReplyPolicy(params: { + cfg: OpenClawConfig; + ctx: SourceReplyDeliveryModeContext & { + ChatType?: string; + WasMentioned?: boolean; + }; + blockStreamingEnabled?: boolean; +}): ChannelInboundReplyPolicy { + const isRoom = params.ctx.ChatType === "group" || params.ctx.ChatType === "channel"; + const sourceReplyDeliveryMode = isRoom + ? resolveChannelSourceReplyDeliveryMode({ + cfg: params.cfg, + ctx: params.ctx, + }) + : undefined; + const sourceRepliesAreToolOnly = sourceReplyDeliveryMode === "message_tool_only"; + return { + sourceReplyDeliveryMode, + disableBlockStreaming: sourceRepliesAreToolOnly + ? true + : typeof params.blockStreamingEnabled === "boolean" + ? !params.blockStreamingEnabled + : undefined, + suppressTyping: + sourceRepliesAreToolOnly && + params.ctx.ChatType === "group" && + params.ctx.WasMentioned !== true, + }; +} + /** Reply pipeline options shared by core channel turns and plugin SDK callers. */ export type ChannelReplyPipeline = ReplyPrefixOptions & { /** Resolves a response prefix against the pipeline's live selected-model context. */ diff --git a/src/plugin-sdk/channel-inbound.test.ts b/src/plugin-sdk/channel-inbound.test.ts index b9ee05231ae9..940574c78f2a 100644 --- a/src/plugin-sdk/channel-inbound.test.ts +++ b/src/plugin-sdk/channel-inbound.test.ts @@ -5,6 +5,8 @@ import { describe, expect, expectTypeOf, it } from "vitest"; import { buildChannelInboundEventContext, type BuildChannelInboundEventContextParams, + type PreparedChannelInbound, + projectPreparedChannelInbound, type PluginHookChannelSenderContext, } from "./channel-inbound.js"; @@ -71,4 +73,97 @@ describe("channel-inbound public helpers", () => { expect(ctx.ChannelContext?.sender?.testUnionId).toBe("union-1"); }); + + it("builds a portable prepared inbound without channel-native types", () => { + const inbound = { + channel: "example", + accountId: "work", + event: { + id: "event-1", + fullId: "example:event-1", + timestamp: 1_710_000_000, + }, + from: "example:user:u1", + sender: { + id: "u1", + name: "Alice", + }, + conversation: { + kind: "group", + id: "room-1", + label: "Example Room", + }, + route: { + agentId: "main", + accountId: "work", + routeSessionKey: "agent:main:example:group:room-1", + }, + reply: { + to: "example:room:room-1", + replyToId: "quoted-1", + }, + message: { + body: "agent body", + bodyForAgent: "agent body", + rawBody: "raw body", + commandBody: "/status", + }, + command: { + kind: "text-slash", + body: "/status", + authorization: { + kind: "denied", + reason: "sender_not_allowed", + }, + }, + media: [ + { + path: "/tmp/example.jpg", + contentType: "image/jpeg", + kind: "image", + }, + ], + context: { + senderE164: "+15550001111", + }, + } satisfies PreparedChannelInbound; + + const projected = projectPreparedChannelInbound({ + inbound, + control: { messageReceivedHooks: "core" }, + }); + expect(projected.input).toEqual({ + id: "event-1", + timestamp: 1_710_000_000, + rawText: "raw body", + textForAgent: "agent body", + textForCommands: "/status", + raw: inbound, + }); + expect(inbound.command.authorization).toEqual({ + kind: "denied", + reason: "sender_not_allowed", + }); + + const ctx = projected.context; + expect(ctx).toMatchObject({ + MessageSid: "event-1", + MessageSidFull: "example:event-1", + BodyForAgent: "agent body", + RawBody: "raw body", + CommandBody: "/status", + ReplyToId: "quoted-1", + CommandAuthorized: false, + ConversationLabel: "Example Room", + GroupSubject: "Example Room", + SenderE164: "+15550001111", + media: [ + { + path: "/tmp/example.jpg", + contentType: "image/jpeg", + kind: "image", + }, + ], + }); + }); }); diff --git a/src/plugin-sdk/channel-inbound.ts b/src/plugin-sdk/channel-inbound.ts index db4419a47b00..c21d66a38e9e 100644 --- a/src/plugin-sdk/channel-inbound.ts +++ b/src/plugin-sdk/channel-inbound.ts @@ -114,6 +114,11 @@ export type { FinalizeChannelInboundContextParams, FinalizeChannelInboundContextResult, }; +export { + projectPreparedChannelInbound, + type PreparedChannelInbound, +} from "../channels/inbound-event/prepared.js"; +export { resolveChannelInboundReplyPolicy } from "../channels/message/index.js"; /** * Deprecated turn-context input alias that still accepts the old `inboundTurnKind` name. *