perf(signal): reuse parsed inbound payload (#127851)

Co-authored-by: Amp <amp@ampcode.com>
This commit is contained in:
Peter Steinberger
2026-08-22 03:27:28 -07:00
committed by GitHub
parent 39cff92880
commit 2d0fff4ac5
3 changed files with 35 additions and 17 deletions
@@ -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}`);
+3 -1
View File
@@ -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();
}
+23 -10
View File
@@ -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> | 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<string, unknown> | null {
@@ -90,8 +98,11 @@ function resolveDataMessage(envelope: SignalIngressEnvelope): Record<string, unk
return isRecord(envelope.editMessage?.dataMessage) ? envelope.editMessage.dataMessage : null;
}
function inspectSignalIngressEvent(event: SignalSseEvent): SignalIngressEventFacts | null {
const envelope = parseReceiveEnvelope(event);
function inspectSignalIngressEvent(
prepared: SignalPreparedIngressEvent,
): SignalIngressEventFacts | null {
const payload = (prepared[1] ??= parseReceivePayload(prepared[0]));
const envelope = isRecord(payload?.envelope) ? (payload.envelope as SignalIngressEnvelope) : null;
if (!envelope || "syncMessage" in envelope) {
return null;
}
@@ -166,16 +177,17 @@ export async function startSignalIngressMonitor(params: {
}
const ingressQueue = queue;
const monitor = createChannelIngressMonitor<
SignalSseEvent,
SignalPreparedIngressEvent,
SignalIngressBody,
SignalIngressPayload
>({
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,