1- import { promises as fsp } from 'node:fs' ;
2-
31import { join } from 'pathe' ;
42
5- import { type QueryOptions } from '@pymodel/minidb' ;
6- import { ClusterDb } from '@pymodel/minidb/cluster' ;
3+ import { classifyStorageError , type QueryOptions } from '@pymodel/minidb' ;
4+ import { ClusterDb , wipeCluster } from '@pymodel/minidb/cluster' ;
75
86import { Disposable , toDisposable } from '#/_base/di/lifecycle' ;
97import { LifecycleScope } from '#/app/scopes' ;
@@ -12,6 +10,7 @@ import { ILogService } from '#/_base/log/log';
1210import { IBootstrapService } from '#/app/bootstrap/bootstrap' ;
1311import {
1412 IQueryStore ,
13+ QueryStoreRebuiltError ,
1514 type Checkpoint ,
1615 type ColumnBounds ,
1716 type ColumnPageQuery ,
@@ -29,6 +28,7 @@ const STORE_SUBDIR = 'query-store';
2928const SHARD_COUNT = 16 ;
3029const LOCK_ACQUIRE_TIMEOUT_MS = 1000 ;
3130const DROP_BATCH_SIZE = 500 ;
31+ const TRANSIENT_ESCALATION_LIMIT = 5 ;
3232
3333function physicalKey ( collection : string , key : string ) : string {
3434 return `${ collection } ${ SEP } ${ key } ` ;
@@ -38,10 +38,6 @@ function indexName(collection: string, name: string): string {
3838 return `${ collection } :${ name } ` ;
3939}
4040
41- function isRebuildable ( error : unknown ) : boolean {
42- return error instanceof SyntaxError || ( error as { name ?: string } ) . name === 'CorruptFrameError' ;
43- }
44-
4541const pendingDisposals = new Set < Promise < void > > ( ) ;
4642
4743export async function drainQueryStoreDisposals ( ) : Promise < void > {
@@ -54,6 +50,9 @@ export class MiniDbQueryStore extends Disposable implements IQueryStore {
5450 private readonly dir : string ;
5551 private dbPromise : Promise < ClusterDb > | undefined ;
5652 private rebuildPromise : Promise < void > | undefined ;
53+ private transientReadFailures = 0 ;
54+ private transientWriteFailures = 0 ;
55+ private storeEpochCounter = 0 ;
5756 private readonly ensuredIndexes = new Set < string > ( ) ;
5857
5958 constructor (
@@ -95,7 +94,7 @@ export class MiniDbQueryStore extends Disposable implements IQueryStore {
9594
9695 private rebuild ( cause : unknown ) : Promise < void > {
9796 this . rebuildPromise ??= ( async ( ) => {
98- this . log . warn ( 'minidb query-store rebuilt after corruption ' , {
97+ this . log . warn ( 'minidb query-store rebuilt after unrecoverable failure ' , {
9998 dir : this . dir ,
10099 error : String ( cause ) ,
101100 } ) ;
@@ -106,18 +105,52 @@ export class MiniDbQueryStore extends Disposable implements IQueryStore {
106105 const db = await previous . catch ( ( ) => undefined ) ;
107106 await db ?. close ( ) . catch ( ( ) => { } ) ;
108107 }
109- await fsp . rm ( this . dir , { recursive : true , force : true } ) ;
108+ const outcome = await wipeCluster ( {
109+ dir : this . dir ,
110+ lockAcquireTimeoutMs : LOCK_ACQUIRE_TIMEOUT_MS ,
111+ } ) ;
112+ if ( outcome === 'locked' ) throw cause ;
113+ this . storeEpochCounter += 1 ;
110114 } ) ( ) ;
111- return this . rebuildPromise ;
115+ const settled = this . rebuildPromise ;
116+ return settled . then (
117+ ( ) => {
118+ if ( this . rebuildPromise === settled ) this . rebuildPromise = undefined ;
119+ } ,
120+ ( error : unknown ) => {
121+ if ( this . rebuildPromise === settled ) this . rebuildPromise = undefined ;
122+ throw error ;
123+ } ,
124+ ) ;
112125 }
113126
114- private async withDb < T > ( op : ( db : ClusterDb ) => Promise < T > ) : Promise < T > {
127+ private async withDb < T > (
128+ op : ( db : ClusterDb ) => Promise < T > ,
129+ kind : 'read' | 'write' ,
130+ expectedStoreEpoch ?: number ,
131+ ) : Promise < T > {
132+ const db = await this . openDb ( ) ;
133+ if ( expectedStoreEpoch !== undefined && expectedStoreEpoch !== this . storeEpochCounter ) {
134+ throw new QueryStoreRebuiltError ( ) ;
135+ }
115136 try {
116- return await op ( await this . openDb ( ) ) ;
137+ const result = await op ( db ) ;
138+ if ( kind === 'write' ) this . transientWriteFailures = 0 ;
139+ else this . transientReadFailures = 0 ;
140+ return result ;
117141 } catch ( error ) {
118- if ( ! isRebuildable ( error ) ) throw error ;
142+ if ( classifyStorageError ( error ) !== 'rebuild' ) {
143+ const failures =
144+ kind === 'write'
145+ ? ( this . transientWriteFailures += 1 )
146+ : ( this . transientReadFailures += 1 ) ;
147+ if ( failures < TRANSIENT_ESCALATION_LIMIT ) throw error ;
148+ }
149+ this . transientReadFailures = 0 ;
150+ this . transientWriteFailures = 0 ;
119151 await this . rebuild ( error ) ;
120- return op ( await this . openDb ( ) ) ;
152+ if ( expectedStoreEpoch !== undefined ) throw new QueryStoreRebuiltError ( ) ;
153+ throw error ;
121154 }
122155 }
123156
@@ -127,41 +160,45 @@ export class MiniDbQueryStore extends Disposable implements IQueryStore {
127160 value : T ,
128161 options ?: { columns ?: Record < string , number > } ,
129162 ) : Promise < void > {
130- await this . withDb ( ( db ) =>
131- db . set ( physicalKey ( collection , key ) , value , { dt : options ?. columns } ) ,
163+ await this . withDb (
164+ ( db ) => db . set ( physicalKey ( collection , key ) , value , { dt : options ?. columns } ) ,
165+ 'write' ,
132166 ) ;
133167 }
134168
135169 async batch ( ops : readonly WriteOp [ ] ) : Promise < void > {
136170 if ( ops . length === 0 ) return ;
137- await this . withDb ( ( db ) =>
138- db . batch (
139- ops . map ( ( op ) =>
140- op . kind === 'put'
141- ? {
142- op : 'set' as const ,
143- key : physicalKey ( op . collection , op . key ) ,
144- value : op . value ,
145- dt : op . columns ,
146- }
147- : { op : 'del' as const , key : physicalKey ( op . collection , op . key ) } ,
171+ await this . withDb (
172+ ( db ) =>
173+ db . batch (
174+ ops . map ( ( op ) =>
175+ op . kind === 'put'
176+ ? {
177+ op : 'set' as const ,
178+ key : physicalKey ( op . collection , op . key ) ,
179+ value : op . value ,
180+ dt : op . columns ,
181+ }
182+ : { op : 'del' as const , key : physicalKey ( op . collection , op . key ) } ,
183+ ) ,
148184 ) ,
149- ) ,
185+ 'write' ,
150186 ) ;
151187 }
152188
153189 async delete ( collection : string , key : string ) : Promise < void > {
154- await this . withDb ( ( db ) => db . del ( physicalKey ( collection , key ) ) ) ;
190+ await this . withDb ( ( db ) => db . del ( physicalKey ( collection , key ) ) , 'write' ) ;
155191 }
156192
157193 async get < T > ( collection : string , key : string ) : Promise < T | undefined > {
158- return this . withDb ( ( db ) => db . get ( physicalKey ( collection , key ) ) as Promise < T | undefined > ) ;
194+ return this . withDb ( ( db ) => db . get ( physicalKey ( collection , key ) ) as Promise < T | undefined > , 'read' ) ;
159195 }
160196
161197 async getMany < T > ( collection : string , keys : readonly string [ ] ) : Promise < Map < string , T > > {
162198 if ( keys . length === 0 ) return new Map ( ) ;
163- const values = await this . withDb ( ( db ) =>
164- db . mget ( keys . map ( ( key ) => physicalKey ( collection , key ) ) ) ,
199+ const values = await this . withDb (
200+ ( db ) => db . mget ( keys . map ( ( key ) => physicalKey ( collection , key ) ) ) ,
201+ 'read' ,
165202 ) ;
166203 const out = new Map < string , T > ( ) ;
167204 values . forEach ( ( value , index ) => {
@@ -172,36 +209,39 @@ export class MiniDbQueryStore extends Disposable implements IQueryStore {
172209
173210 async pageByColumn < T > ( collection : string , query : ColumnPageQuery ) : Promise < Page < T > > {
174211 const dir = query . dir ?? 'asc' ;
175- const rows = ( await this . withDb ( ( db ) =>
176- db . query ( {
177- dt : { [ query . column ] : query . bounds ?? { } } ,
178- filter : query . filter as Record < string , unknown > | undefined ,
179- sort : { [ query . column ] : dir === 'desc' ? - 1 : 1 } ,
180- limit : query . limit ,
181- } ) ,
212+ const rows = ( await this . withDb (
213+ ( db ) =>
214+ db . query ( {
215+ dt : { [ query . column ] : query . bounds ?? { } } ,
216+ filter : query . filter as Record < string , unknown > | undefined ,
217+ sort : { [ query . column ] : dir === 'desc' ? - 1 : 1 } ,
218+ limit : query . limit ,
219+ } ) ,
220+ 'read' ,
182221 ) ) as ReadonlyArray < { value : T } > ;
183222 return { items : rows . map ( ( row ) => row . value ) } ;
184223 }
185224
186225 async listKeys ( collection : string ) : Promise < readonly string [ ] > {
187226 const prefix = `${ collection } ${ SEP } ` ;
188- const entries = await this . withDb ( ( db ) => db . scan ( { prefix } ) ) ;
227+ const entries = await this . withDb ( ( db ) => db . scan ( { prefix } ) , 'read' ) ;
189228 return entries . map ( ( entry ) => entry . key . slice ( prefix . length ) ) ;
190229 }
191230
192231 async dropCollection ( collection : string ) : Promise < void > {
193232 const prefix = `${ collection } ${ SEP } ` ;
194- const entries = await this . withDb ( ( db ) => db . scan ( { prefix } ) ) ;
233+ const entries = await this . withDb ( ( db ) => db . scan ( { prefix } ) , 'read' ) ;
195234 for ( let start = 0 ; start < entries . length ; start += DROP_BATCH_SIZE ) {
196235 const chunk = entries . slice ( start , start + DROP_BATCH_SIZE ) ;
197- await this . withDb ( ( db ) =>
198- db . batch ( chunk . map ( ( entry ) => ( { op : 'del' as const , key : entry . key } ) ) ) ,
236+ await this . withDb (
237+ ( db ) => db . batch ( chunk . map ( ( entry ) => ( { op : 'del' as const , key : entry . key } ) ) ) ,
238+ 'write' ,
199239 ) ;
200240 }
201241 }
202242
203243 query < T > ( collection : string ) : IQuery < T > {
204- return new MiniDbQuery < T > ( ( op ) => this . withDb ( op ) , collection ) ;
244+ return new MiniDbQuery < T > ( ( op ) => this . withDb ( op , 'read' ) , collection ) ;
205245 }
206246
207247 async ensureIndex ( collection : string , def : IndexDef ) : Promise < void > {
@@ -223,16 +263,28 @@ export class MiniDbQueryStore extends Disposable implements IQueryStore {
223263 } catch ( error ) {
224264 if ( ! ( error instanceof Error ) || ! error . message . includes ( 'already exists' ) ) throw error ;
225265 }
226- } ) ;
266+ } , 'write' ) ;
227267 this . ensuredIndexes . add ( guard ) ;
228268 }
229269
230270 async getCheckpoint ( source : string ) : Promise < Checkpoint | undefined > {
231271 return this . get < Checkpoint > ( CHECKPOINT_COLLECTION , source ) ;
232272 }
233273
234- async setCheckpoint ( source : string , checkpoint : Checkpoint ) : Promise < void > {
235- await this . put ( CHECKPOINT_COLLECTION , source , checkpoint ) ;
274+ async setCheckpoint (
275+ source : string ,
276+ checkpoint : Checkpoint ,
277+ expectedStoreEpoch ?: number ,
278+ ) : Promise < void > {
279+ await this . withDb (
280+ ( db ) => db . set ( physicalKey ( CHECKPOINT_COLLECTION , source ) , checkpoint ) ,
281+ 'write' ,
282+ expectedStoreEpoch ,
283+ ) ;
284+ }
285+
286+ storeEpoch ( ) : number {
287+ return this . storeEpochCounter ;
236288 }
237289
238290 async close ( ) : Promise < void > {
0 commit comments