feat(8.0)!: delete fork/branch/commit/history/versions — superseded by the Db API
The COW version-control surface (fork, branches, checkout, commit,
getHistory/streamHistory, asOfCommit, brain.versions) is gone, along with
its subsystems: src/versioning/, the COW object store (CommitLog,
CommitObject, RefManager, TreeObject), HistoricalStorageAdapter, and the
TypeAwareHNSWIndex path it kept alive. The Db API (now/transact/asOf/
with/persist/restore) is the one versioning model in 8.0.
Survivors and replacements:
- BlobStorage survives (the VFS stores file content through it), relocated
to src/storage/blobStorage.ts with binaryDataCodec.ts; its adapter
interface is now BlobStoreAdapter, slimmed to the consumed surface
(write/read/has/delete/getMetadata + MIME-aware compression policy).
- brain.migrate() backup branches are replaced by persist-before-migrate:
MigrateOptions.backupTo persists a hard-link snapshot of the current
generation before any transform runs; MigrationResult.backupPath reports
it, and brain.restore(path) brings it back wholesale.
- CLI: cow.ts (fork/branch/checkout/history/migrate) is replaced by
snapshot.ts — snapshot <path>, restore <path>, history (tx-log),
generation.
- New public read API: brain.transactionLog({limit}) exposes the reified
tx-log (generation/timestamp/meta, newest first) that backs the CLI
history command; TxLogEntry is exported.
Tests: superseded suites deleted; fork/commit blocks excised from shared
suites; BlobStorage tests relocated + reworked against the slimmed store;
migration tests now prove the backupTo snapshot/restore round trip; new
transactionLog coverage in db-mvcc.
This commit is contained in:
parent
431cd64406
commit
8f93add705
64 changed files with 1444 additions and 16313 deletions
498
src/storage/blobStorage.ts
Normal file
498
src/storage/blobStorage.ts
Normal file
|
|
@ -0,0 +1,498 @@
|
|||
/**
|
||||
* @module storage/blobStorage
|
||||
* @description Content-addressed blob store. Backs VFS file content with
|
||||
* SHA-256 addressing (automatic deduplication), reference counting, zstd
|
||||
* compression where it pays (MIME-aware: already-compressed media is stored
|
||||
* raw), and an LRU read cache.
|
||||
*
|
||||
* The store persists through a narrow key-value bridge
|
||||
* ({@link BlobStoreAdapter}) provided by `BaseStorage`, which roots all keys
|
||||
* under the `_cas/` storage area. Key naming is an explicit type contract:
|
||||
* `blob:<hash>` keys hold binary bytes, `blob-meta:<hash>` keys hold JSON
|
||||
* metadata — the key format decides how bytes are encoded, never content
|
||||
* sniffing.
|
||||
*/
|
||||
|
||||
import { createHash } from 'crypto'
|
||||
import { unwrapBinaryData } from './binaryDataCodec.js'
|
||||
|
||||
/**
|
||||
* @description Key-value bridge the blob store persists through. Implemented
|
||||
* by `BaseStorage.initializeBlobStorage()` over the adapter's raw object
|
||||
* primitives.
|
||||
*/
|
||||
export interface BlobStoreAdapter {
|
||||
/** Read the bytes stored under `key`, or `undefined` when absent. */
|
||||
get(key: string): Promise<Buffer | undefined>
|
||||
/** Persist `data` under `key` (overwrites). */
|
||||
put(key: string, data: Buffer): Promise<void>
|
||||
/** Delete the value under `key`. Missing keys are ignored. */
|
||||
delete(key: string): Promise<void>
|
||||
/** List all keys starting with `prefix`. */
|
||||
list(prefix: string): Promise<string[]>
|
||||
}
|
||||
|
||||
/**
|
||||
* @description Metadata persisted alongside each blob (under
|
||||
* `blob-meta:<hash>`).
|
||||
*/
|
||||
export interface BlobMetadata {
|
||||
/** SHA-256 content hash (the blob's identity). */
|
||||
hash: string
|
||||
/** Original (uncompressed) size in bytes. */
|
||||
size: number
|
||||
/** Stored size in bytes (after compression, if any). */
|
||||
compressedSize: number
|
||||
/** Compression applied to the stored bytes. */
|
||||
compression: 'none' | 'zstd'
|
||||
/** Creation timestamp (epoch ms). */
|
||||
createdAt: number
|
||||
/** Number of logical references to this blob (deduplicated writes). */
|
||||
refCount: number
|
||||
}
|
||||
|
||||
/**
|
||||
* @description Options for {@link BlobStorage.write}.
|
||||
*/
|
||||
export interface BlobWriteOptions {
|
||||
/**
|
||||
* Compression strategy. `'auto'` (default) compresses payloads above 1 KB
|
||||
* with zstd unless the MIME type says the bytes are already compressed.
|
||||
* Explicit `'none'`/`'zstd'` is honoured as asserted by the caller.
|
||||
*/
|
||||
compression?: 'none' | 'zstd' | 'auto'
|
||||
/**
|
||||
* Content type of the payload (e.g. `image/jpeg`, `video/mp4`).
|
||||
*
|
||||
* When set on an `auto`-compression write, the store skips zstd for MIME
|
||||
* types that are already heavily compressed (JPEG, PNG, WebP, MP4, WebM,
|
||||
* MP3, ZIP, PDF, etc.). zstd over these formats wastes CPU and rarely
|
||||
* shaves more than a single-digit percent — usually it actually grows the
|
||||
* payload because the entropy is already maximised by the format itself.
|
||||
*
|
||||
* Has no effect when `compression` is `'none'` or `'zstd'` explicitly —
|
||||
* the caller is asserting the choice and the store honours it.
|
||||
*/
|
||||
mimeType?: string
|
||||
}
|
||||
|
||||
/**
|
||||
* MIME types whose payload is already heavily compressed. zstd over these is
|
||||
* almost always a CPU-only loss — the bytes are already near entropy-maximal,
|
||||
* so the output is the same size or slightly larger plus the cost of running
|
||||
* the compressor. Used by `BlobStorage.selectCompression()` in `auto` mode.
|
||||
*
|
||||
* Conservative denylist (well-known formats only). Anything not in this set
|
||||
* goes through the size heuristic. False negatives (compressing something we
|
||||
* should have skipped) waste CPU; false positives (skipping something we
|
||||
* could have compressed) waste a few percent of bytes. The denylist favours
|
||||
* CPU-cycle safety because the formats listed here are the ones where
|
||||
* gzip/zstd is reliably a net loss.
|
||||
*/
|
||||
const ALREADY_COMPRESSED_MIME_TYPES = new Set<string>([
|
||||
// Images
|
||||
'image/jpeg', 'image/jpg', 'image/png', 'image/gif', 'image/webp',
|
||||
'image/avif', 'image/heic', 'image/heif', 'image/jp2',
|
||||
// Video
|
||||
'video/mp4', 'video/webm', 'video/x-matroska', 'video/quicktime',
|
||||
'video/x-msvideo', 'video/mpeg', 'video/3gpp', 'video/x-ms-wmv',
|
||||
// Audio
|
||||
'audio/mpeg', 'audio/mp4', 'audio/aac', 'audio/ogg', 'audio/webm',
|
||||
'audio/opus', 'audio/flac', 'audio/x-ms-wma',
|
||||
// Archives
|
||||
'application/zip', 'application/gzip', 'application/x-gzip',
|
||||
'application/x-bzip2', 'application/x-7z-compressed',
|
||||
'application/x-rar-compressed', 'application/x-xz', 'application/x-zstd',
|
||||
'application/x-compress', 'application/vnd.rar',
|
||||
// Documents with internal compression
|
||||
'application/pdf', 'application/epub+zip',
|
||||
// Office formats (zip-based)
|
||||
'application/vnd.openxmlformats-officedocument.wordprocessingml.document',
|
||||
'application/vnd.openxmlformats-officedocument.spreadsheetml.sheet',
|
||||
'application/vnd.openxmlformats-officedocument.presentationml.presentation',
|
||||
'application/vnd.oasis.opendocument.text',
|
||||
'application/vnd.oasis.opendocument.spreadsheet',
|
||||
'application/vnd.oasis.opendocument.presentation'
|
||||
])
|
||||
|
||||
/**
|
||||
* @description True when the MIME type names a payload format known to be
|
||||
* already heavily compressed. Strips any `;charset=…` / `;boundary=…`
|
||||
* parameters and lowercases the bare type/subtype before lookup, so
|
||||
* `'IMAGE/JPEG; charset=binary'` and `'image/jpeg'` resolve identically.
|
||||
* @param mimeType - MIME type string, or `undefined`.
|
||||
* @returns Whether `auto` compression should skip zstd for this payload.
|
||||
*/
|
||||
export function isAlreadyCompressedMimeType(mimeType: string | undefined): boolean {
|
||||
if (!mimeType) return false
|
||||
const bare = mimeType.split(';', 1)[0].trim().toLowerCase()
|
||||
return ALREADY_COMPRESSED_MIME_TYPES.has(bare)
|
||||
}
|
||||
|
||||
/**
|
||||
* LRU cache entry.
|
||||
*/
|
||||
interface CacheEntry {
|
||||
data: Buffer
|
||||
metadata: BlobMetadata
|
||||
lastAccess: number
|
||||
size: number
|
||||
}
|
||||
|
||||
/**
|
||||
* @description Content-addressed, deduplicating, reference-counted blob
|
||||
* store with MIME-aware zstd compression and an LRU read cache. See the
|
||||
* module doc for the persistence contract.
|
||||
*
|
||||
* @example
|
||||
* const hash = await blobStorage.write(buffer, { mimeType: 'image/png' })
|
||||
* const bytes = await blobStorage.read(hash) // verified against the hash
|
||||
* await blobStorage.delete(hash) // decrements refCount first
|
||||
*/
|
||||
export class BlobStorage {
|
||||
private adapter: BlobStoreAdapter
|
||||
private cache: Map<string, CacheEntry>
|
||||
private cacheMaxSize: number
|
||||
private currentCacheSize: number
|
||||
|
||||
// Compression (lazily loaded)
|
||||
private zstdCompress?: (data: Buffer) => Promise<Buffer>
|
||||
private zstdDecompress?: (data: Buffer) => Promise<Buffer>
|
||||
private compressionReady = false
|
||||
|
||||
// Configuration
|
||||
private readonly CACHE_MAX_SIZE = 100 * 1024 * 1024 // 100MB default
|
||||
private readonly COMPRESSION_THRESHOLD = 1024 // 1KB - don't compress smaller
|
||||
|
||||
/**
|
||||
* @param adapter - Key-value bridge to persist through.
|
||||
* @param options - `cacheMaxSize` bounds the LRU read cache (bytes,
|
||||
* default 100 MB).
|
||||
*/
|
||||
constructor(adapter: BlobStoreAdapter, options?: { cacheMaxSize?: number }) {
|
||||
this.adapter = adapter
|
||||
this.cache = new Map()
|
||||
this.cacheMaxSize = options?.cacheMaxSize ?? this.CACHE_MAX_SIZE
|
||||
this.currentCacheSize = 0
|
||||
}
|
||||
|
||||
/**
|
||||
* Lazy-load the zstd compression module. Falls back to uncompressed
|
||||
* storage when the optional dependency is unavailable.
|
||||
*/
|
||||
private async ensureCompressionReady(): Promise<void> {
|
||||
if (this.compressionReady) return
|
||||
try {
|
||||
// Dynamic import to avoid loading if not needed
|
||||
// @ts-ignore - Optional dependency, gracefully handled if missing
|
||||
const zstd = await import('@mongodb-js/zstd')
|
||||
this.zstdCompress = async (data: Buffer) => {
|
||||
return Buffer.from(await zstd.compress(data, 3)) // Level 3 = fast
|
||||
}
|
||||
this.zstdDecompress = async (data: Buffer) => {
|
||||
return Buffer.from(await zstd.decompress(data))
|
||||
}
|
||||
} catch (error) {
|
||||
console.warn('zstd compression not available, falling back to uncompressed')
|
||||
this.zstdCompress = undefined
|
||||
this.zstdDecompress = undefined
|
||||
}
|
||||
this.compressionReady = true
|
||||
}
|
||||
|
||||
/**
|
||||
* @description Compute the SHA-256 content hash of `data`.
|
||||
* @param data - Bytes to hash.
|
||||
* @returns Hex-encoded SHA-256 hash.
|
||||
*/
|
||||
static hash(data: Buffer): string {
|
||||
return createHash('sha256').update(data).digest('hex')
|
||||
}
|
||||
|
||||
/**
|
||||
* @description Write a blob. Content-addressed: the SHA-256 hash of the
|
||||
* bytes is the storage key, so identical payloads deduplicate (the
|
||||
* existing blob's reference count is incremented instead of rewriting).
|
||||
*
|
||||
* @param data - Blob bytes.
|
||||
* @param options - Compression strategy and MIME hint (see
|
||||
* {@link BlobWriteOptions}).
|
||||
* @returns The blob's SHA-256 hash.
|
||||
*/
|
||||
async write(data: Buffer, options: BlobWriteOptions = {}): Promise<string> {
|
||||
const hash = BlobStorage.hash(data)
|
||||
|
||||
// Deduplication: identical content already stored — just add a reference.
|
||||
if (await this.has(hash)) {
|
||||
await this.incrementRefCount(hash)
|
||||
return hash
|
||||
}
|
||||
|
||||
await this.ensureCompressionReady()
|
||||
|
||||
// Determine compression strategy
|
||||
const compression = this.selectCompression(data, options)
|
||||
|
||||
// Compress if needed
|
||||
let finalData = data
|
||||
let compressedSize = data.length
|
||||
|
||||
if (compression === 'zstd' && this.zstdCompress) {
|
||||
finalData = await this.zstdCompress(data)
|
||||
compressedSize = finalData.length
|
||||
}
|
||||
|
||||
// Record the ACTUAL compression state, not the intended one — prevents
|
||||
// corruption if compression failed to initialize.
|
||||
const actualCompression = finalData === data ? 'none' : compression
|
||||
const metadata: BlobMetadata = {
|
||||
hash,
|
||||
size: data.length,
|
||||
compressedSize,
|
||||
compression: actualCompression,
|
||||
createdAt: Date.now(),
|
||||
refCount: 1
|
||||
}
|
||||
|
||||
await this.adapter.put(`blob:${hash}`, finalData)
|
||||
await this.adapter.put(`blob-meta:${hash}`, Buffer.from(JSON.stringify(metadata)))
|
||||
|
||||
// Write-through cache (caches the ORIGINAL bytes, not the compressed form)
|
||||
this.addToCache(hash, data, metadata)
|
||||
|
||||
return hash
|
||||
}
|
||||
|
||||
/**
|
||||
* @description Read a blob: LRU cache first, then storage with
|
||||
* decompression and integrity verification (the bytes are re-hashed and
|
||||
* compared against the requested hash).
|
||||
*
|
||||
* @param hash - The blob's SHA-256 hash.
|
||||
* @returns The original (decompressed) blob bytes.
|
||||
* @throws Error when the blob is missing or fails integrity verification.
|
||||
*/
|
||||
async read(hash: string): Promise<Buffer> {
|
||||
// Check cache first
|
||||
const cached = this.getFromCache(hash)
|
||||
if (cached) {
|
||||
return cached.data
|
||||
}
|
||||
|
||||
const metadataBuffer = await this.adapter.get(`blob-meta:${hash}`)
|
||||
if (!metadataBuffer) {
|
||||
throw new Error(`Blob metadata not found: ${hash}`)
|
||||
}
|
||||
// Unwrap before parsing (defense-in-depth): metadata should come back as
|
||||
// JSON bytes, but an adapter might return the wrapped binary format.
|
||||
const metadata: BlobMetadata = JSON.parse(unwrapBinaryData(metadataBuffer).toString())
|
||||
|
||||
const data = await this.adapter.get(`blob:${hash}`)
|
||||
if (!data) {
|
||||
throw new Error(`Blob not found: ${hash}`)
|
||||
}
|
||||
|
||||
// Decompress if needed
|
||||
let finalData = data
|
||||
if (metadata.compression === 'zstd') {
|
||||
if (!this.zstdDecompress) {
|
||||
await this.ensureCompressionReady()
|
||||
}
|
||||
if (!this.zstdDecompress) {
|
||||
throw new Error('zstd decompression not available')
|
||||
}
|
||||
finalData = await this.zstdDecompress(data)
|
||||
}
|
||||
|
||||
// Defense-in-depth unwrap: even though the bridge unwraps, verify it
|
||||
// happened and re-unwrap if needed. Hash verification must run on the
|
||||
// original content bytes.
|
||||
const unwrappedData = unwrapBinaryData(finalData)
|
||||
|
||||
// Integrity verification (always on — a content-addressed store that
|
||||
// returns bytes not matching the address is corruption, not a result)
|
||||
if (BlobStorage.hash(unwrappedData) !== hash) {
|
||||
throw new Error(`Blob integrity check failed: ${hash}`)
|
||||
}
|
||||
|
||||
this.addToCache(hash, unwrappedData, metadata)
|
||||
|
||||
return unwrappedData
|
||||
}
|
||||
|
||||
/**
|
||||
* @description Whether a blob with this hash exists (cache or storage).
|
||||
* @param hash - The blob's SHA-256 hash.
|
||||
* @returns True when the blob exists.
|
||||
*/
|
||||
async has(hash: string): Promise<boolean> {
|
||||
if (this.cache.has(hash)) {
|
||||
return true
|
||||
}
|
||||
const exists = await this.adapter.get(`blob:${hash}`)
|
||||
return exists !== undefined
|
||||
}
|
||||
|
||||
/**
|
||||
* @description Drop one reference to the blob. The stored bytes and
|
||||
* metadata are physically deleted only when the reference count reaches
|
||||
* zero — deduplicated content shared by other writers survives.
|
||||
*
|
||||
* @param hash - The blob's SHA-256 hash.
|
||||
*/
|
||||
async delete(hash: string): Promise<void> {
|
||||
const refCount = await this.decrementRefCount(hash)
|
||||
|
||||
// Only delete if no references remain
|
||||
if (refCount > 0) {
|
||||
return
|
||||
}
|
||||
|
||||
await this.adapter.delete(`blob:${hash}`)
|
||||
await this.adapter.delete(`blob-meta:${hash}`)
|
||||
this.removeFromCache(hash)
|
||||
}
|
||||
|
||||
/**
|
||||
* @description Read a blob's metadata without reading its bytes.
|
||||
* @param hash - The blob's SHA-256 hash.
|
||||
* @returns The metadata, or `undefined` when the blob does not exist.
|
||||
*/
|
||||
async getMetadata(hash: string): Promise<BlobMetadata | undefined> {
|
||||
const data = await this.adapter.get(`blob-meta:${hash}`)
|
||||
if (data) {
|
||||
return JSON.parse(unwrapBinaryData(data).toString())
|
||||
}
|
||||
return undefined
|
||||
}
|
||||
|
||||
// ========== PRIVATE METHODS ==========
|
||||
|
||||
/**
|
||||
* Select the compression strategy for a write (see
|
||||
* {@link BlobWriteOptions.compression}).
|
||||
*/
|
||||
private selectCompression(
|
||||
data: Buffer,
|
||||
options: BlobWriteOptions
|
||||
): 'none' | 'zstd' {
|
||||
if (options.compression === 'none') {
|
||||
return 'none'
|
||||
}
|
||||
|
||||
if (options.compression === 'zstd') {
|
||||
return this.zstdCompress ? 'zstd' : 'none'
|
||||
}
|
||||
|
||||
// Auto mode
|
||||
if (data.length < this.COMPRESSION_THRESHOLD) {
|
||||
return 'none' // Too small to benefit
|
||||
}
|
||||
|
||||
// Content-type policy: skip already-compressed media. zstd over
|
||||
// JPEG / MP4 / ZIP etc. is a CPU loss for no measurable byte savings,
|
||||
// and on hot save paths (image / video uploads) it's the difference
|
||||
// between fast and slow. Applies only to `auto`; explicit `'zstd'` is
|
||||
// honoured because the caller is asserting the choice.
|
||||
if (isAlreadyCompressedMimeType(options.mimeType)) {
|
||||
return 'none'
|
||||
}
|
||||
|
||||
return this.zstdCompress ? 'zstd' : 'none'
|
||||
}
|
||||
|
||||
/**
|
||||
* Increment the reference count for an existing blob.
|
||||
*/
|
||||
private async incrementRefCount(hash: string): Promise<number> {
|
||||
const metadata = await this.getMetadata(hash)
|
||||
if (!metadata) {
|
||||
throw new Error(`Cannot increment ref count, blob not found: ${hash}`)
|
||||
}
|
||||
|
||||
metadata.refCount++
|
||||
await this.adapter.put(`blob-meta:${hash}`, Buffer.from(JSON.stringify(metadata)))
|
||||
return metadata.refCount
|
||||
}
|
||||
|
||||
/**
|
||||
* Decrement the reference count for a blob (floored at zero).
|
||||
*/
|
||||
private async decrementRefCount(hash: string): Promise<number> {
|
||||
const metadata = await this.getMetadata(hash)
|
||||
if (!metadata) {
|
||||
return 0
|
||||
}
|
||||
|
||||
metadata.refCount = Math.max(0, metadata.refCount - 1)
|
||||
await this.adapter.put(`blob-meta:${hash}`, Buffer.from(JSON.stringify(metadata)))
|
||||
return metadata.refCount
|
||||
}
|
||||
|
||||
/**
|
||||
* Add a blob to the LRU cache (evicting least-recently-used entries to
|
||||
* stay under the size bound).
|
||||
*/
|
||||
private addToCache(hash: string, data: Buffer, metadata: BlobMetadata): void {
|
||||
if (data.length > this.cacheMaxSize) {
|
||||
return // Blob too large for cache
|
||||
}
|
||||
|
||||
while (
|
||||
this.currentCacheSize + data.length > this.cacheMaxSize &&
|
||||
this.cache.size > 0
|
||||
) {
|
||||
this.evictLRU()
|
||||
}
|
||||
|
||||
this.cache.set(hash, {
|
||||
data,
|
||||
metadata,
|
||||
lastAccess: Date.now(),
|
||||
size: data.length
|
||||
})
|
||||
|
||||
this.currentCacheSize += data.length
|
||||
}
|
||||
|
||||
/**
|
||||
* Get a blob from the cache, refreshing its LRU position.
|
||||
*/
|
||||
private getFromCache(hash: string): CacheEntry | undefined {
|
||||
const entry = this.cache.get(hash)
|
||||
if (entry) {
|
||||
entry.lastAccess = Date.now() // Update LRU
|
||||
}
|
||||
return entry
|
||||
}
|
||||
|
||||
/**
|
||||
* Remove a blob from the cache.
|
||||
*/
|
||||
private removeFromCache(hash: string): void {
|
||||
const entry = this.cache.get(hash)
|
||||
if (entry) {
|
||||
this.cache.delete(hash)
|
||||
this.currentCacheSize -= entry.size
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Evict the least-recently-used cache entry.
|
||||
*/
|
||||
private evictLRU(): void {
|
||||
let oldestHash: string | null = null
|
||||
let oldestTime = Infinity
|
||||
|
||||
for (const [hash, entry] of this.cache.entries()) {
|
||||
if (entry.lastAccess < oldestTime) {
|
||||
oldestTime = entry.lastAccess
|
||||
oldestHash = hash
|
||||
}
|
||||
}
|
||||
|
||||
if (oldestHash) {
|
||||
this.removeFromCache(oldestHash)
|
||||
}
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue