use std::sync::Arc;
use crate::{
schema::definition::{DataType, EdgeMode, Schema, SchemaMode},
store::RocksStorage,
types::StoreError,
};
pub struct SchemaManagement {
store: Arc<RocksStorage>,
schema: Arc<std::sync::RwLock<Schema>>,
base_version: u64,
pending_vertex_labels: Vec<String>,
pending_edge_labels: Vec<String>,
pending_prop_keys: Vec<(String, DataType)>,
pending_edge_mode: Option<EdgeMode>,
pending_schema_mode: Option<SchemaMode>,
}
impl SchemaManagement {
pub(crate) fn new(store: Arc<RocksStorage>, schema: Arc<std::sync::RwLock<Schema>>) -> Self {
let base_version = schema.read().unwrap().version;
Self {
store,
schema,
base_version,
pending_vertex_labels: Vec::new(),
pending_edge_labels: Vec::new(),
pending_prop_keys: Vec::new(),
pending_edge_mode: None,
pending_schema_mode: None,
}
}
pub fn add_vertex_label(&mut self, name: impl Into<String>) -> &mut Self {
self.pending_vertex_labels.push(name.into());
self
}
pub fn add_edge_label(&mut self, name: impl Into<String>) -> &mut Self {
self.pending_edge_labels.push(name.into());
self
}
pub fn add_property_key(&mut self, name: impl Into<String>, data_type: DataType) -> &mut Self {
self.pending_prop_keys.push((name.into(), data_type));
self
}
pub fn set_edge_mode(&mut self, mode: EdgeMode) -> &mut Self {
self.pending_edge_mode = Some(mode);
self
}
pub fn set_schema_mode(&mut self, mode: SchemaMode) -> &mut Self {
self.pending_schema_mode = Some(mode);
self
}
pub fn commit(self) -> Result<(), StoreError> {
use crate::store::rocks::CF_SCHEMA;
use crate::types::kv_codec::{
encode_schema_key, encode_schema_label_value, encode_schema_meta, encode_schema_prop_value,
SCHEMA_KIND_EDGE_LABEL, SCHEMA_KIND_PROP_KEY, SCHEMA_KIND_VERTEX_LABEL, SCHEMA_META_KEY,
};
use rocksdb::WriteBatchWithTransaction;
let mut schema = self.schema.write().map_err(|_| StoreError::LockError)?;
if schema.version != self.base_version {
return Err(StoreError::SchemaConflict(format!(
"Concurrent schema modification: base version {}, current version {}",
self.base_version, schema.version
)));
}
let mut staged = schema.clone();
let mut changed = false;
if let Some(edge_mode) = self.pending_edge_mode {
changed |= staged.edge_mode != edge_mode;
staged.declare_edge_mode(edge_mode)?;
}
if let Some(schema_mode) = self.pending_schema_mode {
changed |= staged.mode != schema_mode;
staged.declare_schema_mode(schema_mode)?;
}
let cf = self.store.db.cf_handle(CF_SCHEMA).ok_or(StoreError::MissingColumnFamily(CF_SCHEMA))?;
let mut batch = WriteBatchWithTransaction::<true>::default();
for name in &self.pending_vertex_labels {
let is_new = staged.vertex_label_id(name).is_none();
let id = staged.declare_vertex_label(name)?;
changed |= is_new;
let key = encode_schema_key(SCHEMA_KIND_VERTEX_LABEL, name);
let val = encode_schema_label_value(id);
batch.put_cf(&cf, key, val);
staged.persisted_vertex_labels.insert(id);
}
for name in &self.pending_edge_labels {
let is_new = staged.edge_label_id(name).is_none();
let id = staged.declare_edge_label(name)?;
changed |= is_new;
let key = encode_schema_key(SCHEMA_KIND_EDGE_LABEL, name);
let val = encode_schema_label_value(id);
batch.put_cf(&cf, key, val);
staged.persisted_edge_labels.insert(id);
}
for (name, data_type) in &self.pending_prop_keys {
let is_new = staged.prop_key_id(name).is_none();
let id = staged.declare_prop_key(name, *data_type)?;
changed |= is_new;
let key = encode_schema_key(SCHEMA_KIND_PROP_KEY, name);
let val = encode_schema_prop_value(id, data_type.to_u8());
batch.put_cf(&cf, key, val);
staged.persisted_prop_keys.insert(id);
}
if !changed {
return Ok(());
}
staged.version += 1;
let meta_bytes = encode_schema_meta(staged.version, staged.edge_mode.to_u8(), staged.mode.to_u8());
batch.put_cf(&cf, SCHEMA_META_KEY, meta_bytes);
self.store.db.write(batch).map_err(StoreError::RocksDb)?;
*schema = staged;
Ok(())
}
}
impl std::fmt::Display for SchemaManagement {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let schema = self.schema.read().map_err(|_| std::fmt::Error)?;
writeln!(f, "=== RocksGraph Schema (Version: {}) ===", schema.version)?;
writeln!(f, "Schema Mode: {:?}", schema.mode)?;
writeln!(f, "Edge Mode: {:?}", schema.edge_mode)?;
writeln!(f, "\nVertex Labels:")?;
let mut vertex_labels: Vec<&str> = schema.vertex_labels.iter().map(|(_, name)| name.as_str()).collect();
vertex_labels.sort();
for label in vertex_labels {
writeln!(f, " - {}", label)?;
}
writeln!(f, "\nEdge Labels:")?;
let mut edge_labels: Vec<&str> = schema.edge_labels.iter().map(|(_, name)| name.as_str()).collect();
edge_labels.sort();
for label in edge_labels {
writeln!(f, " - {}", label)?;
}
writeln!(f, "\nProperty Keys:")?;
let mut prop_keys: Vec<(&str, DataType)> = schema
.prop_keys
.iter()
.map(|(id, name)| {
let data_type = schema.prop_key_types.get(id).map(|cfg| cfg.data_type).unwrap_or(DataType::Bytes);
(name.as_str(), data_type)
})
.collect();
prop_keys.sort_by_key(|(name, _)| *name);
for (name, data_type) in prop_keys {
writeln!(f, " - {} ({:?})", name, data_type)?;
}
write!(f, "======================================")
}
}