millipede_storage_fs/
client.rs1use crate::{
2 FsDataset, FsKeyValueStore, FsRequestQueue,
3 layout::{DATASETS, KEY_VALUE_STORES, REQUEST_QUEUES, store_name, store_path},
4};
5use millipede_core::storage::{Dataset, KeyValueStore, RequestQueue, StorageClient, StorageResult};
6use std::{collections::HashMap, path::Path, path::PathBuf, sync::Arc};
7use tokio::sync::{Mutex, RwLock};
8
9pub struct FsStorageClient {
11 root: PathBuf,
12 operations: Arc<RwLock<()>>,
13 datasets: Mutex<HashMap<String, Arc<FsDataset>>>,
14 kv_stores: Mutex<HashMap<String, Arc<FsKeyValueStore>>>,
15 request_queues: Mutex<HashMap<String, Arc<FsRequestQueue>>>,
16}
17
18impl FsStorageClient {
19 #[must_use]
21 pub fn new(root: impl Into<PathBuf>) -> Self {
22 Self {
23 root: root.into(),
24 operations: Arc::new(RwLock::new(())),
25 datasets: Mutex::new(HashMap::new()),
26 kv_stores: Mutex::new(HashMap::new()),
27 request_queues: Mutex::new(HashMap::new()),
28 }
29 }
30
31 #[must_use]
33 pub fn root(&self) -> &Path {
34 &self.root
35 }
36}
37
38#[async_trait::async_trait]
39impl StorageClient for FsStorageClient {
40 async fn open_dataset(&self, name: Option<&str>) -> StorageResult<Arc<dyn Dataset>> {
41 let name = store_name(name)?;
42 let _operation = self.operations.read().await;
43 let mut datasets = self.datasets.lock().await;
44 if let Some(dataset) = datasets.get(&name) {
45 tokio::fs::create_dir_all(store_path(&self.root, DATASETS, &name)).await?;
46 return Ok(dataset.clone());
47 }
48
49 let path = store_path(&self.root, DATASETS, &name);
50 tokio::fs::create_dir_all(&path).await?;
51 let dataset =
52 Arc::new(FsDataset::open(name.clone(), path, Arc::clone(&self.operations)).await?);
53 datasets.insert(name, dataset.clone());
54 Ok(dataset)
55 }
56
57 async fn open_key_value_store(
58 &self,
59 name: Option<&str>,
60 ) -> StorageResult<Arc<dyn KeyValueStore>> {
61 let name = store_name(name)?;
62 let _operation = self.operations.read().await;
63 let mut stores = self.kv_stores.lock().await;
64 if let Some(store) = stores.get(&name) {
65 tokio::fs::create_dir_all(store_path(&self.root, KEY_VALUE_STORES, &name)).await?;
66 return Ok(store.clone());
67 }
68
69 let path = store_path(&self.root, KEY_VALUE_STORES, &name);
70 tokio::fs::create_dir_all(&path).await?;
71 let store = Arc::new(FsKeyValueStore::open(
72 name.clone(),
73 path,
74 Arc::clone(&self.operations),
75 ));
76 stores.insert(name, store.clone());
77 Ok(store)
78 }
79
80 async fn open_request_queue(&self, name: Option<&str>) -> StorageResult<Arc<dyn RequestQueue>> {
81 let name = store_name(name)?;
82 let _operation = self.operations.read().await;
83 let mut queues = self.request_queues.lock().await;
84 if let Some(queue) = queues.get(&name) {
85 queue.ensure_layout().await?;
86 return Ok(queue.clone());
87 }
88
89 let path = store_path(&self.root, REQUEST_QUEUES, &name);
90 tokio::fs::create_dir_all(path.join("requests")).await?;
91 let queue =
92 Arc::new(FsRequestQueue::open(name.clone(), path, Arc::clone(&self.operations)).await?);
93 queues.insert(name, queue.clone());
94 Ok(queue)
95 }
96
97 async fn purge(&self) -> StorageResult<()> {
104 let _operation = self.operations.write().await;
105 let mut datasets = self.datasets.lock().await;
106 let mut kv_stores = self.kv_stores.lock().await;
107 let mut request_queues = self.request_queues.lock().await;
108 purge_category(&self.root.join(DATASETS), false).await?;
109 for dataset in datasets.values() {
110 dataset.reset_sequence().await;
111 }
112 purge_category(&self.root.join(KEY_VALUE_STORES), true).await?;
113 purge_category(&self.root.join(REQUEST_QUEUES), false).await?;
114 for queue in request_queues.values() {
115 queue.reset().await;
116 }
117 datasets.clear();
118 kv_stores.clear();
119 request_queues.clear();
120 Ok(())
121 }
122}
123
124async fn purge_category(path: &Path, preserve_input: bool) -> StorageResult<()> {
125 tokio::fs::create_dir_all(path).await?;
126 let mut entries = tokio::fs::read_dir(path).await?;
127 while let Some(entry) = entries.next_entry().await? {
128 let entry_path = entry.path();
129 let file_type = entry.file_type().await?;
130 if preserve_input && file_type.is_dir() && entry.file_name() == "default" {
131 purge_default_kvs(&entry_path).await?;
132 } else if file_type.is_dir() {
133 tokio::fs::remove_dir_all(entry_path).await?;
134 } else {
135 tokio::fs::remove_file(entry_path).await?;
136 }
137 }
138 Ok(())
139}
140
141async fn purge_default_kvs(path: &Path) -> StorageResult<()> {
142 let mut entries = tokio::fs::read_dir(path).await?;
143 while let Some(entry) = entries.next_entry().await? {
144 let file_type = entry.file_type().await?;
145 let name = entry.file_name();
146 let preserve = file_type.is_file()
147 && name
148 .to_str()
149 .and_then(|name| name.rsplit_once('.'))
150 .is_some_and(|(key, extension)| key == "INPUT" && !extension.is_empty());
151 if preserve {
152 continue;
153 }
154 if file_type.is_dir() {
155 tokio::fs::remove_dir_all(entry.path()).await?;
156 } else {
157 tokio::fs::remove_file(entry.path()).await?;
158 }
159 }
160 Ok(())
161}