mpackdb
All repositories: gitoria
27.0 KB
import { dirname, basename, extname, resolve } from 'path';import { createReadStream, createWriteStream } from 'fs';import { mkdir, open, stat, readFile, writeFile, rename, unlink } from 'fs/promises';import { serialize, deserialize, uuid, PrimaryKeyType, IndexType } from './mpack.js';import { Cursor } from './Cursor.js';import { IndexManager } from './IndexManager.js';export { PrimaryKeyType, IndexType };/*** MPackDB - A fast, local, append-only JSON database with MessagePack serialization** Features:* - Append-only writes for high performance* - MessagePack binary serialization* - Optional indexes (numeric and lexical)* - File-based locking for concurrent access* - Auto-compaction on startup* - Auto-persistence of indexes*/export class MPackDB {_dbFile = null;_primaryKeyType = null;_primaryKey = null;_classToUse = null;_indexes = [];_indexTypes = {};_uniqueIndexes = new Set();_initPromise = null;_initState = 0;_dataPath = null;_dataStream = null;_lockPath = null;_lockDepth = 0;_indexManager = null;_meta = {nextId: 0,deleted: [],};_debug = false;_indexPersistInterval = 60000; // 60 seconds default_indexPersistThreshold = 1000; // 1000 entries default_processExitHandler = null;/*** Create a new MPackDB instance** @param {string} dbFile - Path to the database file (without extension)* @param {Object} options - Configuration options* @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* @param {PrimaryKeyType} [options.primaryKeyType] - Explicit primary key type (overrides prefix)* @param {string[]} [options.indexes] - Array of field names to index. Use * prefix for numeric, @ for UUID* @param {boolean} [options.debug=false] - Enable debug logging* @param {number} [options.indexPersistInterval=60000] - Milliseconds between automatic index persistence (0 to disable)* @param {number} [options.indexPersistThreshold=1000] - Number of changes before auto-persisting indexes* @param {boolean} [options.compact=true] - Run compaction on init (set false for read-only / secondary instances)** @example* const db = new MPackDB('data/users', {* primaryKey: '*id', // Numeric auto-increment* indexes: ['email', '*age'], // Index email (lexical) and age (numeric)* indexPersistThreshold: 100* });*/_compact = true;constructor(dbFile, { primaryKey, primaryKeyType, indexes, debug, indexPersistInterval, indexPersistThreshold, compact } = {}) {this._dbFile = dbFile || this._dbFile;this._debug = debug || this._debug;if (compact !== undefined) this._compact = compact;if (indexPersistInterval !== undefined) this._indexPersistInterval = indexPersistInterval;if (indexPersistThreshold !== undefined) this._indexPersistThreshold = indexPersistThreshold;// Parse primary key with optional prefixif (primaryKey) {if (primaryKey.startsWith('*')) {// *id = numeric primary keythis._primaryKey = primaryKey.slice(1);this._primaryKeyType = PrimaryKeyType.NUMBER;} else if (primaryKey.startsWith('@')) {// @id = UUID primary keythis._primaryKey = primaryKey.slice(1);this._primaryKeyType = PrimaryKeyType.UUID;} else {// id = string primary key (lexical)this._primaryKey = primaryKey;this._primaryKeyType = PrimaryKeyType.STRING;}// Allow explicit overrideif (primaryKeyType !== undefined) {this._primaryKeyType = primaryKeyType;}}// Primary key is always uniqueif (this._primaryKey) {this._uniqueIndexes.add(this._primaryKey);}const allIndexes = [...new Set([...(this._primaryKey ? [this._primaryKey] : []),...(indexes || []),...(this._indexes || []),])];this._indexes = [];for (let index of allIndexes) {let cleanIndex = index;// !field = unique index (can combine with * and @: !*field, !@field)let isUnique = false;if (cleanIndex.startsWith('!')) {isUnique = true;cleanIndex = cleanIndex.slice(1);}if (cleanIndex.startsWith('*')) {// *field = numeric indexcleanIndex = cleanIndex.slice(1);this._indexTypes[cleanIndex] = IndexType.NUMERIC;} else if (cleanIndex.startsWith('@')) {// @field = UUID index (lexical)cleanIndex = cleanIndex.slice(1);this._indexTypes[cleanIndex] = IndexType.LEXICAL;} else if (cleanIndex === this._primaryKey && this._primaryKeyType === PrimaryKeyType.NUMBER) {// Primary key is numericthis._indexTypes[cleanIndex] = IndexType.NUMERIC;} else {// Default to lexicalthis._indexTypes[cleanIndex] = IndexType.LEXICAL;}if (isUnique) {this._uniqueIndexes.add(cleanIndex);}this._indexes.push(cleanIndex);}}/*** Initialize the database (called automatically by other methods)* Performs compaction and loads metadata** @returns {Promise<MPackDB>} The database instance*/async init() {if (this._initState === 2) return this;if (this._initState === 1) {return this._initPromise;}this._initState = 1;return this._initPromise = new Promise(async success => {if (!this._dbFile) {throw new Error('No database file specified.');}const dbDir = dirname(this._dbFile);const baseName = basename(this._dbFile, extname(this._dbFile));await mkdir(dbDir, { recursive: true });this._dataPath = resolve(dbDir, `${baseName}.mpack`);this._metaPath = resolve(dbDir, `${baseName}.meta.json`);this._lockPath = resolve(dbDir, `${baseName}.lock`);try {this._meta = JSON.parse(await readFile(this._metaPath));} catch (e) {this._meta = {nextId: 0,deleted: [],};}this._initState = 2; // compact causes init to run again thats why we set it to 2 here alreadyif (this._compact) await this.compact();this._dataStream = createWriteStream(this._dataPath, { flags: 'a' }); // needs to be set after compact as it replaces the file with a tmp file// Initialize IndexManager if indexes are specifiedif (this._indexes.length > 0) {const dbDir = dirname(this._dbFile);const baseName = basename(this._dbFile, extname(this._dbFile));this._indexManager = new IndexManager(dbDir, baseName, this._indexes, this._indexTypes, this._primaryKeyType);await this._indexManager.init(this._dataPath);// Start auto-persist with configured interval and thresholdthis._indexManager.startAutoPersist(this._indexPersistInterval, this._indexPersistThreshold);// Register signal handlers for Ctrl+C and kill signalsthis._processExitHandler = async () => {await this.close();process.exit(0);};process.once('SIGINT', this._processExitHandler);process.once('SIGTERM', this._processExitHandler);}this._initPromise = null;success(this);});}/*** Insert a new record into the database** @param {Object} record - The record to insert* @param {Object} [options]* @param {boolean} [options.skipPrimaryKey=false] - Skip auto-generating primary key* @returns {Promise<Object>} The inserted record with primary key** @example* await db.insert({ name: 'Alice', age: 30 });* // Returns: { id: 0, name: 'Alice', age: 30 }*/async insert(record, { skipPrimaryKey = false } = {}) {await this.init();await this._acquireLock('insert', record);try {const recToInsert = this._classToUse? Object.assign(new this._classToUse(), record): { ...record };if (this._primaryKey && !skipPrimaryKey && !recToInsert[this._primaryKey]) {if (this._primaryKeyType === PrimaryKeyType.NUMBER) {recToInsert[this._primaryKey] = this._meta.nextId++;} else if (this._primaryKeyType === PrimaryKeyType.UUID) {recToInsert[this._primaryKey] = uuid();}}// Check unique index constraints (skip auto-generated primary keys — guaranteed unique)if (this._indexManager && this._uniqueIndexes.size > 0) {const autoGenPK = this._primaryKey && !skipPrimaryKey && !record[this._primaryKey]&& (this._primaryKeyType === PrimaryKeyType.NUMBER || this._primaryKeyType === PrimaryKeyType.UUID);const deletedSet = new Set(this._meta.deleted);for (const field of this._uniqueIndexes) {if (autoGenPK && field === this._primaryKey) continue;const value = recToInsert[field];if (value === undefined) continue;const existing = await this._indexManager.get(field, value);const live = existing.filter(e => !deletedSet.has(e.loc[0]));if (live.length > 0) {const err = new Error(`Duplicate key: ${field}=${value}`);err.code = 'DUPLICATE_KEY';err.field = field;err.value = value;throw err;}}}const packedBuffer = serialize(recToInsert);const fileStat = await stat(this._dataPath).catch(() => ({ size: 0 }));const offset = fileStat.size;await new Promise(resolve => this._dataStream.write(packedBuffer, resolve));const loc = [offset, packedBuffer.length];// indexesif (this._indexManager) {this._indexManager.insert(recToInsert, loc);}// Persist meta if we incremented nextIdif (this._primaryKey && this._primaryKeyType === PrimaryKeyType.NUMBER && !skipPrimaryKey && !record[this._primaryKey]) {await this.persistMeta();}return this._primaryKey ? recToInsert[this._primaryKey] : recToInsert;} finally {await this._releaseLock();}}/*** Update records matching a query** @param {string|Function} mixed - Primary key value or query function* @param {Object|Function} dataOrCallback - Data to update or callback function* @param {Object} [options]* @param {boolean} [options.upsert=false] - Insert if no records match* @param {Object|Array} [options.index] - Index hint(s) to narrow disk reads* @returns {Promise<Object[]>} Array of updated records** @example* // Update by primary key* await db.update(0, { age: 31 });** // Update with query function* await db.update(r => r.age > 30, { status: 'senior' });** // Update with callback* await db.update(r => r.age > 30, r => ({ ...r, age: r.age + 1 }));*/async update(mixed, dataOrCallback, { upsert = false, index } = {}) {let insertedRecords = await this.delete(mixed, {index,callback: async record => {return this.insert(typeof dataOrCallback === 'function'? await dataOrCallback(record): dataOrCallback,{ skipPrimaryKey: true });}});if (upsert && insertedRecords.length === 0) {insertedRecords = await this.insert(typeof dataOrCallback === 'function'? await dataOrCallback({}): dataOrCallback);}return insertedRecords;}/*** Update records or insert if not found** @param {string|Function} mixed - Primary key value or query function* @param {Object|Function} dataOrCallback - Data to update/insert or callback function* @param {Object} [options]* @param {Object|Array} [options.index] - Index hint(s) to narrow disk reads* @returns {Promise<Object[]>} Array of updated/inserted records*/async upsert(mixed, dataOrCallback, { index } = {}) {return this.update(mixed, dataOrCallback, { upsert: true, index });}/*** Delete records matching a query** @param {string|Function} mixed - Primary key value or query function* @param {Object} [options]* @param {Object|Array} [options.index] - Index hint(s) to narrow disk reads* @param {Function} [options.callback] - Internal callback per deleted record (used by update)* @returns {Promise<Object[]>} Array of deleted records** @example* // Delete by primary key* await db.delete(0);** // Delete with query function* await db.delete(r => r.age < 18);** // Delete with index hint* await db.delete(r => r.status === 'inactive', { index: { field: 'status', value: 'inactive' } });*/async delete(mixed, { index, callback = async record => record } = {}) {await this.init();await this._acquireLock('delete', mixed);try {const promises = [];for await (const [record, offset] of this.find(mixed, { mode: 'mixed', index })) {this._meta.deleted.push(offset);if (this._indexManager) {await this._indexManager.remove(record, this._primaryKey);}promises.push(callback(record));}await this.persistMeta();return Promise.all(promises);} finally {await this._releaseLock();}}/*** Finds records in the database* @param {undefined|string|number|function} [mixed] - Query: undefined/null for all, primary key value for PK lookup, function for filter* @param {Object} [options] - Query options* @param {Object|Array} [options.index] - Index hint(s) to narrow disk reads before filtering* @returns {Cursor} A cursor for iterating over results* @example* // All records* for await (const user of db.find()) { ... }** // Primary key lookup* const [user] = await db.find(68);** // Filter with index range — only reads records in x 100+* for await (const doc of db.find(r => r.x <= 200, { index: { field: 'x', from: 100 } })) { ... }** // Index intersection — intersects offsets first, then streams matches* for await (const doc of db.find(r => r.x <= 200, {* index: [* { field: 'x', from: 100 },* { field: 'status', value: 'active' }* ]* })) { ... }*/find(mixed, options = {}) {if (typeof mixed === 'undefined' || mixed === null || typeof mixed === 'function') {return new Cursor(this, mixed, options);} else {if (!this._primaryKey) {throw new Error('No primary key specified.');}// Use indexed lookup when availableif (this._indexManager) {return new Cursor(this, null, {...options,_indexLookup: { field: this._primaryKey, value: mixed }});}return new Cursor(this, record => {return record[this._primaryKey] === mixed;}, options);}}/*** Find records within a bounding box defined by 4 corners.* Requires numeric indexes on x and y fields.* Corners can be in any order — min/max are extracted automatically.** @param {Array<{x: number, y: number}>} corners - 4 corner coordinates* @param {function} [filter] - Optional additional filter function* @returns {Cursor}* @example* const results = await db.boundingBox([* { x: -5, y: 5 }, { x: 5, y: 5 },* { x: 5, y: -5 }, { x: -5, y: -5 }* ]);*/boundingBox(corners, filter) {const xs = corners.map(c => c.x);const ys = corners.map(c => c.y);const minX = Math.min(...xs), maxX = Math.max(...xs);const minY = Math.min(...ys), maxY = Math.max(...ys);return this.find(filter || null, {index: [{ field: 'x', from: minX, to: maxX },{ field: 'y', from: minY, to: maxY }]});}/*** Execute a callback while holding the database lock.* The lock is re-entrant: find/insert/delete/update called inside* the callback reuse the same lock instead of deadlocking.** Use this for compound operations that must be atomic, e.g.* find-then-insert (login pattern).** @param {Function} callback - Async function to execute under lock* @returns {Promise<any>} The return value of the callback** @example* const user = await db.withLock(async () => {* const [existing] = await db.find(u => u.email === email);* if (existing) return existing;* return db.insert({ email });* });*/async withLock(callback) {await this.init();await this._acquireLock('withLock');try {return await callback();} finally {await this._releaseLock();}}/*** Generator that yields records from the database file* @param {function|null} [queryFn=null] - Optional filter function* @param {Object} [options] - Generator options* @param {string} [options.mode='record'] - Mode: 'record', 'raw', 'offset', or 'mixed'* @yields {Object|Buffer|Array} Records, buffers, or [record, offset, size] tuples depending on mode*/async *recordGenerator(queryFn = null, { mode = 'record', _indexLookup, index } = {}) {await this.init();await this.refresh();// Indexed primary key lookup — O(log n) instead of full scanif (_indexLookup && this._indexManager) {yield* this._indexedLookup(_indexLookup.field, _indexLookup.value, mode);return;}// Index-based find: collect offsets from index(es), then stream only those recordsif (index && this._indexManager) {yield* this._indexedStream(index, queryFn, mode);return;}// Flush pending writes so reads see all inserted dataif (this._dataStream && this._dataStream.writableLength > 0) {await new Promise(resolve => this._dataStream.once('drain', resolve));}// Check if the data file exists before attempting to read ittry {await stat(this._dataPath);} catch (error) {// File doesn't exist - return empty generator (no records)return;}// Convert deleted array to Set for O(1) lookup instead of O(n)const deletedSet = new Set(this._meta.deleted);const readStream = createReadStream(this._dataPath);let chunks = [];let totalLength = 0;let processedBytes = 0;for await (const chunk of readStream) {chunks.push(chunk);totalLength += chunk.length;while (true) {// Exit 1: Not enough data to even read the 4-byte size header.if (totalLength < 4) {break;}// Safely read the header, even if it's split across chunkslet headerBuffer;if (chunks[0].length >= 4) {headerBuffer = chunks[0];} else {// The header is fragmented, so we must concat just enough to read it.headerBuffer = Buffer.concat(chunks, 4);}const recSize = headerBuffer.readInt32LE(0);// Validate the record size to prevent infinite loops// A record must be at least as large as its header (4 bytes).// A size of 0 or less is invalid and indicates corruption.if (recSize <= 4) {throw new Error(`Invalid record size read from stream: ${recSize}`);}// Exit 2: We have the size, but not the full record yet.if (totalLength < recSize) {break;}// skip deleted records (only after we have the full record)if (deletedSet.has(processedBytes)) {[processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);continue; // Goes back to while (true)}switch (mode) {case 'raw':yield Buffer.concat(chunks, recSize).subarray(0, recSize);break;case 'mixed':case 'record':const recBuffer = Buffer.concat(chunks, recSize).subarray(0, recSize);// Skip the 4-byte size headerconst data = deserialize(recBuffer.subarray(4));const rec = this._classToUse? Object.assign(new this._classToUse(), data): data;if (queryFn) {if (queryFn(rec)) {yield mode === 'mixed' ? [rec, processedBytes, recSize] : rec;}} else {yield mode === 'mixed' ? [rec, processedBytes, recSize] : rec;}break;case 'offset':// if someone uses a queryFn with offset, we have to deserialize it first to perform the queryFnif (queryFn) {const recBuffer = Buffer.concat(chunks, recSize).subarray(0, recSize);const rec = deserialize(recBuffer.subarray(4));if (queryFn(rec)) {yield [processedBytes, recSize];}} else {yield [processedBytes, recSize];}break;default:throw new Error(`Invalid mode: ${mode}`);}[processedBytes, totalLength] = removeChunk(recSize, processedBytes, totalLength, chunks);}}}/*** Indexed lookup — reads records directly by offset from the index.* Uses binary search on the index file for O(log n) lookups.* @private*/async *_indexedLookup(field, value, mode = 'record') {const entries = await this._indexManager.get(field, value);if (entries.length === 0) return;const deletedSet = new Set(this._meta.deleted);const fileHandle = await open(this._dataPath, 'r');try {for (const entry of entries) {const [offset, length] = entry.loc;if (deletedSet.has(offset)) continue;const buffer = Buffer.alloc(length);await fileHandle.read(buffer, 0, length, offset);switch (mode) {case 'raw':yield buffer;break;case 'mixed': {const data = deserialize(buffer.subarray(4));const rec = this._classToUse? Object.assign(new this._classToUse(), data): data;yield [rec, offset, length];break;}case 'record':default: {const data = deserialize(buffer.subarray(4));const rec = this._classToUse? Object.assign(new this._classToUse(), data): data;yield rec;break;}}}} finally {await fileHandle.close();}}/*** Collect offsets from one or more indexes, optionally intersect, then stream records.* @param {Object|Array} index - Single index hint or array of hints* @param {function|null} queryFn - Optional filter function* @param {string} mode - Output mode: 'record', 'raw', or 'mixed'* @private*/async *_indexedStream(index, queryFn, mode = 'record') {const hints = Array.isArray(index) ? index : [index];const deletedSet = new Set(this._meta.deleted);// Collect offset sets from each index hintconst offsetSets = [];for (const hint of hints) {const offsets = new Map(); // offset → [offset, length]if (hint.value !== undefined) {// Exact match via get()const entries = await this._indexManager.get(hint.field, hint.value);for (const e of entries) {if (!deletedSet.has(e.loc[0])) offsets.set(e.loc[0], e.loc);}} else {// Range scan via entries()for await (const e of this._indexManager.entries(hint.field, { from: hint.from, to: hint.to, direction: hint.direction })) {if (!deletedSet.has(e.loc[0])) offsets.set(e.loc[0], e.loc);}}offsetSets.push(offsets);}// Intersect: keep only offsets present in ALL setslet locations;if (offsetSets.length === 1) {locations = Array.from(offsetSets[0].values());} else {// Start with smallest set for efficiencyoffsetSets.sort((a, b) => a.size - b.size);const [smallest, ...rest] = offsetSets;locations = [];for (const [offset, loc] of smallest) {if (rest.every(s => s.has(offset))) locations.push(loc);}}// Stream records from the intersected locationsconst fileHandle = await open(this._dataPath, 'r');try {for (const [offset, length] of locations) {const buffer = Buffer.alloc(length);await fileHandle.read(buffer, 0, length, offset);if (mode === 'raw') {yield buffer;} else {const data = deserialize(buffer.subarray(4));const rec = this._classToUse? Object.assign(new this._classToUse(), data): data;if (queryFn && !queryFn(rec)) continue;if (mode === 'mixed') {yield [rec, offset, length];} else {yield rec;}}}} finally {await fileHandle.close();}}/*** Re-read meta.json from disk so this instance sees changes made by other processes.* Called automatically before every read operation.*/async refresh() {try {this._meta = JSON.parse(await readFile(this._metaPath));} catch (e) {// file missing or corrupt — keep current meta}}async persistMeta() {const metaToSave = { ...this._meta };// Only save nextId if we have a numeric primary keyif (this._primaryKeyType !== PrimaryKeyType.NUMBER) {delete metaToSave.nextId;}await writeFile(this._metaPath, JSON.stringify(metaToSave));}/*** Compact the database by removing deleted records* This rewrites the data file without tombstones** @returns {Promise<void>}*/async compact() {await this._acquireLock('compact');try {const writeStream = createWriteStream(this._dataPath + '.tmp', { flags: 'w' });for await (const binary of this.find(null, { mode: 'raw' })) {writeStream.write(binary);}await new Promise((resolve, reject) => {writeStream.end((err) => err ? reject(err) : resolve());});// rename is atomic, so we can just rename the file and it will replace the old oneawait rename(this._dataPath + '.tmp', this._dataPath);this._meta.deleted = [];await this.persistMeta();// Rebuild indexes — offsets changed after compactionif (this._indexManager) {await this._indexManager.init(this._dataPath, { forceRebuild: true });}} finally {await this._releaseLock();}}async _acquireLock(operation, record, retries = 0) {if (this._lockDepth > 0) {this._lockDepth++;return;}try {await writeFile(this._lockPath, String(process.pid), { flag: 'wx' });this._lockDepth = 1;} catch (e) {if (e.code === 'EEXIST') {this.debug(`Waiting for lock file ${operation} retry #${retries}`, record);await new Promise(resolve => setTimeout(resolve, 100));return this._acquireLock(operation, record, retries + 1);}throw e;}}async _releaseLock() {this._lockDepth--;if (this._lockDepth <= 0) {this._lockDepth = 0;await unlink(this._lockPath).catch(() => { });}}/*** Retrieves a document by its file offset and length* @param {Array} location - [offset, length] tuple* @returns {Promise<Object|null>} The deserialized document or null* @private* @deprecated Currently unused - may be removed in future versions*/async _getDocByLocation([offset, length]) {if (!offset || length === 0) return null;const fileHandle = await open(this._dataPath, 'r');try {const buffer = Buffer.alloc(length);await fileHandle.read(buffer, 0, length, offset);return deserialize(buffer);} finally {await fileHandle.close();}}/*** Close the database and persist all pending changes* Should be called before process exit** @returns {Promise<void>}*/async close() {// Persist indexes before closingif (this._indexManager) {await this._indexManager.close();}// Close data streamif (this._dataStream) {await new Promise((resolve, reject) => {this._dataStream.end((err) => err ? reject(err) : resolve());});}// Remove signal handlersif (this._processExitHandler) {process.off('SIGINT', this._processExitHandler);process.off('SIGTERM', this._processExitHandler);this._processExitHandler = null;}}debug(...args) {if (this._debug) {console.log(...args);}}}/*** Helper function to remove a processed chunk from the buffer* @param {number} recSize - Size of the record to remove* @param {number} processedBytes - Current offset in the file* @param {number} totalLength - Total length of buffered data* @param {Buffer[]} chunks - Array of buffer chunks* @returns {[number, number]} Updated [processedBytes, totalLength]*/function removeChunk(recSize, processedBytes, totalLength, chunks) {processedBytes += recSize;totalLength -= recSize;let bytesToRemove = recSize;while (bytesToRemove > 0 && chunks.length > 0) {const currentChunk = chunks[0];if (bytesToRemove >= currentChunk.length) {bytesToRemove -= currentChunk.length;chunks.shift();} else {chunks[0] = currentChunk.subarray(bytesToRemove);bytesToRemove = 0;}}return [processedBytes, totalLength];}export default MPackDB;
Branches
- mastermain branch
Latest commits
- 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