rocksolid 2.7.0

An ergonomic, high-level RocksDB wrapper for Rust. Features CF-aware optimistic & pessimistic transactions, advanced routing for merge operators and compaction filters, performance tuning profiles, batching, TTL values, and DAO macros.
Documentation
//! Integration tests for `MultiCfBatchWriter` — a typed `WriteBatch` that commits writes
//! spanning multiple Column Families atomically, erasing the need to drop to `db_raw()` +
//! hand-rolled CF-handle juggling.

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") // no-op, still valid
    .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() })
  );
  // Untouched CF_B row survives.
  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(); // removes a1,a2,a3; keeps b1
  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()));
}

/// Append merge operator over MessagePack-serialized `MergeValue<String>` operands
/// (same shape as the one in `cf_integration.rs`).
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");

  // Two merges into CF_A and one into CF_B, all in one atomic multi-CF batch.
  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();
    // writer dropped here without commit() or discard()
  }
  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");
  // commit() consumes the writer, so misuse is prevented at compile time; here we assert the
  // internal guard by driving raw_batch_mut after operations are staged but before commit,
  // then a fresh writer works independently.
  let mut writer = store.batch_writer_multi_cf();
  writer.set_in(CF_A, "k", &"v".to_string()).unwrap();
  writer.commit().unwrap();

  // A second, independent writer still works.
  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() {
  // The RocksDbStore (default-CF convenience wrapper) also exposes batch_writer_multi_cf().
  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())
  );
}