Skip to main content

millipede_storage_fs/
client.rs

1use 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
9/// A file-system storage client using Crawlee-compatible directory layouts.
10pub 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    /// Creates a client rooted at the supplied storage directory.
20    #[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    /// Returns the storage root.
32    #[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    /// Deletes stored contents and clears the opened dataset and store caches.
98    ///
99    /// The default key-value store's `INPUT.<ext>` file is preserved for
100    /// Crawlee migration parity. This matters when `purge_on_start` uses its
101    /// default value of `true`. Previously opened handles remain usable, but
102    /// subsequent opens return fresh instances backed by the purged layout.
103    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}