gitoriaLog in with ident

mpackdb

All repositories: gitoria

ReadmeCodePull requestsReleasesTicketsSettings
Commit7f08da9a7f08da9afixed insert ignoring model definitioncaramboleyo7f08da9a/src/MPackDB.js

18.0 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 = this._classToUse
  179. ? Object.assign(new this._classToUse(), record)
  180. : { ...record };
  181. if (this._primaryKey && !skipPrimaryKey && !recToInsert[this._primaryKey]) {
  182. if (this._primaryKeyType === PrimaryKeyType.NUMBER) {
  183. recToInsert[this._primaryKey] = this._meta.nextId++;
  184. } else if (this._primaryKeyType === PrimaryKeyType.UUID) {
  185. recToInsert[this._primaryKey] = uuid();
  186. }
  187. }
  188. const packedBuffer = serialize(recToInsert);
  189. const fileStat = await stat(this._dataPath).catch(() => ({ size: 0 }));
  190. const offset = fileStat.size;
  191. await new Promise(resolve => this._dataStream.write(packedBuffer, resolve));
  192. const loc = [offset, packedBuffer.length];
  193. // indexes
  194. if (this._indexManager) {
  195. this._indexManager.insert(recToInsert, loc);
  196. }
  197. // Persist meta if we incremented nextId
  198. if (this._primaryKey && this._primaryKeyType === PrimaryKeyType.NUMBER && !skipPrimaryKey && !record[this._primaryKey]) {
  199. await this.persistMeta();
  200. }
  201. return this._primaryKey ? recToInsert[this._primaryKey] : recToInsert;
  202. } finally {
  203. await this._releaseLock();
  204. }
  205. }
  206. /**
  207. * Update records matching a query
  208. *
  209. * @param {string|Function} mixed - Primary key value or query function
  210. * @param {Object|Function} dataOrCallback - Data to update or callback function
  211. * @param {Object} [options]
  212. * @param {boolean} [options.upsert=false] - Insert if no records match
  213. * @returns {Promise<Object[]>} Array of updated records
  214. *
  215. * @example
  216. * // Update by primary key
  217. * await db.update(0, { age: 31 });
  218. *
  219. * // Update with query function
  220. * await db.update(r => r.age > 30, { status: 'senior' });
  221. *
  222. * // Update with callback
  223. * await db.update(r => r.age > 30, r => ({ ...r, age: r.age + 1 }));
  224. */
  225. async update(mixed, dataOrCallback, { upsert = false } = {}) {
  226. let insertedRecords = await this.delete(mixed, async record => {
  227. return this.insert(typeof dataOrCallback === 'function'
  228. ? await dataOrCallback(record)
  229. : dataOrCallback,
  230. { skipPrimaryKey: true }
  231. );
  232. });
  233. console.log('insertedRecords', insertedRecords);
  234. if (upsert && insertedRecords.length === 0) {
  235. insertedRecords = await this.insert(typeof dataOrCallback === 'function'
  236. ? await dataOrCallback({})
  237. : dataOrCallback);
  238. }
  239. return insertedRecords;
  240. }
  241. /**
  242. * Update records or insert if not found
  243. *
  244. * @param {string|Function} mixed - Primary key value or query function
  245. * @param {Object|Function} dataOrCallback - Data to update/insert or callback function
  246. * @returns {Promise<Object[]>} Array of updated/inserted records
  247. */
  248. async upsert(mixed, dataOrCallback) {
  249. return this.update(mixed, dataOrCallback, { upsert: true });
  250. }
  251. /**
  252. * Delete records matching a query
  253. *
  254. * @param {string|Function} mixed - Primary key value or query function
  255. * @param {Function} [callback] - Optional callback to execute for each deleted record
  256. * @returns {Promise<Object[]>} Array of deleted records
  257. *
  258. * @example
  259. * // Delete by primary key
  260. * await db.delete(0);
  261. *
  262. * // Delete with query function
  263. * await db.delete(r => r.age < 18);
  264. */
  265. async delete(mixed, callback = async record => record) {
  266. await this.init();
  267. await this._acquireLock('delete', mixed);
  268. try {
  269. const promises = [];
  270. for await (const [record, offset] of this.find(mixed, { mode: 'mixed' })) {
  271. this._meta.deleted.push(offset);
  272. if (this._indexManager) {
  273. await this._indexManager.remove(record, this._primaryKey);
  274. }
  275. promises.push(callback(record));
  276. }
  277. await this.persistMeta();
  278. return Promise.all(promises);
  279. } finally {
  280. await this._releaseLock();
  281. }
  282. }
  283. /**
  284. * Finds records in the database
  285. * @param {undefined|string|function} [mixed] - Query: undefined/null for all records, string for primary key lookup, function for filter
  286. * @param {Object} [options] - Query options
  287. * @returns {Cursor} A cursor for iterating over results
  288. * @example
  289. * // Get all records
  290. * for await (const user of db.find()) {
  291. * console.log(user.name);
  292. * }
  293. */
  294. find(mixed, options = {}) {
  295. if (typeof mixed === 'undefined' || mixed === null || typeof mixed === 'function') {
  296. return new Cursor(this, mixed, options);
  297. } else {
  298. if (!this._primaryKey) {
  299. throw new Error('No primary key specified.');
  300. }
  301. return new Cursor(this, record => {
  302. return record[this._primaryKey] === mixed;
  303. }, options);
  304. }
  305. }
  306. /**
  307. * Generator that yields records from the database file
  308. * @param {function|null} [queryFn=null] - Optional filter function
  309. * @param {Object} [options] - Generator options
  310. * @param {string} [options.mode='record'] - Mode: 'record', 'raw', 'offset', or 'mixed'
  311. * @yields {Object|Buffer|Array} Records, buffers, or [record, offset, size] tuples depending on mode
  312. */
  313. async *recordGenerator(queryFn = null, { mode = 'record' } = {}) {
  314. await this.init();
  315. // Flush pending writes so reads see all inserted data
  316. if (this._dataStream && this._dataStream.writableLength > 0) {
  317. await new Promise(resolve => this._dataStream.once('drain', resolve));
  318. }
  319. // Check if the data file exists before attempting to read it
  320. try {
  321. await stat(this._dataPath);
  322. } catch (error) {
  323. // File doesn't exist - return empty generator (no records)
  324. return;
  325. }
  326. // Convert deleted array to Set for O(1) lookup instead of O(n)
  327. const deletedSet = new Set(this._meta.deleted);
  328. const readStream = createReadStream(this._dataPath);
  329. let chunks = [];
  330. let totalLength = 0;
  331. let processedBytes = 0;
  332. for await (const chunk of readStream) {
  333. chunks.push(chunk);
  334. totalLength += chunk.length;
  335. while (true) {
  336. // Exit 1: Not enough data to even read the 4-byte size header.
  337. if (totalLength < 4) {
  338. break;
  339. }
  340. // Safely read the header, even if it's split across chunks
  341. let headerBuffer;
  342. if (chunks[0].length >= 4) {
  343. headerBuffer = chunks[0];
  344. } else {
  345. // The header is fragmented, so we must concat just enough to read it.
  346. headerBuffer = Buffer.concat(chunks, 4);
  347. }
  348. const recSize = headerBuffer.readInt32LE(0);
  349. // Validate the record size to prevent infinite loops
  350. // A record must be at least as large as its header (4 bytes).
  351. // A size of 0 or less is invalid and indicates corruption.
  352. if (recSize <= 4) {
  353. throw new Error(`Invalid record size read from stream: ${recSize}`);
  354. }
  355. // Exit 2: We have the size, but not the full record yet.
  356. if (totalLength < recSize) {
  357. break;
  358. }
  359. // skip deleted records (only after we have the full record)
  360. if (deletedSet.has(processedBytes)) {
  361. [processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);
  362. continue; // Goes back to while (true)
  363. }
  364. switch (mode) {
  365. case 'raw':
  366. yield Buffer.concat(chunks, recSize).subarray(0, recSize);
  367. break;
  368. case 'mixed':
  369. case 'record':
  370. const recBuffer = Buffer.concat(chunks, recSize).subarray(0, recSize);
  371. // Skip the 4-byte size header
  372. const data = deserialize(recBuffer.subarray(4));
  373. const rec = this._classToUse
  374. ? Object.assign(new this._classToUse(), data)
  375. : data;
  376. if (queryFn) {
  377. if (queryFn(rec)) {
  378. yield mode === 'mixed' ? [rec, processedBytes, recSize] : rec;
  379. }
  380. } else {
  381. yield mode === 'mixed' ? [rec, processedBytes, recSize] : rec;
  382. }
  383. break;
  384. case 'offset':
  385. // if someone uses a queryFn with offset, we have to deserialize it first to perform the queryFn
  386. if (queryFn) {
  387. const recBuffer = Buffer.concat(chunks, recSize).subarray(0, recSize);
  388. const rec = deserialize(recBuffer.subarray(4));
  389. if (queryFn(rec)) {
  390. yield [processedBytes, recSize];
  391. }
  392. } else {
  393. yield [processedBytes, recSize];
  394. }
  395. break;
  396. default:
  397. throw new Error(`Invalid mode: ${mode}`);
  398. }
  399. [processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);
  400. }
  401. }
  402. }
  403. async persistMeta() {
  404. const metaToSave = { ...this._meta };
  405. // Only save nextId if we have a numeric primary key
  406. if (this._primaryKeyType !== PrimaryKeyType.NUMBER) {
  407. delete metaToSave.nextId;
  408. }
  409. await writeFile(`${this._dbFile}.meta.json`, JSON.stringify(metaToSave));
  410. }
  411. /**
  412. * Compact the database by removing deleted records
  413. * This rewrites the data file without tombstones
  414. *
  415. * @returns {Promise<void>}
  416. */
  417. async compact() {
  418. await this._acquireLock('compact');
  419. try {
  420. const writeStream = createWriteStream(this._dataPath + '.tmp', { flags: 'w' });
  421. for await (const binary of this.find(null, { mode: 'raw' })) {
  422. writeStream.write(binary);
  423. }
  424. writeStream.end();
  425. // rename is atomic, so we can just rename the file and it will replace the old one
  426. await rename(this._dataPath + '.tmp', this._dataPath);
  427. this._meta.deleted = [];
  428. await this.persistMeta();
  429. } finally {
  430. await this._releaseLock();
  431. }
  432. }
  433. async _acquireLock(operation, record, retries = 0) {
  434. try {
  435. await writeFile(this._lockPath, String(process.pid), { flag: 'wx' });
  436. } catch (e) {
  437. if (e.code === 'EEXIST') {
  438. this.debug(`Waiting for lock file ${operation} retry #${retries}`, record);
  439. await new Promise(resolve => setTimeout(resolve, 100));
  440. return this._acquireLock(operation, record, retries + 1);
  441. }
  442. throw e;
  443. }
  444. }
  445. async _releaseLock() {
  446. await unlink(this._lockPath).catch(() => { });
  447. }
  448. /**
  449. * Retrieves a document by its file offset and length
  450. * @param {Array} location - [offset, length] tuple
  451. * @returns {Promise<Object|null>} The deserialized document or null
  452. * @private
  453. * @deprecated Currently unused - may be removed in future versions
  454. */
  455. async _getDocByLocation([offset, length]) {
  456. if (!offset || length === 0) return null;
  457. const fileHandle = await open(this._dataPath, 'r');
  458. try {
  459. const buffer = Buffer.alloc(length);
  460. await fileHandle.read(buffer, 0, length, offset);
  461. return deserialize(buffer);
  462. } finally {
  463. await fileHandle.close();
  464. }
  465. }
  466. /**
  467. * Close the database and persist all pending changes
  468. * Should be called before process exit
  469. *
  470. * @returns {Promise<void>}
  471. */
  472. async close() {
  473. // Persist indexes before closing
  474. if (this._indexManager) {
  475. await this._indexManager.close();
  476. }
  477. // Close data stream
  478. if (this._dataStream) {
  479. await new Promise((resolve, reject) => {
  480. this._dataStream.end((err) => err ? reject(err) : resolve());
  481. });
  482. }
  483. // Remove signal handlers
  484. if (this._processExitHandler) {
  485. process.off('SIGINT', this._processExitHandler);
  486. process.off('SIGTERM', this._processExitHandler);
  487. this._processExitHandler = null;
  488. }
  489. }
  490. debug(...args) {
  491. if (this._debug) {
  492. console.log(...args);
  493. }
  494. }
  495. }
  496. /**
  497. * Helper function to remove a processed chunk from the buffer
  498. * @param {number} recSize - Size of the record to remove
  499. * @param {number} processedBytes - Current offset in the file
  500. * @param {number} totalLength - Total length of buffered data
  501. * @param {Buffer[]} chunks - Array of buffer chunks
  502. * @returns {[number, number]} Updated [processedBytes, totalLength]
  503. */
  504. function removeChunk(recSize, processedBytes, totalLength, chunks) {
  505. processedBytes += recSize;
  506. totalLength -= recSize;
  507. let bytesToRemove = recSize;
  508. while (bytesToRemove > 0 && chunks.length > 0) {
  509. const currentChunk = chunks[0];
  510. if (bytesToRemove >= currentChunk.length) {
  511. bytesToRemove -= currentChunk.length;
  512. chunks.shift();
  513. } else {
  514. chunks[0] = currentChunk.subarray(bytesToRemove);
  515. bytesToRemove = 0;
  516. }
  517. }
  518. return [processedBytes, totalLength];
  519. }
  520. export default MPackDB;

Branches

Latest commits

  • 7f08da9afixed insert ignoring model definitioncaramboleyo
  • 705774a9added flush before findcaramboleyo
  • b4db6391initial commitcaramboleyo