diff --git a/packages/stats/src/aggregator.ts b/packages/stats/src/aggregator.ts index c252e8467..e831fe396 100644 --- a/packages/stats/src/aggregator.ts +++ b/packages/stats/src/aggregator.ts @@ -50,6 +50,7 @@ import type { import { computeUsageWindowStats, fetchUsageSnapshots } from "./usage-windows"; const STATS_SYNC_LOCK_RETRY_MS = 25; +const STATS_SYNC_LOCK_ACQUIRE_STALE_MS = 10 * 1000; const STATS_SYNC_LOCK_STALE_MS = 60 * 60 * 1000; interface StatsLockSnapshot { @@ -103,7 +104,9 @@ async function readStatsLockSnapshot(lockPath: string): Promise(dbPath: string, fn: () => Promise) } catch (error) { if (errorCode(error) !== "EEXIST") throw error; await removeAbandonedStatsLock(lockPath); - const { promise, resolve } = Promise.withResolvers(); - setTimeout(resolve, STATS_SYNC_LOCK_RETRY_MS); - await promise; + await Bun.sleep(STATS_SYNC_LOCK_RETRY_MS); } } diff --git a/packages/stats/test/sync-serial.test.ts b/packages/stats/test/sync-serial.test.ts index f954cda3c..063f82e6f 100644 --- a/packages/stats/test/sync-serial.test.ts +++ b/packages/stats/test/sync-serial.test.ts @@ -110,10 +110,7 @@ describe("stats sync serial mode", () => { return true; }); - vi.spyOn(globalThis, "setTimeout").mockImplementation(callback => { - queueMicrotask(callback as () => void); - return 0; - }); + vi.spyOn(Bun, "sleep").mockResolvedValue(undefined); const result = await withStatsSyncLock(dbPath, async () => "acquired"); expect(result).toBe("acquired"); @@ -121,6 +118,21 @@ describe("stats sync serial mode", () => { expect(await Bun.file(breakerPath).exists()).toBe(false); }); + it("reclaims an unstamped lock after the acquisition grace period", async () => { + const dbPath = path.join(getSessionsDir(), "stats-unstamped-lock.db"); + const lockPath = `${dbPath}.sync.lock`; + await fs.mkdir(path.dirname(dbPath), { recursive: true }); + await Bun.write(lockPath, ""); + const staleTime = new Date(Date.now() - 11_000); + await fs.utimes(lockPath, staleTime, staleTime); + vi.spyOn(Bun, "sleep").mockResolvedValue(undefined); + + const result = await withStatsSyncLock(dbPath, async () => "acquired"); + + expect(result).toBe("acquired"); + expect(await Bun.file(lockPath).exists()).toBe(false); + }); + it("never reclaims a live owner-stamped breaker even when its mtime is stale", async () => { const dbPath = path.join(getSessionsDir(), "stats-live-breaker.db"); const lockPath = `${dbPath}.sync.lock`; @@ -138,17 +150,12 @@ describe("stats sync serial mode", () => { return true; }); const retryScheduled = Promise.withResolvers(); - let resumeRetry: (() => void) | undefined; + const resumeRetry = Promise.withResolvers(); let breakerReleased = false; - vi.spyOn(globalThis, "setTimeout").mockImplementation(callback => { - const run = callback as () => void; - if (breakerReleased) { - queueMicrotask(run); - } else { - resumeRetry = run; - retryScheduled.resolve(); - } - return 0; + vi.spyOn(Bun, "sleep").mockImplementation(async () => { + if (breakerReleased) return; + retryScheduled.resolve(); + await resumeRetry.promise; }); let acquired = false; @@ -162,7 +169,7 @@ describe("stats sync serial mode", () => { } finally { breakerReleased = true; await fs.rm(breakerPath, { force: true }); - resumeRetry?.(); + resumeRetry.resolve(); await pending; } expect(acquired).toBe(true);