Skip to main content

lfsx_server/storage/
backend.rs

1use futures_util::Stream;
2
3use super::s3::S3Store;
4use super::{Budget, CompressReport, DedupeReport, LocalStore, Object, SweepReport, VerifyReport};
5use crate::error::Error;
6use crate::namespace::Namespace;
7use crate::oid::Oid;
8#[cfg(test)]
9use sha2::Digest;
10#[cfg(test)]
11use std::time::Duration;
12
13// Where the objects live. A bucket decouples capacity from the machine, at the
14// price of the things a filesystem gave for nothing: hard links, a directory
15// walk, and a rename that is atomic. Each of those is answered here or refused
16// out loud; none of them is quietly skipped.
17pub struct Store {
18    backend: Backend,
19    usage: super::usage::Usage,
20}
21
22enum Backend {
23    Local(LocalStore),
24    // Even with a bucket the local store stays, because a transfer has to land
25    // somewhere before anyone can tell whether it is the object it claims to be.
26    // It is a write buffer, not the store.
27    // Boxed because a bucket handle beside a local store makes this variant far
28    // larger than the other, and every Store in the process would pay for it.
29    Bucket {
30        bucket: Box<S3Store>,
31        staging: LocalStore,
32        cache: Option<std::sync::Arc<super::cache::Cache>>,
33    },
34}
35
36// One fetch, off the request path, and only from whoever claimed the object:
37// eight parallel transfers of the same cold pack must not become eight copies
38// of the same download. A failure here is a cache that stays cold, which is
39// what it already was, so it is logged at debug and forgotten.
40fn fill_behind(cache: std::sync::Arc<super::cache::Cache>, bucket: S3Store, oid: Oid, size: u64) {
41    // An object that cannot fit is never fetched: filling it would evict the
42    // whole cache and then itself, leaving nothing behind and doing it again on
43    // the next download.
44    if !cache.fits(size) || !cache.claim(&oid) {
45        return;
46    }
47
48    tokio::spawn(async move {
49        let filled = async {
50            use futures_util::StreamExt;
51
52            let chunks = bucket
53                .read(&oid, 0, size)
54                .await?
55                .map(|chunk| chunk.map_err(|error| Error::Storage(std::io::Error::other(error))));
56
57            cache.fill(&oid, size, chunks).await
58        }
59        .await;
60
61        if let Err(error) = filled {
62            tracing::debug!(%error, %oid, "an object could not be copied into the cache");
63        }
64
65        cache.release(&oid);
66    });
67}
68
69impl Store {
70    pub fn local(store: LocalStore) -> Self {
71        Self::over(Backend::Local(store))
72    }
73
74    // Compression and encryption used to be stripped here, because a framed
75    // object was only readable through the file the codec opened and a bucket key
76    // is not one. The codec now reads from a bucket too, so the frames go up as
77    // they are and come back decoded: the header and the index are three ranged
78    // GETs, which is what the format was shaped for.
79    pub fn bucket(bucket: S3Store, staging: LocalStore) -> Self {
80        Self::over(Backend::Bucket {
81            bucket: Box::new(bucket),
82            staging,
83            cache: None,
84        })
85    }
86
87    pub fn with_cache(mut self, disk: Option<super::cache::Cache>) -> Self {
88        if let (Backend::Bucket { cache, .. }, Some(disk)) = (&mut self.backend, disk) {
89            *cache = Some(std::sync::Arc::new(disk));
90        }
91
92        self
93    }
94
95    pub fn cache_stats(&self) -> Option<super::cache::Stats> {
96        match &self.backend {
97            Backend::Bucket {
98                cache: Some(cache), ..
99            } => Some(cache.stats()),
100            _ => None,
101        }
102    }
103
104    fn over(backend: Backend) -> Self {
105        Self {
106            backend,
107            usage: super::usage::Usage::default(),
108        }
109    }
110
111    fn staging(&self) -> &LocalStore {
112        match &self.backend {
113            Backend::Local(store) => store,
114            Backend::Bucket { staging, .. } => staging,
115        }
116    }
117
118    // Everything an interrupted upload can leave behind, wherever it left it. A
119    // bucket deployment still stages locally, so both are swept and the figures
120    // add up to one answer.
121    pub async fn reclaim(&self, older_than: std::time::Duration) -> super::Reclaimed {
122        let mut reclaimed = self.staging().reclaim_staging(older_than).await;
123
124        if let Backend::Bucket { bucket, .. } = &self.backend {
125            match bucket.reclaim_incoming(older_than).await {
126                Ok(theirs) => {
127                    reclaimed.files += theirs.files;
128                    reclaimed.bytes += theirs.bytes;
129                }
130                Err(error) => {
131                    tracing::warn!(%error, "abandoned uploads in the bucket could not be reclaimed");
132                }
133            }
134        }
135
136        reclaimed
137    }
138
139    // Readiness has to ask the backend that actually serves. Once the objects
140    // live in a bucket the volume is a write buffer, and an instance whose
141    // credentials were rotated or whose bucket is gone passes a probe that only
142    // proves its scratch disk works, then fails every transfer it is handed.
143    //
144    // So both are asked, either failing takes the instance out, and they are
145    // named apart: a full disk and a rotated key are not the same afternoon.
146    pub async fn writable(&self) -> Result<(), Error> {
147        self.staging().writable().await.map_err(|error| {
148            Error::Storage(std::io::Error::other(format!(
149                "the staging volume is not writable: {error}"
150            )))
151        })?;
152
153        if let Backend::Bucket { bucket, .. } = &self.backend {
154            bucket.reachable().await?;
155        }
156
157        Ok(())
158    }
159
160    pub fn scans(&self) -> u64 {
161        self.staging().scans()
162    }
163
164    #[tracing::instrument(skip_all, fields(namespace = %ns, oid = %oid))]
165    pub async fn exists(&self, ns: &Namespace, oid: &Oid) -> bool {
166        match &self.backend {
167            Backend::Local(store) => store.exists(ns, oid).await,
168            Backend::Bucket { bucket, .. } => bucket.exists(ns, oid).await,
169        }
170    }
171
172    // Where the client should fetch this object from, when that is somewhere
173    // other than this server. None for a local store, and for a bucket the
174    // operator has not asked to redirect, which is the default, because the
175    // streamed path is the one that counts the bytes and holds the ceiling.
176    //
177    // The caller is responsible for having established that this repository
178    // holds the object. This hands out a signature, not a permission.
179    pub fn redirect(&self, oid: &Oid) -> Option<String> {
180        match &self.backend {
181            Backend::Local(_) => None,
182            // A pre-signed URL hands over whatever sits under that key, and with
183            // a codec in the path that is a frame rather than the object. The
184            // client would hash what arrived, get a digest that is not the one it
185            // asked for, and reject it. So the redirect is given up and the
186            // download streams, which is the only path that can decode.
187            //
188            // Compression is enough on its own, even though it still lets a
189            // client upload straight to the bucket. That asymmetry is the right
190            // way round: an unframed object is a perfectly good entry, so a
191            // direct upload stays safe, while one framed object anywhere in the
192            // store makes every redirect a guess.
193            Backend::Bucket { staging, .. } if staging.frames() => None,
194            Backend::Bucket { bucket, .. } => bucket.presigned_download(oid),
195        }
196    }
197
198    // Where the client should PUT the object, when that is the bucket rather than
199    // this server. None for a local store and for a bucket the operator has not
200    // asked to redirect.
201    pub fn presigned_upload(
202        &self,
203        ns: &Namespace,
204        oid: &Oid,
205        size: u64,
206    ) -> Option<super::s3::Presigned> {
207        match &self.backend {
208            Backend::Local(_) => None,
209            // A client uploading straight to the bucket writes the object as it
210            // is, so a configured key would never touch it and the bucket would
211            // hold plaintext while an operator believed otherwise. Encryption is
212            // a promise about what the storage provider can read; a faster upload
213            // is not worth quietly breaking it. Those transfers keep coming
214            // through the server, which seals them.
215            Backend::Bucket { staging, .. } if staging.encrypts() => None,
216            Backend::Bucket { bucket, .. } => bucket.presigned_upload(ns, oid, size),
217        }
218    }
219
220    // How big an object waiting under this repository's own upload key is. None
221    // when there is nothing waiting, which is every local deployment and every
222    // client that has not used its URL.
223    #[tracing::instrument(skip_all, fields(namespace = %ns, oid = %oid))]
224    pub async fn uploaded_size(&self, ns: &Namespace, oid: &Oid) -> Result<Option<u64>, Error> {
225        match &self.backend {
226            Backend::Local(_) => Ok(None),
227            Backend::Bucket { bucket, .. } => Ok(bucket.uploaded_size(ns, oid).await.ok()),
228        }
229    }
230
231    // Take an upload this repository made into the shared keyspace. Only reachable
232    // for a bucket, because only there does a client write anywhere this server
233    // did not.
234    #[tracing::instrument(skip_all, fields(namespace = %ns, oid = %oid))]
235    pub async fn adopt(&self, ns: &Namespace, oid: &Oid, arrived: u64) -> Result<(), Error> {
236        let outcome = match &self.backend {
237            Backend::Local(_) => Err(Error::Unsupported(
238                "objects are written through this server, so there is nothing to adopt",
239            )),
240            Backend::Bucket { bucket, .. } => bucket.adopt(ns, oid, arrived).await,
241        };
242
243        // Same reason as a write: verify is called once per object, so dropping
244        // what is remembered here would make every one of them re-measure.
245        if outcome.is_ok() {
246            self.usage.stored(ns, arrived).await;
247        }
248
249        outcome
250    }
251
252    #[tracing::instrument(skip_all, fields(namespace = %ns, oid = %oid))]
253    pub async fn open(&self, ns: &Namespace, oid: &Oid) -> Result<Object, Error> {
254        match &self.backend {
255            Backend::Local(store) => store.open(ns, oid).await,
256            Backend::Bucket {
257                bucket,
258                staging,
259                cache,
260            } => {
261                // The marker is the proof of possession and is checked before
262                // anything is read, exactly as a local open checks the link.
263                if !bucket.exists(ns, oid).await {
264                    return Err(Error::NotFound);
265                }
266
267                let size = bucket.size_of(oid).await?;
268
269                // A cached object is read the way the volume backend reads one,
270                // because it is the same bytes in the same form. What is not
271                // cached is served from the bucket as before and filled in
272                // behind the request, so a cold client waits for the bucket and
273                // not for the copy.
274                let cached = match cache {
275                    Some(cache) => cache.open(oid).await,
276                    None => None,
277                };
278
279                let reader = match cached {
280                    Some(file) => super::codec::Reader::File(file),
281                    None => {
282                        if let Some(cache) = cache {
283                            fill_behind(cache.clone(), (**bucket).clone(), oid.to_owned(), size);
284                        }
285
286                        super::codec::Reader::Bucket {
287                            bucket: (**bucket).clone(),
288                            oid: oid.to_owned(),
289                        }
290                    }
291                };
292                let from_cache = matches!(reader, super::codec::Reader::File(_));
293
294                match super::codec::Framed::open(
295                    reader,
296                    size,
297                    staging.keyring().map(AsRef::as_ref),
298                    oid,
299                )
300                .await?
301                {
302                    Some(framed) => Ok(Object::Framed(framed)),
303                    // Not one of ours: the object is the bytes, and streaming
304                    // them straight through costs no extra round trip.
305                    // A raw object that is cached is a plain local file, so
306                    // it streams like one. `Framed::open` consumed the handle
307                    // deciding it was not framed, hence the reopen.
308                    None if from_cache => match cache.as_ref().unwrap().reopen(oid).await {
309                        Some(file) => Ok(Object::Raw { file, size }),
310                        None => Ok(Object::Remote {
311                            bucket: (**bucket).clone(),
312                            oid: oid.to_owned(),
313                            size,
314                        }),
315                    },
316                    None => Ok(Object::Remote {
317                        bucket: (**bucket).clone(),
318                        oid: oid.to_owned(),
319                        size,
320                    }),
321                }
322            }
323        }
324    }
325
326    #[tracing::instrument(skip_all, fields(namespace = %ns, oid = %oid))]
327    pub async fn write<S, E>(
328        &self,
329        ns: &Namespace,
330        oid: &Oid,
331        expected_size: Option<u64>,
332        budget: Option<Budget>,
333        chunks: S,
334    ) -> Result<u64, Error>
335    where
336        S: Stream<Item = Result<axum::body::Bytes, E>> + Unpin,
337        E: std::error::Error + Send + Sync + 'static,
338    {
339        let written = match &self.backend {
340            Backend::Local(store) => store.write(ns, oid, expected_size, budget, chunks).await?,
341            Backend::Bucket {
342                bucket, staging, ..
343            } => {
344                // Asked of the bucket, because the staging store answers about a
345                // local layout a bucket deployment never fills in: it would call
346                // every upload fresh, and re-pushing an object the repository
347                // already holds would grow what is remembered without anything
348                // being stored.
349                let fresh = !bucket.exists(ns, oid).await;
350
351                let staged = staging
352                    .stage(ns, oid, expected_size, budget, chunks)
353                    .await?;
354                let outcome = bucket.store(ns, oid, &staged.path).await;
355
356                // The staging file has served its purpose either way. Leaving it
357                // would be a leak the reclaimer only notices a day later.
358                let _ = tokio::fs::remove_file(&staged.path).await;
359                outcome?;
360
361                super::Written {
362                    bytes: staged.written,
363                    fresh,
364                }
365            }
366        };
367
368        // Added to what is remembered rather than dropping it: a client pushing
369        // a hundred objects would otherwise make the next negotiation measure
370        // the repository again, which on a bucket is what this cache exists to
371        // avoid.
372        if written.fresh {
373            self.usage.stored(ns, written.bytes).await;
374        }
375
376        Ok(written.bytes)
377    }
378
379    // None rather than zero: a bucket has no cheap answer for what the whole
380    // store holds, and building one from a full listing would cost a request per
381    // object on every scrape. Zero would be read as an empty bucket by every
382    // dashboard that averages it, which is the one lie this seam otherwise
383    // refuses to tell: everything else it cannot do answers 501.
384    pub async fn capacity(&self) -> Option<(u64, u64)> {
385        match &self.backend {
386            Backend::Local(store) => Some(store.usage().await),
387            Backend::Bucket { .. } => None,
388        }
389    }
390
391    // Measured at most once a minute per repository, whichever backend is
392    // behind it. A bucket answers this by listing the repository's markers and
393    // asking the size of each, so one uncached call per object in a batch made
394    // a hundred-object push cost a hundred listings: the product, not the sum.
395    pub async fn usage_of(&self, ns: &Namespace) -> (u64, u64) {
396        if let Some(cached) = self.usage.cached(ns).await {
397            return cached;
398        }
399
400        let measured = match &self.backend {
401            Backend::Local(store) => store.measure_of(ns).await,
402            Backend::Bucket { bucket, .. } => bucket.usage_of(ns).await,
403        };
404
405        self.usage.remember(ns, measured.0, measured.1).await;
406
407        measured
408    }
409
410    #[tracing::instrument(skip_all, fields(namespace = %ns, dry_run))]
411    pub async fn sweep(
412        &self,
413        ns: &Namespace,
414        retained: &std::collections::HashSet<String>,
415        grace: std::time::Duration,
416        dry_run: bool,
417    ) -> Result<SweepReport, Error> {
418        match &self.backend {
419            Backend::Local(store) => {
420                let report = store.sweep(ns, retained, grace, dry_run).await;
421
422                // Freeing gigabytes and then answering the next quota check from
423                // the figure measured before is how a client is refused space it
424                // has just been told it reclaimed.
425                self.usage.forget(ns).await;
426
427                report
428            }
429            Backend::Bucket { bucket, .. } => {
430                let report = bucket.sweep(ns, retained, grace, dry_run).await;
431
432                if report.is_ok() && !dry_run {
433                    self.usage.forget(ns).await;
434                }
435
436                report
437            }
438        }
439    }
440
441    #[tracing::instrument(skip_all, fields(namespace = %ns, dry_run))]
442    pub async fn dedupe(&self, ns: &Namespace, dry_run: bool) -> Result<DedupeReport, Error> {
443        match &self.backend {
444            Backend::Local(store) => {
445                let report = store.dedupe(ns, dry_run).await;
446
447                // Freeing gigabytes and then answering the next quota check from
448                // the figure measured before is how a client is refused space it
449                // has just been told it reclaimed.
450                self.usage.forget(ns).await;
451
452                report
453            }
454            // Content addressing already gives this: two repositories pushing the
455            // same object write the same key, and each holds a marker beside it.
456            // There is nothing left to fold in.
457            Backend::Bucket { .. } => Err(Error::Unsupported(
458                "a bucket stores each object once already, so there is nothing to deduplicate",
459            )),
460        }
461    }
462
463    #[tracing::instrument(skip_all, fields(namespace = %ns, dry_run))]
464    pub async fn compress(&self, ns: &Namespace, dry_run: bool) -> Result<CompressReport, Error> {
465        match &self.backend {
466            Backend::Local(store) => {
467                let report = store.compress(ns, dry_run).await;
468
469                // Freeing gigabytes and then answering the next quota check from
470                // the figure measured before is how a client is refused space it
471                // has just been told it reclaimed.
472                self.usage.forget(ns).await;
473
474                report
475            }
476            // Objects arriving now are compressed if the server is configured to;
477            // rewriting the ones already in the bucket means walking it and
478            // reuploading, which is a different piece of work.
479            Backend::Bucket { .. } => Err(Error::Unsupported(
480                "rewriting objects already in a bucket is not implemented",
481            )),
482        }
483    }
484
485    #[tracing::instrument(skip_all, fields(namespace = %ns))]
486    pub async fn verify(&self, ns: &Namespace) -> Result<VerifyReport, Error> {
487        match &self.backend {
488            Backend::Local(store) => store.verify(ns).await,
489            Backend::Bucket { .. } => Err(Error::Unsupported(
490                "verification is not implemented for a bucket yet",
491            )),
492        }
493    }
494}
495
496mod meta;
497
498#[cfg(test)]
499mod tests;