diff --git a/src/datastore/LMDBStoreFactory.ts b/src/datastore/LMDBStoreFactory.ts index 4764e912..cc291b7e 100644 --- a/src/datastore/LMDBStoreFactory.ts +++ b/src/datastore/LMDBStoreFactory.ts @@ -15,6 +15,7 @@ import { encryptionStrategy } from './lmdb/Utils'; const MetricsIntervalMs = 60 * 1000; const CleanupDelayMs = 2 * 60 * 1000; +const CorruptionEscalationWindowMs = 30 * 1000; export class LMDBStoreFactory implements DataStoreFactory { private readonly log = LoggerFactory.getLogger('LMDB.Global'); @@ -29,6 +30,7 @@ export class LMDBStoreFactory implements DataStoreFactory { private initPromise?: Promise; private closed = false; private activeOps = 0; + private lastCorruptionAt = 0; private readonly stores = new Map(); private readonly ownershipTracker: LMDBOwnershipTracker; @@ -189,6 +191,8 @@ export class LMDBStoreFactory implements DataStoreFactory { try { if (msg.includes('MDB_BAD_RSLOT') || msg.includes("doesn't match env pid")) { this.recoverFromFork(); + } else if (isCorruptionError(msg)) { + this.recoverFromCorruption(); } else { this.recoverFromError(); } @@ -230,6 +234,20 @@ export class LMDBStoreFactory implements DataStoreFactory { } } + private recoverFromCorruption(): void { + const now = Date.now(); + const recurring = now - this.lastCorruptionAt < CorruptionEscalationWindowMs; + this.lastCorruptionAt = now; + this.telemetry.count('corruption.detected', 1); + + if (recurring) { + this.telemetry.count('corruption.recreate', 1); + this.deleteAndRecreate(); + } else { + this.recoverFromError(); + } + } + private recoverFromError(): void { this.telemetry.count('error.recover', 1); this.log.warn('Error detected, attempting to reopen LMDB'); @@ -385,3 +403,13 @@ function createEnv(lmdbDir: string) { function createDB(env: RootDatabase, name: string) { return env.openDB({ name, encoding: Encoding }); } + +function isCorruptionError(msg: string): boolean { + return ( + msg.includes('MDB_CORRUPTED') || + msg.includes('MDB_PAGE_NOTFOUND') || + msg.includes('MDB_BAD_VALSIZE') || + msg.includes('MDB_CURSOR_FULL') || + msg.includes('MDB_BAD_TXN') + ); +} diff --git a/src/datastore/lmdb/LMDBStore.ts b/src/datastore/lmdb/LMDBStore.ts index 00c7be3d..aa0f0f6b 100644 --- a/src/datastore/lmdb/LMDBStore.ts +++ b/src/datastore/lmdb/LMDBStore.ts @@ -54,6 +54,10 @@ export class LMDBStore implements DataStore { this.validateDatabase(); return await fn(); } catch (e) { + const commitError = (e as { commitError?: unknown }).commitError; + if (commitError instanceof Promise) { + commitError.catch((cause) => this.onError(cause)); + } this.onError(e); this.telemetry.count(`retry.${op}`, 1); this.validateDatabase(); diff --git a/tst/unit/datastore/LMDB.corruption.test.ts b/tst/unit/datastore/LMDB.corruption.test.ts index 033c2037..4833d80d 100644 --- a/tst/unit/datastore/LMDB.corruption.test.ts +++ b/tst/unit/datastore/LMDB.corruption.test.ts @@ -92,3 +92,55 @@ describe('LMDB error recovery', () => { recreateSpy.mockRestore(); }); }); + +describe('LMDB corruption escalation', () => { + let testDir: string; + let factory: LMDBStoreFactory; + + beforeEach(async () => { + testDir = join( + process.cwd(), + 'node_modules', + '.cache', + 'lmdb-corruption-escalation-test', + `test-${Date.now()}`, + ); + fs.mkdirSync(testDir, { recursive: true }); + factory = new LMDBStoreFactory(testDir); + await factory.initialize(); + }); + + afterEach(async () => { + await factory.close(); + await new Promise((resolve) => setTimeout(resolve, 100)); + if (fs.existsSync(testDir)) { + fs.rmSync(testDir, { recursive: true, force: true }); + } + }); + + it('should reopen without deleting on the first corruption error', () => { + const recoverSpy = vi.spyOn(factory as any, 'recoverFromError').mockImplementation(() => {}); + const deleteSpy = vi.spyOn(factory as any, 'deleteAndRecreate').mockImplementation(() => {}); + + (factory as any).handleError(new Error('MDB_CORRUPTED: Located page was wrong type')); + + expect(recoverSpy).toHaveBeenCalledTimes(1); + expect(deleteSpy).not.toHaveBeenCalled(); + recoverSpy.mockRestore(); + deleteSpy.mockRestore(); + }); + + it('should delete and recreate when corruption recurs within the escalation window', () => { + const recoverSpy = vi.spyOn(factory as any, 'recoverFromError').mockImplementation(() => {}); + const deleteSpy = vi.spyOn(factory as any, 'deleteAndRecreate').mockImplementation(() => {}); + + const handleError = (factory as any).handleError.bind(factory); + handleError(new Error('MDB_CORRUPTED: Located page was wrong type')); + handleError(new Error('MDB_CORRUPTED: Located page was wrong type')); + + expect(recoverSpy).toHaveBeenCalledTimes(1); + expect(deleteSpy).toHaveBeenCalledTimes(1); + recoverSpy.mockRestore(); + deleteSpy.mockRestore(); + }); +}); diff --git a/tst/unit/datastore/LMDB.retry.test.ts b/tst/unit/datastore/LMDB.retry.test.ts index 3a558513..68532cea 100644 --- a/tst/unit/datastore/LMDB.retry.test.ts +++ b/tst/unit/datastore/LMDB.retry.test.ts @@ -141,4 +141,30 @@ describe('LMDB retry after recovery', () => { expect(store.get('key2')).toBe('value2'); }); }); + + describe('lmdb commit rejection handling', () => { + it('should consume the commitError promise and route its cause to recovery', async () => { + const store = factory.get(StoreName.public_schemas); + const realStore = (store as any).store; + + const rootCause = new Error('MDB_CORRUPTED: Located page was wrong type'); + let rejectCommit!: (reason: unknown) => void; + const commitError = new Promise((_resolve, reject) => { + rejectCommit = reject; + }); + const wrapper = Object.assign(new Error('Commit failed (see commitError for details)'), { commitError }); + + const handleErrorSpy = vi.spyOn(factory as any, 'handleError').mockImplementation(() => {}); + realStore.put = () => Promise.reject(wrapper); + + const put = store.put('key', 'value'); + await expect(put).rejects.toBe(wrapper); + + rejectCommit(rootCause); + await new Promise((resolve) => setImmediate(resolve)); + + expect(handleErrorSpy).toHaveBeenCalledWith(rootCause); + handleErrorSpy.mockRestore(); + }); + }); });