gitoriaLog in with ident

mpackdb

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commitb4db6391b4db6391initial commitcaramboleyob4db6391/src/MPackDB.js

17.7 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. _initPromise = null;
  27. _initState = 0;
  28. _dataPath = null;
  29. _dataStream = null;
  30. _lockPath = null;
  31. _indexManager = null;
  32. _meta = {
  33. nextId: 0,
  34. deleted: [],
  35. };
  36. _debug = false;
  37. _indexPersistInterval = 60000; // 60 seconds default
  38. _indexPersistThreshold = 1000; // 1000 entries default
  39. _processExitHandler = null;
  40. /**
  41. * Create a new MPackDB instance
  42. *
  43. * @param {string} dbFile - Path to the database file (without extension)
  44. * @param {Object} options - Configuration options
  45. * @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
  46. * @param {PrimaryKeyType} [options.primaryKeyType] - Explicit primary key type (overrides prefix)
  47. * @param {string[]} [options.indexes] - Array of field names to index. Use * prefix for numeric, @ for UUID
  48. * @param {boolean} [options.debug=false] - Enable debug logging
  49. * @param {number} [options.indexPersistInterval=60000] - Milliseconds between automatic index persistence (0 to disable)
  50. * @param {number} [options.indexPersistThreshold=1000] - Number of changes before auto-persisting indexes
  51. *
  52. * @example
  53. * const db = new MPackDB('data/users', {
  54. * primaryKey: '*id', // Numeric auto-increment
  55. * indexes: ['email', '*age'], // Index email (lexical) and age (numeric)
  56. * indexPersistThreshold: 100
  57. * });
  58. */
  59. constructor(dbFile, { primaryKey, primaryKeyType, indexes, debug, indexPersistInterval, indexPersistThreshold } = {}) {
  60. this._dbFile = dbFile || this._dbFile;
  61. this._debug = debug || this._debug;
  62. if (indexPersistInterval !== undefined) this._indexPersistInterval = indexPersistInterval;
  63. if (indexPersistThreshold !== undefined) this._indexPersistThreshold = indexPersistThreshold;
  64. // Parse primary key with optional prefix
  65. if (primaryKey) {
  66. if (primaryKey.startsWith('*')) {
  67. // *id = numeric primary key
  68. this._primaryKey = primaryKey.slice(1);
  69. this._primaryKeyType = PrimaryKeyType.NUMBER;
  70. } else if (primaryKey.startsWith('@')) {
  71. // @id = UUID primary key
  72. this._primaryKey = primaryKey.slice(1);
  73. this._primaryKeyType = PrimaryKeyType.UUID;
  74. } else {
  75. // id = string primary key (lexical)
  76. this._primaryKey = primaryKey;
  77. this._primaryKeyType = PrimaryKeyType.STRING;
  78. }
  79. // Allow explicit override
  80. if (primaryKeyType !== undefined) {
  81. this._primaryKeyType = primaryKeyType;
  82. }
  83. }
  84. const allIndexes = [...new Set([
  85. ...(this._primaryKey ? [this._primaryKey] : []),
  86. ...(indexes || []),
  87. ...(this._indexes || []),
  88. ])];
  89. this._indexes = [];
  90. for (let index of allIndexes) {
  91. let cleanIndex = index;
  92. if (index.startsWith('*')) {
  93. // *field = numeric index
  94. cleanIndex = index.slice(1);
  95. this._indexTypes[cleanIndex] = IndexType.NUMERIC;
  96. } else if (index.startsWith('@')) {
  97. // @field = UUID index (lexical)
  98. cleanIndex = index.slice(1);
  99. this._indexTypes[cleanIndex] = IndexType.LEXICAL;
  100. } else if (index === this._primaryKey && this._primaryKeyType === PrimaryKeyType.NUMBER) {
  101. // Primary key is numeric
  102. this._indexTypes[index] = IndexType.NUMERIC;
  103. } else {
  104. // Default to lexical
  105. this._indexTypes[index] = IndexType.LEXICAL;
  106. }
  107. this._indexes.push(cleanIndex);
  108. }
  109. }
  110. /**
  111. * Initialize the database (called automatically by other methods)
  112. * Performs compaction and loads metadata
  113. *
  114. * @returns {Promise<MPackDB>} The database instance
  115. */
  116. async init() {
  117. if (this._initState === 2) return this;
  118. if (this._initState === 1) {
  119. return this._initPromise;
  120. }
  121. this._initState = 1;
  122. return this._initPromise = new Promise(async success => {
  123. if (!this._dbFile) {
  124. throw new Error('No database file specified.');
  125. }
  126. const dbDir = dirname(this._dbFile);
  127. const baseName = basename(this._dbFile, extname(this._dbFile));
  128. await mkdir(dbDir, { recursive: true });
  129. this._dataPath = resolve(dbDir, `${baseName}.mpack`);
  130. this._lockPath = resolve(dbDir, `${baseName}.lock`);
  131. try {
  132. this._meta = JSON.parse(await readFile(`${this._dbFile}.meta.json`));
  133. } catch (e) {
  134. this._meta = {
  135. nextId: 0,
  136. deleted: [],
  137. };
  138. }
  139. this._initState = 2; // compact causes init to run again thats why we set it to 2 here already
  140. await this.compact();
  141. this._dataStream = createWriteStream(this._dataPath, { flags: 'a' }); // needs to be set after compact as it replaces the file with a tmp file
  142. // Initialize IndexManager if indexes are specified
  143. if (this._indexes.length > 0) {
  144. const dbDir = dirname(this._dbFile);
  145. const baseName = basename(this._dbFile, extname(this._dbFile));
  146. this._indexManager = new IndexManager(dbDir, baseName, this._indexes, this._indexTypes, this._primaryKeyType);
  147. await this._indexManager.init(this._dataPath);
  148. // Start auto-persist with configured interval and threshold
  149. this._indexManager.startAutoPersist(this._indexPersistInterval, this._indexPersistThreshold);
  150. // Register signal handlers for Ctrl+C and kill signals
  151. this._processExitHandler = async () => {
  152. await this.close();
  153. process.exit(0);
  154. };
  155. process.once('SIGINT', this._processExitHandler);
  156. process.once('SIGTERM', this._processExitHandler);
  157. }
  158. this._initPromise = null;
  159. success(this);
  160. });
  161. }
  162. /**
  163. * Insert a new record into the database
  164. *
  165. * @param {Object} record - The record to insert
  166. * @param {Object} [options]
  167. * @param {boolean} [options.skipPrimaryKey=false] - Skip auto-generating primary key
  168. * @returns {Promise<Object>} The inserted record with primary key
  169. *
  170. * @example
  171. * await db.insert({ name: 'Alice', age: 30 });
  172. * // Returns: { id: 0, name: 'Alice', age: 30 }
  173. */
  174. async insert(record, { skipPrimaryKey = false } = {}) {
  175. await this.init();
  176. await this._acquireLock('insert', record);
  177. try {
  178. const recToInsert = { ...record };
  179. if (this._primaryKey && !skipPrimaryKey && !recToInsert[this._primaryKey]) {
  180. if (this._primaryKeyType === PrimaryKeyType.NUMBER) {
  181. recToInsert[this._primaryKey] = this._meta.nextId++;
  182. } else if (this._primaryKeyType === PrimaryKeyType.UUID) {
  183. recToInsert[this._primaryKey] = uuid();
  184. }
  185. }
  186. const packedBuffer = serialize(recToInsert);
  187. const fileStat = await stat(this._dataPath).catch(() => ({ size: 0 }));
  188. const offset = fileStat.size;
  189. await new Promise(resolve => this._dataStream.write(packedBuffer, resolve));
  190. const loc = [offset, packedBuffer.length];
  191. // indexes
  192. if (this._indexManager) {
  193. this._indexManager.insert(recToInsert, loc);
  194. }
  195. // Persist meta if we incremented nextId
  196. if (this._primaryKey && this._primaryKeyType === PrimaryKeyType.NUMBER && !skipPrimaryKey && !record[this._primaryKey]) {
  197. await this.persistMeta();
  198. }
  199. return this._primaryKey ? recToInsert[this._primaryKey] : recToInsert;
  200. } finally {
  201. await this._releaseLock();
  202. }
  203. }
  204. /**
  205. * Update records matching a query
  206. *
  207. * @param {string|Function} mixed - Primary key value or query function
  208. * @param {Object|Function} dataOrCallback - Data to update or callback function
  209. * @param {Object} [options]
  210. * @param {boolean} [options.upsert=false] - Insert if no records match
  211. * @returns {Promise<Object[]>} Array of updated records
  212. *
  213. * @example
  214. * // Update by primary key
  215. * await db.update(0, { age: 31 });
  216. *
  217. * // Update with query function
  218. * await db.update(r => r.age > 30, { status: 'senior' });
  219. *
  220. * // Update with callback
  221. * await db.update(r => r.age > 30, r => ({ ...r, age: r.age + 1 }));
  222. */
  223. async update(mixed, dataOrCallback, { upsert = false } = {}) {
  224. let insertedRecords = await this.delete(mixed, async record => {
  225. return this.insert(typeof dataOrCallback === 'function'
  226. ? await dataOrCallback(record)
  227. : dataOrCallback,
  228. { skipPrimaryKey: true }
  229. );
  230. });
  231. console.log('insertedRecords', insertedRecords);
  232. if (upsert && insertedRecords.length === 0) {
  233. insertedRecords = await this.insert(typeof dataOrCallback === 'function'
  234. ? await dataOrCallback({})
  235. : dataOrCallback);
  236. }
  237. return insertedRecords;
  238. }
  239. /**
  240. * Update records or insert if not found
  241. *
  242. * @param {string|Function} mixed - Primary key value or query function
  243. * @param {Object|Function} dataOrCallback - Data to update/insert or callback function
  244. * @returns {Promise<Object[]>} Array of updated/inserted records
  245. */
  246. async upsert(mixed, dataOrCallback) {
  247. return this.update(mixed, dataOrCallback, { upsert: true });
  248. }
  249. /**
  250. * Delete records matching a query
  251. *
  252. * @param {string|Function} mixed - Primary key value or query function
  253. * @param {Function} [callback] - Optional callback to execute for each deleted record
  254. * @returns {Promise<Object[]>} Array of deleted records
  255. *
  256. * @example
  257. * // Delete by primary key
  258. * await db.delete(0);
  259. *
  260. * // Delete with query function
  261. * await db.delete(r => r.age < 18);
  262. */
  263. async delete(mixed, callback = async record => record) {
  264. await this.init();
  265. await this._acquireLock('delete', mixed);
  266. try {
  267. const promises = [];
  268. for await (const [record, offset] of this.find(mixed, { mode: 'mixed' })) {
  269. this._meta.deleted.push(offset);
  270. if (this._indexManager) {
  271. await this._indexManager.remove(record, this._primaryKey);
  272. }
  273. promises.push(callback(record));
  274. }
  275. await this.persistMeta();
  276. return Promise.all(promises);
  277. } finally {
  278. await this._releaseLock();
  279. }
  280. }
  281. /**
  282. * Finds records in the database
  283. * @param {undefined|string|function} [mixed] - Query: undefined/null for all records, string for primary key lookup, function for filter
  284. * @param {Object} [options] - Query options
  285. * @returns {Cursor} A cursor for iterating over results
  286. * @example
  287. * // Get all records
  288. * for await (const user of db.find()) {
  289. * console.log(user.name);
  290. * }
  291. */
  292. find(mixed, options = {}) {
  293. if (typeof mixed === 'undefined' || mixed === null || typeof mixed === 'function') {
  294. return new Cursor(this, mixed, options);
  295. } else {
  296. if (!this._primaryKey) {
  297. throw new Error('No primary key specified.');
  298. }
  299. return new Cursor(this, record => {
  300. return record[this._primaryKey] === mixed;
  301. }, options);
  302. }
  303. }
  304. /**
  305. * Generator that yields records from the database file
  306. * @param {function|null} [queryFn=null] - Optional filter function
  307. * @param {Object} [options] - Generator options
  308. * @param {string} [options.mode='record'] - Mode: 'record', 'raw', 'offset', or 'mixed'
  309. * @yields {Object|Buffer|Array} Records, buffers, or [record, offset, size] tuples depending on mode
  310. */
  311. async *recordGenerator(queryFn = null, { mode = 'record' } = {}) {
  312. await this.init();
  313. // Check if the data file exists before attempting to read it
  314. try {
  315. await stat(this._dataPath);
  316. } catch (error) {
  317. // File doesn't exist - return empty generator (no records)
  318. return;
  319. }
  320. // Convert deleted array to Set for O(1) lookup instead of O(n)
  321. const deletedSet = new Set(this._meta.deleted);
  322. const readStream = createReadStream(this._dataPath);
  323. let chunks = [];
  324. let totalLength = 0;
  325. let processedBytes = 0;
  326. for await (const chunk of readStream) {
  327. chunks.push(chunk);
  328. totalLength += chunk.length;
  329. while (true) {
  330. // Exit 1: Not enough data to even read the 4-byte size header.
  331. if (totalLength < 4) {
  332. break;
  333. }
  334. // Safely read the header, even if it's split across chunks
  335. let headerBuffer;
  336. if (chunks[0].length >= 4) {
  337. headerBuffer = chunks[0];
  338. } else {
  339. // The header is fragmented, so we must concat just enough to read it.
  340. headerBuffer = Buffer.concat(chunks, 4);
  341. }
  342. const recSize = headerBuffer.readInt32LE(0);
  343. // Validate the record size to prevent infinite loops
  344. // A record must be at least as large as its header (4 bytes).
  345. // A size of 0 or less is invalid and indicates corruption.
  346. if (recSize <= 4) {
  347. throw new Error(`Invalid record size read from stream: ${recSize}`);
  348. }
  349. // Exit 2: We have the size, but not the full record yet.
  350. if (totalLength < recSize) {
  351. break;
  352. }
  353. // skip deleted records (only after we have the full record)
  354. if (deletedSet.has(processedBytes)) {
  355. [processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);
  356. continue; // Goes back to while (true)
  357. }
  358. switch (mode) {
  359. case 'raw':
  360. yield Buffer.concat(chunks, recSize).subarray(0, recSize);
  361. break;
  362. case 'mixed':
  363. case 'record':
  364. const recBuffer = Buffer.concat(chunks, recSize).subarray(0, recSize);
  365. // Skip the 4-byte size header
  366. const data = deserialize(recBuffer.subarray(4));
  367. const rec = this._classToUse
  368. ? Object.assign(new this._classToUse(), data)
  369. : data;
  370. if (queryFn) {
  371. if (queryFn(rec)) {
  372. yield mode === 'mixed' ? [rec, processedBytes, recSize] : rec;
  373. }
  374. } else {
  375. yield mode === 'mixed' ? [rec, processedBytes, recSize] : rec;
  376. }
  377. break;
  378. case 'offset':
  379. // if someone uses a queryFn with offset, we have to deserialize it first to perform the queryFn
  380. if (queryFn) {
  381. const recBuffer = Buffer.concat(chunks, recSize).subarray(0, recSize);
  382. const rec = deserialize(recBuffer.subarray(4));
  383. if (queryFn(rec)) {
  384. yield [processedBytes, recSize];
  385. }
  386. } else {
  387. yield [processedBytes, recSize];
  388. }
  389. break;
  390. default:
  391. throw new Error(`Invalid mode: ${mode}`);
  392. }
  393. [processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);
  394. }
  395. }
  396. }
  397. async persistMeta() {
  398. const metaToSave = { ...this._meta };
  399. // Only save nextId if we have a numeric primary key
  400. if (this._primaryKeyType !== PrimaryKeyType.NUMBER) {
  401. delete metaToSave.nextId;
  402. }
  403. await writeFile(`${this._dbFile}.meta.json`, JSON.stringify(metaToSave));
  404. }
  405. /**
  406. * Compact the database by removing deleted records
  407. * This rewrites the data file without tombstones
  408. *
  409. * @returns {Promise<void>}
  410. */
  411. async compact() {
  412. await this._acquireLock('compact');
  413. try {
  414. const writeStream = createWriteStream(this._dataPath + '.tmp', { flags: 'w' });
  415. for await (const binary of this.find(null, { mode: 'raw' })) {
  416. writeStream.write(binary);
  417. }
  418. writeStream.end();
  419. // rename is atomic, so we can just rename the file and it will replace the old one
  420. await rename(this._dataPath + '.tmp', this._dataPath);
  421. this._meta.deleted = [];
  422. await this.persistMeta();
  423. } finally {
  424. await this._releaseLock();
  425. }
  426. }
  427. async _acquireLock(operation, record, retries = 0) {
  428. try {
  429. await writeFile(this._lockPath, String(process.pid), { flag: 'wx' });
  430. } catch (e) {
  431. if (e.code === 'EEXIST') {
  432. this.debug(`Waiting for lock file ${operation} retry #${retries}`, record);
  433. await new Promise(resolve => setTimeout(resolve, 100));
  434. return this._acquireLock(operation, record, retries + 1);
  435. }
  436. throw e;
  437. }
  438. }
  439. async _releaseLock() {
  440. await unlink(this._lockPath).catch(() => { });
  441. }
  442. /**
  443. * Retrieves a document by its file offset and length
  444. * @param {Array} location - [offset, length] tuple
  445. * @returns {Promise<Object|null>} The deserialized document or null
  446. * @private
  447. * @deprecated Currently unused - may be removed in future versions
  448. */
  449. async _getDocByLocation([offset, length]) {
  450. if (!offset || length === 0) return null;
  451. const fileHandle = await open(this._dataPath, 'r');
  452. try {
  453. const buffer = Buffer.alloc(length);
  454. await fileHandle.read(buffer, 0, length, offset);
  455. return deserialize(buffer);
  456. } finally {
  457. await fileHandle.close();
  458. }
  459. }
  460. /**
  461. * Close the database and persist all pending changes
  462. * Should be called before process exit
  463. *
  464. * @returns {Promise<void>}
  465. */
  466. async close() {
  467. // Persist indexes before closing
  468. if (this._indexManager) {
  469. await this._indexManager.close();
  470. }
  471. // Close data stream
  472. if (this._dataStream) {
  473. await new Promise((resolve, reject) => {
  474. this._dataStream.end((err) => err ? reject(err) : resolve());
  475. });
  476. }
  477. // Remove signal handlers
  478. if (this._processExitHandler) {
  479. process.off('SIGINT', this._processExitHandler);
  480. process.off('SIGTERM', this._processExitHandler);
  481. this._processExitHandler = null;
  482. }
  483. }
  484. debug(...args) {
  485. if (this._debug) {
  486. console.log(...args);
  487. }
  488. }
  489. }
  490. /**
  491. * Helper function to remove a processed chunk from the buffer
  492. * @param {number} recSize - Size of the record to remove
  493. * @param {number} processedBytes - Current offset in the file
  494. * @param {number} totalLength - Total length of buffered data
  495. * @param {Buffer[]} chunks - Array of buffer chunks
  496. * @returns {[number, number]} Updated [processedBytes, totalLength]
  497. */
  498. function removeChunk(recSize, processedBytes, totalLength, chunks) {
  499. processedBytes += recSize;
  500. totalLength -= recSize;
  501. let bytesToRemove = recSize;
  502. while (bytesToRemove > 0 && chunks.length > 0) {
  503. const currentChunk = chunks[0];
  504. if (bytesToRemove >= currentChunk.length) {
  505. bytesToRemove -= currentChunk.length;
  506. chunks.shift();
  507. } else {
  508. chunks[0] = currentChunk.subarray(bytesToRemove);
  509. bytesToRemove = 0;
  510. }
  511. }
  512. return [processedBytes, totalLength];
  513. }
  514. export default MPackDB;

Branches

Latest commits

  • b4db6391initial commitcaramboleyo