fix(signal): deduplicate durable ingress across sender identity aliases (#118970)

This commit is contained in:
Peter Steinberger
2026-08-03 14:32:22 -07:00
committed by GitHub
parent 386228ad66
commit 5cddc85e55
2 changed files with 284 additions and 6 deletions
+265 -4
View File
@@ -25,23 +25,28 @@ async function startMonitor(queue: SignalIngressQueue, dispatch: SignalIngressDi
}
function signalEvent(params?: {
senderNumber?: string;
senderNumber?: string | null;
senderUuid?: string;
timestamp?: number;
groupId?: string;
message?: string;
reaction?: boolean;
}): SignalSseEvent {
const timestamp = params?.timestamp ?? 1_700_000_000_001;
return {
event: "receive",
data: JSON.stringify({
envelope: {
sourceNumber: params?.senderNumber ?? "+15550001111",
...(params?.senderNumber === null
? {}
: { sourceNumber: params?.senderNumber ?? "+15550001111" }),
...(params?.senderUuid ? { sourceUuid: params.senderUuid } : {}),
timestamp,
dataMessage: {
timestamp,
message: params?.message ?? "hello",
...(params?.reaction
? { reaction: { emoji: "👍", targetSentTimestamp: timestamp - 1 } }
: { message: params?.message ?? "hello" }),
...(params?.groupId ? { groupInfo: { groupId: params.groupId } } : {}),
},
},
@@ -221,10 +226,266 @@ describe("Signal durable ingress", () => {
});
});
it.each([
{ description: "direct phone-only delivery gains a UUID", phoneFirst: true },
{ description: "direct dual-identity delivery loses its UUID", phoneFirst: false },
{
description: "group phone-only delivery gains a UUID",
phoneFirst: true,
groupId: "group-123",
},
{
description: "group dual-identity delivery loses its UUID",
phoneFirst: false,
groupId: "group-123",
},
{
description: "approval reaction delivery gains a UUID",
phoneFirst: true,
reaction: true,
},
{
description: "approval reaction delivery loses its UUID",
phoneFirst: false,
reaction: true,
},
{
description: "group reaction delivery gains a UUID",
phoneFirst: true,
groupId: "group-123",
reaction: true,
},
{
description: "group reaction delivery loses its UUID",
phoneFirst: false,
groupId: "group-123",
reaction: true,
},
])("dedupes after restart when $description", async ({ phoneFirst, groupId, reaction }) => {
await withQueue(async (queue) => {
const shared = {
senderNumber: "+15550002222",
timestamp: 1_700_000_000_099,
...(groupId ? { groupId } : {}),
...(reaction ? { reaction } : {}),
};
const phoneOnly = signalEvent(shared);
const withUuid = signalEvent({
...shared,
senderUuid: "123e4567-e89b-12d3-a456-426614174000",
});
const dispatch = vi.fn().mockResolvedValue(undefined);
const initial = await startMonitor(queue, dispatch);
await initial.monitor.receive(phoneFirst ? phoneOnly : withUuid);
await initial.waitForIdle();
await initial.monitor.stop();
const restarted = await startMonitor(queue, dispatch);
try {
await restarted.monitor.receive(phoneFirst ? withUuid : phoneOnly);
await restarted.waitForIdle();
expect(dispatch).toHaveBeenCalledTimes(1);
} finally {
await restarted.monitor.stop();
}
});
});
it.each([
{
description: "phone-only",
params: { senderNumber: "+15550002222" },
numberAliases: 0,
},
{
description: "UUID-only",
params: {
senderNumber: null,
senderUuid: "123e4567-e89b-12d3-a456-426614174000",
},
numberAliases: 0,
},
{
description: "dual-identity",
params: {
senderNumber: "+15550002222",
senderUuid: "123e4567-e89b-12d3-a456-426614174000",
},
numberAliases: 1,
},
])("bounds completion aliases for $description senders", async ({ params, numberAliases }) => {
await withQueue(async (queue) => {
const complete = vi.spyOn(queue, "complete");
const started = await startMonitor(queue, vi.fn().mockResolvedValue(undefined));
try {
await started.monitor.receive(signalEvent(params));
await started.waitForIdle();
expect(complete.mock.calls.filter(([id]) => typeof id === "string")).toHaveLength(
numberAliases,
);
} finally {
await started.monitor.stop();
}
});
});
it("keeps the original durable message when identity-alias completion fails", async () => {
await withQueue(async (queue) => {
const aliasError = new Error("sqlite alias unavailable");
let failAlias = true;
const failingQueue = {
...queue,
complete: vi.fn<SignalIngressQueue["complete"]>(async (idOrClaim, options) => {
if (typeof idOrClaim === "string" && failAlias) {
failAlias = false;
throw aliasError;
}
return await queue.complete(idOrClaim, options);
}),
} satisfies SignalIngressQueue;
const withUuid = signalEvent({
senderNumber: "+15550002222",
senderUuid: "123e4567-e89b-12d3-a456-426614174000",
});
const dispatch = vi.fn().mockResolvedValue(undefined);
const started = await startMonitor(failingQueue, dispatch);
try {
await expect(started.monitor.receive(withUuid)).rejects.toBe(aliasError);
await started.waitForIdle();
expect(dispatch).toHaveBeenCalledTimes(1);
await started.monitor.receive(withUuid);
await started.monitor.receive(signalEvent({ senderNumber: "+15550002222" }));
await started.waitForIdle();
expect(dispatch).toHaveBeenCalledTimes(1);
} finally {
await started.monitor.stop();
}
});
});
it("retains the prior message window when dual-identity envelopes need two tombstones", async () => {
await withQueue(async (queue) => {
const prune = vi.spyOn(queue, "prune");
const started = await startMonitor(queue, vi.fn().mockResolvedValue(undefined));
try {
await started.monitor.receive(
signalEvent({ senderUuid: "123e4567-e89b-12d3-a456-426614174000" }),
);
await started.waitForIdle();
expect(prune).toHaveBeenCalledWith(
expect.objectContaining({
completedMaxEntries: 2_000,
completedTtlMs: 30 * 24 * 60 * 60 * 1_000,
}),
);
} finally {
await started.monitor.stop();
}
});
});
it.each([true, false])(
"serializes concurrent sender-identity aliases when phone-only arrives first: %s",
async (phoneFirst) => {
await withQueue(async (queue) => {
const phoneOnly = signalEvent({ senderNumber: "+15550002222" });
const withUuid = signalEvent({
senderNumber: "+15550002222",
senderUuid: "123e4567-e89b-12d3-a456-426614174000",
});
const dispatch = vi.fn().mockResolvedValue(undefined);
const started = await startMonitor(queue, dispatch);
try {
await Promise.all(
(phoneFirst ? [phoneOnly, withUuid] : [withUuid, phoneOnly]).map((event) =>
started.monitor.receive(event),
),
);
await started.waitForIdle();
expect(dispatch).toHaveBeenCalledTimes(1);
} finally {
await started.monitor.stop();
}
});
},
);
it.each([true, false])(
"does not double-dispatch an adopted claim when phone-only delivery comes first: %s",
async (phoneFirst) => {
await withQueue(async (queue) => {
const phoneOnly = signalEvent({ senderNumber: "+15550002222" });
const withUuid = signalEvent({
senderNumber: "+15550002222",
senderUuid: "123e4567-e89b-12d3-a456-426614174000",
});
let adopt: (() => void | Promise<void>) | undefined;
const dispatch = vi.fn((_event, lifecycle) => {
adopt = lifecycle.onAdopted;
lifecycle.onDeferred();
return { kind: "deferred" } as const;
});
const started = await startMonitor(queue, dispatch);
try {
await started.monitor.receive(phoneFirst ? phoneOnly : withUuid);
await started.waitForIdle();
expect(await queue.listClaims()).toHaveLength(1);
await started.monitor.receive(phoneFirst ? withUuid : phoneOnly);
await started.waitForIdle();
expect(dispatch).toHaveBeenCalledTimes(1);
expect(await queue.listClaims()).toHaveLength(1);
await adopt?.();
await started.waitForIdle();
expect(dispatch).toHaveBeenCalledTimes(1);
} finally {
await started.monitor.stop();
}
});
},
);
it("keeps identity-alias tombstones scoped to their Signal account", async () => {
await withQueue(async (queue, stateDir) => {
const otherQueue = createChannelIngressQueueForTests<SignalIngressPayload>({
channelId: "signal",
accountId: "other",
stateDir,
});
const firstDispatch = vi.fn().mockResolvedValue(undefined);
const otherDispatch = vi.fn().mockResolvedValue(undefined);
const first = await startMonitor(queue, firstDispatch);
const other = await startSignalIngressMonitor({
accountId: "other",
queue: otherQueue,
dispatch: otherDispatch,
runtime: { error: vi.fn(), log: vi.fn() },
});
try {
await first.monitor.receive(
signalEvent({
senderNumber: "+15550002222",
senderUuid: "123e4567-e89b-12d3-a456-426614174000",
}),
);
await other.receive(signalEvent({ senderNumber: "+15550002222" }));
await first.waitForIdle();
await other.waitForIdle();
expect(firstDispatch).toHaveBeenCalledTimes(1);
expect(otherDispatch).toHaveBeenCalledTimes(1);
} finally {
await first.monitor.stop();
await other.stop();
}
});
});
it("uses a direct-sender or group-conversation lane and stores the raw event", async () => {
await withQueue(async (queue) => {
const direct = signalEvent({ senderUuid: "123e4567-e89b-12d3-a456-426614174000" });
const group = signalEvent({ groupId: "group-123" });
const group = signalEvent({ groupId: "group-123", timestamp: 1_700_000_000_002 });
const dispatch = vi.fn((_event, lifecycle) => {
lifecycle.onDeferred();
return { kind: "deferred" } as const;
+19 -2
View File
@@ -27,6 +27,7 @@ type SignalIngressEnvelope = {
type SignalIngressEventFacts = {
eventId: string;
laneKey: string;
numberAliasEventId?: string;
};
type SignalIngressPayload = {
@@ -125,6 +126,9 @@ function inspectSignalIngressEvent(event: SignalSseEvent): SignalIngressEventFac
return {
eventId: JSON.stringify([senderKey, timestamp]),
laneKey: groupId ? `group:${groupId}` : `direct:${senderKey}`,
...(senderUuid && senderNumber
? { numberAliasEventId: JSON.stringify([`number:${senderNumber}`, timestamp]) }
: {}),
};
}
@@ -157,12 +161,13 @@ export async function startSignalIngressMonitor(params: {
accountId: params.accountId,
});
}
const ingressQueue = queue;
const monitor = createChannelIngressMonitor<
SignalSseEvent,
SignalIngressBody,
SignalIngressPayload
>({
queue,
queue: ingressQueue,
inspect: (event) => inspectSignalIngressEvent(event),
payload: {
version: 1,
@@ -177,11 +182,23 @@ export async function startSignalIngressMonitor(params: {
),
},
deliver: (event, lifecycle) => params.dispatch(event, lifecycle),
onDurableAdmission: async (_event, { facts }) => {
const { numberAliasEventId } = facts as SignalIngressEventFacts;
if (!numberAliasEventId) {
return;
}
// signal-cli can learn or forget a UUID between redeliveries; bridge both
// shipped sender IDs before the monitor releases its admission/claim lock.
if (!(await ingressQueue.complete(numberAliasEventId))) {
await ingressQueue.complete(facts.eventId);
}
},
pollIntervalMs: SIGNAL_INGRESS_DRAIN_INTERVAL_MS,
retention: {
// Signal previously pruned before every enqueue rather than on a timed cadence.
pruneIntervalMs: 0,
completedMaxEntries: 1_000,
// At most two tombstones per message preserve the prior 1,000-message window.
completedMaxEntries: 2_000,
failedMaxEntries: 1_000,
},
appendRetryDelaysMs: [0],