Skip to main content

lfsx_server/storage/
rewrite.rs

1use std::path::Path;
2
3use serde::Serialize;
4use sha2::{Digest, Sha256};
5use tokio::fs;
6use tokio::io::AsyncReadExt;
7
8use super::{LocalStore, shares_bytes_with};
9use super::{codec, crypt};
10use crate::error::Error;
11use crate::namespace::Namespace;
12use crate::oid::Oid;
13
14#[derive(Debug, Default, Serialize, PartialEq, Eq)]
15pub struct CompressReport {
16    pub inspected: u64,
17    pub compressed: u64,
18    pub already: u64,
19    pub left_alone: u64,
20    pub refused: u64,
21    pub before: u64,
22    pub after: u64,
23    pub incomplete: bool,
24    pub dry_run: bool,
25}
26
27impl LocalStore {
28    // Turning compression on only changes what arrives next. This is how a store
29    // that predates it stops paying full price for what it already holds.
30    pub async fn compress(&self, ns: &Namespace, dry_run: bool) -> Result<CompressReport, Error> {
31        let level = self.compression.ok_or(Error::CompressionDisabled)?;
32
33        let walk = self.objects_of(ns).await;
34        let mut report = CompressReport {
35            dry_run,
36            incomplete: !walk.complete,
37            ..CompressReport::default()
38        };
39
40        for found in walk.objects {
41            report.inspected += 1;
42            let on_disk = fs::metadata(&found.path).await?.len();
43            report.before += on_disk;
44
45            if self.is_framed(&found.path, &found.oid, on_disk).await? {
46                report.already += 1;
47                report.after += on_disk;
48                continue;
49            }
50
51            report.after += self
52                .compress_object(&found.path, &found.oid, on_disk, level, &mut report)
53                .await?;
54        }
55
56        if !dry_run && report.compressed > 0 {
57            self.forget_capacity().await;
58        }
59
60        Ok(report)
61    }
62
63    async fn is_framed(&self, path: &Path, oid: &Oid, on_disk: u64) -> Result<bool, Error> {
64        let file = fs::File::open(path).await?;
65
66        Ok(codec::Framed::open(
67            codec::Reader::File(file),
68            on_disk,
69            self.keys.as_deref(),
70            oid,
71        )
72        .await?
73        .is_some())
74    }
75
76    async fn compress_object(
77        &self,
78        path: &Path,
79        oid: &Oid,
80        on_disk: u64,
81        level: i32,
82        report: &mut CompressReport,
83    ) -> Result<u64, Error> {
84        let parent = path.parent().expect("objects live in a fanout directory");
85        let staged = self.staging_path(parent, oid);
86
87        let (digest, compressed) = match self.rewrite(path, &staged, oid, level).await {
88            Ok(measured) => measured,
89            Err(error) => {
90                let _ = fs::remove_file(&staged).await;
91                return Err(error);
92            }
93        };
94
95        // The name is the digest, so this is the same check an operator would
96        // run by hand, and the last chance to run it, since afterwards the file
97        // is no longer the bytes it is named after.
98        if digest != oid.as_str() {
99            tracing::warn!(%oid, %digest, "object does not hash to its own name, leaving it alone");
100            let _ = fs::remove_file(&staged).await;
101            report.refused += 1;
102            return Ok(on_disk);
103        }
104
105        // An object that will not compress costs a header and an index to store
106        // this way. Leaving it is not a failure, it is the right answer.
107        if compressed >= on_disk || report.dry_run {
108            let _ = fs::remove_file(&staged).await;
109            if compressed >= on_disk {
110                report.left_alone += 1;
111                return Ok(on_disk);
112            }
113
114            report.compressed += 1;
115            return Ok(compressed);
116        }
117
118        self.swap_in(path, &staged, oid).await?;
119        report.compressed += 1;
120
121        Ok(compressed)
122    }
123
124    // Every repository holding these bytes has a link to one file, so replacing
125    // this repository's link with a compressed copy would break the sharing the
126    // deduplication just built. The shared copy is what gets replaced, and this
127    // repository is relinked to it. Repositories that have not run this yet keep
128    // the old bytes alive through their own links until their turn.
129    async fn swap_in(&self, path: &Path, staged: &Path, oid: &Oid) -> Result<(), Error> {
130        let content = self.content_path(oid);
131
132        if !shares_bytes_with(path, &content).await {
133            return Ok(fs::rename(staged, path).await?);
134        }
135
136        fs::rename(staged, &content).await?;
137
138        let parent = path.parent().expect("objects live in a fanout directory");
139        let relink = self.staging_path(parent, oid);
140        self.link(&content, &relink).await?;
141
142        Ok(fs::rename(&relink, path).await?)
143    }
144
145    async fn rewrite(
146        &self,
147        path: &Path,
148        staged: &Path,
149        oid: &Oid,
150        level: i32,
151    ) -> Result<(String, u64), Error> {
152        let mut source = fs::File::open(path).await?;
153        let mut writer = codec::Writer::open(
154            fs::File::create(staged).await?,
155            Some(level),
156            self.keys.as_deref().map(crypt::Keyring::writing),
157            oid,
158        )
159        .await?;
160        let mut hasher = Sha256::new();
161        let mut buffer = vec![0u8; 1024 * 1024];
162
163        loop {
164            let read = source.read(&mut buffer).await?;
165            if read == 0 {
166                break;
167            }
168
169            hasher.update(&buffer[..read]);
170            writer.push(&buffer[..read]).await?;
171        }
172
173        writer.finish().await?;
174
175        Ok((
176            hex::encode(hasher.finalize()),
177            fs::metadata(staged).await?.len(),
178        ))
179    }
180}
181
182#[cfg(test)]
183mod tests;