use crate::codec::{frame_value, unframe_value};
use crate::disk::{EdgeRecordDiskV3, EdgeRecordDiskV4};
use crate::error::{storage_err, TopoError};
use crate::ids::{EdgeId, NodeId, Scope};
use crate::op::Op;
use crate::props::{PropValue, Props};
use crate::storage::{EDGES, OPS};
use redb::{ReadableTable, WriteTransaction};
use serde::{Deserialize, Serialize};
use smol_str::SmolStr;
use std::collections::BTreeMap;
#[derive(Debug, Clone, Serialize, Deserialize)]
enum OpDiskV8 {
CreateNode {
id: NodeId,
scope: Scope,
label: SmolStr,
props: Props,
},
SetNodeProps {
id: NodeId,
props: BTreeMap<String, Option<PropValue>>,
},
SetEmbedding {
id: NodeId,
model: String,
vector: Vec<f32>,
},
RemoveNode {
id: NodeId,
},
CreateEdge {
id: EdgeId,
scope: Scope,
ty: SmolStr,
from: NodeId,
to: NodeId,
props: Props,
valid_from: Option<i64>,
},
CloseEdge {
id: EdgeId,
valid_to: Option<i64>,
},
}
fn convert_op(old: OpDiskV8) -> Op {
match old {
OpDiskV8::CreateNode {
id,
scope,
label,
props,
} => Op::CreateNode {
id,
scope,
label,
props,
},
OpDiskV8::SetNodeProps { id, props } => Op::SetNodeProps { id, props },
OpDiskV8::SetEmbedding { id, model, vector } => Op::SetEmbedding { id, model, vector },
OpDiskV8::RemoveNode { id } => Op::RemoveNode { id },
OpDiskV8::CreateEdge {
id,
scope,
ty,
from,
to,
props,
valid_from,
} => Op::CreateEdge {
id,
scope,
ty,
from,
to,
props,
valid_from,
recorded_at: valid_from,
},
OpDiskV8::CloseEdge { id, valid_to } => Op::CloseEdge {
id,
valid_to,
superseded_at: valid_to,
},
}
}
pub(crate) fn bitemporalize(tx: &WriteTransaction) -> Result<(), TopoError> {
{
let mut table = tx.open_table(EDGES).map_err(storage_err)?;
let mut rewrites: Vec<(Vec<u8>, Vec<u8>)> = Vec::new();
for entry in table.iter().map_err(storage_err)? {
let (k, v) = entry.map_err(storage_err)?;
let raw = unframe_value(v.value())?;
let old: EdgeRecordDiskV3 = postcard::from_bytes(&raw)
.map_err(|e| TopoError::Encoding(format!("v9 edges migration: {e}")))?;
let new = EdgeRecordDiskV4 {
id: old.id,
scope: old.scope,
ty: old.ty,
from: old.from,
to: old.to,
props: old.props,
valid_from: old.valid_from,
valid_to: old.valid_to,
recorded_at: old.valid_from,
superseded_at: old.valid_to,
};
let out = postcard::to_allocvec(&new)
.map_err(|e| TopoError::Encoding(format!("v9 edges migration: {e}")))?;
rewrites.push((k.value().to_vec(), frame_value(out)));
}
for (k, v) in rewrites {
table
.insert(k.as_slice(), v.as_slice())
.map_err(storage_err)?;
}
}
{
let mut table = tx.open_table(OPS).map_err(storage_err)?;
let mut rewrites: Vec<(u64, Vec<u8>)> = Vec::new();
for entry in table.iter().map_err(storage_err)? {
let (k, v) = entry.map_err(storage_err)?;
let old: OpDiskV8 = postcard::from_bytes(v.value())
.map_err(|e| TopoError::Encoding(format!("v9 ops migration: {e}")))?;
let new = convert_op(old);
let out = postcard::to_allocvec(&new)
.map_err(|e| TopoError::Encoding(format!("v9 ops migration: {e}")))?;
rewrites.push((k.value(), out));
}
for (k, v) in rewrites {
table.insert(k, v.as_slice()).map_err(storage_err)?;
}
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ids::ScopeId;
#[test]
fn old_op_bytes_error_under_new_decoder() {
let old_create = OpDiskV8::CreateEdge {
id: EdgeId::new(),
scope: Scope::Id(ScopeId::new()),
ty: "knows".into(),
from: NodeId::new(),
to: NodeId::new(),
props: Props::new(),
valid_from: Some(1_000),
};
let bytes = postcard::to_allocvec(&old_create).unwrap();
let decoded: Result<Op, _> = postcard::from_bytes(&bytes);
assert!(
decoded.is_err(),
"old-shaped CreateEdge bytes must NOT decode under the current Op shape"
);
let old_close = OpDiskV8::CloseEdge {
id: EdgeId::new(),
valid_to: Some(2_000),
};
let bytes = postcard::to_allocvec(&old_close).unwrap();
let decoded: Result<Op, _> = postcard::from_bytes(&bytes);
assert!(
decoded.is_err(),
"old-shaped CloseEdge bytes must NOT decode under the current Op shape"
);
}
#[test]
fn old_edge_bytes_error_under_new_decoder() {
let old = EdgeRecordDiskV3 {
id: EdgeId::new(),
scope: 1,
ty: 1,
from: 1,
to: 2,
props: BTreeMap::new(),
valid_from: 1_000,
valid_to: None,
};
let bytes = postcard::to_allocvec(&old).unwrap();
let decoded: Result<EdgeRecordDiskV4, _> = postcard::from_bytes(&bytes);
assert!(
decoded.is_err(),
"old-shaped EdgeRecordDiskV3 bytes must NOT decode under EdgeRecordDiskV4"
);
}
#[test]
fn new_op_bytes_decode_silently_under_old_decoder() {
let new_create = Op::CreateEdge {
id: EdgeId::new(),
scope: Scope::Id(ScopeId::new()),
ty: "knows".into(),
from: NodeId::new(),
to: NodeId::new(),
props: Props::new(),
valid_from: Some(1_000),
recorded_at: Some(1_000),
};
let bytes = postcard::to_allocvec(&new_create).unwrap();
let decoded: OpDiskV8 = postcard::from_bytes(&bytes)
.expect("new-shaped bytes must decode 'successfully' under the old decoder");
match decoded {
OpDiskV8::CreateEdge { valid_from, .. } => assert_eq!(valid_from, Some(1_000)),
other => panic!("unexpected variant: {other:?}"),
}
}
#[test]
fn bitemporalize_backfills_edges_and_rewrites_ops_in_place() {
let dir = tempfile::tempdir().unwrap();
let db = redb::Database::create(dir.path().join("t.redb")).unwrap();
let tx = db.begin_write().unwrap();
let edge_id = EdgeId::new();
{
let mut edges = tx.open_table(EDGES).unwrap();
let old = EdgeRecordDiskV3 {
id: edge_id,
scope: 1,
ty: 7,
from: 1,
to: 2,
props: BTreeMap::new(),
valid_from: 1_000,
valid_to: Some(2_000),
};
let raw = postcard::to_allocvec(&old).unwrap();
edges
.insert(1u64.to_be_bytes().as_slice(), frame_value(raw).as_slice())
.unwrap();
let mut ops = tx.open_table(OPS).unwrap();
let op = OpDiskV8::CreateEdge {
id: edge_id,
scope: Scope::Id(ScopeId::new()),
ty: "knows".into(),
from: NodeId::new(),
to: NodeId::new(),
props: Props::new(),
valid_from: Some(1_000),
};
ops.insert(1u64, postcard::to_allocvec(&op).unwrap().as_slice())
.unwrap();
let close = OpDiskV8::CloseEdge {
id: edge_id,
valid_to: Some(2_000),
};
ops.insert(2u64, postcard::to_allocvec(&close).unwrap().as_slice())
.unwrap();
}
bitemporalize(&tx).unwrap();
{
let edges = tx.open_table(EDGES).unwrap();
let raw = edges.get(1u64.to_be_bytes().as_slice()).unwrap().unwrap();
let unframed = unframe_value(raw.value()).unwrap();
let decoded: EdgeRecordDiskV4 = postcard::from_bytes(&unframed).unwrap();
assert_eq!(decoded.recorded_at, 1_000);
assert_eq!(decoded.superseded_at, Some(2_000));
let ops = tx.open_table(OPS).unwrap();
let raw = ops.get(1u64).unwrap().unwrap();
let decoded: Op = postcard::from_bytes(raw.value()).unwrap();
match decoded {
Op::CreateEdge { recorded_at, .. } => assert_eq!(recorded_at, Some(1_000)),
other => panic!("unexpected op: {other:?}"),
}
let raw = ops.get(2u64).unwrap().unwrap();
let decoded: Op = postcard::from_bytes(raw.value()).unwrap();
match decoded {
Op::CloseEdge { superseded_at, .. } => assert_eq!(superseded_at, Some(2_000)),
other => panic!("unexpected op: {other:?}"),
}
}
tx.commit().unwrap();
}
}