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.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.org()).join(ns.repo())).await
233 }
234
235 pub(super) async fn forget_capacity(&self) {
240 *self.usage.lock().await = None;
241 }
242
243 async fn measure(&self) -> (u64, u64) {
249 #[cfg(unix)]
250 {
251 self.walk_unique(self.root.clone()).await
252 }
253 #[cfg(not(unix))]
254 {
255 let (shared_objects, shared_bytes) = self.walk(self.root.join(".content")).await;
256 let (loose_objects, loose_bytes) = self.walk_unshared(self.root.clone()).await;
257
258 (shared_objects + loose_objects, shared_bytes + loose_bytes)
259 }
260 }
261
262 #[cfg(unix)]
266 async fn walk_unique(&self, from: PathBuf) -> (u64, u64) {
267 use std::collections::HashSet;
268 use std::os::unix::fs::MetadataExt;
269
270 self.scan(from, |metadata, seen: &mut HashSet<(u64, u64)>| {
271 seen.insert((metadata.dev(), metadata.ino()))
272 })
273 .await
274 }
275
276 #[cfg(not(unix))]
280 async fn walk_unshared(&self, from: PathBuf) -> (u64, u64) {
281 let mut objects = 0;
282 let mut bytes = 0;
283 let mut directories = vec![from];
284
285 while let Some(directory) = directories.pop() {
286 let Ok(mut entries) = fs::read_dir(&directory).await else {
287 continue;
288 };
289
290 while let Ok(Some(entry)) = entries.next_entry().await {
291 let name = entry.file_name().to_string_lossy().into_owned();
292 if name.starts_with('.') && directory == self.root {
296 continue;
297 }
298
299 match entry.metadata().await {
300 Ok(metadata) if metadata.is_dir() => directories.push(entry.path()),
301 Ok(metadata)
302 if let Ok(oid) = Oid::parse(&name)
303 && fs::metadata(self.content_path(&oid)).await.is_err() =>
304 {
305 objects += 1;
306 bytes += metadata.len();
307 }
308 _ => {}
309 }
310 }
311 }
312
313 (objects, bytes)
314 }
315
316 async fn walk(&self, from: PathBuf) -> (u64, u64) {
317 self.scan(from, |_, _: &mut ()| true).await
318 }
319
320 async fn scan<S: Default>(
321 &self,
322 from: PathBuf,
323 mut counts: impl FnMut(&std::fs::Metadata, &mut S) -> bool,
324 ) -> (u64, u64) {
325 self.scans.fetch_add(1, Ordering::Relaxed);
326
327 let mut objects = 0;
328 let mut bytes = 0;
329 let mut state = S::default();
330
331 let mut directories = vec![from];
332 while let Some(directory) = directories.pop() {
333 let Ok(mut entries) = fs::read_dir(&directory).await else {
334 continue;
335 };
336
337 while let Ok(Some(entry)) = entries.next_entry().await {
338 let name = entry.file_name().to_string_lossy().into_owned();
339 if name.starts_with('.') {
340 continue;
341 }
342
343 match entry.metadata().await {
344 Ok(metadata) if metadata.is_dir() => directories.push(entry.path()),
345 Ok(metadata) if Oid::parse(&name).is_ok() && counts(&metadata, &mut state) => {
346 objects += 1;
347 bytes += metadata.len();
348 }
349 _ => {}
350 }
351 }
352 }
353
354 (objects, bytes)
355 }
356}