use std::sync::Arc;
use redb::{
backends::InMemoryBackend, Database, ReadableDatabase, ReadableTable, ReadableTableMetadata,
TableDefinition,
};
use crate::errors::AtomicResult;
use super::{
kv_store::{KvIter, KvPair, KvStore},
trees::{Method, Operation, Tree},
};
const TABLE_RESOURCES: TableDefinition<&[u8], &[u8]> = TableDefinition::new("resources_v3");
const TABLE_PROP_VAL_SUB: TableDefinition<&[u8], &[u8]> =
TableDefinition::new("prop_val_sub_index");
const TABLE_VAL_PROP_SUB: TableDefinition<&[u8], &[u8]> =
TableDefinition::new("reference_index_v1");
const TABLE_QUERY_MEMBERS: TableDefinition<&[u8], &[u8]> = TableDefinition::new("members_index_v3");
const TABLE_WATCHED_QUERIES: TableDefinition<&[u8], &[u8]> =
TableDefinition::new("watched_queries_v3");
const TABLE_PLUGIN_META: TableDefinition<&[u8], &[u8]> = TableDefinition::new("plugin_meta");
const TABLE_DRIVE_MAPPING: TableDefinition<&[u8], &[u8]> = TableDefinition::new("drive_mapping");
const TABLE_DID_MAPPING: TableDefinition<&[u8], &[u8]> = TableDefinition::new("did_mapping");
const TABLE_LORO_SNAPSHOTS: TableDefinition<&[u8], &[u8]> = TableDefinition::new("loro_snapshots");
const TABLE_BLOBS: TableDefinition<&[u8], &[u8]> = TableDefinition::new("blobs");
fn table_def(tree: Tree) -> TableDefinition<'static, &'static [u8], &'static [u8]> {
match tree {
Tree::Resources => TABLE_RESOURCES,
Tree::PropValSub => TABLE_PROP_VAL_SUB,
Tree::ValPropSub => TABLE_VAL_PROP_SUB,
Tree::QueryMembers => TABLE_QUERY_MEMBERS,
Tree::WatchedQueries => TABLE_WATCHED_QUERIES,
Tree::PluginMeta => TABLE_PLUGIN_META,
Tree::DriveMapping => TABLE_DRIVE_MAPPING,
Tree::DidMapping => TABLE_DID_MAPPING,
Tree::LoroSnapshots => TABLE_LORO_SNAPSHOTS,
Tree::Blobs => TABLE_BLOBS,
}
}
pub struct RedbStore {
db: Arc<Database>,
batch_buffer: std::sync::Mutex<Option<BatchBuffer>>,
}
#[derive(Default)]
struct BatchBuffer {
ops: Vec<Operation>,
latest: std::collections::HashMap<(String, Vec<u8>), Option<Vec<u8>>>,
}
impl BatchBuffer {
fn push(&mut self, op: Operation) {
let key = (op.tree.to_string(), op.key.clone());
let val = match op.method {
Method::Insert => op.val.clone(),
Method::Delete => None,
};
self.latest.insert(key, val);
self.ops.push(op);
}
fn get(&self, tree: &Tree, key: &[u8]) -> Option<Option<Vec<u8>>> {
self.latest.get(&(tree.to_string(), key.to_vec())).cloned()
}
}
#[cfg(all(feature = "db-redb", not(target_arch = "wasm32")))]
pub fn compact_file(path: &std::path::Path) -> AtomicResult<(u64, u64, bool)> {
let size_before = std::fs::metadata(path).map(|m| m.len()).unwrap_or(0);
let mut db = redb::Database::create(path)
.map_err(|e| format!("Failed to open redb at {}: {e}", path.display()))?;
let did_compact = db
.compact()
.map_err(|e| format!("Compaction failed: {e}"))?;
drop(db);
let size_after = std::fs::metadata(path).map(|m| m.len()).unwrap_or(0);
Ok((size_before, size_after, did_compact))
}
impl RedbStore {
#[cfg(not(target_arch = "wasm32"))]
pub fn new_file(path: &std::path::Path) -> AtomicResult<Self> {
let t = std::time::Instant::now();
let db = Database::create(path)
.map_err(|e| format!("Failed to create redb at {}: {e}", path.display()))?;
tracing::info!("RedbStore::new_file: Database::create in {:?}", t.elapsed());
let t = std::time::Instant::now();
{
let mut tx = db
.begin_write()
.map_err(|e| format!("Failed to begin write tx: {e}"))?;
tx.set_quick_repair(true);
let _ = tx.open_table(TABLE_RESOURCES);
let _ = tx.open_table(TABLE_PROP_VAL_SUB);
let _ = tx.open_table(TABLE_VAL_PROP_SUB);
let _ = tx.open_table(TABLE_QUERY_MEMBERS);
let _ = tx.open_table(TABLE_WATCHED_QUERIES);
let _ = tx.open_table(TABLE_PLUGIN_META);
let _ = tx.open_table(TABLE_DRIVE_MAPPING);
let _ = tx.open_table(TABLE_DID_MAPPING);
let _ = tx.open_table(TABLE_LORO_SNAPSHOTS);
let _ = tx.open_table(TABLE_BLOBS);
tx.commit()
.map_err(|e| format!("Failed to commit initial tables: {e}"))?;
}
tracing::info!("RedbStore::new_file: table-create tx in {:?}", t.elapsed());
Ok(RedbStore {
db: Arc::new(db),
batch_buffer: std::sync::Mutex::new(None),
})
}
pub fn new_memory() -> AtomicResult<Self> {
let backend = InMemoryBackend::new();
let db = Database::builder()
.create_with_backend(backend)
.map_err(|e| format!("Failed to create redb: {e}"))?;
{
let mut tx = db
.begin_write()
.map_err(|e| format!("Failed to begin write tx: {e}"))?;
tx.set_quick_repair(true);
let _ = tx.open_table(TABLE_RESOURCES);
let _ = tx.open_table(TABLE_PROP_VAL_SUB);
let _ = tx.open_table(TABLE_VAL_PROP_SUB);
let _ = tx.open_table(TABLE_QUERY_MEMBERS);
let _ = tx.open_table(TABLE_WATCHED_QUERIES);
let _ = tx.open_table(TABLE_PLUGIN_META);
let _ = tx.open_table(TABLE_DRIVE_MAPPING);
let _ = tx.open_table(TABLE_DID_MAPPING);
let _ = tx.open_table(TABLE_LORO_SNAPSHOTS);
let _ = tx.open_table(TABLE_BLOBS);
tx.commit()
.map_err(|e| format!("Failed to commit initial tables: {e}"))?;
}
Ok(RedbStore {
db: Arc::new(db),
batch_buffer: std::sync::Mutex::new(None),
})
}
#[cfg(target_arch = "wasm32")]
pub async fn new_opfs(filename: &str, encryption_key: Option<&[u8; 32]>) -> AtomicResult<Self> {
let backend = super::opfs_backend::OpfsBackend::open(filename)
.await
.map_err(|e| format!("Failed to open OPFS backend: {:?}", e))?;
let db = match encryption_key {
Some(key) => {
let encrypted = super::encrypted_backend::EncryptedBackend::new(backend, key)
.map_err(|e| format!("Failed to open encrypted OPFS backend: {e}"))?;
Database::builder()
.create_with_backend(encrypted)
.map_err(|e| format!("Failed to create encrypted redb with OPFS: {e}"))?
}
None => Database::builder()
.create_with_backend(backend)
.map_err(|e| format!("Failed to create redb with OPFS: {e}"))?,
};
{
let mut tx = db
.begin_write()
.map_err(|e| format!("Failed to begin write tx: {e}"))?;
tx.set_quick_repair(true);
let _ = tx.open_table(TABLE_RESOURCES);
let _ = tx.open_table(TABLE_PROP_VAL_SUB);
let _ = tx.open_table(TABLE_VAL_PROP_SUB);
let _ = tx.open_table(TABLE_QUERY_MEMBERS);
let _ = tx.open_table(TABLE_WATCHED_QUERIES);
let _ = tx.open_table(TABLE_PLUGIN_META);
let _ = tx.open_table(TABLE_DRIVE_MAPPING);
let _ = tx.open_table(TABLE_DID_MAPPING);
let _ = tx.open_table(TABLE_LORO_SNAPSHOTS);
let _ = tx.open_table(TABLE_BLOBS);
tx.commit()
.map_err(|e| format!("Failed to commit initial tables: {e}"))?;
}
Ok(RedbStore {
db: Arc::new(db),
batch_buffer: std::sync::Mutex::new(None),
})
}
}
fn prefix_upper_bound(prefix: &[u8]) -> Option<Vec<u8>> {
let mut end = prefix.to_vec();
while let Some(last) = end.last_mut() {
if *last < 0xff {
*last += 1;
return Some(end);
}
end.pop();
}
None
}
impl KvStore for RedbStore {
fn get(&self, tree: Tree, key: &[u8]) -> AtomicResult<Option<Vec<u8>>> {
{
let buf = self.batch_buffer.lock().unwrap();
if let Some(buffer) = buf.as_ref() {
if let Some(val) = buffer.get(&tree, key) {
return Ok(val);
}
}
}
let tx = self
.db
.begin_read()
.map_err(|e| format!("redb read tx: {e}"))?;
let table = tx
.open_table(table_def(tree))
.map_err(|e| format!("redb open table: {e}"))?;
let result = table.get(key).map_err(|e| format!("redb get: {e}"))?;
Ok(result.map(|guard| guard.value().to_vec()))
}
fn insert(&self, tree: Tree, key: &[u8], val: &[u8]) -> AtomicResult<()> {
self.apply_batch(&[Operation {
tree,
method: Method::Insert,
key: key.to_vec(),
val: Some(val.to_vec()),
}])
}
fn remove(&self, tree: Tree, key: &[u8]) -> AtomicResult<()> {
self.apply_batch(&[Operation {
tree,
method: Method::Delete,
key: key.to_vec(),
val: None,
}])
}
fn contains_key(&self, tree: Tree, key: &[u8]) -> AtomicResult<bool> {
let tx = self
.db
.begin_read()
.map_err(|e| format!("redb read tx: {e}"))?;
let table = tx
.open_table(table_def(tree))
.map_err(|e| format!("redb open table: {e}"))?;
let result = table
.get(key)
.map_err(|e| format!("redb contains_key: {e}"))?;
Ok(result.is_some())
}
fn scan_prefix(&self, tree: Tree, prefix: &[u8]) -> KvIter {
let tx = match self.db.begin_read() {
Ok(tx) => tx,
Err(e) => return Box::new(std::iter::once(Err(format!("redb read tx: {e}").into()))),
};
let table = match tx.open_table(table_def(tree)) {
Ok(t) => t,
Err(e) => {
return Box::new(std::iter::once(Err(format!("redb open table: {e}").into())))
}
};
let results: Vec<KvPair> = if let Some(end) = prefix_upper_bound(prefix) {
table
.range(prefix..end.as_slice())
.map(|iter| {
iter.filter_map(|r| r.ok())
.map(|(k, v)| (k.value().to_vec(), v.value().to_vec()))
.collect()
})
.unwrap_or_default()
} else {
table
.range(prefix..)
.map(|iter| {
iter.filter_map(|r| r.ok())
.map(|(k, v)| (k.value().to_vec(), v.value().to_vec()))
.collect()
})
.unwrap_or_default()
};
Box::new(results.into_iter().map(Ok))
}
fn range(&self, tree: Tree, start: Vec<u8>, end: Vec<u8>, reverse: bool) -> KvIter {
let tx = match self.db.begin_read() {
Ok(tx) => tx,
Err(e) => return Box::new(std::iter::once(Err(format!("redb read tx: {e}").into()))),
};
let table = match tx.open_table(table_def(tree)) {
Ok(t) => t,
Err(e) => {
return Box::new(std::iter::once(Err(format!("redb open table: {e}").into())))
}
};
let results: Vec<KvPair> = table
.range(start.as_slice()..end.as_slice())
.map(|iter| {
iter.filter_map(|r| r.ok())
.map(|(k, v)| (k.value().to_vec(), v.value().to_vec()))
.collect()
})
.unwrap_or_default();
if reverse {
let mut reversed = results;
reversed.reverse();
Box::new(reversed.into_iter().map(Ok))
} else {
Box::new(results.into_iter().map(Ok))
}
}
fn iter_tree(&self, tree: Tree) -> KvIter {
let tx = match self.db.begin_read() {
Ok(tx) => tx,
Err(e) => return Box::new(std::iter::once(Err(format!("redb read tx: {e}").into()))),
};
let table = match tx.open_table(table_def(tree)) {
Ok(t) => t,
Err(e) => {
return Box::new(std::iter::once(Err(format!("redb open table: {e}").into())))
}
};
let results: Vec<KvPair> = table
.iter()
.map(|iter| {
iter.filter_map(|r| r.ok())
.map(|(k, v)| (k.value().to_vec(), v.value().to_vec()))
.collect()
})
.unwrap_or_default();
Box::new(results.into_iter().map(Ok))
}
fn clear_tree(&self, tree: Tree) -> AtomicResult<()> {
let mut tx = self
.db
.begin_write()
.map_err(|e| format!("redb write tx: {e}"))?;
tx.set_quick_repair(true);
{
let mut table = tx
.open_table(table_def(tree))
.map_err(|e| format!("redb open table: {e}"))?;
let keys: Vec<Vec<u8>> = table
.iter()
.map(|iter| {
iter.filter_map(|r| r.ok())
.map(|(k, _)| k.value().to_vec())
.collect()
})
.unwrap_or_default();
for key in keys {
table
.remove(key.as_slice())
.map_err(|e| format!("redb remove in clear: {e}"))?;
}
}
tx.commit().map_err(|e| format!("redb commit clear: {e}"))?;
Ok(())
}
fn apply_batch(&self, operations: &[Operation]) -> AtomicResult<()> {
if operations.is_empty() {
return Ok(());
}
{
let mut buf = self.batch_buffer.lock().unwrap();
if let Some(buffer) = buf.as_mut() {
for op in operations {
buffer.push(op.clone());
}
return Ok(());
}
}
let mut tx = self
.db
.begin_write()
.map_err(|e| format!("redb write tx: {e}"))?;
tx.set_durability(redb::Durability::None)
.map_err(|e| format!("redb set_durability: {e}"))?;
{
for op in operations {
let mut table = tx
.open_table(table_def(op.tree.clone()))
.map_err(|e| format!("redb open table: {e}"))?;
match op.method {
Method::Insert => {
let val = op.val.as_deref().unwrap_or(b"");
table
.insert(op.key.as_slice(), val)
.map_err(|e| format!("redb batch insert: {e}"))?;
}
Method::Delete => {
table
.remove(op.key.as_slice())
.map_err(|e| format!("redb batch remove: {e}"))?;
}
}
}
}
tx.commit().map_err(|e| format!("redb commit batch: {e}"))?;
Ok(())
}
fn flush(&self) -> AtomicResult<()> {
let mut tx = self
.db
.begin_write()
.map_err(|e| format!("redb flush begin_write: {e}"))?;
tx.set_quick_repair(true);
{
let mut table = tx
.open_table(TABLE_DRIVE_MAPPING)
.map_err(|e| format!("redb flush open table: {e}"))?;
table
.insert(b"__flush_sentinel__".as_slice(), b"".as_slice())
.map_err(|e| format!("redb flush sentinel: {e}"))?;
}
tx.commit().map_err(|e| format!("redb flush commit: {e}"))?;
Ok(())
}
fn len(&self, tree: Tree) -> AtomicResult<usize> {
let tx = self
.db
.begin_read()
.map_err(|e| format!("redb read tx: {e}"))?;
let table = tx
.open_table(table_def(tree))
.map_err(|e| format!("redb open table: {e}"))?;
Ok(table.len().map_err(|e| format!("redb len: {e}"))? as usize)
}
fn begin_batch(&self) {
let mut buf = self.batch_buffer.lock().unwrap();
if buf.is_none() {
*buf = Some(BatchBuffer::default());
}
}
fn commit_batch(&self) -> AtomicResult<()> {
let ops = {
let mut buf = self.batch_buffer.lock().unwrap();
buf.take().map(|b| b.ops).unwrap_or_default()
};
if ops.is_empty() {
return Ok(());
}
let mut tx = self
.db
.begin_write()
.map_err(|e| format!("redb write tx: {e}"))?;
tx.set_durability(redb::Durability::None)
.map_err(|e| format!("redb set_durability: {e}"))?;
{
for op in &ops {
let mut table = tx
.open_table(table_def(op.tree.clone()))
.map_err(|e| format!("redb open table: {e}"))?;
match op.method {
Method::Insert => {
let val = op.val.as_deref().unwrap_or(b"");
table
.insert(op.key.as_slice(), val)
.map_err(|e| format!("redb batch insert: {e}"))?;
}
Method::Delete => {
table
.remove(op.key.as_slice())
.map_err(|e| format!("redb batch remove: {e}"))?;
}
}
}
}
tx.commit().map_err(|e| format!("redb commit batch: {e}"))?;
Ok(())
}
}