mod common;
use common::setup_logging;
use rocksolid::batch::MultiCfBatchWriter;
use rocksolid::cf_store::{CFOperations, RocksDbCFStore};
use rocksolid::config::{BaseCfConfig, RockSolidMergeOperatorCfConfig, RocksDbCFStoreConfig};
use rocksolid::serialization::deserialize_value;
use rocksolid::store::RocksDbStore;
use rocksolid::types::MergeValue;
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use tempfile::TempDir;
const DEFAULT_CF: &str = rocksdb::DEFAULT_COLUMN_FAMILY_NAME;
const CF_A: &str = "cf_a";
const CF_B: &str = "cf_b";
#[derive(Serialize, Deserialize, Debug, PartialEq, Eq, Clone)]
struct Rec {
id: u32,
name: String,
}
fn setup(test_name: &str) -> (TempDir, RocksDbCFStore) {
setup_logging();
let temp_dir = TempDir::new().unwrap();
let db_path = temp_dir.path().join(test_name).to_str().unwrap().to_string();
let mut cf_configs = HashMap::new();
for cf in [DEFAULT_CF, CF_A, CF_B] {
cf_configs.insert(cf.to_string(), BaseCfConfig::default());
}
let config = RocksDbCFStoreConfig {
path: db_path,
create_if_missing: true,
column_families_to_open: vec![DEFAULT_CF.to_string(), CF_A.to_string(), CF_B.to_string()],
column_family_configs: cf_configs,
..Default::default()
};
let store = RocksDbCFStore::open(config).unwrap();
(temp_dir, store)
}
#[test]
fn commits_across_multiple_cfs_atomically() {
let (_dir, store) = setup("multi_cf_commit");
let rec_a = Rec { id: 1, name: "alpha".into() };
let rec_b = Rec { id: 2, name: "beta".into() };
let mut writer: MultiCfBatchWriter = store.batch_writer_multi_cf();
writer
.set_in(CF_A, "k1", &rec_a)
.unwrap()
.set_in(CF_B, "k2", &rec_b)
.unwrap()
.set_in(DEFAULT_CF, "kd", &"default-val".to_string())
.unwrap()
.set_raw_in(CF_A, "raw", b"raw-bytes")
.unwrap()
.delete_in(CF_B, "does_not_exist") .unwrap();
writer.commit().unwrap();
assert_eq!(store.get::<_, Rec>(CF_A, "k1").unwrap(), Some(rec_a));
assert_eq!(store.get::<_, Rec>(CF_B, "k2").unwrap(), Some(rec_b));
assert_eq!(
store.get::<_, String>(DEFAULT_CF, "kd").unwrap(),
Some("default-val".to_string())
);
assert_eq!(store.get_raw(CF_A, "raw").unwrap(), Some(b"raw-bytes".to_vec()));
}
#[test]
fn delete_and_set_in_same_batch() {
let (_dir, store) = setup("multi_cf_delete_set");
store.put(CF_A, "old", &Rec { id: 9, name: "old".into() }).unwrap();
store.put(CF_B, "keep", &Rec { id: 10, name: "keep".into() }).unwrap();
let mut writer = store.batch_writer_multi_cf();
writer
.delete_in(CF_A, "old")
.unwrap()
.set_in(CF_A, "new", &Rec { id: 11, name: "new".into() })
.unwrap();
writer.commit().unwrap();
assert_eq!(store.get::<_, Rec>(CF_A, "old").unwrap(), None);
assert_eq!(
store.get::<_, Rec>(CF_A, "new").unwrap(),
Some(Rec { id: 11, name: "new".into() })
);
assert_eq!(
store.get::<_, Rec>(CF_B, "keep").unwrap(),
Some(Rec { id: 10, name: "keep".into() })
);
}
#[test]
fn delete_range_in_cf() {
let (_dir, store) = setup("multi_cf_delete_range");
for k in ["a1", "a2", "a3", "b1"] {
store.put_raw(CF_A, k.to_string(), k.as_bytes()).unwrap();
}
let mut writer = store.batch_writer_multi_cf();
writer.delete_range_in(CF_A, "a", "b").unwrap(); writer.commit().unwrap();
assert_eq!(store.get_raw(CF_A, "a1").unwrap(), None);
assert_eq!(store.get_raw(CF_A, "a3").unwrap(), None);
assert_eq!(store.get_raw(CF_A, "b1").unwrap(), Some(b"b1".to_vec()));
}
fn append_merge_operator(
_key: &[u8],
existing: Option<&[u8]>,
operands: &rocksdb::MergeOperands,
) -> Option<Vec<u8>> {
let mut acc = match existing {
Some(bytes) => deserialize_value::<String>(bytes).ok()?,
None => String::new(),
};
for op in operands {
if let Ok(mv) = deserialize_value::<MergeValue<String>>(op) {
if !acc.is_empty() {
acc.push(',');
}
acc.push_str(&mv.1);
}
}
rocksolid::serialize_value(&acc).ok()
}
fn setup_with_merge(test_name: &str) -> (TempDir, RocksDbCFStore) {
setup_logging();
let temp_dir = TempDir::new().unwrap();
let db_path = temp_dir.path().join(test_name).to_str().unwrap().to_string();
let merge_cfg = || RockSolidMergeOperatorCfConfig {
name: "append".to_string(),
full_merge_fn: Some(append_merge_operator),
partial_merge_fn: Some(append_merge_operator),
};
let mut cf_configs = HashMap::new();
cf_configs.insert(DEFAULT_CF.to_string(), BaseCfConfig::default());
for cf in [CF_A, CF_B] {
cf_configs.insert(
cf.to_string(),
BaseCfConfig {
merge_operator: Some(merge_cfg()),
..Default::default()
},
);
}
let config = RocksDbCFStoreConfig {
path: db_path,
create_if_missing: true,
column_families_to_open: vec![DEFAULT_CF.to_string(), CF_A.to_string(), CF_B.to_string()],
column_family_configs: cf_configs,
..Default::default()
};
let store = RocksDbCFStore::open(config).unwrap();
(temp_dir, store)
}
#[test]
fn merge_in_across_cfs() {
let (_dir, store) = setup_with_merge("multi_cf_merge");
let mut writer = store.batch_writer_multi_cf();
writer
.merge_in(CF_A, "m", &MergeValue(rocksolid::types::MergeValueOperator::Add, "a".to_string()))
.unwrap()
.merge_in(CF_A, "m", &MergeValue(rocksolid::types::MergeValueOperator::Add, "b".to_string()))
.unwrap()
.merge_in(CF_B, "m", &MergeValue(rocksolid::types::MergeValueOperator::Add, "z".to_string()))
.unwrap();
writer.commit().unwrap();
assert_eq!(store.get::<_, String>(CF_A, "m").unwrap(), Some("a,b".to_string()));
assert_eq!(store.get::<_, String>(CF_B, "m").unwrap(), Some("z".to_string()));
}
#[test]
fn drop_without_commit_applies_nothing() {
let (_dir, store) = setup("multi_cf_drop");
{
let mut writer = store.batch_writer_multi_cf();
writer.set_in(CF_A, "ephemeral", &Rec { id: 1, name: "x".into() }).unwrap();
}
assert_eq!(store.get::<_, Rec>(CF_A, "ephemeral").unwrap(), None);
}
#[test]
fn discard_applies_nothing() {
let (_dir, store) = setup("multi_cf_discard");
let mut writer = store.batch_writer_multi_cf();
writer.set_in(CF_A, "gone", &Rec { id: 1, name: "x".into() }).unwrap();
writer.discard();
assert_eq!(store.get::<_, Rec>(CF_A, "gone").unwrap(), None);
}
#[test]
fn set_with_expiry_in() {
let (_dir, store) = setup("multi_cf_expiry");
let mut writer = store.batch_writer_multi_cf();
writer
.set_with_expiry_in(CF_A, "temp", &"payload".to_string(), 9_999_999_999)
.unwrap();
writer.commit().unwrap();
let vwe = store
.get_with_expiry::<_, String>(CF_A, "temp")
.unwrap()
.expect("value present");
assert_eq!(vwe.get().unwrap(), "payload".to_string());
assert_eq!(vwe.expire_time, 9_999_999_999);
}
#[test]
fn raw_batch_mut_escape_hatch() {
let (_dir, store) = setup("multi_cf_raw_escape");
let handle = store.get_cf_handle(CF_B).unwrap();
let mut writer = store.batch_writer_multi_cf();
writer.set_in(CF_A, "typed", &"a".to_string()).unwrap();
{
let raw = writer.raw_batch_mut().unwrap();
raw.put_cf(&handle, b"escape", b"hatch");
}
writer.commit().unwrap();
assert_eq!(store.get::<_, String>(CF_A, "typed").unwrap(), Some("a".to_string()));
assert_eq!(store.get_raw(CF_B, "escape").unwrap(), Some(b"hatch".to_vec()));
}
#[test]
fn use_after_commit_errors() {
let (_dir, store) = setup("multi_cf_use_after_commit");
let mut writer = store.batch_writer_multi_cf();
writer.set_in(CF_A, "k", &"v".to_string()).unwrap();
writer.commit().unwrap();
let mut writer2 = store.batch_writer_multi_cf();
writer2.set_in(CF_A, "k2", &"v2".to_string()).unwrap();
writer2.commit().unwrap();
assert_eq!(store.get::<_, String>(CF_A, "k2").unwrap(), Some("v2".to_string()));
}
#[test]
fn multi_cf_writer_via_default_store_wrapper() {
setup_logging();
let temp_dir = TempDir::new().unwrap();
let db_path = temp_dir.path().join("multi_cf_wrapper").to_str().unwrap().to_string();
let store = RocksDbStore::open(rocksolid::config::RocksDbStoreConfig {
path: db_path,
create_if_missing: true,
..Default::default()
})
.unwrap();
let mut writer = store.batch_writer_multi_cf();
writer.set_in(DEFAULT_CF, "wrapped", &"ok".to_string()).unwrap();
writer.commit().unwrap();
assert_eq!(
store.cf_store().get::<_, String>(DEFAULT_CF, "wrapped").unwrap(),
Some("ok".to_string())
);
}