diff --git a/packages/triplestore/.npmignore b/packages/triplestore/.npmignore new file mode 100644 index 0000000..db25063 --- /dev/null +++ b/packages/triplestore/.npmignore @@ -0,0 +1,4 @@ +*_test.ts +*_bench.ts +_memory_test.ts +*.map diff --git a/packages/triplestore/README.md b/packages/triplestore/README.md new file mode 100644 index 0000000..976415b --- /dev/null +++ b/packages/triplestore/README.md @@ -0,0 +1,48 @@ +# `@okikio/triplestore` + +Crash-recoverable persistent RDF Dataset backed by a borrowed structural filesystem. + +```ts +import * as rdf from '@okikio/rdf' +import { open } from '@okikio/triplestore' + +await using store = await open(fileSystem, { path: '/knowledge' }) + +await store.add(rdf.quad( + rdf.namedNode('urn:product:1'), + rdf.namedNode('https://schema.org/name'), + rdf.literal('Widget'), +)) +``` + +## Storage capability + +The package does not depend on OPFS directly. It accepts a small filesystem contract so the caller can provide `@okikio/opfs` or another compatible filesystem. + +The filesystem remains caller-owned. Closing the store does not dispose it. + +## Persistence model + +The baseline uses immutable N-Quads segments plus immutable generation records. A generation becomes recoverable only after its commit record and referenced segment can be validated. + +```text +write segment + | + v +publish commit + | + v +new generation visible on reopen +``` + +Recovery scans newest to oldest for a valid committed generation. An interrupted commit publication is ignored. A corrupt newest generation can fall back to an older valid committed generation, but a corrupt only-generation is never treated as an empty database. + +`compact()` writes a snapshot generation so older delta segments are no longer required to reconstruct that state. + +## Current limits + +- one writer per store path +- no cross-process writer coordination yet +- no physical garbage collection contract yet +- in-memory exact-term indexes are rebuilt on open +- persistent-index formats are intentionally deferred until benchmark and durability comparisons justify one diff --git a/packages/triplestore/_memory_test.ts b/packages/triplestore/_memory_test.ts new file mode 100644 index 0000000..f5e7364 --- /dev/null +++ b/packages/triplestore/_memory_test.ts @@ -0,0 +1,96 @@ +import type { DirectoryEntryType, FileStatType, FileSystemType, SignalOptionsType } from './storage.ts' + +/** Injected write fault used by durability tests to simulate a crash during publication. */ +export type WriteFaultType = (path: string, text: string, fs: MemoryFileSystem) => void | Promise + +/** Minimal in-memory filesystem for package-local durability tests. */ +export class MemoryFileSystem implements FileSystemType { + readonly files = new Map() + readonly directories = new Set(['/']) + writeFault: WriteFaultType | undefined + + async exists(path: string, options: SignalOptionsType = {}): Promise { + abort(options.signal) + const normalized = normalize(path) + return this.files.has(normalized) || this.directories.has(normalized) + } + + async ensureDir(path: string, options: SignalOptionsType = {}): Promise { + abort(options.signal) + const parts = normalize(path).split('/').filter(Boolean) + let current = '' + for (const part of parts) { + current += `/${part}` + this.directories.add(current) + } + } + + async *readDir(path: string, options: SignalOptionsType = {}): AsyncGenerator { + abort(options.signal) + const root = normalize(path) + const prefix = root === '/' ? '/' : `${root}/` + const names = new Map() + + for (const file of this.files.keys()) { + if (!file.startsWith(prefix)) continue + const rest = file.slice(prefix.length) + if (rest && !rest.includes('/')) names.set(rest, 'file') + } + for (const directory of this.directories) { + if (!directory.startsWith(prefix) || directory === root) continue + const rest = directory.slice(prefix.length) + if (rest && !rest.includes('/') && !names.has(rest)) names.set(rest, 'directory') + } + + for (const [name, kind] of [...names].sort(([a], [b]) => a.localeCompare(b))) { + abort(options.signal) + yield { name, kind } + } + } + + async 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 + } + + async 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' } + throw new Error(`ENOENT ${path}`) + } + + async writeFile( + path: string, + data: string | Uint8Array, + options: SignalOptionsType & { readonly parents?: boolean } = {}, + ): Promise { + abort(options.signal) + const normalized = normalize(path) + if (options.parents) await this.ensureDir(dirname(normalized), options) + const text = typeof data === 'string' ? data : new TextDecoder().decode(data) + if (this.writeFault) await this.writeFault(normalized, text, this) + this.files.set(normalized, text) + } +} + +/** Normalizes the intentionally small absolute virtual path model used by this fixture. */ +function normalize(path: string): string { + if (!path.startsWith('/')) throw new TypeError('path must be absolute') + return path.length > 1 ? path.replace(/\/+$/g, '') : path +} + +/** Returns the parent path for the fixture's absolute path representation. */ +function dirname(path: string): string { + const index = path.lastIndexOf('/') + return index <= 0 ? '/' : path.slice(0, index) +} + +/** Applies cooperative cancellation consistently across fixture operations. */ +function abort(signal?: AbortSignal): void { + if (signal?.aborted) throw signal.reason ?? new DOMException('Aborted', 'AbortError') +} diff --git a/packages/triplestore/deno.json b/packages/triplestore/deno.json new file mode 100644 index 0000000..4b3cbe2 --- /dev/null +++ b/packages/triplestore/deno.json @@ -0,0 +1,8 @@ +{ + "name": "@okikio/triplestore", + "version": "0.1.0", + "license": "MIT", + "exports": { + ".": "./mod.ts" + } +} diff --git a/packages/triplestore/format.ts b/packages/triplestore/format.ts new file mode 100644 index 0000000..26efaaf --- /dev/null +++ b/packages/triplestore/format.ts @@ -0,0 +1,98 @@ +/** Persistent triplestore format records and validation. @module */ + +/** Current on-disk store format. */ +export const FORMAT_VERSION = 2 as const +/** Current immutable segment format. */ +export const SEGMENT_VERSION = 2 as const + +/** Root format marker written once when a store is created. */ +export interface FormatType { + readonly version: typeof FORMAT_VERSION + readonly store: '@okikio/triplestore' + readonly segment: typeof SEGMENT_VERSION +} + +/** One immutable committed generation. */ +export interface CommitType { + readonly version: typeof FORMAT_VERSION + readonly generation: number + readonly parent: number + readonly mode: 'delta' | 'snapshot' + /** Homogeneous mutation encoded by a delta segment. Snapshots omit it. */ + readonly operation?: 'add' | 'delete' + readonly segment: string + readonly checksum: string + readonly quadCount: number +} + +/** Stable format record for newly created stores. */ +export const FORMAT: FormatType = { + version: FORMAT_VERSION, + store: '@okikio/triplestore', + segment: SEGMENT_VERSION, +} + +/** Parses and validates the root format marker. */ +export function parseFormat(text: string): FormatType { + const value = parseRecord(text, 'format') + if (value.version !== FORMAT_VERSION || value.store !== '@okikio/triplestore' || value.segment !== SEGMENT_VERSION) { + throw new TypeError(`Unsupported triplestore format '${JSON.stringify(value)}'.`) + } + return FORMAT +} + +/** Parses and validates one immutable commit record. */ +export function parseCommit(text: string): CommitType { + const value = parseRecord(text, 'commit') + if (value.version !== FORMAT_VERSION) throw new TypeError(`Unsupported commit version '${String(value.version)}'.`) + if (!isInteger(value.generation) || value.generation < 1) throw new TypeError('Commit generation must be a positive safe integer.') + if (!isInteger(value.parent) || value.parent !== value.generation - 1) throw new TypeError('Commit parent must be the preceding generation.') + if (value.mode !== 'delta' && value.mode !== 'snapshot') throw new TypeError(`Unknown commit mode '${String(value.mode)}'.`) + if (typeof value.segment !== 'string' || !/^segments\/[0-9]{16}(?:\.delta)?\.nq$/.test(value.segment)) throw new TypeError('Commit segment path is invalid.') + if (value.mode === 'delta' && value.operation !== 'add' && value.operation !== 'delete') { + throw new TypeError(`Delta commit operation must be 'add' or 'delete'.`) + } + if (value.mode === 'snapshot' && value.operation !== undefined) { + throw new TypeError('Snapshot commit must not declare an operation.') + } + if (typeof value.checksum !== 'string' || !/^sha256:[0-9a-f]{64}$/.test(value.checksum)) throw new TypeError('Commit checksum is invalid.') + if (!isInteger(value.quadCount) || value.quadCount < 0) throw new TypeError('Commit quadCount must be a non-negative safe integer.') + return { + version: FORMAT_VERSION, + generation: value.generation, + parent: value.parent, + mode: value.mode, + ...(value.mode === 'delta' ? { operation: value.operation as 'add' | 'delete' } : {}), + segment: value.segment, + checksum: value.checksum, + quadCount: value.quadCount, + } +} + +/** Serializes a record with stable key ordering and trailing newline. */ +export function writeRecord(value: FormatType | CommitType): string { + return `${JSON.stringify(value)}\n` +} + +/** Formats a generation as a lexicographically sortable file stem. */ +export function generationName(generation: number): string { + if (!isInteger(generation) || generation < 1) throw new RangeError('generation must be a positive safe integer.') + return generation.toString().padStart(16, '0') +} + +/** Parse record from the current source according to the package grammar contract. */ +function parseRecord(text: string, label: string): Record { + let value: unknown + try { + value = JSON.parse(text) + } catch (cause) { + throw new TypeError(`Invalid triplestore ${label} JSON.`, { cause }) + } + if (typeof value !== 'object' || value === null || Array.isArray(value)) throw new TypeError(`Triplestore ${label} must be an object.`) + return value as Record +} + +/** Returns whether the supplied value satisfies the integer contract. */ +function isInteger(value: unknown): value is number { + return typeof value === 'number' && Number.isSafeInteger(value) +} diff --git a/packages/triplestore/format_test.ts b/packages/triplestore/format_test.ts new file mode 100644 index 0000000..6ac77ba --- /dev/null +++ b/packages/triplestore/format_test.ts @@ -0,0 +1,36 @@ +import { describe, it } from 'node:test' +import { expect } from '@std/expect' +import { FORMAT, generationName, parseCommit, parseFormat, writeRecord } from './format.ts' + +describe('@okikio/triplestore format records', () => { + it('round-trips the root format marker', () => { + expect(parseFormat(writeRecord(FORMAT))).toEqual(FORMAT) + }) + + it('validates homogeneous delta and snapshot commit invariants', () => { + const checksum = `sha256:${'0'.repeat(64)}` + const delta = parseCommit(JSON.stringify({ + version: 2, + generation: 2, + parent: 1, + mode: 'delta', + operation: 'delete', + segment: 'segments/0000000000000002.delta.nq', + checksum, + quadCount: 5, + })) + expect(delta.operation).toBe('delete') + + expect(() => parseCommit(JSON.stringify({ + ...delta, + mode: 'snapshot', + operation: 'add', + segment: 'segments/0000000000000002.nq', + }))).toThrow('must not declare') + }) + + it('uses lexicographically sortable generation names and rejects invalid generations', () => { + expect(generationName(42)).toBe('0000000000000042') + expect(() => generationName(0)).toThrow() + }) +}) diff --git a/packages/triplestore/mod.ts b/packages/triplestore/mod.ts new file mode 100644 index 0000000..86a81db --- /dev/null +++ b/packages/triplestore/mod.ts @@ -0,0 +1,5 @@ +/** Persistent, indexed RDF dataset storage. @module */ + +export * from './format.ts' +export * from './storage.ts' +export * from './store.ts' diff --git a/packages/triplestore/package.json b/packages/triplestore/package.json new file mode 100644 index 0000000..00eee4d --- /dev/null +++ b/packages/triplestore/package.json @@ -0,0 +1,22 @@ +{ + "name": "@okikio/triplestore", + "version": "0.1.0", + "type": "module", + "sideEffects": false, + "exports": { + ".": "./mod.ts" + }, + "description": "Crash-recoverable persistent RDF dataset storage.", + "license": "MIT", + "repository": { + "type": "git", + "url": "git+https://github.com/okikio/sparql-client.git", + "directory": "packages/triplestore" + }, + "publishConfig": { + "access": "public" + }, + "dependencies": { + "@okikio/rdf": "0.1.0" + } +} diff --git a/packages/triplestore/storage.ts b/packages/triplestore/storage.ts new file mode 100644 index 0000000..97b0801 --- /dev/null +++ b/packages/triplestore/storage.ts @@ -0,0 +1,38 @@ +/** Minimal durable file contract required by the triplestore. @module */ + +/** File metadata required for bounded recovery reads. */ +export interface FileStatType { + readonly kind: 'file' + readonly size: number +} + +/** Directory entry required while discovering immutable commit records. */ +export interface DirectoryEntryType { + readonly name: string + readonly kind: 'file' | 'directory' +} + +/** Options shared by storage operations that can be cancelled. */ +export interface SignalOptionsType { + readonly signal?: AbortSignal +} + +/** + * Structural filesystem contract consumed by `@okikio/triplestore`. + * + * `@okikio/opfs` `FileSystemType` satisfies this contract. Keeping the public + * dependency structural lets the store borrow compatible filesystems without + * importing or initializing a concrete runtime adapter. + */ +export interface FileSystemType { + exists(path: string, options?: SignalOptionsType): Promise + ensureDir(path: string, options?: SignalOptionsType): Promise + readDir(path: string, options?: SignalOptionsType): AsyncIterable + readText(path: string, options?: SignalOptionsType): Promise + stat(path: string, options?: SignalOptionsType): Promise + writeFile( + path: string, + data: string | Uint8Array, + options?: SignalOptionsType & { readonly parents?: boolean }, + ): Promise +} diff --git a/packages/triplestore/store.ts b/packages/triplestore/store.ts new file mode 100644 index 0000000..93ca5ad --- /dev/null +++ b/packages/triplestore/store.ts @@ -0,0 +1,532 @@ +/** Crash-recoverable persistent RDF dataset store. @module */ + +import { + Dataset, + iterate, + key, + type Graph, + type ObjectTerm, + type Predicate, + type Quad, + type Subject, +} from '@okikio/rdf' +import { parse as parseNQuads, write as writeNQuads } from '@okikio/rdf/nquads' +import { FORMAT, generationName, parseCommit, parseFormat, writeRecord, type CommitType } from './format.ts' +import type { FileSystemType } from './storage.ts' + +/** Default path used when the caller does not provide an override. */ +const DEFAULT_PATH = '/rdf' +/** Default max segment bytes used when the caller does not provide an override. */ +const DEFAULT_MAX_SEGMENT_BYTES = 256 * 1024 * 1024 +/** Default batch size used when the caller does not provide an override. */ +const DEFAULT_BATCH_SIZE = 10_000 + +/** Persistent store options. */ +export interface OpenOptionsType { + /** Virtual filesystem directory owned by this store. */ + readonly path?: string + /** Maximum immutable segment bytes accepted during recovery. */ + readonly maxSegmentBytes?: number + /** Cooperative cancellation for open and recovery. */ + readonly signal?: AbortSignal +} + +/** Mutation options shared by durable writes. */ +export interface WriteOptionsType { + readonly signal?: AbortSignal +} + +/** Streaming import options. */ +export interface ImportOptionsType extends WriteOptionsType { + /** Maximum quads committed in one immutable delta segment. */ + readonly batchSize?: number +} + +/** One recovery problem that did not prevent opening an earlier valid generation. */ +export interface RecoveryProblemType { + /** Whether the artifact was an incomplete publication or a corrupt committed generation. */ + readonly kind: 'incomplete-commit' | 'corrupt-generation' + readonly path: string + readonly message: string +} + +/** + * Persistent RDF Dataset backed by immutable commit records and segments. + * + * The store borrows the supplied filesystem. Closing the store never closes or + * disposes that filesystem. One store path supports one writer at a time; this + * baseline does not claim cross-process writer coordination. + */ +export class Store implements AsyncDisposable { + readonly #fs: FileSystemType + readonly #path: string + readonly #maxSegmentBytes: number + #dataset: Dataset + #generation: number + #closing = false + #closed = false + #tail: Promise = Promise.resolve() + #recovery: RecoveryProblemType[] + + /** Creates one live store over an already-recovered Dataset; filesystem ownership remains with the caller. */ + private constructor( + fs: FileSystemType, + path: string, + maxSegmentBytes: number, + dataset: Dataset, + generation: number, + recovery: RecoveryProblemType[], + ) { + this.#fs = fs + this.#path = path + this.#maxSegmentBytes = maxSegmentBytes + this.#dataset = dataset + this.#generation = generation + this.#recovery = recovery + } + + /** Number of unique quads in the latest committed generation. */ + get size(): number { + this.#assertOpen() + return this.#dataset.size + } + + /** Latest committed generation, or zero for a new empty store. */ + get generation(): number { + this.#assertOpen() + return this.#generation + } + + /** Non-fatal invalid newer artifacts ignored during recovery. */ + get recovery(): readonly RecoveryProblemType[] { + return this.#recovery + } + + /** Returns whether the latest committed generation contains one equal quad. */ + has(quad: Quad): boolean { + this.#assertOpen() + return this.#dataset.has(quad) + } + + /** Returns an inexpensive cardinality estimate from the current in-memory indexes. */ + estimate(options: { + readonly subject?: Subject | null + readonly predicate?: Predicate | null + readonly object?: ObjectTerm | null + readonly graph?: Graph | null + } = {}): number | undefined { + this.#assertOpen() + return this.#dataset.estimate(options) + } + + /** Lazily reads quads matching an RDF/JS-compatible quad pattern. */ + async *match( + subject: Subject | null = null, + predicate: Predicate | null = null, + object: ObjectTerm | null = null, + graph: Graph | null = null, + ): AsyncGenerator { + this.#assertOpen() + for (const quad of this.#dataset.matchIter({ subject, predicate, object, graph })) yield quad + } + + /** Durably adds one quad. Equal existing quads do not create a generation. */ + async add(quad: Quad, options: WriteOptionsType = {}): Promise { + await this.#mutate(async () => { + if (this.#dataset.has(quad)) return + await this.#commit([{ kind: 'add', quad }], options.signal) + }) + return this + } + + /** Durably removes one quad. Missing quads do not create a generation. */ + async delete(quad: Quad, options: WriteOptionsType = {}): Promise { + await this.#mutate(async () => { + if (!this.#dataset.has(quad)) return + await this.#commit([{ kind: 'delete', quad }], options.signal) + }) + return this + } + + /** Durably adds one materialized batch as a single generation. */ + async addAll(quads: Iterable, options: WriteOptionsType = {}): Promise { + await this.#mutate(async () => { + const operations: OperationType[] = [] + const pending = new Set() + for (const quad of quads) { + abort(options.signal) + const id = key(quad) + if (this.#dataset.has(quad) || pending.has(id)) continue + pending.add(id) + operations.push({ kind: 'add', quad }) + } + if (operations.length > 0) await this.#commit(operations, options.signal) + }) + return this + } + + /** + * Imports a sync/async source in bounded durable batches. + * + * Earlier completed batches remain committed if a later batch fails or the + * caller aborts. This is an explicit streaming-import contract, not an atomic + * transaction over the whole source. + */ + async import(source: Iterable | AsyncIterable, options: ImportOptionsType = {}): Promise { + const batchSize = positive(options.batchSize ?? DEFAULT_BATCH_SIZE, 'batchSize') + let batch: Quad[] = [] + for await (const quad of iterate(source)) { + abort(options.signal) + batch.push(quad) + if (batch.length < batchSize) continue + await this.addAll(batch, options) + batch = [] + } + if (batch.length > 0) await this.addAll(batch, options) + return this + } + + /** Durably removes every quad matching one pattern in a single generation. */ + async deleteMatches( + subject: Subject | null = null, + predicate: Predicate | null = null, + object: ObjectTerm | null = null, + graph: Graph | null = null, + options: WriteOptionsType = {}, + ): Promise { + await this.#mutate(async () => { + const operations = [...this.#dataset.matchIter({ subject, predicate, object, graph })] + .map((quad): OperationType => ({ kind: 'delete', quad })) + if (operations.length > 0) await this.#commit(operations, options.signal) + }) + return this + } + + /** Durably clears the store by publishing an empty snapshot generation. */ + async clear(options: WriteOptionsType = {}): Promise { + await this.#mutate(async () => { + if (this.#dataset.size === 0) return + await this.#snapshot(new Dataset(), options.signal) + }) + return this + } + + /** + * Publishes a complete immutable N-Quads snapshot of the current generation. + * + * The baseline keeps older immutable artifacts. Physical garbage collection + * is intentionally separate because safe deletion needs reader/retention rules. + */ + async compact(options: WriteOptionsType = {}): Promise { + await this.#mutate(async () => { + await this.#snapshot(new Dataset(this.#dataset), options.signal) + }) + return this + } + + /** Returns a detached in-memory snapshot for synchronous RDF Dataset operations. */ + snapshot(): Dataset { + this.#assertOpen() + return new Dataset(this.#dataset) + } + + /** Ends this store handle without disposing the borrowed filesystem. */ + async close(): Promise { + if (this.#closed || this.#closing) { + await this.#tail + return + } + this.#closing = true + await this.#tail + this.#closed = true + } + + /** Allows `await using` to close this store without disposing the borrowed filesystem. */ + async [Symbol.asyncDispose](): Promise { + await this.close() + } + + /** Commit as one isolated step of the Store state machine. */ + async #commit(operations: readonly OperationType[], signal?: AbortSignal): Promise { + abort(signal) + const generation = this.#generation + 1 + const stem = generationName(generation) + const kind = operationKind(operations) + const segment = `segments/${stem}.delta.nq` + const segmentText = writeNQuads(operations.map((operation) => operation.quad)) + const next = new Dataset(this.#dataset) + for (const operation of operations) apply(next, operation) + const commit: CommitType = { + version: 2, + generation, + parent: this.#generation, + mode: 'delta', + operation: kind, + segment, + checksum: await checksum(segmentText), + quadCount: next.size, + } + + await this.#publish(commit, segmentText, signal) + this.#dataset = next + this.#generation = generation + } + + /** Snapshot as one isolated step of the Store state machine. */ + async #snapshot(dataset: Dataset, signal?: AbortSignal): Promise { + abort(signal) + const generation = this.#generation + 1 + const stem = generationName(generation) + const segment = `segments/${stem}.nq` + const segmentText = writeNQuads(dataset) + const commit: CommitType = { + version: 2, + generation, + parent: this.#generation, + mode: 'snapshot', + segment, + checksum: await checksum(segmentText), + quadCount: dataset.size, + } + + await this.#publish(commit, segmentText, signal) + this.#dataset = dataset + this.#generation = generation + } + + /** Publish as one isolated step of the Store state machine. */ + async #publish(commit: CommitType, segmentText: string, signal?: AbortSignal): Promise { + abort(signal) + if (byteLength(segmentText) > this.#maxSegmentBytes) { + throw new RangeError(`Segment exceeds maxSegmentBytes (${this.#maxSegmentBytes}).`) + } + 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.`) + } + + // Publication order is the durability invariant. A segment may be orphaned, + // but a commit never becomes authoritative before its segment is complete. + await this.#fs.writeFile(segmentPath, segmentText, writeOptions(signal)) + abort(signal) + await this.#fs.writeFile(commitPath, writeRecord(commit), writeOptions(signal)) + } + + /** Mutate as one isolated step of the Store state machine. */ + async #mutate(operation: () => Promise): Promise { + this.#assertWritable() + const run = this.#tail.then(operation, operation) + this.#tail = run.then(() => undefined, () => undefined) + return await run + } + + /** Assert open as one isolated step of the Store state machine. */ + #assertOpen(): void { + if (this.#closed) throw new StoreError('closed', 'Triplestore is closed.') + } + + /** Assert writable as one isolated step of the Store state machine. */ + #assertWritable(): void { + if (this.#closing || this.#closed) throw new StoreError('closed', 'Triplestore is closing or closed.') + } + + /** Opens or creates one store and recovers its newest valid generation. */ + static async open(fs: FileSystemType, options: OpenOptionsType = {}): Promise { + const path = normalizePath(options.path ?? DEFAULT_PATH) + const maxSegmentBytes = positive(options.maxSegmentBytes ?? DEFAULT_MAX_SEGMENT_BYTES, 'maxSegmentBytes') + abort(options.signal) + await fs.ensureDir(path, signalOptions(options.signal)) + await fs.ensureDir(join(path, 'segments'), signalOptions(options.signal)) + await fs.ensureDir(join(path, 'commits'), signalOptions(options.signal)) + await ensureFormat(fs, path, options.signal) + + const candidates: number[] = [] + const recovery: RecoveryProblemType[] = [] + for await (const entry of fs.readDir(join(path, 'commits'), signalOptions(options.signal))) { + if (entry.kind !== 'file') continue + const match = /^([0-9]{16})\.json$/.exec(entry.name) + if (match?.[1]) candidates.push(Number.parseInt(match[1], 10)) + } + candidates.sort((a, b) => b - a) + + let committedFailure: unknown + for (const generation of candidates) { + const commitPath = join(path, `commits/${generationName(generation)}.json`) + try { + // Invalid JSON or an incomplete record can be left by an interrupted + // commit-file write. It never became a valid published generation. + parseCommit(await fs.readText(commitPath, signalOptions(options.signal))) + } catch (error) { + recovery.push({ kind: 'incomplete-commit', path: commitPath, message: message(error) }) + continue + } + + try { + const dataset = await recover(fs, path, generation, maxSegmentBytes, options.signal) + return new Store(fs, path, maxSegmentBytes, dataset, generation, recovery) + } catch (error) { + committedFailure ??= error + recovery.push({ kind: 'corrupt-generation', path: commitPath, message: message(error) }) + } + } + + if (committedFailure !== undefined) { + throw new StoreError('corrupt', 'No valid committed triplestore generation could be recovered.', { + cause: committedFailure, + }) + } + return new Store(fs, path, maxSegmentBytes, new Dataset(), 0, recovery) + } +} + +/** Store failure categories stable enough for callers to inspect. */ +export type StoreErrorKind = 'closed' | 'format' | 'corrupt' | 'writer-conflict' + +/** Normalized persistent-store failure. */ +export class StoreError extends Error { + readonly kind: StoreErrorKind + + /** Creates a normalized storage failure with its stable category and underlying cause. */ + constructor(kind: StoreErrorKind, message: string, options?: ErrorOptions) { + super(message, options) + this.name = 'StoreError' + this.kind = kind + } +} + +/** One pending store mutation; generations are homogeneous so publication later groups these by operation kind. */ +type OperationType = { readonly kind: 'add' | 'delete'; readonly quad: Quad } + +/** Creates or opens one persistent RDF store. */ +export async function open(fs: FileSystemType, options: OpenOptionsType = {}): Promise { + return await Store.open(fs, options) +} + +/** Creates or validates the root format marker before any generation is recovered or published. */ +async function ensureFormat(fs: FileSystemType, root: string, signal?: AbortSignal): Promise { + const path = join(root, 'format.json') + if (!await fs.exists(path, signalOptions(signal))) { + await fs.writeFile(path, writeRecord(FORMAT), writeOptions(signal)) + return + } + try { + parseFormat(await fs.readText(path, signalOptions(signal))) + } catch (cause) { + throw new StoreError('format', 'Persistent store format is incompatible or corrupt.', { cause }) + } +} + +/** Finds the newest fully valid committed generation, falling back only to an older valid commit. */ +async function recover( + fs: FileSystemType, + root: string, + head: number, + maxSegmentBytes: number, + signal?: AbortSignal, +): Promise { + const commits = new Map() + let generation = head + while (generation > 0) { + abort(signal) + const path = join(root, `commits/${generationName(generation)}.json`) + const commit = parseCommit(await fs.readText(path, signalOptions(signal))) + if (commit.generation !== generation) throw new StoreError('corrupt', `Commit path generation ${generation} disagrees with record ${commit.generation}.`) + const stem = generationName(generation) + const expectedSegment = commit.mode === 'snapshot' ? `segments/${stem}.nq` : `segments/${stem}.delta.nq` + if (commit.segment !== expectedSegment) throw new StoreError('corrupt', `Commit ${generation} references unexpected segment '${commit.segment}'.`) + commits.set(generation, commit) + if (commit.mode === 'snapshot') break + generation = commit.parent + } + + const dataset = new Dataset() + for (const commit of [...commits.values()].reverse()) { + abort(signal) + const path = join(root, commit.segment) + const stat = await fs.stat(path, signalOptions(signal)) + if (stat.kind !== 'file') throw new StoreError('corrupt', `Segment '${commit.segment}' is not a file.`) + if (stat.size > maxSegmentBytes) throw new StoreError('corrupt', `Segment '${commit.segment}' exceeds maxSegmentBytes.`) + const text = await fs.readText(path, signalOptions(signal)) + if (await checksum(text) !== commit.checksum) throw new StoreError('corrupt', `Segment checksum mismatch for '${commit.segment}'.`) + if (commit.mode === 'snapshot') dataset.clear() + for await (const quad of parseNQuads(text, signalOptions(signal))) { + if (commit.mode === 'snapshot' || commit.operation === 'add') dataset.add(quad) + else dataset.delete(quad) + } + } + + const current = commits.get(head) + if (!current) throw new StoreError('corrupt', `Missing head commit ${head}.`) + if (dataset.size !== current.quadCount) { + throw new StoreError('corrupt', `Recovered ${dataset.size} quads but commit ${head} records ${current.quadCount}.`) + } + return dataset +} + +/** Returns the one mutation kind encoded by a homogeneous delta generation. */ +function operationKind(operations: readonly OperationType[]): OperationType['kind'] { + const first = operations[0] + if (!first) throw new TypeError('Delta generation requires at least one operation.') + for (const operation of operations) { + if (operation.kind !== first.kind) throw new TypeError('Delta generation cannot mix add and delete operations.') + } + return first.kind +} + +/** Applies one homogeneous add/delete delta to the recovered in-memory Dataset. */ +function apply(dataset: Dataset, operation: OperationType): void { + if (operation.kind === 'add') dataset.add(operation.quad) + else dataset.delete(operation.quad) +} + +/** Computes the SHA-256 checksum recorded in immutable commit metadata. */ +async function checksum(text: string): Promise { + const bytes = new TextEncoder().encode(text) + const digest = await crypto.subtle.digest('SHA-256', bytes) + return `sha256:${[...new Uint8Array(digest)].map((byte) => byte.toString(16).padStart(2, '0')).join('')}` +} + +/** Returns the encoded UTF-8 byte length used for persistent record limits. */ +function byteLength(text: string): number { + return new TextEncoder().encode(text).byteLength +} + +/** Normalizes the caller store path without allowing an empty root record path. */ +function normalizePath(value: string): string { + if (!value.startsWith('/')) throw new TypeError('Store path must be absolute.') + const parts = value.split('/').filter(Boolean) + if (parts.some((part) => part === '.' || part === '..')) throw new TypeError('Store path cannot contain dot segments.') + return `/${parts.join('/')}` +} + +/** Joins normalized store-relative path components without introducing duplicate separators. */ +function join(root: string, child: string): string { + return root === '/' ? `/${child}` : `${root}/${child}` +} + +/** Validates a positive safe-integer storage option before it influences persistence limits. */ +function positive(value: number, name: string): number { + if (!Number.isSafeInteger(value) || value < 1) throw new RangeError(`${name} must be a positive safe integer.`) + return value +} + +/** Throws the caller supplied abort reason when cancellation has been requested. */ +function abort(signal?: AbortSignal): void { + if (signal?.aborted) throw signal.reason ?? new DOMException('Aborted', 'AbortError') +} + +/** Adds the operation AbortSignal only when one was supplied. */ +function signalOptions(signal?: AbortSignal): { readonly signal?: AbortSignal } { + return signal === undefined ? {} : { signal } +} + +/** Write options deterministically to the caller-owned output. */ +function writeOptions(signal?: AbortSignal): { readonly signal?: AbortSignal; readonly parents: true } { + return signal === undefined ? { parents: true } : { parents: true, signal } +} + +/** Converts an unknown failure into bounded diagnostic text for recovery records. */ +function message(error: unknown): string { + return error instanceof Error ? error.message : String(error) +} diff --git a/packages/triplestore/store_bench.ts b/packages/triplestore/store_bench.ts new file mode 100644 index 0000000..e05fb42 --- /dev/null +++ b/packages/triplestore/store_bench.ts @@ -0,0 +1,61 @@ +/** Decision benchmark for persistent-store lookup and recovery mechanisms. @module */ + +import { bench, do_not_optimize, group, run } from 'mitata' +import { Dataset, literal, namedNode, quad, type Quad, type Subject } from '@okikio/rdf' +import { MemoryFileSystem } from './_memory_test.ts' +import { open } from './mod.ts' + +const SIZE = 50_000 +const SUBJECTS = 5_000 +const predicate = namedNode('https://example.com/p') +const quads = Array.from({ length: SIZE }, (_, index) => quad( + namedNode(`https://example.com/s/${index % SUBJECTS}`), + predicate, + literal(`value-${index}`), +)) +const target = namedNode('https://example.com/s/1729') +const memory = new Dataset(quads) +const fs = new MemoryFileSystem() +const store = await open(fs, { path: '/db' }) +await store.addAll(quads) +const expected = scan(quads, target) +if (await count(store.match(target)) !== expected) throw new Error('Triplestore lookup oracle failed.') + +/** Full-scan baseline over the same semantic quads. */ +function scan(values: readonly Quad[], subject: Subject): number { + let matches = 0 + for (const value of values) if (value.subject.equals(subject)) matches++ + return matches +} + +/** Counts sync or async RDF values without materialization. */ +async function count(values: Iterable | AsyncIterable): Promise { + let matches = 0 + for await (const _value of values) matches++ + return matches +} + +group('triplestore exact subject lookup: 50k committed quads', () => { + bench('semantic array scan baseline', () => { + do_not_optimize(scan(quads, target)) + }) + + bench('in-memory Dataset exact index', () => { + do_not_optimize([...memory.matchIter({ subject: target })].length) + }) + + bench('persistent Store exact index', async () => { + do_not_optimize(await count(store.match(target))) + }) +}) + +group('triplestore cold recovery from immutable snapshot', () => { + bench('open + checksum + parse + index 50k quads', async () => { + const reopened = await open(fs, { path: '/db' }) + do_not_optimize(reopened.size) + await reopened.close() + }).gc('inner') +}) + +await run() +await store.close() diff --git a/packages/triplestore/store_test.ts b/packages/triplestore/store_test.ts new file mode 100644 index 0000000..1d2e86a --- /dev/null +++ b/packages/triplestore/store_test.ts @@ -0,0 +1,118 @@ +import { describe, it } from 'node:test' +import { expect } from '@std/expect' +import { literal, namedNode, quad } from '@okikio/rdf' +import { open, StoreError } from './mod.ts' +import { MemoryFileSystem } from './_memory_test.ts' + +const predicate = namedNode('https://example.com/p') +const first = quad(namedNode('https://example.com/a'), predicate, literal('one')) +const second = quad(namedNode('https://example.com/b'), predicate, literal('two')) + +describe('@okikio/triplestore', () => { + it('recovers the newest complete committed generation after an interrupted commit publication', async () => { + const fs = new MemoryFileSystem() + const store = await open(fs, { path: '/db' }) + await store.add(first) + + fs.writeFault = (path, text, target) => { + if (!path.endsWith('/commits/0000000000000002.json')) return + target.files.set(path, text.slice(0, Math.max(1, Math.floor(text.length / 2)))) + throw new Error('simulated publication crash') + } + await expect(store.add(second)).rejects.toThrow('simulated publication crash') + fs.writeFault = undefined + + const reopened = await open(fs, { path: '/db' }) + expect(reopened.generation).toBe(1) + expect(reopened.has(first)).toBe(true) + expect(reopened.has(second)).toBe(false) + expect(reopened.recovery[0]?.kind).toBe('incomplete-commit') + }) + + it('never converts a corrupt authoritative committed generation into an empty database', async () => { + const fs = new MemoryFileSystem() + const store = await open(fs, { path: '/db' }) + await store.add(first) + fs.files.set('/db/segments/0000000000000001.delta.nq', 'corrupt\n') + + try { + await open(fs, { path: '/db' }) + throw new Error('Expected corrupt store open to fail.') + } catch (error) { + expect(error instanceof StoreError).toBe(true) + if (!(error instanceof StoreError)) return + expect(error.kind).toBe('corrupt') + } + }) + + it('borrows rather than owns the injected filesystem', async () => { + const fs = new MemoryFileSystem() + const store = await open(fs, { path: '/db' }) + await store.close() + expect(await fs.exists('/db/format.json')).toBe(true) + }) + + it('reopens homogeneous add and delete deltas with exact RDF terms', async () => { + const fs = new MemoryFileSystem() + const store = await open(fs, { path: '/db' }) + await store.addAll([first, second]) + await store.delete(first) + await store.close() + + const reopened = await open(fs, { path: '/db' }) + expect(reopened.generation).toBe(2) + expect(reopened.size).toBe(1) + expect(reopened.has(first)).toBe(false) + expect(reopened.has(second)).toBe(true) + }) + + it('falls back from a corrupt newer generation but never invents an empty authoritative store', async () => { + const fs = new MemoryFileSystem() + const store = await open(fs, { path: '/db' }) + await store.add(first) + await store.add(second) + fs.files.set('/db/segments/0000000000000002.delta.nq', 'corrupt\n') + + const fallback = await open(fs, { path: '/db' }) + expect(fallback.generation).toBe(1) + expect(fallback.has(first)).toBe(true) + expect(fallback.has(second)).toBe(false) + }) + + it('recovers from a compacted snapshot without older delta segments', async () => { + const fs = new MemoryFileSystem() + const store = await open(fs, { path: '/db' }) + await store.addAll([first, second]) + await store.compact() + fs.files.delete('/db/segments/0000000000000001.delta.nq') + + const reopened = await open(fs, { path: '/db' }) + expect(reopened.size).toBe(2) + expect(reopened.has(first)).toBe(true) + expect(reopened.has(second)).toBe(true) + }) + + + it('prevents mutations after close without disposing the borrowed filesystem', async () => { + const fs = new MemoryFileSystem() + const store = await open(fs, { path: '/db' }) + await store.close() + expect(await fs.exists('/db/format.json')).toBe(true) + await expect(store.add(first)).rejects.toThrow() + }) + + it('treats an incomplete first commit as interrupted publication rather than authoritative data', async () => { + const fs = new MemoryFileSystem() + await fs.ensureDir('/db/segments') + await fs.ensureDir('/db/commits') + await fs.writeFile('/db/format.json', '{"version":2,"store":"@okikio/triplestore","segment":2}\n') + await fs.writeFile('/db/segments/0000000000000001.delta.nq', 'orphan\n') + await fs.writeFile('/db/commits/0000000000000001.json', '{"version":2') + + const reopened = await open(fs, { path: '/db' }) + expect(reopened.generation).toBe(0) + expect(reopened.size).toBe(0) + expect(reopened.recovery[0]?.kind).toBe('incomplete-commit') + }) + +})