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
62pub struct Staged {
64 pub path: PathBuf,
65 destination: PathBuf,
66 pub written: u64,
67 fresh: bool,
68}
69
70pub 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#[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#[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 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 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.stored_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 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 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 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 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 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 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 Err(_) => fs::copy(content, final_path).await.map(|_| ()),
474 }
475 }
476}