mpackdb
All repositories: gitoria
5.8 KB
import MPackDB from '../src/MPackDB.js';import { rm } from 'fs/promises';import { spawn } from 'child_process';import { fileURLToPath } from 'url';import { dirname, join } from 'path';const TEST_DB_PATH = 'db/multi-process';const WORKER = join(dirname(fileURLToPath(import.meta.url)), 'helpers/multi-process-worker.js');console.log('--- Multi-Process Write Test ---');console.log('Several separate node processes write to the same database files.');console.log('This is the mpackdb-admin scenario: a second process must be able');console.log('to write while others read/write the same db.\n');await rm(TEST_DB_PATH, { recursive: true, force: true });function runWorker(dbPath, workerId, mode) {return new Promise((resolve, reject) => {const child = spawn(process.execPath, [WORKER, dbPath, String(workerId), mode], { stdio: ['ignore', 'pipe', 'inherit'] });let out = '';child.stdout.on('data', d => out += d);child.on('exit', code => code === 0 ? resolve(out) : reject(new Error(`worker ${workerId} (${mode}) exited with ${code}`)));});}const sleep = ms => new Promise(r => setTimeout(r, ms));let failures = 0;// ============================================================// PHASE 1: 3 concurrent processes insert + delete on a string-PK store// ============================================================console.log('PHASE 1: 3 concurrent worker processes, 15 inserts + 5 deletes each');const utxoPath = `${TEST_DB_PATH}/utxos`;await Promise.all([1, 2, 3].map(w => runWorker(utxoPath, w, 'utxo')));const utxos = new MPackDB(utxoPath, { primaryKey: 'id', indexes: ['symbol', '*blockNumber'], compact: false });for (const w of [1, 2, 3]) {for (let i = 0; i < 15; i++) {const deleted = i % 3 === 0;const found = await utxos.find(`w${w}-${i}`);const expected = deleted ? 0 : 1;if (found.length !== expected) {failures++;console.log(` ✗ w${w}-${i}: indexed find returned ${found.length}, expected ${expected}`);}}// Secondary index must see the other processes' records tooconst bySymbol = await utxos.find(null, { index: { field: 'symbol', value: `SYM${w}` } });if (bySymbol.length !== 10) {failures++;console.log(` ✗ symbol index for SYM${w}: ${bySymbol.length} records, expected 10`);}}const fullScan = await utxos.find(() => true);console.log(` full scan: ${fullScan.length} records (expected: 30)`);if (fullScan.length !== 30) failures++;console.log(` ${failures === 0 ? '✓ PASS' : '✗ FAIL'}\n`);await utxos.close();// ============================================================// PHASE 2: numeric auto-increment PK across processes — no id collisions// ============================================================console.log('PHASE 2: 3 concurrent worker processes on a numeric-PK store');const numPath = `${TEST_DB_PATH}/numeric`;const outputs = await Promise.all([1, 2, 3].map(w => runWorker(numPath, w, 'numeric')));const allIds = outputs.flatMap(o => JSON.parse(o.trim()));const uniqueIds = new Set(allIds);console.log(` ${allIds.length} inserts, ${uniqueIds.size} unique ids (expected: 45 / 45)`);const phase2Fail = allIds.length !== 45 || uniqueIds.size !== 45;if (phase2Fail) failures++;const numDb = new MPackDB(numPath, { primaryKey: '*id', compact: false });const numRecords = await numDb.find();console.log(` records on disk: ${numRecords.length} (expected: 45)`);if (numRecords.length !== 45) failures++;console.log(` ${!phase2Fail && numRecords.length === 45 ? '✓ PASS' : '✗ FAIL'}\n`);await numDb.close();// ============================================================// PHASE 3: live visibility — a worker inserts (never persisting its indexes,// then dying without close) while the parent's already-open instance watches// ============================================================console.log('PHASE 3: live cross-process visibility without index persistence');const livePath = `${TEST_DB_PATH}/live`;const liveDb = new MPackDB(livePath, { primaryKey: 'id', indexes: ['symbol'], compact: false });await liveDb.init();const liveWorker = spawn(process.execPath, [WORKER, livePath, '9', 'live'], { stdio: ['pipe', 'pipe', 'inherit'] });await new Promise((resolve, reject) => {liveWorker.stdout.on('data', d => { if (String(d).includes('INSERTED')) resolve(); });liveWorker.on('exit', () => reject(new Error('live worker died before inserting')));setTimeout(() => reject(new Error('live worker timeout')), 15000);});// Worker inserted 5 records but never persisted its indexes. Our instance must// pick them up via meta/data refresh + index tail scan.let liveSeen = 0;for (let attempt = 0; attempt < 40; attempt++) {const seen = await liveDb.find(null, { index: { field: 'symbol', value: 'LIVE' } });liveSeen = seen.length;if (liveSeen === 5) break;await sleep(250);}console.log(` parent sees ${liveSeen}/5 live records via secondary index`);if (liveSeen !== 5) failures++;const livePk = await liveDb.find('live-3');console.log(` primary key lookup live-3: ${livePk.length} (expected: 1)`);if (livePk.length !== 1) failures++;liveWorker.stdin.write('exit\n');await new Promise(resolve => liveWorker.on('exit', resolve));// Worker died without close() — a FRESH instance must still index everythingconst freshDb = new MPackDB(livePath, { primaryKey: 'id', indexes: ['symbol'], compact: false });const freshSeen = await freshDb.find(null, { index: { field: 'symbol', value: 'LIVE' } });console.log(` fresh instance after worker crash: ${freshSeen.length}/5 via index`);if (freshSeen.length !== 5) failures++;console.log(` ${liveSeen === 5 && livePk.length === 1 && freshSeen.length === 5 ? '✓ PASS' : '✗ FAIL'}\n`);await freshDb.close();await liveDb.close();// ============================================================console.log(failures === 0 ? '✓ All multi-process write tests passed!' : `✗ ${failures} multi-process checks FAILED`);if (failures > 0) throw new Error('multi-process-write tests FAILED');
Branches
- mastermain branch
Latest commits
- 87888725release 1.0.7caramboleyo
- c4cdb9b6node: import prefixes (Deno compat) + pre-existing index-state WIPcaramboleyo
- 0afb8f4bupdate now must be a callbackcaramboleyo
- cde73eb4release 1.0.6caramboleyo
- d01dda02add index hints, intersection, boundingBox; remove findByIndexcaramboleyo
- b8ffc1a0release 1.0.5caramboleyo
- d47876a1reimplemented lost features like indexed find and more testscaramboleyo
- 7f08da9afixed insert ignoring model definitioncaramboleyo
- 705774a9added flush before findcaramboleyo
- b4db6391initial commitcaramboleyo