diff --git a/lib/reconcile.js b/lib/reconcile.js new file mode 100644 index 00000000..0f4b584f --- /dev/null +++ b/lib/reconcile.js @@ -0,0 +1,289 @@ +/** + * Startup scheduler reconciliation (#455) — v3 port of the v2 module + * (PR #462); redis@5 promise API; local submit_type resolves an empty + * snapshot; no own Redis connection — the caller passes its client. + * + * Historically server.js blindly ran `client.del("active_jobs")` on boot, + * orphaning every job that was still live on the scheduler and leaving their + * (TTL-less) hashes behind forever. Instead, reconcile the active_jobs list + * against a single scheduler snapshot taken at startup: + * + * - entries whose scheduler job is still live are KEPT (deduped, exactly + * once — purging the historical duplicate-push backlog); + * - entries already in a terminal state are dropped and their hashes get + * the appropriate retention TTL (#453); + * - entries whose scheduler job no longer exists are zombies: marked + * aborted, given the terminal TTL, and dropped. + * + * A final SCAN sweep applies the same treatment to queued/running job hashes + * that are not in active_jobs at all (pre-fix restarts deleted the list but + * not the hashes). This is sound because server.js gates socket handler + * registration AND the MCP server on reconciliation, so every candidate hash + * predates this restart. + * + * FAIL OPEN: if the scheduler snapshot cannot be taken (scheduler down, + * command missing), active_jobs is left untouched — never mass-abort jobs on + * a scheduler error. + */ + +const { exec } = require("child_process"); +const util = require("util"); +const execP = util.promisify(exec); +const logger = require("./logger").logger; +const config = require("./config"); +// TTL constants only — the redis-client expire helpers are bound to the +// shared client, and this module must operate on whichever client the caller +// passes (tests use isolated connections/databases). +const { + COMPLETED_TTL_SECONDS, + TERMINAL_TTL_SECONDS, +} = require("./redis-client"); + +const TERMINAL_STATUSES = new Set([ + "completed", + "error", + "aborted", + "cancelled", +]); +const SNAPSHOT_TIMEOUT_MS = 15000; + +/** + * Snapshot the ids of every job currently known to the scheduler, mirroring + * lib/jobqueue.js's listing pattern. Reads config.submit_type AT CALL TIME. + * Resolves a Set of live scheduler ids; any exec error (including the + * timeout kill) THROWS to the caller, which must fail open. + */ +async function snapshotSchedulerIds() { + if (config.submit_type === "local") { + // A restart severs self.local_process and the watcher; local in-flight + // jobs are unfinalizable, so an empty snapshot (zombie-marking) is + // correct — and it avoids exec'ing squeue on local-only boxes. + return new Set(); + } + if (config.submit_type === "qsub") { + const { stdout } = await execP("qstat", { timeout: SNAPSHOT_TIMEOUT_MS }); + const live = new Set(); + const lines = (stdout || "").toString().split("\n"); + // qstat output: two header lines, then one job per line, id first + for (let i = 2; i < lines.length; i++) { + const id = lines[i].split(" ")[0].trim(); + if (id) live.add(id); + } + return live; + } + const { stdout } = await execP("squeue --noheader --format=%i", { + timeout: SNAPSHOT_TIMEOUT_MS, + }); + const live = new Set(); + (stdout || "") + .toString() + .split("\n") + .forEach(function (line) { + const id = line.trim(); + if (id) live.add(id); + }); + return live; +} + +function parseTorqueId(raw) { + if (!raw) return null; + try { + const tid = JSON.parse(raw).torque_id; + return tid === undefined || tid === null ? null : String(tid); + } catch (e) { + return null; + } +} + +// Mirror of redis-client.js's expireKeys, but on the injected client. Skips +// falsy keys, logs on failure, never rejects. +async function expireKeys(client, keys, seconds) { + const targets = keys.filter(function (k) { + return k !== undefined && k !== null && k !== ""; + }); + if (targets.length === 0) return; + await Promise.all( + targets.map(function (k) { + return client.expire(k, seconds); + }) + ).catch(function (err) { + logger.error( + "reconcile : expire failed for [" + + targets.join(", ") + + "]: " + + err.message + ); + }); +} + +// Mark an orphaned in-flight hash aborted and apply the terminal retention +// TTL. The hSet is awaited BEFORE the expire (the v2 module fired both from +// a callback; awaiting removes that fire-and-forget race). +async function markZombie(client, id, torque_id) { + await client + .hSet(id, { + status: "aborted", + error: JSON.stringify({ + type: "script error", + error: "orphaned at server restart", + }), + }) + .catch(function (err) { + logger.error(id + " : reconcile : redis hSet failed: " + err.message); + }); + await expireKeys(client, [id, torque_id], TERMINAL_TTL_SECONDS); +} + +// Sweep job hashes that are NOT in active_jobs (seen) but claim to be +// in-flight against the live scheduler snapshot. Resolves the swept count; +// never rejects (a SCAN failure logs and returns the partial count — the +// sweep must never block boot). +async function sweepZombieHashes(client, live, seen) { + let swept = 0; + try { + // node-redis v5 scanIterator yields ARRAYS of keys per batch, not + // single keys. + for await (const keys of client.scanIterator({ COUNT: 500 })) { + for (const key of keys) { + if (seen.has(key)) continue; + // TYPE guard: result blobs / lists live alongside job hashes, and + // HGETALL on a non-hash raises WRONGTYPE (prod-cleanup lesson). + if ((await client.type(key)) !== "hash") continue; + // redis@5 resolves {} for a missing key — guard on emptiness. + const obj = await client.hGetAll(key); + if (Object.keys(obj).length === 0) continue; + if (obj.status === "queued" || obj.status === "running") { + const tid = parseTorqueId(obj.torque_id); + // A scheduler-side hash carries datamonkey_id and is keyed by the + // scheduler id itself. + const liveId = tid || (obj.datamonkey_id ? key : null); + if (!liveId || !live.has(liveId)) { + await markZombie(client, key, tid); + swept++; + } + } + } + } + } catch (err) { + logger.warn( + "reconcile : SCAN failed, skipping zombie-hash sweep: " + err.message + ); + } + return swept; +} + +/** + * Reconcile the active_jobs list (and orphaned in-flight hashes) against a + * single scheduler snapshot. CONTRACT: NEVER rejects — server.js gates + * socket/MCP registration on this promise, so any failure path must resolve. + * Resolves a summary { entries, unique, kept, reaped, swept } for tests. + */ +async function reconcileActiveJobs(client) { + try { + let live; + try { + live = await snapshotSchedulerIds(); + } catch (err) { + // FAIL OPEN: never mass-abort on a scheduler error. + logger.warn( + "reconcile : could not snapshot scheduler queue, leaving active_jobs untouched: " + + err.message + ); + return; + } + + let entries; + try { + entries = (await client.lRange("active_jobs", 0, -1)) || []; + } catch (err) { + logger.warn("reconcile : could not read active_jobs: " + err.message); + return; + } + + const unique = Array.from(new Set(entries)); + const survivors = []; + let reaped = 0; + + for (const id of unique) { + const obj = await client.hGetAll(id); + if (Object.keys(obj).length === 0) { + // no hash — nothing to keep + reaped++; + } else if (TERMINAL_STATUSES.has(obj.status)) { + // already terminal: drop from the list and apply the retention TTL + await expireKeys( + client, + [id, parseTorqueId(obj.torque_id)], + obj.status === "completed" + ? COMPLETED_TTL_SECONDS + : TERMINAL_TTL_SECONDS + ); + reaped++; + } else { + const tid = parseTorqueId(obj.torque_id); + if (tid && live.has(tid)) { + survivors.push(id); + // Defensive #453 invariant pin: an in-flight survivor must be + // TTL-less. PERSIST is a no-op when no TTL is set; it guards the + // resurrect-same-id trap (a TTL'd terminal hash reused in flight). + await client.persist(id).catch(function (err) { + logger.error( + id + " : reconcile : redis persist failed: " + err.message + ); + }); + if (tid !== id) { + await client.persist(tid).catch(function (err) { + logger.error( + tid + " : reconcile : redis persist failed: " + err.message + ); + }); + } + } else { + // unparseable or dead scheduler id -> zombie + await markZombie(client, id, tid); + reaped++; + } + } + } + + // Atomic rebuild: each survivor exactly once, purging the historical + // duplicate backlog. + try { + const m = client.multi().del("active_jobs"); + survivors.forEach(function (id) { + m.rPush("active_jobs", id); + }); + await m.exec(); + } catch (err) { + logger.error("reconcile : active_jobs rebuild failed: " + err.message); + } + + const swept = await sweepZombieHashes(client, live, new Set(unique)); + + logger.warn( + "reconcile : active_jobs entries=" + + entries.length + + " unique=" + + unique.length + + " kept=" + + survivors.length + + " reaped=" + + reaped + + " zombie hashes swept=" + + swept + ); + + return { + entries: entries.length, + unique: unique.length, + kept: survivors.length, + reaped: reaped, + swept: swept, + }; + } catch (err) { + // Belt-and-braces: the boot gate depends on this promise resolving. + logger.error("reconcile : unexpected error: " + err.message); + } +} + +exports.reconcileActiveJobs = reconcileActiveJobs; diff --git a/package.json b/package.json index f625a59b..df7ab910 100644 --- a/package.json +++ b/package.json @@ -8,11 +8,11 @@ "private": true, "scripts": { "test": "npm run test:ci && npm run test:slurm", - "test:ci": "mocha -R spec --exit test/load-smoke.js test/regression/leaks.js test/regression/terminal-packet-id.js test/error-paths/validation.js test/validation/results.js test/mock/lifecycle.js test/golden/qsub-params.js test/golden/factory-parity.js test/golden/hyphy-param-flow.js test/routes/analysis-routes.js test/sanitize-names.test.js test/difFubar.test.js test/difFubar.integration.test.js test/fubar-redirection.test.js && npm run test:axomeme-unit", + "test:ci": "mocha -R spec --exit test/load-smoke.js test/regression/leaks.js test/regression/terminal-packet-id.js test/error-paths/validation.js test/validation/results.js test/mock/lifecycle.js test/golden/qsub-params.js test/golden/factory-parity.js test/golden/hyphy-param-flow.js test/routes/analysis-routes.js test/sanitize-names.test.js test/difFubar.test.js test/difFubar.integration.test.js test/fubar-redirection.test.js test/regression/boot-reconcile-gate.js test/reconcile-edges.js test/reconcile-sweep.js && npm run test:axomeme-unit", "test:axomeme-unit": "mocha -R spec --exit test/axomeme/model-integrity.test.js test/axomeme/session.test.js test/axomeme/predict.test.js \"test/axomeme/vendor/*.test.js\" && mocha -R spec --exit test/axomeme/descriptor.test.js", "test:slurm": "npm run test:analyses && npm run test:slurm-unit && npm run test:integration", "test:analyses": "mocha -R spec --exit test/absrel/absrel.js test/axomeme/axomeme.js test/busted/busted.js test/contrast-fel/contrast-fel.js test/fade/fade.js test/fubar/fubar.js test/gard/gard.js test/hivtrace/hivtrace.js test/meme/meme.js test/multihit/multihit.js test/prime/prime.js test/relax/relax.js test/slac/slac.js", - "test:slurm-unit": "mocha -R spec --exit test/jobstatus.js test/jobqueue.js test/coverage-gaps/job-lifecycle.js test/coverage-gaps/hivtrace-checkjob.js", + "test:slurm-unit": "mocha -R spec --exit test/jobstatus.js test/jobqueue.js test/reconcile.js test/coverage-gaps/job-lifecycle.js test/coverage-gaps/hivtrace-checkjob.js", "test:integration": "mocha -R spec --exit test/integration/slac-completes.js", "test:mcp": "npm run test:mcp-core && npm run test:mcp-axomeme", "test:mcp-core": "mocha -R spec --exit test/mcp/mcp.test.js test/mcp/oauth-redirect.test.js test/mcp/job-notifier.test.js test/mcp/sse-notifications.test.js", diff --git a/server.js b/server.js index ae685704..e8e82d6d 100644 --- a/server.js +++ b/server.js @@ -40,88 +40,102 @@ const io = require("socket.io")(ioPort, ioOptions); // Use the shared redis@5 client factory (see lib/redis-client.js). redis@5 is // promise-native, so commands return promises and are camelCased // (del stays del, hgetall -> hGetAll). -const client = require("./lib/redis-client").client; +const { client, ready } = require("./lib/redis-client"); +const reconcile = require("./lib/reconcile"); +const mcp = require("./lib/mcp"); -// clear active_jobs list -// TODO: we should do more than just clear the active_jobs list -client.del("active_jobs").catch(function(err) { - logger.error("Redis del active_jobs failed: " + err.message); -}); +// Reconcile active_jobs against the live scheduler queue (#455) instead of +// blindly clearing it — jobs still live on the scheduler survive the restart; +// terminal/orphaned entries are dropped and their hashes get retention TTLs +// (#453). Socket handler registration AND the MCP server (v3's second spawn +// surface) are gated on reconciliation so no new submission can race the +// snapshot or the SCAN sweep. reconcileActiveJobs never rejects (fail-open) +// and the snapshot exec carries a 15s timeout, so boot always completes; the +// .catch is belt-and-braces. The io port still binds immediately (same +// v2-shipped contract; clients in the sub-second window auto-reconnect). +ready + .then(function () { + return reconcile.reconcileActiveJobs(client); + }) + .then(registerHandlers) + .catch(function (err) { + logger.error("reconcile : unexpected boot error, registering handlers anyway: " + err.message); + registerHandlers(); + }); -// For every new connection... -io.sockets.on("connection", function(socket) { - //Routes - socket.on("job queue", function(jobs) { - JobQueue(function(jobs) { - socket.emit("job queue", jobs); - socket.disconnect(); +function registerHandlers() { + // For every new connection... + io.sockets.on("connection", function(socket) { + //Routes + socket.on("job queue", function(jobs) { + JobQueue(function(jobs) { + socket.emit("job queue", jobs); + socket.disconnect(); + }); }); - }); - // Query job status by ID (for reconnection after page refresh) - socket.on("job:status", function(params, callback) { - if (!params || !params.jobId) { - if (callback) callback({ status: "error", error: "Missing jobId" }); - return; - } - - // redis@5 hGetAll returns a promise resolving to the hash (an empty object - // when the key is missing), so treat an empty object as "not found". - client.hGetAll(params.jobId).then(function(jobData) { - if (!jobData || Object.keys(jobData).length === 0) { - if (callback) callback({ status: "not_found" }); + // Query job status by ID (for reconnection after page refresh) + socket.on("job:status", function(params, callback) { + if (!params || !params.jobId) { + if (callback) callback({ status: "error", error: "Missing jobId" }); return; } - const response = { - status: jobData.status || "unknown", - torque_id: jobData.torque_id - }; - - if (jobData.status === "completed" && jobData.results) { - // Results are stored as: {"results":"{ stringified JSON }","type":"completed"} - // We need to unwrap and parse the inner results string - try { - const parsedResults = JSON.parse(jobData.results); - if (parsedResults.results && typeof parsedResults.results === "string") { - response.results = JSON.parse(parsedResults.results); - } else { - response.results = parsedResults.results || parsedResults; + // redis@5 hGetAll returns a promise resolving to the hash (an empty object + // when the key is missing), so treat an empty object as "not found". + client.hGetAll(params.jobId).then(function(jobData) { + if (!jobData || Object.keys(jobData).length === 0) { + if (callback) callback({ status: "not_found" }); + return; + } + + const response = { + status: jobData.status || "unknown", + torque_id: jobData.torque_id + }; + + if (jobData.status === "completed" && jobData.results) { + // Results are stored as: {"results":"{ stringified JSON }","type":"completed"} + // We need to unwrap and parse the inner results string + try { + const parsedResults = JSON.parse(jobData.results); + if (parsedResults.results && typeof parsedResults.results === "string") { + response.results = JSON.parse(parsedResults.results); + } else { + response.results = parsedResults.results || parsedResults; + } + } catch (e) { + logger.error("Error parsing job results: " + e.message); + response.results = jobData.results; } - } catch (e) { - logger.error("Error parsing job results: " + e.message); - response.results = jobData.results; } - } - if (jobData.error) { - response.error = jobData.error; - } + if (jobData.error) { + response.error = jobData.error; + } - if (callback) callback(response); - }).catch(function(err) { - logger.error("Redis hGetAll job:status failed: " + err.message); - if (callback) callback({ status: "not_found" }); + if (callback) callback(response); + }).catch(function(err) { + logger.error("Redis hGetAll job:status failed: " + err.message); + if (callback) callback({ status: "not_found" }); + }); }); - }); - const r = new router.io(socket); + const r = new router.io(socket); - // Analysis routes are data-driven — see lib/routes/analysis-routes.js. - // It reproduces the 16 standard spawn/check/resubscribe/cancel blocks plus - // the special hivtrace analysis. (Phase 3, #410) - analysisRoutes.registerAnalysisRoutes(r, socket, { hivtrace: hivtrace }); - - // Acknowledge new connection - socket.emit("connected", { hello: "Ready to serve" }); - -}); + // Analysis routes are data-driven — see lib/routes/analysis-routes.js. + // It reproduces the 16 standard spawn/check/resubscribe/cancel blocks plus + // the special hivtrace analysis. (Phase 3, #410) + analysisRoutes.registerAnalysisRoutes(r, socket, { hivtrace: hivtrace }); + // Acknowledge new connection + socket.emit("connected", { hello: "Ready to serve" }); + }); -// Start MCP server on separate port -const mcp = require("./lib/mcp"); -mcp.startMcpServer(config, client); + // Start MCP server on separate port + mcp.startMcpServer(config, client); +} process.setMaxListeners(20); // bounded; GH #400 removed per-job cancelJob listeners diff --git a/test/reconcile-edges.js b/test/reconcile-edges.js new file mode 100644 index 00000000..971d1005 --- /dev/null +++ b/test/reconcile-edges.js @@ -0,0 +1,357 @@ +// #455 startup-reconciliation edge cases (v3 port of the v2 suite from PR +// #462). Complements test/reconcile.js (happy path against a real SLURM job) +// with the fail-open path, the resolve contract, zombie/reap edges, the +// atomic dedup rebuild, the qsub parsing branch, and two v3-only cases (the +// "local" submit_type snapshot and the 15s snapshot exec timeout) — all +// against PATH-shim scheduler binaries so no real jobs are ever submitted, +// and a live, ISOLATED redis instance (config.json) — never a production +// instance. +// +// server.js gates registerHandlers() (socket routes + the MCP server) on the +// reconcileActiveJobs promise; the "resolves on every path" assertions here +// are what that gate depends on — any rejection would fail these awaits. +// +// config.submit_type is mutated per-case on the require-cached lib/config +// export and ALWAYS restored (before/afterEach/after): test:ci runs every +// suite in one mocha process, so a leaked mutation would poison later suites. +var fs = require("fs"), + os = require("os"), + path = require("path"), + should = require("should"), + redis = require("redis"), + config = require("../lib/config"), + redisClient = require("../lib/redis-client"), + reconcile = require("../lib/reconcile.js"); + +var client = redis.createClient(redisClient.buildClientOptions()); + +describe("startup reconciliation edge cases (#455)", function () { + this.timeout(60000); + + var suffix = Date.now(); + var seeded = []; + var origPath = process.env.PATH; + var origSubmitType = config.submit_type; + + var shimBase = path.join(os.tmpdir(), "reconcile455-shims-" + process.pid); + var failDir = path.join(shimBase, "fail"); + var liveDir = path.join(shimBase, "live"); + var qsubDir = path.join(shimBase, "qsub"); + var hangDir = path.join(shimBase, "hang"); + + function writeShim(dir, name, body) { + var p = path.join(dir, name); + fs.writeFileSync(p, "#!/bin/sh\n" + body + "\n"); + fs.chmodSync(p, 493 /* 0755 */); + } + + // Shim squeue that reports exactly `ids` as the live scheduler snapshot. + function liveSqueue(ids) { + writeShim( + liveDir, + "squeue", + ids + .map(function (i) { + return "echo " + i; + }) + .join("\n") || "true" + ); + process.env.PATH = liveDir + ":" + origPath; + } + + function track(id) { + seeded.push(id); + return id; + } + + before(async function () { + [failDir, liveDir, qsubDir, hangDir].forEach(function (d) { + fs.mkdirSync(d, { recursive: true }); + }); + // scheduler snapshot unavailable: nonzero exit from both listers + writeShim(failDir, "squeue", "exit 2"); + writeShim(failDir, "qstat", "exit 2"); + // scheduler lister that never answers (drives the 15s exec timeout) + writeShim(hangDir, "squeue", "sleep 60"); + // qstat output shape: two header lines, then one job per row, id first + // (printf '%s\n' so the all-dashes separator is not parsed as options) + writeShim( + qsubDir, + "qstat", + "printf '%s\\n' 'Job ID Name User Time S Queue'\n" + + "printf '%s\\n' '------ ---- ---- ---- - -----'\n" + + "printf '%s\\n' '77777.silverback stub sweaver 0 R batch'" + ); + // Baseline every case on the squeue branch — the CI config is + // submit_type "local", which never execs a lister at all. + config.submit_type = "slurm"; + await client.connect(); + }); + + afterEach(function () { + // every test mutates PATH (and some mutate the cached config object); + // always restore the suite baseline so a failing test cannot poison the + // rest of the run + process.env.PATH = origPath; + config.submit_type = "slurm"; + }); + + after(async function () { + process.env.PATH = origPath; + config.submit_type = origSubmitType; + var m = client.multi().del("active_jobs"); + seeded.forEach(function (id) { + m.del(id); + }); + await m.exec(); + client.destroy(); + fs.rmSync(shimBase, { recursive: true, force: true }); + }); + + it("FAIL OPEN: scheduler snapshot failure leaves active_jobs and hashes untouched but still resolves", async function () { + var id = track("test-455e-" + suffix + "-failopen"); + process.env.PATH = failDir; // ONLY the failing shims resolvable + await client + .multi() + .del("active_jobs") + // seed a duplicate on purpose: even the dedup rebuild must not run + .rPush("active_jobs", [id, id]) + .hSet(id, { + status: "running", + torque_id: JSON.stringify({ torque_id: "999999991" }), + }) + .exec(); + await reconcile.reconcileActiveJobs(client); // resolves despite the failure + var entries = await client.lRange("active_jobs", 0, -1); + entries.should.deepEqual([id, id]); // no rebuild, duplicate intact + var obj = await client.hGetAll(id); + obj.status.should.equal("running"); // not zombified + should.not.exist(obj.error); // no script-error write + var ttl = await client.ttl(id); + ttl.should.equal(-1); // no TTL applied + await client.del([id, "active_jobs"]); + }); + + it("FAIL OPEN on a hung lister: the 15s snapshot exec timeout resolves without touching active_jobs", async function () { + // v3-only: the v2 module had no exec timeout, so a hung squeue hung the + // boot gate forever. The promisified exec kills the child after 15s and + // the error takes the same fail-open path as a nonzero exit. + var id = track("test-455e-" + suffix + "-hung"); + process.env.PATH = hangDir + ":" + origPath; + await client + .multi() + .del("active_jobs") + .rPush("active_jobs", id) + .hSet(id, { + status: "running", + torque_id: JSON.stringify({ torque_id: "999999992" }), + }) + .exec(); + var started = Date.now(); + await reconcile.reconcileActiveJobs(client); + var elapsed = Date.now() - started; + elapsed.should.be.aboveOrEqual(14000); // the timeout, not a fast failure + elapsed.should.be.below(45000); // ...and not the shim's full sleep 60 + var entries = await client.lRange("active_jobs", 0, -1); + entries.should.deepEqual([id]); // untouched + (await client.hGet(id, "status")).should.equal("running"); + (await client.ttl(id)).should.equal(-1); + await client.del([id, "active_jobs"]); + }); + + it("submit_type local: empty snapshot with no lister exec — in-flight entries are zombied", async function () { + // v3-only: a restart severs the local child process and its watcher, so + // local in-flight jobs are unfinalizable and zombie-marking is correct. + // PATH holds ONLY the failing shims: if the local branch exec'd any + // lister it would fail open and this test would see status "running". + var id = track("test-455e-" + suffix + "-local"); + config.submit_type = "local"; + process.env.PATH = failDir; + await client + .multi() + .del("active_jobs") + .rPush("active_jobs", id) + .hSet(id, { + status: "running", + torque_id: JSON.stringify({ torque_id: "999999993" }), + }) + .exec(); + await reconcile.reconcileActiveJobs(client); + (await client.hGet(id, "status")).should.equal("aborted"); + var ttl = await client.ttl(id); + ttl.should.be.above(0); + ttl.should.be.belowOrEqual(redisClient.TERMINAL_TTL_SECONDS); + (await client.lRange("active_jobs", 0, -1)).should.deepEqual([]); + await client.del(id); + }); + + it("absent active_jobs list: resolves and unrelated keys are untouched", async function () { + var doneHash = track("test-455e-" + suffix + "-donehash"); + var blobKey = track("test-455e-" + suffix + "-blob"); + liveSqueue([]); + await client + .multi() + .del("active_jobs") + // a terminal hash and a plain-string result blob must survive the + // finalize + zombie sweep untouched (TYPE guard, terminal skip) + .hSet(doneHash, { status: "completed" }) + .set(blobKey, "result-payload") + .exec(); + await reconcile.reconcileActiveJobs(client); + var entries = await client.lRange("active_jobs", 0, -1); + entries.should.deepEqual([]); + (await client.hGet(doneHash, "status")).should.equal("completed"); + (await client.ttl(doneHash)).should.equal(-1); // sweep never TTLs terminal hashes + (await client.get(blobKey)).should.equal("result-payload"); + (await client.ttl(blobKey)).should.equal(-1); + await client.del([doneHash, blobKey]); + }); + + it("zombifies a running entry whose torque_id field is raw non-JSON garbage", async function () { + var id = track("test-455e-" + suffix + "-garbage"); + liveSqueue(["424242001"]); // scheduler is up, just doesn't know this job + await client + .multi() + .del("active_jobs") + .rPush("active_jobs", id) + .hSet(id, { status: "running", torque_id: "@@not json at all@@" }) + .exec(); + await reconcile.reconcileActiveJobs(client); + var obj = await client.hGetAll(id); + obj.status.should.equal("aborted"); + JSON.parse(obj.error).type.should.equal("script error"); + var ttl = await client.ttl(id); + ttl.should.be.above(0); + ttl.should.be.belowOrEqual(redisClient.TERMINAL_TTL_SECONDS); + (await client.lRange("active_jobs", 0, -1)).should.deepEqual([]); // dropped from the rebuilt list + await client.del(id); + }); + + it("zombifies a running entry whose parsed torque_id is null", async function () { + // parseable JSON, but torque_id itself is null -> same zombie path + var id = track("test-455e-" + suffix + "-nulltid"); + liveSqueue(["424242001"]); + await client + .multi() + .del("active_jobs") + .rPush("active_jobs", id) + .hSet(id, { + status: "queued", + torque_id: JSON.stringify({ torque_id: null }), + }) + .exec(); + await reconcile.reconcileActiveJobs(client); + (await client.hGet(id, "status")).should.equal("aborted"); + var ttl = await client.ttl(id); + ttl.should.be.above(0); + ttl.should.be.belowOrEqual(redisClient.TERMINAL_TTL_SECONDS); + await client.del(id); + }); + + it("silently reaps a list entry with no backing hash, keeping a live sibling", async function () { + var ghost = track("test-455e-" + suffix + "-ghost"); + var keeper = track("test-455e-" + suffix + "-keeper"); + liveSqueue(["424242010"]); + await client + .multi() + .del("active_jobs") + .rPush("active_jobs", [ghost, keeper]) // ghost has NO hash + .hSet(keeper, { + status: "running", + torque_id: JSON.stringify({ torque_id: "424242010" }), + }) + .exec(); + await reconcile.reconcileActiveJobs(client); + (await client.lRange("active_jobs", 0, -1)).should.deepEqual([keeper]); + (await client.exists(ghost)).should.equal(0); // no hash conjured for the ghost + await client.del([keeper, "active_jobs"]); + }); + + it("rebuilds [a,b,a,b,a] to exactly [a,b] when both are live (atomic dedup, original order)", async function () { + var a = track("test-455e-" + suffix + "-dup-a"); + var b = track("test-455e-" + suffix + "-dup-b"); + liveSqueue(["424242020", "424242021"]); + await client + .multi() + .del("active_jobs") + .rPush("active_jobs", [a, b, a, b, a]) + .hSet(a, { + status: "running", + torque_id: JSON.stringify({ torque_id: "424242020" }), + }) + .hSet(b, { + status: "running", + torque_id: JSON.stringify({ torque_id: "424242021" }), + }) + .exec(); + await reconcile.reconcileActiveJobs(client); + (await client.lRange("active_jobs", 0, -1)).should.deepEqual([a, b]); // each survivor once, a-then-b order + (await client.ttl(a)).should.equal(-1); // survivors stay TTL-less + (await client.ttl(b)).should.equal(-1); + await client.del([a, b, "active_jobs"]); + }); + + it("drops an already-cancelled entry from the list with the terminal retention ttl", async function () { + // test/reconcile.js covers status=completed -> the completed TTL; this + // is the other terminal branch (cancelled -> terminal TTL) + var id = track("test-455e-" + suffix + "-cancelled"); + liveSqueue([]); + await client + .multi() + .del("active_jobs") + .rPush("active_jobs", id) + .hSet(id, { + status: "cancelled", + torque_id: JSON.stringify({ torque_id: "424242030" }), + }) + .exec(); + await reconcile.reconcileActiveJobs(client); + (await client.lRange("active_jobs", 0, -1)).should.deepEqual([]); + (await client.hGet(id, "status")).should.equal("cancelled"); // status untouched, only TTL'd + var ttl = await client.ttl(id); + ttl.should.be.above(0); + ttl.should.be.belowOrEqual(redisClient.TERMINAL_TTL_SECONDS); + await client.del(id); + }); + + it("sweeps an in-flight hash orphaned outside active_jobs (pre-fix restart residue)", async function () { + var orphan = track("test-455e-" + suffix + "-orphan"); + liveSqueue([]); + await client + .multi() + .del("active_jobs") // orphan is NOT in the list + .hSet(orphan, { + status: "running", + torque_id: JSON.stringify({ torque_id: "424242040" }), + }) + .exec(); + await reconcile.reconcileActiveJobs(client); + (await client.hGet(orphan, "status")).should.equal("aborted"); + var ttl = await client.ttl(orphan); + ttl.should.be.above(0); + ttl.should.be.belowOrEqual(redisClient.TERMINAL_TTL_SECONDS); + await client.del(orphan); + }); + + it("qsub branch: qstat column-1 id keeps the matching in-flight entry alive with no ttl", async function () { + var id = track("test-455e-" + suffix + "-qsub"); + // reconcile reads submit_type from the require-cached config at call + // time, so mutating the shared object routes it down the qstat branch + config.submit_type = "qsub"; + process.env.PATH = qsubDir + ":" + origPath; + await client + .multi() + .del("active_jobs") + .rPush("active_jobs", id) + .hSet(id, { + status: "running", + torque_id: JSON.stringify({ torque_id: "77777.silverback" }), + }) + .exec(); + await reconcile.reconcileActiveJobs(client); + (await client.lRange("active_jobs", 0, -1)).should.deepEqual([id]); // kept alive via qstat parse + (await client.hGet(id, "status")).should.equal("running"); + (await client.ttl(id)).should.equal(-1); // in-flight survivor stays TTL-less + await client.del([id, "active_jobs"]); + }); +}); diff --git a/test/reconcile-sweep.js b/test/reconcile-sweep.js new file mode 100644 index 00000000..d4864fd4 --- /dev/null +++ b/test/reconcile-sweep.js @@ -0,0 +1,346 @@ +// #455 startup reconciliation — zombie SCAN sweep tests (v3 port of the v2 +// suite from PR #462). test/reconcile.js only exercises the active_jobs LIST +// pass; this file covers the trailing SCAN sweep over job hashes that are +// NOT in active_jobs (the pre-fix restart legacy: `del active_jobs` left the +// hashes behind, TTL-less). +// +// Run against a live, ISOLATED redis instance (config.json) — never against +// a production instance. This suite additionally uses redis db 9 so its +// flushDb/seeding cannot collide with the other suites (which use db 0) — +// which works precisely because reconcile only ever touches the client the +// caller passes in. +// +// No real scheduler jobs are submitted: the scheduler snapshot is taken from +// a PATH-shim `squeue` executable that prints a fixed set of "live" ids, so +// the sweep's live/dead decisions are fully deterministic. +var fs = require("fs"), + os = require("os"), + path = require("path"), + should = require("should"), + redis = require("redis"), + config = require("../lib/config"), + logger = require("../lib/logger.js").logger, + redisClient = require("../lib/redis-client"), + reconcile = require("../lib/reconcile.js"); + +// Dedicated db index: the sweep SCANs the entire selected db, so isolation +// keeps the swept/reaped counters in the summary line deterministic. +var client = redis.createClient( + Object.assign({}, redisClient.buildClientOptions(), { database: 9 }) +); + +describe("reconcile zombie-hash SCAN sweep (#455)", function () { + this.timeout(60000); + + var suffix = Date.now(); + + // Sweep candidates (deliberately NOT pushed onto active_jobs): + var sweptRunning = "test-sweep-" + suffix + "-dead-running"; + var sweptQueued = "test-sweep-" + suffix + "-dead-queued"; + var sweptGarbageTid = "test-sweep-" + suffix + "-garbage-tid"; + var untouchedLive = "test-sweep-" + suffix + "-live-running"; + var schedDead = "test-sweep-" + suffix + "-sched-dead"; + var schedLive = "test-sweep-" + suffix + "-sched-live"; + var stringKey = "test-sweep-" + suffix + "-string"; + var listKey = "test-sweep-" + suffix + "-list"; + var terminalCompleted = "test-sweep-" + suffix + "-terminal-completed"; + var terminalAborted = "test-sweep-" + suffix + "-terminal-aborted"; + // The one active_jobs entry — handled by the list pass, must be skipped + // (seen set) by the sweep even though its scheduler job is dead: + var seenZombie = "test-sweep-" + suffix + "-seen-zombie"; + // Filler string keys (> one SCAN batch at COUNT 500) so scanIterator's + // batch-array semantics are actually exercised across multiple batches. + var BULK_KEYS = 600; + + // Fake scheduler ids the PATH-shim squeue reports as live: + var liveTid = "424242"; + + var shimDir = null; + var savedPath = process.env.PATH; + var origSubmitType = config.submit_type; + var warnLines = []; + var summaryLine = null; + // Command spies, installed only for the reconcile call: which keys the + // pass HGETALLed / TYPEd. These pin the sweep's *mechanism* (a guarded key + // is never inspected), not just the surviving-state outcome — the outcome + // alone is reachable even with the guards deleted, because the sweep + // swallows errors and the list pass zombifies seenZombie before the sweep + // reads it (same connection, sequential awaits). + var hgetallKeys = []; + var typedKeys = []; + + before(async function () { + // PATH-shim squeue: prints one "live" id per line — a numeric torque id + // plus the scheduler-side hash key itself (behavior 3's live case). + shimDir = fs.mkdtempSync(path.join(os.tmpdir(), "sweep-shim-")); + fs.writeFileSync( + path.join(shimDir, "squeue"), + "#!/bin/sh\necho " + liveTid + "\necho " + schedLive + "\n", + { mode: 493 } // 0755 + ); + process.env.PATH = shimDir + ":" + savedPath; + // Route reconcile down the squeue branch regardless of the local/CI + // config; restored in after() — test:ci shares one mocha process. + config.submit_type = "slurm"; + + await client.connect(); + // db 9 is exclusively ours — start from a clean slate so the SCAN sweep + // sees exactly the keys seeded below. + await client.flushDb(); + var m = client + .multi() + .rPush("active_jobs", seenZombie) + // 1. in-flight hashes not in active_jobs, scheduler job gone -> swept + .hSet(sweptRunning, { + status: "running", + torque_id: JSON.stringify({ torque_id: "999999901" }), + }) + .hSet(sweptQueued, { + status: "queued", + torque_id: JSON.stringify({ torque_id: "999999902" }), + }) + // adjacent edge: unparseable torque_id, no datamonkey_id -> no live id + // to check -> swept + .hSet(sweptGarbageTid, { + status: "running", + torque_id: "not json at all", + }) + // 2. in-flight, not in active_jobs, but scheduler job still live + .hSet(untouchedLive, { + status: "running", + torque_id: JSON.stringify({ torque_id: liveTid }), + }) + // 3. scheduler-side hashes: keyed by the scheduler id itself, carry + // datamonkey_id, no parseable torque_id + .hSet(schedDead, { status: "running", datamonkey_id: "dm-" + suffix }) + .hSet(schedLive, { status: "running", datamonkey_id: "dm-" + suffix }) + // 4. non-hash keys living alongside job hashes (TYPE guard) + .set(stringKey, "just a result blob") + .rPush(listKey, ["a", "b"]) + // 5. terminal hashes not in active_jobs -> not sweep candidates + .hSet(terminalCompleted, { + status: "completed", + torque_id: JSON.stringify({ torque_id: "999999903" }), + }) + .hSet(terminalAborted, { + status: "aborted", + torque_id: JSON.stringify({ torque_id: "999999904" }), + }) + // 6. the seen entry: dead scheduler job, zombified by the LIST pass + .hSet(seenZombie, { + status: "queued", + torque_id: JSON.stringify({ torque_id: "999999905" }), + }); + // 7. bulk filler so the SCAN spans multiple COUNT-500 batches + for (var i = 0; i < BULK_KEYS; i++) { + m.set("test-sweep-" + suffix + "-bulk-" + i, "x"); + } + await m.exec(); + + // Capture the reconcile summary so the swept counter (which only the + // SCAN sweep increments) is observable. + var realWarn = logger.warn.bind(logger); + var patchedWarn = function (msg) { + warnLines.push(String(msg)); + return realWarn.apply(null, arguments); + }; + logger.warn = patchedWarn; + // Record every key reconcile inspects. The sweep always TYPEs a + // candidate before HGETALL; the list pass HGETALLs without TYPE. + var realHgetall = client.hGetAll; + var realType = client.type; + client.hGetAll = function (key) { + hgetallKeys.push(key); + return realHgetall.apply(client, arguments); + }; + client.type = function (key) { + typedKeys.push(key); + return realType.apply(client, arguments); + }; + try { + // markZombie awaits its hSet + expire, so every terminal TTL has + // landed by the time this resolves — no v2-style polling needed. + await reconcile.reconcileActiveJobs(client); + } finally { + logger.warn = realWarn; + client.hGetAll = realHgetall; + client.type = realType; + } + warnLines.forEach(function (line) { + if (line.indexOf("zombie hashes swept=") !== -1) summaryLine = line; + }); + }); + + after(async function () { + process.env.PATH = savedPath; + config.submit_type = origSubmitType; + if (shimDir) { + fs.rmSync(shimDir, { recursive: true, force: true }); + } + // No real scheduler jobs were submitted, so nothing to scancel. + // db 9 is exclusively this suite's — drop every seeded/created key. + if (client.isOpen) { + await client.flushDb(); + client.destroy(); + } + }); + + async function shouldBeZombified(id) { + var obj = await client.hGetAll(id); + obj.status.should.equal("aborted"); + var parsed = JSON.parse(obj.error); + parsed.type.should.equal("script error"); + parsed.error.should.equal("orphaned at server restart"); + var ttl = await client.ttl(id); + ttl.should.be.above(0); + ttl.should.be.belowOrEqual(redisClient.TERMINAL_TTL_SECONDS); + } + + async function shouldBeUntouchedInFlight(id, expectedStatus) { + var obj = await client.hGetAll(id); + obj.status.should.equal(expectedStatus); + should.not.exist(obj.error); + (await client.ttl(id)).should.equal(-1); + } + + // --- 1. dead in-flight hashes outside active_jobs are swept --- + + it("sweeps a running hash whose scheduler job is gone: aborted + orphan error + terminal ttl", async function () { + await shouldBeZombified(sweptRunning); + }); + + it("sweeps a queued hash whose scheduler job is gone", async function () { + await shouldBeZombified(sweptQueued); + }); + + it("sweeps an in-flight hash with an unparseable torque_id and no datamonkey_id", async function () { + await shouldBeZombified(sweptGarbageTid); + }); + + // --- 2. live in-flight hash outside active_jobs is left alone --- + // Negative test: always green under a full sweep revert; it guards against + // OVER-sweeping and is only meaningful paired with the positive tests above. + + it("leaves a running hash untouched when its torque_id is in the scheduler snapshot", async function () { + await shouldBeUntouchedInFlight(untouchedLive, "running"); + }); + + // --- 3. scheduler-side hashes keyed by the scheduler id itself --- + + it("zombifies a scheduler-side hash (datamonkey_id, no torque_id) whose own key is absent from the snapshot", async function () { + await shouldBeZombified(schedDead); + }); + + // Negative test (over-sweeping guard): green under a sweep revert, paired + // with the schedDead positive above. + it("leaves a scheduler-side hash untouched when its own key is in the snapshot", async function () { + var obj = await client.hGetAll(schedLive); + obj.status.should.equal("running"); + obj.datamonkey_id.should.equal("dm-" + suffix); + should.not.exist(obj.error); + (await client.ttl(schedLive)).should.equal(-1); + }); + + // --- 4. TYPE guard: non-hash keys survive and the sweep still completes --- + // The two survival tests are negative tests: they stay green even with the + // TYPE guard deleted, because a WRONGTYPE reply only fails that key's read + // inside the sweep's try/catch and the keys survive either way. They guard + // against over-sweeping only; the guard *mechanism* is pinned by the spy + // test below, which fails if HGETALL is ever issued on a non-hash key. + + it("leaves a plain string key untouched with no WRONGTYPE failure", async function () { + (await client.get(stringKey)).should.equal("just a result blob"); + (await client.ttl(stringKey)).should.equal(-1); + }); + + it("leaves a list key untouched with no WRONGTYPE failure", async function () { + (await client.lRange(listKey, 0, -1)).should.deepEqual(["a", "b"]); + (await client.ttl(listKey)).should.equal(-1); + }); + + it("completes the sweep and reports the summary despite the non-hash keys", function () { + // reconcileActiveJobs resolved (we are past before()) — and the summary + // line was emitted, proving the sweep did not stall on a WRONGTYPE reply. + should.exist(summaryLine); + }); + + it("never issues HGETALL on the non-hash keys (TYPE-guard mechanism)", function () { + // The sweep visited both keys (TYPE probe recorded) but the guard kept + // HGETALL from ever being issued. Without the guard, HGETALL fires and + // its WRONGTYPE error aborts the batch via the sweep's catch — invisible + // to the survival tests above, but caught here. + typedKeys.should.containEql(stringKey); + typedKeys.should.containEql(listKey); + hgetallKeys.should.not.containEql(stringKey); + hgetallKeys.should.not.containEql(listKey); + }); + + it("iterates every SCAN batch (node-redis v5 yields ARRAYS of keys per batch)", function () { + // With 600+ filler keys the SCAN cannot fit one COUNT-500 batch. Every + // filler is a non-seen string key, so each must show up as a TYPE probe; + // a for-await that mistook the yielded batch array for a single key + // would probe (or sweep) nothing recognizable and fail here. + var bulkTyped = typedKeys.filter(function (key) { + return String(key).indexOf("test-sweep-" + suffix + "-bulk-") === 0; + }); + bulkTyped.length.should.equal(BULK_KEYS); + }); + + // --- 5. terminal hashes outside active_jobs are not sweep candidates --- + // Negative tests: always green under a full sweep revert; they guard + // against over-sweeping (terminal hashes gaining errors/TTLs), paired with + // the behavior-1 positives that catch under-sweeping. + + it("does not touch a completed hash outside active_jobs (no ttl added by the sweep)", async function () { + var obj = await client.hGetAll(terminalCompleted); + obj.status.should.equal("completed"); + should.not.exist(obj.error); + (await client.ttl(terminalCompleted)).should.equal(-1); + }); + + it("does not touch an already-aborted hash outside active_jobs", async function () { + var obj = await client.hGetAll(terminalAborted); + obj.status.should.equal("aborted"); + should.not.exist(obj.error); + (await client.ttl(terminalAborted)).should.equal(-1); + }); + + // --- 6. active_jobs entries (the seen set) are skipped by the sweep --- + + it("handles the active_jobs zombie via the list pass, not the sweep", async function () { + // Sanity: the list pass did zombify it (and dropped it from the list). + (await client.hGet(seenZombie, "status")).should.equal("aborted"); + (await client.lRange("active_jobs", 0, -1)).should.deepEqual([]); + }); + + it("the sweep never inspects the seen entry (seen-set mechanism)", function () { + // Outcome checks alone cannot pin the seen-set guard: the list pass's + // zombifying hSet lands before the sweep's read (sequential awaits on + // the same client), so with the guard deleted, seenZombie is already + // 'aborted' when the sweep reaches it and every state assertion still + // passes. Instead, pin the mechanism: the sweep always TYPEs a candidate + // before HGETALL, and the list pass never calls TYPE — so any TYPE probe + // on seenZombie can only mean the sweep ignored the seen set. + typedKeys.should.not.containEql(seenZombie); + // And the single recorded HGETALL is the list pass's; a second would be + // the sweep's. + hgetallKeys + .filter(function (key) { + return key === seenZombie; + }) + .length.should.equal(1); + }); + + it("counts exactly the four non-list zombies as swept — the seen entry is not double-counted", function () { + // Outcome check (the counter): the swept counter is incremented ONLY + // inside the SCAN sweep. In a clean db 9 the candidates are + // sweptRunning, sweptQueued, sweptGarbageTid and schedDead — four. + // seenZombie's scheduler job is just as dead, but it was reaped by the + // list pass; the mechanism test above pins that the sweep skipped it. + should.exist(summaryLine); + summaryLine.should.match(/zombie hashes swept=4(\s|$)/); + // And the list pass reaped exactly the one active_jobs entry. + summaryLine.should.match(/ reaped=1 /); + summaryLine.should.match(/ kept=0 /); + }); +}); diff --git a/test/reconcile.js b/test/reconcile.js new file mode 100644 index 00000000..0cc41e47 --- /dev/null +++ b/test/reconcile.js @@ -0,0 +1,101 @@ +// #455 startup reconciliation tests (v3 port of the v2 suite from PR #462). +// Requires a live SLURM scheduler (submit_type "slurm") and a live, ISOLATED +// redis instance (config.json) — never against a production instance. +// +// redis@5 port notes: dedicated promise client via the shared +// buildClientOptions() factory (NOT the shared client — this suite owns its +// connection lifetime); markZombie's hSet/expire are awaited inside +// reconcileActiveJobs now, so the v2 setTimeout(500) settle hack is gone. +var util = require("util"), + execP = util.promisify(require("child_process").exec), + should = require("should"), + redis = require("redis"), + config = require("../lib/config"), + redisClient = require("../lib/redis-client"), + reconcile = require("../lib/reconcile.js"); + +var client = redis.createClient(redisClient.buildClientOptions()); + +describe("startup scheduler reconciliation (#455)", function () { + this.timeout(30000); + + var suffix = Date.now(); + var liveId = "test-455-" + suffix + "-live"; + var zombieId = "test-455-" + suffix + "-zombie"; + var terminalId = "test-455-" + suffix + "-term"; + var sbatchId = null; + + before(async function () { + // Only the squeue branch is exercised here; local/qsub deployments have + // no live SLURM scheduler to snapshot. + if (config.submit_type !== "slurm") { + this.skip(); + return; + } + await client.connect(); + // A real scheduler job so liveId survives the reconciliation + var res = await execP( + "sbatch --wrap 'sleep 180' --partition=" + + (config.slurm_partition || "datamonkey") + ); + sbatchId = res.stdout.toString().match(/\d+/)[0]; + await client + .multi() + .del("active_jobs") + .rPush("active_jobs", [liveId, zombieId, zombieId, terminalId]) + .hSet(liveId, { + status: "running", + torque_id: JSON.stringify({ torque_id: sbatchId }), + }) + .hSet(zombieId, { + status: "queued", + torque_id: JSON.stringify({ torque_id: "999999999" }), + }) + .hSet(terminalId, { + status: "completed", + torque_id: JSON.stringify({ torque_id: "888888888" }), + }) + .exec(); + // markZombie awaits its hSet + expire before resolving, so no settle + // window is needed before asserting. + await reconcile.reconcileActiveJobs(client); + }); + + after(async function () { + if (sbatchId) await execP("scancel " + sbatchId).catch(function () {}); + if (client.isOpen) { + await client + .multi() + .lRem("active_jobs", 0, liveId) + .del(liveId) + .del(zombieId) + .del(terminalId) + .exec(); + client.destroy(); + } + }); + + it("keeps only the live job, exactly once", async function () { + var entries = await client.lRange("active_jobs", 0, -1); + entries.should.deepEqual([liveId]); + }); + + it("marks the zombie aborted with the terminal ttl", async function () { + var status = await client.hGet(zombieId, "status"); + status.should.equal("aborted"); + var ttl = await client.ttl(zombieId); + ttl.should.be.above(0); + ttl.should.be.belowOrEqual(redisClient.TERMINAL_TTL_SECONDS); + }); + + it("gives the already-terminal hash the completed retention ttl", async function () { + var ttl = await client.ttl(terminalId); + ttl.should.be.above(0); + ttl.should.be.belowOrEqual(redisClient.COMPLETED_TTL_SECONDS); + }); + + it("leaves the live job hash without a ttl", async function () { + var ttl = await client.ttl(liveId); + ttl.should.equal(-1); + }); +}); diff --git a/test/regression/boot-reconcile-gate.js b/test/regression/boot-reconcile-gate.js new file mode 100644 index 00000000..9088f539 --- /dev/null +++ b/test/regression/boot-reconcile-gate.js @@ -0,0 +1,41 @@ +// #455 static tripwire: server.js must reconcile active_jobs against the +// scheduler at boot instead of blindly deleting the list, and both spawn +// surfaces (socket routes AND the MCP server) must be gated behind that +// reconciliation. A pure source-text check — no redis, no scheduler, no +// server boot — so it runs first-class in CI and fails loudly if a future +// refactor reintroduces the old `client.del("active_jobs")` or hoists a +// spawn surface out of the gate. +var fs = require("fs"), + path = require("path"), + should = require("should"); + +var serverSrc = fs.readFileSync( + path.join(__dirname, "../../server.js"), + "utf8" +); +var reconcileSrc = fs.readFileSync( + path.join(__dirname, "../../lib/reconcile.js"), + "utf8" +); + +describe("boot-time reconciliation gate (#455 tripwire)", function () { + it("server.js never blindly deletes active_jobs", function () { + serverSrc.should.not.containEql('del("active_jobs"'); + serverSrc.should.not.containEql("del('active_jobs'"); + }); + + it("reconciliation gates both spawn surfaces (socket routes and the MCP server)", function () { + var reconcileAt = serverSrc.indexOf("reconcileActiveJobs"); + var socketAt = serverSrc.indexOf('io.sockets.on("connection"'); + var mcpAt = serverSrc.indexOf("startMcpServer("); + reconcileAt.should.be.aboveOrEqual(0); + socketAt.should.be.aboveOrEqual(0); + mcpAt.should.be.aboveOrEqual(0); + reconcileAt.should.be.below(socketAt); + reconcileAt.should.be.below(mcpAt); + }); + + it("lib/reconcile.js keeps the fail-open contract", function () { + reconcileSrc.should.containEql("leaving active_jobs untouched"); + }); +});