Graph time-travel needs an edge's existence recorded per generation so db.asOf(g) hops resolve historically correct endpoints. The metadata layer already threads brainy's commit generation per write; the graph write path did not, leaving a versioned verb-endpoint store unable to answer "which edges existed at generation g". - GraphIndexProvider.addVerb/removeVerb gain a `generation: bigint` parameter (the same watermark the storage layer stamps onto the record). A provider with a per-generation edge chain stamps the edge at that generation; the JS baseline has no such chain and accepts-and-ignores it — graph time-travel is a native-provider capability, and the open-core path serves edges as-of-now (the one documented graph time-travel limitation). - The two graph transaction operations resolve the generation via a thunk at EXECUTE time: the generation store assigns the batch generation only once the commit begins executing, after the operations are planned. The same generation is reused for an operation's rollback half. - All graph-write call sites pass the in-flight generation. Adds a spy-provider test proving the threading, execute-time resolution, and forward/rollback generation reuse. The JS index ignores the value, so behaviour is unchanged: unit 1402/1402, db-mvcc 25/25, bigint-contract relate/unrelate 10/10.
339 lines
11 KiB
TypeScript
339 lines
11 KiB
TypeScript
/**
|
|
* Index Operations with Rollback Support
|
|
*
|
|
* Provides transactional operations for all indexes:
|
|
* - JsHnswVectorIndex (unified vector index)
|
|
* - MetadataIndexManager (roaring bitmap filtering)
|
|
* - GraphAdjacencyIndex (LSM-tree graph storage)
|
|
*
|
|
* Each operation can be executed and rolled back atomically.
|
|
*/
|
|
|
|
import type { JsHnswVectorIndex } from '../../hnsw/hnswIndex.js'
|
|
import type { MetadataIndexManager } from '../../utils/metadataIndex.js'
|
|
import type { GraphIndexProvider } from '../../plugin.js'
|
|
import type { GraphVerb } from '../../coreTypes.js'
|
|
import type { Operation, RollbackAction } from '../types.js'
|
|
|
|
/**
|
|
* Add to HNSW index with rollback support
|
|
*
|
|
* Rollback strategy:
|
|
* - Remove item from index
|
|
|
|
*/
|
|
export class AddToHNSWOperation implements Operation {
|
|
readonly name = 'AddToHNSW'
|
|
|
|
constructor(
|
|
private readonly index: JsHnswVectorIndex,
|
|
private readonly id: string,
|
|
private readonly vector: number[]
|
|
) {}
|
|
|
|
async execute(): Promise<RollbackAction> {
|
|
// Check if item already exists (for rollback decision)
|
|
const existed = await this.itemExists(this.id)
|
|
|
|
// Add to index
|
|
await this.index.addItem({ id: this.id, vector: this.vector })
|
|
|
|
// Return rollback action
|
|
return async () => {
|
|
if (!existed) {
|
|
// Remove newly added item
|
|
await this.index.removeItem(this.id)
|
|
}
|
|
// If item existed before, we don't rollback (update is OK)
|
|
// This prevents index corruption from removing pre-existing items
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Check if item exists in index.
|
|
*
|
|
* `getItem` is an optional, feature-detected provider capability — see the
|
|
* VectorIndexProvider docs; it is intentionally absent from the required
|
|
* contract (Brainy's JS HNSW index omits it). When the capability is
|
|
* missing the answer must be `false`, not `true`: treating unknowable
|
|
* pre-existence as "existed" made every rollback skip removeItem, leaving
|
|
* phantom entries in the index after a failed transaction. The safe default
|
|
* is to remove what this operation added — update flows pair this op with a
|
|
* RemoveFromHNSWOperation whose own rollback restores the prior vector, so
|
|
* reverse-order rollback reconstructs the original state either way.
|
|
*/
|
|
private async itemExists(id: string): Promise<boolean> {
|
|
const index = this.index as JsHnswVectorIndex & {
|
|
getItem?: (id: string) => Promise<unknown>
|
|
}
|
|
if (typeof index.getItem !== 'function') return false
|
|
try {
|
|
const item = await index.getItem(id)
|
|
return item !== undefined && item !== null
|
|
} catch {
|
|
return false
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Remove from HNSW index with rollback support
|
|
*
|
|
* Rollback strategy:
|
|
* - Re-add item to index with original vector
|
|
*
|
|
* Note: Requires storing the vector for rollback
|
|
*/
|
|
export class RemoveFromHNSWOperation implements Operation {
|
|
readonly name = 'RemoveFromHNSW'
|
|
|
|
constructor(
|
|
private readonly index: JsHnswVectorIndex,
|
|
private readonly id: string,
|
|
private readonly vector: number[] // Required for rollback
|
|
) {}
|
|
|
|
async execute(): Promise<RollbackAction> {
|
|
// Remove from index
|
|
await this.index.removeItem(this.id)
|
|
|
|
// Return rollback action
|
|
return async () => {
|
|
// Re-add item with original vector
|
|
await this.index.addItem({ id: this.id, vector: this.vector })
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Add to metadata index with rollback support
|
|
*
|
|
* Rollback strategy:
|
|
* - Remove item from index
|
|
*/
|
|
export class AddToMetadataIndexOperation implements Operation {
|
|
readonly name = 'AddToMetadataIndex'
|
|
|
|
constructor(
|
|
private readonly index: MetadataIndexManager,
|
|
private readonly id: string,
|
|
private readonly entity: any // Entity or metadata structure
|
|
) {}
|
|
|
|
async execute(): Promise<RollbackAction> {
|
|
// Add to metadata index (skipFlush=true for transaction atomicity)
|
|
await this.index.addToIndex(this.id, this.entity, true)
|
|
|
|
// Return rollback action
|
|
return async () => {
|
|
// Remove from metadata index
|
|
await this.index.removeFromIndex(this.id, this.entity)
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Remove from metadata index with rollback support
|
|
*
|
|
* Rollback strategy:
|
|
* - Re-add item to index with original metadata
|
|
*/
|
|
export class RemoveFromMetadataIndexOperation implements Operation {
|
|
readonly name = 'RemoveFromMetadataIndex'
|
|
|
|
constructor(
|
|
private readonly index: MetadataIndexManager,
|
|
private readonly id: string,
|
|
private readonly entity: any // Required for rollback
|
|
) {}
|
|
|
|
async execute(): Promise<RollbackAction> {
|
|
// Remove from metadata index
|
|
await this.index.removeFromIndex(this.id, this.entity)
|
|
|
|
// Return rollback action
|
|
return async () => {
|
|
// Re-add with original metadata (skipFlush=true)
|
|
await this.index.addToIndex(this.id, this.entity, true)
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Add verb to graph index with rollback support
|
|
*
|
|
* Rollback strategy:
|
|
* - Remove verb from graph index
|
|
*
|
|
* 8.0 u64 contract: the coordinator resolves both endpoint ints via the
|
|
* shared idMapper (`getOrAssign`) and passes them alongside the verb; the
|
|
* provider returns the interned verb int, which is surfaced through the
|
|
* optional `onVerbInt` callback so the coordinator can feed its warm cache.
|
|
*
|
|
* Generation: `generationFn` is resolved at execute time (not construction) so
|
|
* the edge is stamped at the transaction's in-flight commit generation — which
|
|
* the generation store only assigns once the batch begins executing. The same
|
|
* generation is reused for the rollback removal, so an add and its undo
|
|
* reference one watermark in a provider's per-generation edge chain.
|
|
*/
|
|
export class AddToGraphIndexOperation implements Operation {
|
|
readonly name = 'AddToGraphIndex'
|
|
|
|
/**
|
|
* @param index - The graph-index provider (JS baseline or native).
|
|
* @param verb - The verb to index (`sourceInt`/`targetInt` mirrored on it).
|
|
* @param sourceInt - The source entity's interned int.
|
|
* @param targetInt - The target entity's interned int.
|
|
* @param generationFn - Resolves the commit generation to stamp this edge at,
|
|
* evaluated when the operation executes (see class note).
|
|
* @param onVerbInt - Optional hook invoked with the interned verb int
|
|
* returned by the provider (feeds the coordinator's verb-int warm cache).
|
|
*/
|
|
constructor(
|
|
private readonly index: GraphIndexProvider,
|
|
private readonly verb: GraphVerb,
|
|
private readonly sourceInt: bigint,
|
|
private readonly targetInt: bigint,
|
|
private readonly generationFn: () => bigint,
|
|
private readonly onVerbInt?: (verbInt: bigint) => void
|
|
) {}
|
|
|
|
async execute(): Promise<RollbackAction> {
|
|
// Stamp this edge at the in-flight commit generation; reuse it for the
|
|
// rollback so add + undo reference the same watermark.
|
|
const generation = this.generationFn()
|
|
const verbInt = await this.index.addVerb(this.verb, this.sourceInt, this.targetInt, generation)
|
|
this.onVerbInt?.(verbInt)
|
|
|
|
// Return rollback action
|
|
return async () => {
|
|
// Remove verb from graph index
|
|
await this.index.removeVerb(this.verb.id, generation)
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Remove verb from graph index with rollback support
|
|
*
|
|
* Rollback strategy:
|
|
* - Re-add verb to graph index
|
|
*
|
|
* 8.0 u64 contract: rollback re-adds through `addVerb(verb, sourceInt,
|
|
* targetInt, generation)`, so the coordinator resolves the endpoint ints up
|
|
* front (while the entity → int mappings are guaranteed to still exist). The
|
|
* removal generation is resolved at execute time and reused for the rollback
|
|
* re-add, so the round trip references one watermark.
|
|
*/
|
|
export class RemoveFromGraphIndexOperation implements Operation {
|
|
readonly name = 'RemoveFromGraphIndex'
|
|
|
|
/**
|
|
* @param index - The graph-index provider (JS baseline or native).
|
|
* @param verb - The verb being removed (required for rollback re-add).
|
|
* @param sourceInt - The source entity's interned int (rollback re-add).
|
|
* @param targetInt - The target entity's interned int (rollback re-add).
|
|
* @param generationFn - Resolves the commit generation for this removal,
|
|
* evaluated when the operation executes.
|
|
*/
|
|
constructor(
|
|
private readonly index: GraphIndexProvider,
|
|
private readonly verb: GraphVerb, // Required for rollback
|
|
private readonly sourceInt: bigint,
|
|
private readonly targetInt: bigint,
|
|
private readonly generationFn: () => bigint
|
|
) {}
|
|
|
|
async execute(): Promise<RollbackAction> {
|
|
// Resolve the removal generation once; reuse it for the rollback re-add.
|
|
const generation = this.generationFn()
|
|
await this.index.removeVerb(this.verb.id, generation)
|
|
|
|
// Return rollback action
|
|
return async () => {
|
|
// Re-add verb with original data
|
|
await this.index.addVerb(this.verb, this.sourceInt, this.targetInt, generation)
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Batch operation: Add multiple items to HNSW index
|
|
*
|
|
* Useful for bulk imports with transaction support.
|
|
* Rolls back all items if any fail.
|
|
*/
|
|
export class BatchAddToHNSWOperation implements Operation {
|
|
readonly name = 'BatchAddToHNSW'
|
|
|
|
private operations: AddToHNSWOperation[]
|
|
|
|
constructor(
|
|
index: JsHnswVectorIndex,
|
|
items: Array<{ id: string; vector: number[] }>
|
|
) {
|
|
this.operations = items.map(
|
|
item => new AddToHNSWOperation(index, item.id, item.vector)
|
|
)
|
|
}
|
|
|
|
async execute(): Promise<RollbackAction> {
|
|
const rollbackActions: RollbackAction[] = []
|
|
|
|
// Execute all operations
|
|
for (const op of this.operations) {
|
|
const rollback = await op.execute()
|
|
if (rollback) {
|
|
rollbackActions.push(rollback)
|
|
}
|
|
}
|
|
|
|
// Return combined rollback action
|
|
return async () => {
|
|
// Execute all rollbacks in reverse order
|
|
for (let i = rollbackActions.length - 1; i >= 0; i--) {
|
|
await rollbackActions[i]()
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Batch operation: Add multiple entities to metadata index
|
|
*
|
|
* Useful for bulk imports with transaction support.
|
|
*/
|
|
export class BatchAddToMetadataIndexOperation implements Operation {
|
|
readonly name = 'BatchAddToMetadataIndex'
|
|
|
|
private operations: AddToMetadataIndexOperation[]
|
|
|
|
constructor(
|
|
index: MetadataIndexManager,
|
|
items: Array<{ id: string; entity: any }>
|
|
) {
|
|
this.operations = items.map(
|
|
item => new AddToMetadataIndexOperation(index, item.id, item.entity)
|
|
)
|
|
}
|
|
|
|
async execute(): Promise<RollbackAction> {
|
|
const rollbackActions: RollbackAction[] = []
|
|
|
|
// Execute all operations
|
|
for (const op of this.operations) {
|
|
const rollback = await op.execute()
|
|
if (rollback) {
|
|
rollbackActions.push(rollback)
|
|
}
|
|
}
|
|
|
|
// Return combined rollback action
|
|
return async () => {
|
|
// Execute all rollbacks in reverse order
|
|
for (let i = rollbackActions.length - 1; i >= 0; i--) {
|
|
await rollbackActions[i]()
|
|
}
|
|
}
|
|
}
|
|
}
|