millipede_storage_memory/
client.rs1use crate::{MemoryDataset, MemoryKeyValueStore, MemoryRequestQueue};
2use millipede_core::storage::{Dataset, KeyValueStore, RequestQueue, StorageClient, StorageResult};
3use std::{
4 collections::HashMap,
5 sync::{Arc, Mutex},
6};
7
8pub struct MemoryStorageClient {
10 datasets: Mutex<HashMap<String, Arc<MemoryDataset>>>,
11 kv_stores: Mutex<HashMap<String, Arc<MemoryKeyValueStore>>>,
12 queues: Mutex<HashMap<String, Arc<dyn RequestQueue>>>,
13}
14
15impl MemoryStorageClient {
16 #[must_use]
18 pub fn new() -> Self {
19 Self {
20 datasets: Mutex::new(HashMap::new()),
21 kv_stores: Mutex::new(HashMap::new()),
22 queues: Mutex::new(HashMap::new()),
23 }
24 }
25}
26
27impl Default for MemoryStorageClient {
28 fn default() -> Self {
29 Self::new()
30 }
31}
32
33#[async_trait::async_trait]
34impl StorageClient for MemoryStorageClient {
35 async fn open_dataset(&self, name: Option<&str>) -> StorageResult<Arc<dyn Dataset>> {
36 let name = name.unwrap_or("default");
37 let mut datasets = self.datasets.lock().expect("datasets mutex poisoned");
38 Ok(datasets
39 .entry(name.to_owned())
40 .or_insert_with(|| Arc::new(MemoryDataset::new(name)))
41 .clone())
42 }
43
44 async fn open_key_value_store(
45 &self,
46 name: Option<&str>,
47 ) -> StorageResult<Arc<dyn KeyValueStore>> {
48 let name = name.unwrap_or("default");
49 let mut stores = self.kv_stores.lock().expect("kv stores mutex poisoned");
50 Ok(stores
51 .entry(name.to_owned())
52 .or_insert_with(|| Arc::new(MemoryKeyValueStore::new(name)))
53 .clone())
54 }
55
56 async fn open_request_queue(&self, name: Option<&str>) -> StorageResult<Arc<dyn RequestQueue>> {
57 let name = name.unwrap_or("default");
58 let mut queues = self.queues.lock().expect("queues mutex poisoned");
59 Ok(queues
60 .entry(name.to_owned())
61 .or_insert_with(|| Arc::new(MemoryRequestQueue::new(name)))
62 .clone())
63 }
64
65 async fn purge(&self) -> StorageResult<()> {
67 let mut datasets = self.datasets.lock().expect("datasets mutex poisoned");
68 for dataset in datasets.values() {
69 dataset.clear();
70 }
71 datasets.clear();
72
73 let mut stores = self.kv_stores.lock().expect("kv stores mutex poisoned");
74 for store in stores.values() {
75 store.clear();
76 }
77 stores.clear();
78 self.queues.lock().expect("queues mutex poisoned").clear();
79 Ok(())
80 }
81}