From 4f202ebefa0d10958c460bebc1bb642b902b4d86 Mon Sep 17 00:00:00 2001 From: lcy Date: Tue, 22 Sep 2026 15:23:50 +0800 Subject: [PATCH 1/2] fix: run profile learning from any active instance under a cross-process lock MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Profile learning was gated on `webServer?.isServerOwner()`. Ownership tracks whichever process bound the web-server port first and is only handed over when that process becomes unreachable (web-server.ts:290) — it does not track which instance the user is actually working in. Two consequences: - With `webServerEnabled: false` there is no server and therefore no owner, so `webServer` stays null (index.ts:244) and automatic profile learning never ran at all. - When the owner stays healthy but stops receiving sessions while another project is active, the prompt queue stalls indefinitely. Learning is scheduled from a `session.idle` event, and the health check only takes over ownership when the server is unreachable, so nothing hands the work over. Any active instance may now learn. Retention cleanup stays owner-only, since it is storage-wide maintenance with no reason to run once per instance. Allowing concurrent learners requires real mutual exclusion. `isLearningRunning` is a module-level boolean and only prevents re-entry within one process, while prompt selection is a plain SELECT (user-prompt-manager.ts:270) whose batch is marked only after the LLM responds. Without cross-process exclusion two instances would analyze the same prompts and the slower writer would clobber the faster one's profile update — `updateProfile` re-reads the current version at write time, so its optimistic check does not detect a concurrent update that landed while the LLM was running. This adds a dedicated lock keyed to the shared storage path, held across the whole learning flow. It deliberately does not reuse `.turso-operation.lock`: that lock makes `assertNoTursoMigrationInProgress` reject ordinary memory writes. Contention skips the round rather than waiting, because the next idle event retries. The lock also respects a write window before reclaiming an unparseable lock file. `writeFileSync` is not atomic, so a reader can observe a file that was created but not yet filled in; reclaiming it on sight would hand the lock to a second process while the first believes it holds it. Cold-buffer staleness had to be addressed for the same reason. The buffer is read once per process and `saveColdBuffers` rewrites the whole file from that in-memory map, so a process holding a stale map would erase entries a peer persisted. The manager now reloads it when the file's mtime changes, keeping the requirement inside the manager rather than relying on callers. Tests run the real plugin in an isolated process (following compaction-agent-preservation.test.ts) instead of string-slicing the handler out of the transpiled source, and the lock is covered by two genuinely concurrent processes synchronised on a start barrier. Both were verified to fail when the corresponding protection is removed. --- src/index.ts | 11 +- src/services/user-memory-learning.ts | 14 + src/services/user-profile/learning-lock.ts | 160 ++++++++++ .../user-profile/user-profile-manager.ts | 29 +- tests/profile-learning-idle.test.ts | 209 +++++++++++++ tests/profile-learning-lock.test.ts | 280 ++++++++++++++++++ 6 files changed, 701 insertions(+), 2 deletions(-) create mode 100644 src/services/user-profile/learning-lock.ts create mode 100644 tests/profile-learning-idle.test.ts create mode 100644 tests/profile-learning-lock.test.ts diff --git a/src/index.ts b/src/index.ts index 5b56c708..854fc933 100644 --- a/src/index.ts +++ b/src/index.ts @@ -955,8 +955,17 @@ export const OpenCodeMemPlugin: Plugin = async (ctx: PluginInput) => { try { await performAutoCapture(ctx, sessionID, directory); + // Prompts are shared across projects, but web-server ownership tracks + // whoever bound the port first and is never handed over while that + // process stays reachable. Gating learning on it stalls the queue + // whenever the owner stops seeing sessions, and disables learning + // outright when the web server is off. Any active instance may learn; + // performUserProfileLearning holds a cross-process lock internally. + await performUserProfileLearning(ctx, directory); + + // Retention cleanup stays owner-only: it is storage-wide maintenance + // that has no reason to run once per active instance. if (webServer?.isServerOwner()) { - await performUserProfileLearning(ctx, directory); const { cleanupService } = await import("./services/cleanup-service.js"); if (await cleanupService.shouldRunCleanup()) await cleanupService.runCleanup(); } diff --git a/src/services/user-memory-learning.ts b/src/services/user-memory-learning.ts index d09ac081..a1dc9f5e 100644 --- a/src/services/user-memory-learning.ts +++ b/src/services/user-memory-learning.ts @@ -8,6 +8,7 @@ import { userProfileManager } from "./user-profile/user-profile-manager.js"; import { sortProfileItems } from "../utils/profile.js"; import type { UserProfile, UserProfileData } from "./user-profile/types.js"; import { loadOpencodeProvider } from "./ai/opencode-provider-loader.js"; +import { tryAcquireProfileLearningLock } from "./user-profile/learning-lock.js"; let isLearningRunning = false; @@ -66,6 +67,18 @@ export async function performUserProfileLearning( }); return; } + + // `isLearningRunning` only guards re-entry inside one process. Prompt selection + // is a plain SELECT and the batch is marked only after the LLM responds, so + // without cross-process exclusion two instances sharing this storage would + // analyze the same prompts and the slower writer would clobber the faster + // one's profile update. Contention skips this round; the next idle retries. + const releaseLearningLock = tryAcquireProfileLearningLock(directory); + if (!releaseLearningLock) { + log("user-profile-learning: skipped (another process holds the learning lock)"); + return; + } + isLearningRunning = true; try { const count = await userPromptManager.countUnanalyzedForUserLearning(); @@ -291,6 +304,7 @@ Rules: throw error; } finally { isLearningRunning = false; + releaseLearningLock(); } } diff --git a/src/services/user-profile/learning-lock.ts b/src/services/user-profile/learning-lock.ts new file mode 100644 index 00000000..225bc2d3 --- /dev/null +++ b/src/services/user-profile/learning-lock.ts @@ -0,0 +1,160 @@ +import { mkdirSync, readFileSync, statSync, unlinkSync, writeFileSync } from "node:fs"; +import { join } from "node:path"; +import { CONFIG } from "../../config.js"; +import { log } from "../logger.js"; + +const LEARNING_LOCK = ".profile-learning.lock"; + +/** + * Upper bound on how long a lock may be held before other processes treat it as + * abandoned. Profile learning issues LLM requests, so this has to exceed a slow + * provider round trip; it only matters when a holder dies without releasing and + * its PID has already been reused by an unrelated process. + */ +const STALE_LOCK_MS = 30 * 60 * 1000; + +/** + * Grace period during which a lock file whose contents cannot be parsed is left + * alone. `writeFileSync` is not atomic, so a reader can observe a file that was + * created but not yet filled in. Deleting it on sight would hand the lock to a + * second process while the first believes it holds it. + */ +const WRITE_WINDOW_MS = 5_000; + +interface LearningLockState { + pid: number; + acquiredAt: number; + directory: string; +} + +function isProcessAlive(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch { + return false; + } +} + +function lockPath(): string { + return join(CONFIG.storagePath || "", LEARNING_LOCK); +} + +function removeLock(path: string): void { + try { + unlinkSync(path); + } catch { + // Already gone (lost the cleanup race) or held open by the OS. Either way + // this process does not own it, so surface nothing. + } +} + +/** + * Returns the state of a lock that is still held, or null when no live holder + * remains. A stale lock is removed as a side effect so the caller can retry. + */ +function readLiveLock(path: string): LearningLockState | null { + let raw: string; + try { + raw = readFileSync(path, "utf-8"); + } catch { + return null; + } + + let state: LearningLockState; + try { + state = JSON.parse(raw) as LearningLockState; + } catch { + // Unparseable: either mid-write by a live acquirer, or genuinely corrupt. + // Respect the write window before reclaiming so we never steal a lock that + // another process is in the middle of taking. + let mtimeMs: number; + try { + mtimeMs = statSync(path).mtimeMs; + } catch { + return null; + } + if (Date.now() - mtimeMs < WRITE_WINDOW_MS) { + return { pid: -1, acquiredAt: mtimeMs, directory: "" }; + } + removeLock(path); + return null; + } + + if (!Number.isInteger(state.pid) || state.pid <= 0) { + removeLock(path); + return null; + } + + if (state.pid !== process.pid && !isProcessAlive(state.pid)) { + log("profile-learning lock: reclaiming lock from dead holder", { pid: state.pid }); + removeLock(path); + return null; + } + + const age = Date.now() - (state.acquiredAt ?? 0); + if (Number.isFinite(age) && age > STALE_LOCK_MS) { + log("profile-learning lock: reclaiming expired lock", { pid: state.pid, ageMs: age }); + removeLock(path); + return null; + } + + return state; +} + +/** + * Serializes profile learning across every process sharing this storage path. + * + * Profile learning selects a batch of prompts with a plain SELECT, issues an LLM + * request, then writes the profile and marks the batch. None of that is atomic, + * so two processes running it concurrently would analyze the same prompts twice + * and the slower writer would clobber the faster one's profile update. + * + * Returns a release function, or null when another process holds the lock — in + * which case the caller must skip this round rather than wait, because the next + * idle event will retry. + */ +export function tryAcquireProfileLearningLock(directory: string): (() => void) | null { + const path = lockPath(); + + try { + mkdirSync(CONFIG.storagePath || "", { recursive: true }); + } catch { + // Storage path already exists, or cannot be created — the write below fails + // loudly enough on its own. + } + + const holder = readLiveLock(path); + if (holder) return null; + + const state: LearningLockState = { + pid: process.pid, + acquiredAt: Date.now(), + directory, + }; + + try { + writeFileSync(path, JSON.stringify(state), { flag: "wx" }); + } catch { + // Lost the race: another process created the file between our check and + // write. `wx` is what makes that detectable rather than silently shared. + return null; + } + + let released = false; + return () => { + if (released) return; + released = true; + try { + const current = JSON.parse(readFileSync(path, "utf-8")) as LearningLockState; + if (current.pid !== process.pid) return; + } catch { + return; + } + removeLock(path); + }; +} + +export function isProfileLearningLockHeld(): boolean { + return readLiveLock(lockPath()) !== null; +} diff --git a/src/services/user-profile/user-profile-manager.ts b/src/services/user-profile/user-profile-manager.ts index 1a4232d4..91377467 100644 --- a/src/services/user-profile/user-profile-manager.ts +++ b/src/services/user-profile/user-profile-manager.ts @@ -1,5 +1,5 @@ import { join } from "node:path"; -import { existsSync, readFileSync, writeFileSync } from "node:fs"; +import { existsSync, readFileSync, statSync, writeFileSync } from "node:fs"; import { tursoConnectionManager } from "../turso/connection-manager.js"; import type { TursoDb } from "../turso/turso-db.js"; import { CONFIG } from "../../config.js"; @@ -78,6 +78,7 @@ export class UserProfileManager { // (COLD_BUFFER_DEFAULT_KEY) only holds items from merges that ran without a profileId. private coldBuffers: Map; private coldBufferPath: string; + private coldBufferMtimeMs: number | null = null; private dedupCheckedCache: Set = new Set(); constructor() { @@ -133,11 +134,34 @@ export class UserProfileManager { return { preferences: [], patterns: [], workflows: [] }; } + /** + * Discards the cached cold buffer when another process has rewritten the file. + * + * The buffer is read once in the constructor and `saveColdBuffers` rewrites the + * whole file from this in-memory map. Now that more than one process can run + * profile learning, a process holding a stale map would erase entries a peer + * persisted in the meantime. Comparing the file's mtime keeps that correctness + * requirement inside the manager instead of relying on every caller to refresh. + */ + private refreshColdBuffersIfChanged(): void { + let mtimeMs: number; + try { + mtimeMs = statSync(this.coldBufferPath).mtimeMs; + } catch { + // No file yet (or unreadable): nothing on disk can be newer than our map. + return; + } + if (this.coldBufferMtimeMs !== null && mtimeMs <= this.coldBufferMtimeMs) return; + this.coldBuffers = this.loadColdBuffers(); + this.coldBufferMtimeMs = mtimeMs; + } + private getColdBuffer(profileId?: string): { preferences: any[]; patterns: any[]; workflows: any[]; } { + this.refreshColdBuffersIfChanged(); const key = profileId || COLD_BUFFER_DEFAULT_KEY; let buf = this.coldBuffers.get(key); if (!buf) { @@ -196,6 +220,9 @@ export class UserProfileManager { } } writeFileSync(this.coldBufferPath, JSON.stringify(obj), "utf-8"); + // Record our own write so the next read does not mistake it for a peer's + // update and reload the map we just persisted. + this.coldBufferMtimeMs = statSync(this.coldBufferPath).mtimeMs; } catch { // Silently ignore disk-full / permission errors. } diff --git a/tests/profile-learning-idle.test.ts b/tests/profile-learning-idle.test.ts new file mode 100644 index 00000000..97d39a17 --- /dev/null +++ b/tests/profile-learning-idle.test.ts @@ -0,0 +1,209 @@ +import { afterAll, describe, expect, it } from "bun:test"; +import { mkdtempSync, rmSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; + +const tempDirs: string[] = []; + +afterAll(() => { + for (const dir of tempDirs) { + try { + rmSync(dir, { recursive: true, force: true }); + } catch { + // Best-effort cleanup. + } + } +}); + +const indexUrl = new URL("../src/index.js", import.meta.url).href; +const clientUrl = new URL("../src/services/client.js", import.meta.url).href; +const configUrl = new URL("../src/config.js", import.meta.url).href; +const tagsUrl = new URL("../src/services/tags.js", import.meta.url).href; +const contextUrl = new URL("../src/services/context.js", import.meta.url).href; +const privacyUrl = new URL("../src/services/privacy.js", import.meta.url).href; +const autoCaptureUrl = new URL("../src/services/auto-capture.js", import.meta.url).href; +const learningUrl = new URL("../src/services/user-memory-learning.js", import.meta.url).href; +const cleanupUrl = new URL("../src/services/cleanup-service.js", import.meta.url).href; +const promptManagerUrl = new URL( + "../src/services/user-prompt/user-prompt-manager.js", + import.meta.url +).href; +const webServerUrl = new URL("../src/services/web-server.js", import.meta.url).href; +const loggerUrl = new URL("../src/services/logger.js", import.meta.url).href; +const languageUrl = new URL("../src/services/language-detector.js", import.meta.url).href; +const tursoReadyUrl = new URL("../src/services/turso/ready.js", import.meta.url).href; + +/** + * Drives the real plugin's `session.idle` handler in an isolated Bun process so + * the module mocks below cannot leak into sibling test files, then reports which + * background jobs the handler actually invoked. + */ +function runIdleScenario(opts: { + owner?: boolean; + webServerEnabled?: boolean; + internalSession?: boolean; + autoCaptureEnabled?: boolean; +}) { + const owner = opts.owner ?? false; + const webServerEnabled = opts.webServerEnabled ?? true; + const internalSession = opts.internalSession ?? false; + const autoCaptureEnabled = opts.autoCaptureEnabled ?? true; + + const dir = mkdtempSync(join(tmpdir(), "opencode-mem-profile-idle-")); + tempDirs.push(dir); + const scriptPath = join(dir, "scenario.mjs"); + + const script = ` +import { mock } from "bun:test"; + +const calls = []; + +mock.module(${JSON.stringify(clientUrl)}, () => ({ + memoryClient: { + warmup: async () => {}, + isReady: async () => true, + close() {}, + }, +})); + +mock.module(${JSON.stringify(configUrl)}, () => ({ + CONFIG: { + autoCaptureEnabled: ${autoCaptureEnabled}, + compaction: { enabled: false, memoryLimit: 10 }, + chatMessage: { enabled: false }, + webServerEnabled: ${webServerEnabled}, + storagePath: ${JSON.stringify(dir)}, + autoCaptureProviderStatus: { ready: true, issues: [] }, + }, + initConfig: () => {}, + isConfigured: () => true, +})); + +mock.module(${JSON.stringify(tagsUrl)}, () => ({ + getTags: () => ({ project: { tag: "project-tag" }, user: { userEmail: "u@example.com" } }), +})); +mock.module(${JSON.stringify(contextUrl)}, () => ({ formatContextForPrompt: () => "" })); +mock.module(${JSON.stringify(privacyUrl)}, () => ({ + stripPrivateContent: (value) => value, + isFullyPrivate: () => false, +})); +mock.module(${JSON.stringify(promptManagerUrl)}, () => ({ userPromptManager: { savePrompt() {} } })); +mock.module(${JSON.stringify(loggerUrl)}, () => ({ log: () => {} })); +mock.module(${JSON.stringify(languageUrl)}, () => ({ getLanguageName: () => "English" })); +// The web server only starts once the Turso readiness gate passes, and that +// gate is what decides whether an owner exists at all. +mock.module(${JSON.stringify(tursoReadyUrl)}, () => ({ ensureTursoReady: async () => {} })); +// The internal-capture check lives inside src/index.ts and resolves the session +// title through the client, so drive it the real way: report the reserved title. +const INTERNAL_TITLE = "opencode-mem capture"; + +mock.module(${JSON.stringify(autoCaptureUrl)}, () => ({ + performAutoCapture: async () => { calls.push("capture"); }, +})); +mock.module(${JSON.stringify(learningUrl)}, () => ({ + performUserProfileLearning: async () => { calls.push("learn"); }, +})); +mock.module(${JSON.stringify(cleanupUrl)}, () => ({ + cleanupService: { + shouldRunCleanup: async () => true, + runCleanup: async () => { calls.push("cleanup"); }, + }, +})); + +mock.module(${JSON.stringify(webServerUrl)}, () => ({ + startWebServer: async () => (${webServerEnabled} + ? { + isServerOwner: () => ${owner}, + getUrl: () => "http://127.0.0.1:4747", + setOnTakeoverCallback: () => {}, + stop: async () => {}, + } + : null), + WebServer: class {}, +})); + +const mockClient = { + session: { + get: async () => ({ data: ${internalSession} ? { title: INTERNAL_TITLE } : { title: "regular work" } }), + messages: async () => ({ data: [] }), + }, + tui: { showToast: async () => ({}) }, +}; + +const { OpenCodeMemPlugin } = await import(${JSON.stringify(indexUrl)}); +const plugin = await OpenCodeMemPlugin({ directory: "/active-project", client: mockClient }); + +// Let the plugin's async web-server bootstrap settle before the idle event, so +// ownership is resolved the same way it is at runtime. +await new Promise((resolve) => setTimeout(resolve, 50)); + +await plugin.event({ + event: { type: "session.idle", properties: { sessionID: "user-session" } }, +}); + +// The handler debounces idle work behind a 10s timer. +await new Promise((resolve) => setTimeout(resolve, 10_600)); + +console.log(JSON.stringify({ calls })); +`; + + writeFileSync(scriptPath, script); + const result = Bun.spawnSync({ + cmd: [process.execPath, scriptPath], + stdout: "pipe", + stderr: "pipe", + }); + + const stdout = Buffer.from(result.stdout).toString("utf8").trim(); + return { + exitCode: result.exitCode, + stderr: Buffer.from(result.stderr).toString("utf8").trim(), + calls: stdout ? (JSON.parse(stdout).calls as string[]) : null, + }; +} + +describe("profile learning trigger on session.idle", () => { + it("learns from an active instance that does not own the web server", () => { + const result = runIdleScenario({ owner: false }); + + expect([result.exitCode, result.stderr]).toEqual([0, ""]); + expect(result.calls).toEqual(["capture", "learn"]); + }, 30_000); + + it("learns when the web server is disabled entirely", () => { + // Regression guard: with no web server there is no owner at all, so gating + // learning on ownership disabled profile learning permanently. + const result = runIdleScenario({ webServerEnabled: false }); + + expect([result.exitCode, result.stderr]).toEqual([0, ""]); + expect(result.calls).toEqual(["capture", "learn"]); + }, 30_000); + + it("keeps retention cleanup owner-only", () => { + const result = runIdleScenario({ owner: true }); + + expect([result.exitCode, result.stderr]).toEqual([0, ""]); + expect(result.calls).toEqual(["capture", "learn", "cleanup"]); + }, 30_000); + + it("runs no cleanup on a non-owner", () => { + const result = runIdleScenario({ owner: false }); + + expect([result.exitCode, result.stderr]).toEqual([0, ""]); + expect(result.calls).not.toContain("cleanup"); + }, 30_000); + + it("skips the plugin's own internal capture sessions", () => { + const result = runIdleScenario({ internalSession: true }); + + expect([result.exitCode, result.stderr]).toEqual([0, ""]); + expect(result.calls).toEqual([]); + }, 30_000); + + it("does nothing when auto-capture is disabled", () => { + const result = runIdleScenario({ autoCaptureEnabled: false }); + + expect([result.exitCode, result.stderr]).toEqual([0, ""]); + expect(result.calls).toEqual([]); + }, 30_000); +}); diff --git a/tests/profile-learning-lock.test.ts b/tests/profile-learning-lock.test.ts new file mode 100644 index 00000000..ca717755 --- /dev/null +++ b/tests/profile-learning-lock.test.ts @@ -0,0 +1,280 @@ +import { afterAll, describe, expect, it } from "bun:test"; +import { existsSync, mkdtempSync, readFileSync, rmSync, utimesSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; + +const tempDirs: string[] = []; + +afterAll(() => { + for (const dir of tempDirs) { + try { + rmSync(dir, { recursive: true, force: true }); + } catch { + // Best-effort cleanup. + } + } +}); + +const lockUrl = new URL("../src/services/user-profile/learning-lock.js", import.meta.url).href; +const configUrl = new URL("../src/config.js", import.meta.url).href; +const loggerUrl = new URL("../src/services/logger.js", import.meta.url).href; + +const LOCK_FILE = ".profile-learning.lock"; + +function storage(): string { + const dir = mkdtempSync(join(tmpdir(), "opencode-mem-learning-lock-")); + tempDirs.push(dir); + return dir; +} + +function preamble(storagePath: string): string { + return ` +import { mock } from "bun:test"; +mock.module(${JSON.stringify(configUrl)}, () => ({ + CONFIG: { storagePath: ${JSON.stringify(storagePath)} }, + initConfig: () => {}, + isConfigured: () => true, +})); +mock.module(${JSON.stringify(loggerUrl)}, () => ({ log: () => {} })); +const lock = await import(${JSON.stringify(lockUrl)}); +`; +} + +function runScript(dir: string, name: string, body: string) { + const scriptPath = join(dir, name); + writeFileSync(scriptPath, body); + const result = Bun.spawnSync({ + cmd: [process.execPath, scriptPath], + stdout: "pipe", + stderr: "pipe", + }); + const stdout = Buffer.from(result.stdout).toString("utf8").trim(); + return { + exitCode: result.exitCode, + stderr: Buffer.from(result.stderr).toString("utf8").trim(), + stdout, + parsed: stdout ? JSON.parse(stdout) : null, + }; +} + +/** + * Spawns two real processes that contend for the lock at the same time. + * + * Both wait on a shared start barrier before acquiring, so the acquisitions + * genuinely overlap. Running two promises in one process would only exercise the + * in-process guard and would pass even with no cross-process lock at all. + */ +function runContention(storagePath: string, barrierPath: string) { + const worker = (label: string) => ` +${preamble(storagePath)} +import { existsSync, writeFileSync, readFileSync } from "node:fs"; + +writeFileSync(${JSON.stringify(barrierPath)} + ".${label}", "ready"); +const peers = ["a", "b"].map((p) => ${JSON.stringify(barrierPath)} + "." + p); +const deadline = Date.now() + 10_000; +while (Date.now() < deadline && !peers.every((p) => existsSync(p))) { + await new Promise((r) => setTimeout(r, 5)); +} + +const release = lock.tryAcquireProfileLearningLock("/project-${label}"); +if (release) { + // Hold the lock long enough that the peer's attempt provably overlaps. + await new Promise((r) => setTimeout(r, 600)); + release(); +} +console.log(JSON.stringify({ acquired: release !== null })); +`; + + const dir = storagePath; + writeFileSync(join(dir, "worker-a.mjs"), worker("a")); + writeFileSync(join(dir, "worker-b.mjs"), worker("b")); + + const spawn = (file: string) => + Bun.spawn({ + cmd: [process.execPath, join(dir, file)], + stdout: "pipe", + stderr: "pipe", + }); + + const a = spawn("worker-a.mjs"); + const b = spawn("worker-b.mjs"); + return Promise.all( + [a, b].map(async (proc) => { + const out = (await new Response(proc.stdout).text()).trim(); + await proc.exited; + return out ? (JSON.parse(out) as { acquired: boolean }) : null; + }) + ); +} + +describe("cross-process profile learning lock", () => { + it("grants the lock to exactly one of two concurrent processes", async () => { + const dir = storage(); + const results = await runContention(dir, join(dir, "barrier")); + + const granted = results.filter((r) => r?.acquired).length; + expect(results).toHaveLength(2); + expect(granted).toBe(1); + }, 45_000); + + it("releases the lock so a later process can acquire it", () => { + const dir = storage(); + + const first = runScript( + dir, + "acquire-release.mjs", + ` +${preamble(dir)} +const release = lock.tryAcquireProfileLearningLock("/project-a"); +release?.(); +const second = lock.tryAcquireProfileLearningLock("/project-a"); +console.log(JSON.stringify({ first: release !== null, second: second !== null })); +second?.(); +` + ); + + expect(first.exitCode).toBe(0); + expect(first.parsed).toEqual({ first: true, second: true }); + expect(existsSync(join(dir, LOCK_FILE))).toBe(false); + }, 20_000); + + it("reclaims a lock whose holder process no longer exists", () => { + const dir = storage(); + // PID 2^22 is above every Linux/macOS pid_max default, so it cannot be live. + writeFileSync( + join(dir, LOCK_FILE), + JSON.stringify({ pid: 4194304, acquiredAt: Date.now(), directory: "/dead" }) + ); + + const result = runScript( + dir, + "reclaim-dead.mjs", + ` +${preamble(dir)} +const release = lock.tryAcquireProfileLearningLock("/project-a"); +console.log(JSON.stringify({ acquired: release !== null })); +release?.(); +` + ); + + expect(result.exitCode).toBe(0); + expect(result.parsed).toEqual({ acquired: true }); + }, 20_000); + + it("reclaims a lock held past the staleness deadline", () => { + const dir = storage(); + writeFileSync( + join(dir, LOCK_FILE), + JSON.stringify({ + pid: process.pid, + acquiredAt: Date.now() - 31 * 60 * 1000, + directory: "/stalled", + }) + ); + + const result = runScript( + dir, + "reclaim-stale.mjs", + ` +${preamble(dir)} +const release = lock.tryAcquireProfileLearningLock("/project-a"); +console.log(JSON.stringify({ acquired: release !== null })); +release?.(); +` + ); + + expect(result.exitCode).toBe(0); + expect(result.parsed).toEqual({ acquired: true }); + }, 20_000); + + it("does not steal a lock file that is still being written", () => { + const dir = storage(); + // A truncated file is what a reader can observe between the `wx` create and + // the content write. Reclaiming it immediately would hand the lock to a + // second process while the first believes it holds it. + writeFileSync(join(dir, LOCK_FILE), "{"); + + const result = runScript( + dir, + "write-window.mjs", + ` +${preamble(dir)} +const release = lock.tryAcquireProfileLearningLock("/project-a"); +console.log(JSON.stringify({ acquired: release !== null })); +release?.(); +` + ); + + expect(result.exitCode).toBe(0); + expect(result.parsed).toEqual({ acquired: false }); + }, 20_000); + + it("reclaims a corrupt lock file once the write window has passed", () => { + const dir = storage(); + const lockFile = join(dir, LOCK_FILE); + writeFileSync(lockFile, "{"); + const stale = new Date(Date.now() - 60_000); + // Backdate past WRITE_WINDOW_MS so the file reads as corrupt, not mid-write. + utimesSync(lockFile, stale, stale); + + const result = runScript( + dir, + "corrupt-expired.mjs", + ` +${preamble(dir)} +const release = lock.tryAcquireProfileLearningLock("/project-a"); +console.log(JSON.stringify({ acquired: release !== null })); +release?.(); +` + ); + + expect(result.exitCode).toBe(0); + expect(result.parsed).toEqual({ acquired: true }); + }, 20_000); + + it("keeps a live holder's lock and records its identity", () => { + const dir = storage(); + + const result = runScript( + dir, + "held.mjs", + ` +${preamble(dir)} +import { readFileSync } from "node:fs"; +const release = lock.tryAcquireProfileLearningLock("/project-a"); +console.log(JSON.stringify({ + held: lock.isProfileLearningLockHeld(), + pid: JSON.parse(readFileSync(${JSON.stringify(join(dir, LOCK_FILE))}, "utf-8")).pid === process.pid, +})); +release?.(); +` + ); + + expect(result.exitCode).toBe(0); + expect(result.parsed).toEqual({ held: true, pid: true }); + }, 20_000); + + it("does not release a lock that now belongs to another process", () => { + const dir = storage(); + const lockFile = join(dir, LOCK_FILE); + + const result = runScript( + dir, + "foreign-release.mjs", + ` +${preamble(dir)} +import { writeFileSync } from "node:fs"; +const release = lock.tryAcquireProfileLearningLock("/project-a"); +// Simulate the lock having been reclaimed by a different live process. +writeFileSync(${JSON.stringify(lockFile)}, JSON.stringify({ + pid: 4194304, acquiredAt: Date.now(), directory: "/other", +})); +release?.(); +console.log(JSON.stringify({ stillPresent: true })); +` + ); + + expect(result.exitCode).toBe(0); + expect(readFileSync(lockFile, "utf-8")).toContain("4194304"); + }, 20_000); +}); From 7e071c5914abd71c0d2836e632bd3f5ac97cb544 Mon Sep 17 00:00:00 2001 From: lcy Date: Thu, 1 Oct 2026 01:44:21 +0800 Subject: [PATCH 2/2] fix: harden cross-process profile learning for real deployments - Replace the FS lock with a standalone SQLite coordination DB: CAS token ownership, no TTL takeover, dead-owner reclaim via pid+boot_id+starttime - Move the in-process guard before the first await; release in a nested finally so the flag stays raised while release is in flight - Read the cold buffer as bytes with a zero-length guard and fail closed on blank/unexpected-schema files instead of treating them as empty - Keep profile parse/merge failures outside the provider fallback catch so storage errors propagate instead of retrying a broken LLM path - Add cross-process tests: real service lock (8), cold buffer (16), storage error (4), lock coordination (14); ordinary memory writes verified to proceed while the learning lock is held - Fix lint: unused import/variable, require() imports -> ESM Combined fork deployment QA: 542 pass / 0 fail / 4 skip. --- src/services/user-memory-learning.ts | 54 +- src/services/user-profile/learning-lock.ts | 460 ++++++++++++---- .../user-profile/user-profile-manager.ts | 287 +++++++--- .../profile-cold-buffer-worker-preload.ts | 51 ++ .../profile-cold-buffer-worker-shortwrite.ts | 69 +++ .../fixtures/profile-learning-lock-worker.mts | 376 +++++++++++++ .../profile-learning-memory-writer.mjs | 132 +++++ .../profile-learning-service-worker.mjs | 226 ++++++++ .../profile-learning-storage-error-worker.mjs | 321 +++++++++++ .../profile-cold-buffer-cross-process.test.ts | 512 ++++++++++++++++++ tests/profile-learning-lock.test.ts | 511 ++++++++++------- tests/profile-learning-service-lock.test.ts | 345 ++++++++++++ tests/profile-learning-storage-error.test.ts | 117 ++++ tests/user-profile-learning-error.test.ts | 6 + 14 files changed, 3075 insertions(+), 392 deletions(-) create mode 100644 tests/fixtures/profile-cold-buffer-worker-preload.ts create mode 100644 tests/fixtures/profile-cold-buffer-worker-shortwrite.ts create mode 100644 tests/fixtures/profile-learning-lock-worker.mts create mode 100644 tests/fixtures/profile-learning-memory-writer.mjs create mode 100644 tests/fixtures/profile-learning-service-worker.mjs create mode 100644 tests/fixtures/profile-learning-storage-error-worker.mjs create mode 100644 tests/profile-cold-buffer-cross-process.test.ts create mode 100644 tests/profile-learning-service-lock.test.ts create mode 100644 tests/profile-learning-storage-error.test.ts diff --git a/src/services/user-memory-learning.ts b/src/services/user-memory-learning.ts index a1dc9f5e..9a5fd660 100644 --- a/src/services/user-memory-learning.ts +++ b/src/services/user-memory-learning.ts @@ -68,19 +68,23 @@ export async function performUserProfileLearning( return; } + // Set before the first await so a same-process re-entry bounces off the flag + // instead of queueing behind this run and running a second analysis. + // // `isLearningRunning` only guards re-entry inside one process. Prompt selection // is a plain SELECT and the batch is marked only after the LLM responds, so // without cross-process exclusion two instances sharing this storage would // analyze the same prompts and the slower writer would clobber the faster // one's profile update. Contention skips this round; the next idle retries. - const releaseLearningLock = tryAcquireProfileLearningLock(directory); - if (!releaseLearningLock) { - log("user-profile-learning: skipped (another process holds the learning lock)"); - return; - } - isLearningRunning = true; + let releaseLearningLock: (() => Promise | void) | null = null; try { + releaseLearningLock = await tryAcquireProfileLearningLock(directory); + if (!releaseLearningLock) { + log("user-profile-learning: skipped (another process holds the learning lock)"); + return; + } + const count = await userPromptManager.countUnanalyzedForUserLearning(); const threshold = CONFIG.userProfileAnalysisInterval; @@ -303,8 +307,17 @@ Rules: log("user-profile-learning: aborted", { error: String(error) }); throw error; } finally { - isLearningRunning = false; - releaseLearningLock(); + // Release only when the lock was actually acquired (contention return, + // acquisition failure, and any throw before acquisition all leave it null). + // The flag resets in a nested finally so it stays raised while the release + // await is in flight — otherwise a same-process re-entry could slip in and + // find the cross-process lock already gone — and it still resets when the + // release itself rejects. + try { + await releaseLearningLock?.(); + } finally { + isLearningRunning = false; + } } } @@ -662,6 +675,14 @@ async function analyzeUserProfile( let opencodeProviderError: unknown; if (CONFIG.opencodeProvider && CONFIG.opencodeModel) { log("user-profile-learning: trying opencode provider"); + // The try/catch boundary is the provider only: LLM client construction, + // the structured-output call, and schema binding. Stored-profile parsing + // and merging happen AFTER a successful LLM response and outside this + // try — a cold-storage read/parse/merge failure must propagate to the + // caller (logged + rethrown there), never be mistaken for a provider + // fault that silently falls back to the external API and then + // re-reads/re-merges the same broken storage. + let rawData: UserProfileData | null = null; try { const { generateStructuredOutput } = await loadOpencodeProvider(); const { getOpenCodeClient } = await import("./ai/profile-llm-client.js"); @@ -708,8 +729,18 @@ Use the update_user_profile tool to save the ${existingProfile ? "updated" : "ne wfCount: result.workflows?.length, }); - const rawData = result as unknown as UserProfileData; + rawData = result as unknown as UserProfileData; + } catch (e) { + opencodeProviderError = e; + log("user-profile-learning: opencode provider failed, falling back to external API", { + error: String(e), + }); + } + // Stored-profile parse/merge runs only after a native success and is + // deliberately OUTSIDE the provider catch: storage faults here are not + // provider faults and must not trigger the external fallback. + if (rawData !== null) { if (existingProfile) { const existingData: UserProfileData = JSON.parse(existingProfile.profileData); const merged = await userProfileManager.mergeProfileData( @@ -721,11 +752,6 @@ Use the update_user_profile tool to save the ${existingProfile ? "updated" : "ne return { raw: rawData, merged }; } return { raw: rawData, merged: null }; - } catch (e) { - opencodeProviderError = e; - log("user-profile-learning: opencode provider failed, falling back to external API", { - error: String(e), - }); } } diff --git a/src/services/user-profile/learning-lock.ts b/src/services/user-profile/learning-lock.ts index 225bc2d3..c058ebd5 100644 --- a/src/services/user-profile/learning-lock.ts +++ b/src/services/user-profile/learning-lock.ts @@ -1,160 +1,412 @@ -import { mkdirSync, readFileSync, statSync, unlinkSync, writeFileSync } from "node:fs"; -import { join } from "node:path"; +import { randomUUID } from "node:crypto"; +import { mkdirSync, readFileSync } from "node:fs"; +import { dirname, resolve } from "node:path"; +import { createClient, type Client } from "@libsql/client"; import { CONFIG } from "../../config.js"; import { log } from "../logger.js"; -const LEARNING_LOCK = ".profile-learning.lock"; - /** - * Upper bound on how long a lock may be held before other processes treat it as - * abandoned. Profile learning issues LLM requests, so this has to exceed a slow - * provider round trip; it only matters when a holder dies without releasing and - * its PID has already been reused by an unrelated process. + * Standalone coordination database (SQLite via @libsql/client) inside + * CONFIG.storagePath. It is deliberately separate from the memory database: + * every statement here is short and autonomous, and no transaction is ever + * held across the LLM round trip, so lock traffic cannot block normal memory + * database work. */ -const STALE_LOCK_MS = 30 * 60 * 1000; +export const PROFILE_LEARNING_COORDINATION_DB = ".profile-learning-coordination.db"; -/** - * Grace period during which a lock file whose contents cannot be parsed is left - * alone. `writeFileSync` is not atomic, so a reader can observe a file that was - * created but not yet filled in. Deleting it on sight would hand the lock to a - * second process while the first believes it holds it. - */ -const WRITE_WINDOW_MS = 5_000; +/** Single logical lock — the table only ever holds this one row. */ +const LOCK_NAME = "profile-learning"; + +/** Bounded wait so a contended coordination DB cannot hang an idle handler. */ +const BUSY_TIMEOUT_MS = 5_000; -interface LearningLockState { +interface ProcessIdentity { pid: number; - acquiredAt: number; - directory: string; + bootId: string | null; + starttime: string | null; } -function isProcessAlive(pid: number): boolean { - try { - process.kill(pid, 0); - return true; - } catch { - return false; - } +interface ValidOwner { + ownerToken: string; + pid: number; + bootId: string | null; + starttime: string | null; } -function lockPath(): string { - return join(CONFIG.storagePath || "", LEARNING_LOCK); +type ProcStatResult = + { status: "ok"; starttime: string } | { status: "absent" } | { status: "unreadable" }; + +/** Linux boot_id is a lowercase UUID printed by the kernel. */ +const BOOT_ID_PATTERN = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i; + +function coordinationDbPath(): string { + return resolve(CONFIG.storagePath || "", PROFILE_LEARNING_COORDINATION_DB); } -function removeLock(path: string): void { +/** + * Opens the coordination database, guarantees the lock table exists, runs the + * callback, and always closes the handle. One handle per operation; nothing is + * cached across the LLM round trip. + */ +async function withCoordinationDb(fn: (client: Client) => Promise): Promise { + const dbPath = coordinationDbPath(); try { - unlinkSync(path); + mkdirSync(dirname(dbPath), { recursive: true }); } catch { - // Already gone (lost the cleanup race) or held open by the OS. Either way - // this process does not own it, so surface nothing. + // Already exists, or the statements below will fail loudly on their own. + } + const client = createClient({ url: `file:${dbPath}` }); + try { + await client.execute(`PRAGMA busy_timeout = ${BUSY_TIMEOUT_MS}`); + await client.execute(`CREATE TABLE IF NOT EXISTS profile_learning_lock ( + name TEXT PRIMARY KEY, + owner_token TEXT NOT NULL, + pid INTEGER NOT NULL, + boot_id TEXT, + starttime TEXT, + acquired_at INTEGER NOT NULL + )`); + return await fn(client); + } finally { + client.close(); } } +function rowsAffected(result: { rowsAffected?: number | bigint }): number { + return Number(result.rowsAffected ?? 0); +} + /** - * Returns the state of a lock that is still held, or null when no live holder - * remains. A stale lock is removed as a side effect so the caller can retry. + * Reads the machine's current boot id. The value is only used as a dead + * signal when it matches the strict kernel format; anything unreadable or + * malformed is treated as "unknown" and never participates in a dead + * decision. */ -function readLiveLock(path: string): LearningLockState | null { - let raw: string; +function readBootId(): string | null { try { - raw = readFileSync(path, "utf-8"); + const value = readFileSync("/proc/sys/kernel/random/boot_id", "utf-8").trim(); + return BOOT_ID_PATTERN.test(value) ? value : null; } catch { return null; } +} - let state: LearningLockState; +/** + * Reads field 22 (starttime) of /proc//stat. The comm field may contain + * spaces and parentheses, so parsing starts after the final ')' in the line; + * after that point field N lives at index N - 3 (state is field 3). + */ +function readProcStat(pid: number): ProcStatResult { + let stat: string; try { - state = JSON.parse(raw) as LearningLockState; - } catch { - // Unparseable: either mid-write by a live acquirer, or genuinely corrupt. - // Respect the write window before reclaiming so we never steal a lock that - // another process is in the middle of taking. - let mtimeMs: number; - try { - mtimeMs = statSync(path).mtimeMs; - } catch { - return null; - } - if (Date.now() - mtimeMs < WRITE_WINDOW_MS) { - return { pid: -1, acquiredAt: mtimeMs, directory: "" }; - } - removeLock(path); - return null; + stat = readFileSync(`/proc/${pid}/stat`, "utf-8"); + } catch (error) { + return (error as NodeJS.ErrnoException)?.code === "ENOENT" + ? { status: "absent" } + : { status: "unreadable" }; + } + const commEnd = stat.lastIndexOf(")"); + if (commEnd === -1) { + return { status: "unreadable" }; } + const fields = stat.slice(commEnd + 2).split(" "); + const starttime = fields[22 - 3]; + return starttime && /^\d+$/.test(starttime) + ? { status: "ok", starttime } + : { status: "unreadable" }; +} - if (!Number.isInteger(state.pid) || state.pid <= 0) { - removeLock(path); - return null; +function currentProcessIdentity(): ProcessIdentity { + const own = readProcStat(process.pid); + return { + pid: process.pid, + bootId: readBootId(), + starttime: own.status === "ok" ? own.starttime : null, + }; +} + +function coercePid(value: unknown): number | null { + if (typeof value === "number" && Number.isInteger(value) && value > 0) { + return value; } + if (typeof value === "bigint" && value > 0n && Number.isSafeInteger(Number(value))) { + return Number(value); + } + if (typeof value === "string" && /^-?\d+$/.test(value)) { + const pid = Number(value); + return Number.isInteger(pid) && pid > 0 ? pid : null; + } + return null; +} - if (state.pid !== process.pid && !isProcessAlive(state.pid)) { - log("profile-learning lock: reclaiming lock from dead holder", { pid: state.pid }); - removeLock(path); - return null; +function coerceOptionalString(value: unknown): string | null { + return typeof value === "string" && value.length > 0 ? value : null; +} + +/** + * Boot id of a stored owner record. Three distinct states: + * - "none" — column is NULL: legitimately recorded on a platform without + * /proc (conservative signal-probe fallback applies); + * - "known" — a well-formed Linux boot id string; + * - invalid — anything else (garbage string, number, object, …). Never + * null-coerced: an owner whose boot id we cannot interpret is + * an untrusted record and the lock is not acquired (fail-closed). + */ +type StoredBootId = { kind: "none" } | { kind: "known"; value: string }; + +function parseStoredBootId(value: unknown): StoredBootId | null { + if (value === null) return { kind: "none" }; + if (typeof value === "string" && BOOT_ID_PATTERN.test(value)) { + return { kind: "known", value }; } + return null; +} - const age = Date.now() - (state.acquiredAt ?? 0); - if (Number.isFinite(age) && age > STALE_LOCK_MS) { - log("profile-learning lock: reclaiming expired lock", { pid: state.pid, ageMs: age }); - removeLock(path); - return null; +/** + * Structural validation of the stored owner record. A malformed record is + * never treated as a free or reclaimable lock (fail-closed): the caller just + * skips this round instead of potentially stealing a live holder's lock. + */ +function parseOwnerRow( + row: Record +): { ok: true; owner: ValidOwner } | { ok: false; reason: string } { + const ownerToken = coerceOptionalString(row["owner_token"]); + if (!ownerToken) return { ok: false, reason: "owner_token missing/empty" }; + + const pid = coercePid(row["pid"]); + if (pid === null) return { ok: false, reason: `invalid pid: ${String(row["pid"])}` }; + + const storedBootId = parseStoredBootId(row["boot_id"]); + if (!storedBootId) return { ok: false, reason: `invalid boot_id: ${String(row["boot_id"])}` }; + const bootId = storedBootId.kind === "known" ? storedBootId.value : null; + + let starttime: string | null = null; + if (row["starttime"] !== null && row["starttime"] !== undefined) { + const raw = row["starttime"]; + if (typeof raw !== "string" || !/^\d+$/.test(raw)) { + return { ok: false, reason: "invalid starttime" }; + } + starttime = raw; + } + + const acquiredAt = row["acquired_at"]; + if ( + typeof acquiredAt !== "number" && + typeof acquiredAt !== "bigint" && + !/^\d+$/.test(String(acquiredAt ?? "")) + ) { + return { ok: false, reason: "invalid acquired_at" }; } - return state; + return { ok: true, owner: { ownerToken, pid, bootId, starttime } }; } /** - * Serializes profile learning across every process sharing this storage path. + * True only when the recorded owner is *definitively* gone from this + * machine's current PID namespace: + * - its boot id is a valid UUID differing from this machine's current valid + * boot id (the holder ran under a previous boot), or + * - /proc//stat is readable and starttime differs (PID reuse), or + * - /proc//stat is absent AND the signal-0 probe returns ESRCH — + * hidepid=2 makes a live holder's /proc entry invisible, so ENOENT alone + * is never proof of death, or + * - (no startup identity recorded) signal 0 returns ESRCH. * - * Profile learning selects a batch of prompts with a plain SELECT, issues an LLM - * request, then writes the profile and marks the batch. None of that is atomic, - * so two processes running it concurrently would analyze the same prompts twice - * and the slower writer would clobber the faster one's profile update. + * EPERM means "alive under another uid" and counts as live; any other + * error means "unknown" and is never a dead signal. Identity we cannot + * establish is never guessed, and no TTL ever overrides this check. * - * Returns a release function, or null when another process holds the lock — in - * which case the caller must skip this round rather than wait, because the next - * idle event will retry. + * Single-machine view: the coordination DB is expected to be shared only by + * processes on one host. A record written by a process in another PID + * namespace is beyond what local /proc can attest and conservatively reads + * as not-dead. */ -export function tryAcquireProfileLearningLock(directory: string): (() => void) | null { - const path = lockPath(); - +/** + * Signal-0 liveness probe. ESRCH is the only "dead" outcome: EPERM means + * alive under another uid (hidepid=2 or another user), and any other error + * means "unknown", which is never a dead signal. + */ +/** + * Signal-0 liveness probe. ESRCH is the only "dead" outcome: EPERM means + * alive under another uid (hidepid=2 or another user), and any other error + * means "unknown", which is never a dead signal. + */ +function isPidGoneBySignal(pid: number): boolean { try { - mkdirSync(CONFIG.storagePath || "", { recursive: true }); - } catch { - // Storage path already exists, or cannot be created — the write below fails - // loudly enough on its own. + process.kill(pid, 0); + return false; + } catch (error) { + return (error as NodeJS.ErrnoException)?.code === "ESRCH"; } +} - const holder = readLiveLock(path); - if (holder) return null; - - const state: LearningLockState = { - pid: process.pid, - acquiredAt: Date.now(), - directory, - }; +function ownerIsDefinitelyDead(owner: ValidOwner, self: ProcessIdentity): boolean { + if (self.bootId !== null && owner.bootId !== null && owner.bootId !== self.bootId) { + return true; + } - try { - writeFileSync(path, JSON.stringify(state), { flag: "wx" }); - } catch { - // Lost the race: another process created the file between our check and - // write. `wx` is what makes that detectable rather than silently shared. - return null; + if (self.starttime !== null && owner.starttime !== null) { + const stat = readProcStat(owner.pid); + if (stat.status === "ok") { + return stat.starttime !== owner.starttime; + } + if (stat.status === "absent") { + // /proc may be hidden (hidepid=2): fall through to the signal probe + // instead of trusting ENOENT as proof of death. + return isPidGoneBySignal(owner.pid); + } + return false; } + return isPidGoneBySignal(owner.pid); +} + +function releaseFn(ownerToken: string): () => Promise { let released = false; - return () => { + return async () => { if (released) return; released = true; try { - const current = JSON.parse(readFileSync(path, "utf-8")) as LearningLockState; - if (current.pid !== process.pid) return; - } catch { - return; + await withCoordinationDb(async (client) => { + await client.execute({ + sql: `DELETE FROM profile_learning_lock WHERE name = ? AND owner_token = ?`, + args: [LOCK_NAME, ownerToken], + }); + }); + } catch (error) { + // The owner row survives; after this process exits another process can + // reclaim it via the dead-owner CAS path. Never throw out of release. + log("profile-learning lock: release failed; row will be reclaimed after exit", { + error: String(error), + }); } - removeLock(path); }; } -export function isProfileLearningLockHeld(): boolean { - return readLiveLock(lockPath()) !== null; +/** + * Serializes profile learning across every process sharing this storage path. + * + * Ownership is a single row in a standalone SQLite coordination database: + * - acquiring is one `INSERT ... ON CONFLICT DO NOTHING` (rowsAffected 1 wins); + * - a conflicting row is only taken over with a conditional + * `UPDATE ... WHERE owner_token = ` (CAS) when the recorded + * owner is provably dead or its original process identity is gone; + * - release is `DELETE ... WHERE owner_token = `, so a stale release + * from a previous owner (even with a reused PID) can never drop someone + * else's lock. + * + * Returns a release function, or null when another process holds the lock or + * the coordination state cannot be trusted (fail-closed) — in both cases the + * caller must skip this round rather than wait, because the next idle event + * retries. + * + * Service wiring note: this function is async and must be `await`ed by the + * caller, with the in-process `isLearningRunning` guard set before the first + * `await`. + */ +export async function tryAcquireProfileLearningLock( + directory: string +): Promise<(() => Promise) | null> { + const identity = currentProcessIdentity(); + const ownerToken = randomUUID(); + + try { + return await withCoordinationDb(async (client) => { + const inserted = await client.execute({ + sql: `INSERT INTO profile_learning_lock + (name, owner_token, pid, boot_id, starttime, acquired_at) + VALUES (?, ?, ?, ?, ?, ?) + ON CONFLICT (name) DO NOTHING`, + args: [ + LOCK_NAME, + ownerToken, + identity.pid, + identity.bootId, + identity.starttime, + Date.now(), + ], + }); + if (rowsAffected(inserted) === 1) { + return releaseFn(ownerToken); + } + + const selected = await client.execute({ + sql: `SELECT owner_token, pid, boot_id, starttime, acquired_at + FROM profile_learning_lock WHERE name = ?`, + args: [LOCK_NAME], + }); + const row = selected.rows[0] as Record | undefined; + if (!row) { + // Released between our failed INSERT and the SELECT. Treat as + // contention: skip this round, the next idle event retries. + return null; + } + + const parsed = parseOwnerRow(row); + if (!parsed.ok) { + log("profile-learning lock: malformed coordination record, refusing to acquire", { + directory, + reason: parsed.reason, + }); + return null; + } + + if (!ownerIsDefinitelyDead(parsed.owner, identity)) { + // Live holder, EPERM, or uncertain identity: never reclaim, never + // guess, and no TTL is allowed to override this. + return null; + } + + const claimed = await client.execute({ + sql: `UPDATE profile_learning_lock + SET owner_token = ?, pid = ?, boot_id = ?, starttime = ?, acquired_at = ? + WHERE name = ? AND owner_token = ?`, + args: [ + ownerToken, + identity.pid, + identity.bootId, + identity.starttime, + Date.now(), + LOCK_NAME, + parsed.owner.ownerToken, + ], + }); + if (rowsAffected(claimed) === 1) { + log("profile-learning lock: reclaimed lock from dead owner", { + directory, + previousPid: parsed.owner.pid, + }); + return releaseFn(ownerToken); + } + // Lost the CAS race to another reclaimer. + return null; + }); + } catch (error) { + log("profile-learning lock: coordination database unavailable, refusing to acquire", { + directory, + error: String(error), + }); + return null; + } +} + +/** + * Whether a coordination record currently exists. Test/inspection helper; a + * database error is reported as "not held" but logged. + */ +export async function isProfileLearningLockHeld(): Promise { + try { + return await withCoordinationDb(async (client) => { + const result = await client.execute({ + sql: `SELECT 1 FROM profile_learning_lock WHERE name = ? LIMIT 1`, + args: [LOCK_NAME], + }); + return result.rows.length > 0; + }); + } catch (error) { + log("profile-learning lock: coordination database unavailable in isProfileLearningLockHeld", { + error: String(error), + }); + return false; + } } diff --git a/src/services/user-profile/user-profile-manager.ts b/src/services/user-profile/user-profile-manager.ts index 91377467..8265861b 100644 --- a/src/services/user-profile/user-profile-manager.ts +++ b/src/services/user-profile/user-profile-manager.ts @@ -1,5 +1,6 @@ import { join } from "node:path"; -import { existsSync, readFileSync, statSync, writeFileSync } from "node:fs"; +import { closeSync, openSync, readFileSync, renameSync, unlinkSync, writeSync } from "node:fs"; +import { randomUUID } from "node:crypto"; import { tursoConnectionManager } from "../turso/connection-manager.js"; import type { TursoDb } from "../turso/turso-db.js"; import { CONFIG } from "../../config.js"; @@ -78,13 +79,20 @@ export class UserProfileManager { // (COLD_BUFFER_DEFAULT_KEY) only holds items from merges that ran without a profileId. private coldBuffers: Map; private coldBufferPath: string; - private coldBufferMtimeMs: number | null = null; + /** + * Set when the last disk read of cold-buffer.json failed in a way that means + * "unknown on-disk state" (permission error, unparseable JSON, bad schema). + * A merge round must then fail closed instead of treating the data as empty + * and overwriting it. The constructor stays fail-safe: a non-critical cache + * read error must never crash the plugin import path, so it is only recorded. + */ + private coldBufferLoadError: Error | null = null; private dedupCheckedCache: Set = new Set(); constructor() { this.dbPath = join(CONFIG.storagePath || "", USER_PROFILES_DB_NAME); this.coldBufferPath = join(CONFIG.storagePath || "", "cold-buffer.json"); - this.coldBuffers = this.loadColdBuffers(); + this.coldBuffers = this.loadColdBuffersFailSafe(); } reset(): void { @@ -92,6 +100,10 @@ export class UserProfileManager { this.initPromise = null; this.dbPath = join(CONFIG.storagePath || "", USER_PROFILES_DB_NAME); this.coldBufferPath = join(CONFIG.storagePath || "", "cold-buffer.json"); + // The new storage starts from its own on-disk snapshot. Keeping the old + // map would resurrect the previous storage's buckets on the next save. + this.coldBuffers = this.loadColdBuffersFailSafe(); + this.coldBufferLoadError = null; } private async initialize(): Promise { @@ -102,7 +114,13 @@ export class UserProfileManager { this.initPromise = (async () => { try { this.dbPath = join(CONFIG.storagePath || "", USER_PROFILES_DB_NAME); - this.coldBufferPath = join(CONFIG.storagePath || "", "cold-buffer.json"); + const newColdBufferPath = join(CONFIG.storagePath || "", "cold-buffer.json"); + if (newColdBufferPath !== this.coldBufferPath) { + // storagePath changed since construction — do not keep draining or + // saving the previous storage's buckets. + this.coldBufferPath = newColdBufferPath; + this.coldBuffers = this.loadColdBuffersFailSafe(); + } this.db = await tursoConnectionManager.getConnection(this.dbPath); await this.initDatabase(); } catch (error) { @@ -135,25 +153,150 @@ export class UserProfileManager { } /** - * Discards the cached cold buffer when another process has rewritten the file. + * Fail-safe wrapper for the constructor / reset() path. A non-critical cache + * read error must not crash plugin import, so the error is recorded and an + * empty map is returned; the next merge round then fails closed. + */ + private loadColdBuffersFailSafe(): Map< + string, + { preferences: any[]; patterns: any[]; workflows: any[] } + > { + try { + return this.loadColdBuffers(); + } catch (e) { + this.coldBufferLoadError = e instanceof Error ? e : new Error(String(e)); + return new Map(); + } + } + + /** + * Reads the complete cold-buffer snapshot from disk, strictly. * - * The buffer is read once in the constructor and `saveColdBuffers` rewrites the - * whole file from this in-memory map. Now that more than one process can run - * profile learning, a process holding a stale map would erase entries a peer - * persisted in the meantime. Comparing the file's mtime keeps that correctness - * requirement inside the manager instead of relying on every caller to refresh. + * A missing file (ENOENT) legitimately means "empty snapshot". Any other + * failure (permissions, unparseable JSON, malformed schema) means the + * on-disk state is unknown and MUST be thrown to the caller — treating it + * as empty would persist an empty snapshot over data we never read. */ - private refreshColdBuffersIfChanged(): void { - let mtimeMs: number; + private loadColdBuffers(): Map< + string, + { preferences: any[]; patterns: any[]; workflows: any[] } + > { + const map = new Map(); + let raw: string; + try { + raw = readFileSync(this.coldBufferPath, "utf-8"); + } catch (e: any) { + if (e?.code === "ENOENT") return map; + throw new Error(`profile cold buffer: unreadable (${this.coldBufferPath}): ${String(e)}`, { + cause: e, + }); + } + let data: any; try { - mtimeMs = statSync(this.coldBufferPath).mtimeMs; - } catch { - // No file yet (or unreadable): nothing on disk can be newer than our map. - return; + data = JSON.parse(raw); + } catch (e) { + throw new Error( + `profile cold buffer: corrupt JSON, refusing to overwrite (${this.coldBufferPath}): ${String(e)}`, + { cause: e } + ); + } + + if (data === null || typeof data !== "object" || Array.isArray(data)) { + throw new Error( + `profile cold buffer: unexpected schema, refusing to overwrite (${this.coldBufferPath})` + ); + } + + // Explicit legacy detection: the flat format is exactly the three category + // keys (each an array, at least one present) and nothing else. Anything + // else containing those keys is an unknown shape, not "legacy". + const legacyKeys = ["preferences", "patterns", "workflows"]; + const presentLegacyKeys = legacyKeys.filter((k) => k in data); + if (presentLegacyKeys.length > 0) { + const onlyLegacyKeys = Object.keys(data).every((k) => legacyKeys.includes(k)); + const allArrays = legacyKeys.every((k) => !(k in data) || Array.isArray(data[k])); + if (onlyLegacyKeys && allArrays) { + // Legacy flat format ({ preferences, patterns, workflows }) is + // cross-user contaminated and cannot be attributed to a profile, so + // it is dropped rather than replayed. This is the documented + // migration path, not a silent data loss: the format predates + // per-profile buckets. + log("profile cold buffer: dropping legacy unattributed buffer", { + keys: presentLegacyKeys, + }); + return map; + } + throw new Error( + `profile cold buffer: legacy keys with unknown shape, refusing to overwrite (${this.coldBufferPath})` + ); + } + + // Per-profile buckets: every key must be a bucket whose category arrays + // hold real items. An invalid entry means the on-disk state is unknown — + // coercing it to [] would persist "known empty" over data we never read. + let loaded = 0; + for (const [pid, v] of Object.entries(data)) { + if (v === null || typeof v !== "object" || Array.isArray(v)) { + throw new Error( + `profile cold buffer: bucket "${pid}" is not an object, refusing to overwrite (${this.coldBufferPath})` + ); + } + const bucket: { preferences: any[]; patterns: any[]; workflows: any[] } = { + preferences: [], + patterns: [], + workflows: [], + }; + for (const category of ["preferences", "patterns", "workflows"] as const) { + const items = v[category]; + if (items === undefined) continue; + if (!Array.isArray(items)) { + throw new Error( + `profile cold buffer: bucket "${pid}.${category}" is not an array, refusing to overwrite (${this.coldBufferPath})` + ); + } + for (const item of items) { + if (item === null || typeof item !== "object" || Array.isArray(item)) { + throw new Error( + `profile cold buffer: bucket "${pid}.${category}" holds a non-object item, refusing to overwrite (${this.coldBufferPath})` + ); + } + if (typeof item.description !== "string") { + throw new Error( + `profile cold buffer: bucket "${pid}.${category}" holds an item without a string description, refusing to overwrite (${this.coldBufferPath})` + ); + } + } + bucket[category] = items; + } + map.set(pid, bucket); + loaded++; + } + if (loaded > 0) { + log("profile cold buffer: loaded from disk", { profiles: loaded }); + } + return map; + } + + /** + * Starts a merge round from the freshest complete on-disk snapshot. + * + * Cooperative learners take the learning lock, then re-read the snapshot at + * the start of every merge round, so a peer's persisted buckets are never + * clobbered by a stale in-memory map (mtime comparison cannot be trusted: + * same-mtime writes and clock rollback both hide peer updates). The loaded + * map is fixed for the whole round: every buffer mutation and save in this + * round goes through the same reference, so pending unsaved increments are + * never dropped by a mid-round reference swap. + */ + private beginColdBufferRound(): void { + try { + this.coldBuffers = this.loadColdBuffers(); + this.coldBufferLoadError = null; + } catch (e) { + // Fail closed: the on-disk state is unknown, so no mutation may proceed. + this.coldBufferLoadError = e instanceof Error ? e : new Error(String(e)); + throw this.coldBufferLoadError; } - if (this.coldBufferMtimeMs !== null && mtimeMs <= this.coldBufferMtimeMs) return; - this.coldBuffers = this.loadColdBuffers(); - this.coldBufferMtimeMs = mtimeMs; } private getColdBuffer(profileId?: string): { @@ -161,7 +304,6 @@ export class UserProfileManager { patterns: any[]; workflows: any[]; } { - this.refreshColdBuffersIfChanged(); const key = profileId || COLD_BUFFER_DEFAULT_KEY; let buf = this.coldBuffers.get(key); if (!buf) { @@ -171,60 +313,56 @@ export class UserProfileManager { return buf; } - private loadColdBuffers(): Map< - string, - { preferences: any[]; patterns: any[]; workflows: any[] } - > { - const map = new Map(); - try { - if (existsSync(this.coldBufferPath)) { - const raw = readFileSync(this.coldBufferPath, "utf-8"); - const data = JSON.parse(raw); - // Legacy flat format ({ preferences, patterns, workflows }) is cross-user - // contaminated and cannot be attributed to a profile, so it is dropped rather - // than replayed. New format is keyed by profileId. - const isLegacyFlat = - data && - typeof data === "object" && - !Array.isArray(data) && - ("preferences" in data || "patterns" in data || "workflows" in data); - if (data && typeof data === "object" && !Array.isArray(data) && !isLegacyFlat) { - let loaded = 0; - for (const [pid, v] of Object.entries(data)) { - map.set(pid, { - preferences: Array.isArray(v?.preferences) ? v.preferences : [], - patterns: Array.isArray(v?.patterns) ? v.patterns : [], - workflows: Array.isArray(v?.workflows) ? v.workflows : [], - }); - loaded++; - } - if (loaded > 0) { - log("profile cold buffer: loaded from disk", { profiles: loaded }); - } - } else if (isLegacyFlat) { - log("profile cold buffer: dropping legacy unattributed buffer"); - } + /** + * Persists the current snapshot atomically: write a unique tmp file with + * O_EXCL (wx) and mode 0600, then rename over the target. Failures are + * thrown to the caller — swallowing them would let a round that never + * reached disk be reported as a successful learning merge. + */ + private saveColdBuffers(): void { + const obj: Record = {}; + for (const [pid, v] of this.coldBuffers.entries()) { + if (v.preferences.length || v.patterns.length || v.workflows.length) { + obj[pid] = v; } - } catch { - // Corrupt or missing file — start with an empty buffer set. } - return map; - } - - private saveColdBuffers(): void { + const tmpPath = `${this.coldBufferPath}.${randomUUID()}.tmp`; + let fd: number | undefined; try { - const obj: Record = {}; - for (const [pid, v] of this.coldBuffers.entries()) { - if (v.preferences.length || v.patterns.length || v.workflows.length) { - obj[pid] = v; + fd = openSync(tmpPath, "wx", 0o600); + // Byte-accurate payload: the string overload's third arg is a file + // position, not a string offset, and payload.length counts UTF-16 + // units — both would corrupt multi-byte content on short writes. + const payload = Buffer.from(JSON.stringify(obj), "utf8"); + let written = 0; + while (written < payload.length) { + // Buffer overload: (fd, buffer, offset, length) — append at the + // current position; n <= 0 means no forward progress is possible. + const n = writeSync(fd!, payload, written, payload.length - written); + if (!(n > 0)) { + throw new Error(`profile cold buffer: write stalled after ${written} bytes`); } + written += n; } - writeFileSync(this.coldBufferPath, JSON.stringify(obj), "utf-8"); - // Record our own write so the next read does not mistake it for a peer's - // update and reload the map we just persisted. - this.coldBufferMtimeMs = statSync(this.coldBufferPath).mtimeMs; - } catch { - // Silently ignore disk-full / permission errors. + closeSync(fd); + fd = undefined; + renameSync(tmpPath, this.coldBufferPath); + } catch (e) { + if (fd !== undefined) { + try { + closeSync(fd); + } catch { + // best effort — the rename below is the cleanup that matters + } + } + try { + unlinkSync(tmpPath); + } catch { + // tmp file may never have been created + } + throw new Error(`profile cold buffer: save failed (${this.coldBufferPath}): ${String(e)}`, { + cause: e, + }); } } @@ -515,6 +653,10 @@ export class UserProfileManager { async deleteProfile(profileId: string): Promise { const db = await this.ready(); await db.run(`DELETE FROM user_profiles WHERE id = ?`, [profileId]); + // Re-read the current on-disk snapshot before deleting the bucket: a + // stale in-memory map (peer wrote meanwhile) would resurrect the deleted + // profile's bucket in the very save that is supposed to remove it. + this.beginColdBufferRound(); if (this.coldBuffers.delete(profileId)) { this.saveColdBuffers(); } @@ -567,6 +709,13 @@ export class UserProfileManager { embedService?: EmbeddingService, profileId?: string ): Promise { + // Cross-process safety: cooperative learners hold the learning lock and + // start every round from the freshest complete on-disk snapshot. Fails + // closed when the on-disk state is unknown (see beginColdBufferRound). + // The snapshot is fixed for the whole round — getColdBuffer/saveColdBuffers + // never swap the map mid-round, so pending unsaved increments survive. + this.beginColdBufferRound(); + const merged: UserProfileData = { preferences: this.ensureArray(existing?.preferences), patterns: this.ensureArray(existing?.patterns), diff --git a/tests/fixtures/profile-cold-buffer-worker-preload.ts b/tests/fixtures/profile-cold-buffer-worker-preload.ts new file mode 100644 index 00000000..02111ecf --- /dev/null +++ b/tests/fixtures/profile-cold-buffer-worker-preload.ts @@ -0,0 +1,51 @@ +/** + * Bun --preload patch for short-write simulation (see + * profile-cold-buffer-worker-shortwrite.ts for the worker itself). + * + * Loaded via `bun --preload` so the patch applies before any ESM import + * binds writeSync. It only limits how many bytes a single low-level + * fs.writeSync call forwards — the production saveColdBuffers() code path, + * the real tmp+rename flow and the real fd bookkeeping all stay live. + * + * mode is read from argv[2] of the *worker* command line: + * short-ascii | short-utf8 | zero | half + */ +// The CJS namespace is the only patchable surface: the ESM namespace binding +// of node:fs is frozen (Object.defineProperty on it throws), while Bun's +// --preload runs before ESM imports bind writeSync, so patching the CJS +// module also intercepts the production ESM import below. +// eslint-disable-next-line @typescript-eslint/no-require-imports +const fs = require("node:fs") as typeof import("node:fs"); +const real = fs.writeSync; +const argv = process.argv as unknown as string[]; +const mode = argv.find((a) => a.startsWith("SWMODE="))?.slice("SWMODE=".length) ?? ""; +let callIndex = 0; + +(fs as any).writeSync = function (fd: number, ...rest: any[]) { + if (fd === 1 || fd === 2) { + return (real as any)(fd, ...rest); + } + callIndex++; + const buffer = rest[0] as Buffer; + if (mode === "zero") { + if (callIndex === 1) return 0; + return (real as any)(fd, ...rest); + } + if (mode === "half") { + // Call 1 writes half the payload; call 2 fails mid-payload with EIO so + // the cleanup path runs with real partial bytes in the tmp file. + if (callIndex === 1) { + const len = rest[2] as number; + return (real as any)(fd, buffer, rest[1], Math.max(0, Math.floor(len / 2))); + } + if (callIndex === 2) { + const err = new Error("simulated EIO mid-payload") as any; + err.code = "EIO"; + throw err; + } + return (real as any)(fd, ...rest); + } + // short-ascii / short-utf8: one byte per call — multi-byte chars are sliced + // mid-codepoint across calls, proving the write loop is byte-accurate. + return (real as any)(fd, buffer, rest[1], 1); +}; diff --git a/tests/fixtures/profile-cold-buffer-worker-shortwrite.ts b/tests/fixtures/profile-cold-buffer-worker-shortwrite.ts new file mode 100644 index 00000000..0bd8070f --- /dev/null +++ b/tests/fixtures/profile-cold-buffer-worker-shortwrite.ts @@ -0,0 +1,69 @@ +/** + * Child-process worker for short-write simulation in the cold-buffer tests. + * Pair with profile-cold-buffer-worker-preload.ts via `bun --preload`. + * + * The preload patch only limits how many bytes a single low-level + * fs.writeSync call forwards — the production saveColdBuffers() path, the + * real tmp+rename flow and the real fd bookkeeping all stay live. + * + * Args: dir profileId description preExisting(""=none) + * Env/argv: SWMODE=short-ascii|short-utf8|zero|half (read by the preload) + * Prints one JSON line: { ok, error, targetBytes, tmpLeft } + */ +import * as fs from "node:fs"; +import { join } from "node:path"; + +// Drop non-positional flags (SWMODE=... is consumed by the preload). +const args = process.argv.slice(2).filter((a) => !a.startsWith("SWMODE=")); +const dir = args[0]!; +const profileId = args[1]!; +const description = args[2]!; +const preExisting = args[3]!; + +const coldBufferPath = join(dir, "cold-buffer.json"); +if (preExisting === "") { + try { + fs.unlinkSync(coldBufferPath); + } catch { + // no file yet + } +} else { + fs.writeFileSync(coldBufferPath, preExisting, "utf-8"); +} + +const { CONFIG } = await import("../../src/config.js"); +CONFIG.storagePath = dir; +CONFIG.userProfileEmbeddingMinDescriptionLength = 5; +delete CONFIG.opencodeProvider; +delete CONFIG.opencodeModel; +delete CONFIG.memoryModel; +delete CONFIG.memoryApiUrl; + +const { UserProfileManager } = + await import("../../src/services/user-profile/user-profile-manager.js"); +const coldEmbed = { isWarmedUp: false } as any; + +const mgr = new UserProfileManager(); +let ok = true; +let error: string | null = null; +try { + await mgr.mergeProfileData( + { preferences: [], patterns: [], workflows: [] }, + { preferences: [{ category: "style", description }] }, + coldEmbed, + profileId + ); +} catch (e) { + ok = false; + error = String(e); +} + +let targetBytes: string | null; +try { + targetBytes = fs.readFileSync(coldBufferPath, "utf-8"); +} catch { + targetBytes = null; +} +const tmpLeft = fs.readdirSync(dir).filter((f) => f.includes(".tmp")).length; + +console.log(JSON.stringify({ ok, error, targetBytes, tmpLeft })); diff --git a/tests/fixtures/profile-learning-lock-worker.mts b/tests/fixtures/profile-learning-lock-worker.mts new file mode 100644 index 00000000..b4a9f380 --- /dev/null +++ b/tests/fixtures/profile-learning-lock-worker.mts @@ -0,0 +1,376 @@ +/** + * Cross-process worker for tests/profile-learning-lock.test.ts. + * + * Runs in a real separate process with CONFIG.storagePath mocked to a temp + * directory, so it never touches the real memory database. All coordination + * with the test parent uses explicit barrier files under the storage dir — + * never sleeps — and every wait is bounded so a broken handshake fails fast + * instead of hanging the suite. + * + * Env contract: + * PLL_STORAGE temp storage path (also hosts the coordination DB + barriers) + * PLL_MODE one of the modes below + * + * Modes: + * hold acquire; signal hs.acquired; wait hs.release; release; + * report held-before/after-release + * try wait hs.acquired (winner confirmed holding), attempt acquire, + * release immediately if won, report outcome + * hold-forever acquire; signal hs.acquired; wait until killed (SIGKILL test) + * probe single acquire attempt (live/EPERM/backdated/reuse/malformed/ + * db-error scenarios); release if acquired + * proc-probe plant owner with pid whose /proc is absent (2^22), then + * force kill(0) outcome via PLL_FORCE_KILL env (EPERM|OK|ESRCH); + * report acquire result. EPERM injected by stubbing + * process.kill in-process — NOT by relying on PID 1 as root. + * boot-variants plant owner rows one at a time (valid-null, invalid-string, + * number, object, same-as-host) and report acquire result for + * each, releasing between attempts. + * cas-race signal hs.ready.