Skip to main content

snerd_rust/
file_store.rs

1use fs3::FileExt;
2use std::collections::HashMap;
3use std::fs::{File, OpenOptions};
4use std::io::{BufRead, BufReader, Write};
5use std::path::{Path, PathBuf};
6use std::sync::atomic::{AtomicBool, Ordering};
7use std::sync::{Arc, Mutex};
8
9use crate::task::RetryableTask;
10
11#[derive(Clone)]
12pub struct FileStore {
13    file_path: Arc<PathBuf>,
14    total_tasks: Arc<Mutex<usize>>,
15    deleted_tasks: Arc<Mutex<usize>>,
16    append_count: Arc<Mutex<usize>>,
17    compacting: Arc<AtomicBool>,
18}
19
20impl FileStore {
21    pub fn new<P: AsRef<Path>>(path: P) -> std::io::Result<Self> {
22        let fs = FileStore {
23            file_path: Arc::new(path.as_ref().to_path_buf()),
24            total_tasks: Arc::new(Mutex::new(0)),
25            deleted_tasks: Arc::new(Mutex::new(0)),
26            append_count: Arc::new(Mutex::new(0)),
27            compacting: Arc::new(AtomicBool::new(false)),
28        };
29        fs.rebuild_metadata()?;
30        Ok(fs)
31    }
32
33    fn rebuild_metadata(&self) -> std::io::Result<()> {
34        if !self.file_path.exists() {
35            return Ok(());
36        }
37
38        let file = File::open(self.file_path.as_ref())?;
39        file.lock_shared()?;
40
41        let mut total = 0;
42        let mut deleted = 0;
43        let mut appended = 0;
44
45        let reader = BufReader::new(&file);
46        for line_str in reader.lines().map_while(Result::ok) {
47            if line_str.trim().is_empty() {
48                continue;
49            }
50            if let Ok(task) = serde_json::from_str::<RetryableTask>(&line_str) {
51                appended += 1;
52                if task.deleted_at.is_some() {
53                    deleted += 1;
54                } else {
55                    total += 1;
56                }
57            }
58        }
59        file.unlock()?;
60
61        *self.total_tasks.lock().unwrap() = total;
62        *self.deleted_tasks.lock().unwrap() = deleted;
63        *self.append_count.lock().unwrap() = appended;
64
65        Ok(())
66    }
67
68    pub fn save_task(&self, task: &RetryableTask) -> std::io::Result<()> {
69        if let Some(parent) = self.file_path.parent() {
70            std::fs::create_dir_all(parent)?;
71        }
72
73        let mut file = OpenOptions::new()
74            .create(true)
75            .append(true)
76            .open(self.file_path.as_ref())?;
77
78        file.lock_exclusive()?;
79
80        let json_str = serde_json::to_string(task)?;
81        writeln!(file, "{}", json_str)?;
82        file.sync_all()?;
83
84        file.unlock()?;
85
86        let is_deleted = task.deleted_at.is_some();
87        {
88            *self.append_count.lock().unwrap() += 1;
89            if is_deleted {
90                *self.deleted_tasks.lock().unwrap() += 1;
91            } else {
92                *self.total_tasks.lock().unwrap() += 1;
93            }
94        }
95
96        if self.should_compact() {
97            let fs_clone = self.clone();
98            tokio::spawn(async move {
99                let _ = fs_clone.compact_log();
100            });
101        }
102
103        Ok(())
104    }
105
106    pub fn read_tasks(&self) -> std::io::Result<Vec<RetryableTask>> {
107        if !self.file_path.exists() {
108            return Ok(vec![]);
109        }
110
111        let file = File::open(self.file_path.as_ref())?;
112        file.lock_shared()?;
113
114        let mut task_map = HashMap::new();
115        let reader = BufReader::new(&file);
116
117        for line_str in reader.lines().map_while(Result::ok) {
118            if line_str.trim().is_empty() {
119                continue;
120            }
121            if let Ok(task) = serde_json::from_str::<RetryableTask>(&line_str) {
122                if task.deleted_at.is_some() {
123                    task_map.remove(&task.task_id);
124                } else {
125                    task_map.insert(task.task_id.clone(), task);
126                }
127            }
128        }
129
130        file.unlock()?;
131
132        Ok(task_map.into_values().collect())
133    }
134
135    pub fn get_latest_task(&self, task_id: &str) -> std::io::Result<Option<RetryableTask>> {
136        if !self.file_path.exists() {
137            return Ok(None);
138        }
139
140        let file = File::open(self.file_path.as_ref())?;
141        file.lock_shared()?;
142
143        let mut latest = None;
144        let reader = BufReader::new(&file);
145
146        for line_str in reader.lines().map_while(Result::ok) {
147            if line_str.trim().is_empty() {
148                continue;
149            }
150            if let Ok(task) = serde_json::from_str::<RetryableTask>(&line_str)
151                && task.task_id == task_id
152            {
153                latest = Some(task);
154            }
155        }
156
157        file.unlock()?;
158
159        Ok(latest)
160    }
161
162    pub fn delete_task(&self, task_id: &str) -> std::io::Result<()> {
163        if let Some(mut task) = self.get_latest_task(task_id)?
164            && task.deleted_at.is_none()
165        {
166            task.mark_deleted();
167            self.save_task(&task)?;
168        }
169        Ok(())
170    }
171
172    fn should_compact(&self) -> bool {
173        if let Ok(metadata) = std::fs::metadata(self.file_path.as_ref())
174            && metadata.len() > 20 * 1024 * 1024
175        {
176            return true;
177        }
178
179        let total = *self.total_tasks.lock().unwrap();
180        let deleted = *self.deleted_tasks.lock().unwrap();
181        if total > 0 && (deleted as f64 / total as f64) > 0.5 {
182            return true;
183        }
184
185        if *self.append_count.lock().unwrap() >= 10000 {
186            return true;
187        }
188
189        false
190    }
191
192    pub fn compact_log(&self) -> std::io::Result<()> {
193        if self
194            .compacting
195            .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
196            .is_err()
197        {
198            return Ok(()); // Already compacting
199        }
200
201        let temp_path = self.file_path.with_extension("tmp");
202
203        let result = (|| -> std::io::Result<()> {
204            let input_file = File::open(self.file_path.as_ref())?;
205            input_file.lock_shared()?;
206
207            let mut temp_file = OpenOptions::new()
208                .create(true)
209                .write(true)
210                .truncate(true)
211                .open(&temp_path)?;
212
213            let mut task_map = HashMap::new();
214            let mut _total_tasks = 0;
215            let mut _deleted_count = 0;
216
217            let reader = BufReader::new(&input_file);
218            for line_str in reader.lines().map_while(Result::ok) {
219                if line_str.trim().is_empty() {
220                    continue;
221                }
222                _total_tasks += 1;
223                if let Ok(task) = serde_json::from_str::<RetryableTask>(&line_str) {
224                    if task.deleted_at.is_some() {
225                        _deleted_count += 1;
226                        task_map.remove(&task.task_id);
227                    } else if let Some(existing) = task_map.get(&task.task_id) {
228                        let existing_task: &RetryableTask = existing;
229                        if task.updated_at > existing_task.updated_at {
230                            task_map.insert(task.task_id.clone(), task);
231                        }
232                    } else {
233                        task_map.insert(task.task_id.clone(), task);
234                    }
235                }
236            }
237
238            for task in task_map.values() {
239                let json_str = serde_json::to_string(task)?;
240                writeln!(temp_file, "{}", json_str)?;
241            }
242
243            temp_file.sync_all()?;
244            input_file.unlock()?;
245
246            std::fs::rename(&temp_path, self.file_path.as_ref())?;
247
248            *self.total_tasks.lock().unwrap() = task_map.len();
249            *self.deleted_tasks.lock().unwrap() = 0;
250            *self.append_count.lock().unwrap() = 0;
251
252            Ok(())
253        })();
254
255        self.compacting.store(false, Ordering::SeqCst);
256        result
257    }
258}