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