fix: 原子化扇出覆盖并对齐事件索引

This commit is contained in:
sunlei 2026-07-24 11:51:30 +08:00
parent 314221518d
commit f18c0dc398
7 changed files with 436 additions and 42 deletions

View File

@ -1143,7 +1143,7 @@ CREATE TABLE IF NOT EXISTS `qqbot_message_event` (
`payload` json NOT NULL, `fanout_status` varchar(32) NOT NULL DEFAULT 'accepted', `fanout_attempt_count` int unsigned NOT NULL DEFAULT 0,
`next_fanout_at` datetime(6) DEFAULT NULL, `fanout_lease_until` datetime(6) DEFAULT NULL, `last_error_code` varchar(64) DEFAULT NULL, `last_error_message` varchar(500) DEFAULT NULL,
`create_time` datetime(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6), `update_time` datetime(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6) ON UPDATE CURRENT_TIMESTAMP(6),
PRIMARY KEY (`id`), UNIQUE KEY `uk_qqbot_message_event_event_id` (`event_id`), KEY `idx_qqbot_message_event_dispatch` (`fanout_status`, `next_fanout_at`), KEY `idx_qqbot_message_event_lease` (`fanout_lease_until`)
PRIMARY KEY (`id`), UNIQUE KEY `uk_qqbot_message_event_event_id` (`event_id`), KEY `idx_qqbot_message_event_dispatch` (`fanout_status`, `next_fanout_at`), KEY `idx_qqbot_message_event_lease` (`fanout_lease_until`), KEY `idx_qqbot_message_event_source_resource_order` (`source_key`, `resource_key`, `occurred_at`, `id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci;
CREATE TABLE IF NOT EXISTS `qqbot_message_delivery` (

View File

@ -1315,7 +1315,8 @@ CREATE TABLE IF NOT EXISTS qqbot_message_event (
update_time DATETIME(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6) ON UPDATE CURRENT_TIMESTAMP(6),
UNIQUE KEY uk_qqbot_message_event_event_id (event_id),
KEY idx_qqbot_message_event_dispatch (fanout_status, next_fanout_at),
KEY idx_qqbot_message_event_lease (fanout_lease_until)
KEY idx_qqbot_message_event_lease (fanout_lease_until),
KEY idx_qqbot_message_event_source_resource_order (source_key, resource_key, occurred_at, id)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci;
CREATE TABLE IF NOT EXISTS qqbot_message_delivery (

View File

@ -1,5 +1,5 @@
import { Injectable } from '@nestjs/common';
import { DataSource, Brackets, type EntityManager } from 'typeorm';
import { DataSource, Brackets, In, type EntityManager } from 'typeorm';
import { KtDateTime } from '@/common';
import {
SystemMessageContractError,
@ -29,7 +29,7 @@ const STUN_MAPPING_PORT_SOURCE = 'network.stun.mapping-port-changed';
const TRANSIENT_ERROR_CODE = 'fanout_transient_error';
const EVENT_EXPIRED_ERROR_CODE = 'fanout_expired';
const EVENT_RESOURCE_MISMATCH_ERROR_CODE = 'event_resource_mismatch';
const SUPERSEDED_STATUSES = new Set(['waiting_ddns', 'pending', 'retry']);
const SUPERSEDED_STATUSES = ['waiting_ddns', 'pending', 'retry'];
interface ClaimToken {
attempt: number;
@ -297,6 +297,8 @@ export class SystemMessageFanoutService {
);
}),
)
.orderBy('newerEvent.occurredAt', 'ASC')
.addOrderBy('newerEvent.id', 'ASC')
.take(1)
.getOne();
return !!newerEvent;
@ -314,15 +316,14 @@ export class SystemMessageFanoutService {
subscriptionId: string,
): Promise<void> {
const deliveries = manager.getRepository(QqbotMessageDelivery);
const candidates = await deliveries.find({
where: { messageEventId: event.id, subscriptionId },
});
for (const delivery of candidates) {
if (SUPERSEDED_STATUSES.has(delivery.status)) {
delivery.status = 'superseded';
await deliveries.save(delivery);
}
}
await deliveries.update(
{
messageEventId: event.id,
status: In(SUPERSEDED_STATUSES),
subscriptionId,
},
{ status: 'superseded' },
);
}
/**
@ -340,24 +341,20 @@ export class SystemMessageFanoutService {
const priorEvents = await events.find({
where: { resourceKey: event.resourceKey, sourceKey: event.sourceKey },
});
const priorIds = new Set(
priorEvents
.filter((candidate) => this.isStrictlyEarlier(candidate, event))
.map((candidate) => candidate.id),
);
if (!priorIds.size) return;
const priorIds = priorEvents
.filter((candidate) => this.isStrictlyEarlier(candidate, event))
.map((candidate) => candidate.id);
if (!priorIds.length) return;
const deliveries = manager.getRepository(QqbotMessageDelivery);
const candidates = await deliveries.find({ where: { subscriptionId } });
for (const delivery of candidates) {
if (
priorIds.has(delivery.messageEventId) &&
SUPERSEDED_STATUSES.has(delivery.status)
) {
delivery.status = 'superseded';
await deliveries.save(delivery);
}
}
await deliveries.update(
{
messageEventId: In(priorIds),
status: In(SUPERSEDED_STATUSES),
subscriptionId,
},
{ status: 'superseded' },
);
}
/**

View File

@ -6,6 +6,7 @@ import type { SystemMessageFanoutStatus, SystemMessageScalar } from '../../../co
@Index('uk_qqbot_message_event_event_id', ['eventId'], { unique: true })
@Index('idx_qqbot_message_event_dispatch', ['fanoutStatus', 'nextFanoutAt'])
@Index('idx_qqbot_message_event_lease', ['fanoutLeaseUntil'])
@Index('idx_qqbot_message_event_source_resource_order', ['sourceKey', 'resourceKey', 'occurredAt', 'id'])
export class QqbotMessageEvent {
@PrimaryColumn({ type: 'bigint' }) id: string;
@Column({ length: 128, name: 'event_id' }) eventId: string;

View File

@ -92,6 +92,7 @@ const persistenceContract: ReadonlyArray<{
{ name: 'uk_qqbot_message_event_event_id', columns: ['eventId'], unique: true },
{ name: 'idx_qqbot_message_event_dispatch', columns: ['fanoutStatus', 'nextFanoutAt'], unique: false },
{ name: 'idx_qqbot_message_event_lease', columns: ['fanoutLeaseUntil'], unique: false },
{ name: 'idx_qqbot_message_event_source_resource_order', columns: ['sourceKey', 'resourceKey', 'occurredAt', 'id'], unique: false },
],
},
{

View File

@ -70,6 +70,7 @@ const messagePushTableDefinitions = {
'primary key (id)', 'unique key uk_qqbot_message_event_event_id (event_id)',
'key idx_qqbot_message_event_dispatch (fanout_status, next_fanout_at)',
'key idx_qqbot_message_event_lease (fanout_lease_until)',
'key idx_qqbot_message_event_source_resource_order (source_key, resource_key, occurred_at, id)',
],
},
qqbot_message_delivery: {

View File

@ -182,6 +182,16 @@ function matches(
): boolean {
return Object.entries(where).every(([key, value]) => {
const actual = row[key];
if (
value &&
typeof value === 'object' &&
'_type' in value &&
(value as { _type?: unknown })._type === 'in'
) {
return (value as unknown as { _value: unknown[] })._value.includes(
actual,
);
}
if (
actual &&
value &&
@ -197,6 +207,29 @@ function matches(
});
}
/** Compares transaction rows without relying on their entity object identity. */
function hasRowChanged(
current: Record<string, unknown> | undefined,
baseline: Record<string, unknown> | undefined,
): boolean {
return JSON.stringify(current) !== JSON.stringify(baseline);
}
/** Normalizes TypeORM's in-memory find operators for independent update assertions. */
function normalizeConditionalWhere(where: Record<string, unknown>) {
return Object.fromEntries(
Object.entries(where).map(([key, value]) => [
key,
value &&
typeof value === 'object' &&
'_type' in value &&
(value as { _type?: unknown })._type === 'in'
? (value as unknown as { _value: unknown[] })._value
: value,
]),
);
}
/** Creates a transaction-faithful in-memory persistence harness rather than fan-out logic. */
function setup(seed: Partial<Store> = {}) {
let state: Store = {
@ -225,10 +258,18 @@ function setup(seed: Partial<Store> = {}) {
predicates: RecordedPredicate[];
take: number;
}> = [];
const newerEventRangeScans: Array<{
eventId: string;
scannedEventIds: string[];
}> = [];
const transactions: string[] = [];
const locked = new Set<string>();
const activeClaimLocks = new Set<string>();
const deliveryReadLocks: Array<null | string> = [];
const deliveryConditionalUpdates: Array<{
values: Record<string, unknown>;
where: Record<string, unknown>;
}> = [];
const savedDeliveryTargets: string[] = [];
let deliverySequence = 1000;
let duplicateRace: 'exact' | 'wrong' | null = null;
@ -243,6 +284,45 @@ function setup(seed: Partial<Store> = {}) {
entered: () => void;
resume: Promise<void>;
} = null;
let beforeDeliveryConditionalUpdate:
| null
| ((deliveries: QqbotMessageDelivery[]) => void) = null;
/** Merges one committed transaction draft without losing separately committed rows. */
const mergeRows = <T extends { id: string }>(
authoritative: T[],
draft: T[],
baseline: T[],
): T[] => {
const merged = new Map(authoritative.map((row) => [row.id, row]));
const baselineById = new Map(baseline.map((row) => [row.id, row]));
for (const row of draft) {
const baselineRow = baselineById.get(row.id);
if (!baselineRow || hasRowChanged(row, baselineRow)) {
merged.set(row.id, row);
}
}
return [...merged.values()];
};
/** Merges locally changed transaction state with authoritative concurrent commits. */
const mergeCommittedStore = (draft: Store, baseline: Store): Store => ({
accounts: mergeRows(state.accounts, draft.accounts, baseline.accounts),
bindings: mergeRows(state.bindings, draft.bindings, baseline.bindings),
deliveries: mergeRows(
state.deliveries,
draft.deliveries,
baseline.deliveries,
),
events: mergeRows(state.events, draft.events, baseline.events),
subscriptions: mergeRows(
state.subscriptions,
draft.subscriptions,
baseline.subscriptions,
),
targets: mergeRows(state.targets, draft.targets, baseline.targets),
templates: mergeRows(state.templates, draft.templates, baseline.templates),
});
/** Returns a repository facade bound to exactly one authoritative or transaction-local store. */
const repository = (
@ -343,8 +423,8 @@ function setup(seed: Partial<Store> = {}) {
const currentOccurredAt = (
occurredAt as { getTime: () => number }
).getTime();
return (
store.events.find((candidate) => {
const candidates = state.events
.filter((candidate) => {
return (
candidate.sourceKey === sourceKey &&
candidate.resourceKey === resourceKey &&
@ -352,8 +432,22 @@ function setup(seed: Partial<Store> = {}) {
(candidate.occurredAt.getTime() === currentOccurredAt &&
BigInt(candidate.id) > BigInt(eventId)))
);
}) ?? null
);
})
.sort((left, right) => {
const occurredAtDifference =
left.occurredAt.getTime() - right.occurredAt.getTime();
if (occurredAtDifference) return occurredAtDifference;
return BigInt(left.id) < BigInt(right.id)
? -1
: BigInt(left.id) > BigInt(right.id)
? 1
: 0;
});
newerEventRangeScans.push({
eventId,
scannedEventIds: candidates.map((candidate) => candidate.id),
});
return candidates[0] ?? null;
}
const claimed =
rows
@ -488,8 +582,10 @@ function setup(seed: Partial<Store> = {}) {
activeClaimLocks.add(id);
transactionLocks?.add(id);
if (pauseEventLock) {
pauseEventLock.entered();
await pauseEventLock.resume;
const pause = pauseEventLock;
pauseEventLock = null;
pause.entered();
await pause.resume;
}
return event;
}
@ -548,6 +644,26 @@ function setup(seed: Partial<Store> = {}) {
where: Record<string, unknown>,
values: Record<string, unknown>,
) => {
if (entity === QqbotMessageDelivery) {
deliveryConditionalUpdates.push({
values: structuredClone(values),
where: normalizeConditionalWhere(where),
});
const mutate = beforeDeliveryConditionalUpdate;
beforeDeliveryConditionalUpdate = null;
mutate?.(state.deliveries);
const matched = state.deliveries.filter((candidate) =>
matches(candidate as unknown as Record<string, unknown>, where),
);
for (const delivery of matched) {
const draftDelivery = rows.find(
(candidate) => candidate.id === delivery.id,
);
if (!draftDelivery) continue;
Object.assign(draftDelivery, structuredClone(delivery), values);
}
return { affected: matched.length };
}
const row = rows.find((candidate) => matches(candidate, where));
if (!row) return { affected: 0 };
Object.assign(row as object, values);
@ -562,6 +678,7 @@ function setup(seed: Partial<Store> = {}) {
getRepository: (entity: unknown) => ReturnType<typeof repository>;
}) => Promise<unknown>,
) => {
const baseline = cloneStore(state);
const draft = cloneStore(state);
const transactionLocks = new Set<string>();
transactions.push('begin');
@ -575,7 +692,7 @@ function setup(seed: Partial<Store> = {}) {
draft.deliveries.push(structuredClone(delivery));
}
}
state = draft;
state = mergeCommittedStore(draft, baseline);
transactions.push('commit');
return result;
} catch (error) {
@ -596,6 +713,7 @@ function setup(seed: Partial<Store> = {}) {
bindings: () => state.bindings,
events: () => state.events,
deliveries: () => state.deliveries,
deliveryConditionalUpdates,
deliveryReadLocks,
failSubscription: (id: null | string) => {
failSubscription = id;
@ -605,6 +723,7 @@ function setup(seed: Partial<Store> = {}) {
},
lock: (id: string) => locked.add(id),
newerEventReads,
newerEventRangeScans,
query,
savedDeliveryTargets: () => savedDeliveryTargets,
pauseNextClaim: () => {
@ -631,6 +750,15 @@ function setup(seed: Partial<Store> = {}) {
pauseEventLock = { entered, resume: wait };
return { reached, resume };
},
runBeforeNextDeliveryConditionalUpdate: (
mutate: (deliveries: QqbotMessageDelivery[]) => void,
) => {
beforeDeliveryConditionalUpdate = mutate;
},
commitEvent: (committed: QqbotMessageEvent) => {
state.events.push(committed);
},
activeEventLocks: () => [...activeClaimLocks],
setDuplicateRace: (value: 'exact' | 'wrong' | null) => {
duplicateRace = value;
},
@ -804,6 +932,109 @@ describe('SystemMessageFanoutService', () => {
).toBe('pending');
});
it('fences earlier-history supersession against authoritative processing and cancelled races', async () => {
const leaseUntil = new KtDateTime(NOW.getTime() + SYSTEM_MESSAGE_LEASE_MS);
const older = event({
fanoutStatus: 'completed',
id: '100',
nextFanoutAt: null,
occurredAt: new KtDateTime(NOW.getTime() - 1),
});
const current = event({
fanoutAttemptCount: 1,
fanoutLeaseUntil: leaseUntil,
fanoutStatus: 'processing',
});
const normal = Object.assign(new QqbotMessageDelivery(), {
attemptCount: 3,
bindingId: 'binding-normal',
expiresAt: new KtDateTime(NOW.getTime() + 1000),
id: 'normal',
lastErrorCode: 'old-error',
lastErrorMessage: 'old message',
messageEventId: older.id,
nextAttemptAt: new KtDateTime(NOW.getTime() + 2000),
processingLeaseUntil: new KtDateTime(NOW.getTime() + 3000),
publishTargetId: 'normal-target',
renderedMessage: 'normal',
selfId: 'bot-a',
sendLogId: 'send-normal',
status: 'waiting_ddns' as const,
subscriptionId: '300',
targetId: 'target-normal',
targetType: 'group' as const,
templateContent: 'normal-template',
templateId: 'template-normal',
variableSnapshot: { endpoint: 'normal' },
});
const racedProcessing = Object.assign(new QqbotMessageDelivery(), {
...structuredClone(normal),
id: 'raced-processing',
publishTargetId: 'processing-target',
status: 'pending' as const,
});
const racedCancelled = Object.assign(new QqbotMessageDelivery(), {
...structuredClone(normal),
id: 'raced-cancelled',
publishTargetId: 'cancelled-target',
status: 'retry' as const,
});
const fixture = setup({
deliveries: [normal, racedProcessing, racedCancelled],
events: [older, current],
});
const expectedProcessing = structuredClone(
racedProcessing,
) as QqbotMessageDelivery;
expectedProcessing.status = 'processing';
const expectedCancelled = structuredClone(
racedCancelled,
) as QqbotMessageDelivery;
expectedCancelled.status = 'cancelled';
fixture.runBeforeNextDeliveryConditionalUpdate((deliveries) => {
deliveries.find(
(delivery) => delivery.id === racedProcessing.id,
)!.status = 'processing';
deliveries.find((delivery) => delivery.id === racedCancelled.id)!.status =
'cancelled';
});
const processClaim = (
fixture.service as unknown as {
processClaim: (token: object, now: Date) => Promise<void>;
}
).processClaim;
await processClaim.call(
fixture.service,
{ attempt: 1, event: current, leaseUntil },
NOW,
);
expect(
fixture.deliveries().find((delivery) => delivery.id === normal.id),
).toMatchObject({ status: 'superseded' });
expect(fixture.deliveryConditionalUpdates).toEqual([
{
values: { status: 'superseded' },
where: {
messageEventId: [older.id],
status: ['waiting_ddns', 'pending', 'retry'],
subscriptionId: '300',
},
},
]);
expect(
fixture
.deliveries()
.find((delivery) => delivery.id === racedProcessing.id),
).toEqual(expectedProcessing);
expect(
fixture
.deliveries()
.find((delivery) => delivery.id === racedCancelled.id),
).toEqual(expectedCancelled);
});
it('does not recreate old A work after committed A-to-B-to-A events return to A', async () => {
const oldLease = new KtDateTime(NOW.getTime() + SYSTEM_MESSAGE_LEASE_MS);
const oldEvent = event({
@ -933,7 +1164,7 @@ describe('SystemMessageFanoutService', () => {
expect(fixture.newerEventReads).toEqual([
{
lock: 'pessimistic_read',
order: [],
order: ['newerEvent.occurredAt', 'newerEvent.id'],
predicates: [
{
expression:
@ -962,6 +1193,172 @@ describe('SystemMessageFanoutService', () => {
]);
});
it('fences current-event supersession against authoritative processing and cancelled races', async () => {
const leaseUntil = new KtDateTime(NOW.getTime() + SYSTEM_MESSAGE_LEASE_MS);
const older = event({
fanoutAttemptCount: 1,
fanoutLeaseUntil: leaseUntil,
fanoutStatus: 'processing',
id: '100',
occurredAt: new KtDateTime(NOW.getTime() - 1),
});
const newer = event({
fanoutStatus: 'completed',
id: '200',
nextFanoutAt: null,
});
const normal = Object.assign(new QqbotMessageDelivery(), {
attemptCount: 4,
bindingId: 'binding-normal',
expiresAt: new KtDateTime(NOW.getTime() + 1000),
id: 'normal',
lastErrorCode: 'old-error',
lastErrorMessage: 'old message',
messageEventId: older.id,
nextAttemptAt: new KtDateTime(NOW.getTime() + 2000),
processingLeaseUntil: new KtDateTime(NOW.getTime() + 3000),
publishTargetId: 'normal-target',
renderedMessage: 'normal',
selfId: 'bot-a',
sendLogId: 'send-normal',
status: 'pending' as const,
subscriptionId: '300',
targetId: 'target-normal',
targetType: 'group' as const,
templateContent: 'normal-template',
templateId: 'template-normal',
variableSnapshot: { endpoint: 'normal' },
});
const racedProcessing = Object.assign(new QqbotMessageDelivery(), {
...structuredClone(normal),
id: 'raced-processing',
publishTargetId: 'processing-target',
status: 'waiting_ddns' as const,
});
const racedCancelled = Object.assign(new QqbotMessageDelivery(), {
...structuredClone(normal),
id: 'raced-cancelled',
publishTargetId: 'cancelled-target',
status: 'retry' as const,
});
const fixture = setup({
deliveries: [normal, racedProcessing, racedCancelled],
events: [older, newer],
});
const expectedProcessing = structuredClone(
racedProcessing,
) as QqbotMessageDelivery;
expectedProcessing.status = 'processing';
const expectedCancelled = structuredClone(
racedCancelled,
) as QqbotMessageDelivery;
expectedCancelled.status = 'cancelled';
fixture.runBeforeNextDeliveryConditionalUpdate((deliveries) => {
deliveries.find(
(delivery) => delivery.id === racedProcessing.id,
)!.status = 'processing';
deliveries.find((delivery) => delivery.id === racedCancelled.id)!.status =
'cancelled';
});
const processClaim = (
fixture.service as unknown as {
processClaim: (token: object, now: Date) => Promise<void>;
}
).processClaim;
await processClaim.call(
fixture.service,
{ attempt: 1, event: older, leaseUntil },
NOW,
);
expect(
fixture.deliveries().find((delivery) => delivery.id === normal.id),
).toMatchObject({ status: 'superseded' });
expect(fixture.deliveryConditionalUpdates).toEqual([
{
values: { status: 'superseded' },
where: {
messageEventId: older.id,
status: ['waiting_ddns', 'pending', 'retry'],
subscriptionId: '300',
},
},
]);
expect(
fixture
.deliveries()
.find((delivery) => delivery.id === racedProcessing.id),
).toEqual(expectedProcessing);
expect(
fixture
.deliveries()
.find((delivery) => delivery.id === racedCancelled.id),
).toEqual(expectedCancelled);
});
it('uses an indexed monotone newer-event range so a newer owner can finish before an older owner observes it', async () => {
const olderLease = new KtDateTime(NOW.getTime() + SYSTEM_MESSAGE_LEASE_MS);
const newerLease = new KtDateTime(NOW.getTime() + SYSTEM_MESSAGE_LEASE_MS);
const older = event({
fanoutAttemptCount: 1,
fanoutLeaseUntil: olderLease,
fanoutStatus: 'processing',
id: '100',
});
const newer = event({
fanoutAttemptCount: 1,
fanoutLeaseUntil: newerLease,
fanoutStatus: 'processing',
id: '200',
nextFanoutAt: null,
occurredAt: new KtDateTime(NOW.getTime() + 1),
});
const fixture = setup({ events: [older] });
const processClaim = (
fixture.service as unknown as {
processClaim: (token: object, now: Date) => Promise<void>;
}
).processClaim;
const pause = fixture.pauseNextEventLock();
const olderRun = processClaim.call(
fixture.service,
{ attempt: 1, event: older, leaseUntil: olderLease },
NOW,
);
await pause.reached;
expect(fixture.activeEventLocks()).toEqual(['100']);
fixture.commitEvent(newer);
await processClaim.call(
fixture.service,
{ attempt: 1, event: newer, leaseUntil: newerLease },
NOW,
);
expect(fixture.adapter.resolveDelivery).toHaveBeenCalledTimes(1);
expect(fixture.newerEventRangeScans).toEqual([
{ eventId: newer.id, scannedEventIds: [] },
]);
pause.resume();
await olderRun;
expect(fixture.adapter.resolveDelivery).toHaveBeenCalledTimes(1);
expect(fixture.newerEventRangeScans).toEqual([
{ eventId: newer.id, scannedEventIds: [] },
{ eventId: older.id, scannedEventIds: [newer.id] },
]);
expect(fixture.newerEventReads).toEqual(
expect.arrayContaining([
expect.objectContaining({
lock: 'pessimistic_read',
order: ['newerEvent.occurredAt', 'newerEvent.id'],
take: 1,
}),
]),
);
});
it('treats same-timestamp BIGINT id 10 as newer than lexical-trap id 9', async () => {
const oldLease = new KtDateTime(NOW.getTime() + SYSTEM_MESSAGE_LEASE_MS);
const oldEvent = event({
@ -1311,11 +1708,7 @@ describe('SystemMessageFanoutService', () => {
'old-failed-target',
'700',
]);
expect(fixture.savedDeliveryTargets()).toEqual([
'700',
'old-failed-target',
'701',
]);
expect(fixture.savedDeliveryTargets()).toEqual(['700', '701']);
expect(fixture.events()[1].fanoutStatus).toBe('retry');
});