feat(persist): HNSW vector indexes survive restart
This commit is contained in:
@@ -17,6 +17,7 @@ const
|
|||||||
SchemaUserPrefix* = "_schema:users:"
|
SchemaUserPrefix* = "_schema:users:"
|
||||||
SchemaPolicyPrefix* = "_schema:policies:"
|
SchemaPolicyPrefix* = "_schema:policies:"
|
||||||
SchemaFtsIndexPrefix* = "_schema:ftsidx:"
|
SchemaFtsIndexPrefix* = "_schema:ftsidx:"
|
||||||
|
SchemaVecIndexPrefix* = "_schema:vecidx:"
|
||||||
## Legacy CREATE TABLE keys (pre-fix) used a migrations: counter suffix
|
## Legacy CREATE TABLE keys (pre-fix) used a migrations: counter suffix
|
||||||
SchemaLegacyCreatePrefix* = "_schema:migrations:"
|
SchemaLegacyCreatePrefix* = "_schema:migrations:"
|
||||||
|
|
||||||
|
|||||||
@@ -1336,6 +1336,10 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu
|
|||||||
docId = docId * 31 + uint64(ord(ch))
|
docId = docId * 31 + uint64(ord(ch))
|
||||||
vengine.insert(hnswIdx, docId, vec, meta)
|
vengine.insert(hnswIdx, docId, vec, meta)
|
||||||
ctx.vectorIndexes[colKey] = hnswIdx
|
ctx.vectorIndexes[colKey] = hnswIdx
|
||||||
|
# Persist reconstructed DDL so restoreEngines can rebuild the index
|
||||||
|
# from table data after a restart (replay re-writes the same key).
|
||||||
|
let vecDdl = "CREATE INDEX " & idxName & " ON " & stmt.ciTarget & " (" & stmt.ciColumns.join(", ") & ") USING HNSW"
|
||||||
|
ctx.db.put(SchemaVecIndexPrefix & colKey, cast[seq[byte]](vecDdl))
|
||||||
return okResult(msg="CREATE INDEX " & idxName & " on " & stmt.ciTarget & " USING HNSW")
|
return okResult(msg="CREATE INDEX " & idxName & " on " & stmt.ciTarget & " USING HNSW")
|
||||||
|
|
||||||
ctx.btrees[colKey] = newBTreeIndex[string, IndexEntry]()
|
ctx.btrees[colKey] = newBTreeIndex[string, IndexEntry]()
|
||||||
@@ -1569,12 +1573,13 @@ proc executeMigrationSql(ctx: ExecutionContext, sql: string): ExecResult =
|
|||||||
return okResult(msg="Empty migration body")
|
return okResult(msg="Empty migration body")
|
||||||
|
|
||||||
proc restoreEngines*(ctx: ExecutionContext) =
|
proc restoreEngines*(ctx: ExecutionContext) =
|
||||||
## Rebuild ephemeral engines (FTS indexes) from persisted schema keys after
|
## Rebuild ephemeral engines (FTS/HNSW indexes) from persisted schema keys
|
||||||
## restoreSchema. Invoked via context.restoreEnginesHook at the end of
|
## after restoreSchema. Invoked via context.restoreEnginesHook at the end of
|
||||||
## newExecutionContext. Replay re-persists the same key, so it is idempotent.
|
## newExecutionContext. Replay re-persists the same key, so it is idempotent.
|
||||||
var ddls: seq[string] = @[]
|
var ddls: seq[string] = @[]
|
||||||
for (key, value) in ctx.db.scanAll():
|
for (key, value) in ctx.db.scanAll():
|
||||||
if not key.startsWith(SchemaFtsIndexPrefix): continue
|
if not key.startsWith(SchemaFtsIndexPrefix) and
|
||||||
|
not key.startsWith(SchemaVecIndexPrefix): continue
|
||||||
let ddl = cast[string](value)
|
let ddl = cast[string](value)
|
||||||
if ddl.len == 0: continue
|
if ddl.len == 0: continue
|
||||||
ddls.add(ddl)
|
ddls.add(ddl)
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ import barabadb/storage/lsm
|
|||||||
import barabadb/query/executor
|
import barabadb/query/executor
|
||||||
import barabadb/query/parser
|
import barabadb/query/parser
|
||||||
import barabadb/fts/engine
|
import barabadb/fts/engine
|
||||||
|
import barabadb/vector/engine as vengine
|
||||||
|
|
||||||
proc execSql(ctx: ExecutionContext, sql: string): ExecResult =
|
proc execSql(ctx: ExecutionContext, sql: string): ExecResult =
|
||||||
let node = parse(sql)
|
let node = parse(sql)
|
||||||
@@ -137,6 +138,35 @@ suite "Schema persistence":
|
|||||||
db2.close()
|
db2.close()
|
||||||
removeDir(dir)
|
removeDir(dir)
|
||||||
|
|
||||||
|
test "HNSW vector index survives reopen":
|
||||||
|
let dir = "/tmp/baradb_schema_persist_vec"
|
||||||
|
removeDir(dir)
|
||||||
|
block:
|
||||||
|
var db = newLSMTree(dir)
|
||||||
|
var ctx = newExecutionContext(db)
|
||||||
|
check execSql(ctx, "CREATE TABLE vecs (id INTEGER PRIMARY KEY, embedding TEXT)").success
|
||||||
|
check execSql(ctx, "INSERT INTO vecs (id, embedding) VALUES (1, '[1.0, 0.0, 0.0]')").success
|
||||||
|
check execSql(ctx, "CREATE INDEX vecs_hnsw ON vecs (embedding) USING HNSW").success
|
||||||
|
check ctx.vectorIndexes.hasKey("vecs.embedding")
|
||||||
|
db.close()
|
||||||
|
# Reopen fresh context (simulates process restart)
|
||||||
|
block:
|
||||||
|
var db2 = newLSMTree(dir)
|
||||||
|
var ctx2 = newExecutionContext(db2)
|
||||||
|
# Index must be rebuilt from the persisted schema key — today it is
|
||||||
|
# silently missing, so vector searches return empty results after reopen.
|
||||||
|
check ctx2.vectorIndexes.hasKey("vecs.embedding")
|
||||||
|
if ctx2.vectorIndexes.hasKey("vecs.embedding"):
|
||||||
|
check vengine.search(ctx2.vectorIndexes["vecs.embedding"],
|
||||||
|
@[1.0'f32, 0.0'f32, 0.0'f32], k = 5).len >= 1
|
||||||
|
# index keeps updating after reopen
|
||||||
|
check execSql(ctx2, "INSERT INTO vecs (id, embedding) VALUES (2, '[0.0, 1.0, 0.0]')").success
|
||||||
|
if ctx2.vectorIndexes.hasKey("vecs.embedding"):
|
||||||
|
check vengine.search(ctx2.vectorIndexes["vecs.embedding"],
|
||||||
|
@[0.0'f32, 1.0'f32, 0.0'f32], k = 5).len >= 1
|
||||||
|
db2.close()
|
||||||
|
removeDir(dir)
|
||||||
|
|
||||||
test "Stable schema key format":
|
test "Stable schema key format":
|
||||||
check tableSchemaKey("users") == "_schema:tables:users"
|
check tableSchemaKey("users") == "_schema:tables:users"
|
||||||
check serializeTableDdl(TableDef(
|
check serializeTableDdl(TableDef(
|
||||||
|
|||||||
Reference in New Issue
Block a user