feat: stable EntityIdMapper — rebuild() no longer renumbers (2.4.0 #1, the foundation) #1
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 {
|
||||
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 (storage as any).readObjectFromPath(deletedPath)
|
||||
const stored = await adapter.readObjectFromPath(deletedPath)
|
||||
if (stored && stored._binary && stored.data) {
|
||||
const buf = Buffer.from(stored.data, 'base64')
|
||||
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, {
|
||||
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,16 +513,22 @@ 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, {
|
||||
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
|
||||
await manifest.save(this.storage)
|
||||
|
|
@ -508,18 +562,36 @@ 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
|
||||
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)) {
|
||||
|
|
@ -527,6 +599,7 @@ export class ColumnStore implements ColumnStoreProvider {
|
|||
} else {
|
||||
return null
|
||||
}
|
||||
}
|
||||
|
||||
const parsed = readSegmentFromBuffer(buf)
|
||||
return new ColumnSegmentCursor(
|
||||
|
|
|
|||
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