gitoriaLog in with ident

mpackdb

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commitc4cdb9b6c4cdb9b6node: import prefixes (Deno compat) + pre-existing index-state WIPcaramboleyoc4cdb9b6/src/MPackDB.js

35.5 KB

  1. import { dirname, basename, extname, resolve } from 'node:path';
  2. import { createReadStream, createWriteStream } from 'node:fs';
  3. import { mkdir, open, stat, readFile, writeFile, rename, unlink } from 'node:fs/promises';
  4. import { AsyncLocalStorage } from 'node:async_hooks';
  5. import { serialize, deserialize, uuid, PrimaryKeyType, IndexType } from './mpack.js';
  6. import { Cursor } from './Cursor.js';
  7. import { IndexManager } from './IndexManager.js';
  8. // Tracks which async call chain currently owns a db's lock, so nested operations
  9. // (e.g. find inside delete, insert inside withLock) re-enter, while unrelated
  10. // concurrent operations queue up instead of walking into the critical section.
  11. const lockContext = new AsyncLocalStorage();
  12. export { PrimaryKeyType, IndexType };
  13. /**
  14. * MPackDB - A fast, local, append-only JSON database with MessagePack serialization
  15. *
  16. * Features:
  17. * - Append-only writes for high performance
  18. * - MessagePack binary serialization
  19. * - Optional indexes (numeric and lexical)
  20. * - File-based locking for concurrent access
  21. * - Auto-compaction on startup
  22. * - Auto-persistence of indexes
  23. */
  24. export class MPackDB {
  25. _dbFile = null;
  26. _primaryKeyType = null;
  27. _primaryKey = null;
  28. _classToUse = null;
  29. _indexes = [];
  30. _indexTypes = {};
  31. _uniqueIndexes = new Set();
  32. _initPromise = null;
  33. _initState = 0;
  34. _dataPath = null;
  35. _dataStream = null;
  36. _dataIno = null;
  37. _lastWrite = null;
  38. _lockPath = null;
  39. _mutexTail = Promise.resolve();
  40. _staleLockTimeout = 30000;
  41. _schema = null;
  42. _indexManager = null;
  43. _meta = {
  44. nextId: 0,
  45. deleted: [],
  46. };
  47. _debug = false;
  48. _indexPersistInterval = 60000; // 60 seconds default
  49. _indexPersistThreshold = 1000; // 1000 entries default
  50. _processExitHandler = null;
  51. /**
  52. * Create a new MPackDB instance
  53. *
  54. * @param {string} dbFile - Path to the database file (without extension)
  55. * @param {Object} options - Configuration options
  56. * @param {string} [options.primaryKey] - Primary key field name. Prefix with * for numeric (e.g., '*id'), @ for UUID (e.g., '@uuid'), or no prefix for string
  57. * @param {PrimaryKeyType} [options.primaryKeyType] - Explicit primary key type (overrides prefix)
  58. * @param {string[]} [options.indexes] - Array of field names to index. Use * prefix for numeric, @ for UUID
  59. * @param {boolean} [options.debug=false] - Enable debug logging
  60. * @param {number} [options.indexPersistInterval=60000] - Milliseconds between automatic index persistence (0 to disable)
  61. * @param {number} [options.indexPersistThreshold=1000] - Number of changes before auto-persisting indexes
  62. * @param {boolean} [options.compact=true] - Run compaction on init (set false whenever another process may have the same files open)
  63. * @param {number} [options.staleLockTimeout=30000] - Take over lock files older than this many ms (crashed holder); 0 disables
  64. *
  65. * @example
  66. * const db = new MPackDB('data/users', {
  67. * primaryKey: '*id', // Numeric auto-increment
  68. * indexes: ['email', '*age'], // Index email (lexical) and age (numeric)
  69. * indexPersistThreshold: 100
  70. * });
  71. */
  72. _compact = true;
  73. _hasPrimaryKeyValue(record) {
  74. if (!this._primaryKey || !record || !Object.prototype.hasOwnProperty.call(record, this._primaryKey)) {
  75. return false;
  76. }
  77. const value = record[this._primaryKey];
  78. return value !== undefined && value !== null;
  79. }
  80. constructor(dbFile, { primaryKey, primaryKeyType, indexes, debug, indexPersistInterval, indexPersistThreshold, compact, staleLockTimeout } = {}) {
  81. this._dbFile = dbFile || this._dbFile;
  82. this._debug = debug || this._debug;
  83. if (compact !== undefined) this._compact = compact;
  84. if (indexPersistInterval !== undefined) this._indexPersistInterval = indexPersistInterval;
  85. if (indexPersistThreshold !== undefined) this._indexPersistThreshold = indexPersistThreshold;
  86. if (staleLockTimeout !== undefined) this._staleLockTimeout = staleLockTimeout;
  87. // Parse primary key with optional prefix
  88. if (primaryKey) {
  89. if (primaryKey.startsWith('*')) {
  90. // *id = numeric primary key
  91. this._primaryKey = primaryKey.slice(1);
  92. this._primaryKeyType = PrimaryKeyType.NUMBER;
  93. } else if (primaryKey.startsWith('@')) {
  94. // @id = UUID primary key
  95. this._primaryKey = primaryKey.slice(1);
  96. this._primaryKeyType = PrimaryKeyType.UUID;
  97. } else {
  98. // id = string primary key (lexical)
  99. this._primaryKey = primaryKey;
  100. this._primaryKeyType = PrimaryKeyType.STRING;
  101. }
  102. // Allow explicit override
  103. if (primaryKeyType !== undefined) {
  104. this._primaryKeyType = primaryKeyType;
  105. }
  106. }
  107. // Primary key is always unique
  108. if (this._primaryKey) {
  109. this._uniqueIndexes.add(this._primaryKey);
  110. }
  111. const allIndexes = [...new Set([
  112. ...(this._primaryKey ? [this._primaryKey] : []),
  113. ...(indexes || []),
  114. ...(this._indexes || []),
  115. ])];
  116. this._indexes = [];
  117. for (let index of allIndexes) {
  118. let cleanIndex = index;
  119. // !field = unique index (can combine with * and @: !*field, !@field)
  120. let isUnique = false;
  121. if (cleanIndex.startsWith('!')) {
  122. isUnique = true;
  123. cleanIndex = cleanIndex.slice(1);
  124. }
  125. if (cleanIndex.startsWith('*')) {
  126. // *field = numeric index
  127. cleanIndex = cleanIndex.slice(1);
  128. this._indexTypes[cleanIndex] = IndexType.NUMERIC;
  129. } else if (cleanIndex.startsWith('@')) {
  130. // @field = UUID index (lexical)
  131. cleanIndex = cleanIndex.slice(1);
  132. this._indexTypes[cleanIndex] = IndexType.LEXICAL;
  133. } else if (cleanIndex === this._primaryKey && this._primaryKeyType === PrimaryKeyType.NUMBER) {
  134. // Primary key is numeric
  135. this._indexTypes[cleanIndex] = IndexType.NUMERIC;
  136. } else {
  137. // Default to lexical
  138. this._indexTypes[cleanIndex] = IndexType.LEXICAL;
  139. }
  140. if (isUnique) {
  141. this._uniqueIndexes.add(cleanIndex);
  142. }
  143. this._indexes.push(cleanIndex);
  144. }
  145. // Schema descriptor persisted into meta.json so other tools (e.g. mpackdb-admin)
  146. // can open this database with the correct primary key and indexes.
  147. if (this._primaryKey || this._indexes.length > 0) {
  148. const pkPrefix = this._primaryKeyType === PrimaryKeyType.NUMBER ? '*'
  149. : this._primaryKeyType === PrimaryKeyType.UUID ? '@' : '';
  150. this._schema = {
  151. primaryKey: this._primaryKey ? pkPrefix + this._primaryKey : undefined,
  152. indexes: this._indexes
  153. .filter(field => field !== this._primaryKey)
  154. .map(field =>
  155. (this._uniqueIndexes.has(field) ? '!' : '')
  156. + (this._indexTypes[field] === IndexType.NUMERIC ? '*' : '')
  157. + field
  158. ),
  159. };
  160. }
  161. }
  162. /**
  163. * Initialize the database (called automatically by other methods)
  164. * Performs compaction and loads metadata
  165. *
  166. * @returns {Promise<MPackDB>} The database instance
  167. */
  168. async init() {
  169. if (this._initState === 2) return this;
  170. if (this._initState === 1) {
  171. return this._initPromise;
  172. }
  173. this._initState = 1;
  174. return this._initPromise = new Promise(async success => {
  175. if (!this._dbFile) {
  176. throw new Error('No database file specified.');
  177. }
  178. const dbDir = dirname(this._dbFile);
  179. const baseName = basename(this._dbFile, extname(this._dbFile));
  180. await mkdir(dbDir, { recursive: true });
  181. this._dataPath = resolve(dbDir, `${baseName}.mpack`);
  182. this._metaPath = resolve(dbDir, `${baseName}.meta.json`);
  183. this._lockPath = resolve(dbDir, `${baseName}.lock`);
  184. try {
  185. this._meta = JSON.parse(await readFile(this._metaPath));
  186. } catch (e) {
  187. this._meta = {
  188. nextId: 0,
  189. deleted: [],
  190. };
  191. }
  192. const didCompact = this._compact;
  193. if (didCompact) await this.compact({ duringInit: true });
  194. this._dataStream = createWriteStream(this._dataPath, { flags: 'a' }); // needs to be set after compact as it replaces the file with a tmp file
  195. this._dataIno = (await stat(this._dataPath).catch(() => null))?.ino ?? null;
  196. // Persist schema into meta.json so other tools (mpackdb-admin) can
  197. // open this db with the right primary key / indexes. The compare
  198. // runs under the lock: refresh() has adopted any newer on-disk meta
  199. // by then, so a concurrent process' meta is never clobbered.
  200. if (this._schema) {
  201. await this._withFileLock('init-schema', async () => {
  202. if (JSON.stringify(this._meta.schema) !== JSON.stringify(this._schema)) {
  203. await this.persistMeta();
  204. }
  205. });
  206. }
  207. // Initialize IndexManager if indexes are specified
  208. if (this._indexes.length > 0) {
  209. const dbDir = dirname(this._dbFile);
  210. const baseName = basename(this._dbFile, extname(this._dbFile));
  211. this._indexManager = new IndexManager(dbDir, baseName, this._indexes, this._indexTypes, this._primaryKeyType, {
  212. // Index persistence rewrites shared files — must hold the db lock.
  213. // lockContext.exit: a threshold-triggered persist inside insert()
  214. // must NOT inherit insert's lock ownership (it outlives it), so it
  215. // queues for its own turn instead.
  216. runExclusive: fn => lockContext.exit(() => this._withFileLock('index-persist', fn)),
  217. });
  218. // Index files are shared with other processes — (re)build them under the lock
  219. await this._withFileLock('index-init', () =>
  220. this._indexManager.init(this._dataPath, { forceRebuild: didCompact })
  221. );
  222. // Start auto-persist with configured interval and threshold
  223. this._indexManager.startAutoPersist(this._indexPersistInterval, this._indexPersistThreshold);
  224. // Register signal handlers for Ctrl+C and kill signals
  225. this._processExitHandler = async () => {
  226. await this.close();
  227. process.exit(0);
  228. };
  229. process.once('SIGINT', this._processExitHandler);
  230. process.once('SIGTERM', this._processExitHandler);
  231. }
  232. this._initState = 2;
  233. this._initPromise = null;
  234. success(this);
  235. });
  236. }
  237. /**
  238. * Insert a new record into the database
  239. *
  240. * @param {Object} record - The record to insert
  241. * @param {Object} [options]
  242. * @param {boolean} [options.skipPrimaryKey=false] - Skip auto-generating primary key
  243. * @returns {Promise<string|number|Object>} The primary-key value when configured, otherwise the inserted record
  244. *
  245. * @example
  246. * const id = await db.insert({ name: 'Alice', age: 30 });
  247. * // Returns: 0 (the generated primary-key value)
  248. */
  249. async insert(record, { skipPrimaryKey = false } = {}) {
  250. await this.init();
  251. return this._withFileLock('insert', async () => {
  252. const recToInsert = this._classToUse
  253. ? Object.assign(new this._classToUse(), record)
  254. : { ...record };
  255. const hasPrimaryKeyValue = this._hasPrimaryKeyValue(recToInsert);
  256. if (this._primaryKey && !skipPrimaryKey && !hasPrimaryKeyValue) {
  257. if (this._primaryKeyType === PrimaryKeyType.NUMBER) {
  258. recToInsert[this._primaryKey] = this._meta.nextId++;
  259. } else if (this._primaryKeyType === PrimaryKeyType.UUID) {
  260. recToInsert[this._primaryKey] = uuid();
  261. }
  262. }
  263. // Check unique index constraints (skip auto-generated primary keys — guaranteed unique)
  264. if (this._indexManager && this._uniqueIndexes.size > 0) {
  265. const autoGenPK = this._primaryKey && !skipPrimaryKey && !hasPrimaryKeyValue
  266. && (this._primaryKeyType === PrimaryKeyType.NUMBER || this._primaryKeyType === PrimaryKeyType.UUID);
  267. const deletedSet = new Set(this._meta.deleted);
  268. for (const field of this._uniqueIndexes) {
  269. if (autoGenPK && field === this._primaryKey) continue;
  270. const value = recToInsert[field];
  271. if (value === undefined) continue;
  272. const existing = await this._indexManager.get(field, value);
  273. const live = existing.filter(e => !deletedSet.has(e.loc[0]));
  274. if (live.length > 0) {
  275. const err = new Error(`Duplicate key: ${field}=${value}`);
  276. err.code = 'DUPLICATE_KEY';
  277. err.field = field;
  278. err.value = value;
  279. throw err;
  280. }
  281. }
  282. }
  283. const packedBuffer = serialize(recToInsert);
  284. const fileStat = await stat(this._dataPath).catch(() => ({ size: 0 }));
  285. const offset = fileStat.size;
  286. const writePromise = new Promise(resolve => this._dataStream.write(packedBuffer, resolve));
  287. this._lastWrite = writePromise;
  288. await writePromise;
  289. const loc = [offset, packedBuffer.length];
  290. // indexes
  291. if (this._indexManager) {
  292. this._indexManager.insert(recToInsert, loc);
  293. }
  294. // Persist meta if we incremented nextId
  295. if (this._primaryKey && this._primaryKeyType === PrimaryKeyType.NUMBER && !skipPrimaryKey && !hasPrimaryKeyValue) {
  296. await this.persistMeta();
  297. }
  298. return this._primaryKey ? recToInsert[this._primaryKey] : recToInsert;
  299. });
  300. }
  301. /**
  302. * Update records matching a query.
  303. * The callback receives the old record and must return the new record.
  304. *
  305. * @param {string|number|Function} mixed - Primary key value or query function
  306. * @param {Function} callback - Receives old record, must return new record
  307. * @param {Object} [options]
  308. * @param {boolean} [options.upsert=false] - Insert if no records match
  309. * @param {Object|Array} [options.index] - Index hint(s) to narrow disk reads
  310. * @returns {Promise<Object[]>} Array of updated records
  311. *
  312. * @example
  313. * await db.update(0, record => {
  314. * record.age = 32;
  315. * delete record.address;
  316. * return record;
  317. * });
  318. */
  319. async update(mixed, callback, { upsert = false, index } = {}) {
  320. if (typeof callback !== 'function') {
  321. throw new Error('update requires a callback function as second argument');
  322. }
  323. let insertedRecords = await this.delete(mixed, {
  324. index,
  325. callback: async record => {
  326. return this.insert(await callback(record), { skipPrimaryKey: true });
  327. }
  328. });
  329. if (upsert && insertedRecords.length === 0) {
  330. insertedRecords = await this.insert(await callback({}));
  331. }
  332. return insertedRecords;
  333. }
  334. /**
  335. * Update records or insert if not found
  336. *
  337. * @param {string|number|Function} mixed - Primary key value or query function
  338. * @param {Function} callback - Receives old record (or {} if inserting), must return new record
  339. * @param {Object} [options]
  340. * @param {Object|Array} [options.index] - Index hint(s) to narrow disk reads
  341. * @returns {Promise<Object[]>} Array of updated/inserted records
  342. */
  343. async upsert(mixed, callback, { index } = {}) {
  344. return this.update(mixed, callback, { upsert: true, index });
  345. }
  346. /**
  347. * Delete records matching a query
  348. *
  349. * @param {string|Function} mixed - Primary key value or query function
  350. * @param {Object} [options]
  351. * @param {Object|Array} [options.index] - Index hint(s) to narrow disk reads
  352. * @param {Function} [options.callback] - Internal callback per deleted record (used by update)
  353. * @returns {Promise<Object[]>} Array of deleted records
  354. *
  355. * @example
  356. * // Delete by primary key
  357. * await db.delete(0);
  358. *
  359. * // Delete with query function
  360. * await db.delete(r => r.age < 18);
  361. *
  362. * // Delete with index hint
  363. * await db.delete(r => r.status === 'inactive', { index: { field: 'status', value: 'inactive' } });
  364. */
  365. async delete(mixed, { index, callback = async record => record } = {}) {
  366. await this.init();
  367. return this._withFileLock('delete', async () => {
  368. const promises = [];
  369. for await (const [record, offset] of this.find(mixed, { mode: 'mixed', index })) {
  370. this._meta.deleted.push(offset);
  371. if (this._indexManager) {
  372. await this._indexManager.remove(record, this._primaryKey);
  373. }
  374. promises.push(callback(record));
  375. }
  376. await this.persistMeta();
  377. return Promise.all(promises);
  378. });
  379. }
  380. /**
  381. * Finds records in the database
  382. * @param {undefined|string|number|function} [mixed] - Query: undefined/null for all, primary key value for PK lookup, function for filter
  383. * @param {Object} [options] - Query options
  384. * @param {Object|Array} [options.index] - Index hint(s) to narrow disk reads before filtering
  385. * @returns {Cursor} A cursor for iterating over results
  386. * @example
  387. * // All records
  388. * for await (const user of db.find()) { ... }
  389. *
  390. * // Primary key lookup
  391. * const [user] = await db.find(68);
  392. *
  393. * // Filter with index range — only reads records in x 100+
  394. * for await (const doc of db.find(r => r.x <= 200, { index: { field: 'x', from: 100 } })) { ... }
  395. *
  396. * // Index intersection — intersects offsets first, then streams matches
  397. * for await (const doc of db.find(r => r.x <= 200, {
  398. * index: [
  399. * { field: 'x', from: 100 },
  400. * { field: 'status', value: 'active' }
  401. * ]
  402. * })) { ... }
  403. */
  404. find(mixed, options = {}) {
  405. if (typeof mixed === 'undefined' || mixed === null || typeof mixed === 'function') {
  406. return new Cursor(this, mixed, options);
  407. } else {
  408. if (!this._primaryKey) {
  409. throw new Error('No primary key specified.');
  410. }
  411. // Use indexed lookup when available
  412. if (this._indexManager) {
  413. return new Cursor(this, null, {
  414. ...options,
  415. _indexLookup: { field: this._primaryKey, value: mixed }
  416. });
  417. }
  418. return new Cursor(this, record => {
  419. return record[this._primaryKey] === mixed;
  420. }, options);
  421. }
  422. }
  423. /**
  424. * Find records within a bounding box defined by 4 corners.
  425. * Requires numeric indexes on x and y fields.
  426. * Corners can be in any order — min/max are extracted automatically.
  427. *
  428. * @param {Array<{x: number, y: number}>} corners - 4 corner coordinates
  429. * @param {function} [filter] - Optional additional filter function
  430. * @returns {Cursor}
  431. * @example
  432. * const results = await db.boundingBox([
  433. * { x: -5, y: 5 }, { x: 5, y: 5 },
  434. * { x: 5, y: -5 }, { x: -5, y: -5 }
  435. * ]);
  436. */
  437. boundingBox(corners, filter) {
  438. const xs = corners.map(c => c.x);
  439. const ys = corners.map(c => c.y);
  440. const minX = Math.min(...xs), maxX = Math.max(...xs);
  441. const minY = Math.min(...ys), maxY = Math.max(...ys);
  442. return this.find(filter || null, {
  443. index: [
  444. { field: 'x', from: minX, to: maxX },
  445. { field: 'y', from: minY, to: maxY }
  446. ]
  447. });
  448. }
  449. /**
  450. * Execute a callback while holding the database lock.
  451. * The lock is re-entrant: find/insert/delete/update called inside
  452. * the callback reuse the same lock instead of deadlocking.
  453. *
  454. * Use this for compound operations that must be atomic, e.g.
  455. * find-then-insert (login pattern).
  456. *
  457. * @param {Function} callback - Async function to execute under lock
  458. * @returns {Promise<any>} The return value of the callback
  459. *
  460. * @example
  461. * const user = await db.withLock(async () => {
  462. * const [existing] = await db.find(u => u.email === email);
  463. * if (existing) return existing;
  464. * return db.insert({ email });
  465. * });
  466. */
  467. async withLock(callback) {
  468. await this.init();
  469. return this._withFileLock('withLock', callback);
  470. }
  471. /**
  472. * Generator that yields records from the database file
  473. * @param {function|null} [queryFn=null] - Optional filter function
  474. * @param {Object} [options] - Generator options
  475. * @param {string} [options.mode='record'] - Mode: 'record', 'raw', 'offset', or 'mixed'
  476. * @yields {Object|Buffer|Array} Records, buffers, or [record, offset, size] tuples depending on mode
  477. */
  478. async *recordGenerator(queryFn = null, { mode = 'record', _indexLookup, index } = {}) {
  479. await this.init();
  480. await this.refresh();
  481. // Indexed primary key lookup — O(log n) instead of full scan
  482. if (_indexLookup && this._indexManager) {
  483. yield* this._indexedLookup(_indexLookup.field, _indexLookup.value, mode);
  484. return;
  485. }
  486. // Index-based find: collect offsets from index(es), then stream only those records
  487. if (index && this._indexManager) {
  488. yield* this._indexedStream(index, queryFn, mode);
  489. return;
  490. }
  491. // Flush pending writes so reads see all inserted data.
  492. // (Waiting for 'drain' here could hang forever — 'drain' only fires
  493. // after a write() returned false, which small writes never do.)
  494. if (this._lastWrite) {
  495. await this._lastWrite;
  496. }
  497. // Check if the data file exists before attempting to read it
  498. try {
  499. await stat(this._dataPath);
  500. } catch (error) {
  501. // File doesn't exist - return empty generator (no records)
  502. return;
  503. }
  504. // Convert deleted array to Set for O(1) lookup instead of O(n)
  505. const deletedSet = new Set(this._meta.deleted);
  506. const readStream = createReadStream(this._dataPath);
  507. let chunks = [];
  508. let totalLength = 0;
  509. let processedBytes = 0;
  510. for await (const chunk of readStream) {
  511. chunks.push(chunk);
  512. totalLength += chunk.length;
  513. while (true) {
  514. // Exit 1: Not enough data to even read the 4-byte size header.
  515. if (totalLength < 4) {
  516. break;
  517. }
  518. // Safely read the header, even if it's split across chunks
  519. let headerBuffer;
  520. if (chunks[0].length >= 4) {
  521. headerBuffer = chunks[0];
  522. } else {
  523. // The header is fragmented, so we must concat just enough to read it.
  524. headerBuffer = Buffer.concat(chunks, 4);
  525. }
  526. const recSize = headerBuffer.readInt32LE(0);
  527. // Validate the record size to prevent infinite loops
  528. // A record must be at least as large as its header (4 bytes).
  529. // A size of 0 or less is invalid and indicates corruption.
  530. if (recSize <= 4) {
  531. throw new Error(`Invalid record size read from stream: ${recSize}`);
  532. }
  533. // Exit 2: We have the size, but not the full record yet.
  534. if (totalLength < recSize) {
  535. break;
  536. }
  537. // skip deleted records (only after we have the full record)
  538. if (deletedSet.has(processedBytes)) {
  539. [processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);
  540. continue; // Goes back to while (true)
  541. }
  542. switch (mode) {
  543. case 'raw':
  544. yield Buffer.concat(chunks, recSize).subarray(0, recSize);
  545. break;
  546. case 'mixed':
  547. case 'record':
  548. const recBuffer = Buffer.concat(chunks, recSize).subarray(0, recSize);
  549. // Skip the 4-byte size header
  550. const data = deserialize(recBuffer.subarray(4));
  551. const rec = this._classToUse
  552. ? Object.assign(new this._classToUse(), data)
  553. : data;
  554. if (queryFn) {
  555. if (queryFn(rec)) {
  556. yield mode === 'mixed' ? [rec, processedBytes, recSize] : rec;
  557. }
  558. } else {
  559. yield mode === 'mixed' ? [rec, processedBytes, recSize] : rec;
  560. }
  561. break;
  562. case 'offset':
  563. // if someone uses a queryFn with offset, we have to deserialize it first to perform the queryFn
  564. if (queryFn) {
  565. const recBuffer = Buffer.concat(chunks, recSize).subarray(0, recSize);
  566. const rec = deserialize(recBuffer.subarray(4));
  567. if (queryFn(rec)) {
  568. yield [processedBytes, recSize];
  569. }
  570. } else {
  571. yield [processedBytes, recSize];
  572. }
  573. break;
  574. default:
  575. throw new Error(`Invalid mode: ${mode}`);
  576. }
  577. [processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);
  578. }
  579. }
  580. }
  581. /**
  582. * Indexed lookup — reads records directly by offset from the index.
  583. * Uses binary search on the index file for O(log n) lookups.
  584. * @private
  585. */
  586. async *_indexedLookup(field, value, mode = 'record') {
  587. const entries = await this._indexManager.get(field, value);
  588. if (entries.length === 0) return;
  589. const deletedSet = new Set(this._meta.deleted);
  590. const fileHandle = await open(this._dataPath, 'r');
  591. try {
  592. for (const entry of entries) {
  593. const [offset, length] = entry.loc;
  594. if (deletedSet.has(offset)) continue;
  595. const buffer = Buffer.alloc(length);
  596. await fileHandle.read(buffer, 0, length, offset);
  597. switch (mode) {
  598. case 'raw':
  599. yield buffer;
  600. break;
  601. case 'mixed': {
  602. const data = deserialize(buffer.subarray(4));
  603. const rec = this._classToUse
  604. ? Object.assign(new this._classToUse(), data)
  605. : data;
  606. yield [rec, offset, length];
  607. break;
  608. }
  609. case 'record':
  610. default: {
  611. const data = deserialize(buffer.subarray(4));
  612. const rec = this._classToUse
  613. ? Object.assign(new this._classToUse(), data)
  614. : data;
  615. yield rec;
  616. break;
  617. }
  618. }
  619. }
  620. } finally {
  621. await fileHandle.close();
  622. }
  623. }
  624. /**
  625. * Collect offsets from one or more indexes, optionally intersect, then stream records.
  626. * @param {Object|Array} index - Single index hint or array of hints
  627. * @param {function|null} queryFn - Optional filter function
  628. * @param {string} mode - Output mode: 'record', 'raw', or 'mixed'
  629. * @private
  630. */
  631. async *_indexedStream(index, queryFn, mode = 'record') {
  632. const hints = Array.isArray(index) ? index : [index];
  633. const deletedSet = new Set(this._meta.deleted);
  634. // Collect offset sets from each index hint
  635. const offsetSets = [];
  636. for (const hint of hints) {
  637. const offsets = new Map(); // offset → [offset, length]
  638. if (hint.value !== undefined) {
  639. // Exact match via get()
  640. const entries = await this._indexManager.get(hint.field, hint.value);
  641. for (const e of entries) {
  642. if (!deletedSet.has(e.loc[0])) offsets.set(e.loc[0], e.loc);
  643. }
  644. } else {
  645. // Range scan via entries()
  646. for await (const e of this._indexManager.entries(hint.field, { from: hint.from, to: hint.to, direction: hint.direction })) {
  647. if (!deletedSet.has(e.loc[0])) offsets.set(e.loc[0], e.loc);
  648. }
  649. }
  650. offsetSets.push(offsets);
  651. }
  652. // Intersect: keep only offsets present in ALL sets
  653. let locations;
  654. if (offsetSets.length === 1) {
  655. locations = Array.from(offsetSets[0].values());
  656. } else {
  657. // Start with smallest set for efficiency
  658. offsetSets.sort((a, b) => a.size - b.size);
  659. const [smallest, ...rest] = offsetSets;
  660. locations = [];
  661. for (const [offset, loc] of smallest) {
  662. if (rest.every(s => s.has(offset))) locations.push(loc);
  663. }
  664. }
  665. // Stream records from the intersected locations
  666. const fileHandle = await open(this._dataPath, 'r');
  667. try {
  668. for (const [offset, length] of locations) {
  669. const buffer = Buffer.alloc(length);
  670. await fileHandle.read(buffer, 0, length, offset);
  671. if (mode === 'raw') {
  672. yield buffer;
  673. } else {
  674. const data = deserialize(buffer.subarray(4));
  675. const rec = this._classToUse
  676. ? Object.assign(new this._classToUse(), data)
  677. : data;
  678. if (queryFn && !queryFn(rec)) continue;
  679. if (mode === 'mixed') {
  680. yield [rec, offset, length];
  681. } else {
  682. yield rec;
  683. }
  684. }
  685. }
  686. } finally {
  687. await fileHandle.close();
  688. }
  689. }
  690. /**
  691. * Pick up changes made by OTHER processes. Called automatically before every
  692. * read operation and after acquiring the write lock.
  693. *
  694. * The in-memory meta is authoritative for this process — it may hold
  695. * un-persisted mutations of an in-flight write. It is only replaced when the
  696. * on-disk copy carries a NEWER version, i.e. another process persisted under
  697. * the file lock. (Unconditionally reloading here was the root cause of the
  698. * lost-tombstone/lost-insert corruption under concurrent reads + writes.)
  699. *
  700. * Also detects external data file changes: appended records are indexed via
  701. * a tail scan, a replaced file (external compaction) triggers a full reopen.
  702. */
  703. async refresh() {
  704. try {
  705. const diskMeta = JSON.parse(await readFile(this._metaPath));
  706. if ((diskMeta.version || 0) > (this._meta.version || 0)) {
  707. this._meta = { nextId: 0, deleted: [], ...diskMeta };
  708. }
  709. } catch (e) {
  710. // file missing or torn write in progress — keep current meta
  711. }
  712. const dataStat = await stat(this._dataPath).catch(() => null);
  713. if (dataStat) {
  714. if (this._dataIno !== null && dataStat.ino !== this._dataIno) {
  715. await this._reopenAfterExternalReplace(dataStat);
  716. } else if (this._indexManager) {
  717. await this._indexManager.catchUp(this._dataPath, dataStat.size);
  718. }
  719. }
  720. }
  721. /**
  722. * The data file was replaced by another process (external compaction):
  723. * all offsets changed and our write stream points at the orphaned inode.
  724. * Reopen the stream, force-reload meta, and reset the index state.
  725. * @private
  726. */
  727. async _reopenAfterExternalReplace(dataStat) {
  728. this.debug(`Data file was replaced externally (compaction), reopening ${this._dataPath}`);
  729. if (this._dataStream) {
  730. await new Promise((resolve, reject) => {
  731. this._dataStream.end((err) => err ? reject(err) : resolve());
  732. });
  733. this._dataStream = createWriteStream(this._dataPath, { flags: 'a' });
  734. }
  735. this._dataIno = dataStat.ino;
  736. try {
  737. this._meta = { nextId: 0, deleted: [], ...JSON.parse(await readFile(this._metaPath)) };
  738. } catch (e) { }
  739. if (this._indexManager) {
  740. await this._indexManager.resetAfterReplace(this._dataPath);
  741. }
  742. }
  743. async persistMeta() {
  744. // Monotonic version: refresh() in other processes reloads only when it
  745. // sees a version newer than its in-memory one
  746. this._meta.version = (this._meta.version || 0) + 1;
  747. const metaToSave = { ...this._meta };
  748. // Only save nextId if we have a numeric primary key
  749. if (this._primaryKeyType !== PrimaryKeyType.NUMBER) {
  750. delete metaToSave.nextId;
  751. }
  752. if (this._schema) {
  753. this._meta.schema = this._schema;
  754. metaToSave.schema = this._schema;
  755. }
  756. // Atomic write (tmp + rename) so concurrent readers never parse a torn file
  757. const tmpPath = `${this._metaPath}.${process.pid}.tmp`;
  758. await writeFile(tmpPath, JSON.stringify(metaToSave));
  759. await rename(tmpPath, this._metaPath);
  760. }
  761. /**
  762. * Compact the database by removing deleted records
  763. * This rewrites the data file without tombstones
  764. *
  765. * @returns {Promise<void>}
  766. */
  767. async compact({ duringInit = false } = {}) {
  768. return this._withFileLock('compact', async () => {
  769. const writeStream = createWriteStream(this._dataPath + '.tmp', { flags: 'w' });
  770. const records = duringInit
  771. ? this._rawRecordsForCompaction()
  772. : this.find(null, { mode: 'raw' });
  773. for await (const binary of records) {
  774. writeStream.write(binary);
  775. }
  776. await new Promise((resolve, reject) => {
  777. writeStream.end((err) => err ? reject(err) : resolve());
  778. });
  779. // The old write stream points at the inode that rename() is about to
  780. // orphan — writes to it would be silently lost. Close it first.
  781. if (this._dataStream) {
  782. await new Promise((resolve, reject) => {
  783. this._dataStream.end((err) => err ? reject(err) : resolve());
  784. });
  785. }
  786. // rename is atomic, so we can just rename the file and it will replace the old one
  787. await rename(this._dataPath + '.tmp', this._dataPath);
  788. if (this._dataStream) {
  789. this._dataStream = createWriteStream(this._dataPath, { flags: 'a' });
  790. }
  791. this._dataIno = (await stat(this._dataPath).catch(() => null))?.ino ?? null;
  792. this._meta.deleted = [];
  793. await this.persistMeta();
  794. // Rebuild indexes — offsets changed after compaction
  795. if (!duringInit && this._indexManager) {
  796. await this._indexManager.init(this._dataPath, { forceRebuild: true });
  797. }
  798. });
  799. }
  800. async *_rawRecordsForCompaction() {
  801. const deletedSet = new Set(this._meta.deleted);
  802. try {
  803. await stat(this._dataPath);
  804. } catch (error) {
  805. return;
  806. }
  807. const readStream = createReadStream(this._dataPath);
  808. let chunks = [];
  809. let totalLength = 0;
  810. let processedBytes = 0;
  811. for await (const chunk of readStream) {
  812. chunks.push(chunk);
  813. totalLength += chunk.length;
  814. while (true) {
  815. if (totalLength < 4) break;
  816. const headerBuffer = chunks[0].length >= 4
  817. ? chunks[0]
  818. : Buffer.concat(chunks, 4);
  819. const recSize = headerBuffer.readInt32LE(0);
  820. if (recSize <= 4) {
  821. throw new Error(`Invalid record size read from stream: ${recSize}`);
  822. }
  823. if (totalLength < recSize) break;
  824. const record = Buffer.concat(chunks, recSize).subarray(0, recSize);
  825. if (!deletedSet.has(processedBytes)) {
  826. yield record;
  827. }
  828. [processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);
  829. }
  830. }
  831. }
  832. /**
  833. * Run a function exclusively: serialized against other operations in this
  834. * process (FIFO queue) and against other processes (lock file).
  835. *
  836. * Re-entrant only for operations called INSIDE the callback's async call
  837. * chain (find inside delete, insert inside withLock, ...). Operations
  838. * started concurrently from outside wait their turn.
  839. *
  840. * On acquiring the file lock, the on-disk meta is refreshed so writes in
  841. * other processes (fresh nextId, tombstones) are visible before mutating.
  842. * @private
  843. */
  844. async _withFileLock(operation, fn) {
  845. if (lockContext.getStore()?.owner === this) {
  846. // Causal re-entry: our call chain already holds this db's lock
  847. return fn();
  848. }
  849. // In-process FIFO queue — avoids the 100ms lock file polling between
  850. // concurrent operations of the same instance
  851. const prev = this._mutexTail;
  852. let releaseQueue;
  853. this._mutexTail = new Promise(r => releaseQueue = r);
  854. await prev;
  855. try {
  856. await this._acquireFileLock(operation);
  857. try {
  858. // Sync with other processes' writes before entering the critical
  859. // section — also during init (paths are set before any locked op)
  860. if (this._metaPath) await this.refresh();
  861. return await lockContext.run({ owner: this }, fn);
  862. } finally {
  863. await unlink(this._lockPath).catch(() => { });
  864. }
  865. } finally {
  866. releaseQueue();
  867. }
  868. }
  869. async _acquireFileLock(operation, retries = 0) {
  870. while (true) {
  871. try {
  872. await writeFile(this._lockPath, String(process.pid), { flag: 'wx' });
  873. return;
  874. } catch (e) {
  875. if (e.code !== 'EEXIST') throw e;
  876. // Stale lock recovery: a crashed holder never unlinks its lock file.
  877. // If the lock file is older than staleLockTimeout, take it over.
  878. if (this._staleLockTimeout > 0) {
  879. const lockStat = await stat(this._lockPath).catch(() => null);
  880. if (lockStat && Date.now() - lockStat.mtimeMs > this._staleLockTimeout) {
  881. const stalePath = `${this._lockPath}.stale-${process.pid}`;
  882. try {
  883. await rename(this._lockPath, stalePath);
  884. const staleStat = await stat(stalePath);
  885. if (Date.now() - staleStat.mtimeMs > this._staleLockTimeout) {
  886. console.warn(`[mpackdb] removed stale lock ${this._lockPath} (held > ${this._staleLockTimeout}ms, holder presumed dead)`);
  887. await unlink(stalePath).catch(() => { });
  888. } else {
  889. // raced a fresh lock — put it back
  890. await rename(stalePath, this._lockPath).catch(() => { });
  891. }
  892. } catch (e2) {
  893. // another waiter beat us to the takeover
  894. }
  895. continue;
  896. }
  897. }
  898. this.debug(`Waiting for lock file ${operation} retry #${retries}`);
  899. await new Promise(resolve => setTimeout(resolve, 25));
  900. retries++;
  901. }
  902. }
  903. }
  904. /**
  905. * Retrieves a document by its file offset and length
  906. * @param {Array} location - [offset, length] tuple
  907. * @returns {Promise<Object|null>} The deserialized document or null
  908. * @private
  909. * @deprecated Currently unused - may be removed in future versions
  910. */
  911. async _getDocByLocation([offset, length]) {
  912. if (!offset || length === 0) return null;
  913. const fileHandle = await open(this._dataPath, 'r');
  914. try {
  915. const buffer = Buffer.alloc(length);
  916. await fileHandle.read(buffer, 0, length, offset);
  917. return deserialize(buffer);
  918. } finally {
  919. await fileHandle.close();
  920. }
  921. }
  922. /**
  923. * Close the database and persist all pending changes
  924. * Should be called before process exit
  925. *
  926. * @returns {Promise<void>}
  927. */
  928. async close() {
  929. // Persist indexes before closing
  930. if (this._indexManager) {
  931. await this._indexManager.close();
  932. }
  933. // Close data stream
  934. if (this._dataStream) {
  935. await new Promise((resolve, reject) => {
  936. this._dataStream.end((err) => err ? reject(err) : resolve());
  937. });
  938. }
  939. // Remove signal handlers
  940. if (this._processExitHandler) {
  941. process.off('SIGINT', this._processExitHandler);
  942. process.off('SIGTERM', this._processExitHandler);
  943. this._processExitHandler = null;
  944. }
  945. }
  946. debug(...args) {
  947. if (this._debug) {
  948. console.log(...args);
  949. }
  950. }
  951. }
  952. /**
  953. * Helper function to remove a processed chunk from the buffer
  954. * @param {number} recSize - Size of the record to remove
  955. * @param {number} processedBytes - Current offset in the file
  956. * @param {number} totalLength - Total length of buffered data
  957. * @param {Buffer[]} chunks - Array of buffer chunks
  958. * @returns {[number, number]} Updated [processedBytes, totalLength]
  959. */
  960. function removeChunk(recSize, processedBytes, totalLength, chunks) {
  961. processedBytes += recSize;
  962. totalLength -= recSize;
  963. let bytesToRemove = recSize;
  964. while (bytesToRemove > 0 && chunks.length > 0) {
  965. const currentChunk = chunks[0];
  966. if (bytesToRemove >= currentChunk.length) {
  967. bytesToRemove -= currentChunk.length;
  968. chunks.shift();
  969. } else {
  970. chunks[0] = currentChunk.subarray(bytesToRemove);
  971. bytesToRemove = 0;
  972. }
  973. }
  974. return [processedBytes, totalLength];
  975. }
  976. export default MPackDB;

Branches

Latest commits

  • 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