gitoriaLog in with ident

mpackdb

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commitb8ffc1a0b8ffc1a0release 1.0.5caramboleyob8ffc1a0/src/MPackDB.js

24.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. *
  248. * @param {string|Function} mixed - Primary key value or query function
  249. * @param {Object|Function} dataOrCallback - Data to update or callback function
  250. * @param {Object} [options]
  251. * @param {boolean} [options.upsert=false] - Insert if no records match
  252. * @returns {Promise<Object[]>} Array of updated records
  253. *
  254. * @example
  255. * // Update by primary key
  256. * await db.update(0, { age: 31 });
  257. *
  258. * // Update with query function
  259. * await db.update(r => r.age > 30, { status: 'senior' });
  260. *
  261. * // Update with callback
  262. * await db.update(r => r.age > 30, r => ({ ...r, age: r.age + 1 }));
  263. */
  264. async update(mixed, dataOrCallback, { upsert = false } = {}) {
  265. let insertedRecords = await this.delete(mixed, async record => {
  266. return this.insert(typeof dataOrCallback === 'function'
  267. ? await dataOrCallback(record)
  268. : dataOrCallback,
  269. { skipPrimaryKey: true }
  270. );
  271. });
  272. if (upsert && insertedRecords.length === 0) {
  273. insertedRecords = await this.insert(typeof dataOrCallback === 'function'
  274. ? await dataOrCallback({})
  275. : dataOrCallback);
  276. }
  277. return insertedRecords;
  278. }
  279. /**
  280. * Update records or insert if not found
  281. *
  282. * @param {string|Function} mixed - Primary key value or query function
  283. * @param {Object|Function} dataOrCallback - Data to update/insert or callback function
  284. * @returns {Promise<Object[]>} Array of updated/inserted records
  285. */
  286. async upsert(mixed, dataOrCallback) {
  287. return this.update(mixed, dataOrCallback, { upsert: true });
  288. }
  289. /**
  290. * Delete records matching a query
  291. *
  292. * @param {string|Function} mixed - Primary key value or query function
  293. * @param {Function} [callback] - Optional callback to execute for each deleted record
  294. * @returns {Promise<Object[]>} Array of deleted records
  295. *
  296. * @example
  297. * // Delete by primary key
  298. * await db.delete(0);
  299. *
  300. * // Delete with query function
  301. * await db.delete(r => r.age < 18);
  302. */
  303. async delete(mixed, callback = async record => record) {
  304. await this.init();
  305. await this._acquireLock('delete', mixed);
  306. try {
  307. const promises = [];
  308. for await (const [record, offset] of this.find(mixed, { mode: 'mixed' })) {
  309. this._meta.deleted.push(offset);
  310. if (this._indexManager) {
  311. await this._indexManager.remove(record, this._primaryKey);
  312. }
  313. promises.push(callback(record));
  314. }
  315. await this.persistMeta();
  316. return Promise.all(promises);
  317. } finally {
  318. await this._releaseLock();
  319. }
  320. }
  321. /**
  322. * Finds records in the database
  323. * @param {undefined|string|function} [mixed] - Query: undefined/null for all records, string for primary key lookup, function for filter
  324. * @param {Object} [options] - Query options
  325. * @returns {Cursor} A cursor for iterating over results
  326. * @example
  327. * // Get all records
  328. * for await (const user of db.find()) {
  329. * console.log(user.name);
  330. * }
  331. */
  332. find(mixed, options = {}) {
  333. if (typeof mixed === 'undefined' || mixed === null || typeof mixed === 'function') {
  334. return new Cursor(this, mixed, options);
  335. } else {
  336. if (!this._primaryKey) {
  337. throw new Error('No primary key specified.');
  338. }
  339. // Use indexed lookup when available
  340. if (this._indexManager) {
  341. return new Cursor(this, null, {
  342. ...options,
  343. _indexLookup: { field: this._primaryKey, value: mixed }
  344. });
  345. }
  346. return new Cursor(this, record => {
  347. return record[this._primaryKey] === mixed;
  348. }, options);
  349. }
  350. }
  351. /**
  352. * Walk an index in sorted order with optional range and filter.
  353. * Returns a Cursor (async iterable + awaitable).
  354. *
  355. * @param {string} field - Indexed field name
  356. * @param {Object} [options]
  357. * @param {any} [options.from] - Start key (inclusive), uses binary search to skip ahead
  358. * @param {number} [options.limit] - Max records to return
  359. * @param {Function} [options.filter] - Filter function applied per record during iteration (memory efficient)
  360. * @returns {Cursor}
  361. *
  362. * @example
  363. * // All users sorted by age
  364. * for await (const user of db.findByIndex('age')) { ... }
  365. *
  366. * // Age 18+
  367. * for await (const user of db.findByIndex('age', { from: 18 })) { ... }
  368. *
  369. * // Age 18+, only admins, max 10
  370. * for await (const user of db.findByIndex('age', { from: 18, filter: u => u.role === 'admin', limit: 10 })) { ... }
  371. */
  372. findByIndex(field, options = {}) {
  373. return new Cursor(this, null, { ...options, _indexWalk: { field } });
  374. }
  375. /**
  376. * Execute a callback while holding the database lock.
  377. * The lock is re-entrant: find/insert/delete/update called inside
  378. * the callback reuse the same lock instead of deadlocking.
  379. *
  380. * Use this for compound operations that must be atomic, e.g.
  381. * find-then-insert (login pattern).
  382. *
  383. * @param {Function} callback - Async function to execute under lock
  384. * @returns {Promise<any>} The return value of the callback
  385. *
  386. * @example
  387. * const user = await db.withLock(async () => {
  388. * const [existing] = await db.find(u => u.email === email);
  389. * if (existing) return existing;
  390. * return db.insert({ email });
  391. * });
  392. */
  393. async withLock(callback) {
  394. await this.init();
  395. await this._acquireLock('withLock');
  396. try {
  397. return await callback();
  398. } finally {
  399. await this._releaseLock();
  400. }
  401. }
  402. /**
  403. * Generator that yields records from the database file
  404. * @param {function|null} [queryFn=null] - Optional filter function
  405. * @param {Object} [options] - Generator options
  406. * @param {string} [options.mode='record'] - Mode: 'record', 'raw', 'offset', or 'mixed'
  407. * @yields {Object|Buffer|Array} Records, buffers, or [record, offset, size] tuples depending on mode
  408. */
  409. async *recordGenerator(queryFn = null, { mode = 'record', _indexLookup, _indexWalk, from, limit, filter } = {}) {
  410. await this.init();
  411. await this.refresh();
  412. // Indexed primary key lookup — O(log n) instead of full scan
  413. if (_indexLookup && this._indexManager) {
  414. yield* this._indexedLookup(_indexLookup.field, _indexLookup.value, mode);
  415. return;
  416. }
  417. // Sorted index walk with optional range/filter
  418. if (_indexWalk && this._indexManager) {
  419. yield* this._indexWalk(_indexWalk.field, { from, limit, filter, mode });
  420. return;
  421. }
  422. // Flush pending writes so reads see all inserted data
  423. if (this._dataStream && this._dataStream.writableLength > 0) {
  424. await new Promise(resolve => this._dataStream.once('drain', resolve));
  425. }
  426. // Check if the data file exists before attempting to read it
  427. try {
  428. await stat(this._dataPath);
  429. } catch (error) {
  430. // File doesn't exist - return empty generator (no records)
  431. return;
  432. }
  433. // Convert deleted array to Set for O(1) lookup instead of O(n)
  434. const deletedSet = new Set(this._meta.deleted);
  435. const readStream = createReadStream(this._dataPath);
  436. let chunks = [];
  437. let totalLength = 0;
  438. let processedBytes = 0;
  439. for await (const chunk of readStream) {
  440. chunks.push(chunk);
  441. totalLength += chunk.length;
  442. while (true) {
  443. // Exit 1: Not enough data to even read the 4-byte size header.
  444. if (totalLength < 4) {
  445. break;
  446. }
  447. // Safely read the header, even if it's split across chunks
  448. let headerBuffer;
  449. if (chunks[0].length >= 4) {
  450. headerBuffer = chunks[0];
  451. } else {
  452. // The header is fragmented, so we must concat just enough to read it.
  453. headerBuffer = Buffer.concat(chunks, 4);
  454. }
  455. const recSize = headerBuffer.readInt32LE(0);
  456. // Validate the record size to prevent infinite loops
  457. // A record must be at least as large as its header (4 bytes).
  458. // A size of 0 or less is invalid and indicates corruption.
  459. if (recSize <= 4) {
  460. throw new Error(`Invalid record size read from stream: ${recSize}`);
  461. }
  462. // Exit 2: We have the size, but not the full record yet.
  463. if (totalLength < recSize) {
  464. break;
  465. }
  466. // skip deleted records (only after we have the full record)
  467. if (deletedSet.has(processedBytes)) {
  468. [processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);
  469. continue; // Goes back to while (true)
  470. }
  471. switch (mode) {
  472. case 'raw':
  473. yield Buffer.concat(chunks, recSize).subarray(0, recSize);
  474. break;
  475. case 'mixed':
  476. case 'record':
  477. const recBuffer = Buffer.concat(chunks, recSize).subarray(0, recSize);
  478. // Skip the 4-byte size header
  479. const data = deserialize(recBuffer.subarray(4));
  480. const rec = this._classToUse
  481. ? Object.assign(new this._classToUse(), data)
  482. : data;
  483. if (queryFn) {
  484. if (queryFn(rec)) {
  485. yield mode === 'mixed' ? [rec, processedBytes, recSize] : rec;
  486. }
  487. } else {
  488. yield mode === 'mixed' ? [rec, processedBytes, recSize] : rec;
  489. }
  490. break;
  491. case 'offset':
  492. // if someone uses a queryFn with offset, we have to deserialize it first to perform the queryFn
  493. if (queryFn) {
  494. const recBuffer = Buffer.concat(chunks, recSize).subarray(0, recSize);
  495. const rec = deserialize(recBuffer.subarray(4));
  496. if (queryFn(rec)) {
  497. yield [processedBytes, recSize];
  498. }
  499. } else {
  500. yield [processedBytes, recSize];
  501. }
  502. break;
  503. default:
  504. throw new Error(`Invalid mode: ${mode}`);
  505. }
  506. [processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);
  507. }
  508. }
  509. }
  510. /**
  511. * Indexed lookup — reads records directly by offset from the index.
  512. * Uses binary search on the index file for O(log n) lookups.
  513. * @private
  514. */
  515. async *_indexedLookup(field, value, mode = 'record') {
  516. const entries = await this._indexManager.get(field, value);
  517. if (entries.length === 0) return;
  518. const deletedSet = new Set(this._meta.deleted);
  519. const fileHandle = await open(this._dataPath, 'r');
  520. try {
  521. for (const entry of entries) {
  522. const [offset, length] = entry.loc;
  523. if (deletedSet.has(offset)) continue;
  524. const buffer = Buffer.alloc(length);
  525. await fileHandle.read(buffer, 0, length, offset);
  526. switch (mode) {
  527. case 'raw':
  528. yield buffer;
  529. break;
  530. case 'mixed': {
  531. const data = deserialize(buffer.subarray(4));
  532. const rec = this._classToUse
  533. ? Object.assign(new this._classToUse(), data)
  534. : data;
  535. yield [rec, offset, length];
  536. break;
  537. }
  538. case 'record':
  539. default: {
  540. const data = deserialize(buffer.subarray(4));
  541. const rec = this._classToUse
  542. ? Object.assign(new this._classToUse(), data)
  543. : data;
  544. yield rec;
  545. break;
  546. }
  547. }
  548. }
  549. } finally {
  550. await fileHandle.close();
  551. }
  552. }
  553. /**
  554. * Walk an index in sorted order, reading records by offset.
  555. * @private
  556. */
  557. async *_indexWalk(field, { from, limit, filter, mode = 'record' } = {}) {
  558. const deletedSet = new Set(this._meta.deleted);
  559. const fileHandle = await open(this._dataPath, 'r');
  560. let count = 0;
  561. try {
  562. for await (const entry of this._indexManager.entries(field, { from })) {
  563. const [offset, length] = entry.loc;
  564. if (deletedSet.has(offset)) continue;
  565. const buffer = Buffer.alloc(length);
  566. await fileHandle.read(buffer, 0, length, offset);
  567. if (mode === 'raw') {
  568. yield buffer;
  569. } else {
  570. const data = deserialize(buffer.subarray(4));
  571. const rec = this._classToUse
  572. ? Object.assign(new this._classToUse(), data)
  573. : data;
  574. if (filter && !filter(rec)) continue;
  575. if (mode === 'mixed') {
  576. yield [rec, offset, length];
  577. } else {
  578. yield rec;
  579. }
  580. }
  581. count++;
  582. if (limit !== undefined && count >= limit) return;
  583. }
  584. } finally {
  585. await fileHandle.close();
  586. }
  587. }
  588. /**
  589. * Re-read meta.json from disk so this instance sees changes made by other processes.
  590. * Called automatically before every read operation.
  591. */
  592. async refresh() {
  593. try {
  594. this._meta = JSON.parse(await readFile(this._metaPath));
  595. } catch (e) {
  596. // file missing or corrupt — keep current meta
  597. }
  598. }
  599. async persistMeta() {
  600. const metaToSave = { ...this._meta };
  601. // Only save nextId if we have a numeric primary key
  602. if (this._primaryKeyType !== PrimaryKeyType.NUMBER) {
  603. delete metaToSave.nextId;
  604. }
  605. await writeFile(this._metaPath, JSON.stringify(metaToSave));
  606. }
  607. /**
  608. * Compact the database by removing deleted records
  609. * This rewrites the data file without tombstones
  610. *
  611. * @returns {Promise<void>}
  612. */
  613. async compact() {
  614. await this._acquireLock('compact');
  615. try {
  616. const writeStream = createWriteStream(this._dataPath + '.tmp', { flags: 'w' });
  617. for await (const binary of this.find(null, { mode: 'raw' })) {
  618. writeStream.write(binary);
  619. }
  620. await new Promise((resolve, reject) => {
  621. writeStream.end((err) => err ? reject(err) : resolve());
  622. });
  623. // rename is atomic, so we can just rename the file and it will replace the old one
  624. await rename(this._dataPath + '.tmp', this._dataPath);
  625. this._meta.deleted = [];
  626. await this.persistMeta();
  627. // Rebuild indexes — offsets changed after compaction
  628. if (this._indexManager) {
  629. await this._indexManager.init(this._dataPath, { forceRebuild: true });
  630. }
  631. } finally {
  632. await this._releaseLock();
  633. }
  634. }
  635. async _acquireLock(operation, record, retries = 0) {
  636. if (this._lockDepth > 0) {
  637. this._lockDepth++;
  638. return;
  639. }
  640. try {
  641. await writeFile(this._lockPath, String(process.pid), { flag: 'wx' });
  642. this._lockDepth = 1;
  643. } catch (e) {
  644. if (e.code === 'EEXIST') {
  645. this.debug(`Waiting for lock file ${operation} retry #${retries}`, record);
  646. await new Promise(resolve => setTimeout(resolve, 100));
  647. return this._acquireLock(operation, record, retries + 1);
  648. }
  649. throw e;
  650. }
  651. }
  652. async _releaseLock() {
  653. this._lockDepth--;
  654. if (this._lockDepth <= 0) {
  655. this._lockDepth = 0;
  656. await unlink(this._lockPath).catch(() => { });
  657. }
  658. }
  659. /**
  660. * Retrieves a document by its file offset and length
  661. * @param {Array} location - [offset, length] tuple
  662. * @returns {Promise<Object|null>} The deserialized document or null
  663. * @private
  664. * @deprecated Currently unused - may be removed in future versions
  665. */
  666. async _getDocByLocation([offset, length]) {
  667. if (!offset || length === 0) return null;
  668. const fileHandle = await open(this._dataPath, 'r');
  669. try {
  670. const buffer = Buffer.alloc(length);
  671. await fileHandle.read(buffer, 0, length, offset);
  672. return deserialize(buffer);
  673. } finally {
  674. await fileHandle.close();
  675. }
  676. }
  677. /**
  678. * Close the database and persist all pending changes
  679. * Should be called before process exit
  680. *
  681. * @returns {Promise<void>}
  682. */
  683. async close() {
  684. // Persist indexes before closing
  685. if (this._indexManager) {
  686. await this._indexManager.close();
  687. }
  688. // Close data stream
  689. if (this._dataStream) {
  690. await new Promise((resolve, reject) => {
  691. this._dataStream.end((err) => err ? reject(err) : resolve());
  692. });
  693. }
  694. // Remove signal handlers
  695. if (this._processExitHandler) {
  696. process.off('SIGINT', this._processExitHandler);
  697. process.off('SIGTERM', this._processExitHandler);
  698. this._processExitHandler = null;
  699. }
  700. }
  701. debug(...args) {
  702. if (this._debug) {
  703. console.log(...args);
  704. }
  705. }
  706. }
  707. /**
  708. * Helper function to remove a processed chunk from the buffer
  709. * @param {number} recSize - Size of the record to remove
  710. * @param {number} processedBytes - Current offset in the file
  711. * @param {number} totalLength - Total length of buffered data
  712. * @param {Buffer[]} chunks - Array of buffer chunks
  713. * @returns {[number, number]} Updated [processedBytes, totalLength]
  714. */
  715. function removeChunk(recSize, processedBytes, totalLength, chunks) {
  716. processedBytes += recSize;
  717. totalLength -= recSize;
  718. let bytesToRemove = recSize;
  719. while (bytesToRemove > 0 && chunks.length > 0) {
  720. const currentChunk = chunks[0];
  721. if (bytesToRemove >= currentChunk.length) {
  722. bytesToRemove -= currentChunk.length;
  723. chunks.shift();
  724. } else {
  725. chunks[0] = currentChunk.subarray(bytesToRemove);
  726. bytesToRemove = 0;
  727. }
  728. }
  729. return [processedBytes, totalLength];
  730. }
  731. export default MPackDB;

Branches

Latest commits

  • b8ffc1a0release 1.0.5caramboleyo
  • d47876a1reimplemented lost features like indexed find and more testscaramboleyo
  • 7f08da9afixed insert ignoring model definitioncaramboleyo
  • 705774a9added flush before findcaramboleyo
  • b4db6391initial commitcaramboleyo