gitoriaLog in with ident

mpackdb

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commit0afb8f4b0afb8f4bupdate now must be a callbackcaramboleyo0afb8f4b/src/MPackDB.js

26.9 KB

  1. import { dirname, basename, extname, resolve } from 'path';
  2. import { createReadStream, createWriteStream } from 'fs';
  3. import { mkdir, open, stat, readFile, writeFile, rename, unlink } from 'fs/promises';
  4. import { serialize, deserialize, uuid, PrimaryKeyType, IndexType } from './mpack.js';
  5. import { Cursor } from './Cursor.js';
  6. import { IndexManager } from './IndexManager.js';
  7. export { PrimaryKeyType, IndexType };
  8. /**
  9. * MPackDB - A fast, local, append-only JSON database with MessagePack serialization
  10. *
  11. * Features:
  12. * - Append-only writes for high performance
  13. * - MessagePack binary serialization
  14. * - Optional indexes (numeric and lexical)
  15. * - File-based locking for concurrent access
  16. * - Auto-compaction on startup
  17. * - Auto-persistence of indexes
  18. */
  19. export class MPackDB {
  20. _dbFile = null;
  21. _primaryKeyType = null;
  22. _primaryKey = null;
  23. _classToUse = null;
  24. _indexes = [];
  25. _indexTypes = {};
  26. _uniqueIndexes = new Set();
  27. _initPromise = null;
  28. _initState = 0;
  29. _dataPath = null;
  30. _dataStream = null;
  31. _lockPath = null;
  32. _lockDepth = 0;
  33. _indexManager = null;
  34. _meta = {
  35. nextId: 0,
  36. deleted: [],
  37. };
  38. _debug = false;
  39. _indexPersistInterval = 60000; // 60 seconds default
  40. _indexPersistThreshold = 1000; // 1000 entries default
  41. _processExitHandler = null;
  42. /**
  43. * Create a new MPackDB instance
  44. *
  45. * @param {string} dbFile - Path to the database file (without extension)
  46. * @param {Object} options - Configuration options
  47. * @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
  48. * @param {PrimaryKeyType} [options.primaryKeyType] - Explicit primary key type (overrides prefix)
  49. * @param {string[]} [options.indexes] - Array of field names to index. Use * prefix for numeric, @ for UUID
  50. * @param {boolean} [options.debug=false] - Enable debug logging
  51. * @param {number} [options.indexPersistInterval=60000] - Milliseconds between automatic index persistence (0 to disable)
  52. * @param {number} [options.indexPersistThreshold=1000] - Number of changes before auto-persisting indexes
  53. * @param {boolean} [options.compact=true] - Run compaction on init (set false for read-only / secondary instances)
  54. *
  55. * @example
  56. * const db = new MPackDB('data/users', {
  57. * primaryKey: '*id', // Numeric auto-increment
  58. * indexes: ['email', '*age'], // Index email (lexical) and age (numeric)
  59. * indexPersistThreshold: 100
  60. * });
  61. */
  62. _compact = true;
  63. constructor(dbFile, { primaryKey, primaryKeyType, indexes, debug, indexPersistInterval, indexPersistThreshold, compact } = {}) {
  64. this._dbFile = dbFile || this._dbFile;
  65. this._debug = debug || this._debug;
  66. if (compact !== undefined) this._compact = compact;
  67. if (indexPersistInterval !== undefined) this._indexPersistInterval = indexPersistInterval;
  68. if (indexPersistThreshold !== undefined) this._indexPersistThreshold = indexPersistThreshold;
  69. // Parse primary key with optional prefix
  70. if (primaryKey) {
  71. if (primaryKey.startsWith('*')) {
  72. // *id = numeric primary key
  73. this._primaryKey = primaryKey.slice(1);
  74. this._primaryKeyType = PrimaryKeyType.NUMBER;
  75. } else if (primaryKey.startsWith('@')) {
  76. // @id = UUID primary key
  77. this._primaryKey = primaryKey.slice(1);
  78. this._primaryKeyType = PrimaryKeyType.UUID;
  79. } else {
  80. // id = string primary key (lexical)
  81. this._primaryKey = primaryKey;
  82. this._primaryKeyType = PrimaryKeyType.STRING;
  83. }
  84. // Allow explicit override
  85. if (primaryKeyType !== undefined) {
  86. this._primaryKeyType = primaryKeyType;
  87. }
  88. }
  89. // Primary key is always unique
  90. if (this._primaryKey) {
  91. this._uniqueIndexes.add(this._primaryKey);
  92. }
  93. const allIndexes = [...new Set([
  94. ...(this._primaryKey ? [this._primaryKey] : []),
  95. ...(indexes || []),
  96. ...(this._indexes || []),
  97. ])];
  98. this._indexes = [];
  99. for (let index of allIndexes) {
  100. let cleanIndex = index;
  101. // !field = unique index (can combine with * and @: !*field, !@field)
  102. let isUnique = false;
  103. if (cleanIndex.startsWith('!')) {
  104. isUnique = true;
  105. cleanIndex = cleanIndex.slice(1);
  106. }
  107. if (cleanIndex.startsWith('*')) {
  108. // *field = numeric index
  109. cleanIndex = cleanIndex.slice(1);
  110. this._indexTypes[cleanIndex] = IndexType.NUMERIC;
  111. } else if (cleanIndex.startsWith('@')) {
  112. // @field = UUID index (lexical)
  113. cleanIndex = cleanIndex.slice(1);
  114. this._indexTypes[cleanIndex] = IndexType.LEXICAL;
  115. } else if (cleanIndex === this._primaryKey && this._primaryKeyType === PrimaryKeyType.NUMBER) {
  116. // Primary key is numeric
  117. this._indexTypes[cleanIndex] = IndexType.NUMERIC;
  118. } else {
  119. // Default to lexical
  120. this._indexTypes[cleanIndex] = IndexType.LEXICAL;
  121. }
  122. if (isUnique) {
  123. this._uniqueIndexes.add(cleanIndex);
  124. }
  125. this._indexes.push(cleanIndex);
  126. }
  127. }
  128. /**
  129. * Initialize the database (called automatically by other methods)
  130. * Performs compaction and loads metadata
  131. *
  132. * @returns {Promise<MPackDB>} The database instance
  133. */
  134. async init() {
  135. if (this._initState === 2) return this;
  136. if (this._initState === 1) {
  137. return this._initPromise;
  138. }
  139. this._initState = 1;
  140. return this._initPromise = new Promise(async success => {
  141. if (!this._dbFile) {
  142. throw new Error('No database file specified.');
  143. }
  144. const dbDir = dirname(this._dbFile);
  145. const baseName = basename(this._dbFile, extname(this._dbFile));
  146. await mkdir(dbDir, { recursive: true });
  147. this._dataPath = resolve(dbDir, `${baseName}.mpack`);
  148. this._metaPath = resolve(dbDir, `${baseName}.meta.json`);
  149. this._lockPath = resolve(dbDir, `${baseName}.lock`);
  150. try {
  151. this._meta = JSON.parse(await readFile(this._metaPath));
  152. } catch (e) {
  153. this._meta = {
  154. nextId: 0,
  155. deleted: [],
  156. };
  157. }
  158. this._initState = 2; // compact causes init to run again thats why we set it to 2 here already
  159. if (this._compact) await this.compact();
  160. this._dataStream = createWriteStream(this._dataPath, { flags: 'a' }); // needs to be set after compact as it replaces the file with a tmp file
  161. // Initialize IndexManager if indexes are specified
  162. if (this._indexes.length > 0) {
  163. const dbDir = dirname(this._dbFile);
  164. const baseName = basename(this._dbFile, extname(this._dbFile));
  165. this._indexManager = new IndexManager(dbDir, baseName, this._indexes, this._indexTypes, this._primaryKeyType);
  166. await this._indexManager.init(this._dataPath);
  167. // Start auto-persist with configured interval and threshold
  168. this._indexManager.startAutoPersist(this._indexPersistInterval, this._indexPersistThreshold);
  169. // Register signal handlers for Ctrl+C and kill signals
  170. this._processExitHandler = async () => {
  171. await this.close();
  172. process.exit(0);
  173. };
  174. process.once('SIGINT', this._processExitHandler);
  175. process.once('SIGTERM', this._processExitHandler);
  176. }
  177. this._initPromise = null;
  178. success(this);
  179. });
  180. }
  181. /**
  182. * Insert a new record into the database
  183. *
  184. * @param {Object} record - The record to insert
  185. * @param {Object} [options]
  186. * @param {boolean} [options.skipPrimaryKey=false] - Skip auto-generating primary key
  187. * @returns {Promise<Object>} The inserted record with primary key
  188. *
  189. * @example
  190. * await db.insert({ name: 'Alice', age: 30 });
  191. * // Returns: { id: 0, name: 'Alice', age: 30 }
  192. */
  193. async insert(record, { skipPrimaryKey = false } = {}) {
  194. await this.init();
  195. await this._acquireLock('insert', record);
  196. try {
  197. const recToInsert = this._classToUse
  198. ? Object.assign(new this._classToUse(), record)
  199. : { ...record };
  200. if (this._primaryKey && !skipPrimaryKey && !recToInsert[this._primaryKey]) {
  201. if (this._primaryKeyType === PrimaryKeyType.NUMBER) {
  202. recToInsert[this._primaryKey] = this._meta.nextId++;
  203. } else if (this._primaryKeyType === PrimaryKeyType.UUID) {
  204. recToInsert[this._primaryKey] = uuid();
  205. }
  206. }
  207. // Check unique index constraints (skip auto-generated primary keys — guaranteed unique)
  208. if (this._indexManager && this._uniqueIndexes.size > 0) {
  209. const autoGenPK = this._primaryKey && !skipPrimaryKey && !record[this._primaryKey]
  210. && (this._primaryKeyType === PrimaryKeyType.NUMBER || this._primaryKeyType === PrimaryKeyType.UUID);
  211. const deletedSet = new Set(this._meta.deleted);
  212. for (const field of this._uniqueIndexes) {
  213. if (autoGenPK && field === this._primaryKey) continue;
  214. const value = recToInsert[field];
  215. if (value === undefined) continue;
  216. const existing = await this._indexManager.get(field, value);
  217. const live = existing.filter(e => !deletedSet.has(e.loc[0]));
  218. if (live.length > 0) {
  219. const err = new Error(`Duplicate key: ${field}=${value}`);
  220. err.code = 'DUPLICATE_KEY';
  221. err.field = field;
  222. err.value = value;
  223. throw err;
  224. }
  225. }
  226. }
  227. const packedBuffer = serialize(recToInsert);
  228. const fileStat = await stat(this._dataPath).catch(() => ({ size: 0 }));
  229. const offset = fileStat.size;
  230. await new Promise(resolve => this._dataStream.write(packedBuffer, resolve));
  231. const loc = [offset, packedBuffer.length];
  232. // indexes
  233. if (this._indexManager) {
  234. this._indexManager.insert(recToInsert, loc);
  235. }
  236. // Persist meta if we incremented nextId
  237. if (this._primaryKey && this._primaryKeyType === PrimaryKeyType.NUMBER && !skipPrimaryKey && !record[this._primaryKey]) {
  238. await this.persistMeta();
  239. }
  240. return this._primaryKey ? recToInsert[this._primaryKey] : recToInsert;
  241. } finally {
  242. await this._releaseLock();
  243. }
  244. }
  245. /**
  246. * Update records matching a query.
  247. * The callback receives the old record and must return the new record.
  248. *
  249. * @param {string|number|Function} mixed - Primary key value or query function
  250. * @param {Function} callback - Receives old record, must return new record
  251. * @param {Object} [options]
  252. * @param {boolean} [options.upsert=false] - Insert if no records match
  253. * @param {Object|Array} [options.index] - Index hint(s) to narrow disk reads
  254. * @returns {Promise<Object[]>} Array of updated records
  255. *
  256. * @example
  257. * await db.update(0, record => {
  258. * record.age = 32;
  259. * delete record.address;
  260. * return record;
  261. * });
  262. */
  263. async update(mixed, callback, { upsert = false, index } = {}) {
  264. if (typeof callback !== 'function') {
  265. throw new Error('update requires a callback function as second argument');
  266. }
  267. let insertedRecords = await this.delete(mixed, {
  268. index,
  269. callback: async record => {
  270. return this.insert(await callback(record), { skipPrimaryKey: true });
  271. }
  272. });
  273. if (upsert && insertedRecords.length === 0) {
  274. insertedRecords = await this.insert(await callback({}));
  275. }
  276. return insertedRecords;
  277. }
  278. /**
  279. * Update records or insert if not found
  280. *
  281. * @param {string|number|Function} mixed - Primary key value or query function
  282. * @param {Function} callback - Receives old record (or {} if inserting), must return new record
  283. * @param {Object} [options]
  284. * @param {Object|Array} [options.index] - Index hint(s) to narrow disk reads
  285. * @returns {Promise<Object[]>} Array of updated/inserted records
  286. */
  287. async upsert(mixed, callback, { index } = {}) {
  288. return this.update(mixed, callback, { upsert: true, index });
  289. }
  290. /**
  291. * Delete records matching a query
  292. *
  293. * @param {string|Function} mixed - Primary key value or query function
  294. * @param {Object} [options]
  295. * @param {Object|Array} [options.index] - Index hint(s) to narrow disk reads
  296. * @param {Function} [options.callback] - Internal callback per deleted record (used by update)
  297. * @returns {Promise<Object[]>} Array of deleted records
  298. *
  299. * @example
  300. * // Delete by primary key
  301. * await db.delete(0);
  302. *
  303. * // Delete with query function
  304. * await db.delete(r => r.age < 18);
  305. *
  306. * // Delete with index hint
  307. * await db.delete(r => r.status === 'inactive', { index: { field: 'status', value: 'inactive' } });
  308. */
  309. async delete(mixed, { index, callback = async record => record } = {}) {
  310. await this.init();
  311. await this._acquireLock('delete', mixed);
  312. try {
  313. const promises = [];
  314. for await (const [record, offset] of this.find(mixed, { mode: 'mixed', index })) {
  315. this._meta.deleted.push(offset);
  316. if (this._indexManager) {
  317. await this._indexManager.remove(record, this._primaryKey);
  318. }
  319. promises.push(callback(record));
  320. }
  321. await this.persistMeta();
  322. return Promise.all(promises);
  323. } finally {
  324. await this._releaseLock();
  325. }
  326. }
  327. /**
  328. * Finds records in the database
  329. * @param {undefined|string|number|function} [mixed] - Query: undefined/null for all, primary key value for PK lookup, function for filter
  330. * @param {Object} [options] - Query options
  331. * @param {Object|Array} [options.index] - Index hint(s) to narrow disk reads before filtering
  332. * @returns {Cursor} A cursor for iterating over results
  333. * @example
  334. * // All records
  335. * for await (const user of db.find()) { ... }
  336. *
  337. * // Primary key lookup
  338. * const [user] = await db.find(68);
  339. *
  340. * // Filter with index range — only reads records in x 100+
  341. * for await (const doc of db.find(r => r.x <= 200, { index: { field: 'x', from: 100 } })) { ... }
  342. *
  343. * // Index intersection — intersects offsets first, then streams matches
  344. * for await (const doc of db.find(r => r.x <= 200, {
  345. * index: [
  346. * { field: 'x', from: 100 },
  347. * { field: 'status', value: 'active' }
  348. * ]
  349. * })) { ... }
  350. */
  351. find(mixed, options = {}) {
  352. if (typeof mixed === 'undefined' || mixed === null || typeof mixed === 'function') {
  353. return new Cursor(this, mixed, options);
  354. } else {
  355. if (!this._primaryKey) {
  356. throw new Error('No primary key specified.');
  357. }
  358. // Use indexed lookup when available
  359. if (this._indexManager) {
  360. return new Cursor(this, null, {
  361. ...options,
  362. _indexLookup: { field: this._primaryKey, value: mixed }
  363. });
  364. }
  365. return new Cursor(this, record => {
  366. return record[this._primaryKey] === mixed;
  367. }, options);
  368. }
  369. }
  370. /**
  371. * Find records within a bounding box defined by 4 corners.
  372. * Requires numeric indexes on x and y fields.
  373. * Corners can be in any order — min/max are extracted automatically.
  374. *
  375. * @param {Array<{x: number, y: number}>} corners - 4 corner coordinates
  376. * @param {function} [filter] - Optional additional filter function
  377. * @returns {Cursor}
  378. * @example
  379. * const results = await db.boundingBox([
  380. * { x: -5, y: 5 }, { x: 5, y: 5 },
  381. * { x: 5, y: -5 }, { x: -5, y: -5 }
  382. * ]);
  383. */
  384. boundingBox(corners, filter) {
  385. const xs = corners.map(c => c.x);
  386. const ys = corners.map(c => c.y);
  387. const minX = Math.min(...xs), maxX = Math.max(...xs);
  388. const minY = Math.min(...ys), maxY = Math.max(...ys);
  389. return this.find(filter || null, {
  390. index: [
  391. { field: 'x', from: minX, to: maxX },
  392. { field: 'y', from: minY, to: maxY }
  393. ]
  394. });
  395. }
  396. /**
  397. * Execute a callback while holding the database lock.
  398. * The lock is re-entrant: find/insert/delete/update called inside
  399. * the callback reuse the same lock instead of deadlocking.
  400. *
  401. * Use this for compound operations that must be atomic, e.g.
  402. * find-then-insert (login pattern).
  403. *
  404. * @param {Function} callback - Async function to execute under lock
  405. * @returns {Promise<any>} The return value of the callback
  406. *
  407. * @example
  408. * const user = await db.withLock(async () => {
  409. * const [existing] = await db.find(u => u.email === email);
  410. * if (existing) return existing;
  411. * return db.insert({ email });
  412. * });
  413. */
  414. async withLock(callback) {
  415. await this.init();
  416. await this._acquireLock('withLock');
  417. try {
  418. return await callback();
  419. } finally {
  420. await this._releaseLock();
  421. }
  422. }
  423. /**
  424. * Generator that yields records from the database file
  425. * @param {function|null} [queryFn=null] - Optional filter function
  426. * @param {Object} [options] - Generator options
  427. * @param {string} [options.mode='record'] - Mode: 'record', 'raw', 'offset', or 'mixed'
  428. * @yields {Object|Buffer|Array} Records, buffers, or [record, offset, size] tuples depending on mode
  429. */
  430. async *recordGenerator(queryFn = null, { mode = 'record', _indexLookup, index } = {}) {
  431. await this.init();
  432. await this.refresh();
  433. // Indexed primary key lookup — O(log n) instead of full scan
  434. if (_indexLookup && this._indexManager) {
  435. yield* this._indexedLookup(_indexLookup.field, _indexLookup.value, mode);
  436. return;
  437. }
  438. // Index-based find: collect offsets from index(es), then stream only those records
  439. if (index && this._indexManager) {
  440. yield* this._indexedStream(index, queryFn, mode);
  441. return;
  442. }
  443. // Flush pending writes so reads see all inserted data
  444. if (this._dataStream && this._dataStream.writableLength > 0) {
  445. await new Promise(resolve => this._dataStream.once('drain', resolve));
  446. }
  447. // Check if the data file exists before attempting to read it
  448. try {
  449. await stat(this._dataPath);
  450. } catch (error) {
  451. // File doesn't exist - return empty generator (no records)
  452. return;
  453. }
  454. // Convert deleted array to Set for O(1) lookup instead of O(n)
  455. const deletedSet = new Set(this._meta.deleted);
  456. const readStream = createReadStream(this._dataPath);
  457. let chunks = [];
  458. let totalLength = 0;
  459. let processedBytes = 0;
  460. for await (const chunk of readStream) {
  461. chunks.push(chunk);
  462. totalLength += chunk.length;
  463. while (true) {
  464. // Exit 1: Not enough data to even read the 4-byte size header.
  465. if (totalLength < 4) {
  466. break;
  467. }
  468. // Safely read the header, even if it's split across chunks
  469. let headerBuffer;
  470. if (chunks[0].length >= 4) {
  471. headerBuffer = chunks[0];
  472. } else {
  473. // The header is fragmented, so we must concat just enough to read it.
  474. headerBuffer = Buffer.concat(chunks, 4);
  475. }
  476. const recSize = headerBuffer.readInt32LE(0);
  477. // Validate the record size to prevent infinite loops
  478. // A record must be at least as large as its header (4 bytes).
  479. // A size of 0 or less is invalid and indicates corruption.
  480. if (recSize <= 4) {
  481. throw new Error(`Invalid record size read from stream: ${recSize}`);
  482. }
  483. // Exit 2: We have the size, but not the full record yet.
  484. if (totalLength < recSize) {
  485. break;
  486. }
  487. // skip deleted records (only after we have the full record)
  488. if (deletedSet.has(processedBytes)) {
  489. [processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);
  490. continue; // Goes back to while (true)
  491. }
  492. switch (mode) {
  493. case 'raw':
  494. yield Buffer.concat(chunks, recSize).subarray(0, recSize);
  495. break;
  496. case 'mixed':
  497. case 'record':
  498. const recBuffer = Buffer.concat(chunks, recSize).subarray(0, recSize);
  499. // Skip the 4-byte size header
  500. const data = deserialize(recBuffer.subarray(4));
  501. const rec = this._classToUse
  502. ? Object.assign(new this._classToUse(), data)
  503. : data;
  504. if (queryFn) {
  505. if (queryFn(rec)) {
  506. yield mode === 'mixed' ? [rec, processedBytes, recSize] : rec;
  507. }
  508. } else {
  509. yield mode === 'mixed' ? [rec, processedBytes, recSize] : rec;
  510. }
  511. break;
  512. case 'offset':
  513. // if someone uses a queryFn with offset, we have to deserialize it first to perform the queryFn
  514. if (queryFn) {
  515. const recBuffer = Buffer.concat(chunks, recSize).subarray(0, recSize);
  516. const rec = deserialize(recBuffer.subarray(4));
  517. if (queryFn(rec)) {
  518. yield [processedBytes, recSize];
  519. }
  520. } else {
  521. yield [processedBytes, recSize];
  522. }
  523. break;
  524. default:
  525. throw new Error(`Invalid mode: ${mode}`);
  526. }
  527. [processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);
  528. }
  529. }
  530. }
  531. /**
  532. * Indexed lookup — reads records directly by offset from the index.
  533. * Uses binary search on the index file for O(log n) lookups.
  534. * @private
  535. */
  536. async *_indexedLookup(field, value, mode = 'record') {
  537. const entries = await this._indexManager.get(field, value);
  538. if (entries.length === 0) return;
  539. const deletedSet = new Set(this._meta.deleted);
  540. const fileHandle = await open(this._dataPath, 'r');
  541. try {
  542. for (const entry of entries) {
  543. const [offset, length] = entry.loc;
  544. if (deletedSet.has(offset)) continue;
  545. const buffer = Buffer.alloc(length);
  546. await fileHandle.read(buffer, 0, length, offset);
  547. switch (mode) {
  548. case 'raw':
  549. yield buffer;
  550. break;
  551. case 'mixed': {
  552. const data = deserialize(buffer.subarray(4));
  553. const rec = this._classToUse
  554. ? Object.assign(new this._classToUse(), data)
  555. : data;
  556. yield [rec, offset, length];
  557. break;
  558. }
  559. case 'record':
  560. default: {
  561. const data = deserialize(buffer.subarray(4));
  562. const rec = this._classToUse
  563. ? Object.assign(new this._classToUse(), data)
  564. : data;
  565. yield rec;
  566. break;
  567. }
  568. }
  569. }
  570. } finally {
  571. await fileHandle.close();
  572. }
  573. }
  574. /**
  575. * Collect offsets from one or more indexes, optionally intersect, then stream records.
  576. * @param {Object|Array} index - Single index hint or array of hints
  577. * @param {function|null} queryFn - Optional filter function
  578. * @param {string} mode - Output mode: 'record', 'raw', or 'mixed'
  579. * @private
  580. */
  581. async *_indexedStream(index, queryFn, mode = 'record') {
  582. const hints = Array.isArray(index) ? index : [index];
  583. const deletedSet = new Set(this._meta.deleted);
  584. // Collect offset sets from each index hint
  585. const offsetSets = [];
  586. for (const hint of hints) {
  587. const offsets = new Map(); // offset → [offset, length]
  588. if (hint.value !== undefined) {
  589. // Exact match via get()
  590. const entries = await this._indexManager.get(hint.field, hint.value);
  591. for (const e of entries) {
  592. if (!deletedSet.has(e.loc[0])) offsets.set(e.loc[0], e.loc);
  593. }
  594. } else {
  595. // Range scan via entries()
  596. for await (const e of this._indexManager.entries(hint.field, { from: hint.from, to: hint.to, direction: hint.direction })) {
  597. if (!deletedSet.has(e.loc[0])) offsets.set(e.loc[0], e.loc);
  598. }
  599. }
  600. offsetSets.push(offsets);
  601. }
  602. // Intersect: keep only offsets present in ALL sets
  603. let locations;
  604. if (offsetSets.length === 1) {
  605. locations = Array.from(offsetSets[0].values());
  606. } else {
  607. // Start with smallest set for efficiency
  608. offsetSets.sort((a, b) => a.size - b.size);
  609. const [smallest, ...rest] = offsetSets;
  610. locations = [];
  611. for (const [offset, loc] of smallest) {
  612. if (rest.every(s => s.has(offset))) locations.push(loc);
  613. }
  614. }
  615. // Stream records from the intersected locations
  616. const fileHandle = await open(this._dataPath, 'r');
  617. try {
  618. for (const [offset, length] of locations) {
  619. const buffer = Buffer.alloc(length);
  620. await fileHandle.read(buffer, 0, length, offset);
  621. if (mode === 'raw') {
  622. yield buffer;
  623. } else {
  624. const data = deserialize(buffer.subarray(4));
  625. const rec = this._classToUse
  626. ? Object.assign(new this._classToUse(), data)
  627. : data;
  628. if (queryFn && !queryFn(rec)) continue;
  629. if (mode === 'mixed') {
  630. yield [rec, offset, length];
  631. } else {
  632. yield rec;
  633. }
  634. }
  635. }
  636. } finally {
  637. await fileHandle.close();
  638. }
  639. }
  640. /**
  641. * Re-read meta.json from disk so this instance sees changes made by other processes.
  642. * Called automatically before every read operation.
  643. */
  644. async refresh() {
  645. try {
  646. this._meta = JSON.parse(await readFile(this._metaPath));
  647. } catch (e) {
  648. // file missing or corrupt — keep current meta
  649. }
  650. }
  651. async persistMeta() {
  652. const metaToSave = { ...this._meta };
  653. // Only save nextId if we have a numeric primary key
  654. if (this._primaryKeyType !== PrimaryKeyType.NUMBER) {
  655. delete metaToSave.nextId;
  656. }
  657. await writeFile(this._metaPath, JSON.stringify(metaToSave));
  658. }
  659. /**
  660. * Compact the database by removing deleted records
  661. * This rewrites the data file without tombstones
  662. *
  663. * @returns {Promise<void>}
  664. */
  665. async compact() {
  666. await this._acquireLock('compact');
  667. try {
  668. const writeStream = createWriteStream(this._dataPath + '.tmp', { flags: 'w' });
  669. for await (const binary of this.find(null, { mode: 'raw' })) {
  670. writeStream.write(binary);
  671. }
  672. await new Promise((resolve, reject) => {
  673. writeStream.end((err) => err ? reject(err) : resolve());
  674. });
  675. // rename is atomic, so we can just rename the file and it will replace the old one
  676. await rename(this._dataPath + '.tmp', this._dataPath);
  677. this._meta.deleted = [];
  678. await this.persistMeta();
  679. // Rebuild indexes — offsets changed after compaction
  680. if (this._indexManager) {
  681. await this._indexManager.init(this._dataPath, { forceRebuild: true });
  682. }
  683. } finally {
  684. await this._releaseLock();
  685. }
  686. }
  687. async _acquireLock(operation, record, retries = 0) {
  688. if (this._lockDepth > 0) {
  689. this._lockDepth++;
  690. return;
  691. }
  692. try {
  693. await writeFile(this._lockPath, String(process.pid), { flag: 'wx' });
  694. this._lockDepth = 1;
  695. } catch (e) {
  696. if (e.code === 'EEXIST') {
  697. this.debug(`Waiting for lock file ${operation} retry #${retries}`, record);
  698. await new Promise(resolve => setTimeout(resolve, 100));
  699. return this._acquireLock(operation, record, retries + 1);
  700. }
  701. throw e;
  702. }
  703. }
  704. async _releaseLock() {
  705. this._lockDepth--;
  706. if (this._lockDepth <= 0) {
  707. this._lockDepth = 0;
  708. await unlink(this._lockPath).catch(() => { });
  709. }
  710. }
  711. /**
  712. * Retrieves a document by its file offset and length
  713. * @param {Array} location - [offset, length] tuple
  714. * @returns {Promise<Object|null>} The deserialized document or null
  715. * @private
  716. * @deprecated Currently unused - may be removed in future versions
  717. */
  718. async _getDocByLocation([offset, length]) {
  719. if (!offset || length === 0) return null;
  720. const fileHandle = await open(this._dataPath, 'r');
  721. try {
  722. const buffer = Buffer.alloc(length);
  723. await fileHandle.read(buffer, 0, length, offset);
  724. return deserialize(buffer);
  725. } finally {
  726. await fileHandle.close();
  727. }
  728. }
  729. /**
  730. * Close the database and persist all pending changes
  731. * Should be called before process exit
  732. *
  733. * @returns {Promise<void>}
  734. */
  735. async close() {
  736. // Persist indexes before closing
  737. if (this._indexManager) {
  738. await this._indexManager.close();
  739. }
  740. // Close data stream
  741. if (this._dataStream) {
  742. await new Promise((resolve, reject) => {
  743. this._dataStream.end((err) => err ? reject(err) : resolve());
  744. });
  745. }
  746. // Remove signal handlers
  747. if (this._processExitHandler) {
  748. process.off('SIGINT', this._processExitHandler);
  749. process.off('SIGTERM', this._processExitHandler);
  750. this._processExitHandler = null;
  751. }
  752. }
  753. debug(...args) {
  754. if (this._debug) {
  755. console.log(...args);
  756. }
  757. }
  758. }
  759. /**
  760. * Helper function to remove a processed chunk from the buffer
  761. * @param {number} recSize - Size of the record to remove
  762. * @param {number} processedBytes - Current offset in the file
  763. * @param {number} totalLength - Total length of buffered data
  764. * @param {Buffer[]} chunks - Array of buffer chunks
  765. * @returns {[number, number]} Updated [processedBytes, totalLength]
  766. */
  767. function removeChunk(recSize, processedBytes, totalLength, chunks) {
  768. processedBytes += recSize;
  769. totalLength -= recSize;
  770. let bytesToRemove = recSize;
  771. while (bytesToRemove > 0 && chunks.length > 0) {
  772. const currentChunk = chunks[0];
  773. if (bytesToRemove >= currentChunk.length) {
  774. bytesToRemove -= currentChunk.length;
  775. chunks.shift();
  776. } else {
  777. chunks[0] = currentChunk.subarray(bytesToRemove);
  778. bytesToRemove = 0;
  779. }
  780. }
  781. return [processedBytes, totalLength];
  782. }
  783. export default MPackDB;

Branches

Latest commits

  • 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