gitoriaLog in with ident

mpackdb

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commit705774a9705774a9added flush before findcaramboleyo705774a9/src/MPackDB.js

17.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. _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. // Flush pending writes so reads see all inserted data
  314. if (this._dataStream && this._dataStream.writableLength > 0) {
  315. await new Promise(resolve => this._dataStream.once('drain', resolve));
  316. }
  317. // Check if the data file exists before attempting to read it
  318. try {
  319. await stat(this._dataPath);
  320. } catch (error) {
  321. // File doesn't exist - return empty generator (no records)
  322. return;
  323. }
  324. // Convert deleted array to Set for O(1) lookup instead of O(n)
  325. const deletedSet = new Set(this._meta.deleted);
  326. const readStream = createReadStream(this._dataPath);
  327. let chunks = [];
  328. let totalLength = 0;
  329. let processedBytes = 0;
  330. for await (const chunk of readStream) {
  331. chunks.push(chunk);
  332. totalLength += chunk.length;
  333. while (true) {
  334. // Exit 1: Not enough data to even read the 4-byte size header.
  335. if (totalLength < 4) {
  336. break;
  337. }
  338. // Safely read the header, even if it's split across chunks
  339. let headerBuffer;
  340. if (chunks[0].length >= 4) {
  341. headerBuffer = chunks[0];
  342. } else {
  343. // The header is fragmented, so we must concat just enough to read it.
  344. headerBuffer = Buffer.concat(chunks, 4);
  345. }
  346. const recSize = headerBuffer.readInt32LE(0);
  347. // Validate the record size to prevent infinite loops
  348. // A record must be at least as large as its header (4 bytes).
  349. // A size of 0 or less is invalid and indicates corruption.
  350. if (recSize <= 4) {
  351. throw new Error(`Invalid record size read from stream: ${recSize}`);
  352. }
  353. // Exit 2: We have the size, but not the full record yet.
  354. if (totalLength < recSize) {
  355. break;
  356. }
  357. // skip deleted records (only after we have the full record)
  358. if (deletedSet.has(processedBytes)) {
  359. [processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);
  360. continue; // Goes back to while (true)
  361. }
  362. switch (mode) {
  363. case 'raw':
  364. yield Buffer.concat(chunks, recSize).subarray(0, recSize);
  365. break;
  366. case 'mixed':
  367. case 'record':
  368. const recBuffer = Buffer.concat(chunks, recSize).subarray(0, recSize);
  369. // Skip the 4-byte size header
  370. const data = deserialize(recBuffer.subarray(4));
  371. const rec = this._classToUse
  372. ? Object.assign(new this._classToUse(), data)
  373. : data;
  374. if (queryFn) {
  375. if (queryFn(rec)) {
  376. yield mode === 'mixed' ? [rec, processedBytes, recSize] : rec;
  377. }
  378. } else {
  379. yield mode === 'mixed' ? [rec, processedBytes, recSize] : rec;
  380. }
  381. break;
  382. case 'offset':
  383. // if someone uses a queryFn with offset, we have to deserialize it first to perform the queryFn
  384. if (queryFn) {
  385. const recBuffer = Buffer.concat(chunks, recSize).subarray(0, recSize);
  386. const rec = deserialize(recBuffer.subarray(4));
  387. if (queryFn(rec)) {
  388. yield [processedBytes, recSize];
  389. }
  390. } else {
  391. yield [processedBytes, recSize];
  392. }
  393. break;
  394. default:
  395. throw new Error(`Invalid mode: ${mode}`);
  396. }
  397. [processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);
  398. }
  399. }
  400. }
  401. async persistMeta() {
  402. const metaToSave = { ...this._meta };
  403. // Only save nextId if we have a numeric primary key
  404. if (this._primaryKeyType !== PrimaryKeyType.NUMBER) {
  405. delete metaToSave.nextId;
  406. }
  407. await writeFile(`${this._dbFile}.meta.json`, JSON.stringify(metaToSave));
  408. }
  409. /**
  410. * Compact the database by removing deleted records
  411. * This rewrites the data file without tombstones
  412. *
  413. * @returns {Promise<void>}
  414. */
  415. async compact() {
  416. await this._acquireLock('compact');
  417. try {
  418. const writeStream = createWriteStream(this._dataPath + '.tmp', { flags: 'w' });
  419. for await (const binary of this.find(null, { mode: 'raw' })) {
  420. writeStream.write(binary);
  421. }
  422. writeStream.end();
  423. // rename is atomic, so we can just rename the file and it will replace the old one
  424. await rename(this._dataPath + '.tmp', this._dataPath);
  425. this._meta.deleted = [];
  426. await this.persistMeta();
  427. } finally {
  428. await this._releaseLock();
  429. }
  430. }
  431. async _acquireLock(operation, record, retries = 0) {
  432. try {
  433. await writeFile(this._lockPath, String(process.pid), { flag: 'wx' });
  434. } catch (e) {
  435. if (e.code === 'EEXIST') {
  436. this.debug(`Waiting for lock file ${operation} retry #${retries}`, record);
  437. await new Promise(resolve => setTimeout(resolve, 100));
  438. return this._acquireLock(operation, record, retries + 1);
  439. }
  440. throw e;
  441. }
  442. }
  443. async _releaseLock() {
  444. await unlink(this._lockPath).catch(() => { });
  445. }
  446. /**
  447. * Retrieves a document by its file offset and length
  448. * @param {Array} location - [offset, length] tuple
  449. * @returns {Promise<Object|null>} The deserialized document or null
  450. * @private
  451. * @deprecated Currently unused - may be removed in future versions
  452. */
  453. async _getDocByLocation([offset, length]) {
  454. if (!offset || length === 0) return null;
  455. const fileHandle = await open(this._dataPath, 'r');
  456. try {
  457. const buffer = Buffer.alloc(length);
  458. await fileHandle.read(buffer, 0, length, offset);
  459. return deserialize(buffer);
  460. } finally {
  461. await fileHandle.close();
  462. }
  463. }
  464. /**
  465. * Close the database and persist all pending changes
  466. * Should be called before process exit
  467. *
  468. * @returns {Promise<void>}
  469. */
  470. async close() {
  471. // Persist indexes before closing
  472. if (this._indexManager) {
  473. await this._indexManager.close();
  474. }
  475. // Close data stream
  476. if (this._dataStream) {
  477. await new Promise((resolve, reject) => {
  478. this._dataStream.end((err) => err ? reject(err) : resolve());
  479. });
  480. }
  481. // Remove signal handlers
  482. if (this._processExitHandler) {
  483. process.off('SIGINT', this._processExitHandler);
  484. process.off('SIGTERM', this._processExitHandler);
  485. this._processExitHandler = null;
  486. }
  487. }
  488. debug(...args) {
  489. if (this._debug) {
  490. console.log(...args);
  491. }
  492. }
  493. }
  494. /**
  495. * Helper function to remove a processed chunk from the buffer
  496. * @param {number} recSize - Size of the record to remove
  497. * @param {number} processedBytes - Current offset in the file
  498. * @param {number} totalLength - Total length of buffered data
  499. * @param {Buffer[]} chunks - Array of buffer chunks
  500. * @returns {[number, number]} Updated [processedBytes, totalLength]
  501. */
  502. function removeChunk(recSize, processedBytes, totalLength, chunks) {
  503. processedBytes += recSize;
  504. totalLength -= recSize;
  505. let bytesToRemove = recSize;
  506. while (bytesToRemove > 0 && chunks.length > 0) {
  507. const currentChunk = chunks[0];
  508. if (bytesToRemove >= currentChunk.length) {
  509. bytesToRemove -= currentChunk.length;
  510. chunks.shift();
  511. } else {
  512. chunks[0] = currentChunk.subarray(bytesToRemove);
  513. bytesToRemove = 0;
  514. }
  515. }
  516. return [processedBytes, totalLength];
  517. }
  518. export default MPackDB;

Branches

Latest commits

  • 705774a9added flush before findcaramboleyo
  • b4db6391initial commitcaramboleyo