fix: 阻止迟到旧事件重建投递

This commit is contained in:
sunlei 2026-07-24 11:27:00 +08:00
parent fb50b0b8c4
commit 314221518d
2 changed files with 348 additions and 15 deletions

View File

@ -236,6 +236,15 @@ export class SystemMessageFanoutService {
}); });
if (!event || !this.ownsClaim(event, token)) return 'stale_claim'; if (!event || !this.ownsClaim(event, token)) return 'stale_claim';
if (await this.hasStrictlyNewerEvent(manager, event)) {
await this.supersedeCurrentEventDeliveries(
manager,
event,
subscriptionId,
);
return 'handled';
}
const subscriptions = manager.getRepository(QqbotMessageSubscription); const subscriptions = manager.getRepository(QqbotMessageSubscription);
const subscription = await subscriptions.findOne({ const subscription = await subscriptions.findOne({
where: { id: subscriptionId }, where: { id: subscriptionId },
@ -258,6 +267,64 @@ export class SystemMessageFanoutService {
return 'handled'; return 'handled';
} }
/**
* Uses a locking current read to detect committed newer facts before any subscription work.
* @param manager - Current subscription transaction manager that already owns the event lock.
* @param event - Locked claimed event that must not create work after a later fact exists.
* @returns Whether an exact source/resource event is strictly newer by occurrence then BIGINT ID.
*/
private async hasStrictlyNewerEvent(
manager: EntityManager,
event: QqbotMessageEvent,
): Promise<boolean> {
const newerEvent = await manager
.getRepository(QqbotMessageEvent)
.createQueryBuilder('newerEvent')
.setLock('pessimistic_read')
.where(
'newerEvent.sourceKey = :sourceKey AND newerEvent.resourceKey = :resourceKey',
{ resourceKey: event.resourceKey, sourceKey: event.sourceKey },
)
.andWhere(
new Brackets((where) => {
where
.where('newerEvent.occurredAt > :occurredAt', {
occurredAt: event.occurredAt,
})
.orWhere(
'newerEvent.occurredAt = :occurredAt AND newerEvent.id > :eventId',
{ eventId: event.id, occurredAt: event.occurredAt },
);
}),
)
.take(1)
.getOne();
return !!newerEvent;
}
/**
* Supersedes only the current old event's mutable rows after a newer fact is committed.
* @param manager - Current subscription transaction manager.
* @param event - Older locked event whose existing unfinished rows are being retired.
* @param subscriptionId - Exact subscription scope that excludes newer and unrelated work.
*/
private async supersedeCurrentEventDeliveries(
manager: EntityManager,
event: QqbotMessageEvent,
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);
}
}
}
/** /**
* Supersedes only unfinished deliveries of events strictly earlier than the current event. * Supersedes only unfinished deliveries of events strictly earlier than the current event.
* @param manager - Current subscription transaction manager. * @param manager - Current subscription transaction manager.

View File

@ -219,6 +219,12 @@ function setup(seed: Partial<Store> = {}) {
predicates: [] as RecordedPredicate[], predicates: [] as RecordedPredicate[],
take: 0, take: 0,
}; };
const newerEventReads: Array<{
lock: string;
order: string[];
predicates: RecordedPredicate[];
take: number;
}> = [];
const transactions: string[] = []; const transactions: string[] = [];
const locked = new Set<string>(); const locked = new Set<string>();
const activeClaimLocks = new Set<string>(); const activeClaimLocks = new Set<string>();
@ -262,9 +268,16 @@ function setup(seed: Partial<Store> = {}) {
return { return {
create: (input: object) => create: (input: object) =>
Object.assign(new (entity as new () => object)(), input), Object.assign(new (entity as new () => object)(), input),
createQueryBuilder: () => { createQueryBuilder: (alias = 'event') => {
let builderLock = ''; let builderLock = '';
let queryNow: Date | null = null; let queryNow: Date | null = null;
const newerEventRead = {
lock: '',
order: [] as string[],
predicates: [] as RecordedPredicate[],
take: 0,
};
const newerEventParameters: Record<string, unknown> = {};
/** Records the nested TypeORM predicate structure and its bound values. */ /** Records the nested TypeORM predicate structure and its bound values. */
const recordBrackets = (value: { const recordBrackets = (value: {
whereFactory?: (where: { whereFactory?: (where: {
@ -281,6 +294,9 @@ function setup(seed: Partial<Store> = {}) {
if (typeof expression === 'string') { if (typeof expression === 'string') {
const now = (parameters as { now?: unknown }).now; const now = (parameters as { now?: unknown }).now;
if (now instanceof Date) queryNow = now; if (now instanceof Date) queryNow = now;
if (alias === 'newerEvent') {
Object.assign(newerEventParameters, parameters);
}
predicates.push({ expression, parameters }); predicates.push({ expression, parameters });
return recorder; return recorder;
} }
@ -303,10 +319,42 @@ function setup(seed: Partial<Store> = {}) {
}; };
const builder = { const builder = {
addOrderBy: (field: string) => { addOrderBy: (field: string) => {
query.order.push(field); if (alias === 'newerEvent') newerEventRead.order.push(field);
else query.order.push(field);
return builder; return builder;
}, },
getOne: async () => { getOne: async () => {
if (alias === 'newerEvent') {
newerEventReads.push(structuredClone(newerEventRead));
const sourceKey = newerEventParameters.sourceKey;
const resourceKey = newerEventParameters.resourceKey;
const occurredAt = newerEventParameters.occurredAt;
const eventId = newerEventParameters.eventId;
if (
typeof sourceKey !== 'string' ||
typeof resourceKey !== 'string' ||
!occurredAt ||
typeof (occurredAt as { getTime?: unknown }).getTime !==
'function' ||
typeof eventId !== 'string'
) {
return null;
}
const currentOccurredAt = (
occurredAt as { getTime: () => number }
).getTime();
return (
store.events.find((candidate) => {
return (
candidate.sourceKey === sourceKey &&
candidate.resourceKey === resourceKey &&
(candidate.occurredAt.getTime() > currentOccurredAt ||
(candidate.occurredAt.getTime() === currentOccurredAt &&
BigInt(candidate.id) > BigInt(eventId)))
);
}) ?? null
);
}
const claimed = const claimed =
rows rows
.filter((row) => { .filter((row) => {
@ -357,11 +405,13 @@ function setup(seed: Partial<Store> = {}) {
return claimed; return claimed;
}, },
orderBy: (field: string) => { orderBy: (field: string) => {
query.order.push(field); if (alias === 'newerEvent') newerEventRead.order.push(field);
else query.order.push(field);
return builder; return builder;
}, },
setLock: (value: string) => { setLock: (value: string) => {
query.lock = value; if (alias === 'newerEvent') newerEventRead.lock = value;
else query.lock = value;
builderLock = value; builderLock = value;
return builder; return builder;
}, },
@ -370,12 +420,27 @@ function setup(seed: Partial<Store> = {}) {
return builder; return builder;
}, },
take: (value: number) => { take: (value: number) => {
query.take = value; if (alias === 'newerEvent') newerEventRead.take = value;
else query.take = value;
return builder; return builder;
}, },
where: (value: unknown) => { where: (value: unknown, parameters?: object) => {
query.brackets = value.constructor.name === 'Brackets'; const predicates =
query.predicates.push( alias === 'newerEvent'
? newerEventRead.predicates
: query.predicates;
if (typeof value === 'string') {
predicates.push({
expression: value,
parameters: parameters ?? {},
});
Object.assign(newerEventParameters, parameters);
return builder;
}
if (alias !== 'newerEvent') {
query.brackets = value.constructor.name === 'Brackets';
}
predicates.push(
recordBrackets( recordBrackets(
value as { value as {
whereFactory?: (where: { whereFactory?: (where: {
@ -387,6 +452,9 @@ function setup(seed: Partial<Store> = {}) {
); );
return builder; return builder;
}, },
andWhere: (value: unknown, parameters?: object) => {
return builder.where(value, parameters);
},
}; };
return builder; return builder;
}, },
@ -536,6 +604,7 @@ function setup(seed: Partial<Store> = {}) {
failDeliveryTarget = id; failDeliveryTarget = id;
}, },
lock: (id: string) => locked.add(id), lock: (id: string) => locked.add(id),
newerEventReads,
query, query,
savedDeliveryTargets: () => savedDeliveryTargets, savedDeliveryTargets: () => savedDeliveryTargets,
pauseNextClaim: () => { pauseNextClaim: () => {
@ -695,11 +764,6 @@ describe('SystemMessageFanoutService', () => {
id: '100', id: '100',
occurredAt: new KtDateTime(NOW.getTime() - 1), occurredAt: new KtDateTime(NOW.getTime() - 1),
}); });
const newer = event({
id: '201',
occurredAt: new KtDateTime(NOW.getTime() + 1),
nextFanoutAt: new KtDateTime(NOW.getTime() + 1),
});
const makeDelivery = ( const makeDelivery = (
id: string, id: string,
messageEventId: string, messageEventId: string,
@ -713,7 +777,7 @@ describe('SystemMessageFanoutService', () => {
subscriptionId: '300', subscriptionId: '300',
}); });
const fixture = setup({ const fixture = setup({
events: [older, event(), newer], events: [older, event()],
deliveries: [ deliveries: [
makeDelivery('1', '100', 'pending'), makeDelivery('1', '100', 'pending'),
makeDelivery('2', '100', 'retry'), makeDelivery('2', '100', 'retry'),
@ -729,7 +793,7 @@ describe('SystemMessageFanoutService', () => {
.deliveries() .deliveries()
.filter((item) => item.messageEventId === '100') .filter((item) => item.messageEventId === '100')
.map((item) => item.status), .map((item) => item.status),
).toEqual(['superseded', 'superseded', 'processing', 'superseded']); ).toEqual(['superseded', 'superseded', 'processing']);
expect( expect(
fixture.deliveries().find((item) => item.messageEventId === '201') fixture.deliveries().find((item) => item.messageEventId === '201')
?.status, ?.status,
@ -740,6 +804,208 @@ describe('SystemMessageFanoutService', () => {
).toBe('pending'); ).toBe('pending');
}); });
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({
fanoutAttemptCount: 1,
fanoutLeaseUntil: oldLease,
fanoutStatus: 'processing',
id: '9',
occurredAt: new KtDateTime(NOW.getTime() - 2),
payload: { endpoint: 'endpoint-a', portForwardId: RESOURCE_KEY },
});
const middleEvent = event({
fanoutStatus: 'completed',
id: '10',
nextFanoutAt: null,
occurredAt: new KtDateTime(NOW.getTime() - 1),
payload: { endpoint: 'endpoint-b', portForwardId: RESOURCE_KEY },
});
const newestEvent = event({
fanoutStatus: 'completed',
id: '11',
nextFanoutAt: null,
payload: { endpoint: 'endpoint-a', portForwardId: RESOURCE_KEY },
});
const unrelatedSourceEvent = event({
fanoutStatus: 'completed',
id: '12',
nextFanoutAt: null,
resourceKey: 'other-resource',
sourceKey: 'other.source',
});
const oldRows = [
'waiting_ddns',
'pending',
'retry',
'processing',
'success',
].map((status, index) =>
Object.assign(new QqbotMessageDelivery(), {
id: `old-${index}`,
messageEventId: oldEvent.id,
publishTargetId: `old-target-${index}`,
status,
subscriptionId: '300',
}),
);
const newestDelivery = Object.assign(new QqbotMessageDelivery(), {
id: 'newest',
messageEventId: newestEvent.id,
publishTargetId: '700',
renderedMessage: 'endpoint=endpoint-a',
status: 'pending',
subscriptionId: '300',
});
const unrelated = Object.assign(new QqbotMessageDelivery(), {
id: 'unrelated',
messageEventId: unrelatedSourceEvent.id,
publishTargetId: 'other-target',
status: 'pending',
subscriptionId: '301',
});
const unrelatedSourceResource = Object.assign(new QqbotMessageDelivery(), {
id: 'unrelated-source-resource',
messageEventId: unrelatedSourceEvent.id,
publishTargetId: 'other-source-resource-target',
status: 'pending',
subscriptionId: '300',
});
const fixture = setup({
deliveries: [
...oldRows,
newestDelivery,
unrelated,
unrelatedSourceResource,
],
events: [oldEvent, middleEvent, newestEvent, unrelatedSourceEvent],
});
const processClaim = (
fixture.service as unknown as {
processClaim: (token: object, now: Date) => Promise<void>;
}
).processClaim;
const newestSnapshot = structuredClone(newestDelivery);
await processClaim.call(
fixture.service,
{ attempt: 1, event: oldEvent, leaseUntil: oldLease },
NOW,
);
expect(fixture.newerEventReads).toHaveLength(1);
expect(fixture.adapter.resolveDelivery).not.toHaveBeenCalled();
expect(
fixture
.deliveries()
.filter((item) => item.messageEventId === oldEvent.id),
).toHaveLength(5);
expect(
fixture
.deliveries()
.filter((item) => item.messageEventId === oldEvent.id)
.map((item) => item.status),
).toEqual([
'superseded',
'superseded',
'superseded',
'processing',
'success',
]);
expect(
fixture.deliveries().find((item) => item.id === newestDelivery.id),
).toEqual(newestSnapshot);
expect(
fixture.deliveries().find((item) => item.id === unrelated.id),
).toEqual(unrelated);
expect(
fixture
.deliveries()
.find((item) => item.id === unrelatedSourceResource.id),
).toEqual(unrelatedSourceResource);
expect(
fixture.events().find((item) => item.id === oldEvent.id),
).toMatchObject({
fanoutAttemptCount: 1,
fanoutLeaseUntil: null,
fanoutStatus: 'completed',
});
expect(fixture.newerEventReads).toEqual([
{
lock: 'pessimistic_read',
order: [],
predicates: [
{
expression:
'newerEvent.sourceKey = :sourceKey AND newerEvent.resourceKey = :resourceKey',
parameters: { sourceKey: SOURCE_KEY, resourceKey: RESOURCE_KEY },
},
{
brackets: [
{
expression: 'newerEvent.occurredAt > :occurredAt',
parameters: { occurredAt: oldEvent.occurredAt },
},
{
expression:
'newerEvent.occurredAt = :occurredAt AND newerEvent.id > :eventId',
parameters: {
occurredAt: oldEvent.occurredAt,
eventId: oldEvent.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({
fanoutAttemptCount: 1,
fanoutLeaseUntil: oldLease,
fanoutStatus: 'processing',
id: '9',
payload: { endpoint: 'endpoint-a', portForwardId: RESOURCE_KEY },
});
const newestEvent = event({
fanoutStatus: 'completed',
id: '10',
nextFanoutAt: null,
payload: { endpoint: 'endpoint-a', portForwardId: RESOURCE_KEY },
});
const oldDelivery = Object.assign(new QqbotMessageDelivery(), {
id: 'old',
messageEventId: oldEvent.id,
publishTargetId: 'old-target',
status: 'pending',
subscriptionId: '300',
});
const fixture = setup({
deliveries: [oldDelivery],
events: [oldEvent, newestEvent],
});
const processClaim = (
fixture.service as unknown as {
processClaim: (token: object, now: Date) => Promise<void>;
}
).processClaim;
await processClaim.call(
fixture.service,
{ attempt: 1, event: oldEvent, leaseUntil: oldLease },
NOW,
);
expect(fixture.adapter.resolveDelivery).not.toHaveBeenCalled();
expect(fixture.deliveries()).toEqual([
expect.objectContaining({ id: 'old', status: 'superseded' }),
]);
});
it('accepts exact event-target replay and verifies only an exact duplicate-key race', async () => { it('accepts exact event-target replay and verifies only an exact duplicate-key race', async () => {
const existing = Object.assign(new QqbotMessageDelivery(), { const existing = Object.assign(new QqbotMessageDelivery(), {
id: 'old', id: 'old',