diff --git a/docs/SQLITE_CONCURRENCY.md b/docs/SQLITE_CONCURRENCY.md index 448c5ed..fac4ebe 100644 --- a/docs/SQLITE_CONCURRENCY.md +++ b/docs/SQLITE_CONCURRENCY.md @@ -1,8 +1,64 @@ -# SQLite contention regression (#162 / #163) +# SQLite contention recovery (#162 / #163) + +## Ordinary writes and maintenance + +The local compatibility bridge installs recovery on the native session store +after successful migration. The upstream store remains the owner of SQL, +transactions, IDs and session state. No schema changes or whole-turn retries are +introduced. + +- Audited input, message, session, permission and todo operations may retry + `SQLITE_BUSY` (including its extended codes). An operation retries only after + its transaction has rolled back, or when its existing ID-based writes are + idempotent. `SQLITE_LOCKED`, constraint, filesystem and other errors propagate. +- Writes through the bridge serialize by database path within one process. + Conversation writes precede queued usage maintenance; ordering within each + class is FIFO. Existing externally owned transactions retain native behavior + and are never replayed by the bridge. Counter/goal and legacy workflow writes + are serialized but not replayed because a partial commit is not idempotent. +- Retryable operations use a 25 ms native busy wait, scoped to each synchronous + SQLite call. The original connection timeout is restored before returning + from that call, including failures. Async backoff uses jitter and a maximum + delay of 200 ms. The total budget is 30 seconds, including time spent in the + local queue. Busy usage writes serve pending conversation writes between + attempts while later usage writes remain ordered behind them. + A synchronous SQL statement already executing cannot be preempted. +- Closing a store cancels its queued writes and backoff. Exhaustion preserves + the SQLite cause and reports the operation, attempts and elapsed time. Model + calls, shell commands and completed tools are outside the retry boundary. +- Automatic usage cleanup runs at most once per connection per five minutes, + after checking whether any rows are expired. A busy cleanup is deferred for + one second and does not fail an already committed usage fact. Explicit + `beforeTime` cleanup requests retain the upstream behavior and error handling. + Retention remains 30 days, and successful usage writes are never buffered. + +The synchronous dynamic-workflow journal and synchronous permission-mode API +retain their native contract and timeout. They do not acquire an async retry +capability through this bridge. Cross-process contention is coordinated by +SQLite and bounded backoff, not by the process-local queue. + +The operation allowlist and maintenance SQL were reviewed against upstream +`zai-org/ZCode` commit `29628c9acdb81b703bbd4080c207a0e7ce5e276e`. The required +`sqlite-write-recovery` patch validates both post-migration hooks and the native +cleanup transaction. Missing, ambiguous or partially patched anchors stop +synchronization. The implementation is in `src/runtime-sqlite-recovery.ts`, +with scoped native waits and cleanup admission in its two SQLite helpers; +`vendor/cli-config.cjs` is their existing distribution boundary. + +Required regression coverage includes both store-open paths, recovered and +exhausted locks, responsive timers during backoff, same-process connections, +input-promotion rollback, ordered updates, close during recovery, preserved +non-BUSY errors, and sustained writes from 10–15 independent processes. Cleanup +tests must verify retention and that repeated usage writes do not each open a +cleanup transaction. All storage regressions run on real Node SQLite against +the extracted release artifact. + +## Original timeout mitigation The CLI's shared session database uses SQLite WAL. Readers can overlap a writer, -but different CLI processes still serialize writes to this file. After successful -store initialization, each connection now waits up to 10 seconds for a write lock. +but different CLI processes still serialize writes to this file. The original +#163 mitigation sets a 10-second connection timeout after initialization. This +remains the default outside the scoped retryable calls described above. Both the synchronous constructor and asynchronous `openStartup()` path apply this setting **after** migrations finish. Startup migration lock budgets, short busy waits, backoff, rollback, and failure cleanup are unchanged. @@ -29,17 +85,28 @@ requests. They cover: - the effective `PRAGMA busy_timeout` after sync open, async startup, and reopen; - a native session write blocked by another process for 6.5 seconds, longer than the previous 5-second timeout, then succeeding exactly once; -- a deliberately shortened timeout that reports `SQLITE_BUSY`, writes no session, - and permits a later write after the lock is released; +- recovery after a 10.5-second lock, with timers progressing during backoff; +- a deliberately shortened recovery budget that reports `SQLITE_BUSY`, writes + no session, and permits a later write after the lock is released; - four independent processes migrating the same fresh database and persisting distinct sessions without duplicate migrations or integrity errors; - killing a writer with uncommitted changes, then reopening without those changes - and successfully writing another session. + and successfully writing another session; +- same-process connections completing an asynchronous native transaction without + starving its commit, plus caller-owned transactions retaining rollback ownership; +- a failure after message/part insertion that rolls back before input promotion + retries, and an idempotent partial part write that cannot overtake a newer update; +- non-BUSY errors and partially committed goal counters never being replayed; +- closing a store cancelling both backoff and pending writes; +- 100 usage facts requiring one cleanup transaction when expired rows exist, + with recent records retained and explicit cleanup still honored; +- retrying usage writes yielding priority to conversation writes; +- twelve independent processes persisting 300 input promotions and matching + messages, parts and usage records, followed by integrity and foreign-key checks. -The hold-and-release test needs a separate process: `DatabaseSync` blocks the -calling event loop, so a timer in that same process cannot release the lock. -Each test creates its own database and holder rather than reusing the remainder -of a lock window from a previous measurement. +Lock holders use independent processes. One test schedules their release from +the writer's event loop to verify that the new recovery yields instead of +blocking that timer. Each test creates its own database and holder. CI builds and packs once, then tests that exact artifact on Node 22.19.0, 24, and 26 on Linux, plus Node 24 on macOS. Bun's `node:sqlite` compatibility implementation @@ -47,10 +114,11 @@ is not used as a substitute for Node's SQLite driver. ## Follow-up boundary -If real workloads still exceed the wait budget, investigate long transactions and -add bounded retries only at persistence boundaries that can be safely rolled back -and replayed. Never retry a whole agent turn and repeat completed external tools. -Per-session databases or a shared writer service are separate architectural changes -requiring discovery, lifecycle, and migration design. The contention regression -tests also pass against the locked Desktop 3.14.0 runtime. They do not establish -that sustained high-concurrency workloads are free of contention. +The regressions use the locked Desktop 3.14.3 / runtime 0.16.9 artifact. They do +not establish that every production workload is free of contention. Keep #162 +open for validation on the reported workload. On recovery exhaustion, the error +includes `operation`, `attempts`, `elapsedMs` and the original SQLite cause; +use those fields to identify the next persistence boundary needing investigation. +Synchronous journal writes, non-idempotent counter operations, larger transactions +and cross-process fairness remain separate follow-up areas. Per-session databases +or a shared writer service require lifecycle and migration design. diff --git a/scripts/check-runtime.ts b/scripts/check-runtime.ts index a2590eb..a441b47 100755 --- a/scripts/check-runtime.ts +++ b/scripts/check-runtime.ts @@ -18,6 +18,7 @@ import { hasRuntimeHttpNoContentGuard, hasRuntimeNetworkRetryGuard, hasRuntimeSqliteBusyTimeout, + hasRuntimeSqliteWriteRecovery, hasRuntimeStreamEofFinishGuard, patchRuntimeGoalFailurePause, patchRuntimeHttpNoContent, @@ -25,6 +26,7 @@ import { patchRuntimeNetworkRetryClassification, patchRuntimeOfficialMcpAvailability, patchRuntimeSqliteBusyTimeout, + patchRuntimeSqliteWriteRecovery, patchRuntimeStreamEofFinishGuard, parseRuntimePatchReports, runtimePatchPlan, @@ -85,6 +87,8 @@ if (patchRuntimeLoginModelDefaults(runtimeSource) !== runtimeSource || !hasRuntimeNetworkRetryGuard(runtimeSource) || patchRuntimeSqliteBusyTimeout(runtimeSource) !== runtimeSource || !hasRuntimeSqliteBusyTimeout(runtimeSource) + || patchRuntimeSqliteWriteRecovery(runtimeSource) !== runtimeSource + || !hasRuntimeSqliteWriteRecovery(runtimeSource) || patchRuntimeStreamEofFinishGuard(runtimeSource) !== runtimeSource || !hasRuntimeStreamEofFinishGuard(runtimeSource) || (patchEnabled("cli-help-contract") && !hasRuntimeCliHelpContract(runtimeSource)) diff --git a/scripts/runtime-sqlite-patches.ts b/scripts/runtime-sqlite-patches.ts new file mode 100644 index 0000000..fcbce96 --- /dev/null +++ b/scripts/runtime-sqlite-patches.ts @@ -0,0 +1,88 @@ +export const sqliteBusyTimeoutMs = 10_000; + +const pragma = `pragma busy_timeout = ${sqliteBusyTimeoutMs}`; +const helper = 'require(require("node:path").join(__dirname,"cli-config.cjs"))'; +const identifier = "[A-Za-z_$][\\w$]*"; +const boundedIdentifier = "[A-Za-z_$][\\w$]{0,80}"; + +function escape(value: string): string { + return value.replace(/[.*+?^${}()|[\]\\]/gu, "\\$&"); +} + +function count(source: string, pattern: RegExp): number { + return [...source.matchAll(pattern)].length; +} + +/** Both open paths must retain their steady-state timeout after migration. */ +export function hasRuntimeSqliteBusyTimeout(runtime: string): boolean { + const recovery = escape(`,${helper}.installSqliteWriteRecovery(`); + const sync = new RegExp(`try\\{${identifier}!==${identifier}&&\\(${identifier}\\(this\\.db,this\\.dbPath,${identifier}\\),this\\.db\\.exec\\("${pragma}"\\)(?:${recovery}this\\))?\\)\\}catch`, "gu"); + const startup = new RegExp(`try\\{return await ${identifier}\\((${identifier})\\.db,\\1\\.dbPath,${identifier}\\),\\1\\.db\\.exec\\("${pragma}"\\)(?:${recovery}\\1\\))?,\\1\\}catch`, "gu"); + return count(runtime, sync) === 1 && count(runtime, startup) === 1; +} + +export function patchRuntimeSqliteBusyTimeout(runtime: string): string { + if (hasRuntimeSqliteBusyTimeout(runtime)) return runtime; + const sync = /try\{([A-Za-z_$][\w$]*)!==([A-Za-z_$][\w$]*)&&([A-Za-z_$][\w$]*)\(this\.db,this\.dbPath,([A-Za-z_$][\w$]*)\)\}catch/gu; + const startup = /try\{return await ([A-Za-z_$][\w$]*)\(([A-Za-z_$][\w$]*)\.db,\2\.dbPath,([A-Za-z_$][\w$]*)\),\2\}catch/gu; + if (count(runtime, sync) !== 1 || count(runtime, startup) !== 1) { + throw new Error("ZCode runtime is incompatible with the SQLite busy-timeout patch (store migration anchors missing or ambiguous)."); + } + const patched = runtime.replace(sync, (_match, mode, deferred, migrate, timeout) => + `try{${mode}!==${deferred}&&(${migrate}(this.db,this.dbPath,${timeout}),this.db.exec("${pragma}"))}catch` + ).replace(startup, (_match, migrate, store, options) => + `try{return await ${migrate}(${store}.db,${store}.dbPath,${options}),${store}.db.exec("${pragma}"),${store}}catch` + ); + if (!hasRuntimeSqliteBusyTimeout(patched)) throw new Error("SQLite busy-timeout patch failed postcondition verification."); + return patched; +} + +export function hasRuntimeSqliteWriteRecovery(runtime: string): boolean { + const sync = `${helper}.installSqliteWriteRecovery(this)`; + const startup = new RegExp(`${escape(`.db.exec("${pragma}"),${helper}.installSqliteWriteRecovery(`)}(${identifier})\\),\\1\\}catch`, "gu"); + return hasRuntimeSqliteBusyTimeout(runtime) + && runtime.split(sync).length === 2 + && count(runtime, startup) === 1 + && runtime.split(`${helper}.pruneSqliteUsage(`).length === 2; +} + +/** Keep the native SQL body and transaction; only the maintenance admission moves into our helper. */ +function patchUsageCleanup(runtime: string): string { + const labels = [...runtime.matchAll(new RegExp(`${boundedIdentifier}\\((${boundedIdentifier}),"pruneUsage"\\)`, "gu"))]; + if (labels.length !== 1) throw new Error("SQLite recovery pruneUsage symbol is missing or ambiguous."); + const name = labels[0]![1]!; + const prefix = new RegExp(`async function ${escape(name)}\\((${boundedIdentifier}),(${boundedIdentifier})=\\{\\}\\)\\{(?:let|const) (${boundedIdentifier})=\\2\\.beforeTime\\?\\?Date\\.now\\(\\)-[^;]{1,100};`, "gu"); + const starts = [...runtime.matchAll(prefix)]; + if (starts.length !== 1) throw new Error("SQLite recovery cleanup boundary is incompatible."); + const start = starts[0]!; + const [, db, options, cutoff] = start; + const bodyStart = start.index! + start[0].length; + const tail = new RegExp(`\\}catch\\((${boundedIdentifier})\\)\\{throw ${escape(db!)}\\.exec\\("rollback"\\),\\1\\}\\}`, "u"); + const end = tail.exec(runtime.slice(bodyStart, bodyStart + 4_000)); + if (!end) throw new Error("SQLite recovery cleanup transaction is incompatible."); + const bodyEnd = bodyStart + end.index + end[0].length - 1; + const body = runtime.slice(bodyStart, bodyEnd); + if (/\bawait\b|\byield\b/u.test(body) + || !body.startsWith(`${db}.exec("begin immediate");try{`) + || !["model_usage", "turn_usage", "tool_usage"].every(table => + body.includes(`${db}.prepare("delete from ${table} where started_at < ?").run(${cutoff})`)) + || !body.includes(`${db}.exec("commit")`)) { + throw new Error("SQLite recovery cleanup SQL changed; review the native transaction before patching."); + } + const replacement = `${start[0]}return ${helper}.pruneSqliteUsage(${db},${options},${cutoff},()=>{${body}})}`; + return runtime.slice(0, start.index) + replacement + runtime.slice(bodyEnd + 1); +} + +export function patchRuntimeSqliteWriteRecovery(runtime: string): string { + if (hasRuntimeSqliteWriteRecovery(runtime)) return runtime; + if (runtime.includes(".installSqliteWriteRecovery(") || runtime.includes(".pruneSqliteUsage(")) { + throw new Error("SQLite recovery patch is partial or ambiguous."); + } + if (!hasRuntimeSqliteBusyTimeout(runtime)) throw new Error("SQLite recovery requires the post-migration timeout patch."); + const timeout = new RegExp(`(${identifier})\\.db\\.exec\\("${pragma}"\\)`, "gu"); + if (count(runtime, timeout) !== 2) throw new Error("SQLite recovery open boundaries are ambiguous."); + let patched = runtime.replace(timeout, (call, store) => `${call},${helper}.installSqliteWriteRecovery(${store})`); + patched = patchUsageCleanup(patched); + if (!hasRuntimeSqliteWriteRecovery(patched)) throw new Error("SQLite recovery patch failed postcondition verification."); + return patched; +} diff --git a/scripts/sync-runtime.ts b/scripts/sync-runtime.ts index 7e90a36..8733422 100755 --- a/scripts/sync-runtime.ts +++ b/scripts/sync-runtime.ts @@ -14,6 +14,14 @@ import { } from "../src/runtime-capabilities.ts"; import { parseReleaseVersion, syncedReleaseVersion } from "./release-version.ts"; import { markRuntimeModified } from "./runtime-attribution.ts"; +import { + hasRuntimeSqliteBusyTimeout, hasRuntimeSqliteWriteRecovery, + patchRuntimeSqliteBusyTimeout, patchRuntimeSqliteWriteRecovery +} from "./runtime-sqlite-patches.ts"; +export { + hasRuntimeSqliteBusyTimeout, hasRuntimeSqliteWriteRecovery, + patchRuntimeSqliteBusyTimeout, patchRuntimeSqliteWriteRecovery, sqliteBusyTimeoutMs +} from "./runtime-sqlite-patches.ts"; const root = resolve(dirname(fileURLToPath(import.meta.url)), ".."); const cdnRoot = "https://cdn-zcode.z.ai/zcode/electron/releases"; @@ -1105,44 +1113,6 @@ export function patchRuntimeHttpNoContent(runtime: string): string { return changed ? patched : runtime; } -export const sqliteBusyTimeoutMs = 10_000; - -const runtimeSqliteBusyTimeoutPragma = `pragma busy_timeout = ${sqliteBusyTimeoutMs}`; - -/** Verify both successful store-open paths, not an unrelated pragma string. */ -export function hasRuntimeSqliteBusyTimeout(runtime: string): boolean { - const pragma = escapeRegExpName(runtimeSqliteBusyTimeoutPragma); - const sync = new RegExp(`try\\{[A-Za-z_$][\\w$]*!==[A-Za-z_$][\\w$]*&&\\([A-Za-z_$][\\w$]*\\(this\\.db,this\\.dbPath,[A-Za-z_$][\\w$]*\\),this\\.db\\.exec\\("${pragma}"\\)\\)\\}catch`, "gu"); - const startup = new RegExp(`try\\{return await [A-Za-z_$][\\w$]*\\(([A-Za-z_$][\\w$]*)\\.db,\\1\\.dbPath,[A-Za-z_$][\\w$]*\\),\\1\\.db\\.exec\\("${pragma}"\\),\\1\\}catch`, "gu"); - return countRegExpMatches(runtime, sync) === 1 && countRegExpMatches(runtime, startup) === 1; -} - -/** - * Concurrent zcode processes share ~/.zcode/cli/db/db.sqlite. The runtime opens - * it with a 5s timeout. Async startup temporarily uses a short timeout for - * migration retries and resets it in finally. Apply the steady-state timeout - * AFTER either migration path succeeds, preserving startup's lock budget and - * cleanup. This bounds ordinary write contention; it is not a transaction or - * whole-turn retry, and cannot guarantee success under sustained contention. - */ -export function patchRuntimeSqliteBusyTimeout(runtime: string): string { - if (hasRuntimeSqliteBusyTimeout(runtime)) return runtime; - const sync = /try\{([A-Za-z_$][\w$]*)!==([A-Za-z_$][\w$]*)&&([A-Za-z_$][\w$]*)\(this\.db,this\.dbPath,([A-Za-z_$][\w$]*)\)\}catch/gu; - const startup = /try\{return await ([A-Za-z_$][\w$]*)\(([A-Za-z_$][\w$]*)\.db,\2\.dbPath,([A-Za-z_$][\w$]*)\),\2\}catch/gu; - if (countRegExpMatches(runtime, sync) !== 1 || countRegExpMatches(runtime, startup) !== 1) { - throw new Error("ZCode runtime is incompatible with the SQLite busy-timeout patch (store migration anchors missing or ambiguous)."); - } - const patched = runtime.replace( - sync, - (_match, mode: string, deferred: string, migrate: string, timeout: string) => `try{${mode}!==${deferred}&&(${migrate}(this.db,this.dbPath,${timeout}),this.db.exec("${runtimeSqliteBusyTimeoutPragma}"))}catch` - ).replace( - startup, - (_match, migrate: string, store: string, options: string) => `try{return await ${migrate}(${store}.db,${store}.dbPath,${options}),${store}.db.exec("${runtimeSqliteBusyTimeoutPragma}"),${store}}catch` - ); - if (!hasRuntimeSqliteBusyTimeout(patched)) throw new Error("SQLite busy-timeout patch failed postcondition verification."); - return patched; -} - function escapeRegExpName(value: string): string { return value.replace(/[.*+?^${}()|[\]\\]/gu, "\\$&"); } @@ -1629,6 +1599,12 @@ export const runtimePatchPlan: readonly RuntimePatchDefinition[] = [ apply: patchRuntimeSqliteBusyTimeout, verify: hasRuntimeSqliteBusyTimeout }, + { + id: "sqlite-write-recovery", + requirement: "required", + apply: patchRuntimeSqliteWriteRecovery, + verify: hasRuntimeSqliteWriteRecovery + }, { id: "oauth-http-errors", requirement: "optional", diff --git a/src/runtime-config-bridge.ts b/src/runtime-config-bridge.ts index c60d63a..5e786bf 100644 --- a/src/runtime-config-bridge.ts +++ b/src/runtime-config-bridge.ts @@ -6,6 +6,7 @@ import { cliSettingsPath, legacyCliConfigPath, providerConfigPath, providerMigra export { assertSessionModelReady, readSessionModelState } from "./session-model-recovery.ts"; export { readTuiRuntimeProjection, sendTuiBackgroundTaskMessage } from "./runtime-tui-bridge.ts"; export { restoreTuiBackgroundTasks } from "./runtime-background-restore.ts"; +export { installSqliteWriteRecovery, pruneSqliteUsage, sqliteRecoveryStats } from "./runtime-sqlite-recovery.ts"; function record(value: unknown): Record | undefined { return value && typeof value === "object" && !Array.isArray(value) ? value as Record : undefined; diff --git a/src/runtime-sqlite-connection.ts b/src/runtime-sqlite-connection.ts new file mode 100644 index 0000000..d894f16 --- /dev/null +++ b/src/runtime-sqlite-connection.ts @@ -0,0 +1,74 @@ +import { AsyncLocalStorage } from "node:async_hooks"; + +export interface RuntimeSqliteStatement { + get(...values: unknown[]): unknown; +} + +export interface RuntimeSqliteConnection { + readonly isTransaction: boolean; + exec(sql: string): unknown; + prepare(sql: string): RuntimeSqliteStatement; +} + +export interface SqliteOperationScope { + db: RuntimeSqliteConnection; + active: boolean; + nativeBusyMs?: number; + deadline: number; +} + +export const sqliteOperationScope = new AsyncLocalStorage(); + +export function isSqliteBusyError(error: unknown): error is Error & { errcode: number } { + if (!error || typeof error !== "object" || !("errcode" in error)) return false; + const code = error.errcode; + return typeof code === "number" && Number.isInteger(code) && (code & 0xff) === 5; +} + +const guardedConnections = new WeakSet(); +const statementExecutors = new Set(["run", "get", "all"]); + +/** Apply short waits only inside an audited operation, without changing other clients of the connection. */ +export function installSqliteConnectionWait(db: RuntimeSqliteConnection): void { + if (guardedConnections.has(db)) return; + const execute = db.exec.bind(db); + const prepare = db.prepare.bind(db); + const readTimeout = prepare("pragma busy_timeout"); + + const call = (operation: () => T): T => { + const scope = sqliteOperationScope.getStore(); + if (!scope?.active || scope.db !== db || scope.nativeBusyMs === undefined) return operation(); + const row = readTimeout.get() as { timeout: number }; + const previous = row.timeout; + const timeout = Math.min(previous, scope.nativeBusyMs, Math.max(0, Math.floor(scope.deadline - performance.now()))); + if (timeout === previous) return operation(); + execute(`pragma busy_timeout = ${timeout}`); + let failed = false; + try { + return operation(); + } catch (error) { + failed = true; + throw error; + } finally { + // #162:不能把短等待留在共享连接上,也不能让恢复 PRAGMA 的错误覆盖原始写入错误。 + try { execute(`pragma busy_timeout = ${previous}`); } + catch (error) { if (!failed) throw error; } + } + }; + + db.exec = sql => call(() => execute(sql)); + db.prepare = sql => { + const statement = call(() => prepare(sql)); + return new Proxy(statement, { + get(target, key) { + const value = Reflect.get(target, key, target); + if (typeof value !== "function") return value; + // Native StatementSync methods require the original receiver, not the proxy. + return statementExecutors.has(key) + ? (...args: unknown[]) => call(() => Reflect.apply(value, target, args)) + : value.bind(target); + } + }); + }; + guardedConnections.add(db); +} diff --git a/src/runtime-sqlite-recovery.ts b/src/runtime-sqlite-recovery.ts new file mode 100644 index 0000000..e50f77a --- /dev/null +++ b/src/runtime-sqlite-recovery.ts @@ -0,0 +1,280 @@ +import { realpathSync } from "node:fs"; +import { resolve } from "node:path"; +import { + installSqliteConnectionWait, + isSqliteBusyError, + sqliteOperationScope, + type RuntimeSqliteConnection, + type SqliteOperationScope +} from "./runtime-sqlite-connection.ts"; + +export { pruneSqliteUsage } from "./runtime-sqlite-usage.ts"; +export { isSqliteBusyError } from "./runtime-sqlite-connection.ts"; + +export const sqliteWriteRecoveryPolicy = { + nativeBusyMs: 25, + writeBudgetMs: 30_000, + usageBudgetMs: 30_000, + initialDelayMs: 10, + maximumDelayMs: 200 +} as const; + +// Audited against ZCode 29628c9: native transactions roll back before throwing; +// nontransactional methods upsert/delete stable IDs and preserve input/message sequences. +export const sqliteReplayableWrites = [ + "createSession", "updateSession", "createForkedSessionWithMetadata", "commitForkBundle", + "commitSharedContextImportBundle", "transitionSharedContextImport", "saveMessage", + "removeMessage", "savePart", "removePart", "saveSessionEntry", "saveSessionInput", + "updateSessionInputs", "promoteSessionInput", "markSessionInputPromoted", "settleSessionInput", + "commitPermissionFullAccess", "updateTodos", "recordInputHistory", "saveProjectPermission", + "setRevert", "clearRevert" +] as const; + +const usageWrites = ["recordModelUsage", "upsertTurnUsage", "upsertToolUsage"] as const; +// These can increment counters or create identities before a later write fails. +// Order them with other writes, but never replay their possibly committed effects. +const orderedWrites = [ + "setTarget", "cloneTargetForFork", "createTarget", "updateTargetStatus", "startTargetRun", + "heartbeatTargetRun", "finishTargetRun", "recoverInterruptedTargetRun", "accountTargetUsage", + "updateTargetSummaryTitle", "clearTarget", "claimLegacySessionWorkspace", + "repairLegacyRemoteSessionWorkspace", "repairRemoteSessionPaths", "upsertScriptWorkflowDefinition", + "createScriptWorkflowRun", "updateScriptWorkflowRun", "createScriptWorkflowActivity", + "updateScriptWorkflowActivity", "appendScriptWorkflowEvent", "createSessionTaskLink", "pruneUsage" +] as const; + +interface RuntimeSqliteStore { + db: RuntimeSqliteConnection; + getDatabasePath(): string; + close(): void; +} + +type StoreMethod = (...args: unknown[]) => unknown; +type Priority = "conversation" | "usage"; + +interface WriteJob { + owner: InstalledRecovery; + priority: Priority; + run(): Promise; + resolve(value: unknown): void; + reject(error: unknown): void; +} + +interface WriteQueue { + jobs: WriteJob[]; + active: boolean; + users: number; +} + +export interface SqliteRecoveryStats { + recoveredWrites: number; + retries: number; + exhaustedWrites: number; + lastRecovery?: { operation: string; attempts: number; elapsedMs: number }; +} + +interface InstalledRecovery { + controller: AbortController; + stats: SqliteRecoveryStats; +} + +export interface SqliteRecoveryOptions { + /** Internal test seam; no new user configuration or environment variables. */ + writeBudgetMs?: number; + usageBudgetMs?: number; + nativeBusyMs?: number; +} + +const queues = new Map(); +const installed = new WeakMap(); +const activeConnections = new WeakMap(); + +export class SqliteWriteRecoveryError extends Error { + readonly code = "ERR_SQLITE_ERROR"; + readonly errcode: number; + + constructor( + cause: Error & { errcode: number }, + readonly operation: string, + readonly attempts: number, + readonly elapsedMs: number + ) { + super(`${cause.message}; SQLite ${operation} exhausted recovery after ${attempts} attempts (${elapsedMs} ms)`, { cause }); + this.name = "SqliteWriteRecoveryError"; + this.errcode = cause.errcode; + } +} + +export function sqliteRecoveryStats(store: RuntimeSqliteStore): SqliteRecoveryStats | undefined { + const stats = installed.get(store)?.stats; + return stats ? { ...stats, ...(stats.lastRecovery ? { lastRecovery: { ...stats.lastRecovery } } : {}) } : undefined; +} + +function queueKey(store: RuntimeSqliteStore): string | RuntimeSqliteConnection { + const path = store.getDatabasePath(); + if (path === ":memory:") return store.db; + try { return realpathSync(path); } + catch { return resolve(path); } +} + +function delay(ms: number, signal: AbortSignal): Promise { + if (signal.aborted) return Promise.reject(signal.reason); + return new Promise((resolve, reject) => { + const finish = () => { signal.removeEventListener("abort", abort); resolve(); }; + const timer = setTimeout(finish, ms); + const abort = () => { clearTimeout(timer); signal.removeEventListener("abort", abort); reject(signal.reason); }; + signal.addEventListener("abort", abort, { once: true }); + }); +} + +async function drain(queue: WriteQueue, key: string | RuntimeSqliteConnection): Promise { + queue.active = true; + try { + while (queue.jobs.length) { + const preferred = queue.jobs.findIndex(job => job.priority === "conversation"); + const [job] = queue.jobs.splice(preferred < 0 ? 0 : preferred, 1); + try { job!.resolve(await job!.run()); } + catch (error) { job!.reject(error); } + } + } finally { + queue.active = false; + if (queue.users === 0) queues.delete(key); + } +} + +async function serviceConversationWrites(queue: WriteQueue): Promise { + let index: number; + while ((index = queue.jobs.findIndex(job => job.priority === "conversation")) >= 0) { + const [job] = queue.jobs.splice(index, 1); + try { job!.resolve(await job!.run()); } + catch (error) { job!.reject(error); } + } +} + +/** Install once, after migration; do not replace SQL or the native transaction owner. */ +export function installSqliteWriteRecovery(store: RuntimeSqliteStore, options: SqliteRecoveryOptions = {}): void { + if (installed.has(store)) return; + const methods = store as unknown as Record; + const specifications = [ + ...sqliteReplayableWrites.map(name => ({ name, retry: true, priority: "conversation" as const })), + ...usageWrites.map(name => ({ name, retry: true, priority: "usage" as const })), + ...orderedWrites.map(name => ({ name, retry: false, priority: name === "pruneUsage" ? "usage" as const : "conversation" as const })) + ]; + const originals = specifications.map(spec => { + const method = methods[spec.name]; + if (typeof method !== "function") throw new Error(`SQLite recovery requires native ${spec.name}.`); + return { ...spec, method: method as StoreMethod }; + }); + const key = queueKey(store); + const queue = queues.get(key) ?? { jobs: [], active: false, users: 0 }; + const recovery: InstalledRecovery = { + controller: new AbortController(), + stats: { recoveredWrites: 0, retries: 0, exhaustedWrites: 0 } + }; + installSqliteConnectionWait(store.db); + queue.users++; + queues.set(key, queue); + installed.set(store, recovery); + + const close = store.close.bind(store); + store.close = () => { + if (!recovery.controller.signal.aborted) { + const reason = new Error("SQLite session store closed while waiting for a write."); + recovery.controller.abort(reason); + const cancelled = queue.jobs.filter(job => job.owner === recovery); + queue.jobs = queue.jobs.filter(job => job.owner !== recovery); + for (const job of cancelled) job.reject(reason); + queue.users--; + if (!queue.active && queue.users === 0) queues.delete(key); + } + close(); + }; + + for (const { name, method, retry, priority } of originals) { + methods[name] = (...args: unknown[]) => { + const signal = recovery.controller.signal; + if (signal.aborted) return Promise.reject(signal.reason); + const scope = sqliteOperationScope.getStore(); + // Calls within an owned native operation belong to its transaction/retry boundary. + if (scope?.active && scope.db === store.db) return method.apply(store, args); + // A transaction opened outside the bridge cannot safely be replayed or queued inside itself. + if (store.db.isTransaction && !activeConnections.has(store.db)) return method.apply(store, args); + const startedAt = performance.now(); + const budget = priority === "usage" + ? options.usageBudgetMs ?? sqliteWriteRecoveryPolicy.usageBudgetMs + : options.writeBudgetMs ?? sqliteWriteRecoveryPolicy.writeBudgetMs; + return new Promise((resolve, reject) => { + queue.jobs.push({ owner: recovery, priority, resolve, reject, run: () => runWrite({ + store, recovery, operation: name, retry, startedAt, deadline: startedAt + budget, + nativeBusyMs: options.nativeBusyMs ?? sqliteWriteRecoveryPolicy.nativeBusyMs, + ...(priority === "usage" ? { serviceConversation: () => serviceConversationWrites(queue) } : {}), + call: () => method.apply(store, args) + }) }); + if (!queue.active) void drain(queue, key); + }); + }; + } +} + +async function runWrite(input: { + store: RuntimeSqliteStore; + recovery: InstalledRecovery; + operation: string; + retry: boolean; + startedAt: number; + deadline: number; + nativeBusyMs: number; + serviceConversation?(): Promise; + call(): unknown; +}): Promise { + const { store, recovery, operation, startedAt, deadline } = input; + const signal = recovery.controller.signal; + let attempts = 0; + let backoff: number = sqliteWriteRecoveryPolicy.initialDelayMs; + let lastBusy: (Error & { errcode: number }) | undefined; + while (true) { + if (signal.aborted) throw signal.reason; + if (performance.now() >= deadline) { + recovery.stats.exhaustedWrites++; + const elapsed = Math.round(performance.now() - startedAt); + if (lastBusy) throw new SqliteWriteRecoveryError(lastBusy, operation, attempts, elapsed); + throw new Error(`SQLite ${operation} waited ${elapsed} ms for preceding writes; no write was attempted.`); + } + // A caller can open a transaction after this job was queued. Do not silently + // join it and acknowledge data that its later rollback would remove. + if (store.db.isTransaction) { + await delay(Math.min(backoff, deadline - performance.now()), signal); + backoff = Math.min(sqliteWriteRecoveryPolicy.maximumDelayMs, backoff * 2); + await input.serviceConversation?.(); + continue; + } + attempts++; + const scope: SqliteOperationScope = { + db: store.db, active: true, deadline, + ...(input.retry ? { nativeBusyMs: input.nativeBusyMs } : {}) + }; + activeConnections.set(store.db, scope); + try { + const result = await sqliteOperationScope.run(scope, input.call); + if (attempts > 1) { + recovery.stats.recoveredWrites++; + recovery.stats.lastRecovery = { operation, attempts, elapsedMs: Math.round(performance.now() - startedAt) }; + } + return result; + } catch (error) { + // #162:只有已回滚/幂等的存储动作可重试;仍持事务或其他错误必须原样上抛。 + if (!input.retry || !isSqliteBusyError(error) || store.db.isTransaction) throw error; + lastBusy = error; + } finally { + scope.active = false; + activeConnections.delete(store.db); + } + recovery.stats.retries++; + const remaining = deadline - performance.now(); + if (remaining <= 0) continue; + await delay(Math.min(remaining, backoff * (0.5 + Math.random() * 0.5)), signal); + // A busy usage fact releases priority between attempts, after leaving its transaction. + // Later usage writes remain queued so stale facts cannot overtake newer updates. + await input.serviceConversation?.(); + backoff = Math.min(sqliteWriteRecoveryPolicy.maximumDelayMs, backoff * 2); + } +} diff --git a/src/runtime-sqlite-usage.ts b/src/runtime-sqlite-usage.ts new file mode 100644 index 0000000..5665e58 --- /dev/null +++ b/src/runtime-sqlite-usage.ts @@ -0,0 +1,39 @@ +import { isSqliteBusyError, type RuntimeSqliteConnection } from "./runtime-sqlite-connection.ts"; + +export const sqliteUsageCleanupIntervalMs = 5 * 60_000; +export const sqliteUsageCleanupRetryMs = 1_000; + +const nextCleanup = new WeakMap(); +const expiredUsageQuery = `select + exists(select 1 from model_usage where started_at < ?) or + exists(select 1 from turn_usage where started_at < ?) or + exists(select 1 from tool_usage where started_at < ?) as expired`; + +/** The native function still supplies the cutoff and owns the cleanup transaction. */ +export function pruneSqliteUsage( + db: RuntimeSqliteConnection, + options: { beforeTime?: number }, + beforeTime: number, + prune: () => void +): void { + if (options.beforeTime !== undefined) { + prune(); + return; + } + const now = performance.now(); + if (now < (nextCleanup.get(db) ?? 0)) return; + // A caller-owned transaction must not gain a nested maintenance transaction. + if (db.isTransaction) { + nextCleanup.set(db, now + sqliteUsageCleanupRetryMs); + return; + } + try { + const row = db.prepare(expiredUsageQuery).get(beforeTime, beforeTime, beforeTime) as { expired: number }; + if (row.expired) prune(); + nextCleanup.set(db, performance.now() + sqliteUsageCleanupIntervalMs); + } catch (error) { + // #162:用量 UPSERT 已提交;清理抢锁失败不能重放这个事实或拖住主回合。 + if (!isSqliteBusyError(error) || db.isTransaction) throw error; + nextCleanup.set(db, performance.now() + sqliteUsageCleanupRetryMs); + } +} diff --git a/test/fixtures/sqlite-contention.cjs b/test/fixtures/sqlite-contention.cjs new file mode 100644 index 0000000..62d9099 --- /dev/null +++ b/test/fixtures/sqlite-contention.cjs @@ -0,0 +1,41 @@ +const assert = require("node:assert/strict"); +const { fork } = require("node:child_process"); +const { mkdtemp, rm } = require("node:fs/promises"); +const { tmpdir } = require("node:os"); +const path = require("node:path"); + +async function withDirectory(run) { + const directory = await mkdtemp(path.join(tmpdir(), "zcode-sqlite-runtime-")); + try { await run(directory, path.join(directory, "sessions.sqlite")); } + finally { await rm(directory, { recursive: true, force: true }); } +} + +async function startWorker(mode, options) { + const child = fork(path.join(__dirname, "sqlite-session-store.cjs"), [mode, JSON.stringify(options)], { + execArgv: [], stdio: ["ignore", "ignore", "pipe", "ipc"], timeout: 45_000 + }); + let stderr = ""; + child.stderr.on("data", chunk => { stderr += chunk; }); + const exited = new Promise(resolve => child.once("close", (code, signal) => resolve({ code, signal }))); + try { + await new Promise((resolve, reject) => { + child.once("error", reject); + child.once("message", message => message?.type === "ready" ? resolve() : reject(new Error("Unexpected worker message"))); + child.once("close", () => reject(new Error(`Worker exited before ready: ${stderr}`))); + }); + } catch (error) { + child.kill(); + await exited; + throw error; + } + return { + send(message) { if (child.connected) child.send(message); }, + async wait() { + const result = await exited; + assert.equal(result.code, 0, `Worker failed (${result.signal}): ${stderr}`); + }, + async stop(signal = "SIGTERM") { if (child.connected) child.kill(signal); await exited; } + }; +} + +module.exports = { withDirectory, startWorker }; diff --git a/test/fixtures/sqlite-session-store.cjs b/test/fixtures/sqlite-session-store.cjs index 95584fe..5305266 100644 --- a/test/fixtures/sqlite-session-store.cjs +++ b/test/fixtures/sqlite-session-store.cjs @@ -4,7 +4,7 @@ const fs = require("node:fs"); const Module = require("node:module"); const path = require("node:path"); -function loadSessionStore() { +function loadSessionStore(options = {}) { assert.equal(process.versions.bun, undefined, "SQLite runtime tests must execute in real Node.js"); const file = path.resolve(process.env.ZCODE_TEST_RUNTIME || path.join(__dirname, "../../vendor/zcode.cjs")); let source = fs.readFileSync(file, "utf8"); @@ -17,6 +17,16 @@ function loadSessionStore() { const runtime = new Module(file, module); runtime.filename = file; runtime.paths = Module._nodeModulePaths(path.dirname(file)); + if (options.recoveryOptions) { + const nativeRequire = runtime.require.bind(runtime); + const helperPath = path.join(path.dirname(file), "cli-config.cjs"); + const helper = nativeRequire(helperPath); + const testHelper = { + ...helper, + installSqliteWriteRecovery: store => helper.installSqliteWriteRecovery(store, options.recoveryOptions) + }; + runtime.require = id => id === helperPath ? testHelper : nativeRequire(id); + } runtime._compile(source, file); assert.equal(typeof runtime.exports.openStartup, "function"); return runtime.exports; @@ -52,7 +62,7 @@ async function worker(mode, options) { process.send({ type: "ready" }); return; } - assert.equal(mode, "startup"); + assert.ok(mode === "startup" || mode === "stress"); await new Promise(resolve => { process.once("message", resolve); process.send({ type: "ready" }); @@ -62,6 +72,18 @@ async function worker(mode, options) { try { assert.equal(store.db.prepare("PRAGMA busy_timeout").get().timeout, 10_000); await store.createSession(sessionInput(path.dirname(options.dbPath), options.id)); + if (mode === "stress") { + for (let index = 0; index < options.turns; index++) { + const id = `${options.id}-${index}`; + const now = Date.now(); + await store.saveSessionInput({ id, sessionID: options.id, kind: "prompt", delivery: "queue", payload: { text: id } }); + await store.promoteSessionInput({ id, sessionID: options.id, + message: { id: `message-${id}`, sessionID: options.id, role: "user", time: { created: now } }, + parts: [{ id: `part-${id}`, messageID: `message-${id}`, sessionID: options.id, type: "text", text: id + "x".repeat(4096) }] + }); + await store.upsertToolUsage({ id: `usage-${id}`, sessionID: options.id, toolCallID: `tool-${id}`, toolName: "Read", status: "completed", approvalStatus: "none", startedAt: now }); + } + } } finally { store.close(); } diff --git a/test/node/sqlite-recovery.test.cjs b/test/node/sqlite-recovery.test.cjs new file mode 100644 index 0000000..5ba9047 --- /dev/null +++ b/test/node/sqlite-recovery.test.cjs @@ -0,0 +1,360 @@ +const assert = require("node:assert/strict"); +const path = require("node:path"); +const { test } = require("node:test"); +const { loadSessionStore, sessionInput } = require("../fixtures/sqlite-session-store.cjs"); +const { withDirectory, startWorker } = require("../fixtures/sqlite-contention.cjs"); + +assert.equal(process.versions.bun, undefined); +const Store = loadSessionStore(); +const runtimeFile = path.resolve(process.env.ZCODE_TEST_RUNTIME || path.join(__dirname, "../../vendor/zcode.cjs")); +const { sqliteRecoveryStats } = require(path.join(path.dirname(runtimeFile), "cli-config.cjs")); + +function input(sessionID, id) { + return { id, sessionID, kind: "prompt", delivery: "queue", payload: { text: id } }; +} + +function promotion(sessionID, id, text = id) { + return { + id, sessionID, + message: { id: `message-${id}`, sessionID, role: "user", time: { created: Date.now() } }, + parts: [{ id: `part-${id}`, sessionID, messageID: `message-${id}`, type: "text", text }] + }; +} + +function usage(id, startedAt = Date.now()) { + return { id, sessionID: "session", toolCallID: id, toolName: "Read", status: "completed", approvalStatus: "none", startedAt }; +} + +function interceptRuns(db, beforeRun) { + const prepare = db.prepare.bind(db); + db.prepare = sql => { + const statement = prepare(sql); + const run = statement.run.bind(statement); + statement.run = (...args) => { beforeRun(sql.replace(/\s+/gu, " ").trim().toLowerCase(), args); return run(...args); }; + return statement; + }; + return () => { db.prepare = prepare; }; +} + +test("lock recovery yields to timers and restores the native timeout", { timeout: 10_000 }, async () => { + await withDirectory(async (directory, dbPath) => { + const store = await Store.openStartup({ dbPath }); + let holder, release, heartbeat; + try { + holder = await startWorker("hold", { dbPath }); + let ticks = 0; + heartbeat = setInterval(() => { ticks++; }, 10); + release = setTimeout(() => holder.send("release"), 180); + await store.createSession(sessionInput(directory, "responsive")); + assert.ok(ticks >= 3, `Event loop blocked during contention: ${ticks} ticks`); + assert.equal(store.db.prepare("pragma busy_timeout").get().timeout, 10_000); + assert.ok(sqliteRecoveryStats(store).retries > 0); + assert.equal(sqliteRecoveryStats(store).lastRecovery.operation, "createSession"); + assert.equal(store.db.prepare("select count(*) AS n from session").get().n, 1); + } finally { + clearTimeout(release); clearInterval(heartbeat); + await holder?.stop(); store.close(); + } + }); +}); + +test("a write recovers after contention exceeds the former ten-second timeout", { timeout: 25_000 }, async () => { + await withDirectory(async (directory, dbPath) => { + const store = await Store.openStartup({ dbPath }); + let holder; + try { + holder = await startWorker("hold", { dbPath, holdMs: 10_500 }); + holder.send("wait"); + const startedAt = performance.now(); + await store.createSession(sessionInput(directory, "long-contention")); + assert.ok(performance.now() - startedAt >= 10_000); + assert.equal(store.db.prepare("select count(*) AS n from session").get().n, 1); + assert.ok(sqliteRecoveryStats(store).lastRecovery.attempts > 1); + await holder.wait(); + } finally { await holder?.stop(); store.close(); } + }); +}); + +test("connections in one process serialize an async native transaction without blocking its commit", { timeout: 10_000 }, async () => { + await withDirectory(async (directory, dbPath) => { + const first = await Store.openStartup({ dbPath }); + const second = await Store.openStartup({ dbPath }); + try { + const id = "imported", now = Date.now(); + const importing = first.commitSharedContextImportBundle({ + session: sessionInput(directory, id), + contextMessage: { info: { id: "context", sessionID: id, role: "user", visibility: "model-only", source: "shared_context", time: { created: now } }, parts: [] }, + provenance: { id: `provenance:${id}`, sessionID: id, type: "v4/shared_context_import", time: { created: now, updated: now }, data: { contextId: id, status: "pending" } } + }); + assert.equal(first.db.isTransaction, true); + await Promise.all([importing, second.createSession(sessionInput(directory, "concurrent"))]); + assert.equal(first.db.isTransaction, false); + assert.equal(first.db.prepare("select count(*) AS n from session").get().n, 2); + assert.equal(sqliteRecoveryStats(second).retries, 0); + } finally { first.close(); second.close(); } + }); +}); + +test("promotion rolls back a partial write before retrying and commits each message once", async () => { + await withDirectory(async (directory, dbPath) => { + const store = await Store.openStartup({ dbPath }); + try { + await store.createSession(sessionInput(directory, "session")); + await store.saveSessionInput(input("session", "prompt")); + let failed = false, sawPartialWrite = false, rollbacks = 0; + const execute = store.db.exec.bind(store.db); + store.db.exec = sql => { + const result = execute(sql); + if (sql === "rollback") { + rollbacks++; + assert.equal(store.db.prepare("select count(*) AS n from message").get().n, 0); + assert.equal(store.db.prepare("select count(*) AS n from part").get().n, 0); + } + return result; + }; + interceptRuns(store.db, sql => { + if (!failed && sql.startsWith("update session_input")) { + failed = true; + sawPartialWrite = store.db.prepare("select count(*) AS n from part").get().n === 1; + throw Object.assign(new Error("injected SQLITE_BUSY"), { errcode: 5, code: "ERR_SQLITE_ERROR" }); + } + }); + await store.promoteSessionInput(promotion("session", "prompt")); + assert.equal(sawPartialWrite, true); + assert.equal(rollbacks, 1); + assert.equal(store.db.isTransaction, false); + assert.equal(store.db.prepare("select count(*) AS n from message").get().n, 1); + assert.equal(store.db.prepare("select count(*) AS n from part").get().n, 1); + assert.equal(store.db.prepare("select status from session_input").get().status, "promoted"); + assert.equal(sqliteRecoveryStats(store).lastRecovery.attempts, 2); + } finally { store.close(); } + }); +}); + +test("idempotent part recovery preserves sequence and cannot overwrite a newer queued update", async () => { + await withDirectory(async (directory, dbPath) => { + const store = await Store.openStartup({ dbPath }); + try { + await store.createSession(sessionInput(directory, "session")); + const data = promotion("session", "prompt"); + await store.saveMessage(data.message); + let failed = false; + interceptRuns(store.db, sql => { + if (!failed && sql.startsWith("update session set time_updated")) { + failed = true; + throw Object.assign(new Error("busy after part upsert"), { errcode: 5 }); + } + }); + await Promise.all([ + store.savePart({ ...data.parts[0], text: "older" }), + store.savePart({ ...data.parts[0], text: "newer" }) + ]); + const rows = store.db.prepare("select sequence, data from part").all(); + assert.equal(rows.length, 1); + assert.equal(rows[0].sequence, 0); + assert.equal(JSON.parse(rows[0].data).text, "newer"); + assert.equal(sqliteRecoveryStats(store).lastRecovery.operation, "savePart"); + } finally { store.close(); } + }); +}); + +test("non-BUSY failures and partially committed goal counters are not replayed", async () => { + await withDirectory(async (directory, dbPath) => { + const store = await Store.openStartup({ dbPath }); + try { + await store.createSession(sessionInput(directory, "session")); + const fault = Object.assign(new Error("injected constraint failure"), { errcode: 787 }); + const restore = interceptRuns(store.db, sql => { if (sql.startsWith("insert into message")) throw fault; }); + await assert.rejects(store.saveMessage(promotion("session", "prompt").message), error => error === fault); + assert.equal(sqliteRecoveryStats(store).retries, 0); + restore(); + const goal = await store.setTarget({ sessionID: "session", objective: "Test accounting", status: "active" }); + const busy = Object.assign(new Error("busy after committed accounting"), { errcode: 5 }); + interceptRuns(store.db, sql => { if (sql.startsWith("update session set time_updated")) throw busy; }); + await assert.rejects(store.accountTargetUsage({ sessionID: "session", targetID: goal.targetID, tokensUsedDelta: 7 }), error => error === busy); + assert.equal(store.db.prepare("select tokens_used from session_target").get().tokens_used, 7); + assert.equal(sqliteRecoveryStats(store).retries, 0); + } finally { store.close(); } + }); +}); + +test("caller-owned transactions retain rollback ownership and are never replayed", async () => { + await withDirectory(async (directory, dbPath) => { + const store = await Store.openStartup({ dbPath }); + try { + await store.createSession(sessionInput(directory, "session")); + const busy = Object.assign(new Error("busy inside caller transaction"), { errcode: 5 }); + const restore = interceptRuns(store.db, sql => { if (sql.startsWith("insert into session_input")) throw busy; }); + store.db.exec("begin immediate"); + await assert.rejects(store.saveSessionInput(input("session", "prompt")), error => error === busy); + assert.equal(store.db.isTransaction, true); + assert.equal(sqliteRecoveryStats(store).retries, 0); + store.db.exec("rollback"); + restore(); + await store.saveSessionInput(input("session", "after-rollback")); + assert.equal(store.db.prepare("select count(*) AS n from session_input").get().n, 1); + } finally { if (store.db.isTransaction) store.db.exec("rollback"); store.close(); } + }); +}); + +test("a caller-owned transaction can finish while another connection waits in the write queue", { timeout: 5_000 }, async () => { + await withDirectory(async (directory, dbPath) => { + const owner = await Store.openStartup({ dbPath }); + const waiting = await Store.openStartup({ dbPath }); + try { + await owner.createSession(sessionInput(directory, "session")); + owner.db.exec("begin immediate"); + const write = waiting.createSession(sessionInput(directory, "waiting")); + await owner.saveSessionInput(input("session", "inside-transaction")); + owner.db.exec("commit"); + await write; + assert.equal(owner.db.prepare("select count(*) AS n from session").get().n, 2); + assert.equal(owner.db.prepare("select count(*) AS n from session_input").get().n, 1); + } finally { if (owner.db.isTransaction) owner.db.exec("rollback"); owner.close(); waiting.close(); } + }); +}); + +test("a queued write does not join a transaction opened after its submission", { timeout: 5_000 }, async () => { + await withDirectory(async (directory, dbPath) => { + const store = await Store.openStartup({ dbPath }); + try { + await store.createSession(sessionInput(directory, "session")); + const preceding = store.upsertToolUsage(usage("preceding")); + const queued = store.saveSessionInput(input("session", "queued-before-transaction")); + store.db.exec("begin immediate"); + const rollback = new Promise(resolve => setTimeout(() => { store.db.exec("rollback"); resolve(); }, 60)); + await Promise.all([preceding, queued, rollback]); + assert.equal(store.db.prepare("select count(*) AS n from session_input").get().n, 1); + } finally { if (store.db.isTransaction) store.db.exec("rollback"); store.close(); } + }); +}); + +test("closing a store cancels backoff and queued writes without persisting them", { timeout: 10_000 }, async () => { + await withDirectory(async (directory, dbPath) => { + const store = await Store.openStartup({ dbPath }); + let holder, closed = false; + try { + holder = await startWorker("hold", { dbPath }); + const writes = Promise.allSettled([ + store.createSession(sessionInput(directory, "first")), + store.createSession(sessionInput(directory, "queued")) + ]); + await new Promise(resolve => setTimeout(resolve, 60)); + assert.ok(sqliteRecoveryStats(store).retries > 0); + store.close(); closed = true; + const results = await writes; + assert.ok(results.every(result => result.status === "rejected" && /closed/u.test(result.reason.message))); + holder.send("release"); await holder.wait(); + const reopened = await Store.openStartup({ dbPath }); + try { assert.equal(reopened.db.prepare("select count(*) AS n from session").get().n, 0); } + finally { reopened.close(); } + } finally { await holder?.stop(); if (!closed) store.close(); } + }); +}); + +test("closing one store immediately rejects its queued writes behind a different recovering store", { timeout: 10_000 }, async () => { + await withDirectory(async (directory, dbPath) => { + const active = await Store.openStartup({ dbPath }); + const queued = await Store.openStartup({ dbPath }); + let holder, closed = false, deadline; + try { + holder = await startWorker("hold", { dbPath }); + const activeWrite = active.createSession(sessionInput(directory, "active")); + const pending = queued.createSession(sessionInput(directory, "cancelled")); + const observed = pending.then(() => "written", error => /closed/u.test(error.message) ? "closed" : "unexpected error"); + queued.close(); closed = true; + assert.equal(await Promise.race([ + observed, + new Promise(resolve => { deadline = setTimeout(() => resolve("still queued"), 150); }) + ]), "closed"); + holder.send("release"); + await activeWrite; + assert.deepEqual(active.db.prepare("select id from session").all().map(row => row.id), ["active"]); + } finally { clearTimeout(deadline); await holder?.stop(); if (!closed) queued.close(); active.close(); } + }); +}); + +test("usage cleanup retains recent data without a cleanup transaction for every fact", async () => { + await withDirectory(async (directory, dbPath) => { + const store = await Store.openStartup({ dbPath }); + try { + await store.createSession(sessionInput(directory, "session")); + // Seed an expired row using the explicit native schema; automatic cleanup is still due. + store.db.prepare("insert into tool_usage (id,session_id,tool_call_id,tool_name,approval_status,status,started_at) values (?,?,?,?,?,?,?)") + .run("expired", "session", "expired", "Read", "none", "completed", 1); + let transactions = 0, deletes = 0; + const execute = store.db.exec.bind(store.db); + store.db.exec = sql => { if (sql === "begin immediate") transactions++; return execute(sql); }; + interceptRuns(store.db, sql => { if (/^delete from (model|turn|tool)_usage/u.test(sql)) deletes++; }); + for (let index = 0; index < 100; index++) await store.upsertToolUsage(usage(`current-${index}`)); + assert.equal(transactions, 1); + assert.equal(deletes, 3); + assert.equal(store.db.prepare("select count(*) AS n from tool_usage").get().n, 100); + await store.pruneUsage({ beforeTime: Date.now() + 1 }); + assert.equal(transactions, 2); + assert.equal(store.db.prepare("select count(*) AS n from tool_usage").get().n, 0); + } finally { store.close(); } + }); +}); + +test("queued conversation writes precede usage maintenance and each class remains ordered", async () => { + await withDirectory(async (directory, dbPath) => { + const store = await Store.openStartup({ dbPath }); + try { + await store.createSession(sessionInput(directory, "session")); + const order = []; + interceptRuns(store.db, (sql, args) => { + if (sql.startsWith("insert into tool_usage") || sql.startsWith("insert into session (")) order.push(args[0]); + }); + await Promise.all([ + store.upsertToolUsage(usage("usage-1")), + store.upsertToolUsage(usage("usage-2")), + store.createSession(sessionInput(directory, "conversation")) + ]); + assert.deepEqual(order, ["usage-1", "conversation", "usage-2"]); + } finally { store.close(); } + }); +}); + +test("a retrying usage write yields to conversation writes without reordering usage facts", async () => { + await withDirectory(async (directory, dbPath) => { + const store = await Store.openStartup({ dbPath }); + try { + await store.createSession(sessionInput(directory, "session")); + const order = []; + interceptRuns(store.db, (sql, args) => { + if (sql.startsWith("insert into tool_usage")) { + if (!order.includes("conversation")) throw Object.assign(new Error("busy usage"), { errcode: 5 }); + order.push(args[0]); + } else if (sql.startsWith("insert into session (")) order.push(args[0]); + }); + await Promise.all([ + store.upsertToolUsage(usage("usage-1")), + store.upsertToolUsage(usage("usage-2")), + store.createSession(sessionInput(directory, "conversation")) + ]); + assert.deepEqual(order, ["conversation", "usage-1", "usage-2"]); + assert.equal(sqliteRecoveryStats(store).lastRecovery.operation, "upsertToolUsage"); + } finally { store.close(); } + }); +}); + +test("twelve independent processes preserve inputs, transcripts and usage under sustained writes", { timeout: 60_000 }, async () => { + await withDirectory(async (_directory, dbPath) => { + const workers = [], processes = 12, turns = 25; + try { + for (let index = 0; index < processes; index++) workers.push(await startWorker("stress", { dbPath, id: `writer-${index}`, turns })); + for (const worker of workers) worker.send("go"); + await Promise.all(workers.map(worker => worker.wait())); + const store = await Store.openStartup({ dbPath }); + try { + for (const table of ["session_input", "message", "part", "tool_usage"]) { + assert.equal(store.db.prepare(`select count(*) AS n from ${table}`).get().n, processes * turns, table); + } + assert.equal(store.db.prepare("select count(*) AS n from session_input where status = 'promoted'").get().n, processes * turns); + assert.deepEqual(store.db.prepare("pragma foreign_key_check").all(), []); + assert.equal(store.db.prepare("pragma integrity_check").get().integrity_check, "ok"); + } finally { store.close(); } + } finally { await Promise.all(workers.map(worker => worker.stop())); } + }); +}); diff --git a/test/node/sqlite-session-store.test.cjs b/test/node/sqlite-session-store.test.cjs index 79ea6fb..40bf53a 100644 --- a/test/node/sqlite-session-store.test.cjs +++ b/test/node/sqlite-session-store.test.cjs @@ -1,47 +1,11 @@ const assert = require("node:assert/strict"); -const { fork } = require("node:child_process"); -const { mkdtemp, rm } = require("node:fs/promises"); -const { tmpdir } = require("node:os"); -const path = require("node:path"); const { test } = require("node:test"); const { loadSessionStore, sessionInput } = require("../fixtures/sqlite-session-store.cjs"); +const { withDirectory, startWorker } = require("../fixtures/sqlite-contention.cjs"); assert.equal(process.versions.bun, undefined, "Run with npm run test:node / node --test, not bun test"); const Store = loadSessionStore(); - -async function withDirectory(run) { - const directory = await mkdtemp(path.join(tmpdir(), "zcode-sqlite-runtime-")); - try { await run(directory, path.join(directory, "sessions.sqlite")); } - finally { await rm(directory, { recursive: true, force: true }); } -} - -async function startWorker(mode, options) { - const child = fork(path.join(__dirname, "../fixtures/sqlite-session-store.cjs"), [mode, JSON.stringify(options)], { - execArgv: [], stdio: ["ignore", "ignore", "pipe", "ipc"], timeout: 20_000 - }); - let stderr = ""; - child.stderr.on("data", chunk => { stderr += chunk; }); - const exited = new Promise(resolve => child.once("close", (code, signal) => resolve({ code, signal }))); - try { - await new Promise((resolve, reject) => { - child.once("error", reject); - child.once("message", message => message?.type === "ready" ? resolve() : reject(new Error("Unexpected worker message"))); - child.once("close", () => reject(new Error(`Worker exited before ready: ${stderr}`))); - }); - } catch (error) { - child.kill(); - await exited; - throw error; - } - return { - send(message) { if (child.connected) child.send(message); }, - async wait() { - const result = await exited; - assert.equal(result.code, 0, `Worker failed (${result.signal}): ${stderr}`); - }, - async stop(signal = "SIGTERM") { if (child.connected) child.kill(signal); await exited; } - }; -} +const BoundedStore = loadSessionStore({ recoveryOptions: { writeBudgetMs: 250 } }); test("real runtime keeps the write timeout after sync open, async startup and reopen", { timeout: 15_000 }, async t => { await withDirectory(async (_directory, dbPath) => { @@ -66,7 +30,7 @@ test("a native session write survives a lock held beyond the old five-second win try { assert.equal(store.db.prepare("PRAGMA busy_timeout").get().timeout, 10_000); holder = await startWorker("hold", { dbPath, holdMs: 6500 }); - // The independent process releases the lock even while DatabaseSync blocks this event loop. + // Keep the held interval independent of the writer's retry scheduling. holder.send("wait"); const startedAt = performance.now(); await store.createSession(sessionInput(directory, "waited-session")); @@ -82,10 +46,10 @@ test("a native session write survives a lock held beyond the old five-second win test("lock timeout is bounded, writes nothing, and the same connection works after release", { timeout: 15_000 }, async () => { await withDirectory(async (directory, dbPath) => { - const store = await Store.openStartup({ dbPath }); + const store = await BoundedStore.openStartup({ dbPath }); let holder; try { - // Use a short test budget; the production budget is asserted separately above. + // Shorten both budgets for this failure test; production recovery is exercised above 10 seconds. store.db.exec("PRAGMA busy_timeout = 250"); holder = await startWorker("hold", { dbPath }); const startedAt = performance.now(); diff --git a/test/runtime-sqlite-patches.test.ts b/test/runtime-sqlite-patches.test.ts new file mode 100644 index 0000000..80d46fc --- /dev/null +++ b/test/runtime-sqlite-patches.test.ts @@ -0,0 +1,76 @@ +import { describe, expect, test } from "bun:test"; +import { + hasRuntimeSqliteBusyTimeout, + hasRuntimeSqliteWriteRecovery, + patchRuntimeSqliteBusyTimeout, + patchRuntimeSqliteWriteRecovery +} from "../scripts/runtime-sqlite-patches.ts"; + +function fixture(): string { + return [ + 'const deferred=Symbol(),retention=30*86400000;class Store{constructor(t={},n){this.dbPath="test.sqlite";let r=5000;this.db=t.db;', + 'try{n!==deferred&&migrateSync(this.db,this.dbPath,r)}catch(e){this.db.close();throw e}}', + 'static async openStartup(t={},n={}){let o=new Store(t,deferred);try{return await migrateAsync(o.db,o.dbPath,n),o}catch(e){try{o.close()}catch{}throw e}}', + 'close(){this.db.close()}}', + 'async function prune(db,options={}){let cutoff=options.beforeTime??Date.now()-retention;db.exec("begin immediate");try{', + 'db.prepare("delete from model_usage where started_at < ?").run(cutoff),', + 'db.prepare("delete from turn_usage where started_at < ?").run(cutoff),', + 'db.prepare("delete from tool_usage where started_at < ?").run(cutoff),db.exec("commit")', + '}catch(error){throw db.exec("rollback"),error}}label(prune,"pruneUsage");' + ].join(""); +} + +describe("SQLite recovery bundle boundaries", () => { + test("installs after both migrations and preserves the native cleanup cutoff and transaction", async () => { + const source = patchRuntimeSqliteBusyTimeout(fixture()); + const patched = patchRuntimeSqliteWriteRecovery(source); + expect(hasRuntimeSqliteBusyTimeout(patched)).toBe(true); + expect(hasRuntimeSqliteWriteRecovery(patched)).toBe(true); + expect(patchRuntimeSqliteWriteRecovery(patched)).toBe(patched); + expect(patchRuntimeSqliteBusyTimeout(patched)).toBe(patched); + const events: unknown[] = []; + const db = { + exec(sql: string) { events.push(sql); }, + prepare(sql: string) { return { run(value: unknown) { events.push([sql, value]); } }; }, + close() { events.push("close"); } + }; + const helper = { + installSqliteWriteRecovery(store: { db: unknown }) { expect(store.db).toBe(db); events.push("install"); }, + pruneSqliteUsage(connection: unknown, options: unknown, cutoff: number, run: () => void) { + expect(connection).toBe(db); + events.push(["cleanup", options, cutoff]); + run(); + } + }; + const run = new Function("require", "__dirname", "migrateSync", "migrateAsync", "label", `${patched};return {Store,prune};`)( + (name: string) => name === "node:path" ? { join: (...parts: string[]) => parts.join("/") } : helper, + "/fixture", () => events.push("sync"), async () => events.push("async"), () => {} + ); + new run.Store({ db }); + expect(events).toEqual(["sync", "pragma busy_timeout = 10000", "install"]); + events.length = 0; + await run.Store.openStartup({ db }); + expect(events).toEqual(["async", "pragma busy_timeout = 10000", "install"]); + events.length = 0; + await run.prune(db, { beforeTime: 42 }); + expect(events).toEqual([ + ["cleanup", { beforeTime: 42 }, 42], "begin immediate", + ["delete from model_usage where started_at < ?", 42], + ["delete from turn_usage where started_at < ?", 42], + ["delete from tool_usage where started_at < ?", 42], "commit" + ]); + }); + + test("rejects partial patches and unknown cleanup semantics", () => { + const source = patchRuntimeSqliteBusyTimeout(fixture()); + expect(hasRuntimeSqliteWriteRecovery(source)).toBe(false); + expect(() => patchRuntimeSqliteWriteRecovery(fixture())).toThrow("post-migration"); + expect(() => patchRuntimeSqliteWriteRecovery(source + source)).toThrow(); + expect(() => patchRuntimeSqliteWriteRecovery(source.replace('"pruneUsage"', '"renamed"'))).toThrow("symbol"); + expect(() => patchRuntimeSqliteWriteRecovery(source.replace("delete from turn_usage", "delete from different_table"))).toThrow("SQL changed"); + const patched = patchRuntimeSqliteWriteRecovery(source); + const partial = patched.replace(".installSqliteWriteRecovery(this)", ".other(this)"); + expect(hasRuntimeSqliteWriteRecovery(partial)).toBe(false); + expect(() => patchRuntimeSqliteWriteRecovery(partial)).toThrow("partial"); + }); +}); diff --git a/test/runtime-sqlite-usage.test.ts b/test/runtime-sqlite-usage.test.ts new file mode 100644 index 0000000..99dd8f7 --- /dev/null +++ b/test/runtime-sqlite-usage.test.ts @@ -0,0 +1,36 @@ +import { expect, test } from "bun:test"; +import { isSqliteBusyError } from "../src/runtime-sqlite-connection.ts"; +import { pruneSqliteUsage } from "../src/runtime-sqlite-usage.ts"; + +test("only the SQLite BUSY family is retryable", () => { + for (const errcode of [5, 261, 517, 773]) expect(isSqliteBusyError({ errcode })).toBe(true); + for (const errcode of [6, 19, 787, 10, "5", undefined, NaN, 5.1]) expect(isSqliteBusyError({ errcode })).toBe(false); + expect(isSqliteBusyError(new Error("database is locked"))).toBe(false); +}); + +test("automatic cleanup skips empty tables, coalesces checks, and preserves explicit requests", () => { + let checks = 0, cleanups = 0; + const db = { isTransaction: false, exec() {}, prepare() { return { get() { checks++; return { expired: 0 }; } }; } }; + const prune = () => { cleanups++; }; + for (let i = 0; i < 100; i++) pruneSqliteUsage(db, {}, 10, prune); + expect(checks).toBe(1); + expect(cleanups).toBe(0); + pruneSqliteUsage(db, { beforeTime: 20 }, 20, prune); + expect(cleanups).toBe(1); +}); + +test("maintenance BUSY is deferred, but transaction and non-BUSY failures remain visible", () => { + const busy = Object.assign(new Error("database is locked"), { errcode: 5 }); + const createDb = () => ({ isTransaction: false, exec() {}, prepare() { return { get() { return { expired: 1 }; } }; } }); + const db = createDb(); + let calls = 0; + const fail = () => { calls++; throw busy; }; + expect(() => pruneSqliteUsage(db, {}, 10, fail)).not.toThrow(); + pruneSqliteUsage(db, {}, 10, fail); + expect(calls).toBe(1); + expect(() => pruneSqliteUsage(db, { beforeTime: 10 }, 10, fail)).toThrow(busy); + const io = Object.assign(new Error("disk I/O error"), { errcode: 10 }); + expect(() => pruneSqliteUsage(createDb(), {}, 10, () => { throw io; })).toThrow(io); + const open = createDb(); + expect(() => pruneSqliteUsage(open, {}, 10, () => { open.isTransaction = true; throw busy; })).toThrow(busy); +});