use std::sync::Arc;
use zenoh::{session::Session, Result as ZResult};
use zenoh_backend_traits::config::StorageConfig;
pub use super::replica::{Replica, StorageService};
pub enum StorageMessage {
Stop,
GetStatus(tokio::sync::mpsc::Sender<serde_json::Value>),
}
pub(crate) async fn start_storage(
store_intercept: super::StoreIntercept,
config: StorageConfig,
admin_key: String,
zenoh: Arc<Session>,
) -> ZResult<flume::Sender<StorageMessage>> {
let parts: Vec<&str> = admin_key.split('/').collect();
let uuid = parts[2];
let storage_name = parts[7];
let name = format!("{uuid}/{storage_name}");
tracing::trace!("Start storage '{}' on keyexpr '{}'", name, config.key_expr);
let (tx, rx) = flume::bounded(1);
tokio::task::spawn(async move {
if config.replica_config.is_some() {
Replica::start(zenoh.clone(), store_intercept, config, &name, rx).await;
} else {
StorageService::start(zenoh.clone(), config, &name, store_intercept, rx, None).await;
}
});
Ok(tx)
}