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;