fix(coding-agent): reconcile archived session stats
This commit is contained in:
@@ -1,5 +1,6 @@
|
||||
import * as fs from "node:fs";
|
||||
import { workerHostEntry } from "@oh-my-pi/pi-utils";
|
||||
import * as path from "node:path";
|
||||
import { getStatsDbPath, workerHostEntry } from "@oh-my-pi/pi-utils";
|
||||
import {
|
||||
getRecentErrors as dbGetRecentErrors,
|
||||
getRecentRequests as dbGetRecentRequests,
|
||||
@@ -48,6 +49,209 @@ import type {
|
||||
} from "./types";
|
||||
import { computeUsageWindowStats, fetchUsageSnapshots } from "./usage-windows";
|
||||
|
||||
const STATS_SYNC_LOCK_RETRY_MS = 25;
|
||||
const STATS_SYNC_LOCK_STALE_MS = 60 * 60 * 1000;
|
||||
|
||||
interface StatsLockSnapshot {
|
||||
dev: number;
|
||||
ino: number;
|
||||
size: number;
|
||||
mtimeMs: number;
|
||||
text: string;
|
||||
}
|
||||
|
||||
function errorCode(error: unknown): string | undefined {
|
||||
if (!error || typeof error !== "object" || !("code" in error)) return undefined;
|
||||
return typeof error.code === "string" ? error.code : undefined;
|
||||
}
|
||||
|
||||
function sameStatsLock(left: StatsLockSnapshot, right: StatsLockSnapshot): boolean {
|
||||
return (
|
||||
left.dev === right.dev &&
|
||||
left.ino === right.ino &&
|
||||
left.size === right.size &&
|
||||
left.mtimeMs === right.mtimeMs &&
|
||||
left.text === right.text
|
||||
);
|
||||
}
|
||||
|
||||
async function readStatsLockSnapshot(lockPath: string): Promise<StatsLockSnapshot | null> {
|
||||
try {
|
||||
const before = await fs.promises.stat(lockPath);
|
||||
const text = await fs.promises.readFile(lockPath, "utf8");
|
||||
const after = await fs.promises.stat(lockPath);
|
||||
const first = {
|
||||
dev: before.dev,
|
||||
ino: before.ino,
|
||||
size: before.size,
|
||||
mtimeMs: before.mtimeMs,
|
||||
text,
|
||||
};
|
||||
const second = {
|
||||
dev: after.dev,
|
||||
ino: after.ino,
|
||||
size: after.size,
|
||||
mtimeMs: after.mtimeMs,
|
||||
text,
|
||||
};
|
||||
return sameStatsLock(first, second) ? second : null;
|
||||
} catch (error) {
|
||||
if (errorCode(error) === "ENOENT") return null;
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
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;
|
||||
try {
|
||||
process.kill(pid, 0);
|
||||
return true;
|
||||
} catch (error) {
|
||||
return errorCode(error) !== "ESRCH";
|
||||
}
|
||||
}
|
||||
|
||||
function createStatsLockToken(): string {
|
||||
return `${process.pid}\n${Date.now()}\n${Math.random()}\n`;
|
||||
}
|
||||
|
||||
async function removeAbandonedStatsLockFile(lockPath: string): Promise<boolean> {
|
||||
const snapshot = await readStatsLockSnapshot(lockPath);
|
||||
if (!snapshot || statsLockOwnerIsRunning(snapshot)) return false;
|
||||
|
||||
const current = await readStatsLockSnapshot(lockPath);
|
||||
if (!current || !sameStatsLock(snapshot, current) || statsLockOwnerIsRunning(current)) return false;
|
||||
try {
|
||||
await fs.promises.unlink(lockPath);
|
||||
return true;
|
||||
} catch (error) {
|
||||
if (errorCode(error) === "ENOENT") return false;
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async function removeOwnedStatsLockFile(lockPath: string, owned: StatsLockSnapshot): Promise<void> {
|
||||
const current = await readStatsLockSnapshot(lockPath);
|
||||
if (!current || !sameStatsLock(owned, current)) return;
|
||||
try {
|
||||
await fs.promises.unlink(lockPath);
|
||||
} catch (error) {
|
||||
if (errorCode(error) !== "ENOENT") throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async function removeAbandonedStatsLock(lockPath: string): Promise<boolean> {
|
||||
const snapshot = await readStatsLockSnapshot(lockPath);
|
||||
if (!snapshot || statsLockOwnerIsRunning(snapshot)) return false;
|
||||
|
||||
const breakerPath = `${lockPath}.break`;
|
||||
let breaker: fs.promises.FileHandle;
|
||||
try {
|
||||
breaker = await fs.promises.open(breakerPath, "wx");
|
||||
} catch (error) {
|
||||
if (errorCode(error) === "EEXIST") {
|
||||
await removeAbandonedStatsLockFile(breakerPath);
|
||||
return false;
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
|
||||
const breakerToken = createStatsLockToken();
|
||||
let breakerSnapshot: StatsLockSnapshot | null = null;
|
||||
let breakerTokenWritten = false;
|
||||
let removed = false;
|
||||
let cleanupError: unknown;
|
||||
try {
|
||||
await breaker.writeFile(breakerToken);
|
||||
breakerTokenWritten = true;
|
||||
breakerSnapshot = await readStatsLockSnapshot(breakerPath);
|
||||
if (!breakerSnapshot || breakerSnapshot.text !== breakerToken) {
|
||||
throw new Error(`Stats lock breaker changed while acquiring ${breakerPath}`);
|
||||
}
|
||||
|
||||
const current = await readStatsLockSnapshot(lockPath);
|
||||
if (current && sameStatsLock(snapshot, current) && !statsLockOwnerIsRunning(current)) {
|
||||
try {
|
||||
await fs.promises.unlink(lockPath);
|
||||
removed = true;
|
||||
} catch (error) {
|
||||
if (errorCode(error) !== "ENOENT") throw error;
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
try {
|
||||
await breaker.close();
|
||||
} catch (error) {
|
||||
cleanupError = error;
|
||||
}
|
||||
try {
|
||||
if (!breakerSnapshot && breakerTokenWritten) {
|
||||
const current = await readStatsLockSnapshot(breakerPath);
|
||||
if (current?.text === breakerToken) breakerSnapshot = current;
|
||||
}
|
||||
if (breakerSnapshot) await removeOwnedStatsLockFile(breakerPath, breakerSnapshot);
|
||||
} catch (error) {
|
||||
cleanupError ??= error;
|
||||
}
|
||||
}
|
||||
if (cleanupError) throw cleanupError;
|
||||
return removed;
|
||||
}
|
||||
|
||||
/**
|
||||
* Serialize stats ingestion and archive reconciliation across processes.
|
||||
* The lock covers file discovery, parsing, and the final SQLite write so a
|
||||
* parse result for a session moved by GC can never commit after cleanup.
|
||||
*/
|
||||
export async function withStatsSyncLock<T>(dbPath: string, fn: () => Promise<T>): Promise<T> {
|
||||
const lockPath = `${dbPath}.sync.lock`;
|
||||
await fs.promises.mkdir(path.dirname(dbPath), { recursive: true });
|
||||
let handle: fs.promises.FileHandle | undefined;
|
||||
while (!handle) {
|
||||
try {
|
||||
handle = await fs.promises.open(lockPath, "wx");
|
||||
} 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;
|
||||
}
|
||||
}
|
||||
|
||||
const token = createStatsLockToken();
|
||||
let result: T | undefined;
|
||||
let operationError: unknown;
|
||||
let operationFailed = false;
|
||||
let tokenWritten = false;
|
||||
try {
|
||||
await handle.writeFile(token);
|
||||
tokenWritten = true;
|
||||
result = await fn();
|
||||
} catch (error) {
|
||||
operationFailed = true;
|
||||
operationError = error;
|
||||
}
|
||||
|
||||
let cleanupError: unknown;
|
||||
try {
|
||||
await handle.close();
|
||||
} catch (error) {
|
||||
cleanupError = error;
|
||||
}
|
||||
try {
|
||||
const current = await fs.promises.readFile(lockPath, "utf8");
|
||||
if (!tokenWritten || current === token) await fs.promises.unlink(lockPath);
|
||||
} catch (error) {
|
||||
if (errorCode(error) !== "ENOENT") cleanupError ??= error;
|
||||
}
|
||||
|
||||
if (operationFailed) throw operationError;
|
||||
if (cleanupError) throw cleanupError;
|
||||
return result as T;
|
||||
}
|
||||
|
||||
/**
|
||||
* Apply a freshly parsed result to the database. Runs entirely on the
|
||||
* main thread so the single SQLite handle owns every write.
|
||||
@@ -218,6 +422,10 @@ export async function smokeTestSyncWorker({ timeoutMs = 5_000 }: { timeoutMs?: n
|
||||
* bar walks at a steady rate).
|
||||
*/
|
||||
export async function syncAllSessions(opts?: SyncOptions): Promise<{ processed: number; files: number }> {
|
||||
return withStatsSyncLock(getStatsDbPath(), () => syncAllSessionsLocked(opts));
|
||||
}
|
||||
|
||||
async function syncAllSessionsLocked(opts?: SyncOptions): Promise<{ processed: number; files: number }> {
|
||||
await initDb();
|
||||
|
||||
const files = await listAllSessionFiles();
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import { afterEach, describe, expect, it, vi } from "bun:test";
|
||||
import * as fs from "node:fs/promises";
|
||||
import * as path from "node:path";
|
||||
import { syncAllSessions } from "@oh-my-pi/omp-stats/aggregator";
|
||||
import { syncAllSessions, withStatsSyncLock } from "@oh-my-pi/omp-stats/aggregator";
|
||||
import { getOverallStats } from "@oh-my-pi/omp-stats/db";
|
||||
import { getSessionsDir } from "@oh-my-pi/pi-utils";
|
||||
import { installStatsTestIsolation } from "./helpers/temp-agent";
|
||||
@@ -93,4 +93,78 @@ describe("stats sync serial mode", () => {
|
||||
await expect(syncAllSessions({ workers: 2 })).rejects.toBe(workerProbe);
|
||||
expect(workerSpy).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("reclaims a dead owner's abandoned stale breaker", async () => {
|
||||
const dbPath = path.join(getSessionsDir(), "stats-lock.db");
|
||||
const lockPath = `${dbPath}.sync.lock`;
|
||||
const breakerPath = `${lockPath}.break`;
|
||||
const deadPid = 424_242;
|
||||
await fs.mkdir(path.dirname(dbPath), { recursive: true });
|
||||
await Bun.write(lockPath, `${deadPid}\n1\nprimary\n`);
|
||||
await Bun.write(breakerPath, `${deadPid}\n1\nbreaker\n`);
|
||||
const staleTime = new Date(Date.now() - 2 * 60 * 60 * 1000);
|
||||
await fs.utimes(lockPath, staleTime, staleTime);
|
||||
await fs.utimes(breakerPath, staleTime, staleTime);
|
||||
vi.spyOn(process, "kill").mockImplementation(pid => {
|
||||
if (pid === deadPid) throw Object.assign(new Error("dead owner"), { code: "ESRCH" });
|
||||
return true;
|
||||
});
|
||||
|
||||
vi.spyOn(globalThis, "setTimeout").mockImplementation(callback => {
|
||||
queueMicrotask(callback as () => void);
|
||||
return 0;
|
||||
});
|
||||
const result = await withStatsSyncLock(dbPath, async () => "acquired");
|
||||
|
||||
expect(result).toBe("acquired");
|
||||
expect(await Bun.file(lockPath).exists()).toBe(false);
|
||||
expect(await Bun.file(breakerPath).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`;
|
||||
const breakerPath = `${lockPath}.break`;
|
||||
const deadPid = 424_243;
|
||||
const liveBreakerToken = `${process.pid}\n1\nlive-breaker\n`;
|
||||
await fs.mkdir(path.dirname(dbPath), { recursive: true });
|
||||
await Bun.write(lockPath, `${deadPid}\n1\nprimary\n`);
|
||||
await Bun.write(breakerPath, liveBreakerToken);
|
||||
const staleTime = new Date(Date.now() - 2 * 60 * 60 * 1000);
|
||||
await fs.utimes(lockPath, staleTime, staleTime);
|
||||
await fs.utimes(breakerPath, staleTime, staleTime);
|
||||
vi.spyOn(process, "kill").mockImplementation(pid => {
|
||||
if (pid === deadPid) throw Object.assign(new Error("dead owner"), { code: "ESRCH" });
|
||||
return true;
|
||||
});
|
||||
const retryScheduled = Promise.withResolvers<void>();
|
||||
let resumeRetry: (() => void) | undefined;
|
||||
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;
|
||||
});
|
||||
|
||||
let acquired = false;
|
||||
const pending = withStatsSyncLock(dbPath, async () => {
|
||||
acquired = true;
|
||||
});
|
||||
try {
|
||||
await retryScheduled.promise;
|
||||
expect(acquired).toBe(false);
|
||||
expect(await Bun.file(breakerPath).text()).toBe(liveBreakerToken);
|
||||
} finally {
|
||||
breakerReleased = true;
|
||||
await fs.rm(breakerPath, { force: true });
|
||||
resumeRetry?.();
|
||||
await pending;
|
||||
}
|
||||
expect(acquired).toBe(true);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user