use crate::access_control::acl_simple::SimpleAccessController;
use crate::access_control::traits::AccessController;
use crate::address::Address;
use crate::data_store::Datastore;
use crate::events::EventEmitter;
use crate::guardian::error::{GuardianError, Result};
use crate::log::identity::Identity;
use crate::log::lamport_clock::LamportClock;
use crate::p2p::EventBus;
use crate::p2p::network::core::docs::WillowDocs;
use crate::stores::operation::Operation;
use crate::traits::{KeyValueStore, NewStoreOptions, Store, StoreIndex, TracerWrapper};
use bytes::Bytes;
use iroh_docs::{AuthorId, Capability, api::Doc, store::Query};
use opentelemetry::trace::{TracerProvider, noop::NoopTracerProvider};
use parking_lot::RwLock;
use std::collections::HashMap;
use std::sync::Arc;
use tracing::{Span, debug, info, instrument, warn};
pub mod index;
const NAMESPACE_CACHE_KEY: &[u8] = b"_iroh_docs_namespace_id";
const WRITABLE_CACHE_KEY: &[u8] = b"_iroh_docs_writable";
pub struct KeyValueIndex {
index: Arc<RwLock<HashMap<String, Vec<u8>>>>,
}
impl Default for KeyValueIndex {
fn default() -> Self {
Self::new()
}
}
impl KeyValueIndex {
pub fn new() -> Self {
Self {
index: Arc::new(RwLock::new(HashMap::new())),
}
}
pub fn get_value(&self, key: &str) -> Option<Vec<u8>> {
let guard = self.index.read();
guard.get(key).cloned()
}
pub fn get_all(&self) -> HashMap<String, Vec<u8>> {
let guard = self.index.read();
guard.clone()
}
pub fn len(&self) -> usize {
let guard = self.index.read();
guard.len()
}
pub fn is_empty(&self) -> bool {
let guard = self.index.read();
guard.is_empty()
}
pub fn insert(&self, key: String, value: Vec<u8>) {
let mut guard = self.index.write();
guard.insert(key, value);
}
pub fn remove(&self, key: &str) {
let mut guard = self.index.write();
guard.remove(key);
}
pub fn clear_all(&self) {
let mut guard = self.index.write();
guard.clear();
}
}
async fn refresh_kv_index(
docs: &WillowDocs,
doc: &Doc,
client: &Arc<crate::p2p::network::client::IrohClient>,
index: &Arc<KeyValueIndex>,
) -> Result<usize> {
let entries = docs
.get_many(doc, Query::single_latest_per_key().build())
.await?;
index.clear_all();
let mut count = 0;
for entry in &entries {
let key = String::from_utf8_lossy(entry.key()).to_string();
if entry.content_len() == 0 {
continue;
}
let hash_str = entry.content_hash().to_hex();
match client.cat_bytes(&hash_str).await {
Ok(value) => {
index.insert(key, value);
count += 1;
}
Err(e) => {
warn!("Failed to read content for key from iroh-docs: {:?}", e);
}
}
}
debug!(
"KeyValue index synchronized from iroh-docs: {} entries",
count
);
Ok(count)
}
impl StoreIndex for KeyValueIndex {
type Error = GuardianError;
fn contains_key(&self, key: &str) -> std::result::Result<bool, Self::Error> {
let guard = self.index.read();
Ok(guard.contains_key(key))
}
fn get_bytes(&self, key: &str) -> std::result::Result<Option<Vec<u8>>, Self::Error> {
let guard = self.index.read();
Ok(guard.get(key).cloned())
}
fn keys(&self) -> std::result::Result<Vec<String>, Self::Error> {
let guard = self.index.read();
Ok(guard.keys().cloned().collect())
}
fn len(&self) -> std::result::Result<usize, Self::Error> {
let guard = self.index.read();
Ok(guard.len())
}
fn is_empty(&self) -> std::result::Result<bool, Self::Error> {
let guard = self.index.read();
Ok(guard.is_empty())
}
fn update_index(
&mut self,
_log: &crate::log::Log,
_entries: &[crate::log::entry::Entry],
) -> std::result::Result<(), Self::Error> {
Ok(())
}
fn clear(&mut self) -> std::result::Result<(), Self::Error> {
let mut guard = self.index.write();
guard.clear();
Ok(())
}
}
pub struct GuardianDBKeyValue {
docs: WillowDocs,
doc_handle: Doc,
author_id: AuthorId,
access_controller: Arc<dyn AccessController>,
event_bus: Arc<EventBus>,
client: Arc<crate::p2p::network::client::IrohClient>,
identity: Arc<Identity>,
cached_address: Arc<dyn Address + Send + Sync>,
db_name: String,
cache: Arc<dyn Datastore>,
index: Arc<KeyValueIndex>,
span: Span,
tracer: Arc<TracerWrapper>,
emitter_interface: Arc<dyn crate::events::EmitterInterface + Send + Sync>,
empty_log: Arc<RwLock<crate::log::Log>>,
writable: bool,
}
#[async_trait::async_trait]
impl Store for GuardianDBKeyValue {
type Error = GuardianError;
#[allow(deprecated)]
fn events(&self) -> &dyn crate::events::EmitterInterface {
self.emitter_interface.as_ref()
}
async fn close(&self) -> Result<()> {
debug!("Starting KeyValue store close operation (iroh-docs backend)");
if let Err(e) = self.docs.close_doc(&self.doc_handle).await {
warn!("Failed to close iroh-docs document: {:?}", e);
}
debug!("KeyValue store close completed");
Ok(())
}
fn address(&self) -> &dyn Address {
self.cached_address.as_ref()
}
fn index(&self) -> Box<dyn StoreIndex<Error = GuardianError> + Send + Sync> {
Box::new(KeyValueIndex {
index: self.index.index.clone(),
})
}
fn store_type(&self) -> &str {
"keyvalue"
}
fn cache(&self) -> Arc<dyn Datastore> {
self.cache.clone()
}
async fn drop(&self) -> Result<()> {
debug!("Starting KeyValue store drop operation (iroh-docs backend)");
self.index.clear_all();
let namespace_id = self.doc_handle.id();
if let Err(e) = self.docs.drop_doc(namespace_id).await {
warn!("Failed to drop iroh-docs document: {:?}", e);
}
if let Err(e) = self.cache.delete(NAMESPACE_CACHE_KEY).await {
warn!("Failed to remove namespace from cache: {:?}", e);
}
debug!("KeyValue store drop completed");
Ok(())
}
async fn load(&self, _amount: usize) -> Result<()> {
self.sync_index_from_docs().await?;
Ok(())
}
async fn sync(&self, _heads: Vec<crate::log::entry::Entry>) -> Result<()> {
self.sync_index_from_docs().await?;
Ok(())
}
async fn load_more_from(&self, _amount: u64, _entries: Vec<crate::log::entry::Entry>) {
}
async fn load_from_snapshot(&self) -> Result<()> {
self.sync_index_from_docs().await?;
Ok(())
}
fn op_log(&self) -> Arc<RwLock<crate::log::Log>> {
self.empty_log.clone()
}
fn client(&self) -> Arc<crate::p2p::network::client::IrohClient> {
self.client.clone()
}
fn db_name(&self) -> &str {
&self.db_name
}
fn identity(&self) -> &Identity {
&self.identity
}
fn access_controller(&self) -> &dyn crate::access_control::traits::AccessController {
self.access_controller.as_ref()
}
async fn add_operation(
&self,
op: Operation,
_on_progress_callback: Option<tokio::sync::mpsc::Sender<crate::log::entry::Entry>>,
) -> Result<crate::log::entry::Entry> {
self.ensure_writable()?;
let key = op.key().cloned().unwrap_or_default();
match op.op() {
"PUT" => {
let value = op.value().to_vec();
self.docs
.set_bytes(
&self.doc_handle,
self.author_id,
Bytes::from(key.clone().into_bytes()),
Bytes::from(value.clone()),
)
.await?;
self.index.insert(key, value);
}
"DEL" => {
self.docs
.del(
&self.doc_handle,
self.author_id,
Bytes::from(key.clone().into_bytes()),
)
.await?;
self.index.remove(&key);
}
other => {
return Err(GuardianError::Store(format!(
"Unknown operation: {}",
other
)));
}
}
let payload = crate::guardian::serializer::serialize(&op).unwrap_or_default();
let clock = LamportClock::new(self.identity.pub_key());
let entry_arc = crate::log::entry::Entry::create(
&self.client,
(*self.identity).clone(),
"",
&payload,
&[],
Some(clock),
);
let entry = (*entry_arc).clone();
Ok(entry)
}
fn span(&self) -> Arc<tracing::Span> {
Arc::new(self.span.clone())
}
fn tracer(&self) -> Arc<TracerWrapper> {
self.tracer.clone()
}
fn event_bus(&self) -> Arc<EventBus> {
self.event_bus.clone()
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}
impl GuardianDBKeyValue {
pub fn span(&self) -> &Span {
&self.span
}
pub fn namespace_id(&self) -> iroh_docs::NamespaceId {
self.doc_handle.id()
}
pub fn author_id(&self) -> AuthorId {
self.author_id
}
}
#[async_trait::async_trait]
impl KeyValueStore for GuardianDBKeyValue {
fn all(&self) -> HashMap<String, Vec<u8>> {
self.index.get_all()
}
async fn put(&self, key: &str, value: Vec<u8>) -> Result<Operation> {
self.put_impl(key, value).await
}
async fn delete(&self, key: &str) -> Result<Operation> {
self.delete_impl(key).await
}
async fn get(&self, key: &str) -> Result<Option<Vec<u8>>> {
self.get_impl(key).await
}
async fn share_ticket(&self) -> Result<String> {
let ticket = self.docs.share_doc(&self.doc_handle, true).await?;
Ok(ticket.to_string())
}
}
impl GuardianDBKeyValue {
pub async fn share_tickets(&self) -> Result<(String, String)> {
let read_ticket = self.docs.share_doc(&self.doc_handle, false).await?;
let write_ticket = self.docs.share_doc(&self.doc_handle, true).await?;
Ok((read_ticket.to_string(), write_ticket.to_string()))
}
pub fn is_writable(&self) -> bool {
self.writable
}
fn ensure_writable(&self) -> Result<()> {
if self.writable {
Ok(())
} else {
Err(GuardianError::Store(
"store is read-only: this replica cannot originate writes".to_string(),
))
}
}
async fn persist_writable(cache: &dyn Datastore, writable: bool) {
if let Err(e) = cache.put(WRITABLE_CACHE_KEY, &[writable as u8]).await {
warn!("Failed to persist writability flag: {:?}", e);
}
}
async fn load_writable(cache: &dyn Datastore) -> bool {
match cache.get(WRITABLE_CACHE_KEY).await {
Ok(Some(bytes)) if !bytes.is_empty() => bytes[0] != 0,
_ => true,
}
}
pub fn len(&self) -> usize {
self.index.len()
}
pub fn is_empty(&self) -> bool {
self.index.is_empty()
}
pub fn contains_key(&self, key: &str) -> bool {
self.index.get_value(key).is_some()
}
pub fn keys(&self) -> Vec<String> {
self.index.get_all().keys().cloned().collect()
}
pub fn all(&self) -> HashMap<String, Vec<u8>> {
self.index.get_all()
}
pub async fn sync_index_from_docs(&self) -> Result<usize> {
refresh_kv_index(&self.docs, &self.doc_handle, &self.client, &self.index).await
}
fn spawn_live_index_sync(&self) {
let docs = self.docs.clone();
let doc = self.doc_handle.clone();
let client = self.client.clone();
let index = self.index.clone();
tokio::spawn(async move {
let mut stream = match doc.subscribe().await {
Ok(s) => s,
Err(e) => {
warn!("Failed to subscribe to iroh-docs doc events (KV): {:?}", e);
return;
}
};
use futures::StreamExt;
use iroh_docs::engine::LiveEvent;
while let Some(event) = stream.next().await {
let is_remote = matches!(
event,
Ok(LiveEvent::InsertRemote { .. })
| Ok(LiveEvent::ContentReady { .. })
| Ok(LiveEvent::PendingContentReady)
| Ok(LiveEvent::SyncFinished(_))
);
if is_remote && let Err(e) = refresh_kv_index(&docs, &doc, &client, &index).await {
warn!("Failed to update KV index via live sync: {:?}", e);
}
}
debug!("Live index sync terminated for KV store");
});
}
#[instrument(level = "debug", skip(self, value))]
pub async fn put_impl(&self, key: &str, value: Vec<u8>) -> Result<Operation> {
self.ensure_writable()?;
if key.is_empty() {
return Err(GuardianError::Store("The key cannot be empty".to_string()));
}
if value.is_empty() {
return Err(GuardianError::Store(
"The value cannot be empty".to_string(),
));
}
self.docs
.set_bytes(
&self.doc_handle,
self.author_id,
Bytes::from(key.as_bytes().to_vec()),
Bytes::from(value.clone()),
)
.await
.map_err(|e| {
GuardianError::Store(format!("Error writing key '{}' to iroh-docs: {}", key, e))
})?;
self.index.insert(key.to_string(), value.clone());
debug!("PUT key='{}' ({} bytes) via iroh-docs", key, value.len());
Ok(Operation::new(
Some(key.to_string()),
"PUT".to_string(),
Some(value),
))
}
#[instrument(level = "debug", skip(self))]
pub async fn delete_impl(&self, key: &str) -> Result<Operation> {
self.ensure_writable()?;
if key.is_empty() {
return Err(GuardianError::Store("The key cannot be empty".to_string()));
}
if !self.contains_key(key) {
return Err(GuardianError::Store(format!("Key '{}' not found", key)));
}
let deleted = self
.docs
.del(
&self.doc_handle,
self.author_id,
Bytes::from(key.as_bytes().to_vec()),
)
.await
.map_err(|e| {
GuardianError::Store(format!("Error deleting key '{}' in iroh-docs: {}", key, e))
})?;
self.index.remove(key);
debug!(
"DEL key='{}' ({} entries removed) via iroh-docs",
key, deleted
);
Ok(Operation::new(
Some(key.to_string()),
"DEL".to_string(),
None,
))
}
#[instrument(level = "debug", skip(self))]
pub async fn get_impl(&self, key: &str) -> Result<Option<Vec<u8>>> {
if key.is_empty() {
return Err(GuardianError::Store("The key cannot be empty".to_string()));
}
Ok(self.index.get_value(key))
}
pub fn get_type(&self) -> &'static str {
"keyvalue"
}
#[instrument(level = "debug", skip(client, identity, addr, options))]
pub async fn new(
client: Arc<crate::p2p::network::client::IrohClient>,
identity: Arc<Identity>,
addr: Arc<dyn Address + Send + Sync>,
options: Option<NewStoreOptions>,
) -> Result<Self> {
let opts = options.unwrap_or_default();
if !client.has_docs_client().await {
client.init_docs().await.map_err(|e| {
GuardianError::Store(format!("Failed to initialize iroh-docs: {}", e))
})?;
}
let mut docs = client.docs_client().await.ok_or_else(|| {
GuardianError::Store("iroh-docs not available after initialization".to_string())
})?;
let author_id = docs.get_or_init_author().await.map_err(|e| {
GuardianError::Store(format!("Failed to initialize iroh-docs author: {}", e))
})?;
let db_name = addr.get_path().to_string();
let span = tracing::info_span!("keyvalue_store", address = %addr.to_string());
let event_bus = opts.event_bus.unwrap_or_default();
let event_bus = Arc::new(event_bus);
let access_controller = opts.access_controller.unwrap_or_else(|| {
let mut default_access = HashMap::new();
default_access.insert("write".to_string(), vec!["*".to_string()]);
Arc::new(SimpleAccessController::new(Some(default_access))) as Arc<dyn AccessController>
});
let access_controller_for_registry = access_controller.clone();
let tracer = opts.tracer.unwrap_or_else(|| {
Arc::new(TracerWrapper::Noop(
NoopTracerProvider::new().tracer("berty.guardian-db"),
))
});
let cache: Arc<dyn Datastore> = if let Some(cache) = opts.cache {
cache
} else {
let cache_dir = if opts.directory.is_empty() {
format!("./GuardianDB/{}/cache", addr)
} else {
format!("{}/cache", opts.directory)
};
Self::create_cache(addr.as_ref(), &cache_dir)?
};
let emitter_interface: Arc<dyn crate::events::EmitterInterface + Send + Sync> =
Arc::new(EventEmitter::new());
let store_key = addr
.to_string()
.rsplit('/')
.next()
.unwrap_or_default()
.to_string();
let requested_read_only = opts.read_only.unwrap_or(false);
let resolved_ticket: Option<String> = match opts.doc_ticket.clone() {
Some(t) => Some(t),
None => client.backend().resolve_shared_ticket(&store_key).await,
};
let (doc_handle, doc_is_writable) = if let Some(ticket_str) = resolved_ticket.as_ref() {
let ticket = ticket_str
.parse::<iroh_docs::DocTicket>()
.map_err(|e| GuardianError::Store(format!("Invalid DocTicket: {}", e)))?;
let ticket_writable = matches!(ticket.capability, Capability::Write(_));
let doc = docs.import_doc(ticket).await?;
let ns_id = doc.id();
cache
.put(NAMESPACE_CACHE_KEY, ns_id.as_bytes())
.await
.map_err(|e| {
GuardianError::Store(format!("Failed to persist imported NamespaceId: {}", e))
})?;
Self::persist_writable(cache.as_ref(), ticket_writable).await;
info!(
writable = ticket_writable,
"Imported shared iroh-docs document via ticket: {:?}", ns_id
);
(doc, ticket_writable)
} else {
match cache.get(NAMESPACE_CACHE_KEY).await {
Ok(Some(namespace_bytes)) if namespace_bytes.len() == 32 => {
let mut ns_bytes = [0u8; 32];
ns_bytes.copy_from_slice(&namespace_bytes);
let namespace_id = iroh_docs::NamespaceId::from(ns_bytes);
match docs.open_doc(namespace_id).await? {
Some(doc) => {
let writable = Self::load_writable(cache.as_ref()).await;
info!(
writable,
"Reopened existing iroh-docs document: {:?}", namespace_id
);
(doc, writable)
}
None if requested_read_only => {
return Err(GuardianError::Store(format!(
"Read-only store '{}' cannot create a namespace and the cached \
namespace {:?} was not found; no ticket available to import",
store_key, namespace_id
)));
}
None => {
warn!(
"Cached namespace {:?} not found, creating new document",
namespace_id
);
let doc = docs.create_doc().await?;
let ns_id = doc.id();
cache
.put(NAMESPACE_CACHE_KEY, ns_id.as_bytes())
.await
.map_err(|e| {
GuardianError::Store(format!(
"Failed to persist NamespaceId: {}",
e
))
})?;
Self::persist_writable(cache.as_ref(), true).await;
info!("Created new iroh-docs document: {:?}", ns_id);
(doc, true)
}
}
}
_ if requested_read_only => {
return Err(GuardianError::Store(format!(
"Read-only store '{}' cannot create a namespace and none was available \
to import (no ticket, no cached namespace)",
store_key
)));
}
_ => {
let doc = docs.create_doc().await?;
let ns_id = doc.id();
cache
.put(NAMESPACE_CACHE_KEY, ns_id.as_bytes())
.await
.map_err(|e| {
GuardianError::Store(format!("Failed to persist NamespaceId: {}", e))
})?;
Self::persist_writable(cache.as_ref(), true).await;
info!("Created new iroh-docs document: {:?}", ns_id);
(doc, true)
}
}
};
let writable = doc_is_writable && !requested_read_only;
let empty_log = {
use crate::log::{AdHocAccess, Log, LogOptions};
let log_opts = LogOptions {
id: Some(&db_name),
access: AdHocAccess,
entries: &[],
heads: &[],
clock: None,
sort_fn: None,
};
Arc::new(RwLock::new(Log::new(
client.clone(),
(*identity).clone(),
log_opts,
)))
};
let index = Arc::new(KeyValueIndex::new());
let cached_address = addr.clone();
let store = GuardianDBKeyValue {
docs,
doc_handle,
author_id,
access_controller,
event_bus,
client,
identity,
cached_address,
db_name,
cache,
index,
span,
tracer,
emitter_interface,
empty_log,
writable,
};
match store.sync_index_from_docs().await {
Ok(count) => {
if count > 0 {
info!(
"KeyValue store initialized with {} entries from iroh-docs",
count
);
}
}
Err(e) => {
warn!(
"Failed to sync index on initialization: {:?}. Store will start empty.",
e
);
}
}
store.spawn_live_index_sync();
match store.share_tickets().await {
Ok((read_ticket, write_ticket)) => {
store
.client
.backend()
.register_ticket_provider(
store_key,
read_ticket,
write_ticket,
access_controller_for_registry,
)
.await;
}
Err(e) => {
warn!(
"Failed to generate share tickets, store not registered for exchange: {:?}",
e
);
}
}
info!(
"GuardianDBKeyValue initialized with iroh-docs backend (namespace={:?}, author={:?})",
store.doc_handle.id(),
store.author_id
);
Ok(store)
}
fn create_cache(address: &dyn Address, cache_dir: &str) -> Result<Arc<dyn Datastore>> {
use crate::cache::level_down::LevelDownCache;
use crate::cache::{Cache, CacheMode, Options};
let cache_options = Options {
span: None,
max_cache_size: Some(100 * 1024 * 1024),
cache_mode: CacheMode::Auto,
};
let cache_manager = LevelDownCache::new(Some(&cache_options));
let address_string = address.to_string();
let parsed_address = crate::address::parse(&address_string)
.map_err(|e| GuardianError::Store(format!("Failed to parse address: {}", e)))?;
let boxed_datastore = cache_manager
.load(cache_dir, &parsed_address)
.map_err(|e| GuardianError::Store(format!("Failed to create cache: {}", e)))?;
struct DatastoreWrapper {
inner: Box<dyn Datastore + Send + Sync>,
}
#[async_trait::async_trait]
impl Datastore for DatastoreWrapper {
async fn get(&self, key: &[u8]) -> crate::guardian::error::Result<Option<Vec<u8>>> {
self.inner.get(key).await
}
async fn put(&self, key: &[u8], value: &[u8]) -> crate::guardian::error::Result<()> {
self.inner.put(key, value).await
}
async fn has(&self, key: &[u8]) -> crate::guardian::error::Result<bool> {
self.inner.has(key).await
}
async fn delete(&self, key: &[u8]) -> crate::guardian::error::Result<()> {
self.inner.delete(key).await
}
async fn query(
&self,
query: &crate::data_store::Query,
) -> crate::guardian::error::Result<crate::data_store::Results> {
self.inner.query(query).await
}
async fn list_keys(
&self,
prefix: &[u8],
) -> crate::guardian::error::Result<Vec<crate::data_store::Key>> {
self.inner.list_keys(prefix).await
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
}
Ok(Arc::new(DatastoreWrapper {
inner: boxed_datastore,
}))
}
}