diff --git a/src/backend/services/events/workerInvoker.integration.test.ts b/src/backend/services/events/workerInvoker.integration.test.ts index e9fef6506..98c874bfd 100644 --- a/src/backend/services/events/workerInvoker.integration.test.ts +++ b/src/backend/services/events/workerInvoker.integration.test.ts @@ -39,7 +39,6 @@ import { describe, expect, it, - type MockInstance, vi, } from 'vitest'; import { @@ -317,14 +316,15 @@ const heldForMs = async (subId: string): Promise => { }; /** - * 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 => - waitUntil(() => - expect( - deferred.mock.settledResults.filter((r) => r.type === 'fulfilled'), - ).toHaveLength(count), +const failuresRecorded = (subId: string, count: number): Promise => + 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([]);