lfsx_server/storage/
sweep.rs1use std::collections::HashSet;
2use std::path::{Path, PathBuf};
3use std::sync::atomic::Ordering;
4use std::time::{Duration, Instant, SystemTime};
5
6use serde::Serialize;
7use tokio::fs;
8
9use super::LocalStore;
10use crate::error::Error;
11use crate::namespace::Namespace;
12use crate::oid::Oid;
13
14#[derive(Debug, Default, Serialize)]
15pub struct SweepReport {
16 pub swept: usize,
17 pub bytes: u64,
18 pub within_grace: usize,
19 pub incomplete: bool,
20 pub dry_run: bool,
21}
22
23impl LocalStore {
24 pub async fn sweep(
34 &self,
35 ns: &Namespace,
36 retained: &HashSet<String>,
37 grace: Duration,
38 dry_run: bool,
39 ) -> Result<SweepReport, Error> {
40 let walk = self.objects_of(ns).await;
41 let mut report = SweepReport {
42 dry_run,
43 incomplete: !walk.complete,
44 ..SweepReport::default()
45 };
46
47 let elsewhere = self.other_repositories(ns).await;
52
53 for found in walk.objects {
54 if retained.contains(found.oid.as_str()) {
55 continue;
56 }
57
58 let metadata = fs::metadata(&found.path).await?;
59 if age(&metadata) < grace {
60 report.within_grace += 1;
61 continue;
62 }
63
64 self.collect(
65 &found.path,
66 &found.oid,
67 metadata.len(),
68 &elsewhere,
69 &mut report,
70 )
71 .await?;
72 }
73
74 if !dry_run {
75 self.forget_capacity().await;
76 }
77
78 Ok(report)
79 }
80
81 async fn collect(
86 &self,
87 path: &Path,
88 oid: &Oid,
89 held: u64,
90 elsewhere: &[PathBuf],
91 report: &mut SweepReport,
92 ) -> Result<(), Error> {
93 report.swept += 1;
94
95 if report.dry_run {
96 if !self.referenced_elsewhere(oid, elsewhere, 1).await {
99 report.bytes += held;
100 }
101
102 return Ok(());
103 }
104
105 fs::remove_file(path).await?;
106
107 if self.referenced_elsewhere(oid, elsewhere, 0).await {
108 return Ok(());
109 }
110
111 let content = self.content_path(oid);
112 let size = fs::metadata(&content)
113 .await
114 .map(|shared| shared.len())
115 .unwrap_or(held);
116
117 if fs::remove_file(&content).await.is_ok() {
122 report.bytes += size;
123 }
124
125 Ok(())
126 }
127
128 async fn referenced_elsewhere(&self, oid: &Oid, elsewhere: &[PathBuf], ours: u64) -> bool {
136 if let Some(links) = links_to(&self.content_path(oid)).await {
137 return links > 1 + ours;
138 }
139
140 for repo in elsewhere {
144 let (first, second) = oid.fanout();
145 let candidate = repo.join(first).join(second).join(oid.as_str());
146 if fs::metadata(candidate).await.is_ok() {
147 return true;
148 }
149 }
150
151 false
152 }
153
154 async fn other_repositories(&self, sweeping: &Namespace) -> Vec<PathBuf> {
156 let mut out = Vec::new();
157 let Ok(mut orgs) = fs::read_dir(&self.root).await else {
158 return out;
159 };
160
161 while let Ok(Some(org)) = orgs.next_entry().await {
162 let org_name = org.file_name().to_string_lossy().into_owned();
163 if org_name.starts_with('.') {
164 continue;
165 }
166
167 let Ok(mut repos) = fs::read_dir(org.path()).await else {
168 continue;
169 };
170
171 while let Ok(Some(repo)) = repos.next_entry().await {
172 let repo_name = repo.file_name().to_string_lossy().into_owned();
173 if org_name == sweeping.stored_org() && repo_name == sweeping.repo() {
174 continue;
175 }
176
177 out.push(repo.path());
178 }
179 }
180
181 out
182 }
183}
184
185#[cfg(unix)]
188async fn links_to(content: &Path) -> Option<u64> {
189 use std::os::unix::fs::MetadataExt;
190
191 fs::metadata(content)
192 .await
193 .ok()
194 .map(|shared| shared.nlink())
195}
196
197#[cfg(not(unix))]
198async fn links_to(_content: &Path) -> Option<u64> {
199 None
200}
201
202pub(super) fn age(metadata: &std::fs::Metadata) -> Duration {
203 metadata
204 .modified()
205 .ok()
206 .and_then(|modified| SystemTime::now().duration_since(modified).ok())
207 .unwrap_or_default()
208}
209
210const USAGE_TTL: Duration = Duration::from_secs(60);
211
212impl LocalStore {
213 pub async fn usage(&self) -> (u64, u64) {
214 let mut cached = self.usage.lock().await;
215
216 if let Some((measured_at, objects, bytes)) = *cached
217 && measured_at.elapsed() < USAGE_TTL
218 {
219 return (objects, bytes);
220 }
221
222 let measured = self.measure().await;
223 *cached = Some((Instant::now(), measured.0, measured.1));
224
225 measured
226 }
227
228 pub(super) async fn measure_of(&self, ns: &Namespace) -> (u64, u64) {
232 self.walk(self.root.join(ns.stored_org()).join(ns.repo()))
233 .await
234 }
235
236 pub(super) async fn forget_capacity(&self) {
241 *self.usage.lock().await = None;
242 }
243
244 async fn measure(&self) -> (u64, u64) {
250 #[cfg(unix)]
251 {
252 self.walk_unique(self.root.clone()).await
253 }
254 #[cfg(not(unix))]
255 {
256 let (shared_objects, shared_bytes) = self.walk(self.root.join(".content")).await;
257 let (loose_objects, loose_bytes) = self.walk_unshared(self.root.clone()).await;
258
259 (shared_objects + loose_objects, shared_bytes + loose_bytes)
260 }
261 }
262
263 #[cfg(unix)]
267 async fn walk_unique(&self, from: PathBuf) -> (u64, u64) {
268 use std::collections::HashSet;
269 use std::os::unix::fs::MetadataExt;
270
271 self.scan(from, |metadata, seen: &mut HashSet<(u64, u64)>| {
272 seen.insert((metadata.dev(), metadata.ino()))
273 })
274 .await
275 }
276
277 #[cfg(not(unix))]
281 async fn walk_unshared(&self, from: PathBuf) -> (u64, u64) {
282 let mut objects = 0;
283 let mut bytes = 0;
284 let mut directories = vec![from];
285
286 while let Some(directory) = directories.pop() {
287 let Ok(mut entries) = fs::read_dir(&directory).await else {
288 continue;
289 };
290
291 while let Ok(Some(entry)) = entries.next_entry().await {
292 let name = entry.file_name().to_string_lossy().into_owned();
293 if name.starts_with('.') && directory == self.root {
297 continue;
298 }
299
300 match entry.metadata().await {
301 Ok(metadata) if metadata.is_dir() => directories.push(entry.path()),
302 Ok(metadata)
303 if let Ok(oid) = Oid::parse(&name)
304 && fs::metadata(self.content_path(&oid)).await.is_err() =>
305 {
306 objects += 1;
307 bytes += metadata.len();
308 }
309 _ => {}
310 }
311 }
312 }
313
314 (objects, bytes)
315 }
316
317 async fn walk(&self, from: PathBuf) -> (u64, u64) {
318 self.scan(from, |_, _: &mut ()| true).await
319 }
320
321 async fn scan<S: Default>(
322 &self,
323 from: PathBuf,
324 mut counts: impl FnMut(&std::fs::Metadata, &mut S) -> bool,
325 ) -> (u64, u64) {
326 self.scans.fetch_add(1, Ordering::Relaxed);
327
328 let mut objects = 0;
329 let mut bytes = 0;
330 let mut state = S::default();
331
332 let mut directories = vec![from];
333 while let Some(directory) = directories.pop() {
334 let Ok(mut entries) = fs::read_dir(&directory).await else {
335 continue;
336 };
337
338 while let Ok(Some(entry)) = entries.next_entry().await {
339 let name = entry.file_name().to_string_lossy().into_owned();
340 if name.starts_with('.') {
341 continue;
342 }
343
344 match entry.metadata().await {
345 Ok(metadata) if metadata.is_dir() => directories.push(entry.path()),
346 Ok(metadata) if Oid::parse(&name).is_ok() && counts(&metadata, &mut state) => {
347 objects += 1;
348 bytes += metadata.len();
349 }
350 _ => {}
351 }
352 }
353 }
354
355 (objects, bytes)
356 }
357}