diff --git a/extensions/signal/src/monitor/event-handler.ts b/extensions/signal/src/monitor/event-handler.ts index 0509959e9272..af15c1d4a8b5 100644 --- a/extensions/signal/src/monitor/event-handler.ts +++ b/extensions/signal/src/monitor/event-handler.ts @@ -946,17 +946,20 @@ export function createSignalEventHandler(deps: SignalEventHandlerDeps) { return async ( event: { event?: string; data?: string }, turnAdoptionLifecycle?: SignalIngressLifecycle, + preparedPayload?: SignalReceivePayload, ): Promise<{ kind: "deferred" } | { kind: "failed-retryable"; error: unknown } | void> => { if (event.event !== "receive" || !event.data) { return; } - let payload: SignalReceivePayload | null; - try { - payload = JSON.parse(event.data) as SignalReceivePayload; - } catch (err) { - deps.runtime.error?.(`failed to parse event: ${String(err)}`); - return; + let payload: SignalReceivePayload | null = preparedPayload ?? null; + if (!preparedPayload) { + try { + payload = JSON.parse(event.data) as SignalReceivePayload; + } catch (err) { + deps.runtime.error?.(`failed to parse event: ${String(err)}`); + return; + } } if (payload?.exception?.message) { deps.runtime.error?.(`receive exception: ${payload.exception.message}`); diff --git a/extensions/signal/src/signal-ingress.test.ts b/extensions/signal/src/signal-ingress.test.ts index a35bbfa48e67..244a403d7548 100644 --- a/extensions/signal/src/signal-ingress.test.ts +++ b/extensions/signal/src/signal-ingress.test.ts @@ -119,7 +119,9 @@ describe("Signal durable ingress", () => { try { await recovered.waitForIdle(); expect(recoveredDispatch).toHaveBeenCalledTimes(1); - expect(recoveredDispatch).toHaveBeenCalledWith(event, expect.any(Object)); + const [recoveredEvent, recoveredLifecycle] = recoveredDispatch.mock.calls[0] ?? []; + expect(recoveredEvent).toEqual(event); + expect(recoveredLifecycle).toEqual(expect.any(Object)); } finally { await recovered.monitor.stop(); } diff --git a/extensions/signal/src/signal-ingress.ts b/extensions/signal/src/signal-ingress.ts index 3e0daeb5f942..6e1c5ddd0569 100644 --- a/extensions/signal/src/signal-ingress.ts +++ b/extensions/signal/src/signal-ingress.ts @@ -13,6 +13,7 @@ import { normalizeNullableString as normalizeRawString, } from "openclaw/plugin-sdk/string-coerce-runtime"; import type { SignalSseEvent } from "./client-adapter.js"; +import type { SignalReceivePayload } from "./monitor/event-handler.types.js"; import { getOptionalSignalRuntime } from "./runtime.js"; const SIGNAL_INGRESS_DRAIN_INTERVAL_MS = 1_000; @@ -33,6 +34,11 @@ type SignalIngressEventFacts = { numberAliasEventId?: string; }; +type SignalPreparedIngressEvent = [ + event: SignalSseEvent, + parsedPayload: SignalReceivePayload | null | undefined, +]; + type SignalIngressPayload = { version: 1; receivedAt: number; @@ -48,6 +54,7 @@ type SignalIngressDispatchResult = ChannelIngressMonitorDeliveryResult; type SignalIngressDispatch = ( event: SignalSseEvent, lifecycle: SignalIngressLifecycle, + parsedPayload: SignalReceivePayload, ) => Promise | SignalIngressDispatchResult | void; const SignalIngressPermanentError = createChannelIngressError< @@ -58,7 +65,7 @@ function normalizeTimestamp(value: unknown): number | null { return asPositiveSafeInteger(value) ?? null; } -function parseReceiveEnvelope(event: SignalSseEvent): SignalIngressEnvelope | null { +function parseReceivePayload(event: SignalSseEvent): SignalReceivePayload | null { if (event.event !== "receive" || !event.data) { return null; } @@ -80,7 +87,8 @@ function parseReceiveEnvelope(event: SignalSseEvent): SignalIngressEnvelope | nu "Signal receive event must contain a JSON object", ); } - return isRecord(parsed.envelope) ? (parsed.envelope as SignalIngressEnvelope) : null; + // SAFETY: SignalReceivePayload has only optional fields; downstream code validates each field. + return parsed as SignalReceivePayload; } function resolveDataMessage(envelope: SignalIngressEnvelope): Record | null { @@ -90,8 +98,11 @@ function resolveDataMessage(envelope: SignalIngressEnvelope): Record({ queue: ingressQueue, - inspect: (event) => inspectSignalIngressEvent(event), + inspect: (prepared) => inspectSignalIngressEvent(prepared), payload: { version: 1, - serialize: (event, { receivedAt }) => ({ receivedAt, event }), - deserialize: (body) => body.event, + // Parsed JSON remains transient; durable rows retain the exact raw event shape. + serialize: ([event], { receivedAt }) => ({ receivedAt, event }), + deserialize: (body) => [body.event, undefined], encode: ({ body }) => ({ version: 1, ...body }), decode: (payload) => ({ version: payload.version, body: payload }), createClaimError: (_kind, claim) => @@ -184,7 +196,8 @@ export async function startSignalIngressMonitor(params: { `Signal ingress row ${claim.id} has an invalid payload`, ), }, - deliver: (event, lifecycle) => params.dispatch(event, lifecycle), + deliver: ([event, parsedPayload], lifecycle) => + parsedPayload ? params.dispatch(event, lifecycle, parsedPayload) : undefined, onDurableAdmission: async (_event, { facts }) => { const { numberAliasEventId } = facts as SignalIngressEventFacts; if (!numberAliasEventId) { @@ -215,7 +228,7 @@ export async function startSignalIngressMonitor(params: { return { receive: async (event) => { - await monitor.admit(event); + await monitor.admit([event, undefined]); await monitor.waitForPumpIdle(); }, stop: monitor.stop,