From 099320af2971aa4c9570d88740efc4b54cdf47a9 Mon Sep 17 00:00:00 2001 From: ScreenTinker Date: Mon, 6 Jul 2026 23:40:11 -0500 Subject: [PATCH] fix(db): WAL checkpointer worker-death handling (respawn + inline-autocheckpoint fallback) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Close the disk-fill trap: with wal_autocheckpoint=0 a dead worker means nothing checkpoints. Controller now respawns an unexpectedly-dead worker (bounded: RespawnMax/RespawnWindowMs + backoff); on exhaustion it re-arms a conservative inline autocheckpoint (FallbackPages) on the main connection + reclaims the backlog, logging loudly. Clean stopWalCheckpointer() teardown is distinguished via a 'stopping' flag so SIGTERM never triggers respawn. Env-gated worker fault-injection (WAL_CKPT_FAIL_START) for tests. Local only — no bump/tag. --- server/config.js | 9 +++ server/db/wal-checkpointer-worker.js | 4 + server/db/wal-checkpointer.js | 112 +++++++++++++++++++++------ 3 files changed, 100 insertions(+), 25 deletions(-) diff --git a/server/config.js b/server/config.js index be1ae69..384f2e9 100644 --- a/server/config.js +++ b/server/config.js @@ -267,6 +267,15 @@ module.exports = { // ...or escalate if the WAL grew across this many consecutive PASSIVE runs (PASSIVE not // keeping up even below the high-water). Belt-and-suspenders with the MB bound above. walCheckpointStarvationRuns: parseInt(process.env.WAL_CHECKPOINT_STARVATION_RUNS) || 3, + // Worker-death handling: with autocheckpoint=0 a dead worker means nothing checkpoints and + // the WAL grows until the disk fills. An unexpectedly-dead worker is respawned up to + // RespawnMax times per RespawnWindowMs (with a small backoff); if that's exhausted we + // re-arm a conservative inline autocheckpoint of FallbackPages on the main connection + // (degraded-but-safe: occasional inline stall beats unbounded WAL growth). + walCheckpointRespawnMax: parseInt(process.env.WAL_CHECKPOINT_RESPAWN_MAX) || 5, + walCheckpointRespawnWindowMs: parseInt(process.env.WAL_CHECKPOINT_RESPAWN_WINDOW_MS) || 60000, + walCheckpointRespawnBackoffMs: parseInt(process.env.WAL_CHECKPOINT_RESPAWN_BACKOFF_MS) || 1000, + walCheckpointFallbackPages: parseInt(process.env.WAL_CHECKPOINT_FALLBACK_PAGES) || 1000, // #146 device_status_log write batching (lib/status-log-writer.js). Status // transitions are buffered and coalesced to the NET state per device per flush, // so a flapping device writes ~1 row/flush instead of a row per transition — diff --git a/server/db/wal-checkpointer-worker.js b/server/db/wal-checkpointer-worker.js index a5aad68..7378884 100644 --- a/server/db/wal-checkpointer-worker.js +++ b/server/db/wal-checkpointer-worker.js @@ -12,6 +12,10 @@ const Database = require('better-sqlite3'); const { dbPath, intervalMs, highWaterBytes, starvationRuns } = workerData; +// Fault injection for TESTS ONLY (env-gated; inert in prod). Exits immediately on start so +// the controller's respawn / autocheckpoint-fallback path can be exercised deterministically. +if (process.env.WAL_CKPT_FAIL_START) process.exit(1); + // Fresh, worker-owned connection (NOT the main handle). const db = new Database(dbPath); db.pragma('busy_timeout = 5000'); // wait (on THIS worker thread) through the main writer's brief locks diff --git a/server/db/wal-checkpointer.js b/server/db/wal-checkpointer.js index 7468ce0..d6a9c60 100644 --- a/server/db/wal-checkpointer.js +++ b/server/db/wal-checkpointer.js @@ -2,49 +2,111 @@ // // SQLite's default auto-checkpoint runs a synchronous, fsync-heavy checkpoint inline on the // write that trips the 1000-page threshold; on slow storage that blocks the event loop for -// ~600-750ms on a ~60s beat (the periodic p99 spike). Here we disable inline auto-checkpoint -// on the MAIN connection and delegate checkpointing to a worker_threads worker that opens its +// ~600-750ms on a ~60s beat (the periodic p99 spike). We disable inline auto-checkpoint on +// the MAIN connection and delegate checkpointing to a worker_threads worker that opens its // OWN connection (see wal-checkpointer-worker.js) so the fsync blocks the worker, not the loop. +// +// FAILURE MODE this file closes: with wal_autocheckpoint=0, if the worker dies NOTHING +// checkpoints and the WAL grows until the disk fills. So an unexpectedly-dead worker is +// respawned (bounded retry); if it can't be kept alive, we re-enable a conservative inline +// autocheckpoint on the main connection as a degraded-but-safe fallback (occasional inline +// stall << unbounded WAL growth). const path = require('path'); const { Worker } = require('worker_threads'); const config = require('../config'); let worker = null; +let mainDb = null; +let mainDbPath = null; +let stopping = false; // true only during our own stopWalCheckpointer() teardown +let fallbackEngaged = false; // true once we've given up on the worker and re-armed inline autocheckpoint +const respawnAt = []; // timestamps (ms) of recent respawns, for the rate window -// Call ONCE at boot, after the DB is open + migrated. `db` is the main connection (used only -// to flip the pragma + do a one-time handoff checkpoint on the main thread at boot). `dbPath` -// is the STRING the worker uses to open its own handle — the main handle is never shared. -function startWalCheckpointer(db, dbPath) { - if (worker) return worker; - - // From now on the main thread NEVER inline-checkpoints (removes the loop-blocking fsync). - db.pragma('wal_autocheckpoint = 0'); - // Hand the worker a clean WAL (one-time, at boot, on a small WAL — cheap). Explicit - // checkpoints are independent of wal_autocheckpoint, so this still works with it at 0. - try { db.pragma('wal_checkpoint(TRUNCATE)'); } catch (_) { /* best-effort */ } - - worker = new Worker(path.join(__dirname, 'wal-checkpointer-worker.js'), { +function spawnWorker() { + const w = new Worker(path.join(__dirname, 'wal-checkpointer-worker.js'), { workerData: { - dbPath, // string only (thread-safe handoff) + dbPath: mainDbPath, // string only (thread-safe handoff) intervalMs: config.walCheckpointIntervalMs, highWaterBytes: config.walCheckpointHighWaterMB * 1024 * 1024, starvationRuns: config.walCheckpointStarvationRuns, }, }); - worker.on('message', (m) => { if (m && m.log) console.log('[wal-checkpoint] ' + m.log); }); - worker.on('error', (e) => console.error('[wal-checkpoint] worker error:', e && e.message)); - worker.on('exit', (code) => { if (code !== 0) console.warn(`[wal-checkpoint] worker exited (code ${code})`); worker = null; }); - // A worker thread cannot outlive its process, but unref() also ensures it never KEEPS the + w.on('message', (m) => { if (m && m.log) console.log('[wal-checkpoint] ' + m.log); }); + w.on('error', (e) => console.error('[wal-checkpoint] worker error:', e && e.message)); // 'exit' handles recovery + w.on('exit', (code) => onWorkerExit(code)); + // A worker thread cannot outlive its process; unref() also ensures it never KEEPS the // process alive during shutdown — so there's no orphaned worker/connection either way. - worker.unref(); + w.unref(); + return w; +} - console.log(`[wal-checkpoint] off-thread checkpointer started (every ${config.walCheckpointIntervalMs}ms; escalate >${config.walCheckpointHighWaterMB}MB or ${config.walCheckpointStarvationRuns} growing runs)`); +// Called on every worker 'exit'. Distinguishes our intentional teardown (stopping=true — +// stay silent, no respawn) from an unexpected death (respawn, then fall back if exhausted). +function onWorkerExit(code) { + worker = null; + if (stopping || fallbackEngaged) return; // clean stop, or we've already given up — no noise + console.warn(`[wal-checkpoint] worker died unexpectedly (code ${code}) — attempting respawn`); + scheduleRespawn(); +} + +function scheduleRespawn() { + if (stopping || fallbackEngaged) return; + const now = Date.now(); + while (respawnAt.length && now - respawnAt[0] > config.walCheckpointRespawnWindowMs) respawnAt.shift(); + if (respawnAt.length >= config.walCheckpointRespawnMax) { + engageFallback(); // too many respawns in the window -> give up + return; + } + respawnAt.push(now); + const t = setTimeout(() => { + if (stopping || fallbackEngaged) return; + try { + worker = spawnWorker(); + console.warn(`[wal-checkpoint] worker respawned (${respawnAt.length}/${config.walCheckpointRespawnMax} in window)`); + } catch (e) { + console.error('[wal-checkpoint] respawn spawn failed: ' + (e && e.message)); + scheduleRespawn(); // count this failure too + } + }, config.walCheckpointRespawnBackoffMs); + if (t.unref) t.unref(); // never let the backoff timer hold the process open +} + +// Degraded-but-safe: re-arm a conservative inline autocheckpoint on the MAIN connection so +// the WAL can never grow unbounded, and reclaim the backlog the dead worker left behind. +function engageFallback() { + if (fallbackEngaged) return; + fallbackEngaged = true; + try { mainDb.pragma(`wal_autocheckpoint = ${config.walCheckpointFallbackPages}`); } catch (_) {} + try { mainDb.pragma('wal_checkpoint(TRUNCATE)'); } catch (_) {} // one-time reclaim of the dead-worker backlog + console.error('[wal-checkpoint] worker unrecoverable — re-enabled inline autocheckpoint as fallback'); +} + +// Call ONCE at boot, after the DB is open + migrated. `db` is the main connection (used to +// flip the pragma, do the one-time handoff checkpoint, and arm the fallback if needed). +// `dbPath` is the STRING the worker uses to open its own handle — the main handle is never shared. +function startWalCheckpointer(db, dbPath) { + if (worker) return worker; + mainDb = db; + mainDbPath = dbPath; + stopping = false; + fallbackEngaged = false; + respawnAt.length = 0; + + // From now on the main thread NEVER inline-checkpoints (removes the loop-blocking fsync). + db.pragma('wal_autocheckpoint = 0'); + // Hand the worker a clean WAL (one-time, at boot; explicit checkpoints are independent of + // wal_autocheckpoint, so this still works at 0). Also reclaims any WAL a prior crash left. + try { db.pragma('wal_checkpoint(TRUNCATE)'); } catch (_) { /* best-effort */ } + + worker = spawnWorker(); + console.log(`[wal-checkpoint] off-thread checkpointer started (every ${config.walCheckpointIntervalMs}ms; escalate >${config.walCheckpointHighWaterMB}MB or ${config.walCheckpointStarvationRuns} growing runs; respawn max ${config.walCheckpointRespawnMax}/${config.walCheckpointRespawnWindowMs}ms)`); return worker; } -// Graceful teardown: ask the worker to stop (clears its timer + closes its connection), then -// force-terminate as a backstop. Safe to call when not started. +// Graceful teardown: mark intentional (so onWorkerExit stays silent), ask the worker to stop +// (clears its timer + closes its connection), then force-terminate as a backstop. Safe when not started. async function stopWalCheckpointer() { + stopping = true; if (!worker) return; const w = worker; worker = null; @@ -53,4 +115,4 @@ async function stopWalCheckpointer() { try { await w.terminate(); } catch (_) {} } -module.exports = { startWalCheckpointer, stopWalCheckpointer }; +module.exports = { startWalCheckpointer, stopWalCheckpointer, _getWorker: () => worker };