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..9a5fd660 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,8 +67,24 @@ 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. 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; @@ -290,7 +307,17 @@ Rules: log("user-profile-learning: aborted", { error: String(error) }); throw error; } finally { - isLearningRunning = false; + // 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; + } } } @@ -648,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"); @@ -694,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( @@ -707,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 new file mode 100644 index 00000000..c058ebd5 --- /dev/null +++ b/src/services/user-profile/learning-lock.ts @@ -0,0 +1,412 @@ +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"; + +/** + * 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. + */ +export const PROFILE_LEARNING_COORDINATION_DB = ".profile-learning-coordination.db"; + +/** 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 ProcessIdentity { + pid: number; + bootId: string | null; + starttime: string | null; +} + +interface ValidOwner { + ownerToken: string; + pid: number; + bootId: string | null; + starttime: string | null; +} + +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); +} + +/** + * 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 { + mkdirSync(dirname(dbPath), { recursive: true }); + } catch { + // 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); +} + +/** + * 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 readBootId(): string | null { + try { + const value = readFileSync("/proc/sys/kernel/random/boot_id", "utf-8").trim(); + return BOOT_ID_PATTERN.test(value) ? value : null; + } catch { + return null; + } +} + +/** + * 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 { + 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" }; +} + +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; +} + +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; +} + +/** + * 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 { ok: true, owner: { ownerToken, pid, bootId, starttime } }; +} + +/** + * 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. + * + * 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. + * + * 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. + */ +/** + * 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 { + process.kill(pid, 0); + return false; + } catch (error) { + return (error as NodeJS.ErrnoException)?.code === "ESRCH"; + } +} + +function ownerIsDefinitelyDead(owner: ValidOwner, self: ProcessIdentity): boolean { + if (self.bootId !== null && owner.bootId !== null && owner.bootId !== self.bootId) { + return true; + } + + 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 async () => { + if (released) return; + released = true; + try { + 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), + }); + } + }; +} + +/** + * 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 1a4232d4..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, 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,12 +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; + /** + * 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 { @@ -91,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 { @@ -101,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) { @@ -133,71 +152,217 @@ export class UserProfileManager { return { preferences: [], patterns: [], workflows: [] }; } - private getColdBuffer(profileId?: string): { - preferences: any[]; - patterns: any[]; - workflows: any[]; - } { - const key = profileId || COLD_BUFFER_DEFAULT_KEY; - let buf = this.coldBuffers.get(key); - if (!buf) { - buf = this.emptyColdBuffer(); - this.coldBuffers.set(key, buf); + /** + * 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(); } - return buf; } + /** + * Reads the complete cold-buffer snapshot from disk, strictly. + * + * 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 loadColdBuffers(): Map< string, { preferences: any[]; patterns: any[]; workflows: any[] } > { const map = new Map(); + let raw: string; 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++; + 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 { + 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 (loaded > 0) { - log("profile cold buffer: loaded from disk", { profiles: loaded }); + 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})` + ); } - } else if (isLegacyFlat) { - log("profile cold buffer: dropping legacy unattributed buffer"); } + bucket[category] = items; } - } catch { - // Corrupt or missing file — start with an empty buffer set. + 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; + } + } + + private getColdBuffer(profileId?: string): { + preferences: any[]; + patterns: any[]; + workflows: any[]; + } { + const key = profileId || COLD_BUFFER_DEFAULT_KEY; + let buf = this.coldBuffers.get(key); + if (!buf) { + buf = this.emptyColdBuffer(); + this.coldBuffers.set(key, buf); + } + return buf; + } + + /** + * 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; + } + } + 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; + } + 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 } } - writeFileSync(this.coldBufferPath, JSON.stringify(obj), "utf-8"); - } catch { - // Silently ignore disk-full / permission errors. + 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, + }); } } @@ -488,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(); } @@ -540,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.