Skip to main content

snerd_rust/
file_store.rs

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