gitoriaLog in with ident

mpackdb

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Main branchmaster87888725release 1.0.7caramboleyomaster/tests/multi-process-write.test.js

5.8 KB

  1. import MPackDB from '../src/MPackDB.js';
  2. import { rm } from 'fs/promises';
  3. import { spawn } from 'child_process';
  4. import { fileURLToPath } from 'url';
  5. import { dirname, join } from 'path';
  6. const TEST_DB_PATH = 'db/multi-process';
  7. const WORKER = join(dirname(fileURLToPath(import.meta.url)), 'helpers/multi-process-worker.js');
  8. console.log('--- Multi-Process Write Test ---');
  9. console.log('Several separate node processes write to the same database files.');
  10. console.log('This is the mpackdb-admin scenario: a second process must be able');
  11. console.log('to write while others read/write the same db.\n');
  12. await rm(TEST_DB_PATH, { recursive: true, force: true });
  13. function runWorker(dbPath, workerId, mode) {
  14. return new Promise((resolve, reject) => {
  15. const child = spawn(process.execPath, [WORKER, dbPath, String(workerId), mode], { stdio: ['ignore', 'pipe', 'inherit'] });
  16. let out = '';
  17. child.stdout.on('data', d => out += d);
  18. child.on('exit', code => code === 0 ? resolve(out) : reject(new Error(`worker ${workerId} (${mode}) exited with ${code}`)));
  19. });
  20. }
  21. const sleep = ms => new Promise(r => setTimeout(r, ms));
  22. let failures = 0;
  23. // ============================================================
  24. // PHASE 1: 3 concurrent processes insert + delete on a string-PK store
  25. // ============================================================
  26. console.log('PHASE 1: 3 concurrent worker processes, 15 inserts + 5 deletes each');
  27. const utxoPath = `${TEST_DB_PATH}/utxos`;
  28. await Promise.all([1, 2, 3].map(w => runWorker(utxoPath, w, 'utxo')));
  29. const utxos = new MPackDB(utxoPath, { primaryKey: 'id', indexes: ['symbol', '*blockNumber'], compact: false });
  30. for (const w of [1, 2, 3]) {
  31. for (let i = 0; i < 15; i++) {
  32. const deleted = i % 3 === 0;
  33. const found = await utxos.find(`w${w}-${i}`);
  34. const expected = deleted ? 0 : 1;
  35. if (found.length !== expected) {
  36. failures++;
  37. console.log(` ✗ w${w}-${i}: indexed find returned ${found.length}, expected ${expected}`);
  38. }
  39. }
  40. // Secondary index must see the other processes' records too
  41. const bySymbol = await utxos.find(null, { index: { field: 'symbol', value: `SYM${w}` } });
  42. if (bySymbol.length !== 10) {
  43. failures++;
  44. console.log(` ✗ symbol index for SYM${w}: ${bySymbol.length} records, expected 10`);
  45. }
  46. }
  47. const fullScan = await utxos.find(() => true);
  48. console.log(` full scan: ${fullScan.length} records (expected: 30)`);
  49. if (fullScan.length !== 30) failures++;
  50. console.log(` ${failures === 0 ? '✓ PASS' : '✗ FAIL'}\n`);
  51. await utxos.close();
  52. // ============================================================
  53. // PHASE 2: numeric auto-increment PK across processes — no id collisions
  54. // ============================================================
  55. console.log('PHASE 2: 3 concurrent worker processes on a numeric-PK store');
  56. const numPath = `${TEST_DB_PATH}/numeric`;
  57. const outputs = await Promise.all([1, 2, 3].map(w => runWorker(numPath, w, 'numeric')));
  58. const allIds = outputs.flatMap(o => JSON.parse(o.trim()));
  59. const uniqueIds = new Set(allIds);
  60. console.log(` ${allIds.length} inserts, ${uniqueIds.size} unique ids (expected: 45 / 45)`);
  61. const phase2Fail = allIds.length !== 45 || uniqueIds.size !== 45;
  62. if (phase2Fail) failures++;
  63. const numDb = new MPackDB(numPath, { primaryKey: '*id', compact: false });
  64. const numRecords = await numDb.find();
  65. console.log(` records on disk: ${numRecords.length} (expected: 45)`);
  66. if (numRecords.length !== 45) failures++;
  67. console.log(` ${!phase2Fail && numRecords.length === 45 ? '✓ PASS' : '✗ FAIL'}\n`);
  68. await numDb.close();
  69. // ============================================================
  70. // PHASE 3: live visibility — a worker inserts (never persisting its indexes,
  71. // then dying without close) while the parent's already-open instance watches
  72. // ============================================================
  73. console.log('PHASE 3: live cross-process visibility without index persistence');
  74. const livePath = `${TEST_DB_PATH}/live`;
  75. const liveDb = new MPackDB(livePath, { primaryKey: 'id', indexes: ['symbol'], compact: false });
  76. await liveDb.init();
  77. const liveWorker = spawn(process.execPath, [WORKER, livePath, '9', 'live'], { stdio: ['pipe', 'pipe', 'inherit'] });
  78. await new Promise((resolve, reject) => {
  79. liveWorker.stdout.on('data', d => { if (String(d).includes('INSERTED')) resolve(); });
  80. liveWorker.on('exit', () => reject(new Error('live worker died before inserting')));
  81. setTimeout(() => reject(new Error('live worker timeout')), 15000);
  82. });
  83. // Worker inserted 5 records but never persisted its indexes. Our instance must
  84. // pick them up via meta/data refresh + index tail scan.
  85. let liveSeen = 0;
  86. for (let attempt = 0; attempt < 40; attempt++) {
  87. const seen = await liveDb.find(null, { index: { field: 'symbol', value: 'LIVE' } });
  88. liveSeen = seen.length;
  89. if (liveSeen === 5) break;
  90. await sleep(250);
  91. }
  92. console.log(` parent sees ${liveSeen}/5 live records via secondary index`);
  93. if (liveSeen !== 5) failures++;
  94. const livePk = await liveDb.find('live-3');
  95. console.log(` primary key lookup live-3: ${livePk.length} (expected: 1)`);
  96. if (livePk.length !== 1) failures++;
  97. liveWorker.stdin.write('exit\n');
  98. await new Promise(resolve => liveWorker.on('exit', resolve));
  99. // Worker died without close() — a FRESH instance must still index everything
  100. const freshDb = new MPackDB(livePath, { primaryKey: 'id', indexes: ['symbol'], compact: false });
  101. const freshSeen = await freshDb.find(null, { index: { field: 'symbol', value: 'LIVE' } });
  102. console.log(` fresh instance after worker crash: ${freshSeen.length}/5 via index`);
  103. if (freshSeen.length !== 5) failures++;
  104. console.log(` ${liveSeen === 5 && livePk.length === 1 && freshSeen.length === 5 ? '✓ PASS' : '✗ FAIL'}\n`);
  105. await freshDb.close();
  106. await liveDb.close();
  107. // ============================================================
  108. console.log(failures === 0 ? '✓ All multi-process write tests passed!' : `✗ ${failures} multi-process checks FAILED`);
  109. if (failures > 0) throw new Error('multi-process-write tests FAILED');

Branches

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