feat(persist): graphs survive restart (rebuild from backing tables)
This commit is contained in:
@@ -18,6 +18,7 @@ const
|
|||||||
SchemaPolicyPrefix* = "_schema:policies:"
|
SchemaPolicyPrefix* = "_schema:policies:"
|
||||||
SchemaFtsIndexPrefix* = "_schema:ftsidx:"
|
SchemaFtsIndexPrefix* = "_schema:ftsidx:"
|
||||||
SchemaVecIndexPrefix* = "_schema:vecidx:"
|
SchemaVecIndexPrefix* = "_schema:vecidx:"
|
||||||
|
SchemaGraphsPrefix* = "_schema:graphs:"
|
||||||
## 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:"
|
||||||
|
|
||||||
|
|||||||
@@ -903,6 +903,9 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu
|
|||||||
ctx.tables.del(name & "_nodes")
|
ctx.tables.del(name & "_nodes")
|
||||||
ctx.graphs.del(name)
|
ctx.graphs.del(name)
|
||||||
return errResult("Failed to create graph edges table: " & edgesRes.message)
|
return errResult("Failed to create graph edges table: " & edgesRes.message)
|
||||||
|
# Persist a marker so restoreEngines can rebuild the Graph from the
|
||||||
|
# backing tables after a restart. Written only on the success path.
|
||||||
|
ctx.db.put(SchemaGraphsPrefix & name, cast[seq[byte]]("CREATE GRAPH " & name))
|
||||||
return okResult(msg="CREATE GRAPH " & name)
|
return okResult(msg="CREATE GRAPH " & name)
|
||||||
|
|
||||||
of nkDropGraph:
|
of nkDropGraph:
|
||||||
@@ -912,6 +915,7 @@ proc executeQueryImpl(ctx: ExecutionContext, astNode: Node, params: seq[WireValu
|
|||||||
return okResult()
|
return okResult()
|
||||||
return errResult("Graph '" & name & "' does not exist")
|
return errResult("Graph '" & name & "' does not exist")
|
||||||
ctx.graphs.del(name)
|
ctx.graphs.del(name)
|
||||||
|
ctx.db.delete(SchemaGraphsPrefix & name)
|
||||||
var dropNodesSql = "DROP TABLE " & name & "_nodes"
|
var dropNodesSql = "DROP TABLE " & name & "_nodes"
|
||||||
var dropEdgesSql = "DROP TABLE " & name & "_edges"
|
var dropEdgesSql = "DROP TABLE " & name & "_edges"
|
||||||
let nodesTokens = qlex.tokenize(dropNodesSql)
|
let nodesTokens = qlex.tokenize(dropNodesSql)
|
||||||
@@ -1573,9 +1577,10 @@ 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/HNSW indexes) from persisted schema keys
|
## Rebuild ephemeral engines (FTS/HNSW indexes, graphs) from persisted
|
||||||
## after restoreSchema. Invoked via context.restoreEnginesHook at the end of
|
## schema keys after restoreSchema. Invoked via context.restoreEnginesHook
|
||||||
## newExecutionContext. Replay re-persists the same key, so it is idempotent.
|
## at the end of newExecutionContext. Index 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) and
|
if not key.startsWith(SchemaFtsIndexPrefix) and
|
||||||
@@ -1591,6 +1596,57 @@ proc restoreEngines*(ctx: ExecutionContext) =
|
|||||||
except CatchableError as e:
|
except CatchableError as e:
|
||||||
warn("restoreEngines: replay raised for DDL '" & ddl & "': " & e.msg)
|
warn("restoreEngines: replay raised for DDL '" & ddl & "': " & e.msg)
|
||||||
|
|
||||||
|
# Graphs cannot be replayed via CREATE GRAPH (the backing tables already
|
||||||
|
# exist after restart), so rebuild each Graph object from the rows of its
|
||||||
|
# <name>_nodes / <name>_edges backing tables. Row mapping mirrors the
|
||||||
|
# INSERT path in exec/dml.nim.
|
||||||
|
var graphNames: seq[string] = @[]
|
||||||
|
for (key, _) in ctx.db.scanAll():
|
||||||
|
if key.startsWith(SchemaGraphsPrefix):
|
||||||
|
let name = key[SchemaGraphsPrefix.len..^1]
|
||||||
|
if name.len > 0: graphNames.add(name)
|
||||||
|
for name in graphNames:
|
||||||
|
if name in ctx.graphs: continue
|
||||||
|
try:
|
||||||
|
var g = gengine.newGraph()
|
||||||
|
for row in execScan(ctx, name & "_nodes"):
|
||||||
|
try:
|
||||||
|
if "id" notin row: continue
|
||||||
|
let idStr = valueToString(row["id"])
|
||||||
|
if idStr.len == 0: continue
|
||||||
|
let nid = gengine.NodeId(parseUInt(idStr))
|
||||||
|
var label = ""
|
||||||
|
var props = initTable[string, string]()
|
||||||
|
for col, val in row:
|
||||||
|
if col == "node_label":
|
||||||
|
label = valueToString(val)
|
||||||
|
elif col != "id" and col != "properties" and
|
||||||
|
col != "$key" and col != "$value":
|
||||||
|
props[col] = valueToString(val)
|
||||||
|
gengine.addNodeWithId(g, nid, label, props)
|
||||||
|
except CatchableError:
|
||||||
|
discard
|
||||||
|
for row in execScan(ctx, name & "_edges"):
|
||||||
|
try:
|
||||||
|
if "source_id" notin row or "dest_id" notin row: continue
|
||||||
|
let srcStr = valueToString(row["source_id"])
|
||||||
|
let dstStr = valueToString(row["dest_id"])
|
||||||
|
if srcStr.len == 0 or dstStr.len == 0: continue
|
||||||
|
var label = ""
|
||||||
|
var weight = 1.0
|
||||||
|
if "edge_label" in row:
|
||||||
|
label = valueToString(row["edge_label"])
|
||||||
|
if "weight" in row:
|
||||||
|
try: weight = parseFloat(valueToString(row["weight"]))
|
||||||
|
except CatchableError: discard
|
||||||
|
gengine.addEdgeWithId(g, gengine.NodeId(parseUInt(srcStr)),
|
||||||
|
gengine.NodeId(parseUInt(dstStr)), label, weight)
|
||||||
|
except CatchableError:
|
||||||
|
discard
|
||||||
|
ctx.graphs[name] = g
|
||||||
|
except CatchableError as e:
|
||||||
|
warn("restoreEngines: graph rebuild failed for '" & name & "': " & e.msg)
|
||||||
|
|
||||||
# ----------------------------------------------------------------------
|
# ----------------------------------------------------------------------
|
||||||
# Hook wiring — breaks the module cycle between executor and the exec/*
|
# Hook wiring — breaks the module cycle between executor and the exec/*
|
||||||
# submodules: eval.nim calls back into the engine for subqueries, hybrid
|
# submodules: eval.nim calls back into the engine for subqueries, hybrid
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ 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
|
import barabadb/vector/engine as vengine
|
||||||
|
import barabadb/graph/engine as gengine
|
||||||
|
|
||||||
proc execSql(ctx: ExecutionContext, sql: string): ExecResult =
|
proc execSql(ctx: ExecutionContext, sql: string): ExecResult =
|
||||||
let node = parse(sql)
|
let node = parse(sql)
|
||||||
@@ -167,6 +168,31 @@ suite "Schema persistence":
|
|||||||
db2.close()
|
db2.close()
|
||||||
removeDir(dir)
|
removeDir(dir)
|
||||||
|
|
||||||
|
test "Graph survives reopen":
|
||||||
|
let dir = "/tmp/baradb_schema_persist_graph"
|
||||||
|
removeDir(dir)
|
||||||
|
block:
|
||||||
|
var db = newLSMTree(dir)
|
||||||
|
var ctx = newExecutionContext(db)
|
||||||
|
check execSql(ctx, "CREATE GRAPH social").success
|
||||||
|
check execSql(ctx, "INSERT INTO social_nodes (id, node_label) VALUES (1, 'person')").success
|
||||||
|
check execSql(ctx, "INSERT INTO social_nodes (id, node_label) VALUES (2, 'person')").success
|
||||||
|
check execSql(ctx, "INSERT INTO social_edges (source_id, dest_id, edge_label, weight) VALUES (1, 2, 'knows', 1.0)").success
|
||||||
|
check "social" in ctx.graphs
|
||||||
|
db.close()
|
||||||
|
# Reopen fresh context (simulates process restart)
|
||||||
|
block:
|
||||||
|
var db2 = newLSMTree(dir)
|
||||||
|
var ctx2 = newExecutionContext(db2)
|
||||||
|
# Graph must be rebuilt from the backing tables — today it is
|
||||||
|
# silently missing after reopen.
|
||||||
|
check "social" in ctx2.graphs
|
||||||
|
if "social" in ctx2.graphs:
|
||||||
|
check gengine.nodeCount(ctx2.graphs["social"]) == 2
|
||||||
|
check gengine.edgeCount(ctx2.graphs["social"]) == 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