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