diff --git a/packages/ai/src/stream.ts b/packages/ai/src/stream.ts index 8daf7610a..7c5cfc79a 100644 --- a/packages/ai/src/stream.ts +++ b/packages/ai/src/stream.ts @@ -177,6 +177,8 @@ const PROVIDER_INFLIGHT_LOCK_STALE_MS = 10_000; const PROVIDER_INFLIGHT_LEASE_STALE_MS = 30_000; const PROVIDER_INFLIGHT_HEARTBEAT_MS = 5_000; const PROVIDER_INFLIGHT_SIGNAL_FALLBACK_MS = 250; +const PROVIDER_INFLIGHT_HEARTBEAT_FLUSH_TIMEOUT_MS = 1_000; +const PROVIDER_INFLIGHT_RELEASE_TIMEOUT_MS = 5_000; let configuredProviderMaxInFlightRequests: Record = {}; let providerInFlightRootOverride: string | undefined; @@ -519,9 +521,32 @@ async function removeProviderInFlightLeaseDir(leasePath: string): Promise // write `.wakeup` into an unrelated provider directory. async function releaseProviderInFlightLease(lease: ProviderInFlightLease): Promise { clearInterval(lease.heartbeat); - await lease.flushHeartbeat(); - await removeProviderInFlightLeaseDir(lease.path); - await signalProviderInFlightWaitersInDir(path.dirname(lease.path)); + const flushTimeout = Promise.withResolvers<"timeout">(); + const flushTimer = setTimeout(() => flushTimeout.resolve("timeout"), PROVIDER_INFLIGHT_HEARTBEAT_FLUSH_TIMEOUT_MS); + flushTimer.unref?.(); + try { + const outcome = await Promise.race([lease.flushHeartbeat().then(() => "flushed" as const), flushTimeout.promise]); + if (outcome === "timeout") { + logger.warn("Provider in-flight heartbeat flush timed out; forcing lease cleanup", { path: lease.path }); + } + } finally { + clearTimeout(flushTimer); + } + + const releaseTimeout = Promise.withResolvers(); + const releaseTimer = setTimeout( + () => releaseTimeout.reject(new Error("Provider in-flight lease cleanup timed out")), + PROVIDER_INFLIGHT_RELEASE_TIMEOUT_MS, + ); + releaseTimer.unref?.(); + try { + await Promise.race([removeProviderInFlightLeaseDir(lease.path), releaseTimeout.promise]); + } finally { + clearTimeout(releaseTimer); + } + // Wake-up is an optimization: waiters also poll every 250 ms. Do not let a + // notification-file stall keep a completed provider request open. + void signalProviderInFlightWaitersInDir(path.dirname(lease.path)); } async function acquireProviderInFlightSlot( @@ -585,11 +610,11 @@ function withProviderInFlightLimit { let release: (() => Promise) | undefined; - let released = false; - const releaseOnce = async () => { - if (!release || released) return; - released = true; - await release(); + let releasePromise: Promise | undefined; + const releaseOnce = () => { + if (!release) return Promise.resolve(); + releasePromise ??= release(); + return releasePromise; }; try { const startedWaitingAt = Date.now(); @@ -601,18 +626,39 @@ function withProviderInFlightLimit { expect(mock.calls).toHaveLength(2); }); + test("releases its provider lease before reporting stream completion", async () => { + registerMockApi(); + const mock = createMockModel({ provider: "tests", responses: [{ content: ["reply"] }] }); + + const stream = streamSimple(mock.model, context(), { maxInFlightRequests: { tests: 1 } }); + const result = await stream.result(); + + expect(result.content).toEqual([{ type: "text", text: "reply" }]); + const entries = await fs.readdir(limiterDir("tests"), { withFileTypes: true }); + expect(entries.filter(entry => entry.isDirectory())).toHaveLength(0); + }); + test("removes an aborted queued request without dispatching it", async () => { registerMockApi(); const firstStarted = Promise.withResolvers();