lfsx_server/storage/
sweep.rs1use std::collections::HashSet;
2use std::path::Path;
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;
12
13#[derive(Debug, Default, Serialize)]
14pub struct SweepReport {
15 pub swept: usize,
16 pub bytes: u64,
17 pub within_grace: usize,
18 pub dry_run: bool,
19}
20
21impl LocalStore {
22 pub async fn sweep(
23 &self,
24 ns: &Namespace,
25 retained: &HashSet<String>,
26 grace: Duration,
27 dry_run: bool,
28 ) -> Result<SweepReport, Error> {
29 let mut report = SweepReport {
30 dry_run,
31 ..SweepReport::default()
32 };
33
34 let Ok(mut prefixes) = fs::read_dir(self.root.join(ns.org()).join(ns.repo())).await else {
35 return Ok(report);
36 };
37
38 while let Some(prefix) = prefixes.next_entry().await? {
39 let Ok(mut fanouts) = fs::read_dir(prefix.path()).await else {
40 continue;
41 };
42
43 while let Some(fanout) = fanouts.next_entry().await? {
44 self.sweep_directory(&fanout.path(), retained, grace, &mut report)
45 .await?;
46 }
47
48 if !dry_run {
49 let _ = fs::remove_dir(prefix.path()).await;
50 }
51 }
52
53 Ok(report)
54 }
55
56 async fn sweep_directory(
57 &self,
58 directory: &Path,
59 retained: &HashSet<String>,
60 grace: Duration,
61 report: &mut SweepReport,
62 ) -> Result<(), Error> {
63 let Ok(mut entries) = fs::read_dir(directory).await else {
64 return Ok(());
65 };
66
67 while let Some(entry) = entries.next_entry().await? {
68 let name = entry.file_name().to_string_lossy().into_owned();
69 if Self::validate_oid(&name).is_err() || retained.contains(&name) {
70 continue;
71 }
72
73 let metadata = entry.metadata().await?;
74 if age(&metadata) < grace {
75 report.within_grace += 1;
76 continue;
77 }
78
79 report.swept += 1;
80 report.bytes += metadata.len();
81
82 if !report.dry_run {
83 fs::remove_file(entry.path()).await?;
84 }
85 }
86
87 if !report.dry_run {
88 let _ = fs::remove_dir(directory).await;
89 }
90
91 Ok(())
92 }
93}
94
95fn age(metadata: &std::fs::Metadata) -> Duration {
96 metadata
97 .modified()
98 .ok()
99 .and_then(|modified| SystemTime::now().duration_since(modified).ok())
100 .unwrap_or_default()
101}
102
103const USAGE_TTL: Duration = Duration::from_secs(60);
104
105impl LocalStore {
106 pub async fn usage(&self) -> (u64, u64) {
107 let mut cached = self.usage.lock().await;
108
109 if let Some((measured_at, objects, bytes)) = *cached
110 && measured_at.elapsed() < USAGE_TTL
111 {
112 return (objects, bytes);
113 }
114
115 let measured = self.measure().await;
116 *cached = Some((Instant::now(), measured.0, measured.1));
117
118 measured
119 }
120
121 async fn measure(&self) -> (u64, u64) {
122 self.scans.fetch_add(1, Ordering::Relaxed);
123
124 let mut objects = 0;
125 let mut bytes = 0;
126
127 let mut directories = vec![self.root.clone()];
128 while let Some(directory) = directories.pop() {
129 let Ok(mut entries) = fs::read_dir(&directory).await else {
130 continue;
131 };
132
133 while let Ok(Some(entry)) = entries.next_entry().await {
134 let name = entry.file_name().to_string_lossy().into_owned();
135 if name.starts_with('.') {
136 continue;
137 }
138
139 match entry.metadata().await {
140 Ok(metadata) if metadata.is_dir() => directories.push(entry.path()),
141 Ok(metadata) if LocalStore::validate_oid(&name).is_ok() => {
142 objects += 1;
143 bytes += metadata.len();
144 }
145 _ => {}
146 }
147 }
148 }
149
150 (objects, bytes)
151 }
152}