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