test(events): wait for the recorded failure, not any future lease score

The in-flight claim's lease (2x the invoke timeout) is also a future score, so
polling heldForMs > 0 could pass before the failed attempt's hold was written:
the backoff test read 600 instead of 2000, and the 404/hang/429/suspended-owner
tests ended while their attempt was still in flight, letting it land in the
next test. Wait on the subscription's failure counter instead, which only moves
after the hold or gap marker is written.
This commit is contained in:
Daniel Salazar committed 2026-10-08 11:30:53 -07:00
1 parent 47a1be7f70
commit 06f8e0c5cc
1 file changed
+22 -32
@@ -39,7 +39,6 @@ import {
describe,
expect,
it,
type MockInstance,
vi,
} from 'vitest';
import {
@@ -317,14 +316,15 @@ const heldForMs = async (subId: string): Promise<number> => {
};
/**
* Wait until `count` failed attempts have had their holds written. A future
* score alone can't tell a hold from the in-flight claim's lease.
* Wait until this subscription has `count` failed attempts on record. The
* counter moves only after the attempt's hold or gap marker is written, so it
* can't be the in-flight claim's lease or another test's late attempt.
*/
const holdsWritten = (deferred: MockInstance, count: number): Promise<void> =>
waitUntil(() =>
expect(
deferred.mock.settledResults.filter((r) => r.type === 'fulfilled'),
).toHaveLength(count),
const failuresRecorded = (subId: string, count: number): Promise<void> =>
waitUntil(async () =>
expect(await env.server.clients.redis.get(`ev:qf:{${subId}}`)).toBe(
String(count),
),
);
beforeAll(async () => {
@@ -599,11 +599,7 @@ describe('what each answer does to the delivery', () => {
// It counted: a refusal is still a handler that did not work. Counted
// only once the marker is in, so it is also what says the answer was
// handled — the depth alone reads 1 while the event is still leased.
await waitUntil(async () =>
expect(await env.server.clients.redis.get(`ev:qf:{${subId}}`)).toBe(
'1',
),
);
await failuresRecorded(subId, 1);
// The event is gone and a marker stands in its place, so the
// subscription learns there was one rather than reading silence.
@@ -626,16 +622,14 @@ describe('what each answer does to the delivery', () => {
await touch('unmarked.txt');
await invoked(1);
await waitUntil(async () =>
expect(await heldForMs(subId)).toBeGreaterThan(0),
);
await failuresRecorded(subId, 1);
expect(await heldForMs(subId)).toBeGreaterThan(0);
expect(await pending().depth(subId)).toBe(1);
});
it('holds one it could not answer, for longer each time', async () => {
answer = 500;
const subId = await subscribe();
const deferred = vi.spyOn(pending(), 'deferAfterFailure');
// Freeze the clock so every hold is read against the same instant it
// was written from.
jump(0);
@@ -649,7 +643,7 @@ describe('what each answer does to the delivery', () => {
attempt < EVENTS_CONSECUTIVE_FAILURES;
attempt++
) {
await holdsWritten(deferred, attempt);
await failuresRecorded(subId, attempt);
waits.push(await heldForMs(subId));
// Nothing may take it while it is held.
@@ -672,7 +666,6 @@ describe('what each answer does to the delivery', () => {
),
);
expect(calls).toHaveLength(EVENTS_CONSECUTIVE_FAILURES);
deferred.mockRestore();
});
it('tells the developer their handler stopped working', async () => {
@@ -686,7 +679,7 @@ describe('what each answer does to the delivery', () => {
attempt < EVENTS_CONSECUTIVE_FAILURES;
attempt++
) {
await invoked(attempt);
await failuresRecorded(subId, attempt);
jump(deliveryBackoffMs(attempt) + 50);
await events().sweepPending();
}
@@ -710,9 +703,8 @@ describe('what each answer does to the delivery', () => {
await invoked(1);
// Still owed, and held rather than dropped: nobody said no.
await waitUntil(async () =>
expect(await heldForMs(subId)).toBeGreaterThan(0),
);
await failuresRecorded(subId, 1);
expect(await heldForMs(subId)).toBeGreaterThan(0);
expect(await pending().depth(subId)).toBe(1);
answer = 200;
});
@@ -724,9 +716,8 @@ describe('what each answer does to the delivery', () => {
await touch('busy.txt');
await invoked(1);
await waitUntil(async () =>
expect(await heldForMs(subId)).toBeGreaterThan(0),
);
await failuresRecorded(subId, 1);
expect(await heldForMs(subId)).toBeGreaterThan(0);
expect(await pending().depth(subId)).toBe(1);
});
});
@@ -794,7 +785,6 @@ describe('with no events worker to address', () => {
// Unpublished after the fact: the app's set now hashes to nothing, so
// there is no script to name and nowhere to send the delivery.
await env.server.stores.eventHandler.remove(appUid, HANDLER);
const deferred = vi.spyOn(pending(), 'deferAfterFailure');
try {
await touch('unresolved.txt');
@@ -803,7 +793,7 @@ describe('with no events worker to address', () => {
attempt < EVENTS_CONSECUTIVE_FAILURES;
attempt++
) {
await holdsWritten(deferred, attempt);
await failuresRecorded(subId, attempt);
jump(deliveryBackoffMs(attempt) + 50);
await events().sweepPending();
}
@@ -816,7 +806,6 @@ describe('with no events worker to address', () => {
// Nothing was ever called: there was nowhere to call.
expect(calls).toEqual([]);
} finally {
deferred.mockRestore();
await env.server.stores.eventHandler.publish({
appUid,
name: HANDLER,
@@ -835,9 +824,10 @@ describe('with no events worker to address', () => {
try {
await touch('suspended-owner.txt');
await waitUntil(async () =>
expect(await heldForMs(subId)).toBeGreaterThan(0),
);
// Not just the claim's lease: un-suspending the owner while the
// attempt is still in flight would let it address and call out.
await failuresRecorded(subId, 1);
expect(await heldForMs(subId)).toBeGreaterThan(0);
expect(await pending().depth(subId)).toBe(1);
// A suspended owner is nobody to address: nothing was called.
expect(calls).toEqual([]);