use std::collections::BTreeMap;
use std::sync::Arc;
use uqa_core::{DocId, FieldName, PathSegment, Value};
use crate::backend::{StorageBackendError, StorageBackendResult};
pub type Document = BTreeMap<FieldName, Value>;
const MISSING_SHARED_SLOT: usize = usize::MAX;
static SHARED_NULL_VALUE: Value = Value::Null;
#[derive(Debug, Clone, PartialEq)]
pub struct SharedDocumentRow {
values: Arc<Vec<Value>>,
projection: Arc<[usize]>,
}
impl SharedDocumentRow {
pub(crate) fn new(values: Arc<Vec<Value>>, projection: Arc<[usize]>) -> Self {
debug_assert!(projection
.iter()
.all(|slot| *slot == MISSING_SHARED_SLOT || *slot < values.len()));
Self { values, projection }
}
pub fn project<'a>(&'a self, output: &mut Vec<&'a Value>) {
output.clear();
output.extend(self.projection.iter().map(|slot| {
if *slot == MISSING_SHARED_SLOT {
&SHARED_NULL_VALUE
} else {
&self.values[*slot]
}
}));
}
pub fn with_projected<R>(&self, visitor: impl FnOnce(&[&Value]) -> R) -> R {
const INLINE_FIELDS: usize = 32;
if self.projection.len() <= INLINE_FIELDS {
let mut projected = [&SHARED_NULL_VALUE; INLINE_FIELDS];
for (output, slot) in projected.iter_mut().zip(self.projection.iter()) {
if *slot != MISSING_SHARED_SLOT {
*output = &self.values[*slot];
}
}
visitor(&projected[..self.projection.len()])
} else {
let projected = self
.projection
.iter()
.map(|slot| {
if *slot == MISSING_SHARED_SLOT {
&SHARED_NULL_VALUE
} else {
&self.values[*slot]
}
})
.collect::<Vec<_>>();
visitor(&projected)
}
}
pub fn indexed_values(&self) -> (&[Value], &[usize]) {
(&self.values, &self.projection)
}
pub fn into_parts(self) -> (Arc<Vec<Value>>, Arc<[usize]>) {
(self.values, self.projection)
}
}
pub trait DocumentStore: Send + Sync {
fn put(&mut self, doc_id: DocId, document: Document) -> StorageBackendResult<()>;
fn get(&self, doc_id: DocId) -> StorageBackendResult<Option<Document>>;
fn contains_doc_id(&self, doc_id: DocId) -> StorageBackendResult<bool> {
Ok(self.get(doc_id)?.is_some())
}
fn delete(&mut self, doc_id: DocId) -> StorageBackendResult<()>;
fn clear(&mut self) -> StorageBackendResult<()>;
fn get_field(&self, doc_id: DocId, field: &str) -> StorageBackendResult<Option<Value>> {
Ok(self
.get(doc_id)?
.and_then(|document| document.get(field).cloned()))
}
fn find_doc_id_by_field(
&self,
field: &str,
value: &Value,
) -> StorageBackendResult<Option<DocId>> {
for doc_id in self.doc_ids()? {
if self.get_field(doc_id, field)?.as_ref() == Some(value) {
return Ok(Some(doc_id));
}
}
Ok(None)
}
fn patch_fields(
&mut self,
doc_id: DocId,
updates: &BTreeMap<String, Value>,
) -> StorageBackendResult<bool> {
let Some(mut document) = self.get(doc_id)? else {
return Ok(false);
};
for (field, value) in updates {
if matches!(value, Value::Null) {
document.remove(field);
} else {
document.insert(field.clone(), value.clone());
}
}
self.put(doc_id, document)?;
Ok(true)
}
fn get_many(&self, doc_ids: &[DocId]) -> StorageBackendResult<BTreeMap<DocId, Document>> {
let mut out = BTreeMap::new();
for doc_id in doc_ids {
if let Some(document) = self.get(*doc_id)? {
out.insert(*doc_id, document);
}
}
Ok(out)
}
fn get_fields_multi(
&self,
doc_ids: &[DocId],
fields: &[&str],
) -> StorageBackendResult<BTreeMap<DocId, Vec<Value>>> {
let mut out = BTreeMap::new();
for doc_id in doc_ids {
let Some(document) = self.get(*doc_id)? else {
continue;
};
let values = fields
.iter()
.map(|field| document.get(*field).cloned().unwrap_or(Value::Null))
.collect();
out.insert(*doc_id, values);
}
Ok(out)
}
fn for_each_fields_multi(
&self,
doc_ids: &[DocId],
fields: &[&str],
visitor: &mut dyn FnMut(DocId, Vec<Value>) -> bool,
) -> StorageBackendResult<()> {
let mut projected = self.get_fields_multi(doc_ids, fields)?;
for doc_id in doc_ids {
let values = projected
.remove(doc_id)
.unwrap_or_else(|| vec![Value::Null; fields.len()]);
if !visitor(*doc_id, values) {
break;
}
}
Ok(())
}
fn for_each_fields_multi_ref(
&self,
doc_ids: &[DocId],
fields: &[&str],
visitor: &mut dyn FnMut(DocId, &[&Value]) -> bool,
) -> StorageBackendResult<()> {
self.for_each_fields_multi(doc_ids, fields, &mut |doc_id, values| {
let references: Vec<&Value> = values.iter().collect();
visitor(doc_id, &references)
})
}
fn for_each_fields_multi_ref_with_presence(
&self,
doc_ids: &[DocId],
fields: &[&str],
visitor: &mut dyn FnMut(DocId, bool, &[&Value]) -> bool,
) -> StorageBackendResult<()> {
if fields.is_empty() {
for doc_id in doc_ids {
if !visitor(*doc_id, self.contains_doc_id(*doc_id)?, &[]) {
break;
}
}
return Ok(());
}
let projected = self.get_fields_multi(doc_ids, fields)?;
let null = Value::Null;
let missing = vec![&null; fields.len()];
for doc_id in doc_ids {
let Some(values) = projected.get(doc_id) else {
if !visitor(*doc_id, false, &missing) {
break;
}
continue;
};
let references = values.iter().collect::<Vec<_>>();
if !visitor(*doc_id, true, &references) {
break;
}
}
Ok(())
}
fn get_shared_fields(
&self,
_doc_ids: &[DocId],
_fields: &[&str],
) -> StorageBackendResult<Option<Vec<Option<SharedDocumentRow>>>> {
Ok(None)
}
fn get_fields_bulk(
&self,
doc_ids: &[DocId],
field: &str,
) -> StorageBackendResult<BTreeMap<DocId, Value>> {
let mut out = BTreeMap::new();
for doc_id in doc_ids {
out.insert(
*doc_id,
self.get_field(*doc_id, field)?.unwrap_or(Value::Null),
);
}
Ok(out)
}
fn has_value(&self, field: &str, value: &Value) -> StorageBackendResult<bool> {
for doc_id in self.doc_ids()? {
if self.get_field(doc_id, field)?.as_ref() == Some(value) {
return Ok(true);
}
}
Ok(false)
}
fn find_doc_id_by_fields(
&self,
fields: &[String],
values: &[Value],
) -> StorageBackendResult<Option<DocId>> {
if fields.is_empty() || fields.len() != values.len() {
return Ok(None);
}
for doc_id in self.doc_ids()? {
let mut matches = true;
for (field, value) in fields.iter().zip(values) {
if self.get_field(doc_id, field)?.unwrap_or(Value::Null) != *value {
matches = false;
break;
}
}
if matches {
return Ok(Some(doc_id));
}
}
Ok(None)
}
fn eval_path(
&self,
doc_id: DocId,
path: &[PathSegment],
) -> StorageBackendResult<Option<Value>> {
let Some(document) = self.get(doc_id)? else {
return Ok(None);
};
Ok(eval_path_in_document(&document, path))
}
fn doc_ids(&self) -> StorageBackendResult<Vec<DocId>>;
fn next_doc_id(&self, after: Option<DocId>) -> StorageBackendResult<Option<DocId>> {
Ok(self
.doc_ids()?
.into_iter()
.filter(|doc_id| after.is_none_or(|after| *doc_id > after))
.min())
}
fn next_doc_ids(&self, after: Option<DocId>, limit: usize) -> StorageBackendResult<Vec<DocId>> {
if limit == 0 {
return Ok(Vec::new());
}
let mut doc_ids = self.doc_ids()?;
doc_ids.sort_unstable();
Ok(doc_ids
.into_iter()
.filter(|doc_id| after.is_none_or(|after| *doc_id > after))
.take(limit)
.collect())
}
fn next_shared_fields(
&self,
_after: Option<DocId>,
_limit: usize,
_fields: &[&str],
) -> StorageBackendResult<Option<Vec<(DocId, SharedDocumentRow)>>> {
Ok(None)
}
fn for_each_next_fields(
&self,
_after: Option<DocId>,
_limit: usize,
_fields: &[&str],
_visitor: &mut dyn FnMut(DocId, &[&Value]) -> bool,
) -> StorageBackendResult<Option<usize>> {
Ok(None)
}
fn max_doc_id(&self) -> StorageBackendResult<DocId> {
Ok(self.doc_ids()?.into_iter().max().unwrap_or(0))
}
fn len(&self) -> StorageBackendResult<usize>;
fn is_empty(&self) -> StorageBackendResult<bool> {
Ok(self.len()? == 0)
}
fn iter_all(&self) -> StorageBackendResult<Box<dyn Iterator<Item = (DocId, Document)> + '_>> {
let mut ids = self.doc_ids()?;
ids.sort_unstable();
let snapshot = self.snapshot()?;
let mut rows = Vec::with_capacity(ids.len());
for doc_id in ids {
if let Some(document) = snapshot.get(doc_id)? {
rows.push((doc_id, document));
}
}
Ok(Box::new(rows.into_iter()))
}
fn snapshot(&self) -> StorageBackendResult<Arc<dyn DocumentStore>>;
fn writable_snapshot(&self) -> StorageBackendResult<Box<dyn DocumentStore>> {
Err(StorageBackendError::Other(
"writable document-store snapshots are not supported by this backend".into(),
))
}
}
pub fn eval_path_in_document(doc: &Document, path: &[PathSegment]) -> Option<Value> {
let mut current: Value = match path.first()? {
PathSegment::Key(k) => doc.get(k)?.clone(),
PathSegment::Index(_) => return None,
};
for seg in path.iter().skip(1) {
current = match (current, seg) {
(Value::Map(m), PathSegment::Key(k)) => m.get(k)?.clone(),
(Value::List(items), PathSegment::Index(i)) => items.get(*i)?.clone(),
(Value::List(items), PathSegment::Key(k)) => {
let collected: Vec<Value> = items
.into_iter()
.filter_map(|v| match v {
Value::Map(m) => m.get(k).cloned(),
_ => None,
})
.collect();
Value::List(collected)
}
_ => return None,
};
}
Some(current)
}
mod memory;
pub use memory::MemoryDocumentStore;