Skip to main content

millipede_storage_memory/
client.rs

1use crate::{MemoryDataset, MemoryKeyValueStore, MemoryRequestQueue};
2use millipede_core::storage::{Dataset, KeyValueStore, RequestQueue, StorageClient, StorageResult};
3use std::{
4    collections::HashMap,
5    sync::{Arc, Mutex},
6};
7
8/// An in-process storage client that shares named stores across open calls.
9pub 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    /// Creates an empty storage client.
17    #[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    /// Empties existing handles, then detaches all datasets, stores, and queues.
66    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}