fix(stats): avoided macos stats sync worker abort

Kept the documented workers: 1 path inline and defaulted macOS stats sync to that serial parser path so /stats dashboard launches do not re-enter Bun workers on macOS.

Fixes #3733
This commit is contained in:
roboomp
2026-06-28 16:21:59 +00:00
parent 2622ab13fe
commit b48d488f2c
3 changed files with 144 additions and 33 deletions
+4
View File
@@ -2,6 +2,10 @@
## [Unreleased]
### Fixed
- Kept stats session sync on the serial parser path for `workers: 1` and macOS defaults, avoiding Bun worker re-entry aborts when launching `/stats` ([#3733](https://github.com/can1357/oh-my-pi/issues/3733)).
## [16.2.3] - 2026-06-28
### Added
+49 -33
View File
@@ -23,7 +23,7 @@ import {
setFileOffset,
updateUserMessageLinks,
} from "./db";
import { getSessionEntry, listAllSessionFiles, type ParseSessionResult } from "./parser";
import { getSessionEntry, listAllSessionFiles, type ParseSessionResult, parseSessionFile } from "./parser";
import type { SyncWorkerRequest, SyncWorkerResponse } from "./sync-worker";
// Coding-agent binary/bundle workers route through the CLI entrypoint with a
// hidden argv mode, so the compiled binary and npm bundle only need one
@@ -68,6 +68,9 @@ export interface SyncOptions {
}
function defaultWorkerCount(): number {
// Bun 1.3.x can abort the macOS process when stats sync workers re-enter
// the compiled `omp` binary. Keep macOS on the documented serial path.
if (process.platform === "darwin") return 1;
// `navigator.hardwareConcurrency` is the portable answer in Bun; fall
// back to a small fixed pool if it's somehow unavailable.
const hw = typeof navigator !== "undefined" ? (navigator.hardwareConcurrency ?? 0) : 0;
@@ -183,11 +186,11 @@ export async function smokeTestSyncWorker({ timeoutMs = 5_000 }: { timeoutMs?: n
/**
* Sync all session files to the database.
*
* Parsing fans out across a worker pool (one in-flight job per worker)
* while DB writes and offset bookkeeping stay on the calling thread so the
* single SQLite handle stays uncontended. `onProgress` fires once per
* completed file (skipped files included so the bar walks at a steady
* rate).
* `workers: 1` parses inline. Larger pools fan parsing out across workers
* (one in-flight job per worker) while DB writes and offset bookkeeping stay on
* the calling thread so the single SQLite handle stays uncontended.
* `onProgress` fires once per completed file (skipped files included so the
* bar walks at a steady rate).
*/
export async function syncAllSessions(opts?: SyncOptions): Promise<{ processed: number; files: number }> {
await initDb();
@@ -200,10 +203,6 @@ export async function syncAllSessions(opts?: SyncOptions): Promise<{ processed:
let completed = 0;
let cursor = 0;
const poolSize = Math.max(1, Math.min(files.length, opts?.workers ?? defaultWorkerCount()));
const handles: WorkerHandle[] = [];
for (let i = 0; i < poolSize; i++) handles.push(spawnWorker());
const report = (sessionFile: string) => {
completed++;
opts?.onProgress?.({
@@ -214,34 +213,51 @@ export async function syncAllSessions(opts?: SyncOptions): Promise<{ processed:
});
};
const processFile = async (
sessionFile: string,
parse: (sessionFile: string, fromOffset: number) => Promise<ParseSessionResult>,
): Promise<void> => {
let fileStats: fs.Stats;
try {
fileStats = await fs.promises.stat(sessionFile);
} catch {
report(sessionFile);
return;
}
const lastModified = fileStats.mtimeMs;
const stored = getFileOffset(sessionFile);
if (stored && stored.lastModified >= lastModified) {
report(sessionFile);
return;
}
const fromOffset = stored?.offset ?? 0;
const result = await parse(sessionFile, fromOffset);
const inserted = applyParseResult(sessionFile, lastModified, result);
if (inserted > 0) {
totalProcessed += inserted;
filesProcessed++;
}
report(sessionFile);
};
const poolSize = Math.max(1, Math.min(files.length, opts?.workers ?? defaultWorkerCount()));
if (poolSize === 1) {
for (const sessionFile of files) {
await processFile(sessionFile, parseSessionFile);
}
return { processed: totalProcessed, files: filesProcessed };
}
const handles: WorkerHandle[] = [];
for (let i = 0; i < poolSize; i++) handles.push(spawnWorker());
async function drain(handle: WorkerHandle): Promise<void> {
while (true) {
const idx = cursor++;
if (idx >= files.length) return;
const sessionFile = files[idx];
let fileStats: fs.Stats;
try {
fileStats = await fs.promises.stat(sessionFile);
} catch {
report(sessionFile);
continue;
}
const lastModified = fileStats.mtimeMs;
const stored = getFileOffset(sessionFile);
if (stored && stored.lastModified >= lastModified) {
report(sessionFile);
continue;
}
const fromOffset = stored?.offset ?? 0;
const result = await dispatch(handle, { sessionFile, fromOffset });
const inserted = applyParseResult(sessionFile, lastModified, result);
if (inserted > 0) {
totalProcessed += inserted;
filesProcessed++;
}
report(sessionFile);
await processFile(sessionFile, (file, fromOffset) => dispatch(handle, { sessionFile: file, fromOffset }));
}
}
+91
View File
@@ -0,0 +1,91 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test";
import * as fs from "node:fs/promises";
import * as os from "node:os";
import * as path from "node:path";
import { syncAllSessions } from "@oh-my-pi/omp-stats/aggregator";
import { closeDb, getOverallStats } from "@oh-my-pi/omp-stats/db";
import { getAgentDir, getSessionsDir, setAgentDir, TempDir } from "@oh-my-pi/pi-utils";
const originalConfigDir = process.env.PI_CONFIG_DIR;
const originalAgentDir = getAgentDir();
let tempDir: TempDir | null = null;
beforeEach(() => {
tempDir = TempDir.createSync("@pi-stats-sync-serial-");
const configDir = path.relative(os.homedir(), tempDir.join("config"));
process.env.PI_CONFIG_DIR = configDir;
setAgentDir(tempDir.join("agent"));
});
afterEach(() => {
vi.restoreAllMocks();
closeDb();
if (originalConfigDir === undefined) {
delete process.env.PI_CONFIG_DIR;
} else {
process.env.PI_CONFIG_DIR = originalConfigDir;
}
setAgentDir(originalAgentDir);
tempDir?.removeSync();
tempDir = null;
});
async function writeSessionFile(): Promise<void> {
const sessionDir = path.join(getSessionsDir(), "--tmp--sync-serial");
await fs.mkdir(sessionDir, { recursive: true });
const timestamp = new Date().toISOString();
const sessionFile = path.join(sessionDir, "session.jsonl");
const assistant = {
type: "message",
id: "assistant-1",
parentId: null,
timestamp,
message: {
role: "assistant",
content: [{ type: "text", text: "ok" }],
api: "openai-responses",
provider: "openai",
model: "gpt-5.4",
usage: {
input: 1,
output: 2,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 3,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
},
stopReason: "stop",
timestamp: Date.now(),
duration: 10,
ttft: 5,
},
};
await Bun.write(sessionFile, `${JSON.stringify(assistant)}\n`);
}
describe("stats sync serial mode", () => {
it("honors workers: 1 without spawning a worker", async () => {
await writeSessionFile();
const workerSpy = vi.spyOn(globalThis, "Worker");
const synced = await syncAllSessions({ workers: 1 });
const overall = getOverallStats();
expect(synced.files).toBe(1);
expect(overall.totalRequests).toBe(1);
expect(workerSpy).not.toHaveBeenCalled();
});
it("uses the serial parser by default on macOS", async () => {
await writeSessionFile();
vi.spyOn(process, "platform", "get").mockReturnValue("darwin");
const workerSpy = vi.spyOn(globalThis, "Worker");
const synced = await syncAllSessions();
const overall = getOverallStats();
expect(synced.files).toBe(1);
expect(overall.totalRequests).toBe(1);
expect(workerSpy).not.toHaveBeenCalled();
});
});