fix(stats): reclaim interrupted sync locks promptly
This commit is contained in:
@@ -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<StatsLockSnapsho
|
||||
|
||||
function statsLockOwnerIsRunning(snapshot: StatsLockSnapshot): boolean {
|
||||
const pid = Number.parseInt(snapshot.text.split(/\r?\n/, 1)[0] ?? "", 10);
|
||||
if (!Number.isSafeInteger(pid) || pid <= 0) return Date.now() - snapshot.mtimeMs < STATS_SYNC_LOCK_STALE_MS;
|
||||
if (!Number.isSafeInteger(pid) || pid <= 0) {
|
||||
return Date.now() - snapshot.mtimeMs < STATS_SYNC_LOCK_ACQUIRE_STALE_MS;
|
||||
}
|
||||
try {
|
||||
process.kill(pid, 0);
|
||||
return true;
|
||||
@@ -214,9 +217,7 @@ export async function withStatsSyncLock<T>(dbPath: string, fn: () => Promise<T>)
|
||||
} catch (error) {
|
||||
if (errorCode(error) !== "EEXIST") throw error;
|
||||
await removeAbandonedStatsLock(lockPath);
|
||||
const { promise, resolve } = Promise.withResolvers<void>();
|
||||
setTimeout(resolve, STATS_SYNC_LOCK_RETRY_MS);
|
||||
await promise;
|
||||
await Bun.sleep(STATS_SYNC_LOCK_RETRY_MS);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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<void>();
|
||||
let resumeRetry: (() => void) | undefined;
|
||||
const resumeRetry = Promise.withResolvers<void>();
|
||||
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);
|
||||
|
||||
Reference in New Issue
Block a user