mpackdb
All repositories: gitoria
20.2 KB
import { writeFile, readFile, open, rename, stat } from 'node:fs/promises';import { createWriteStream } from 'node:fs';import { resolve } from 'node:path';import { deserialize } from './mpack.js';import { PrimaryKeyType, IndexType } from './MPackDB.js';const BLOCK_SIZE = 4096; // 4KB/*** Manages indexes for the database (in-memory deltas + on-disk persistence)*/export class IndexManager {_dbDir;_dbName;_indexes;_indexTypes;_primaryKeyType;_persistIntervalId = null;_indexPaths = {};_deltaIndexes = {};_tombstones = {};_persistThreshold = 1000;_totalDeltaCount = 0;_coveredBytes = 0;_idxStatePath = null;_catchUpPromise = null;_runExclusive = fn => fn();/*** Creates a new IndexManager* @param {string} dbDir - Database directory path* @param {string} dbName - Database name* @param {string[]} indexes - Array of field names to index* @param {Object} indexTypes - Map of field names to IndexType* @param {number} primaryKeyType - Primary key type* @param {Object} [options]* @param {Function} [options.runExclusive] - Wraps index file writes in the db's lock*/constructor(dbDir, dbName, indexes, indexTypes, primaryKeyType, { runExclusive } = {}) {this._dbDir = dbDir;this._dbName = dbName;this._indexes = indexes;this._indexTypes = indexTypes;this._primaryKeyType = primaryKeyType;if (runExclusive) this._runExclusive = runExclusive;// Index paths are deterministic and do not depend on async initialization.// Populate them immediately, then repair a missing slot on access so a// declared index can never leak readFile(undefined) to the caller.for (const field of this._indexes) this._getIndexPath(field);}_getIndexPath(field) {if (!this._indexes.includes(field)) {const configured = this._indexes.length > 0 ? this._indexes.join(', ') : '(none)';const error = new Error(`No index configured for field "${String(field)}". Configured indexes: ${configured}`);error.code = 'INDEX_NOT_FOUND';error.field = field;error.indexes = [...this._indexes];throw error;}const expectedPath = resolve(this._dbDir, `${this._dbName}.${field}.txt`);if (this._indexPaths[field] !== expectedPath) {this._indexPaths[field] = expectedPath;}return expectedPath;}/*** Initializes the index manager (rebuilds indexes from data file)* @param {string} dataPath - Path to the data file* @returns {Promise<void>}*/async init(dataPath, { forceRebuild = false } = {}) {this._idxStatePath = resolve(this._dbDir, `${this._dbName}.idxstate.json`);let needRebuild = forceRebuild;for (const field of this._indexes) {this._deltaIndexes[field] = [];this._tombstones[field] = new Map();const indexPath = this._getIndexPath(field);if (needRebuild) continue;try {const stats = await stat(indexPath);if (stats.size === 0) {const dataStats = await stat(dataPath).catch(() => ({ size: 0 }));if (dataStats.size > 0) needRebuild = true;}} catch (e) {if (e.code === 'ENOENT') {const dataStats = await stat(dataPath).catch(() => ({ size: 0 }));if (dataStats.size > 0) needRebuild = true;} else {throw e;}}}// coveredBytes: how far into the data file the on-disk indexes reach.// Without it we cannot trust existing index files (another process may// have appended records it never persisted) — rebuild once to establish it.if (!needRebuild) {const state = await this._readIdxState();if (state === null) {needRebuild = true;} else {this._coveredBytes = state.coveredBytes || 0;}}if (needRebuild) {for (const field of this._indexes) {await this._rebuildIndex(dataPath, field);}this._coveredBytes = (await stat(dataPath).catch(() => ({ size: 0 }))).size;await this._writeIdxState();}}async _readIdxState() {try {return JSON.parse(await readFile(this._idxStatePath, 'utf-8'));} catch (e) {return null;}}async _writeIdxState() {const tmpPath = `${this._idxStatePath}.${process.pid}.tmp`;await writeFile(tmpPath, JSON.stringify({ coveredBytes: this._coveredBytes }), 'utf-8');await rename(tmpPath, this._idxStatePath);}/*** Index records another process appended to the data file since our last* known coverage point. Called from refresh(). Append-only data makes this* a cheap incremental scan of just the new tail.* @param {string} dataPath - Path to the data file* @param {number} fileSize - Current data file size (from the caller's stat)*/async catchUp(dataPath, fileSize) {if (this._idxStatePath === null) return; // not initialized yetif (fileSize <= this._coveredBytes) return;if (this._catchUpPromise) return this._catchUpPromise;this._catchUpPromise = (async () => {const from = this._coveredBytes;for await (const { doc, loc } of this._getDocsAndLocationsFromDataFile(dataPath, from)) {for (const field of this._indexes) {if (doc && doc[field] !== undefined) {this._deltaIndexes[field].push({ key: doc[field], loc });this._totalDeltaCount++;}}this._coveredBytes = Math.max(this._coveredBytes, loc[0] + loc[1]);}if (this._totalDeltaCount >= this._persistThreshold) {this.persist().catch(err => console.error('Index persist error:', err));}})();try {return await this._catchUpPromise;} finally {this._catchUpPromise = null;}}/*** The data file was replaced (external compaction): all offsets are invalid.* Drop delta/tombstone state and re-baseline from the compactor's rebuilt* index files, or rebuild if its state file is unreadable.* @param {string} dataPath - Path to the (new) data file*/async resetAfterReplace(dataPath) {for (const field of this._indexes) {this._deltaIndexes[field] = [];this._tombstones[field].clear();}this._totalDeltaCount = 0;const state = await this._readIdxState();if (state !== null) {this._coveredBytes = state.coveredBytes || 0;} else {for (const field of this._indexes) {await this._rebuildIndex(dataPath, field);}this._coveredBytes = (await stat(dataPath).catch(() => ({ size: 0 }))).size;await this._writeIdxState();}}/*** Rebuilds an index from the data file* @param {string} dataPath - Path to the data file* @param {string} field - Field name to index* @returns {Promise<void>}* @private*/async _rebuildIndex(dataPath, field) {const indexPath = this._getIndexPath(field);const tempIndexPath = `${indexPath}.${Date.now()}.tmp`;const writeStream = createWriteStream(tempIndexPath, { flags: 'w', encoding: 'utf-8' });const docIterator = this._getDocsAndLocationsFromDataFile(dataPath);const entries = [];for await (const { doc, loc } of docIterator) {if (doc && doc[field] !== undefined) {const key = doc[field];entries.push({ key, loc });}}// Sort entries before writing to the index fileentries.sort((a, b) => this._compareKeys(a.key, b.key));for (const entry of entries) {writeStream.write(`${entry.key},${entry.loc[0]},${entry.loc[1]}\n`);}await new Promise(resolve => writeStream.end(resolve));await rename(tempIndexPath, indexPath);}/*** Generator that yields documents and their locations from the data file* @param {string} dataPath - Path to the data file* @yields {{doc: Object, loc: Array}} Document and location tuple* @private*/async* _getDocsAndLocationsFromDataFile(dataPath, startOffset = 0) {let fileHandle;try {fileHandle = await open(dataPath, 'r');const stats = await fileHandle.stat();let offset = startOffset;while (offset < stats.size) {// Read size headerconst sizeBuffer = Buffer.alloc(4);await fileHandle.read(sizeBuffer, 0, 4, offset);const size = sizeBuffer.readInt32LE(0);if (size <= 0 || offset + size > stats.size) {break; // Corrupted or incomplete record}const docBuffer = Buffer.alloc(size);await fileHandle.read(docBuffer, 0, size, offset);const doc = deserialize(docBuffer.subarray(4));yield { doc, loc: [offset, size] };offset += size;}} catch (e) {if (e.code !== 'ENOENT') {throw e;}// if file does not exist, do nothing} finally {await fileHandle?.close();}}/*** Compares two keys for sorting* @param {any} keyA - First key* @param {any} keyB - Second key* @returns {number} -1 if keyA < keyB, 0 if equal, 1 if keyA > keyB* @private*/_compareKeys(keyA, keyB) {if (typeof keyA === 'number' && typeof keyB === 'number') {return keyA - keyB;}return String(keyA).localeCompare(String(keyB));}/*** Starts automatic periodic persistence of indexes* @param {number} interval - Interval in milliseconds between auto-persists* @param {number} threshold - Number of operations before triggering auto-persist*/startAutoPersist(interval, threshold) {this._persistThreshold = threshold;// Start periodic persistif (interval > 0) {this._persistIntervalId = setInterval(async () => {if (this._totalDeltaCount > 0) {await this.persist();}}, interval);// Don't keep process alive just for this timerthis._persistIntervalId.unref();}}/*** Closes the index manager (persists indexes and clears interval)* @returns {Promise<void>}*/async close() {if (this._persistIntervalId) {clearInterval(this._persistIntervalId);this._persistIntervalId = null;}await this.persist();}/*** Persists all in-memory index deltas to disk.* Runs under the db lock (via runExclusive): the read-merge-write below* would otherwise lose entries another process persisted in between.* @returns {Promise<void>}*/async persist() {return this._runExclusive(() => this._persistLocked());}async _persistLocked() {for (const field of this._indexes) {let onDiskEntries = [];const indexPath = this._getIndexPath(field);try {const indexContent = await readFile(indexPath, 'utf-8');onDiskEntries = indexContent.trim().split('\n').filter(Boolean).map(line => {const [key, offset, length] = line.split(',');return { key: this._parseKey(key, field), loc: [parseInt(offset), parseInt(length)] };});} catch (e) {if (e.code !== 'ENOENT') throw e;}const fieldTombstones = this._tombstones[field];const validDiskEntries = onDiskEntries.filter(e => !fieldTombstones.has(e.loc[0]));// Dedupe by location — catch-up scans and cross-process persists can// both have picked up the same recordconst merged = new Map();for (const entry of [...validDiskEntries, ...(this._deltaIndexes[field] || [])]) {merged.set(`${entry.loc[0]}:${entry.loc[1]}`, entry);}const finalIndexData = Array.from(merged.values());finalIndexData.sort((a, b) => this._compareKeys(a.key, b.key));if (finalIndexData.length > 0) {const indexContent = finalIndexData.map(e => `${e.key},${e.loc[0]},${e.loc[1]}`).join('\n') + '\n';await writeFile(indexPath, indexContent, 'utf-8');} else {await writeFile(indexPath, '', 'utf-8');}// Clear the in-memory changes now that they are persistedthis._deltaIndexes[field] = [];this._tombstones[field].clear();}this._totalDeltaCount = 0;// Coverage can only grow: another process's state may reach further than oursconst state = await this._readIdxState();this._coveredBytes = Math.max(this._coveredBytes, state?.coveredBytes || 0);await this._writeIdxState();}/*** Parses a key based on the field's index type* @param {any} key - The key to parse* @param {string} field - Field name* @returns {any} Parsed key (number or string)* @private*/_parseKey(key, field) {const indexType = this._indexTypes[field];// Use index type if specified, otherwise fall back to primary key type logicif (indexType === IndexType.NUMERIC) {const num = parseInt(key, 10);if (!isNaN(num)) return num;} else if (indexType === IndexType.LEXICAL) {return key;}// Fallback to primary key type for fields without explicit index typeif (this._primaryKeyType === PrimaryKeyType.NUMBER) {const num = parseInt(key, 10);if (!isNaN(num)) return num;}return key;}_parseIndexLines(indexContent, field) {return indexContent.trim().split('\n').filter(Boolean).map(line => {const [keyStr, offset, length] = line.split(',');const locOffset = Number.parseInt(offset, 10);const locLength = Number.parseInt(length, 10);if (!Number.isSafeInteger(locOffset) || !Number.isSafeInteger(locLength) || locLength <= 0) {return null;}return {key: this._parseKey(keyStr, field),loc: [locOffset, locLength],};}).filter(Boolean);}async _readIndexEntries(field) {const indexPath = this._getIndexPath(field);try {return this._parseIndexLines(await readFile(indexPath, 'utf-8'), field);} catch (e) {if (e.code === 'ENOENT') return [];throw e;}}/*** Performs a binary search on the blocks of an index file on disk.* This is a direct port of Joshua Bloch's famously correct binary search* algorithm, adapted for file blocks.** @param {string} field The index field to search.* @param {string|number} key The key to search for.* @returns {Promise<number>} The starting offset in the file for the linear scan.*/async _binarySearchOnDisk(field, key) {return (await this._readIndexEntries(field)).length === 0 ? -1 : 0;}/*** Finds the last entry for a given field and key (checks both delta and disk)* @param {string} field - Field name* @param {any} key - Key value to search for* @returns {Promise<Object|null>} Entry object with key and loc, or null if not found*/async findLastEntry(field, key) {const parsedKey = this._parseKey(key, field);const deltaIndex = this._deltaIndexes[field] || [];for (let i = deltaIndex.length - 1; i >= 0; i--) {const entry = deltaIndex[i];if (this._compareKeys(entry.key, parsedKey) === 0) {return entry;}}const onDiskEntry = await this._findLastEntryOnDisk(field, parsedKey);if (!onDiskEntry) return null;const isTombstoned = this._tombstones[field].has(onDiskEntry.loc[0]);return isTombstoned ? null : onDiskEntry;}/*** Finds the last entry for a key on disk (linear scan from binary search position)* @param {string} field - Field name* @param {any} key - Key to search for* @returns {Promise<Object|null>} Entry object or null* @private*/async _findLastEntryOnDisk(field, key) {let lastMatch = null;for (const entry of await this._readIndexEntries(field)) {if (this._compareKeys(entry.key, key) === 0) {lastMatch = entry;}}return lastMatch;}/*** Adds a document to the indexes* @param {Object} doc - The document to index* @param {Array} loc - [offset, length] location in data file*/insert(doc, loc) {for (const field of this._indexes) {if (doc[field] !== undefined) {this._deltaIndexes[field].push({ key: doc[field], loc });this._totalDeltaCount++;}}this._coveredBytes = Math.max(this._coveredBytes, loc[0] + loc[1]);// Auto-persist if threshold reachedif (this._totalDeltaCount >= this._persistThreshold) {// Persist async without blockingthis.persist().catch(err => console.error('Index persist error:', err));}}/*** Removes a document from the indexes* @param {Object} doc - The document to remove* @param {string} primaryKeyField - Primary key field name* @returns {Promise<void>}*/async remove(doc, primaryKeyField) {const pkValue = doc[primaryKeyField];const deltaPkIndex = this._deltaIndexes[primaryKeyField] || [];const indexInDelta = deltaPkIndex.findIndex(e => this._compareKeys(e.key, pkValue) === 0);if (indexInDelta !== -1) {const offsetToRemove = deltaPkIndex[indexInDelta].loc[0];for (const field of this._indexes) {this._deltaIndexes[field] = (this._deltaIndexes[field] || []).filter(e => e.loc[0] !== offsetToRemove);}} else {const onDiskEntry = await this._findLastEntryOnDisk(primaryKeyField, pkValue);if (onDiskEntry) {for (const field of this._indexes) {if (doc[field] !== undefined) {this._tombstones[field].set(onDiskEntry.loc[0], true);}}}}}/*** Gets all document locations from the first index (used for full scans)* @returns {Promise<Array[]>} Array of all [offset, length] locations*/async getAllLocations() {const field = this._indexes[0];const indexPath = this._getIndexPath(field);let onDiskEntries = [];try {const indexContent = await readFile(indexPath, 'utf-8');onDiskEntries = indexContent.trim().split('\n').filter(Boolean).map(line => {const [, offset, length] = line.split(',');return { loc: [parseInt(offset), parseInt(length)] };});} catch (e) {if (e.code !== 'ENOENT') throw e;}const fieldTombstones = this._tombstones[field] || new Map();const validDiskEntries = onDiskEntries.filter(e => !fieldTombstones.has(e.loc[0]));const deltaEntries = this._deltaIndexes[field] || [];const finalEntries = new Map();for (const entry of validDiskEntries) {finalEntries.set(`${entry.loc[0]}:${entry.loc[1]}`, entry.loc);}for (const entry of deltaEntries) {finalEntries.set(`${entry.loc[0]}:${entry.loc[1]}`, entry.loc);}return Array.from(finalEntries.values());}/*** Gets document locations for a given field and key* @param {string} field - Field name* @param {any} key - Key value to search for* @returns {Promise<Array[]>} Array of [offset, length] locations*/async get(field, key) {const parsedKey = this._parseKey(key, field);const onDiskResults = await this._getOnDisk(field, parsedKey);const deltaResults = (this._deltaIndexes[field] || []).filter(e => this._compareKeys(e.key, parsedKey) === 0);const tombstonedOffsets = this._tombstones[field] || new Map();const finalEntries = new Map();for (const entry of onDiskResults) {if (!tombstonedOffsets.has(entry.loc[0])) {finalEntries.set(`${entry.loc[0]}:${entry.loc[1]}`, entry);}}for (const entry of deltaResults) {finalEntries.set(`${entry.loc[0]}:${entry.loc[1]}`, entry);}return Array.from(finalEntries.values());}/*** Gets all entries for a key from disk* @param {string} field - Field name* @param {any} key - Key to search for* @returns {Promise<Array>} Array of entry objects* @private*/async _getOnDisk(field, key) {const results = [];for (const entry of await this._readIndexEntries(field)) {if (this._compareKeys(entry.key, key) === 0) {results.push(entry);}}return results;}/*** Yields index entries in sorted order by streaming blocks from disk.* Merges with in-memory deltas, excludes tombstones.* @param {string} field - Field name* @param {Object} [options]* @param {any} [options.from] - Start key (inclusive), uses binary search to skip ahead* @param {any} [options.to] - End key (inclusive), stops scan when exceeded* @param {'asc'|'desc'} [options.direction='asc'] - Scan direction* @yields {{key: any, loc: [number, number]}}*/async *entries(field, { from, to, direction = 'asc' } = {}) {const desc = direction === 'desc';const tombstones = this._tombstones[field] || new Map();const seen = new Set();const withinRange = entry => {if (from !== undefined && this._compareKeys(entry.key, from) < 0) return false;if (to !== undefined && this._compareKeys(entry.key, to) > 0) return false;return true;};const entries = [...(await this._readIndexEntries(field)).filter(e => !tombstones.has(e.loc[0])),...(this._deltaIndexes[field] || []).filter(e => !tombstones.has(e.loc[0])),].filter(withinRange).sort((a, b) => this._compareKeys(a.key, b.key));if (desc) entries.reverse();for (const entry of entries) {const locKey = `${entry.loc[0]}:${entry.loc[1]}`;if (seen.has(locKey)) continue;seen.add(locKey);yield entry;}}async* _getDocsByLocation(dataPath, locations) {const fileHandle = await open(dataPath, 'r');try {for (const loc of locations) {if (loc && loc.length === 2 && loc[1] > 0) {const [offset, length] = loc;const buffer = Buffer.alloc(length);await fileHandle.read(buffer, 0, length, offset);yield deserialize(buffer.subarray(4));}}} finally {await fileHandle.close();}}}export default IndexManager;
Branches
- mastermain branch
Latest commits
- 87888725release 1.0.7caramboleyo
- c4cdb9b6node: import prefixes (Deno compat) + pre-existing index-state WIPcaramboleyo
- 0afb8f4bupdate now must be a callbackcaramboleyo
- cde73eb4release 1.0.6caramboleyo
- 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