Skip to main content

lfsx_server/storage/
mod.rs

1mod backend;
2pub(crate) mod cache;
3pub(crate) mod codec;
4pub mod crypt;
5mod dedupe;
6mod rewrite;
7pub mod s3;
8mod staging;
9mod sweep;
10mod usage;
11mod verify;
12mod walk;
13
14#[cfg(test)]
15mod tests;
16
17use std::path::{Path, PathBuf};
18use std::sync::atomic::{AtomicU64, Ordering};
19use std::time::Instant;
20
21use tokio::sync::Mutex;
22
23use futures_util::{Stream, StreamExt};
24use sha2::{Digest, Sha256};
25use tokio::fs;
26use tokio::io::AsyncWriteExt;
27
28use crate::error::Error;
29use crate::namespace::Namespace;
30use crate::oid::Oid;
31
32pub use backend::Store;
33pub use dedupe::DedupeReport;
34use dedupe::shares_bytes_with;
35pub use rewrite::CompressReport;
36pub use verify::VerifyReport;
37
38enum Sink {
39    Raw(fs::File),
40    Framed(Box<codec::Writer>),
41}
42
43impl Sink {
44    async fn write(&mut self, chunk: &[u8]) -> Result<(), Error> {
45        match self {
46            Self::Raw(file) => Ok(file.write_all(chunk).await?),
47            Self::Framed(writer) => writer.push(chunk).await,
48        }
49    }
50
51    async fn finish(self) -> Result<(), Error> {
52        match self {
53            Self::Raw(mut file) => {
54                file.flush().await?;
55                Ok(file.sync_all().await?)
56            }
57            Self::Framed(writer) => writer.finish().await,
58        }
59    }
60}
61
62// A transfer that has passed every check and is waiting to be put somewhere.
63pub struct Staged {
64    pub path: PathBuf,
65    destination: PathBuf,
66    pub written: u64,
67    fresh: bool,
68}
69
70// What a download reads from, whether or not the bytes on disk are the object.
71pub enum Object {
72    Raw {
73        file: fs::File,
74        size: u64,
75    },
76    Framed(codec::Framed),
77    Remote {
78        bucket: s3::S3Store,
79        oid: Oid,
80        size: u64,
81    },
82}
83
84impl Object {
85    pub fn size(&self) -> u64 {
86        match self {
87            Self::Raw { size, .. } => *size,
88            Self::Framed(framed) => framed.plaintext(),
89            Self::Remote { size, .. } => *size,
90        }
91    }
92
93    pub async fn stream(
94        self,
95        start: u64,
96        length: u64,
97    ) -> Result<futures_util::stream::BoxStream<'static, Result<axum::body::Bytes, Error>>, Error>
98    {
99        use futures_util::StreamExt;
100        use tokio::io::AsyncSeekExt;
101
102        match self {
103            Self::Raw { mut file, .. } => {
104                file.seek(std::io::SeekFrom::Start(start)).await?;
105                let reader =
106                    tokio_util::io::ReaderStream::new(tokio::io::AsyncReadExt::take(file, length));
107
108                Ok(reader.map(|chunk| chunk.map_err(Error::from)).boxed())
109            }
110            Self::Framed(framed) => Ok(framed.stream(start, length).boxed()),
111            Self::Remote { bucket, oid, .. } => {
112                let chunks = bucket.read(&oid, start, length).await?;
113
114                Ok(chunks
115                    .map(|chunk| {
116                        chunk.map_err(|error| Error::Storage(std::io::Error::other(error)))
117                    })
118                    .boxed())
119            }
120        }
121    }
122}
123pub use staging::{Reclaimed, reclaim};
124pub use sweep::SweepReport;
125
126// What is left of a repository's budget for one transfer. It travels with the
127// upload because a client that skips negotiation may also skip declaring a
128// size, and a budget checked once against a number the client chose is not a
129// budget.
130#[derive(Debug, Clone, Copy)]
131pub struct Budget {
132    pub used: u64,
133    pub limit: u64,
134}
135
136impl Budget {
137    pub fn exceeded_by(&self, arriving: u64) -> bool {
138        self.used + arriving > self.limit
139    }
140
141    pub fn refusal(&self) -> Error {
142        Error::OverQuota {
143            used: self.used,
144            limit: self.limit,
145        }
146    }
147}
148
149// What a write did, which the caller needs and the byte count alone does not
150// say: an object the repository already held costs it no room, so counting it
151// again would push a repository over a quota it never grew into.
152#[derive(Debug)]
153pub struct Written {
154    pub bytes: u64,
155    pub fresh: bool,
156}
157
158pub struct LocalStore {
159    root: PathBuf,
160    counter: AtomicU64,
161    usage: Mutex<Option<(Instant, u64, u64)>>,
162    scans: AtomicU64,
163    max_object_size: Option<u64>,
164    compression: Option<i32>,
165    // Shared rather than cloned: the staging store and the store proper are two
166    // handles onto one deployment's keys.
167    keys: Option<std::sync::Arc<crypt::Keyring>>,
168}
169
170impl LocalStore {
171    pub fn new(root: impl Into<PathBuf>) -> Self {
172        Self {
173            root: root.into(),
174            counter: AtomicU64::new(0),
175            usage: Mutex::new(None),
176            scans: AtomicU64::new(0),
177            max_object_size: None,
178            compression: None,
179            keys: None,
180        }
181    }
182
183    pub fn with_compression(mut self, level: Option<i32>) -> Self {
184        self.compression = level;
185        self
186    }
187
188    pub fn with_max_object_size(mut self, limit: Option<u64>) -> Self {
189        self.max_object_size = limit;
190        self
191    }
192
193    pub fn with_encryption(mut self, keys: Option<std::sync::Arc<crypt::Keyring>>) -> Self {
194        self.keys = keys;
195        self
196    }
197
198    pub(super) fn keyring(&self) -> Option<&std::sync::Arc<crypt::Keyring>> {
199        self.keys.as_ref()
200    }
201
202    pub fn encrypts(&self) -> bool {
203        self.keys.is_some()
204    }
205
206    // Whether an object in this store is a framing of the bytes rather than the
207    // bytes. Compression and encryption both put a frame under the plaintext
208    // digest, and only the codec turns one back into the object a client asked
209    // for by that name.
210    pub fn frames(&self) -> bool {
211        self.compression.is_some() || self.keys.is_some()
212    }
213
214    fn object_path(&self, ns: &Namespace, oid: &Oid) -> PathBuf {
215        let (first, second) = oid.fanout();
216        self.root
217            .join(ns.org())
218            .join(ns.repo())
219            .join(first)
220            .join(second)
221            .join(oid.as_str())
222    }
223
224    fn content_path(&self, oid: &Oid) -> PathBuf {
225        let (first, second) = oid.fanout();
226        self.root
227            .join(".content")
228            .join(first)
229            .join(second)
230            .join(oid.as_str())
231    }
232
233    pub fn scans(&self) -> u64 {
234        self.scans.load(Ordering::Relaxed)
235    }
236
237    pub async fn writable(&self) -> Result<(), Error> {
238        fs::create_dir_all(&self.root).await?;
239
240        let ticket = self.counter.fetch_add(1, Ordering::Relaxed);
241        let probe = self.root.join(format!(".readiness.{ticket}"));
242
243        fs::write(&probe, b"").await?;
244        fs::remove_file(&probe).await?;
245
246        Ok(())
247    }
248
249    pub async fn exists(&self, ns: &Namespace, oid: &Oid) -> bool {
250        fs::metadata(self.object_path(ns, oid)).await.is_ok()
251    }
252
253    pub async fn open(&self, ns: &Namespace, oid: &Oid) -> Result<Object, Error> {
254        let path = self.object_path(ns, oid);
255        let file = fs::File::open(&path).await.map_err(|_| Error::NotFound)?;
256        let on_disk = file.metadata().await?.len();
257
258        match codec::Framed::open(
259            codec::Reader::File(file),
260            on_disk,
261            self.keys.as_deref(),
262            oid,
263        )
264        .await?
265        {
266            Some(framed) => Ok(Object::Framed(framed)),
267            None => Ok(Object::Raw {
268                file: fs::File::open(&path).await.map_err(|_| Error::NotFound)?,
269                size: on_disk,
270            }),
271        }
272    }
273
274    // Everything a transfer has to survive before it counts as an object: the
275    // digest it claims, the size it declared, the ceiling on a single object and
276    // the repository's remaining budget. It ends on local disk whatever the
277    // backend is, because a bucket cannot be asked to hold bytes that might turn
278    // out to be the wrong ones.
279    pub async fn stage<S, E>(
280        &self,
281        ns: &Namespace,
282        oid: &Oid,
283        expected_size: Option<u64>,
284        budget: Option<Budget>,
285        mut chunks: S,
286    ) -> Result<Staged, Error>
287    where
288        S: Stream<Item = Result<axum::body::Bytes, E>> + Unpin,
289        E: std::error::Error + Send + Sync + 'static,
290    {
291        if let Some(limit) = self.max_object_size
292            && expected_size.is_some_and(|declared| declared > limit)
293        {
294            return Err(Error::TooLarge { limit });
295        }
296
297        let path = self.object_path(ns, oid);
298        let parent = path.parent().expect("object paths always have a parent");
299        fs::create_dir_all(parent).await?;
300
301        // A retried transfer of an object this repository already holds costs it
302        // no room, so it must not count against the budget a second time.
303        let fresh = fs::metadata(&path).await.is_err();
304        let staged = self.staging_path(parent, oid);
305
306        match self.stream_to(&staged, oid, budget, &mut chunks).await {
307            Ok((digest, written)) => {
308                if let Err(error) = Self::agrees(oid, expected_size, &digest, written) {
309                    let _ = fs::remove_file(&staged).await;
310                    return Err(error);
311                }
312
313                Ok(Staged {
314                    path: staged,
315                    destination: path,
316                    written,
317                    fresh,
318                })
319            }
320            Err(error) => {
321                let _ = fs::remove_file(&staged).await;
322                Err(error)
323            }
324        }
325    }
326
327    fn agrees(
328        oid: &Oid,
329        expected_size: Option<u64>,
330        digest: &str,
331        written: u64,
332    ) -> Result<(), Error> {
333        if let Some(declared) = expected_size.filter(|declared| *declared != written) {
334            return Err(Error::SizeMismatch {
335                declared,
336                actual: written,
337            });
338        }
339
340        if digest != oid.as_str() {
341            return Err(Error::OidMismatch {
342                declared: oid.to_string(),
343                actual: digest.to_owned(),
344            });
345        }
346
347        Ok(())
348    }
349
350    pub async fn write<S, E>(
351        &self,
352        ns: &Namespace,
353        oid: &Oid,
354        expected_size: Option<u64>,
355        budget: Option<Budget>,
356        chunks: S,
357    ) -> Result<Written, Error>
358    where
359        S: Stream<Item = Result<axum::body::Bytes, E>> + Unpin,
360        E: std::error::Error + Send + Sync + 'static,
361    {
362        let staged = self.stage(ns, oid, expected_size, budget, chunks).await?;
363
364        self.link_or_move(&staged.path, &staged.destination, oid)
365            .await?;
366
367        Ok(Written {
368            bytes: staged.written,
369            fresh: staged.fresh,
370        })
371    }
372
373    fn staging_path(&self, parent: &Path, oid: &Oid) -> PathBuf {
374        let ticket = self.counter.fetch_add(1, Ordering::Relaxed);
375        parent.join(format!("{oid}.{ticket}.part"))
376    }
377
378    async fn stream_to<S, E>(
379        &self,
380        staged: &Path,
381        oid: &Oid,
382        budget: Option<Budget>,
383        chunks: &mut S,
384    ) -> Result<(String, u64), Error>
385    where
386        S: Stream<Item = Result<axum::body::Bytes, E>> + Unpin,
387        E: std::error::Error + Send + Sync + 'static,
388    {
389        let file = fs::File::create(staged).await?;
390        // The digest, the declared size and the budget are all counted on the
391        // plaintext going past, whatever the bytes look like once they land,
392        // so compression is a different sink, not a different path.
393        let mut sink = match (self.compression, self.keys.as_deref()) {
394            (None, None) => Sink::Raw(file),
395            (level, keys) => Sink::Framed(Box::new(
396                codec::Writer::open(file, level, keys.map(crypt::Keyring::writing), oid).await?,
397            )),
398        };
399        let mut hasher = Sha256::new();
400        let mut written = 0u64;
401
402        while let Some(chunk) = chunks.next().await {
403            let chunk = chunk.map_err(std::io::Error::other)?;
404            hasher.update(&chunk);
405            written += chunk.len() as u64;
406
407            // The declared size is a claim by the client, so the ceiling has to
408            // hold against a body that ignores it. Stopping at the chunk that
409            // crosses the line is the point: reading to the end to find out how
410            // big it was would be the outage this limit exists to prevent.
411            if let Some(limit) = self.max_object_size.filter(|limit| written > *limit) {
412                return Err(Error::TooLarge { limit });
413            }
414
415            if let Some(budget) = budget.filter(|budget| budget.exceeded_by(written)) {
416                return Err(budget.refusal());
417            }
418
419            sink.write(&chunk).await?;
420        }
421
422        sink.finish().await?;
423
424        Ok((hex::encode(hasher.finalize()), written))
425    }
426
427    // One copy of the bytes under .content, and a hard link per repository that
428    // holds them. Two projects sharing an asset pack cost the disk once, and the
429    // link count is the reference count: the filesystem does the bookkeeping, so
430    // nothing can leak a repository's contents to another and nothing needs a
431    // migration: objects already sitting at their repository path keep working as
432    // ordinary files with one link.
433    async fn link_or_move(&self, staged: &Path, final_path: &Path, oid: &Oid) -> Result<(), Error> {
434        let content = self.content_path(oid);
435        let parent = content.parent().expect("content paths have a parent");
436        fs::create_dir_all(parent).await?;
437
438        if fs::metadata(&content).await.is_err() {
439            fs::rename(staged, &content).await?;
440        }
441
442        match self.link(&content, final_path).await {
443            // The content was collected between finding it and linking to it:
444            // a concurrent retain on another repository dropped its last other
445            // reference. The staged copy is still here precisely for this, so
446            // put it back and link again rather than failing a push that did
447            // nothing wrong.
448            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
449                fs::rename(staged, &content).await?;
450                self.link(&content, final_path).await?;
451            }
452            outcome => outcome?,
453        }
454
455        let _ = fs::remove_file(staged).await;
456        Ok(())
457    }
458
459    async fn link(&self, content: &Path, final_path: &Path) -> Result<(), std::io::Error> {
460        let from = content.to_path_buf();
461        let to = final_path.to_path_buf();
462        let linked = tokio::task::spawn_blocking(move || std::fs::hard_link(&from, &to))
463            .await
464            .map_err(std::io::Error::other)?;
465
466        match linked {
467            Ok(()) => Ok(()),
468            Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => Ok(()),
469            Err(error) if error.kind() == std::io::ErrorKind::NotFound => Err(error),
470            // A filesystem without hard links, or one crossing a device
471            // boundary: fall back to a full copy so the transfer still
472            // succeeds. The disk pays for it, the client never notices.
473            Err(_) => fs::copy(content, final_path).await.map(|_| ()),
474        }
475    }
476}