use crate::guardian::core::{GuardianDB as BaseGuardianDB, NewGuardianDBOptions};
use crate::guardian::error::{GuardianError, Result};
use crate::p2p::network::client::IrohClient;
use crate::traits::{
AsyncDocumentFilter, BaseGuardianDB as BaseGuardianDBTrait, CreateDBOptions, Document,
DocumentStore, EventLogStore, GuardianDBKVStoreProvider, KeyValueStore, ProgressCallback,
Store,
};
use parking_lot::RwLock;
#[cfg(feature = "odm")]
use std::collections::BTreeSet;
use std::sync::Arc;
pub mod core;
pub mod error;
pub mod serializer;
pub struct GuardianDB {
base: BaseGuardianDB,
#[cfg(feature = "odm")]
collection_names: Arc<RwLock<BTreeSet<String>>>,
}
impl GuardianDB {
pub async fn new(client: IrohClient, options: Option<NewGuardianDBOptions>) -> Result<Self> {
use crate::log::identity::{DefaultIdentificator, Identificator};
let directory = options
.as_ref()
.and_then(|o| o.directory.clone())
.unwrap_or_else(|| std::path::PathBuf::from("./GuardianDB"));
let identity_path = directory.join("identity.json");
let identity = if identity_path.exists() {
match std::fs::read_to_string(&identity_path) {
Ok(data) => match serde_json::from_str::<crate::log::identity::Identity>(&data) {
Ok(id) => {
tracing::debug!("Loaded persisted identity from {:?}", identity_path);
id
}
Err(e) => {
tracing::warn!("Failed to deserialize identity file, creating new: {}", e);
let mut identificator = DefaultIdentificator::new();
let id = identificator.create(&client.node_id().to_string());
Self::save_identity(&identity_path, &id);
id
}
},
Err(e) => {
tracing::warn!("Failed to read identity file, creating new: {}", e);
let mut identificator = DefaultIdentificator::new();
let id = identificator.create(&client.node_id().to_string());
Self::save_identity(&identity_path, &id);
id
}
}
} else {
let mut identificator = DefaultIdentificator::new();
let id = identificator.create(&client.node_id().to_string());
Self::save_identity(&identity_path, &id);
id
};
let base = BaseGuardianDB::new_guardian_db(client, identity, options).await?;
Ok(GuardianDB {
base,
#[cfg(feature = "odm")]
collection_names: Arc::new(RwLock::new(BTreeSet::new())),
})
}
fn save_identity(path: &std::path::Path, identity: &crate::log::identity::Identity) {
if let Some(parent) = path.parent()
&& let Err(e) = std::fs::create_dir_all(parent)
{
tracing::warn!("Failed to create identity directory: {}", e);
return;
}
match serde_json::to_string_pretty(identity) {
Ok(data) => {
if let Err(e) = std::fs::write(path, data) {
tracing::warn!("Failed to write identity file: {}", e);
} else {
tracing::debug!("Identity persisted to {:?}", path);
}
}
Err(e) => {
tracing::warn!("Failed to serialize identity: {}", e);
}
}
}
pub async fn log(
&self,
address: &str,
options: Option<CreateDBOptions>,
) -> Result<Arc<dyn EventLogStore<Error = GuardianError>>> {
let mut opts = options.unwrap_or_default();
opts.create = Some(true);
opts.store_type = Some("eventlog".to_string());
if opts.event_bus.is_none() {
opts.event_bus = Some((*self.base.event_bus()).clone());
}
tracing::debug!(
address = address,
"GuardianDB::log - Creating EventLogStore"
);
let store = match self
.base
.create(address, "eventlog", Some(opts.clone()))
.await
{
Ok(store) => store,
Err(GuardianError::DatabaseAlreadyExists(existing_addr)) => {
tracing::debug!(address = address, full_addr = %existing_addr, "EventLogStore already exists, opening existing one");
self.base.open(&existing_addr, opts).await?
}
Err(e) => return Err(e),
};
tracing::debug!(
address = address,
store_type = store.store_type(),
has_index = store.index().supports_entry_queries(),
"EventLogStore created"
);
if store.store_type() == "eventlog" {
Ok(Arc::new(EventLogStoreWrapper::new(store)))
} else {
Err(GuardianError::Store(format!(
"Incorrect store type. Expected: eventlog, found: {}",
store.store_type()
)))
}
}
pub async fn key_value(
&self,
address: &str,
options: Option<CreateDBOptions>,
) -> Result<Arc<dyn KeyValueStore<Error = GuardianError>>> {
let mut opts = options.unwrap_or_default();
opts.create = Some(true);
opts.store_type = Some("keyvalue".to_string());
if opts.event_bus.is_none() {
opts.event_bus = Some((*self.base.event_bus()).clone());
}
let store = match self
.base
.create(address, "keyvalue", Some(opts.clone()))
.await
{
Ok(store) => store,
Err(GuardianError::DatabaseAlreadyExists(existing_addr)) => {
tracing::debug!(address = address, full_addr = %existing_addr, "KeyValueStore already exists, opening existing one");
self.base.open(&existing_addr, opts).await?
}
Err(e) => return Err(e),
};
if store.store_type() == "keyvalue" {
Ok(Arc::new(KeyValueStoreWrapper::new(store)))
} else {
Err(GuardianError::Store(format!(
"Incorrect store type. Expected: keyvalue, found: {}",
store.store_type()
)))
}
}
pub async fn docs(
&self,
address: &str,
options: Option<CreateDBOptions>,
) -> Result<Arc<dyn DocumentStore<Error = GuardianError>>> {
let mut opts = options.unwrap_or_default();
opts.create = Some(true);
opts.store_type = Some("document".to_string());
if opts.event_bus.is_none() {
opts.event_bus = Some((*self.base.event_bus()).clone());
}
let open_options = opts.clone();
let store = match self.base.create(address, "document", Some(opts)).await {
Ok(store) => store,
Err(GuardianError::DatabaseAlreadyExists(existing_addr)) => {
tracing::debug!(
address = address,
full_addr = %existing_addr,
"DocumentStore already exists, opening existing one"
);
self.base.open(&existing_addr, open_options).await?
}
Err(error) => return Err(error),
};
if store.store_type() == "document" {
#[cfg(feature = "odm")]
self.collection_names.write().insert(address.to_string());
Ok(Arc::new(DocumentStoreWrapper::new(store)))
} else {
Err(GuardianError::Store(format!(
"Incorrect store type. Expected: document, found: {}",
store.store_type()
)))
}
}
#[cfg(feature = "odm")]
pub async fn init_collection(&self, name: &str) -> crate::odm::Result<crate::odm::Collection> {
let store = self.docs(name, None).await?;
let storage = Arc::new(crate::odm::DocumentStoreStorage::new(store));
crate::odm::Collection::schemaless(name, storage).await
}
#[cfg(feature = "odm")]
pub async fn init_collection_with_schema(
&self,
name: &str,
schema: crate::odm::ModelSchema,
) -> crate::odm::Result<crate::odm::Collection> {
let store = self.docs(name, None).await?;
let storage = Arc::new(crate::odm::DocumentStoreStorage::new(store));
crate::odm::Collection::new(name, schema, storage).await
}
#[cfg(feature = "odm")]
pub async fn model_collection<M: crate::odm::Model>(
&self,
) -> crate::odm::Result<crate::odm::TypedCollection<M>> {
let schema = M::schema();
let store = self.docs(schema.collection(), None).await?;
let storage = Arc::new(crate::odm::DocumentStoreStorage::new(store));
crate::odm::TypedCollection::new(storage).await
}
#[cfg(feature = "odm")]
pub fn list_collections(&self) -> Vec<String> {
self.collection_names.read().iter().cloned().collect()
}
pub fn base(&self) -> &BaseGuardianDB {
&self.base
}
pub fn list_stores(
&self,
) -> Vec<(String, Arc<dyn Store<Error = GuardianError> + Send + Sync>)> {
self.base.list_stores()
}
pub fn register_access_control_type_with_name(
&self,
controller_type: &str,
constructor: crate::traits::AccessControllerConstructor,
) -> Result<()> {
self.base
.register_access_control_type_with_name(controller_type, constructor)
}
pub async fn register_access_control_type(
&self,
constructor: crate::traits::AccessControllerConstructor,
) -> Result<()> {
self.base.register_access_control_type(constructor).await
}
pub fn get_access_control_type(
&self,
controller_type: &str,
) -> Option<crate::traits::AccessControllerConstructor> {
self.base.get_access_control_type(controller_type)
}
pub fn access_control_types_names(&self) -> Vec<String> {
self.base.access_control_types_names()
}
pub async fn register_default_access_control_types(&self) -> Result<()> {
self.base.register_default_access_control_types().await
}
pub async fn connect_to_peer(&self, peer_id: iroh::EndpointId) -> Result<()> {
self.base.connect_to_peer(peer_id).await
}
}
pub struct EventLogStoreWrapper {
store: Arc<dyn Store<Error = GuardianError> + Send + Sync>,
}
impl EventLogStoreWrapper {
fn new(store: Arc<dyn Store<Error = GuardianError> + Send + Sync>) -> Self {
Self { store }
}
pub fn inner_store(&self) -> &Arc<dyn Store<Error = GuardianError> + Send + Sync> {
&self.store
}
pub fn try_get_basestore(&self) -> Option<&crate::stores::base_store::BaseStore> {
if let Some(event_log_store) =
self.store
.as_any()
.downcast_ref::<crate::stores::event_log_store::GuardianDBEventLogStore>()
{
return Some(event_log_store.basestore());
}
None
}
pub async fn connect_to_peer(&self, peer_id: iroh::EndpointId) -> Result<()> {
if let Some(base_store) = self.try_get_basestore() {
base_store.exchange_heads(peer_id).await
} else {
Err(GuardianError::Store(
"Could not access BaseStore for synchronization".to_string(),
))
}
}
fn query_from_index(
&self,
options: &crate::traits::StreamOptions,
) -> Result<Vec<crate::log::entry::Entry>> {
let index = self.store.index();
let is_simple_amount_query = options.gt.is_none()
&& options.gte.is_none()
&& options.lt.is_none()
&& options.lte.is_none();
if is_simple_amount_query {
let amount = match options.amount {
Some(a) if a > 0 => a as usize,
Some(-1) | None => {
match index.len() {
Ok(len) => len,
Err(_) => return self.query_from_oplog(options), }
}
_ => 0,
};
if let Some(entries) = index.get_last_entries(amount) {
return Ok(entries);
}
}
if let Some(hash) = options.gte.as_ref()
&& options.amount == Some(1)
&& options.gt.is_none()
&& options.lt.is_none()
&& options.lte.is_none()
{
if let Some(entry) = index.get_entry_by_hash(hash) {
return Ok(vec![entry]);
} else {
return Ok(Vec::new()); }
}
self.query_from_oplog(options)
}
fn query_from_oplog(
&self,
options: &crate::traits::StreamOptions,
) -> Result<Vec<crate::log::entry::Entry>> {
let oplog = self.store.op_log();
let oplog_guard = oplog.read();
let mut all_entries: Vec<_> = oplog_guard
.values()
.iter()
.map(|arc_entry| arc_entry.as_ref().clone())
.collect();
all_entries.sort_by_key(|b| b.clock().time());
let mut filtered_entries = all_entries;
if let Some(hash) = &options.gte {
if let Some(start_idx) = filtered_entries.iter().position(|e| e.hash() == hash) {
filtered_entries = filtered_entries.into_iter().skip(start_idx).collect();
} else {
return Ok(Vec::new()); }
}
if let Some(hash) = &options.gt {
if let Some(start_idx) = filtered_entries.iter().position(|e| e.hash() == hash) {
filtered_entries = filtered_entries.into_iter().skip(start_idx + 1).collect();
} else {
return Ok(Vec::new()); }
}
if let Some(hash) = &options.lte {
if let Some(end_idx) = filtered_entries.iter().position(|e| e.hash() == hash) {
filtered_entries = filtered_entries.into_iter().take(end_idx + 1).collect();
} else {
return Ok(Vec::new()); }
}
if let Some(hash) = &options.lt {
if let Some(end_idx) = filtered_entries.iter().position(|e| e.hash() == hash) {
filtered_entries = filtered_entries.into_iter().take(end_idx).collect();
} else {
return Ok(Vec::new()); }
}
let amount = match options.amount {
Some(a) if a > 0 => a as usize,
Some(-1) | None => filtered_entries.len(), _ => 0,
};
filtered_entries.truncate(amount);
Ok(filtered_entries)
}
}
#[async_trait::async_trait]
impl Store for EventLogStoreWrapper {
type Error = GuardianError;
#[allow(deprecated)]
fn events(&self) -> &dyn crate::events::EmitterInterface {
self.store.events()
}
async fn close(&self) -> std::result::Result<(), Self::Error> {
self.store.close().await
}
fn address(&self) -> &dyn crate::address::Address {
self.store.address()
}
fn index(&self) -> Box<dyn crate::traits::StoreIndex<Error = GuardianError> + Send + Sync> {
self.store.index()
}
fn store_type(&self) -> &str {
self.store.store_type()
}
fn cache(&self) -> Arc<dyn crate::data_store::Datastore> {
self.store.cache()
}
async fn drop(&self) -> std::result::Result<(), Self::Error> {
Ok(())
}
async fn load(&self, amount: usize) -> std::result::Result<(), Self::Error> {
self.store.load(amount).await
}
async fn sync(
&self,
heads: Vec<crate::log::entry::Entry>,
) -> std::result::Result<(), Self::Error> {
self.store.sync(heads).await
}
async fn load_more_from(&self, _amount: u64, entries: Vec<crate::log::entry::Entry>) {
self.store.load_more_from(_amount, entries).await
}
async fn load_from_snapshot(&self) -> std::result::Result<(), Self::Error> {
self.store.load_from_snapshot().await
}
fn op_log(&self) -> Arc<RwLock<crate::log::Log>> {
self.store.op_log()
}
fn client(&self) -> Arc<IrohClient> {
unimplemented!("Adaptation between iroh client types pending")
}
fn db_name(&self) -> &str {
self.store.db_name()
}
fn identity(&self) -> &crate::log::identity::Identity {
self.store.identity()
}
fn access_controller(&self) -> &dyn crate::access_control::traits::AccessController {
self.store.access_controller()
}
async fn add_operation(
&self,
op: crate::stores::operation::Operation,
on_progress_callback: Option<ProgressCallback>,
) -> std::result::Result<crate::log::entry::Entry, Self::Error> {
self.store.add_operation(op, on_progress_callback).await
}
fn span(&self) -> Arc<tracing::Span> {
self.store.span()
}
fn tracer(&self) -> Arc<crate::traits::TracerWrapper> {
self.store.tracer()
}
fn event_bus(&self) -> Arc<crate::p2p::EventBus> {
self.store.event_bus()
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}
#[async_trait::async_trait]
impl EventLogStore for EventLogStoreWrapper {
async fn add(
&self,
data: Vec<u8>,
) -> std::result::Result<crate::stores::operation::Operation, Self::Error> {
let operation =
crate::stores::operation::Operation::new(None, "ADD".to_string(), Some(data));
let _entry = self.add_operation(operation.clone(), None).await?;
Ok(operation)
}
async fn get(
&self,
hash: &iroh_blobs::Hash,
) -> std::result::Result<crate::stores::operation::Operation, Self::Error> {
if let Some(entry) = self.store.index().get_entry_by_hash(hash) {
let operation = crate::stores::operation::parse_operation(entry)
.map_err(|e| GuardianError::Store(format!("Failed to parse entry: {}", e)))?;
return Ok(operation);
}
let oplog = self.store.op_log();
let oplog_guard = oplog.read();
for arc_entry in oplog_guard.values() {
if arc_entry.hash() == hash {
let entry = arc_entry.as_ref().clone();
let operation = crate::stores::operation::parse_operation(entry)
.map_err(|e| GuardianError::Store(format!("Failed to parse entry: {}", e)))?;
return Ok(operation);
}
}
Err(GuardianError::Store(format!(
"Operation not found for Hash: {}",
hex::encode(hash.as_bytes())
)))
}
async fn list(
&self,
options: Option<crate::traits::StreamOptions>,
) -> std::result::Result<Vec<crate::stores::operation::Operation>, Self::Error> {
let options = options.unwrap_or_default();
let entries = if self.store.index().supports_entry_queries() {
self.query_from_index(&options)?
} else {
self.query_from_oplog(&options)?
};
let mut operations = Vec::with_capacity(entries.len());
for entry in entries {
match crate::stores::operation::parse_operation(entry) {
Ok(operation) => operations.push(operation),
Err(e) => {
eprintln!("Warning: Failed to parse entry: {}", e);
}
}
}
Ok(operations)
}
}
struct KeyValueStoreWrapper {
store: Arc<dyn Store<Error = GuardianError> + Send + Sync>,
}
impl KeyValueStoreWrapper {
fn new(store: Arc<dyn Store<Error = GuardianError> + Send + Sync>) -> Self {
Self { store }
}
}
#[async_trait::async_trait]
impl Store for KeyValueStoreWrapper {
type Error = GuardianError;
#[allow(deprecated)]
fn events(&self) -> &dyn crate::events::EmitterInterface {
self.store.events()
}
async fn close(&self) -> std::result::Result<(), Self::Error> {
self.store.close().await
}
fn address(&self) -> &dyn crate::address::Address {
self.store.address()
}
fn index(&self) -> Box<dyn crate::traits::StoreIndex<Error = GuardianError> + Send + Sync> {
self.store.index()
}
fn store_type(&self) -> &str {
self.store.store_type()
}
fn cache(&self) -> Arc<dyn crate::data_store::Datastore> {
self.store.cache()
}
async fn drop(&self) -> std::result::Result<(), Self::Error> {
Ok(())
}
async fn load(&self, amount: usize) -> std::result::Result<(), Self::Error> {
self.store.load(amount).await
}
async fn sync(
&self,
heads: Vec<crate::log::entry::Entry>,
) -> std::result::Result<(), Self::Error> {
self.store.sync(heads).await
}
async fn load_more_from(&self, _amount: u64, entries: Vec<crate::log::entry::Entry>) {
self.store.load_more_from(_amount, entries).await
}
async fn load_from_snapshot(&self) -> std::result::Result<(), Self::Error> {
self.store.load_from_snapshot().await
}
fn op_log(&self) -> Arc<RwLock<crate::log::Log>> {
self.store.op_log()
}
fn client(&self) -> Arc<IrohClient> {
unimplemented!("Adaptation between iroh client types pending")
}
fn db_name(&self) -> &str {
self.store.db_name()
}
fn identity(&self) -> &crate::log::identity::Identity {
self.store.identity()
}
fn access_controller(&self) -> &dyn crate::access_control::traits::AccessController {
self.store.access_controller()
}
async fn add_operation(
&self,
op: crate::stores::operation::Operation,
on_progress_callback: Option<ProgressCallback>,
) -> std::result::Result<crate::log::entry::Entry, Self::Error> {
self.store.add_operation(op, on_progress_callback).await
}
fn span(&self) -> Arc<tracing::Span> {
self.store.span()
}
fn tracer(&self) -> Arc<crate::traits::TracerWrapper> {
self.store.tracer()
}
fn event_bus(&self) -> Arc<crate::p2p::EventBus> {
self.store.event_bus()
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}
#[async_trait::async_trait]
impl KeyValueStore for KeyValueStoreWrapper {
async fn get(&self, key: &str) -> std::result::Result<Option<Vec<u8>>, Self::Error> {
let index = self.store.index();
if let Ok(Some(bytes)) = index.get_bytes(key) {
return Ok(Some(bytes));
}
let oplog = self.store.op_log();
let oplog_guard = oplog.read();
let mut latest_value: Option<Vec<u8>> = None;
let mut latest_time = 0;
for arc_entry in oplog_guard.values() {
let entry = arc_entry.as_ref().clone();
if let Ok(operation) = crate::stores::operation::parse_operation(entry.clone()) {
if let Some(op_key) = operation.key()
&& op_key == key
{
let entry_time = entry.clock().time();
let op_str = operation.op();
if op_str == "PUT" {
if entry_time > latest_time {
latest_time = entry_time;
latest_value = Some(operation.value().to_vec());
}
} else if op_str == "DEL" {
if entry_time > latest_time {
latest_time = entry_time;
latest_value = None; }
}
}
}
}
Ok(latest_value)
}
async fn put(
&self,
key: &str,
value: Vec<u8>,
) -> std::result::Result<crate::stores::operation::Operation, Self::Error> {
let operation = crate::stores::operation::Operation::new(
Some(key.to_string()),
"PUT".to_string(),
Some(value),
);
self.add_operation(operation.clone(), None).await?;
Ok(operation)
}
async fn delete(
&self,
key: &str,
) -> std::result::Result<crate::stores::operation::Operation, Self::Error> {
let operation = crate::stores::operation::Operation::new(
Some(key.to_string()),
"DEL".to_string(),
None,
);
self.add_operation(operation.clone(), None).await?;
Ok(operation)
}
fn all(&self) -> std::collections::HashMap<String, Vec<u8>> {
let mut result = std::collections::HashMap::new();
let index = self.store.index();
if let Ok(keys) = index.keys() {
for key in keys {
if let Ok(Some(bytes)) = index.get_bytes(&key) {
result.insert(key, bytes);
}
}
}
if result.is_empty() {
let oplog = self.store.op_log();
let oplog_guard = oplog.read();
let mut key_operations: std::collections::HashMap<
String,
(u64, String, Option<Vec<u8>>),
> = std::collections::HashMap::new();
for arc_entry in oplog_guard.values() {
let entry = arc_entry.as_ref().clone();
if let Ok(operation) = crate::stores::operation::parse_operation(entry.clone())
&& let Some(op_key) = operation.key()
{
let timestamp = entry.clock().time();
let op_type = operation.op().to_string();
let value = if !operation.value().is_empty() {
Some(operation.value().to_vec())
} else {
None
};
let key_clone = op_key.clone();
if let Some((existing_time, _, _)) = key_operations.get(&key_clone) {
if timestamp > *existing_time {
key_operations.insert(key_clone, (timestamp, op_type, value));
}
} else {
key_operations.insert(key_clone, (timestamp, op_type, value));
}
}
}
for (key, (_timestamp, op_type, value)) in key_operations {
let op_str = op_type.as_str();
if op_str == "PUT" {
if let Some(val) = value {
result.insert(key, val);
}
} else if op_str == "DEL" {
result.remove(&key);
} else {
if let Some(val) = value {
result.insert(key, val);
}
}
}
}
result
}
async fn share_ticket(&self) -> std::result::Result<String, Self::Error> {
self.store
.as_any()
.downcast_ref::<crate::stores::kv_store::GuardianDBKeyValue>()
.ok_or_else(|| {
GuardianError::Store(
"share_ticket: underlying store is not a GuardianDBKeyValue".to_string(),
)
})?
.share_ticket()
.await
}
}
struct DocumentStoreWrapper {
store: Arc<dyn Store<Error = GuardianError> + Send + Sync>,
}
impl DocumentStoreWrapper {
fn new(store: Arc<dyn Store<Error = GuardianError> + Send + Sync>) -> Self {
Self { store }
}
fn search_documents_by_key(
&self,
key: &str,
opts: &crate::traits::DocumentStoreGetOptions,
) -> Result<Vec<Document>> {
let index = self.store.index();
let mut key_for_search = key.to_string();
let has_multiple_terms = key.contains(' ');
if has_multiple_terms {
key_for_search = key_for_search.replace('.', " ");
}
if opts.case_insensitive {
key_for_search = key_for_search.to_lowercase();
}
let mut documents = Vec::new();
let all_keys = index.keys().unwrap_or_default();
for index_key in all_keys {
let mut index_key_for_search = index_key.clone();
if opts.case_insensitive {
index_key_for_search = index_key_for_search.to_lowercase();
}
let matches = if opts.partial_matches {
index_key_for_search.contains(&key_for_search)
} else {
index_key_for_search == key_for_search
};
if matches {
if let Ok(Some(doc_bytes)) = index.get_bytes(&index_key) {
match serde_json::from_slice::<serde_json::Value>(&doc_bytes) {
Ok(json_value) => {
let doc: Document = Box::new(json_value);
documents.push(doc);
}
Err(e) => {
eprintln!(
"Warning: Failed to deserialize document for key '{}': {}",
index_key, e
);
}
}
} else {
eprintln!(
"Warning: key '{}' found but without a corresponding value",
index_key
);
}
}
}
Ok(documents)
}
fn search_documents_from_oplog(
&self,
key: &str,
opts: &crate::traits::DocumentStoreGetOptions,
) -> Result<Vec<Document>> {
let oplog = self.store.op_log();
let oplog_guard = oplog.read();
let mut documents = Vec::new();
let mut processed_keys = std::collections::HashSet::new();
let entries: Vec<Arc<crate::log::entry::Entry>> =
oplog_guard.values().into_iter().collect();
for arc_entry in entries.iter().rev() {
let entry: crate::log::entry::Entry = (**arc_entry).clone();
if let Ok(operation) = crate::stores::operation::parse_operation(entry) {
if let Some(op_key) = operation.key() {
if processed_keys.contains(op_key) {
continue;
}
let mut op_key_search = op_key.clone();
let mut key_search = key.to_string();
if opts.case_insensitive {
op_key_search = op_key_search.to_lowercase();
key_search = key_search.to_lowercase();
}
let matches = if opts.partial_matches {
op_key_search.contains(&key_search)
} else {
op_key_search == key_search
};
if matches {
processed_keys.insert(op_key.clone());
if operation.op() == "DEL" {
continue;
}
if operation.op() == "PUT" && !operation.value().is_empty() {
match serde_json::from_slice::<serde_json::Value>(operation.value()) {
Ok(json_value) => {
let doc: Document = Box::new(json_value);
documents.push(doc);
}
Err(_) => {
let simple_doc = serde_json::json!({
"key": op_key,
"value": String::from_utf8_lossy(operation.value()),
"op_type": operation.op()
});
let doc: Document = Box::new(simple_doc);
documents.push(doc);
}
}
}
}
}
}
}
Ok(documents)
}
fn get_all_documents_from_index(&self) -> Result<Vec<Document>> {
let index = self.store.index();
let mut documents = Vec::new();
let all_keys = index.keys().unwrap_or_default();
eprintln!(
"DEBUG: get_all_documents_from_index - Total keys in the index: {}",
all_keys.len()
);
for key in all_keys {
eprintln!(
"DEBUG: get_all_documents_from_index - Processing key: {}",
key
);
if let Ok(Some(doc_bytes)) = index.get_bytes(&key) {
eprintln!(
"DEBUG: get_all_documents_from_index - Bytes retrieved for key '{}': {} bytes",
key,
doc_bytes.len()
);
match serde_json::from_slice::<serde_json::Value>(&doc_bytes) {
Ok(json_value) => {
eprintln!(
"DEBUG: get_all_documents_from_index - Document deserialized successfully: {:?}",
json_value
);
let doc: Document = Box::new(json_value);
documents.push(doc);
}
Err(e) => {
eprintln!(
"Warning: Failed to deserialize document for key '{}': {}",
key, e
);
}
}
} else {
eprintln!(
"DEBUG: get_all_documents_from_index - No bytes found for key: {}",
key
);
}
}
eprintln!(
"DEBUG: get_all_documents_from_index - Total documents collected: {}",
documents.len()
);
Ok(documents)
}
}
#[async_trait::async_trait]
impl Store for DocumentStoreWrapper {
type Error = GuardianError;
#[allow(deprecated)]
fn events(&self) -> &dyn crate::events::EmitterInterface {
self.store.events()
}
async fn close(&self) -> std::result::Result<(), Self::Error> {
self.store.close().await
}
fn address(&self) -> &dyn crate::address::Address {
self.store.address()
}
fn index(&self) -> Box<dyn crate::traits::StoreIndex<Error = GuardianError> + Send + Sync> {
self.store.index()
}
fn store_type(&self) -> &str {
self.store.store_type()
}
fn cache(&self) -> Arc<dyn crate::data_store::Datastore> {
self.store.cache()
}
async fn drop(&self) -> std::result::Result<(), Self::Error> {
Ok(())
}
async fn load(&self, amount: usize) -> std::result::Result<(), Self::Error> {
self.store.load(amount).await
}
async fn sync(
&self,
heads: Vec<crate::log::entry::Entry>,
) -> std::result::Result<(), Self::Error> {
self.store.sync(heads).await
}
async fn load_more_from(&self, _amount: u64, entries: Vec<crate::log::entry::Entry>) {
self.store.load_more_from(_amount, entries).await
}
async fn load_from_snapshot(&self) -> std::result::Result<(), Self::Error> {
self.store.load_from_snapshot().await
}
fn op_log(&self) -> Arc<RwLock<crate::log::Log>> {
self.store.op_log()
}
fn client(&self) -> Arc<IrohClient> {
unimplemented!("Adaptation between iroh client types pending")
}
fn db_name(&self) -> &str {
self.store.db_name()
}
fn identity(&self) -> &crate::log::identity::Identity {
self.store.identity()
}
fn access_controller(&self) -> &dyn crate::access_control::traits::AccessController {
self.store.access_controller()
}
async fn add_operation(
&self,
op: crate::stores::operation::Operation,
on_progress_callback: Option<ProgressCallback>,
) -> std::result::Result<crate::log::entry::Entry, Self::Error> {
self.store.add_operation(op, on_progress_callback).await
}
fn span(&self) -> Arc<tracing::Span> {
self.store.span()
}
fn tracer(&self) -> Arc<crate::traits::TracerWrapper> {
self.store.tracer()
}
fn event_bus(&self) -> Arc<crate::p2p::EventBus> {
self.store.event_bus()
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}
#[async_trait::async_trait]
impl DocumentStore for DocumentStoreWrapper {
async fn put(
&self,
document: Document,
) -> std::result::Result<crate::stores::operation::Operation, Self::Error> {
let key = if let Some(json_val) = document.downcast_ref::<serde_json::Value>() {
json_val
.get("_id")
.or_else(|| json_val.get("id"))
.or_else(|| json_val.get("key"))
.and_then(|v| v.as_str())
.map(|s| s.to_string())
} else {
None
};
let data = if let Some(json_val) = document.downcast_ref::<serde_json::Value>() {
serde_json::to_vec(json_val).map_err(|e| {
GuardianError::Store(format!("Failed to serialize JSON document: {}", e))
})?
} else if let Some(bytes) = document.downcast_ref::<Vec<u8>>() {
bytes.clone()
} else {
format!("{:?}", document).into_bytes()
};
let operation =
crate::stores::operation::Operation::new(key, "PUT".to_string(), Some(data));
self.add_operation(operation.clone(), None).await?;
Ok(operation)
}
async fn delete(
&self,
key: &str,
) -> std::result::Result<crate::stores::operation::Operation, Self::Error> {
let operation = crate::stores::operation::Operation::new(
Some(key.to_string()),
"DEL".to_string(),
None,
);
self.add_operation(operation.clone(), None).await?;
Ok(operation)
}
async fn put_batch(
&self,
values: Vec<Document>,
) -> std::result::Result<crate::stores::operation::Operation, Self::Error> {
if values.is_empty() {
return Err(GuardianError::InvalidArgument(
"Nothing to add to the store".to_string(),
));
}
let mut last_operation = None;
for document in values {
let op = self.put(document).await?;
last_operation = Some(op);
}
Ok(last_operation.unwrap())
}
async fn put_all(
&self,
values: Vec<Document>,
) -> std::result::Result<crate::stores::operation::Operation, Self::Error> {
if values.is_empty() {
return Err(GuardianError::InvalidArgument(
"Nothing to add to the store".to_string(),
));
}
let mut last_operation = None;
for document in values {
let op = self.put(document).await?;
last_operation = Some(op);
}
Ok(last_operation.unwrap())
}
async fn get(
&self,
key: &str,
opts: Option<crate::traits::DocumentStoreGetOptions>,
) -> std::result::Result<Vec<Document>, Self::Error> {
let opts = opts.unwrap_or_default();
let documents_from_index = self.search_documents_by_key(key, &opts)?;
if !documents_from_index.is_empty() {
return Ok(documents_from_index);
}
let documents_from_oplog = self.search_documents_from_oplog(key, &opts)?;
Ok(documents_from_oplog)
}
async fn query(
&self,
filter: AsyncDocumentFilter,
) -> std::result::Result<Vec<Document>, Self::Error> {
let all_documents = self.get_all_documents_from_index()?;
let mut filtered_documents = Vec::new();
for document in all_documents {
let filter_future = filter(&document);
match filter_future.await {
Ok(true) => {
filtered_documents.push(document);
}
Ok(false) => {
continue;
}
Err(e) => {
eprintln!("Warning: Error applying filter to the document: {}", e);
continue;
}
}
}
Ok(filtered_documents)
}
async fn share_ticket(&self) -> std::result::Result<String, Self::Error> {
self.store
.as_any()
.downcast_ref::<crate::stores::document_store::GuardianDBDocumentStore>()
.ok_or_else(|| {
GuardianError::Store(
"share_ticket: underlying store is not a GuardianDBDocumentStore".to_string(),
)
})?
.share_ticket()
.await
}
}
#[async_trait::async_trait]
impl BaseGuardianDBTrait for GuardianDB {
type Error = GuardianError;
async fn open(
&self,
address: &str,
options: &mut CreateDBOptions,
) -> std::result::Result<Arc<dyn Store<Error = GuardianError>>, Self::Error> {
let opts = options.clone();
let result = self.base.open(address, opts).await?;
Ok(result as Arc<dyn Store<Error = GuardianError>>)
}
async fn determine_address(
&self,
name: &str,
store_type: &str,
options: &crate::traits::DetermineAddressOptions,
) -> std::result::Result<Box<dyn crate::address::Address>, Self::Error> {
let opts = Some(options.clone());
let result = self.base.determine_address(name, store_type, opts).await?;
Ok(Box::new(result))
}
fn client(&self) -> Arc<crate::p2p::network::client::IrohClient> {
Arc::new(self.base.client().clone())
}
fn identity(&self) -> Arc<crate::log::identity::Identity> {
Arc::new(self.base.identity().clone())
}
fn get_store(&self, address: &str) -> Option<Arc<dyn Store<Error = GuardianError>>> {
self.base
.get_store(address)
.map(|store| store as Arc<dyn Store<Error = GuardianError>>)
}
async fn create(
&self,
name: &str,
store_type: &str,
options: &mut CreateDBOptions,
) -> std::result::Result<Arc<dyn Store<Error = GuardianError>>, Self::Error> {
let opts = Some(options.clone());
let result = self.base.create(name, store_type, opts).await?;
Ok(result as Arc<dyn Store<Error = GuardianError>>)
}
fn register_store_type(
&mut self,
store_type: &str,
constructor: crate::traits::StoreConstructor,
) {
self.base
.register_store_type(store_type.to_string(), constructor);
}
fn unregister_store_type(&mut self, store_type: &str) {
self.base.unregister_store_type(store_type);
}
fn register_access_controller_type(
&mut self,
constructor: crate::traits::AccessControllerConstructor,
) -> std::result::Result<(), Self::Error> {
self.base
.register_access_control_type_with_name("simple", constructor)
}
fn unregister_access_controller_type(&mut self, controller_type: &str) {
self.base.unregister_access_control_type(controller_type);
}
fn get_access_controller_type(
&self,
controller_type: &str,
) -> Option<crate::traits::AccessControllerConstructor> {
self.base.get_access_controller_type(controller_type)
}
fn event_bus(&self) -> crate::p2p::EventBus {
(*self.base.event_bus()).clone()
}
fn span(&self) -> &tracing::Span {
self.base.span()
}
fn tracer(&self) -> Arc<crate::traits::TracerWrapper> {
let boxed_tracer = self.base.tracer();
Arc::new(crate::traits::TracerWrapper::new_opentelemetry(
boxed_tracer,
))
}
}
#[async_trait::async_trait]
impl GuardianDBKVStoreProvider for GuardianDB {
type Error = GuardianError;
async fn key_value(
&self,
address: &str,
options: &mut CreateDBOptions,
) -> std::result::Result<Box<dyn KeyValueStore<Error = GuardianError>>, Self::Error> {
let opts_clone = options.clone();
let arc_store = self.key_value(address, Some(opts_clone)).await?;
Ok(Box::new(KeyValueStoreBoxWrapper::new(arc_store)))
}
}
pub struct KeyValueStoreBoxWrapper {
inner: Arc<dyn KeyValueStore<Error = GuardianError>>,
}
impl KeyValueStoreBoxWrapper {
pub fn new(inner: Arc<dyn KeyValueStore<Error = GuardianError>>) -> Self {
Self { inner }
}
}
#[async_trait::async_trait]
impl Store for KeyValueStoreBoxWrapper {
type Error = GuardianError;
fn address(&self) -> &dyn crate::address::Address {
self.inner.address()
}
fn store_type(&self) -> &str {
self.inner.store_type()
}
async fn close(&self) -> std::result::Result<(), Self::Error> {
self.inner.close().await
}
async fn drop(&self) -> std::result::Result<(), Self::Error> {
self.inner.close().await
}
fn events(&self) -> &dyn crate::events::EmitterInterface {
unimplemented!("events() is deprecated, use event_bus() instead")
}
fn index(&self) -> Box<dyn crate::traits::StoreIndex<Error = Self::Error> + Send + Sync> {
self.inner.index()
}
fn cache(&self) -> Arc<dyn crate::data_store::Datastore> {
self.inner.cache()
}
async fn load(&self, amount: usize) -> std::result::Result<(), Self::Error> {
self.inner.load(amount).await
}
async fn sync(
&self,
heads: Vec<crate::log::entry::Entry>,
) -> std::result::Result<(), Self::Error> {
self.inner.sync(heads).await
}
async fn load_more_from(&self, _amount: u64, entries: Vec<crate::log::entry::Entry>) {
self.inner.load_more_from(_amount, entries).await
}
async fn load_from_snapshot(&self) -> std::result::Result<(), Self::Error> {
self.inner.load_from_snapshot().await
}
fn op_log(&self) -> Arc<parking_lot::RwLock<crate::log::Log>> {
self.inner.op_log()
}
fn client(&self) -> Arc<crate::p2p::network::client::IrohClient> {
self.inner.client()
}
fn db_name(&self) -> &str {
self.inner.db_name()
}
fn identity(&self) -> &crate::log::identity::Identity {
self.inner.identity()
}
fn access_controller(&self) -> &dyn crate::access_control::traits::AccessController {
self.inner.access_controller()
}
async fn add_operation(
&self,
op: crate::stores::operation::Operation,
on_progress_callback: Option<crate::traits::ProgressCallback>,
) -> std::result::Result<crate::log::entry::Entry, Self::Error> {
self.inner.add_operation(op, on_progress_callback).await
}
fn span(&self) -> Arc<tracing::Span> {
self.inner.span()
}
fn tracer(&self) -> Arc<crate::traits::TracerWrapper> {
self.inner.tracer()
}
fn event_bus(&self) -> Arc<crate::p2p::EventBus> {
self.inner.event_bus()
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}
#[async_trait::async_trait]
impl KeyValueStore for KeyValueStoreBoxWrapper {
async fn put(
&self,
key: &str,
value: Vec<u8>,
) -> std::result::Result<crate::stores::operation::Operation, Self::Error> {
self.inner.put(key, value).await
}
async fn get(&self, key: &str) -> std::result::Result<Option<Vec<u8>>, Self::Error> {
self.inner.get(key).await
}
async fn delete(
&self,
key: &str,
) -> std::result::Result<crate::stores::operation::Operation, Self::Error> {
self.inner.delete(key).await
}
fn all(&self) -> std::collections::HashMap<String, Vec<u8>> {
self.inner.all()
}
async fn share_ticket(&self) -> std::result::Result<String, Self::Error> {
self.inner.share_ticket().await
}
}