use tracing::debug;
use super::btree::{DOCUMENTS, INDEXES, SparseEngine, coll_prefix, redb_err, tenant_prefix};
impl SparseEngine {
pub fn scan_documents(
&self,
database_id: u64,
tenant_id: u64,
collection: &str,
limit: usize,
) -> crate::Result<Vec<(String, Vec<u8>)>> {
let prefix = coll_prefix(database_id, tenant_id, collection);
let end = format!("{prefix}\u{ffff}");
let read_txn = self.db.begin_read().map_err(|e| redb_err("read txn", e))?;
let table = read_txn
.open_table(DOCUMENTS)
.map_err(|e| redb_err("open table", e))?;
let range = table
.range(prefix.as_str()..end.as_str())
.map_err(|e| redb_err("doc range", e))?;
let mut results = Vec::with_capacity(limit.min(256));
for entry in range {
if results.len() >= limit {
break;
}
let entry = entry.map_err(|e| redb_err("doc entry", e))?;
let key = entry.0.value().to_string();
let doc_id = key.strip_prefix(&prefix).unwrap_or(&key).to_string();
let value = entry.1.value().to_vec();
results.push((doc_id, value));
}
debug!(collection, count = results.len(), "document scan");
Ok(results)
}
pub fn scan_documents_for_each<F>(
&self,
database_id: u64,
tenant_id: u64,
collection: &str,
limit: usize,
mut f: F,
) -> crate::Result<()>
where
F: FnMut(&str, &[u8]) -> crate::Result<()>,
{
let prefix = coll_prefix(database_id, tenant_id, collection);
let end = format!("{prefix}\u{ffff}");
let read_txn = self.db.begin_read().map_err(|e| redb_err("read txn", e))?;
let table = read_txn
.open_table(DOCUMENTS)
.map_err(|e| redb_err("open table", e))?;
let range = table
.range(prefix.as_str()..end.as_str())
.map_err(|e| redb_err("doc range", e))?;
let mut count = 0usize;
for entry in range {
if count >= limit {
break;
}
let entry = entry.map_err(|e| redb_err("doc entry", e))?;
let key = entry.0.value().to_string();
let doc_id = key.strip_prefix(&prefix).unwrap_or(&key);
let value = entry.1.value();
f(doc_id, value)?;
count += 1;
}
debug!(collection, count, "streaming document scan");
Ok(())
}
pub fn scan_documents_chunked<F>(
&self,
database_id: u64,
tenant_id: u64,
collection: &str,
total_limit: usize,
chunk_size: usize,
mut handler: F,
) -> crate::Result<usize>
where
F: FnMut(&[(String, Vec<u8>)]),
{
let prefix = coll_prefix(database_id, tenant_id, collection);
let end = format!("{prefix}\u{ffff}");
let read_txn = self.db.begin_read().map_err(|e| redb_err("read txn", e))?;
let table = read_txn
.open_table(DOCUMENTS)
.map_err(|e| redb_err("open table", e))?;
let range = table
.range(prefix.as_str()..end.as_str())
.map_err(|e| redb_err("doc range", e))?;
let mut chunk = Vec::with_capacity(chunk_size);
let mut total = 0usize;
for entry in range {
if total >= total_limit {
break;
}
let entry = entry.map_err(|e| redb_err("doc entry", e))?;
let key = entry.0.value().to_string();
let doc_id = key.strip_prefix(&prefix).unwrap_or(&key).to_string();
let value = entry.1.value().to_vec();
chunk.push((doc_id, value));
total += 1;
if chunk.len() >= chunk_size {
handler(&chunk);
chunk.clear();
}
}
if !chunk.is_empty() {
handler(&chunk);
}
debug!(collection, total, chunk_size, "chunked document scan");
Ok(total)
}
pub fn scan_index_groups(
&self,
database_id: u64,
tenant_id: u64,
collection: &str,
field: &str,
) -> crate::Result<Vec<(String, usize)>> {
let prefix = format!(
"{}{field}:",
coll_prefix(database_id, tenant_id, collection)
);
let end = format!("{prefix}\u{ffff}");
let read_txn = self.db.begin_read().map_err(|e| redb_err("read txn", e))?;
let table = read_txn
.open_table(INDEXES)
.map_err(|e| redb_err("open table", e))?;
let range = table
.range(prefix.as_str()..end.as_str())
.map_err(|e| redb_err("index range", e))?;
let mut groups: std::collections::HashMap<String, usize> = std::collections::HashMap::new();
for entry in range {
let entry = entry.map_err(|e| redb_err("index entry", e))?;
let key = entry.0.value().to_string();
if let Some(rest) = key.strip_prefix(&prefix)
&& let Some(colon_pos) = rest.rfind(':')
{
let value = &rest[..colon_pos];
*groups.entry(value.to_string()).or_default() += 1;
}
}
let mut result: Vec<(String, usize)> = groups.into_iter().collect();
result.sort_by(|a, b| a.0.cmp(&b.0));
debug!(collection, field, groups = result.len(), "index group scan");
Ok(result)
}
pub fn scan_index_groups_filtered(
&self,
database_id: u64,
tenant_id: u64,
collection: &str,
field: &str,
doc_ids: &std::collections::HashSet<String>,
) -> crate::Result<Vec<(String, usize)>> {
let prefix = format!(
"{}{field}:",
coll_prefix(database_id, tenant_id, collection)
);
let end = format!("{prefix}\u{ffff}");
let read_txn = self.db.begin_read().map_err(|e| redb_err("read txn", e))?;
let table = read_txn
.open_table(INDEXES)
.map_err(|e| redb_err("open table", e))?;
let range = table
.range(prefix.as_str()..end.as_str())
.map_err(|e| redb_err("index range", e))?;
let mut groups: std::collections::HashMap<String, usize> = std::collections::HashMap::new();
for entry in range {
let entry = entry.map_err(|e| redb_err("index entry", e))?;
let key = entry.0.value().to_string();
if let Some(rest) = key.strip_prefix(&prefix)
&& let Some(colon_pos) = rest.rfind(':')
{
let value = &rest[..colon_pos];
let doc_id = &rest[colon_pos + 1..];
if doc_ids.contains(doc_id) {
*groups.entry(value.to_string()).or_default() += 1;
}
}
}
let mut result: Vec<(String, usize)> = groups.into_iter().collect();
result.sort_by_key(|r| std::cmp::Reverse(r.1)); debug!(
collection,
field,
groups = result.len(),
"filtered index group scan"
);
Ok(result)
}
pub fn scan_documents_filtered(
&self,
database_id: u64,
tenant_id: u64,
collection: &str,
limit: usize,
predicate: &dyn Fn(&[u8]) -> bool,
) -> crate::Result<Vec<(String, Vec<u8>)>> {
let prefix = coll_prefix(database_id, tenant_id, collection);
let end = format!("{prefix}\u{ffff}");
let read_txn = self.db.begin_read().map_err(|e| redb_err("read txn", e))?;
let table = read_txn
.open_table(DOCUMENTS)
.map_err(|e| redb_err("open table", e))?;
let range = table
.range(prefix.as_str()..end.as_str())
.map_err(|e| redb_err("doc range", e))?;
let mut results = Vec::with_capacity(limit.min(256));
for entry in range {
if results.len() >= limit {
break;
}
let entry = entry.map_err(|e| redb_err("doc entry", e))?;
let value_bytes = entry.1.value();
if !predicate(value_bytes) {
continue;
}
let key = entry.0.value().to_string();
let doc_id = key.strip_prefix(&prefix).unwrap_or(&key).to_string();
results.push((doc_id, value_bytes.to_vec()));
}
debug!(collection, count = results.len(), "filtered document scan");
Ok(results)
}
pub fn export_documents(&self) -> crate::Result<Vec<(String, Vec<u8>)>> {
let txn = self.db.begin_read().map_err(|e| redb_err("read txn", e))?;
let table = txn
.open_table(DOCUMENTS)
.map_err(|e| redb_err("open docs", e))?;
let mut pairs = Vec::new();
let iter = table
.range::<&str>(..)
.map_err(|e| redb_err("iter docs", e))?;
for entry in iter {
let entry = entry.map_err(|e| redb_err("read doc entry", e))?;
pairs.push((entry.0.value().to_string(), entry.1.value().to_vec()));
}
Ok(pairs)
}
pub fn export_indexes(&self) -> crate::Result<Vec<(String, Vec<u8>)>> {
let txn = self.db.begin_read().map_err(|e| redb_err("read txn", e))?;
let table = txn
.open_table(INDEXES)
.map_err(|e| redb_err("open indexes", e))?;
let mut pairs = Vec::new();
let iter = table
.range::<&str>(..)
.map_err(|e| redb_err("iter indexes", e))?;
for entry in iter {
let entry = entry.map_err(|e| redb_err("read index entry", e))?;
pairs.push((entry.0.value().to_string(), entry.1.value().to_vec()));
}
Ok(pairs)
}
pub fn import_documents(&self, pairs: &[(String, Vec<u8>)]) -> crate::Result<()> {
let txn = self
.db
.begin_write()
.map_err(|e| redb_err("write txn", e))?;
{
let mut table = txn
.open_table(DOCUMENTS)
.map_err(|e| redb_err("open docs", e))?;
for (key, value) in pairs {
table
.insert(key.as_str(), value.as_slice())
.map_err(|e| redb_err("insert doc", e))?;
}
}
txn.commit().map_err(|e| redb_err("commit", e))?;
Ok(())
}
pub fn import_indexes(&self, pairs: &[(String, Vec<u8>)]) -> crate::Result<()> {
let txn = self
.db
.begin_write()
.map_err(|e| redb_err("write txn", e))?;
{
let mut table = txn
.open_table(INDEXES)
.map_err(|e| redb_err("open indexes", e))?;
for (key, value) in pairs {
table
.insert(key.as_str(), value.as_slice())
.map_err(|e| redb_err("insert idx", e))?;
}
}
txn.commit().map_err(|e| redb_err("commit", e))?;
Ok(())
}
pub fn scan_all_for_tenant(
&self,
database_id: u64,
tenant_id: u64,
) -> crate::Result<Vec<(String, Vec<u8>)>> {
self.scan_table_for_tenant(DOCUMENTS, database_id, tenant_id, "tenant scan")
}
pub fn scan_indexes_for_tenant(
&self,
database_id: u64,
tenant_id: u64,
) -> crate::Result<Vec<(String, Vec<u8>)>> {
self.scan_table_for_tenant(INDEXES, database_id, tenant_id, "index scan")
}
fn scan_table_for_tenant(
&self,
table_def: redb::TableDefinition<&str, &[u8]>,
database_id: u64,
tenant_id: u64,
label: &str,
) -> crate::Result<Vec<(String, Vec<u8>)>> {
let prefix = tenant_prefix(database_id, tenant_id);
let end = format!("{prefix}\u{ffff}");
let read_txn = self.db.begin_read().map_err(|e| redb_err("read txn", e))?;
let table = read_txn
.open_table(table_def)
.map_err(|e| redb_err("open table", e))?;
let range = table
.range(prefix.as_str()..end.as_str())
.map_err(|e| redb_err(label, e))?;
let mut results = Vec::new();
for entry in range {
let entry = entry.map_err(|e| redb_err("entry", e))?;
results.push((entry.0.value().to_string(), entry.1.value().to_vec()));
}
Ok(results)
}
pub fn put_index_raw(&self, key: &str, value: &[u8]) -> crate::Result<()> {
let write_txn = self
.db
.begin_write()
.map_err(|e| redb_err("write txn", e))?;
{
let mut table = write_txn
.open_table(INDEXES)
.map_err(|e| redb_err("open index table", e))?;
table
.insert(key, value)
.map_err(|e| redb_err("index insert", e))?;
}
write_txn.commit().map_err(|e| redb_err("commit", e))?;
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
fn open_temp() -> (SparseEngine, tempfile::TempDir) {
let dir = tempfile::tempdir().unwrap();
let engine = SparseEngine::open(&dir.path().join("sparse.redb")).unwrap();
(engine, dir)
}
#[test]
fn for_each_matches_scan_documents() {
let (engine, _dir) = open_temp();
engine.put(0, 1, "users", "u1", b"alice").unwrap();
engine.put(0, 1, "users", "u2", b"bob").unwrap();
engine.put(0, 1, "users", "u3", b"carol").unwrap();
let materialized = engine.scan_documents(0, 1, "users", usize::MAX).unwrap();
let mut streamed: Vec<(String, Vec<u8>)> = Vec::new();
engine
.scan_documents_for_each(0, 1, "users", usize::MAX, |doc_id, bytes| {
streamed.push((doc_id.to_string(), bytes.to_vec()));
Ok(())
})
.unwrap();
assert_eq!(materialized, streamed);
}
#[test]
fn for_each_respects_limit() {
let (engine, _dir) = open_temp();
engine.put(0, 1, "users", "u1", b"alice").unwrap();
engine.put(0, 1, "users", "u2", b"bob").unwrap();
engine.put(0, 1, "users", "u3", b"carol").unwrap();
let materialized = engine.scan_documents(0, 1, "users", 2).unwrap();
let mut streamed: Vec<(String, Vec<u8>)> = Vec::new();
engine
.scan_documents_for_each(0, 1, "users", 2, |doc_id, bytes| {
streamed.push((doc_id.to_string(), bytes.to_vec()));
Ok(())
})
.unwrap();
assert_eq!(materialized.len(), 2);
assert_eq!(materialized, streamed);
}
#[test]
fn for_each_propagates_callback_error() {
let (engine, _dir) = open_temp();
engine.put(0, 1, "users", "u1", b"alice").unwrap();
engine.put(0, 1, "users", "u2", b"bob").unwrap();
let mut seen = 0usize;
let result =
engine.scan_documents_for_each(0, 1, "users", usize::MAX, |_doc_id, _bytes| {
seen += 1;
Err(crate::Error::Internal {
detail: "stop".to_string(),
})
});
assert!(result.is_err());
assert_eq!(seen, 1);
}
}