Skip to main content

lfsx_server/storage/
dedupe.rs

1use std::path::Path;
2
3use serde::Serialize;
4use sha2::{Digest, Sha256};
5use tokio::fs;
6use tokio::io::AsyncReadExt;
7
8use super::LocalStore;
9use crate::error::Error;
10use crate::namespace::Namespace;
11use crate::oid::Oid;
12
13#[derive(Debug, Default, Serialize, PartialEq, Eq)]
14pub struct DedupeReport {
15    pub inspected: u64,
16    pub already_shared: u64,
17    pub adopted: u64,
18    pub linked: u64,
19    pub reclaimed: u64,
20    pub refused: u64,
21    pub incomplete: bool,
22    pub dry_run: bool,
23}
24
25impl LocalStore {
26    // Objects written before the shared store existed are ordinary files with a
27    // single link. They serve correctly and they never collapse, so a server
28    // that predates deduplication keeps paying full price for every pack two
29    // projects share. This folds them in, one repository at a time.
30    pub async fn dedupe(&self, ns: &Namespace, dry_run: bool) -> Result<DedupeReport, Error> {
31        let walk = self.objects_of(ns).await;
32        let mut report = DedupeReport {
33            dry_run,
34            incomplete: !walk.complete,
35            ..DedupeReport::default()
36        };
37
38        for found in walk.objects {
39            report.inspected += 1;
40            let content = self.content_path(&found.oid);
41
42            if shares_bytes_with(&found.path, &content).await {
43                report.already_shared += 1;
44                continue;
45            }
46
47            match fs::metadata(&content).await {
48                Ok(shared) => {
49                    self.adopt(&found.path, &content, &found.oid, shared.len(), &mut report)
50                        .await?
51                }
52                Err(_) => {
53                    self.promote(&found.path, &content, &found.oid, &mut report)
54                        .await?
55                }
56            }
57        }
58
59        if !dry_run && (report.adopted > 0 || report.linked > 0) {
60            self.forget_capacity().await;
61        }
62
63        Ok(report)
64    }
65
66    // The shared store already holds these bytes. Replacing this repository's
67    // copy with a link to them is what frees the disk, but only after checking
68    // that what is there really is this object: linking to a corrupt entry would
69    // spread it to a repository that had a good copy of its own.
70    async fn adopt(
71        &self,
72        path: &Path,
73        content: &Path,
74        oid: &Oid,
75        size: u64,
76        report: &mut DedupeReport,
77    ) -> Result<(), Error> {
78        if report.dry_run {
79            report.linked += 1;
80            report.reclaimed += size;
81            return Ok(());
82        }
83
84        if !hashes_to(content, oid).await {
85            tracing::warn!(
86                %oid,
87                "shared copy does not hash to its own name, leaving the repository's own file alone"
88            );
89            report.refused += 1;
90            return Ok(());
91        }
92
93        let parent = path.parent().expect("objects live in a fanout directory");
94        let staged = self.staging_path(parent, oid);
95
96        self.link(content, &staged).await?;
97        // Rename over the original rather than removing it first: a crash here
98        // leaves either the old file or the new link, never a gap where the
99        // repository has no object at all.
100        fs::rename(&staged, path).await?;
101
102        report.linked += 1;
103        report.reclaimed += size;
104
105        Ok(())
106    }
107
108    // Nothing shares these bytes yet, so this repository's copy becomes the
109    // shared one and gets a link back in its place. Nothing is freed today; the
110    // next repository to hold the same object is the one that stops paying.
111    async fn promote(
112        &self,
113        path: &Path,
114        content: &Path,
115        oid: &Oid,
116        report: &mut DedupeReport,
117    ) -> Result<(), Error> {
118        if report.dry_run {
119            report.adopted += 1;
120            return Ok(());
121        }
122
123        if !hashes_to(path, oid).await {
124            tracing::warn!(
125                %oid,
126                "object does not hash to its own name, leaving it out of the shared store"
127            );
128            report.refused += 1;
129            return Ok(());
130        }
131
132        let parent = content.parent().expect("content paths have a parent");
133        fs::create_dir_all(parent).await?;
134
135        fs::rename(path, content).await?;
136        if let Err(error) = self.link(content, path).await {
137            // Put it back where the repository expects it rather than leave the
138            // object reachable only from the shared store.
139            fs::rename(content, path).await?;
140            return Err(error.into());
141        }
142
143        report.adopted += 1;
144
145        Ok(())
146    }
147}
148
149// Two paths that resolve to the same inode are already one set of bytes with
150// two names, which is exactly what deduplication produces, so this is how a
151// second run knows there is nothing left to do.
152#[cfg(unix)]
153pub(super) async fn shares_bytes_with(path: &Path, content: &Path) -> bool {
154    use std::os::unix::fs::MetadataExt;
155
156    let (Ok(one), Ok(other)) = (fs::metadata(path).await, fs::metadata(content).await) else {
157        return false;
158    };
159
160    (one.dev(), one.ino()) == (other.dev(), other.ino())
161}
162
163// Without inode numbers there is no way to tell a link from a copy, so every
164// run relinks. The result is the same, the work is repeated. The server ships
165// on Linux; this keeps the tests honest everywhere else.
166#[cfg(not(unix))]
167pub(super) async fn shares_bytes_with(_path: &Path, _content: &Path) -> bool {
168    false
169}
170
171async fn hashes_to(path: &Path, oid: &Oid) -> bool {
172    let Ok(mut file) = fs::File::open(path).await else {
173        return false;
174    };
175
176    let mut hasher = Sha256::new();
177    let mut buffer = vec![0u8; 128 * 1024];
178
179    loop {
180        match file.read(&mut buffer).await {
181            Ok(0) => break,
182            Ok(read) => hasher.update(&buffer[..read]),
183            Err(_) => return false,
184        }
185    }
186
187    hex::encode(hasher.finalize()) == oid.as_str()
188}
189
190#[cfg(test)]
191mod tests;