use std::collections::HashMap;
use std::fs::File;
use std::io::{BufReader, BufWriter, Read, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::time::Instant;
use rocksdb::{
ColumnFamilyDescriptor, IngestExternalFileOptions, OptimisticTransactionDB, Options, SstFileWriter,
WriteBatchWithTransaction,
};
use crate::{
schema::{
definition::{EdgeMode, SchemaMode},
DataType, GraphOptions,
},
store::rocks::cf_options,
types::{
gvalue::Primitive,
keys::{CanonicalEdgeKey, LabelId, Rank, VertexKey},
kv_codec::{
self, 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,
},
prop_codec, StoreError,
},
};
use super::{bulk_sort::ExternalSorter, CF_EDGES_IN, CF_EDGES_OUT, CF_SCHEMA, CF_VERTEX_DEGREE, CF_VERTICES};
pub(crate) const BULK_LOAD_IN_PROGRESS_KEY: &[u8] = b"_bulk_load_in_progress";
const DEFAULT_MAX_SST_SIZE: usize = 58 * 1024 * 1024;
const DEFAULT_MAX_MEMORY_BYTES: usize = 512 * 1024 * 1024;
#[derive(Debug)]
pub struct BulkSchema {
pub vertex_labels: Vec<String>,
pub edge_labels: Vec<String>,
pub prop_keys: Vec<(String, DataType)>,
}
pub struct BulkVertex {
pub id: VertexKey,
pub label: String,
pub props: HashMap<String, Primitive>,
}
pub struct BulkEdge {
pub src: VertexKey,
pub dst: VertexKey,
pub label: String,
pub props: HashMap<String, Primitive>,
pub rank: Option<Rank>,
}
#[derive(Debug)]
pub struct BulkLoadStats {
pub vertices_written: u64,
pub edges_written: u64,
pub sst_files: usize,
}
pub struct SstBulkLoader {
db_path: PathBuf,
work_dir: PathBuf,
max_sst_size: usize,
max_memory_bytes: usize,
}
impl SstBulkLoader {
pub fn new(db_path: impl Into<PathBuf>, work_dir: impl Into<PathBuf>) -> Self {
Self {
db_path: db_path.into(),
work_dir: work_dir.into(),
max_sst_size: DEFAULT_MAX_SST_SIZE,
max_memory_bytes: DEFAULT_MAX_MEMORY_BYTES,
}
}
pub fn with_max_sst_size(mut self, bytes: usize) -> Self {
self.max_sst_size = bytes;
self
}
pub fn with_max_memory(mut self, bytes: usize) -> Self {
self.max_memory_bytes = bytes;
self
}
}
struct ResolvedSchema {
vertex_label_ids: HashMap<String, LabelId>,
edge_label_ids: HashMap<String, LabelId>,
prop_key_ids: HashMap<String, u16>,
}
impl ResolvedSchema {
fn resolve_vertex_label(&self, name: &str) -> Result<LabelId, StoreError> {
self.vertex_label_ids
.get(name)
.copied()
.ok_or_else(|| StoreError::SchemaViolation(format!("unknown vertex label '{name}'")))
}
fn resolve_edge_label(&self, name: &str) -> Result<LabelId, StoreError> {
self.edge_label_ids
.get(name)
.copied()
.ok_or_else(|| StoreError::SchemaViolation(format!("unknown edge label '{name}'")))
}
fn resolve_prop_key(&self, name: &str) -> Result<u16, StoreError> {
self.prop_key_ids
.get(name)
.copied()
.ok_or_else(|| StoreError::SchemaViolation(format!("unknown property key '{name}'")))
}
fn encode_props(&self, props: &HashMap<String, Primitive>) -> Vec<u8> {
let id_props: HashMap<u16, Primitive> =
props.iter().filter_map(|(name, val)| self.prop_key_ids.get(name).map(|&id| (id, val.clone()))).collect();
prop_codec::encode_props(&id_props)
}
}
fn resolve_schema(schema: &BulkSchema) -> ResolvedSchema {
let vertex_label_ids: HashMap<String, LabelId> =
schema.vertex_labels.iter().enumerate().map(|(i, name)| (name.clone(), (i + 1) as LabelId)).collect();
let edge_label_ids: HashMap<String, LabelId> =
schema.edge_labels.iter().enumerate().map(|(i, name)| (name.clone(), (i + 1) as LabelId)).collect();
let prop_key_ids: HashMap<String, u16> =
schema.prop_keys.iter().enumerate().map(|(i, (name, _))| (name.clone(), (i + 4) as u16)).collect();
ResolvedSchema { vertex_label_ids, edge_label_ids, prop_key_ids }
}
fn write_schema_cf(
db: &OptimisticTransactionDB,
schema: &BulkSchema,
graph_opts: &GraphOptions,
) -> Result<ResolvedSchema, StoreError> {
let resolved = resolve_schema(schema);
let cf = db.cf_handle(CF_SCHEMA).ok_or(StoreError::MissingColumnFamily(CF_SCHEMA))?;
let mut batch = WriteBatchWithTransaction::<true>::default();
let meta = encode_schema_meta(1, graph_opts.edge_mode.to_u8(), graph_opts.mode.to_u8());
batch.put_cf(&cf, SCHEMA_META_KEY, meta);
batch.put_cf(&cf, BULK_LOAD_IN_PROGRESS_KEY, [1u8]);
for (name, &id) in &resolved.vertex_label_ids {
batch.put_cf(&cf, encode_schema_key(SCHEMA_KIND_VERTEX_LABEL, name), encode_schema_label_value(id));
}
for (name, &id) in &resolved.edge_label_ids {
batch.put_cf(&cf, encode_schema_key(SCHEMA_KIND_EDGE_LABEL, name), encode_schema_label_value(id));
}
let reserved: [(&str, u16, DataType); 3] =
[("id", 1, DataType::Int64), ("label", 2, DataType::Int32), ("rank", 3, DataType::UInt16)];
for (name, id, dt) in &reserved {
batch.put_cf(&cf, encode_schema_key(SCHEMA_KIND_PROP_KEY, name), encode_schema_prop_value(*id, dt.to_u8()));
}
for (i, (name, dt)) in schema.prop_keys.iter().enumerate() {
let id = (i + 4) as u16;
batch.put_cf(&cf, encode_schema_key(SCHEMA_KIND_PROP_KEY, name), encode_schema_prop_value(id, dt.to_u8()));
}
db.write(batch).map_err(StoreError::RocksDb)?;
Ok(resolved)
}
fn bulk_log(start: Instant, msg: &str) {
eprintln!("[bulk {:>7.1}s] {}", start.elapsed().as_secs_f64(), msg);
}
struct WorkDirGuard {
path: PathBuf,
active: bool,
}
impl Drop for WorkDirGuard {
fn drop(&mut self) {
if self.active {
let _ = std::fs::remove_dir_all(&self.path);
}
}
}
#[derive(Debug)]
struct SortedLabelFile {
path: PathBuf,
count: u64,
}
impl SortedLabelFile {
fn write_from(sorter: ExternalSorter, path: &Path) -> Result<Self, StoreError> {
let file = File::create(path).map_err(StoreError::Io)?;
let mut w = BufWriter::new(file);
w.write_all(&0u64.to_le_bytes()).map_err(StoreError::Io)?; let mut count = 0u64;
let mut last: Option<(VertexKey, LabelId)> = None;
for item in sorter.finish()? {
let (key, val) = item?;
let vid = VertexKey::from_be_bytes(
key.try_into().map_err(|_| StoreError::CorruptData("label sorter: key must be 8 bytes"))?,
);
let lid = LabelId::from_be_bytes(
val.try_into().map_err(|_| StoreError::CorruptData("label sorter: value must be 4 bytes"))?,
);
if let Some((lv, ll)) = last {
if lv == vid {
if ll != lid {
return Err(StoreError::SchemaViolation(format!(
"vertex {vid} appears with conflicting labels in input"
)));
}
continue; }
}
last = Some((vid, lid));
w.write_all(&vid.to_be_bytes()).map_err(StoreError::Io)?;
w.write_all(&lid.to_be_bytes()).map_err(StoreError::Io)?;
count += 1;
}
w.flush().map_err(StoreError::Io)?;
let mut file = w.into_inner().map_err(|e| StoreError::Io(e.into_error()))?;
file.seek(SeekFrom::Start(0)).map_err(StoreError::Io)?;
file.write_all(&count.to_le_bytes()).map_err(StoreError::Io)?;
Ok(Self { path: path.to_owned(), count })
}
fn reader(&self) -> Result<LabelFileIter, StoreError> {
let file = File::open(&self.path).map_err(StoreError::Io)?;
let mut reader = BufReader::new(file);
let mut buf = [0u8; 8];
reader.read_exact(&mut buf).map_err(StoreError::Io)?;
Ok(LabelFileIter { reader, remaining: u64::from_le_bytes(buf) })
}
}
impl Drop for SortedLabelFile {
fn drop(&mut self) {
let _ = std::fs::remove_file(&self.path);
}
}
struct LabelFileIter {
reader: BufReader<File>,
remaining: u64,
}
impl Iterator for LabelFileIter {
type Item = Result<(VertexKey, LabelId), StoreError>;
fn next(&mut self) -> Option<Self::Item> {
if self.remaining == 0 {
return None;
}
self.remaining -= 1;
let mut buf = [0u8; 12];
if let Err(e) = self.reader.read_exact(&mut buf) {
return Some(Err(StoreError::Io(e)));
}
Some(Ok((
VertexKey::from_be_bytes(buf[0..8].try_into().unwrap()),
LabelId::from_be_bytes(buf[8..12].try_into().unwrap()),
)))
}
}
struct DegreeCounter<I: Iterator<Item = Result<(Vec<u8>, Vec<u8>), StoreError>>> {
iter: I,
head: Option<VertexKey>,
}
impl<I: Iterator<Item = Result<(Vec<u8>, Vec<u8>), StoreError>>> DegreeCounter<I> {
fn new(mut iter: I) -> Result<Self, StoreError> {
let head = Self::advance(&mut iter)?;
Ok(Self { iter, head })
}
fn advance(iter: &mut I) -> Result<Option<VertexKey>, StoreError> {
match iter.next() {
None => Ok(None),
Some(Err(e)) => Err(e),
Some(Ok((key, _))) => Ok(Some(VertexKey::from_be_bytes(
key.try_into().map_err(|_| StoreError::CorruptData("degree sorter: key must be 8 bytes"))?,
))),
}
}
fn count_for(&mut self, vid: VertexKey) -> Result<u32, StoreError> {
let mut count = 0u32;
loop {
match self.head {
None => return Ok(count),
Some(cur) if cur < vid => {
self.head = Self::advance(&mut self.iter)?;
}
Some(cur) if cur == vid => {
count += 1;
self.head = Self::advance(&mut self.iter)?;
}
_ => return Ok(count),
}
}
}
}
fn annotate_edges(
annot_iter: impl Iterator<Item = Result<(Vec<u8>, Vec<u8>), StoreError>>,
label_file: &SortedLabelFile,
out_sorter: &mut ExternalSorter,
) -> Result<(), StoreError> {
let mut label_iter = label_file.reader()?;
let mut cur: Option<(VertexKey, LabelId)> = label_iter.next().transpose()?;
let mut cached: Option<(VertexKey, LabelId)> = None;
for item in annot_iter {
let (key, props) = item?;
if key.len() != 30 {
return Err(StoreError::CorruptData("annotation key must be 30 bytes"));
}
let lookup_id = VertexKey::from_be_bytes(key[0..8].try_into().unwrap());
let edge_key = key[8..30].to_vec();
let label = if cached.map(|(v, _)| v) == Some(lookup_id) {
cached.unwrap().1
} else {
loop {
match cur {
None => {
return Err(StoreError::SchemaViolation(format!(
"edge references vertex {lookup_id} not in vertex set"
)));
}
Some((vid, _)) if vid < lookup_id => {
cur = label_iter.next().transpose()?;
}
Some((vid, lid)) if vid == lookup_id => {
cached = Some((vid, lid));
break lid;
}
Some((vid, _)) => {
return Err(StoreError::SchemaViolation(format!(
"edge references vertex {lookup_id} not in vertex set (next in file: {vid})"
)));
}
}
}
};
out_sorter.push(edge_key, kv_codec::EdgeValue { end_vertex_label: label, property_blob: props }.encode())?;
}
Ok(())
}
fn write_degree_sst<I1, I2>(
label_file: &SortedLabelFile,
out_deg_iter: I1,
in_deg_iter: I2,
work_dir: &Path,
max_sst_size: usize,
cf_opts: &Options,
) -> Result<Vec<PathBuf>, StoreError>
where
I1: Iterator<Item = Result<(Vec<u8>, Vec<u8>), StoreError>>,
I2: Iterator<Item = Result<(Vec<u8>, Vec<u8>), StoreError>>,
{
if label_file.count == 0 {
return Ok(Vec::new());
}
let mut out_ctr = DegreeCounter::new(out_deg_iter)?;
let mut in_ctr = DegreeCounter::new(in_deg_iter)?;
let mut files = Vec::new();
let mut chunk = 0usize;
let mut path = work_dir.join(format!("bulk_vertex_degree_{chunk}.sst"));
let mut writer = SstFileWriter::create(cf_opts);
writer.open(&path).map_err(StoreError::RocksDb)?;
let mut written = 0usize;
for label_item in label_file.reader()? {
let (vid, lid) = label_item?;
let out_cnt = out_ctr.count_for(vid)?;
let in_cnt = in_ctr.count_for(vid)?;
let key = kv_codec::encode_vertex_key(vid);
let val = kv_codec::VertexDegree { vertex_label_id: lid, out_e_cnt: out_cnt, in_e_cnt: in_cnt }.encode();
if writer.file_size() >= max_sst_size as u64 {
writer.finish().map_err(StoreError::RocksDb)?;
files.push(path);
chunk += 1;
path = work_dir.join(format!("bulk_vertex_degree_{chunk}.sst"));
writer = SstFileWriter::create(cf_opts);
writer.open(&path).map_err(StoreError::RocksDb)?;
}
writer.put(key, val).map_err(StoreError::RocksDb)?;
written += 1;
}
if written > 0 {
writer.finish().map_err(StoreError::RocksDb)?;
files.push(path);
}
Ok(files)
}
struct SstPaths {
vertices: Vec<PathBuf>,
degree: Vec<PathBuf>,
edges_out: Vec<PathBuf>,
edges_in: Vec<PathBuf>,
}
impl SstPaths {
fn total_files(&self) -> usize {
self.vertices.len() + self.degree.len() + self.edges_out.len() + self.edges_in.len()
}
}
#[allow(clippy::type_complexity)]
fn write_sst_from_iter(
cf_name: &str,
iter: &mut dyn Iterator<Item = Result<(Vec<u8>, Vec<u8>), StoreError>>,
work_dir: &Path,
max_sst_size: usize,
cf_opts: &Options,
) -> Result<Vec<PathBuf>, StoreError> {
let mut files = Vec::new();
let mut chunk = 0usize;
let mut path = work_dir.join(format!("bulk_{cf_name}_{chunk}.sst"));
let mut writer = SstFileWriter::create(cf_opts);
writer.open(&path).map_err(StoreError::RocksDb)?;
let mut count = 0usize;
#[allow(clippy::while_let_loop)]
loop {
let result = match iter.next() {
Some(r) => r,
None => break,
};
let (key, val) = result?;
count += 1;
if writer.file_size() >= max_sst_size as u64 {
writer.finish().map_err(StoreError::RocksDb)?;
files.push(path);
chunk += 1;
path = work_dir.join(format!("bulk_{cf_name}_{chunk}.sst"));
writer = SstFileWriter::create(cf_opts);
writer.open(&path).map_err(StoreError::RocksDb)?;
}
writer.put(key, val).map_err(StoreError::RocksDb)?;
}
if count > 0 {
writer.finish().map_err(StoreError::RocksDb)?;
files.push(path);
}
Ok(files)
}
fn write_sst_from_iter_dedup(
cf_name: &str,
iter: impl Iterator<Item = Result<(Vec<u8>, Vec<u8>), StoreError>>,
work_dir: &Path,
max_sst_size: usize,
cf_opts: &Options,
) -> Result<Vec<PathBuf>, StoreError> {
let mut last_key: Option<Vec<u8>> = None;
let mut deduped = iter.map(move |r| {
let (key, val) = r?;
if last_key.as_deref() == Some(&key[..]) {
let cek = kv_codec::decode_edge_key(&key, crate::types::keys::Direction::OUT)
.map(|ek| CanonicalEdgeKey {
src_id: ek.primary_id,
label_id: ek.label_id,
dst_id: ek.secondary_id,
rank: ek.rank,
})
.ok_or(StoreError::CorruptData("duplicate edge key could not be decoded"))?;
return Err(StoreError::DuplicateEdge(cek));
}
last_key = Some(key.clone());
Ok((key, val))
});
write_sst_from_iter(cf_name, &mut deduped, work_dir, max_sst_size, cf_opts)
}
impl SstBulkLoader {
pub fn load_initial(
self,
schema: BulkSchema,
vertices: impl Iterator<Item = BulkVertex>,
edges: impl Iterator<Item = BulkEdge>,
graph_opts: GraphOptions,
rocks_opts: &crate::store::rocks::store::RocksOptions,
) -> Result<BulkLoadStats, StoreError> {
let t0 = Instant::now();
bulk_log(t0, "starting bulk load");
if self.work_dir.exists() {
std::fs::remove_dir_all(&self.work_dir).map_err(StoreError::Io)?;
}
std::fs::create_dir_all(&self.work_dir).map_err(StoreError::Io)?;
let mut work_dir_guard = WorkDirGuard { path: self.work_dir.clone(), active: true };
let mut db_opts = Options::default();
db_opts.create_if_missing(true);
db_opts.create_missing_column_families(true);
db_opts.set_max_background_jobs(rocks_opts.max_background_jobs);
let v_bo = cf_options::vertex_block_opts(rocks_opts);
let e_bo = cf_options::edge_block_opts(rocks_opts);
let cfs = vec![
ColumnFamilyDescriptor::new(CF_VERTICES, cf_options::vertex_cf_opts(rocks_opts, &v_bo)),
ColumnFamilyDescriptor::new(CF_VERTEX_DEGREE, cf_options::vertex_cf_opts(rocks_opts, &v_bo)),
ColumnFamilyDescriptor::new(CF_EDGES_OUT, cf_options::edge_cf_opts(rocks_opts, &e_bo)),
ColumnFamilyDescriptor::new(CF_EDGES_IN, cf_options::edge_cf_opts(rocks_opts, &e_bo)),
ColumnFamilyDescriptor::new(CF_SCHEMA, Options::default()),
];
let db =
OptimisticTransactionDB::open_cf_descriptors(&db_opts, &self.db_path, cfs).map_err(StoreError::RocksDb)?;
bulk_log(t0, "writing schema CF");
let resolved = write_schema_cf(&db, &schema, &graph_opts)?;
let budget_v = self.max_memory_bytes / 4;
let mut vertex_sorter = ExternalSorter::new(self.work_dir.join("sv"), budget_v);
let mut label_sorter = ExternalSorter::new(self.work_dir.join("sl"), budget_v);
let mut vcount = 0u64;
bulk_log(t0, "phase 1 — streaming vertices");
for v in vertices {
let lid = resolved.resolve_vertex_label(&v.label)?;
if graph_opts.mode == SchemaMode::Strict {
for k in v.props.keys() {
resolved.resolve_prop_key(k)?;
}
}
let blob = resolved.encode_props(&v.props);
vertex_sorter.push(
kv_codec::encode_vertex_key(v.id).to_vec(),
kv_codec::VertexValue { label_id: lid, property_blob: blob }.encode(),
)?;
label_sorter.push(v.id.to_be_bytes().to_vec(), lid.to_be_bytes().to_vec())?;
vcount += 1;
}
let label_file_path = self.work_dir.join("vertex_labels.bin");
let label_file = SortedLabelFile::write_from(label_sorter, &label_file_path)?;
bulk_log(t0, &format!("phase 1 — {vcount} vertices; label file written ({} unique)", label_file.count));
let budget_a = self.max_memory_bytes / 4;
let budget_d = self.max_memory_bytes / 8;
let mut dst_annot = ExternalSorter::new(self.work_dir.join("ea_dst"), budget_a);
let mut src_annot = ExternalSorter::new(self.work_dir.join("ea_src"), budget_a);
let mut out_deg = ExternalSorter::new(self.work_dir.join("deg_out"), budget_d);
let mut in_deg = ExternalSorter::new(self.work_dir.join("deg_in"), budget_d);
let ecount;
match graph_opts.edge_mode {
EdgeMode::Single => {
bulk_log(t0, "phase 1 — streaming edges (Single mode)");
let mut n = 0u64;
let mut last_report = Instant::now();
for edge in edges {
let lid = resolved.resolve_edge_label(&edge.label)?;
if graph_opts.mode == SchemaMode::Strict {
for k in edge.props.keys() {
resolved.resolve_prop_key(k)?;
}
}
let blob = resolved.encode_props(&edge.props);
let cek = CanonicalEdgeKey { src_id: edge.src, label_id: lid, dst_id: edge.dst, rank: 0 };
let mut dk = [0u8; 30];
dk[0..8].copy_from_slice(&edge.dst.to_be_bytes());
dk[8..30].copy_from_slice(&kv_codec::encode_edge_key(&cek.out_key()));
dst_annot.push(dk.to_vec(), blob.clone())?;
let mut sk = [0u8; 30];
sk[0..8].copy_from_slice(&edge.src.to_be_bytes());
sk[8..30].copy_from_slice(&kv_codec::encode_edge_key(&cek.in_key()));
src_annot.push(sk.to_vec(), blob)?;
out_deg.push(edge.src.to_be_bytes().to_vec(), vec![])?;
in_deg.push(edge.dst.to_be_bytes().to_vec(), vec![])?;
n += 1;
if n % 1_000_000 == 0 && last_report.elapsed().as_secs_f64() >= 5.0 {
bulk_log(t0, &format!(" {n} edges streamed"));
last_report = Instant::now();
}
}
ecount = n;
}
EdgeMode::Multi => {
let pre_budget = self.max_memory_bytes / 4;
let mut pre_sorter = ExternalSorter::new(self.work_dir.join("sm_pre"), pre_budget);
let mut n = 0u64;
let mut last_report = Instant::now();
bulk_log(t0, "phase 1 — streaming edges into pre-sorter (Multi mode)");
for edge in edges {
let lid = resolved.resolve_edge_label(&edge.label)?;
if graph_opts.mode == SchemaMode::Strict {
for k in edge.props.keys() {
resolved.resolve_prop_key(k)?;
}
}
let rank_for_sort = edge.rank.unwrap_or(Rank::MAX);
let mut sort_key = [0u8; 22];
sort_key[0..8].copy_from_slice(&edge.src.to_be_bytes());
sort_key[8..12].copy_from_slice(&lid.to_be_bytes());
sort_key[12..20].copy_from_slice(&edge.dst.to_be_bytes());
sort_key[20..22].copy_from_slice(&rank_for_sort.to_be_bytes());
pre_sorter.push(sort_key.to_vec(), resolved.encode_props(&edge.props))?;
n += 1;
if n % 1_000_000 == 0 && last_report.elapsed().as_secs_f64() >= 5.0 {
bulk_log(t0, &format!(" {n} edges streamed"));
last_report = Instant::now();
}
}
bulk_log(t0, &format!("phase 1 — {n} edges pre-sorted; assigning ranks"));
let mut last_prefix: Option<[u8; 20]> = None;
let mut next_rank: Rank = 0;
let mut last_explicit: Option<Rank> = None;
for item in pre_sorter.finish()? {
let (sort_key, blob) = item?;
let key22: [u8; 22] = sort_key
.try_into()
.map_err(|_| StoreError::CorruptData("corrupt pre-sort key in Multi mode"))?;
let prefix: [u8; 20] = key22[0..20].try_into().unwrap();
let rank_from_key = Rank::from_be_bytes(key22[20..22].try_into().unwrap());
let src = VertexKey::from_be_bytes(prefix[0..8].try_into().unwrap());
let lid = LabelId::from_be_bytes(prefix[8..12].try_into().unwrap());
let dst = VertexKey::from_be_bytes(prefix[12..20].try_into().unwrap());
if Some(prefix) != last_prefix {
next_rank = 0;
last_explicit = None;
last_prefix = Some(prefix);
}
let rank = if rank_from_key == Rank::MAX {
let r = next_rank;
next_rank = next_rank.checked_add(1).ok_or_else(|| {
StoreError::SchemaViolation(format!("rank overflow for edge ({src}->{dst})"))
})?;
r
} else {
if last_explicit == Some(rank_from_key) {
return Err(StoreError::DuplicateEdge(CanonicalEdgeKey {
src_id: src,
label_id: lid,
dst_id: dst,
rank: rank_from_key,
}));
}
last_explicit = Some(rank_from_key);
next_rank = rank_from_key.checked_add(1).ok_or_else(|| {
StoreError::SchemaViolation(format!("rank overflow for edge ({src}->{dst})"))
})?;
rank_from_key
};
let cek = CanonicalEdgeKey { src_id: src, label_id: lid, dst_id: dst, rank };
let mut dk = [0u8; 30];
dk[0..8].copy_from_slice(&dst.to_be_bytes());
dk[8..30].copy_from_slice(&kv_codec::encode_edge_key(&cek.out_key()));
dst_annot.push(dk.to_vec(), blob.clone())?;
let mut sk = [0u8; 30];
sk[0..8].copy_from_slice(&src.to_be_bytes());
sk[8..30].copy_from_slice(&kv_codec::encode_edge_key(&cek.in_key()));
src_annot.push(sk.to_vec(), blob)?;
out_deg.push(src.to_be_bytes().to_vec(), vec![])?;
in_deg.push(dst.to_be_bytes().to_vec(), vec![])?;
}
ecount = n;
}
}
bulk_log(t0, &format!("phase 1 done — {ecount} edges"));
let v_bo = cf_options::vertex_block_opts(rocks_opts);
let v_opts = cf_options::vertex_cf_opts(rocks_opts, &v_bo);
let e_bo = cf_options::edge_block_opts(rocks_opts);
let e_opts = cf_options::edge_cf_opts(rocks_opts, &e_bo);
bulk_log(t0, &format!("phase 2 — writing vertex SSTs ({vcount} vertices)"));
let vert_files =
write_sst_from_iter("vertices", &mut vertex_sorter.finish()?, &self.work_dir, self.max_sst_size, &v_opts)?;
bulk_log(t0, "phase 2 — writing degree SSTs");
let deg_files = write_degree_sst(
&label_file,
out_deg.finish()?,
in_deg.finish()?,
&self.work_dir,
self.max_sst_size,
&v_opts,
)?;
bulk_log(t0, "phase 2 — annotating + writing edges_out SSTs");
let mut out_edge_sorter = ExternalSorter::new(self.work_dir.join("eo"), self.max_memory_bytes);
annotate_edges(dst_annot.finish()?, &label_file, &mut out_edge_sorter)?;
let out_files = write_sst_from_iter_dedup(
"edges_out",
out_edge_sorter.finish()?,
&self.work_dir,
self.max_sst_size,
&e_opts,
)?;
bulk_log(t0, "phase 2 — annotating + writing edges_in SSTs");
let mut in_edge_sorter = ExternalSorter::new(self.work_dir.join("ei"), self.max_memory_bytes);
annotate_edges(src_annot.finish()?, &label_file, &mut in_edge_sorter)?;
let in_files =
write_sst_from_iter("edges_in", &mut in_edge_sorter.finish()?, &self.work_dir, self.max_sst_size, &e_opts)?;
let sst_paths = SstPaths { vertices: vert_files, degree: deg_files, edges_out: out_files, edges_in: in_files };
bulk_log(t0, &format!("phase 3 — ingesting {} SST files (atomic)", sst_paths.total_files()));
let mut ingest_opts = IngestExternalFileOptions::default();
ingest_opts.set_move_files(true);
macro_rules! ingest {
($paths:expr, $cf_name:expr) => {
if !$paths.is_empty() {
let cf = db.cf_handle($cf_name).ok_or(StoreError::MissingColumnFamily($cf_name))?;
db.ingest_external_file_cf_opts(&cf, &ingest_opts, $paths.to_vec()).map_err(StoreError::RocksDb)?;
}
};
}
ingest!(&sst_paths.vertices, CF_VERTICES);
ingest!(&sst_paths.degree, CF_VERTEX_DEGREE);
ingest!(&sst_paths.edges_out, CF_EDGES_OUT);
ingest!(&sst_paths.edges_in, CF_EDGES_IN);
let cf_sch = db.cf_handle(CF_SCHEMA).ok_or(StoreError::MissingColumnFamily(CF_SCHEMA))?;
let mut cleanup = WriteBatchWithTransaction::<true>::default();
cleanup.delete_cf(&cf_sch, BULK_LOAD_IN_PROGRESS_KEY);
db.write(cleanup).map_err(StoreError::RocksDb)?;
let n_files = sst_paths.total_files();
bulk_log(t0, &format!("done — {vcount} vertices, {ecount} edges, {n_files} SST files"));
bulk_log(t0, "compacting all CFs (moves L0 SSTs into deeper levels for fast scans)");
for cf_name in [CF_VERTICES, CF_VERTEX_DEGREE, CF_EDGES_OUT, CF_EDGES_IN] {
if let Some(cf) = db.cf_handle(cf_name) {
db.compact_range_cf(&cf, None::<&[u8]>, None::<&[u8]>);
}
}
bulk_log(t0, "compaction done");
work_dir_guard.active = false; Ok(BulkLoadStats { vertices_written: vcount, edges_written: ecount, sst_files: n_files })
}
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use tempfile::tempdir;
use crate::{
api::Graph,
gremlin::{traversal::TraversalBuilder, value::Value},
schema::{
definition::{EdgeMode, SchemaMode},
GraphOptions,
},
store::rocks::store::RocksOptions,
};
use super::*;
fn small_schema() -> BulkSchema {
BulkSchema { vertex_labels: vec!["Person".into()], edge_labels: vec!["Knows".into()], prop_keys: vec![] }
}
fn small_vertices() -> Vec<BulkVertex> {
(1..=5).map(|i| BulkVertex { id: i, label: "Person".into(), props: HashMap::new() }).collect()
}
fn small_edges() -> Vec<BulkEdge> {
vec![
BulkEdge { src: 1, dst: 2, label: "Knows".into(), props: HashMap::new(), rank: None },
BulkEdge { src: 1, dst: 3, label: "Knows".into(), props: HashMap::new(), rank: None },
BulkEdge { src: 2, dst: 3, label: "Knows".into(), props: HashMap::new(), rank: None },
BulkEdge { src: 3, dst: 4, label: "Knows".into(), props: HashMap::new(), rank: None },
BulkEdge { src: 4, dst: 5, label: "Knows".into(), props: HashMap::new(), rank: None },
BulkEdge { src: 5, dst: 1, label: "Knows".into(), props: HashMap::new(), rank: None },
]
}
#[test]
fn test_load_initial_small() {
let dir = tempdir().unwrap();
let db_path = dir.path().join("db");
let loader = SstBulkLoader::new(&db_path, dir.path().join("_bulk_work"));
let stats = loader
.load_initial(
small_schema(),
small_vertices().into_iter(),
small_edges().into_iter(),
GraphOptions::default(),
&RocksOptions::default(),
)
.unwrap();
assert_eq!(stats.vertices_written, 5);
assert_eq!(stats.edges_written, 6);
assert!(stats.sst_files >= 4);
let graph = Graph::open(&db_path).unwrap();
let mut snap = graph.read();
let v_count = snap.g().V([]).count().next().unwrap().unwrap();
assert_eq!(v_count, Value::Int64(5));
let out_count = snap.g().V([1_i64]).out(["Knows"]).count().next().unwrap().unwrap();
assert_eq!(out_count, Value::Int64(2));
graph.close().unwrap();
}
#[test]
fn test_duplicate_edge_single_mode() {
let dir = tempdir().unwrap();
let vertices: Vec<BulkVertex> =
(1..=3).map(|i| BulkVertex { id: i, label: "Person".into(), props: HashMap::new() }).collect();
let edges = vec![
BulkEdge { src: 1, dst: 2, label: "Knows".into(), props: HashMap::new(), rank: None },
BulkEdge { src: 1, dst: 2, label: "Knows".into(), props: HashMap::new(), rank: None },
];
let err = SstBulkLoader::new(dir.path().join("db"), dir.path().join("_bulk_work"))
.load_initial(
small_schema(),
vertices.into_iter(),
edges.into_iter(),
GraphOptions { edge_mode: EdgeMode::Single, ..Default::default() },
&RocksOptions::default(),
)
.unwrap_err();
assert!(matches!(err, StoreError::DuplicateEdge(_)));
}
#[test]
fn test_strict_mode_unknown_label() {
let dir = tempdir().unwrap();
let vertices = vec![BulkVertex { id: 1, label: "Unknown".into(), props: HashMap::new() }];
let err = SstBulkLoader::new(dir.path().join("db"), dir.path().join("_bulk_work"))
.load_initial(
small_schema(),
vertices.into_iter(),
std::iter::empty(),
GraphOptions { mode: SchemaMode::Strict, ..Default::default() },
&RocksOptions::default(),
)
.unwrap_err();
assert!(matches!(err, StoreError::SchemaViolation(_)));
}
#[test]
fn test_crash_marker_detection() {
let dir = tempdir().unwrap();
let db_path = dir.path().join("db");
SstBulkLoader::new(&db_path, dir.path().join("_bulk_work"))
.load_initial(
small_schema(),
small_vertices().into_iter(),
small_edges().into_iter(),
GraphOptions::default(),
&RocksOptions::default(),
)
.unwrap();
{
use super::super::{CF_EDGES_IN, CF_EDGES_OUT, CF_SCHEMA, CF_VERTEX_DEGREE, CF_VERTICES};
use crate::store::rocks::cf_options;
use rocksdb::{
ColumnFamilyDescriptor, MultiThreaded, OptimisticTransactionDB, Options, WriteBatchWithTransaction,
};
let rocks_opts = RocksOptions::default();
let v_bo = cf_options::vertex_block_opts(&rocks_opts);
let e_bo = cf_options::edge_block_opts(&rocks_opts);
let mut dbo = Options::default();
dbo.create_if_missing(false);
let cfs = vec![
ColumnFamilyDescriptor::new(CF_VERTICES, cf_options::vertex_cf_opts(&rocks_opts, &v_bo)),
ColumnFamilyDescriptor::new(CF_VERTEX_DEGREE, cf_options::vertex_cf_opts(&rocks_opts, &v_bo)),
ColumnFamilyDescriptor::new(CF_EDGES_OUT, cf_options::edge_cf_opts(&rocks_opts, &e_bo)),
ColumnFamilyDescriptor::new(CF_EDGES_IN, cf_options::edge_cf_opts(&rocks_opts, &e_bo)),
ColumnFamilyDescriptor::new(CF_SCHEMA, Options::default()),
];
let db: OptimisticTransactionDB<MultiThreaded> =
OptimisticTransactionDB::open_cf_descriptors(&dbo, &db_path, cfs).unwrap();
let cf = db.cf_handle(CF_SCHEMA).unwrap();
let mut batch = WriteBatchWithTransaction::<true>::default();
batch.put_cf(&cf, BULK_LOAD_IN_PROGRESS_KEY, [1u8]);
db.write(batch).unwrap();
}
let graph = Graph::open(&db_path).unwrap();
let mut snap = graph.read();
let v_count = snap.g().V([]).count().next().unwrap().unwrap();
assert_eq!(v_count, Value::Int64(5));
graph.close().unwrap();
}
#[test]
fn test_multiple_edges_multi_mode_auto_rank() {
let dir = tempdir().unwrap();
let vertices: Vec<BulkVertex> =
(1..=2).map(|i| BulkVertex { id: i, label: "Person".into(), props: HashMap::new() }).collect();
let edges = vec![
BulkEdge { src: 1, dst: 2, label: "Knows".into(), props: HashMap::new(), rank: None },
BulkEdge { src: 1, dst: 2, label: "Knows".into(), props: HashMap::new(), rank: None },
];
let stats = SstBulkLoader::new(dir.path().join("db"), dir.path().join("_bulk_work"))
.load_initial(
small_schema(),
vertices.into_iter(),
edges.into_iter(),
GraphOptions { edge_mode: EdgeMode::Multi, ..Default::default() },
&RocksOptions::default(),
)
.unwrap();
assert_eq!(stats.edges_written, 2);
let graph = Graph::open(dir.path().join("db")).unwrap();
let mut snap = graph.read();
let edges_count = snap.g().V([1_i64]).outE(["Knows"]).count().next().unwrap().unwrap();
assert_eq!(edges_count, Value::Int64(2));
graph.close().unwrap();
}
#[test]
fn test_edge_referencing_unknown_vertex() {
let dir = tempdir().unwrap();
let vertices = vec![BulkVertex { id: 1, label: "Person".into(), props: HashMap::new() }];
let edges = vec![BulkEdge { src: 1, dst: 2, label: "Knows".into(), props: HashMap::new(), rank: None }];
let err = SstBulkLoader::new(dir.path().join("db"), dir.path().join("_bulk_work"))
.load_initial(
small_schema(),
vertices.into_iter(),
edges.into_iter(),
GraphOptions::default(),
&RocksOptions::default(),
)
.unwrap_err();
assert!(matches!(err, StoreError::SchemaViolation(_)));
}
#[test]
fn test_load_initial_external_sort() {
let dir = tempdir().unwrap();
let db_path = dir.path().join("db");
let loader = SstBulkLoader::new(&db_path, dir.path().join("_bulk_work")).with_max_memory(1);
let stats = loader
.load_initial(
small_schema(),
small_vertices().into_iter(),
small_edges().into_iter(),
GraphOptions::default(),
&RocksOptions::default(),
)
.unwrap();
assert_eq!(stats.vertices_written, 5);
assert_eq!(stats.edges_written, 6);
let graph = Graph::open(&db_path).unwrap();
let mut snap = graph.read();
let v_count = snap.g().V([]).count().next().unwrap().unwrap();
assert_eq!(v_count, Value::Int64(5));
let out_count = snap.g().V([1_i64]).out(["Knows"]).count().next().unwrap().unwrap();
assert_eq!(out_count, Value::Int64(2));
graph.close().unwrap();
}
#[test]
fn test_dedup_iter_returns_correct_duplicate_edge() {
use crate::{
store::rocks::cf_options,
types::{keys::CanonicalEdgeKey, kv_codec},
};
let dir = tempdir().unwrap();
let rocks_opts = RocksOptions::default();
let e_bo = cf_options::edge_block_opts(&rocks_opts);
let e_opts = cf_options::edge_cf_opts(&rocks_opts, &e_bo);
let cek = CanonicalEdgeKey { src_id: 10, label_id: 1, dst_id: 20, rank: 0 };
let key = kv_codec::encode_edge_key(&cek.out_key()).to_vec();
let val = kv_codec::EdgeValue { end_vertex_label: 1, property_blob: vec![] }.encode();
let pairs = vec![Ok((key.clone(), val.clone())), Ok((key.clone(), val.clone()))];
let err = write_sst_from_iter_dedup("test_edges", pairs.into_iter(), dir.path(), 64 * 1024 * 1024, &e_opts)
.unwrap_err();
match err {
StoreError::DuplicateEdge(detected) => {
assert_eq!(detected.src_id, 10);
assert_eq!(detected.label_id, 1);
assert_eq!(detected.dst_id, 20);
assert_eq!(detected.rank, 0);
}
other => panic!("expected DuplicateEdge, got: {other}"),
}
}
#[test]
fn test_sorted_label_file_basic_and_dedup() {
let dir = tempdir().unwrap();
let path = dir.path().join("labels.bin");
let mut sorter = ExternalSorter::new(dir.path().join("sort"), 1024 * 1024);
sorter.push(20i64.to_be_bytes().to_vec(), 2i32.to_be_bytes().to_vec()).unwrap();
sorter.push(10i64.to_be_bytes().to_vec(), 1i32.to_be_bytes().to_vec()).unwrap();
sorter.push(10i64.to_be_bytes().to_vec(), 1i32.to_be_bytes().to_vec()).unwrap(); sorter.push(30i64.to_be_bytes().to_vec(), 3i32.to_be_bytes().to_vec()).unwrap();
let file = SortedLabelFile::write_from(sorter, &path).unwrap();
assert_eq!(file.count, 3);
let reader: Vec<_> = file.reader().unwrap().map(|r| r.unwrap()).collect();
assert_eq!(reader, vec![(10, 1), (20, 2), (30, 3)]);
}
#[test]
fn test_sorted_label_file_conflicting_labels() {
let dir = tempdir().unwrap();
let path = dir.path().join("labels.bin");
let mut sorter = ExternalSorter::new(dir.path().join("sort"), 1024 * 1024);
sorter.push(10i64.to_be_bytes().to_vec(), 1i32.to_be_bytes().to_vec()).unwrap();
sorter.push(10i64.to_be_bytes().to_vec(), 2i32.to_be_bytes().to_vec()).unwrap();
let err = SortedLabelFile::write_from(sorter, &path).unwrap_err();
assert!(matches!(err, StoreError::SchemaViolation(_)));
}
#[test]
#[allow(clippy::type_complexity)]
fn test_degree_counter() {
let items: Vec<Result<(Vec<u8>, Vec<u8>), StoreError>> = vec![
Ok((10i64.to_be_bytes().to_vec(), vec![])),
Ok((10i64.to_be_bytes().to_vec(), vec![])),
Ok((10i64.to_be_bytes().to_vec(), vec![])),
Ok((20i64.to_be_bytes().to_vec(), vec![])),
Ok((30i64.to_be_bytes().to_vec(), vec![])),
Ok((30i64.to_be_bytes().to_vec(), vec![])),
];
let mut counter = DegreeCounter::new(items.into_iter()).unwrap();
assert_eq!(counter.count_for(10).unwrap(), 3);
assert_eq!(counter.count_for(15).unwrap(), 0); assert_eq!(counter.count_for(20).unwrap(), 1);
assert_eq!(counter.count_for(30).unwrap(), 2);
assert_eq!(counter.count_for(40).unwrap(), 0); }
#[test]
#[allow(clippy::type_complexity)]
fn test_annotate_edges_mismatched_vertex() {
let dir = tempdir().unwrap();
let path = dir.path().join("labels.bin");
let mut sorter = ExternalSorter::new(dir.path().join("sort"), 1024 * 1024);
sorter.push(10i64.to_be_bytes().to_vec(), 1i32.to_be_bytes().to_vec()).unwrap();
let file = SortedLabelFile::write_from(sorter, &path).unwrap();
let annot_item: Vec<Result<(Vec<u8>, Vec<u8>), StoreError>> = vec![Ok((
{
let mut k = vec![0u8; 30];
k[0..8].copy_from_slice(&20i64.to_be_bytes());
k
},
vec![],
))];
let mut out_sorter = ExternalSorter::new(dir.path().join("out"), 1024 * 1024);
let err = annotate_edges(annot_item.into_iter(), &file, &mut out_sorter).unwrap_err();
assert!(matches!(err, StoreError::SchemaViolation(_)));
}
#[test]
fn test_empty_input() {
let dir = tempdir().unwrap();
let db_path = dir.path().join("db");
let stats = SstBulkLoader::new(&db_path, dir.path().join("_w"))
.load_initial(
small_schema(),
std::iter::empty(),
std::iter::empty(),
GraphOptions::default(),
&RocksOptions::default(),
)
.unwrap();
assert_eq!(stats.vertices_written, 0);
assert_eq!(stats.edges_written, 0);
let graph = Graph::open(&db_path).unwrap();
let mut snap = graph.read();
assert_eq!(snap.g().V([]).count().next().unwrap().unwrap(), Value::Int64(0));
graph.close().unwrap();
}
#[test]
fn test_vertices_only_no_edges() {
let dir = tempdir().unwrap();
let db_path = dir.path().join("db");
let vertices: Vec<BulkVertex> =
(1..=3).map(|i| BulkVertex { id: i, label: "Person".into(), props: HashMap::new() }).collect();
let stats = SstBulkLoader::new(&db_path, dir.path().join("_w"))
.load_initial(
small_schema(),
vertices.into_iter(),
std::iter::empty(),
GraphOptions::default(),
&RocksOptions::default(),
)
.unwrap();
assert_eq!(stats.vertices_written, 3);
assert_eq!(stats.edges_written, 0);
let graph = Graph::open(&db_path).unwrap();
let mut snap = graph.read();
assert_eq!(snap.g().V([]).count().next().unwrap().unwrap(), Value::Int64(3));
assert_eq!(snap.g().V([1_i64]).out([]).count().next().unwrap().unwrap(), Value::Int64(0));
graph.close().unwrap();
}
#[test]
fn test_properties_roundtrip() {
use crate::{schema::DataType, Primitive};
let dir = tempdir().unwrap();
let db_path = dir.path().join("db");
let schema = BulkSchema {
vertex_labels: vec!["Person".into()],
edge_labels: vec!["Knows".into()],
prop_keys: vec![("age".into(), DataType::Int64), ("score".into(), DataType::Int64)],
};
let vertices = vec![
BulkVertex { id: 1, label: "Person".into(), props: [("age".into(), Primitive::Int64(42))].into() },
BulkVertex { id: 2, label: "Person".into(), props: HashMap::new() },
];
let edges = vec![BulkEdge {
src: 1,
dst: 2,
label: "Knows".into(),
props: [("score".into(), Primitive::Int64(100))].into(),
rank: None,
}];
SstBulkLoader::new(&db_path, dir.path().join("_w"))
.load_initial(
schema,
vertices.into_iter(),
edges.into_iter(),
GraphOptions::default(),
&RocksOptions::default(),
)
.unwrap();
let graph = Graph::open(&db_path).unwrap();
let mut snap = graph.read();
assert_eq!(snap.g().V([1_i64]).values(["age"]).next().unwrap().unwrap(), Value::Int64(42));
assert_eq!(snap.g().V([1_i64]).outE(["Knows"]).values(["score"]).next().unwrap().unwrap(), Value::Int64(100));
graph.close().unwrap();
}
#[test]
fn test_multi_mode_explicit_ranks() {
let dir = tempdir().unwrap();
let db_path = dir.path().join("db");
let vertices: Vec<BulkVertex> =
(1..=2).map(|i| BulkVertex { id: i, label: "Person".into(), props: HashMap::new() }).collect();
let edges = vec![
BulkEdge { src: 1, dst: 2, label: "Knows".into(), props: HashMap::new(), rank: Some(5) },
BulkEdge { src: 1, dst: 2, label: "Knows".into(), props: HashMap::new(), rank: Some(10) },
];
let stats = SstBulkLoader::new(&db_path, dir.path().join("_w"))
.load_initial(
small_schema(),
vertices.into_iter(),
edges.into_iter(),
GraphOptions { edge_mode: EdgeMode::Multi, ..Default::default() },
&RocksOptions::default(),
)
.unwrap();
assert_eq!(stats.edges_written, 2);
let graph = Graph::open(&db_path).unwrap();
let mut snap = graph.read();
assert_eq!(snap.g().V([1_i64]).outE(["Knows"]).count().next().unwrap().unwrap(), Value::Int64(2));
graph.close().unwrap();
}
#[test]
fn test_multi_mode_explicit_rank_duplicate() {
let dir = tempdir().unwrap();
let vertices: Vec<BulkVertex> =
(1..=2).map(|i| BulkVertex { id: i, label: "Person".into(), props: HashMap::new() }).collect();
let edges = vec![
BulkEdge { src: 1, dst: 2, label: "Knows".into(), props: HashMap::new(), rank: Some(3) },
BulkEdge { src: 1, dst: 2, label: "Knows".into(), props: HashMap::new(), rank: Some(3) },
];
let err = SstBulkLoader::new(dir.path().join("db"), dir.path().join("_w"))
.load_initial(
small_schema(),
vertices.into_iter(),
edges.into_iter(),
GraphOptions { edge_mode: EdgeMode::Multi, ..Default::default() },
&RocksOptions::default(),
)
.unwrap_err();
assert!(matches!(err, StoreError::DuplicateEdge(_)));
}
#[test]
fn test_multi_mode_rank_overflow() {
let dir = tempdir().unwrap();
let vertices: Vec<BulkVertex> =
(1..=2).map(|i| BulkVertex { id: i, label: "Person".into(), props: HashMap::new() }).collect();
let edges = vec![
BulkEdge { src: 1, dst: 2, label: "Knows".into(), props: HashMap::new(), rank: Some(65534) },
BulkEdge { src: 1, dst: 2, label: "Knows".into(), props: HashMap::new(), rank: None },
];
let err = SstBulkLoader::new(dir.path().join("db"), dir.path().join("_w"))
.load_initial(
small_schema(),
vertices.into_iter(),
edges.into_iter(),
GraphOptions { edge_mode: EdgeMode::Multi, ..Default::default() },
&RocksOptions::default(),
)
.unwrap_err();
assert!(matches!(err, StoreError::SchemaViolation(_)));
}
#[test]
fn test_sst_file_splitting() {
let dir = tempdir().unwrap();
let db_path = dir.path().join("db");
let stats = SstBulkLoader::new(&db_path, dir.path().join("_w"))
.with_max_sst_size(1)
.load_initial(
small_schema(),
small_vertices().into_iter(),
small_edges().into_iter(),
GraphOptions::default(),
&RocksOptions::default(),
)
.unwrap();
assert!(stats.sst_files >= 4, "expected at least one file per CF, got {}", stats.sst_files);
let graph = Graph::open(&db_path).unwrap();
let mut snap = graph.read();
assert_eq!(snap.g().V([]).count().next().unwrap().unwrap(), Value::Int64(5));
assert_eq!(snap.g().V([1_i64]).out(["Knows"]).count().next().unwrap().unwrap(), Value::Int64(2));
graph.close().unwrap();
}
#[test]
fn test_work_dir_cleaned_up_on_error() {
let dir = tempdir().unwrap();
let work_dir = dir.path().join("_w");
let vertices: Vec<BulkVertex> =
(1..=2).map(|i| BulkVertex { id: i, label: "Person".into(), props: HashMap::new() }).collect();
let edges = vec![BulkEdge { src: 1, dst: 2, label: "Unknown".into(), props: HashMap::new(), rank: None }];
let err = SstBulkLoader::new(dir.path().join("db"), work_dir.clone())
.load_initial(
small_schema(),
vertices.into_iter(),
edges.into_iter(),
GraphOptions::default(),
&RocksOptions::default(),
)
.unwrap_err();
assert!(matches!(err, StoreError::SchemaViolation(_)));
assert!(!work_dir.exists(), "WorkDirGuard should have removed work_dir on error");
}
#[test]
fn test_empty_sorted_label_file() {
let dir = tempdir().unwrap();
let path = dir.path().join("labels.bin");
let sorter = ExternalSorter::new(dir.path().join("sort"), 1024 * 1024);
let file = SortedLabelFile::write_from(sorter, &path).unwrap();
assert_eq!(file.count, 0);
let items: Vec<_> = file.reader().unwrap().collect();
assert!(items.is_empty());
}
#[test]
fn test_strict_mode_undeclared_edge_property() {
use crate::{schema::DataType, Primitive};
let dir = tempdir().unwrap();
let schema = BulkSchema {
vertex_labels: vec!["Person".into()],
edge_labels: vec!["Knows".into()],
prop_keys: vec![("age".into(), DataType::Int64)], };
let vertices: Vec<BulkVertex> =
(1..=2).map(|i| BulkVertex { id: i, label: "Person".into(), props: HashMap::new() }).collect();
let edges = vec![BulkEdge {
src: 1,
dst: 2,
label: "Knows".into(),
props: [("weight".into(), Primitive::Float64(1.5))].into(),
rank: None,
}];
let err = SstBulkLoader::new(dir.path().join("db"), dir.path().join("_w"))
.load_initial(
schema,
vertices.into_iter(),
edges.into_iter(),
GraphOptions { mode: SchemaMode::Strict, ..Default::default() },
&RocksOptions::default(),
)
.unwrap_err();
assert!(matches!(err, StoreError::SchemaViolation(_)));
}
}