From d8d9613a6d5405ed0adab3cb780f347a1cf6746d Mon Sep 17 00:00:00 2001 From: can1357 Date: Fri, 26 Jun 2026 11:43:38 +0200 Subject: [PATCH] test(ai): verified cleanup of first-item watchdog timer - Implement `drainMicrotasksUntil` helper to manage async race conditions during tests. - Add fake timer validation to ensure the first-item watchdog is cleared when the upstream source rejects. - Restore real timers in `afterEach` to prevent test contamination. --- .../ai/test/stream-timeout-defaults.test.ts | 47 +++++++++++++------ 1 file changed, 33 insertions(+), 14 deletions(-) diff --git a/packages/ai/test/stream-timeout-defaults.test.ts b/packages/ai/test/stream-timeout-defaults.test.ts index 912a904c8..0ed90405d 100644 --- a/packages/ai/test/stream-timeout-defaults.test.ts +++ b/packages/ai/test/stream-timeout-defaults.test.ts @@ -1,4 +1,4 @@ -import { afterEach, beforeEach, describe, expect, it } from "bun:test"; +import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test"; import { getOpenAIStreamFirstEventTimeoutMs, getOpenAIStreamIdleTimeoutMs, @@ -42,6 +42,7 @@ afterEach(() => { Bun.env[key] = prior; } } + vi.useRealTimers(); }); describe("getStreamIdleTimeoutMs(fallbackMs)", () => { @@ -144,6 +145,14 @@ async function expectRejectsWithMessage(run: () => Promise, message: strin expect((caught as Error).message).toBe(message); } +async function drainMicrotasksUntil(predicate: () => boolean, errorMessage: string): Promise { + for (let i = 0; i < 1000; i++) { + if (predicate()) return; + await Promise.resolve(); + } + throw new Error(errorMessage); +} + describe("iterateWithIdleTimeout", () => { it("does not reset the first-progress deadline for no-progress items", async () => { const abortController = new AbortController(); @@ -177,29 +186,39 @@ describe("iterateWithIdleTimeout", () => { } }); - it("cleans first-item timers when the source throws before progress", async () => { - let firstItemTimedOut = false; + it("clears the first-item watchdog timer when the source rejects before progress", async () => { + vi.useFakeTimers(); + const baselineTimers = vi.getTimerCount(); - // biome-ignore lint/correctness/useYield: intentionally yields nothing — the test exercises the path where the source generator throws before its first yield. - async function* failingStream(): AsyncGenerator { - throw new Error("stream failed"); - } + // The first pull stays pending until we reject it, so the watchdog can be + // observed armed before the source errors. Rejecting a deferred promise is a + // microtask resolution, so the source error wins the watchdog race + // deterministically — with fake timers the watchdog never fires unless time is + // advanced, which this test never does. + const firstPull = Promise.withResolvers>(); + const failingStream: AsyncIterable = { + [Symbol.asyncIterator]: () => ({ next: () => firstPull.promise }), + }; - await expectRejectsWithMessage(async () => { - for await (const _item of iterateWithIdleTimeout(failingStream(), { + const settled = expectRejectsWithMessage(async () => { + for await (const _item of iterateWithIdleTimeout(failingStream, { firstItemTimeoutMs: 10, errorMessage: "idle timeout", firstItemErrorMessage: "first progress timeout", - onFirstItemTimeout: () => { - firstItemTimedOut = true; - }, })) { // Unreachable. } }, "stream failed"); - await Bun.sleep(20); - expect(firstItemTimedOut).toBe(false); + // The wrapper arms a first-item watchdog timer while it waits for the first pull. + await drainMicrotasksUntil(() => vi.getTimerCount() > baselineTimers, "first-item watchdog timer was not armed"); + + firstPull.reject(new Error("stream failed")); + await settled; + + // Surfacing the source error must clear the watchdog timer rather than leave it + // pending to fire after the stream has already failed. + expect(vi.getTimerCount()).toBe(baselineTimers); }); it("cleans first-item timers when the consumer returns before progress", async () => {