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(()); }
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}