diff --git a/packages/triplestore/_memory_test.ts b/packages/triplestore/_memory_test.ts index 606fba8..7cffab2 100644 --- a/packages/triplestore/_memory_test.ts +++ b/packages/triplestore/_memory_test.ts @@ -18,13 +18,13 @@ export class MemoryFileSystem implements FileSystemType { readonly directories = new Set(['/']) writeFault: WriteFaultType | undefined - async exists(path: string, options: SignalOptionsType = {}): Promise { + exists(path: string, options: SignalOptionsType = {}): Promise { abort(options.signal) const normalized = normalize(path) - return this.files.has(normalized) || this.directories.has(normalized) + return Promise.resolve(this.files.has(normalized) || this.directories.has(normalized)) } - async ensureDir(path: string, options: SignalOptionsType = {}): Promise { + ensureDir(path: string, options: SignalOptionsType = {}): Promise { abort(options.signal) const parts = normalize(path).split('/').filter(Boolean) let current = '' @@ -32,6 +32,7 @@ export class MemoryFileSystem implements FileSystemType { current += `/${part}` this.directories.add(current) } + return Promise.resolve() } async *readDir( @@ -60,22 +61,22 @@ export class MemoryFileSystem implements FileSystemType { } } - async readText(path: string, options: SignalOptionsType = {}): Promise { + readText(path: string, options: SignalOptionsType = {}): Promise { abort(options.signal) const value = this.files.get(normalize(path)) if (value === undefined) throw new Error(`ENOENT ${path}`) - return value + return Promise.resolve(value) } - async stat( + stat( path: string, options: SignalOptionsType = {}, ): Promise { abort(options.signal) const normalized = normalize(path) const file = this.files.get(normalized) - if (file !== undefined) return { kind: 'file', size: new TextEncoder().encode(file).byteLength } - if (this.directories.has(normalized)) return { kind: 'directory' } + if (file !== undefined) return Promise.resolve({ kind: 'file', size: new TextEncoder().encode(file).byteLength }) + if (this.directories.has(normalized)) return Promise.resolve({ kind: 'directory' }) throw new Error(`ENOENT ${path}`) } diff --git a/packages/triplestore/package.json b/packages/triplestore/package.json index 00eee4d..ed85db0 100644 --- a/packages/triplestore/package.json +++ b/packages/triplestore/package.json @@ -17,6 +17,6 @@ "access": "public" }, "dependencies": { - "@okikio/rdf": "0.1.0" + "@okikio/rdf": "workspace:*" } } diff --git a/packages/triplestore/store.ts b/packages/triplestore/store.ts index 54d2e69..0285fae 100644 --- a/packages/triplestore/store.ts +++ b/packages/triplestore/store.ts @@ -329,10 +329,21 @@ export class Store implements AsyncDisposable { const segmentPath = join(this.#path, commit.segment) const commitPath = join(this.#path, `commits/${generationName(commit.generation)}.json`) if (await this.#fs.exists(commitPath, signalOptions(signal))) { - throw new StoreError( - 'writer-conflict', - `Commit generation ${commit.generation} already exists.`, - ) + const text = await this.#fs.readText(commitPath, signalOptions(signal)) + let existing: CommitType | undefined + try { + existing = parseCommit(text) + } catch { + // An interrupted commit-file write can leave invalid JSON or an incomplete + // record at this exact generation. Recovery already classifies that file + // as unpublished, so a live handle may replace the same debris on retry. + } + if (existing !== undefined) { + throw new StoreError( + 'writer-conflict', + `Commit generation ${commit.generation} already exists.`, + ) + } } // Publication order is the durability invariant. A segment may be orphaned, diff --git a/packages/triplestore/store_test.ts b/packages/triplestore/store_test.ts index 14f9843..d3c26c3 100644 --- a/packages/triplestore/store_test.ts +++ b/packages/triplestore/store_test.ts @@ -29,6 +29,39 @@ describe('@okikio/triplestore', () => { expect(reopened.recovery[0]?.kind).toBe('incomplete-commit') }) + it('retries a torn commit publication on the same live handle', async () => { + const fs = new MemoryFileSystem() + const store = await open(fs, { path: '/db' }) + let faulted = false + + fs.writeFault = (path, text, target) => { + if (faulted || !path.endsWith('/commits/0000000000000001.json')) return + faulted = true + target.files.set(path, text.slice(0, Math.max(1, Math.floor(text.length / 2)))) + throw new Error('simulated torn commit') + } + + await expect(store.add(first)).rejects.toThrow('simulated torn commit') + fs.writeFault = undefined + await store.add(first) + + expect(store.generation).toBe(1) + expect(store.has(first)).toBe(true) + const reopened = await open(fs, { path: '/db' }) + expect(reopened.generation).toBe(1) + expect(reopened.has(first)).toBe(true) + }) + + it('keeps a valid same-generation commit as a competing-writer conflict', async () => { + const fs = new MemoryFileSystem() + const left = await open(fs, { path: '/db' }) + const right = await open(fs, { path: '/db' }) + + await left.add(first) + await expect(right.add(second)).rejects.toMatchObject({ kind: 'writer-conflict' }) + expect(right.generation).toBe(0) + }) + it('never converts a corrupt authoritative committed generation into an empty database', async () => { const fs = new MemoryFileSystem() const store = await open(fs, { path: '/db' })