feat: column-store JS↔native interchange — raw-blob unify (2.4.0 #4)
Brainy's JS ColumnStore.flushBuffer now writes segment payloads + the
per-field DELETED bitmap via the binary-blob primitive (saveBinaryBlob) at
the cortex-shared key convention `_column_index/<field>/L<level>-NNNNNN`
(suffix-free; adapter appends its own). Byte-for-byte identical to what
cortex's NativeColumnStore writes — JS and native engines now read each
other's segments with no envelope re-encoding.
Closes the gap that the cortex 2.3.1 read-side fallback opened: cortex was
reading both formats but the JS engine was still WRITING the legacy
{_binary, base64} envelope, so a fresh JS write always required the cortex
fallback to kick in (and migrate lazily on first read). With this commit
JS writes natively in the unified format, and cortex's fallback only
fires on indexes persisted by older brainy releases.
Backward compat:
- ColumnStore.loadSegmentCursor tries the raw blob first, then falls back
to the legacy `{_binary, base64}` envelope at the `.cidx` object-path.
Indexes written by pre-2.4.0 brainy keep loading correctly.
- ColumnStore.init's DELETED-bitmap load has the same dual-format read.
- Adapters without the binary-blob primitive (custom adapters that didn't
follow the 7.25.0 surface) fall through to the legacy envelope writer
too, so writes still succeed there.
Tests (1433 total, +5 vs prior tip):
- tests/unit/indexes/columnStore/column-store-interchange.test.ts —
pins down the contract: (1) flush writes the raw blob, NO legacy
envelope; (2) DELETED bitmap likewise; (3) a legacy-format on-disk
segment loads correctly via the fallback; (4) legacy DELETED bitmap
ditto; (5) round-trip in the new format.
- All 101 existing ColumnStore tests pass — the new write path is
exercised by the existing lifecycle tests (MemoryStorage has the blob
primitive, so the new branch fires).
This commit is contained in:
parent
d4cb26c604
commit
71bc30b829
2 changed files with 293 additions and 29 deletions
|
|
@ -125,12 +125,29 @@ export class ColumnStore implements ColumnStoreProvider {
|
|||
this.manifests.set(fieldName, manifest)
|
||||
this.fieldTypes.set(fieldName, manifest.valueType)
|
||||
|
||||
// Load global deleted bitmap if it exists
|
||||
// Load global deleted bitmap if it exists. Raw blob preferred
|
||||
// (2.4.0 #4 cortex-shared format); legacy envelope fallback for
|
||||
// indexes written before format unification.
|
||||
try {
|
||||
const deletedPath = `${this.basePath}/${fieldName}/DELETED.bin`
|
||||
const stored = await (storage as any).readObjectFromPath(deletedPath)
|
||||
if (stored && stored._binary && stored.data) {
|
||||
const buf = Buffer.from(stored.data, 'base64')
|
||||
let buf: Buffer | null = null
|
||||
|
||||
const adapter = storage as unknown as {
|
||||
loadBinaryBlob?: (key: string) => Promise<Buffer | null>
|
||||
readObjectFromPath: (path: string) => Promise<any>
|
||||
}
|
||||
if (typeof adapter.loadBinaryBlob === 'function') {
|
||||
const key = `${this.basePath}/${fieldName}/DELETED`
|
||||
const blob = await adapter.loadBinaryBlob(key)
|
||||
if (blob && blob.length > 0) buf = blob
|
||||
}
|
||||
if (!buf) {
|
||||
const deletedPath = `${this.basePath}/${fieldName}/DELETED.bin`
|
||||
const stored = await adapter.readObjectFromPath(deletedPath)
|
||||
if (stored && stored._binary && stored.data) {
|
||||
buf = Buffer.from(stored.data, 'base64')
|
||||
}
|
||||
}
|
||||
if (buf) {
|
||||
this.deletedEntities.set(fieldName, RoaringBitmap32.deserialize(buf, true))
|
||||
}
|
||||
} catch {
|
||||
|
|
@ -408,6 +425,23 @@ export class ColumnStore implements ColumnStoreProvider {
|
|||
|
||||
/**
|
||||
* Flush a single field's tail buffer to a new L0 segment.
|
||||
*
|
||||
* Storage path (2.4.0 #4 / cortex-interchange contract):
|
||||
* When the storage adapter exposes the binary-blob primitive
|
||||
* (`saveBinaryBlob` — every brainy adapter ≥7.25.0 does), the segment
|
||||
* bytes are written as a raw blob at the shared cortex key
|
||||
* `_column_index/<field>/L<level>-NNNNNN` (suffix-free; the adapter
|
||||
* appends its own). This matches byte-for-byte what cortex's
|
||||
* `NativeColumnStore` writes, so JS- and native-written indexes
|
||||
* interchange without re-encoding.
|
||||
*
|
||||
* When the adapter doesn't (older custom adapters that didn't follow
|
||||
* the 7.25.0 primitive surface), we fall back to the legacy
|
||||
* `writeObjectToPath` envelope — `{ _binary, data: <base64> }` at
|
||||
* `_column_index/<field>/L<level>-NNNNNN.cidx`. Cortex's 2.3.1
|
||||
* read-side fallback handles the legacy envelope on its end, so the
|
||||
* format mismatch is transparent in both directions during the
|
||||
* transition.
|
||||
*/
|
||||
private async flushBuffer(field: string, buffer: ColumnTailBuffer): Promise<void> {
|
||||
const { values, entityIds, pendingRemovals } = buffer.drain()
|
||||
|
|
@ -416,10 +450,15 @@ export class ColumnStore implements ColumnStoreProvider {
|
|||
const manifest = this.manifests.get(field)
|
||||
if (!manifest) return
|
||||
|
||||
const storage = this.storage as unknown as {
|
||||
saveBinaryBlob?: (key: string, data: Buffer) => Promise<void>
|
||||
writeObjectToPath: (path: string, value: unknown) => Promise<void>
|
||||
}
|
||||
const canUseBlob = typeof storage.saveBinaryBlob === 'function'
|
||||
|
||||
// Write segment if we have values
|
||||
if (values.length > 0) {
|
||||
const segId = manifest.nextSegmentId()
|
||||
const segPath = manifest.segmentPath(0, segId)
|
||||
|
||||
const segBuffer = writeSegmentToBuffer(
|
||||
{
|
||||
|
|
@ -435,13 +474,22 @@ export class ColumnStore implements ColumnStoreProvider {
|
|||
entityIds
|
||||
)
|
||||
|
||||
// Store segment as binary data
|
||||
await (this.storage as any).writeObjectToPath(segPath, {
|
||||
_binary: true,
|
||||
data: segBuffer.toString('base64')
|
||||
})
|
||||
if (canUseBlob) {
|
||||
// Raw blob, cortex-shared key convention.
|
||||
const key = `${this.basePath}/${field}/L0-${String(segId).padStart(6, '0')}`
|
||||
await storage.saveBinaryBlob!(key, segBuffer)
|
||||
} else {
|
||||
// Legacy envelope path for adapters without the binary-blob primitive.
|
||||
const segPath = manifest.segmentPath(0, segId)
|
||||
await storage.writeObjectToPath(segPath, {
|
||||
_binary: true,
|
||||
data: segBuffer.toString('base64')
|
||||
})
|
||||
}
|
||||
|
||||
// Register in manifest
|
||||
// Register in manifest. The `.cidx` file name in the manifest is the
|
||||
// legacy-format file the cortex 2.3.1 read-side fallback looks for; the
|
||||
// raw-blob key derives from segId at read time. Both paths converge.
|
||||
manifest.addSegment({
|
||||
id: segId,
|
||||
level: 0,
|
||||
|
|
@ -465,15 +513,21 @@ export class ColumnStore implements ColumnStoreProvider {
|
|||
for (const id of pendingRemovals) deleted.add(id)
|
||||
}
|
||||
|
||||
// Persist the global deleted bitmap alongside the manifest
|
||||
// Persist the global deleted bitmap alongside the manifest. Same
|
||||
// raw-blob preference as segments — shared cortex key `<base>/<field>/DELETED`.
|
||||
const deleted = this.deletedEntities.get(field)
|
||||
if (deleted && deleted.size > 0) {
|
||||
const deletedPath = `${this.basePath}/${field}/DELETED.bin`
|
||||
const serialized = deleted.serialize(true)
|
||||
await (this.storage as any).writeObjectToPath(deletedPath, {
|
||||
_binary: true,
|
||||
data: Buffer.from(serialized).toString('base64')
|
||||
})
|
||||
if (canUseBlob) {
|
||||
const deletedKey = `${this.basePath}/${field}/DELETED`
|
||||
await storage.saveBinaryBlob!(deletedKey, Buffer.from(serialized))
|
||||
} else {
|
||||
const deletedPath = `${this.basePath}/${field}/DELETED.bin`
|
||||
await storage.writeObjectToPath(deletedPath, {
|
||||
_binary: true,
|
||||
data: Buffer.from(serialized).toString('base64')
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// Save manifest
|
||||
|
|
@ -508,24 +562,43 @@ export class ColumnStore implements ColumnStoreProvider {
|
|||
|
||||
/**
|
||||
* Load a segment from storage and create a cursor.
|
||||
*
|
||||
* Reads the raw-blob layout first (the 2.4.0 #4 shared-with-cortex format),
|
||||
* then falls back to the legacy `{ _binary, base64 }` envelope at the
|
||||
* `.cidx` object-path so indexes written before the format unification keep
|
||||
* loading correctly. Mirror of cortex's 2.3.1 read-side fallback.
|
||||
*/
|
||||
private async loadSegmentCursor(field: string, seg: SegmentMeta): Promise<ColumnSegmentCursor | null> {
|
||||
const manifest = this.manifests.get(field)
|
||||
if (!manifest) return null
|
||||
|
||||
try {
|
||||
const segPath = manifest.segmentPath(seg.level, seg.id)
|
||||
const stored = await (this.storage as any).readObjectFromPath(segPath)
|
||||
if (!stored) return null
|
||||
let buf: Buffer | null = null
|
||||
|
||||
// Handle binary storage format
|
||||
let buf: Buffer
|
||||
if (stored._binary && stored.data) {
|
||||
buf = Buffer.from(stored.data, 'base64')
|
||||
} else if (Buffer.isBuffer(stored)) {
|
||||
buf = stored
|
||||
} else {
|
||||
return null
|
||||
const storage = this.storage as unknown as {
|
||||
loadBinaryBlob?: (key: string) => Promise<Buffer | null>
|
||||
readObjectFromPath: (path: string) => Promise<any>
|
||||
}
|
||||
|
||||
// Preferred: raw blob at the cortex-shared key.
|
||||
if (typeof storage.loadBinaryBlob === 'function') {
|
||||
const key = `${this.basePath}/${field}/L${seg.level}-${String(seg.id).padStart(6, '0')}`
|
||||
const blob = await storage.loadBinaryBlob(key)
|
||||
if (blob && blob.length > 0) buf = blob
|
||||
}
|
||||
|
||||
// Legacy fallback: `{ _binary, base64 }` envelope at the .cidx object-path.
|
||||
if (!buf) {
|
||||
const segPath = manifest.segmentPath(seg.level, seg.id)
|
||||
const stored = await storage.readObjectFromPath(segPath)
|
||||
if (!stored) return null
|
||||
if (stored._binary && stored.data) {
|
||||
buf = Buffer.from(stored.data, 'base64')
|
||||
} else if (Buffer.isBuffer(stored)) {
|
||||
buf = stored
|
||||
} else {
|
||||
return null
|
||||
}
|
||||
}
|
||||
|
||||
const parsed = readSegmentFromBuffer(buf)
|
||||
|
|
|
|||
191
tests/unit/indexes/columnStore/column-store-interchange.test.ts
Normal file
191
tests/unit/indexes/columnStore/column-store-interchange.test.ts
Normal file
|
|
@ -0,0 +1,191 @@
|
|||
/**
|
||||
* @module column-store-interchange.test
|
||||
* @description 2.4.0 #4 — column-store JS↔native interchange. Pins down that
|
||||
* brainy's JS `ColumnStore` writes segments + DELETED bitmaps via the binary-
|
||||
* blob primitive at the cortex-shared key convention, AND that indexes
|
||||
* persisted under the pre-2.4.0 envelope (`{ _binary, base64 }` at the
|
||||
* `.cidx` object-path) still load correctly via the legacy-read fallback.
|
||||
*
|
||||
* Together these properties make the JS and native engines byte-compatible
|
||||
* at the segment payload, so writing through one and reading through the
|
||||
* other is a no-op format-wise — closing the gap the cortex 2.3.1 read-side
|
||||
* fallback opened on the cortex side.
|
||||
*/
|
||||
|
||||
import { describe, it, expect, beforeEach, afterEach } from 'vitest'
|
||||
import { ColumnStore } from '../../../../src/indexes/columnStore/ColumnStore.js'
|
||||
import { MemoryStorage } from '../../../../src/storage/adapters/memoryStorage.js'
|
||||
import { EntityIdMapper } from '../../../../src/utils/entityIdMapper.js'
|
||||
import { writeSegmentToBuffer } from '../../../../src/indexes/columnStore/ColumnSegmentFormat.js'
|
||||
import { ColumnManifest } from '../../../../src/indexes/columnStore/ColumnManifest.js'
|
||||
import { RoaringBitmap32 } from '../../../../src/utils/roaring/index.js'
|
||||
|
||||
describe('ColumnStore — JS↔native interchange (2.4.0 #4)', () => {
|
||||
let storage: MemoryStorage
|
||||
let idMapper: EntityIdMapper
|
||||
let store: ColumnStore
|
||||
|
||||
beforeEach(async () => {
|
||||
storage = new MemoryStorage()
|
||||
await storage.init()
|
||||
idMapper = new EntityIdMapper({ storage, storageKey: 'test:idMapper' })
|
||||
await idMapper.init()
|
||||
store = new ColumnStore({ flushThreshold: 10 })
|
||||
await store.init(storage, idMapper)
|
||||
})
|
||||
|
||||
afterEach(async () => {
|
||||
await store.close()
|
||||
})
|
||||
|
||||
// -----------------------------------------------------------------------
|
||||
// WRITE side: raw blob at the cortex-shared key, no legacy envelope
|
||||
// -----------------------------------------------------------------------
|
||||
|
||||
it('flushBuffer writes segment bytes as a raw blob at the cortex-shared key', async () => {
|
||||
for (let i = 0; i < 5; i++) {
|
||||
store.addEntity(idMapper.getOrAssign(`u-${i}`), { score: i * 10 })
|
||||
}
|
||||
await store.flush()
|
||||
|
||||
// The raw blob exists at the suffix-free cortex key.
|
||||
const key = '_column_index/score/L0-000001'
|
||||
const raw = await (storage as any).loadBinaryBlob(key)
|
||||
expect(Buffer.isBuffer(raw)).toBe(true)
|
||||
expect(raw!.length).toBeGreaterThan(0)
|
||||
|
||||
// The legacy envelope at the `.cidx` object-path is NOT written (no
|
||||
// 33% base64 bloat, no dual-write).
|
||||
const legacyPath = '_column_index/score/L0-000001.cidx'
|
||||
const legacyEnvelope = await (storage as any).readObjectFromPath(legacyPath)
|
||||
expect(legacyEnvelope).toBeFalsy()
|
||||
})
|
||||
|
||||
it('flushBuffer writes the DELETED bitmap as a raw blob too (not an envelope)', async () => {
|
||||
for (let i = 0; i < 3; i++) {
|
||||
store.addEntity(idMapper.getOrAssign(`u-${i}`), { score: i })
|
||||
}
|
||||
store.removeEntity(idMapper.getInt('u-1')!)
|
||||
await store.flush()
|
||||
|
||||
const rawDel = await (storage as any).loadBinaryBlob('_column_index/score/DELETED')
|
||||
expect(Buffer.isBuffer(rawDel)).toBe(true)
|
||||
expect(rawDel!.length).toBeGreaterThan(0)
|
||||
|
||||
const legacyEnvelope = await (storage as any).readObjectFromPath(
|
||||
'_column_index/score/DELETED.bin'
|
||||
)
|
||||
expect(legacyEnvelope).toBeFalsy()
|
||||
})
|
||||
|
||||
// -----------------------------------------------------------------------
|
||||
// READ side: pre-2.4.0 legacy envelope still loads via the fallback
|
||||
// -----------------------------------------------------------------------
|
||||
|
||||
it('loads a segment persisted as the pre-2.4.0 envelope (legacy fallback)', async () => {
|
||||
// Seed a legacy-format index DIRECTLY on storage — the on-disk shape
|
||||
// brainy 7.20.0..7.25.0 wrote before this commit.
|
||||
const legacyStore = new ColumnStore()
|
||||
await legacyStore.init(storage, idMapper)
|
||||
|
||||
// Build a segment buffer using the shared cidx format.
|
||||
const values: number[] = [10, 20, 30]
|
||||
const entityIds: number[] = [
|
||||
idMapper.getOrAssign('legacy-a'),
|
||||
idMapper.getOrAssign('legacy-b'),
|
||||
idMapper.getOrAssign('legacy-c')
|
||||
]
|
||||
const segBuffer = writeSegmentToBuffer(
|
||||
{
|
||||
fieldName: 'score',
|
||||
fieldNameLength: 5,
|
||||
valueType: 0,
|
||||
level: 0,
|
||||
codec: 0,
|
||||
flags: 0,
|
||||
count: values.length
|
||||
},
|
||||
values,
|
||||
entityIds
|
||||
)
|
||||
|
||||
// Persist via the LEGACY envelope path (writeObjectToPath, base64 wrap).
|
||||
await (storage as any).writeObjectToPath('_column_index/score/L0-000001.cidx', {
|
||||
_binary: true,
|
||||
data: segBuffer.toString('base64')
|
||||
})
|
||||
const manifest = new ColumnManifest('score', '_column_index')
|
||||
manifest.valueType = 0
|
||||
manifest.multiValue = false
|
||||
manifest.addSegment({
|
||||
id: 1,
|
||||
level: 0,
|
||||
count: values.length,
|
||||
minValue: values[0],
|
||||
maxValue: values[values.length - 1],
|
||||
file: 'L0-000001.cidx'
|
||||
})
|
||||
await manifest.save(storage)
|
||||
await legacyStore.close()
|
||||
|
||||
// NEW ColumnStore re-discovers the manifest at init, then the legacy
|
||||
// segment loads via the readObjectFromPath fallback (no raw blob exists
|
||||
// at the cortex key, so it falls through). sortTopK returns the entities
|
||||
// in the order the segment recorded — proves the bytes round-tripped.
|
||||
const readStore = new ColumnStore()
|
||||
await readStore.init(storage, idMapper)
|
||||
const sortedInts = await readStore.sortTopK('score', 'asc', 10)
|
||||
const sortedUuids = sortedInts.map(i => idMapper.getUuid(i))
|
||||
expect(sortedUuids).toEqual(['legacy-a', 'legacy-b', 'legacy-c'])
|
||||
await readStore.close()
|
||||
})
|
||||
|
||||
it('loads a DELETED bitmap persisted as the pre-2.4.0 envelope (legacy fallback)', async () => {
|
||||
// Manually seed a legacy DELETED.bin envelope, then re-init and confirm
|
||||
// the entity it marks gone is excluded from filter/sort results.
|
||||
const bm = new RoaringBitmap32()
|
||||
bm.add(idMapper.getOrAssign('dead'))
|
||||
const serialized = Buffer.from(bm.serialize(true))
|
||||
await (storage as any).writeObjectToPath('_column_index/score/DELETED.bin', {
|
||||
_binary: true,
|
||||
data: serialized.toString('base64')
|
||||
})
|
||||
|
||||
// Add a couple of live entities so a manifest exists (init only loads
|
||||
// DELETED for fields with a manifest on disk).
|
||||
store.addEntity(idMapper.getOrAssign('alive-1'), { score: 100 })
|
||||
store.addEntity(idMapper.getOrAssign('alive-2'), { score: 200 })
|
||||
await store.flush()
|
||||
await store.close()
|
||||
|
||||
const readStore = new ColumnStore()
|
||||
await readStore.init(storage, idMapper)
|
||||
const sortedInts = await readStore.sortTopK('score', 'asc', 10)
|
||||
const sortedUuids = sortedInts.map(i => idMapper.getUuid(i))
|
||||
|
||||
// The dead entity must NOT appear; the alive ones must.
|
||||
expect(sortedUuids).not.toContain('dead')
|
||||
expect(sortedUuids).toContain('alive-1')
|
||||
expect(sortedUuids).toContain('alive-2')
|
||||
await readStore.close()
|
||||
})
|
||||
|
||||
// -----------------------------------------------------------------------
|
||||
// ROUND-TRIP: write-then-read in the new format
|
||||
// -----------------------------------------------------------------------
|
||||
|
||||
it('write (new format) → re-init → read returns the same data', async () => {
|
||||
for (let i = 0; i < 5; i++) {
|
||||
store.addEntity(idMapper.getOrAssign(`u-${i}`), { score: i * 10 })
|
||||
}
|
||||
await store.flush()
|
||||
await store.close()
|
||||
|
||||
const readStore = new ColumnStore()
|
||||
await readStore.init(storage, idMapper)
|
||||
const sortedInts = await readStore.sortTopK('score', 'asc', 10)
|
||||
const sortedUuids = sortedInts.map(i => idMapper.getUuid(i))
|
||||
expect(sortedUuids).toEqual(['u-0', 'u-1', 'u-2', 'u-3', 'u-4'])
|
||||
await readStore.close()
|
||||
})
|
||||
})
|
||||
Loading…
Add table
Add a link
Reference in a new issue