Skip to content

Commit d24c87b

Browse files
refactor(pipeline): thread run context via state, remove shared globals
CRITICAL from the bug hunt: _activeRunDir / _activeDbRunId were module-level singletons in pipeline.js, read by saveResults (crash-safe run-dir copy) and upsertResult (DB write gate). Two overlapping runs — a continuous-loop tail-write racing a concurrent /api/start takeover — could alias each other's run dir / DB row and write results into the wrong run (the Test-#11 contamination class). Remove the shared mutable state by threading the run context through the per-run `state` object (already passed to every runner; state.activeDbRunId already existed): - saveResults(state) reads state.activeRunDir; upsertResult(state, result, snip) gates on state.activeDbRunId. All 37 call sites updated to pass the in-scope state. - Deleted _activeRunDir/_activeDbRunId + setActiveRunDir/setActiveDbRunId/ getActiveDbRunId (grep-zero remaining refs). server.js assigns state.activeRunDir/state.activeDbRunId directly where it used to call the setters. continuous.js sets st.activeDbRunId=null/st.activeRunDir=null on its own loopState — the getActiveDbRunId()/restore dance (a workaround for the shared global) is gone. Single-run behavior is identical; overlapping runs can no longer alias. Two-stage reviewed (spec-compliant + code-quality APPROVED): every call site passes a populated state, undefined/null guard semantics match the old globals, no DB-write or crash-safe-copy is silently skipped, TEST RUN + failure-log untouched. Test suite green (188/31/45/35).
1 parent 781abd0 commit d24c87b

3 files changed

Lines changed: 82 additions & 81 deletions

File tree

audit/continuous.js

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@ import { EventEmitter } from 'events';
2323
import { writeFileSync, readFileSync, existsSync, mkdirSync } from 'fs';
2424
import path from 'path';
2525
import { fileURLToPath } from 'url';
26-
import { createState, setActiveDbRunId, getActiveDbRunId, triggerPipelineStop } from './pipeline.js';
26+
import { createState, triggerPipelineStop } from './pipeline.js';
2727
import { sleep } from '../protocol/speedtest.js';
2828

2929
const __dirname = path.dirname(fileURLToPath(import.meta.url));
@@ -304,12 +304,17 @@ async function _runOnePass(loopState, batchId, frozenNodes = null) {
304304
batchBroadcast,
305305
)
306306
: (st) => {
307-
const prev = getActiveDbRunId();
308-
setActiveDbRunId(null);
307+
// The continuous loop persists per-node to batch_results separately;
308+
// null this loop-state's run dir + db id so runAudit's results-table /
309+
// run-dir writes stay out of the runs/results tables. `st` is the loop's
310+
// own createState() instance, isolated from any direct run's state —
311+
// no restore needed (nothing else reads these fields off loopState).
312+
st.activeDbRunId = null;
313+
st.activeRunDir = null;
309314
return pipeline.runAudit(false, st, batchBroadcast, frozenNodes, {
310315
testRun: !!_ctrl.testRun,
311316
pricingMode: _ctrl.pricingMode || null,
312-
}).finally(() => setActiveDbRunId(prev));
317+
});
313318
};
314319
}
315320

audit/pipeline.js

Lines changed: 57 additions & 53 deletions
Original file line numberDiff line numberDiff line change
@@ -166,15 +166,9 @@ export function getResults() { return results; }
166166
// ─── Crash-Safe Results Persistence ────────────────────────────────────────
167167
// Write to temp file then rename (atomic on most filesystems).
168168
// Also continuously save to the active run directory so a kill never loses data.
169-
let _activeRunDir = null;
170-
let _activeDbRunId = null;
171-
172-
export function setActiveRunDir(dir) { _activeRunDir = dir; }
173-
174-
/** Set the SQLite run_id for the current audit run (called from server.js). */
175-
export function setActiveDbRunId(id) { _activeDbRunId = id; }
176-
/** Get the current SQLite run_id (null if not yet set). */
177-
export function getActiveDbRunId() { return _activeDbRunId; }
169+
// Run context (run dir + SQLite run_id) lives on the per-run `state` object
170+
// (`state.activeRunDir` / `state.activeDbRunId`) so two overlapping runs can never
171+
// alias each other's destination — there is no shared module global anymore.
178172

179173
// ─── On-chain reporter (per-run state) ───────────────────────────────────────
180174
// Buffers per-node records during a run and self-sends a memo TX with the
@@ -276,15 +270,17 @@ async function _finalizeOnchainReporter() {
276270
_onchainReporter = null;
277271
}
278272

279-
export function saveResults() {
273+
export function saveResults(state) {
280274
const data = JSON.stringify(results, null, 2);
281275
const tmpFile = RESULTS_FILE + '.tmp';
282276
writeFileSync(tmpFile, data, 'utf8');
283277
try { renameSync(tmpFile, RESULTS_FILE); } catch { writeFileSync(RESULTS_FILE, data, 'utf8'); }
284278

285-
// Also save to the active run directory (crash-safe copy)
286-
if (_activeRunDir) {
287-
const runFile = path.join(_activeRunDir, 'results.json');
279+
// Also save to the active run directory (crash-safe copy). Run dir is carried
280+
// on the per-run state so overlapping runs can't write to each other's dir.
281+
const _runDir = state && state.activeRunDir;
282+
if (_runDir) {
283+
const runFile = path.join(_runDir, 'results.json');
288284
const runTmp = runFile + '.tmp';
289285
try {
290286
writeFileSync(runTmp, data, 'utf8');
@@ -322,7 +318,7 @@ function _sanitizeSnippet(raw) {
322318
return s.length > _MAX_SNIPPET ? s.slice(-_MAX_SNIPPET) : s;
323319
}
324320

325-
function upsertResult(result, logSnippet = null) {
321+
function upsertResult(state, result, logSnippet = null) {
326322
const idx = results.findIndex(r => r.address === result.address);
327323
if (idx !== -1) results[idx] = result;
328324
else results.push(result);
@@ -333,9 +329,12 @@ function upsertResult(result, logSnippet = null) {
333329
}
334330

335331
// ─── SQLite persistence (non-blocking — failure must not stop the audit) ─
336-
if (_activeDbRunId != null) {
332+
// DB run_id lives on the per-run state so overlapping runs can't write into
333+
// each other's results table.
334+
const _dbRunId = state && state.activeDbRunId;
335+
if (_dbRunId != null) {
337336
try {
338-
const resultId = _dbInsertResult(_activeDbRunId, result);
337+
const resultId = _dbInsertResult(_dbRunId, result);
339338
// For ANY failed test (no speed measured OR error/errorCode set), write a
340339
// detailed error_log row. The popup contract from CLAUDE.md says every
341340
// failure MUST have a copyable log — never fall through silently.
@@ -407,6 +406,11 @@ export function createState() {
407406
// with that exact node first instead of whatever order the parallel
408407
// online-scan happens to return.
409408
resumeHeadAddr: null,
409+
// Run context — set by server.js (startFreshRun / loadRunIntoState) and by
410+
// the continuous loop. Read by saveResults() (crash-safe run-dir copy) and
411+
// upsertResult() (SQLite run_id gate). Per-run so overlapping runs can't alias.
412+
activeRunDir: null,
413+
activeDbRunId: null,
410414
};
411415
}
412416

@@ -628,7 +632,7 @@ export async function runAudit(resume, state, broadcast, preloadedNodes = null,
628632
state.passedBaseline = 0;
629633
state.nodeSpeedHistory = [];
630634
state.baselineHistory = [];
631-
saveResults();
635+
saveResults(state);
632636
}
633637

634638
broadcast('log', { msg: `Balance: ${state.balance}` });
@@ -933,9 +937,9 @@ export async function runAudit(resume, state, broadcast, preloadedNodes = null,
933937
state.failedNodes++;
934938
}
935939
commitBaselineSample(node.address); // 1 baseline sample per recorded node
936-
upsertResult(result); // success — no snippet needed
940+
upsertResult(state, result); // success — no snippet needed
937941
state.resumeHeadAddr = null;
938-
saveResults();
942+
saveResults(state);
939943
broadcast('result', { result, state });
940944
if (result.actualMbps != null) {
941945
const ts = new Date().toLocaleTimeString('en-US', { hour12: false });
@@ -1002,9 +1006,9 @@ export async function runAudit(resume, state, broadcast, preloadedNodes = null,
10021006
const failResult = buildFailResult(node, status, state, errMsg, error?.diag || {});
10031007
state.failedNodes++;
10041008
commitBaselineSample(node.address); // 1 baseline sample per recorded node
1005-
upsertResult(failResult, _logSnippet);
1009+
upsertResult(state, failResult, _logSnippet);
10061010
state.resumeHeadAddr = null;
1007-
saveResults();
1011+
saveResults(state);
10081012
broadcast('result', { result: failResult, state });
10091013
const retryLabel = retried > 0 ? ` (${retried} retries)` : '';
10101014
const label = /timeout/i.test(errMsg) ? '⏱ Timeout' : /already exists/i.test(errMsg) ? '🚫 Node bug' : 'FAIL';
@@ -1128,22 +1132,22 @@ export async function runAudit(resume, state, broadcast, preloadedNodes = null,
11281132
state.failedNodes = Math.max(0, state.failedNodes - 1);
11291133
state.testedNodes++;
11301134
if (result.pass10mbps) state.passed10++;
1131-
upsertResult(result);
1132-
saveResults();
1135+
upsertResult(state, result);
1136+
saveResults(state);
11331137
broadcast('result', { result, state });
11341138
broadcast('log', { msg: ` ✓ Internet-recovery retest PASS: ${result.actualMbps} Mbps` });
11351139
} else if (result) {
11361140
// Truthy result but null mbps — still a failure (it was already counted
11371141
// as failed before this retest). Persist it without touching counters.
1138-
upsertResult(result, _sanitizeSnippet(_irSnippet));
1139-
saveResults();
1142+
upsertResult(state, result, _sanitizeSnippet(_irSnippet));
1143+
saveResults(state);
11401144
broadcast('result', { result, state });
11411145
broadcast('log', { msg: ` ✗ Internet-recovery retest FAIL: ${result.errorCode || 'no speed'}` });
11421146
} else {
11431147
const errMsg = error?.message || 'Unknown';
11441148
const failResult = buildFailResult(node, status, state, errMsg, error?.diag || {});
1145-
upsertResult(failResult, _sanitizeSnippet(_irSnippet));
1146-
saveResults();
1149+
upsertResult(state, failResult, _sanitizeSnippet(_irSnippet));
1150+
saveResults(state);
11471151
broadcast('result', { result: failResult, state });
11481152
broadcast('log', { msg: ` ✗ Internet-recovery retest FAIL: ${errMsg.slice(0, 80)}` });
11491153
}
@@ -1204,23 +1208,23 @@ export async function runAudit(resume, state, broadcast, preloadedNodes = null,
12041208
state.failedNodes = Math.max(0, state.failedNodes - 1);
12051209
state.testedNodes++;
12061210
if (result.pass10mbps) state.passed10++;
1207-
upsertResult(result);
1208-
saveResults();
1211+
upsertResult(state, result);
1212+
saveResults(state);
12091213
broadcast('result', { result, state });
12101214
broadcast('log', { msg: ` ✓ Retest PASS: ${result.actualMbps} Mbps` });
12111215
} else if (result) {
12121216
// Truthy result but null mbps (e.g. SESSION_UNMAPPED) — still a failure (it
12131217
// was already counted as failed before this retest). Persist without touching
12141218
// counters; flipping to "tested" here would wrongly clear a still-failed node.
1215-
upsertResult(result, _sanitizeSnippet(_ironSnippet));
1216-
saveResults();
1219+
upsertResult(state, result, _sanitizeSnippet(_ironSnippet));
1220+
saveResults(state);
12171221
broadcast('result', { result, state });
12181222
broadcast('log', { msg: ` ✗ Retest FAIL: ${result.errorCode || 'no speed'}` });
12191223
} else {
12201224
const errMsg = error?.message || 'Unknown';
12211225
const failResult = buildFailResult(node, status, state, errMsg, error?.diag || {});
1222-
upsertResult(failResult, _sanitizeSnippet(_ironSnippet));
1223-
saveResults();
1226+
upsertResult(state, failResult, _sanitizeSnippet(_ironSnippet));
1227+
saveResults(state);
12241228
broadcast('result', { result: failResult, state });
12251229
broadcast('log', { msg: ` ✗ Retest FAIL: ${errMsg.slice(0, 80)}` });
12261230
}
@@ -1401,8 +1405,8 @@ export async function runRetestSkips(skipAddrs, state, broadcast) {
14011405

14021406
if (result && result.actualMbps != null) {
14031407
recomputeCounters(state);
1404-
upsertResult(result);
1405-
saveResults();
1408+
upsertResult(state, result);
1409+
saveResults(state);
14061410
broadcast('result', { result, state });
14071411
state.retestPassed++;
14081412
const sla = result.actualMbps >= 10 ? 'SLA:PASS' : 'SLA:FAIL';
@@ -1413,8 +1417,8 @@ export async function runRetestSkips(skipAddrs, state, broadcast) {
14131417
const errMsg = error?.message || result?.error || 'Unknown';
14141418
const failResult = result || buildFailResult(node, null, state, errMsg, error?.diag || {});
14151419
recomputeCounters(state);
1416-
upsertResult(failResult, _sanitizeSnippet(_rrSnippet));
1417-
saveResults();
1420+
upsertResult(state, failResult, _sanitizeSnippet(_rrSnippet));
1421+
saveResults(state);
14181422
broadcast('result', { result: failResult, state });
14191423
state.retestFailed++;
14201424
broadcast('log', { msg: ` ✗ #${testNum} FAIL: ${errMsg.slice(0, 80)}` });
@@ -1686,8 +1690,8 @@ export async function runPlanTest(planId, state, broadcast) {
16861690
const errResult = buildFailResult(node, status, state, `plan-session: ${_sanitizeSnippet(err.message)}`, { planId, subscriptionId });
16871691
errResult.inPlan = true;
16881692
errResult.planIds = [planId];
1689-
upsertResult(errResult, _sanitizeSnippet(_ptSnippetSess));
1690-
saveResults();
1693+
upsertResult(state, errResult, _sanitizeSnippet(_ptSnippetSess));
1694+
saveResults(state);
16911695
broadcast('result', { result: errResult, state });
16921696
continue;
16931697
}
@@ -1717,8 +1721,8 @@ export async function runPlanTest(planId, state, broadcast) {
17171721
if (result.slaApplicable && result.pass15mbps) state.passed15++;
17181722
if (result.pass10mbps) state.passed10++;
17191723
if (result.passBaseline) state.passedBaseline++;
1720-
upsertResult(result);
1721-
saveResults();
1724+
upsertResult(state, result);
1725+
saveResults(state);
17221726
broadcast('result', { result, state });
17231727
if (result.actualMbps != null) {
17241728
planPassed++;
@@ -1737,8 +1741,8 @@ export async function runPlanTest(planId, state, broadcast) {
17371741
const failResult = buildFailResult(node, status, state, `plan-test: ${errMsg}`, error?.diag || {});
17381742
failResult.inPlan = true;
17391743
failResult.planIds = [planId];
1740-
upsertResult(failResult, _sanitizeSnippet(_ptSnippet));
1741-
saveResults();
1744+
upsertResult(state, failResult, _sanitizeSnippet(_ptSnippet));
1745+
saveResults(state);
17421746
broadcast('result', { result: failResult, state });
17431747
broadcast('log', { msg: ` ✗ Test error: ${errMsg}` });
17441748
}
@@ -2058,9 +2062,9 @@ export async function runSubPlanTest(planId, subscriptionId, granterAddr, state,
20582062
evictResult.diag = evictResult.diag || {};
20592063
evictResult.diag.viaSubscription = true;
20602064
evictResult.diag.evicted = true;
2061-
upsertResult(evictResult, _sanitizeSnippet(`PLAN_EVICTED: node ${node.address} not in plan ${planId} (refreshed ${_refreshAge}ms ago)`));
2065+
upsertResult(state, evictResult, _sanitizeSnippet(`PLAN_EVICTED: node ${node.address} not in plan ${planId} (refreshed ${_refreshAge}ms ago)`));
20622066
state.resumeHeadAddr = null;
2063-
saveResults();
2067+
saveResults(state);
20642068
broadcast('result', { result: evictResult, state });
20652069
continue;
20662070
}
@@ -2076,9 +2080,9 @@ export async function runSubPlanTest(planId, subscriptionId, granterAddr, state,
20762080
evictResult.diag = evictResult.diag || {};
20772081
evictResult.diag.viaSubscription = true;
20782082
evictResult.diag.evicted = true;
2079-
upsertResult(evictResult, _sanitizeSnippet(`PLAN_EVICTED: chain refresh confirmed node ${node.address} not in plan ${planId}`));
2083+
upsertResult(state, evictResult, _sanitizeSnippet(`PLAN_EVICTED: chain refresh confirmed node ${node.address} not in plan ${planId}`));
20802084
state.resumeHeadAddr = null;
2081-
saveResults();
2085+
saveResults(state);
20822086
broadcast('result', { result: evictResult, state });
20832087
continue;
20842088
}
@@ -2160,9 +2164,9 @@ export async function runSubPlanTest(planId, subscriptionId, granterAddr, state,
21602164
errResult.diag.selfGranter = _selfGranter;
21612165
errResult.diag.granter = granterAddr;
21622166
if (_isEviction) errResult.diag.evicted = true;
2163-
upsertResult(errResult, _sanitizeSnippet(_spSnippetSess));
2167+
upsertResult(state, errResult, _sanitizeSnippet(_spSnippetSess));
21642168
state.resumeHeadAddr = null;
2165-
saveResults();
2169+
saveResults(state);
21662170
broadcast('result', { result: errResult, state });
21672171
continue;
21682172
}
@@ -2194,9 +2198,9 @@ export async function runSubPlanTest(planId, subscriptionId, granterAddr, state,
21942198
if (result.slaApplicable && result.pass15mbps) state.passed15++;
21952199
if (result.pass10mbps) state.passed10++;
21962200
if (result.passBaseline) state.passedBaseline++;
2197-
upsertResult(result);
2201+
upsertResult(state, result);
21982202
state.resumeHeadAddr = null;
2199-
saveResults();
2203+
saveResults(state);
22002204
broadcast('result', { result, state });
22012205
if (result.actualMbps != null) {
22022206
subPassed++;
@@ -2220,9 +2224,9 @@ export async function runSubPlanTest(planId, subscriptionId, granterAddr, state,
22202224
failResult.diag.feeGranted = !_selfGranter;
22212225
failResult.diag.selfGranter = _selfGranter;
22222226
failResult.diag.granter = granterAddr;
2223-
upsertResult(failResult, _sanitizeSnippet(_spSnippet));
2227+
upsertResult(state, failResult, _sanitizeSnippet(_spSnippet));
22242228
state.resumeHeadAddr = null;
2225-
saveResults();
2229+
saveResults(state);
22262230
broadcast('result', { result: failResult, state });
22272231
broadcast('log', { msg: ` ✗ Test error: ${errMsg}` });
22282232
}

0 commit comments

Comments
 (0)