Skip to main content

lfsx_server/storage/
sweep.rs

1use 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    // Objects go, the fanout directories they lived in stay. Removing an
25    // emptied directory raced every upload: a push creates its fanout, and for
26    // the moment between that and the staging file appearing the directory is
27    // empty, so a collection running alongside took it and the push failed on a
28    // directory that had just been made for it. Nothing here can hold a lock the
29    // filesystem would honour, so the fix is to stop competing: the shared
30    // .content tree has never pruned its directories either. What is left is an
31    // inode and a block per prefix, reused by the next object that hashes into
32    // it, against a push failing for a reason no operator could act on.
33    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        // Resolved once. The set of repositories does not change while a sweep
48        // runs, and listing every organisation again for each object made
49        // collection cost objects times repositories: invisible on a small store,
50        // and the whole runtime on the one where an operator finally needs it.
51        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    // The bytes live once under .content, linked from each repository that holds
82    // them. Dropping this repository's link frees nothing until the last one
83    // goes, so only then is it counted as freed, and a dry run that counted
84    // otherwise would promise space it cannot deliver.
85    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            // This repository's link is still there, so it is one of the ones the
97            // count includes.
98            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        // Count the bytes only if this call is the one that removed them. Two
118        // repositories dropping their last reference at the same time would
119        // otherwise each claim the same space, and two reports would add up to
120        // more than the disk ever held.
121        if fs::remove_file(&content).await.is_ok() {
122            report.bytes += size;
123        }
124
125        Ok(())
126    }
127
128    // Is this object still linked from a repository other than the one being
129    // swept? `ours` is how many of the links belong to the repository being
130    // swept: one before its link is removed, none after.
131    //
132    // The filesystem already keeps this count, so on Unix it is one stat rather
133    // than a walk of the store: the shared entry under .content, plus one link
134    // per repository holding the object.
135    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        // No shared entry to count, which is what an object written before the
141        // shared store looks like, or a platform without link counts. Probing the
142        // repositories resolved at the start of the sweep is what is left.
143        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    // Every repository in the store except the one being swept, as directories.
155    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// How many names the bytes have. None where the filesystem cannot say, which is
186// every non-Unix target and any object with no shared entry to count from.
187#[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    // Uncached on purpose: what a repository holds right now. Remembering it is
229    // the seam's job, because a bucket needs exactly the same policy and used
230    // to go without.
231    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    // The whole-store figure, because it just became wrong. Freeing gigabytes
237    // and then reporting the old total for another minute is how an operator
238    // concludes the collection (or the migration) did nothing. What the
239    // repository holds is remembered a layer up and dropped there.
240    pub(super) async fn forget_capacity(&self) {
241        *self.usage.lock().await = None;
242    }
243
244    // What the disk actually holds, which is not the sum of what the
245    // repositories logically hold: an object shared by three projects is three
246    // links to one set of bytes. Counting per-repository paths would report the
247    // pre-deduplication total and grow every time another project links the
248    // same pack: the opposite of what the number is for.
249    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    // Every hard link to one object reports the same inode, so counting each
264    // inode once measures bytes on disk exactly, including the copy fallback,
265    // which really does duplicate them and really should be counted twice.
266    #[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    // Without inode numbers, count the shared store plus anything a repository
278    // holds that has no counterpart there. A copy made by the fallback path is
279    // undercounted, which needs a filesystem with no hard links to happen at all.
280    #[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                // Only the root carries dot-directories worth skipping: .content
294                // is counted through the links that point into it, and .locks
295                // holds no objects at all.
296                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}