1use std::collections::HashMap;
5use std::fmt;
6use std::fs::{self, File, OpenOptions};
7use std::io::{self, ErrorKind, Read as _, Seek as _, SeekFrom, Write as _};
8use std::path::{Path, PathBuf};
9use std::pin::Pin;
10use std::sync::{Arc, Mutex, Weak};
11use std::task::{Context, Poll};
12use std::time::{Duration, SystemTime};
13
14use bytes::{Bytes, BytesMut};
15use futures_core::Stream;
16use mkit_core::hash::{Hash, Hasher};
17use mkit_transport_file::{create_dir_all_durably, sync_dir, temp_path};
18
19use super::{io_error, unavailable};
20use crate::store::{
21 BlobBody, BlobKey, BlobMeta, BlobStore, ByteRange, CommitOutcome, MAX_BLOB_PIECE_BYTES,
22 PackSink, StoreError,
23};
24
25pub(super) const READ_BLOCK: usize = 64 * 1024;
27
28type MultipartLocks = Arc<Mutex<HashMap<[u8; 32], Weak<tokio::sync::Mutex<()>>>>>;
29
30#[derive(Debug, Clone)]
48pub struct FsBlobStore {
49 pub(super) root: PathBuf,
50 keyspace: &'static str,
51 pub(super) multipart_locks: MultipartLocks,
52}
53
54impl FsBlobStore {
55 #[must_use]
58 pub fn new(root: impl Into<PathBuf>) -> Self {
59 Self::with_keyspace(root, "packs")
60 }
61
62 #[must_use]
70 pub fn with_keyspace(root: impl Into<PathBuf>, keyspace: &'static str) -> Self {
71 assert!(
72 !keyspace.is_empty() && !keyspace.starts_with('.') && !keyspace.contains(['/', '\\']),
73 "a keyspace is one plain path component: {keyspace:?}"
74 );
75 assert!(
76 !crate::store::is_reserved_pack_keyspace(keyspace),
77 "a keyspace must not alias a sibling namespace: {keyspace:?}"
78 );
79 Self {
80 root: root.into(),
81 keyspace,
82 multipart_locks: Arc::new(Mutex::new(HashMap::new())),
83 }
84 }
85
86 #[must_use]
88 pub fn root(&self) -> &Path {
89 &self.root
90 }
91
92 #[must_use]
94 pub fn keyspace(&self) -> &'static str {
95 self.keyspace
96 }
97
98 fn dir(&self) -> PathBuf {
99 self.root.join(self.keyspace)
100 }
101
102 fn path(&self, key: &BlobKey) -> Result<PathBuf, StoreError> {
103 Ok(self.root.join(key.relative_path(self.keyspace)?))
104 }
105
106 pub fn sweep_stale_uploads(&self, min_age: Duration) -> io::Result<usize> {
128 let now = SystemTime::now();
129 let mut removed = 0;
130 for dir in [
131 self.dir(),
132 self.root.join("upload-markers/v1"),
133 self.root.join("objects"),
134 self.root.join("object-offsets/v1"),
135 ] {
136 let entries = match fs::read_dir(dir) {
137 Ok(entries) => entries,
138 Err(e) if e.kind() == ErrorKind::NotFound => continue,
139 Err(e) => return Err(e),
140 };
141 for entry in entries {
142 let Ok(entry) = entry else { continue };
143 let name = entry.file_name();
144 if !name.to_str().is_some_and(is_upload_temp_name) {
145 continue;
146 }
147 let Ok(meta) = entry.metadata() else { continue };
149 let stale = meta.is_file()
150 && meta
151 .modified()
152 .ok()
153 .and_then(|m| now.duration_since(m).ok())
154 .is_some_and(|age| age >= min_age);
155 if stale && fs::remove_file(entry.path()).is_ok() {
156 removed += 1;
157 }
158 }
159 }
160 Ok(removed + super::multipart::sweep_sessions(&self.root, now)?)
161 }
162}
163
164pub(super) fn is_upload_temp_name(name: &str) -> bool {
167 let decimal = |s: &str, max: usize| {
168 !s.is_empty() && s.len() <= max && s.bytes().all(|b| b.is_ascii_digit())
169 };
170 let Some(rest) = name.strip_prefix('.') else {
171 return false;
172 };
173 let Some((hex, rest)) = rest.split_at_checked(64) else {
174 return false;
175 };
176 let Some((pid, seq)) = rest.strip_prefix(".tmp.").and_then(|r| r.split_once('.')) else {
177 return false;
178 };
179 hex.bytes().all(|b| matches!(b, b'0'..=b'9' | b'a'..=b'f'))
180 && decimal(pid, 10)
181 && decimal(seq, 20)
182}
183
184pub struct FsPackSink {
187 file: Option<File>,
189 tmp: Option<PathBuf>,
191 dest: PathBuf,
192 key: BlobKey,
193 declared: u64,
194 written: u64,
195 hasher: Hasher,
196}
197
198impl fmt::Debug for FsPackSink {
199 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
200 f.debug_struct("FsPackSink")
201 .field("tmp", &self.tmp)
202 .field("dest", &self.dest)
203 .field("declared", &self.declared)
204 .field("written", &self.written)
205 .finish_non_exhaustive()
206 }
207}
208
209impl FsPackSink {
210 fn publish(&mut self, root: Option<Hash>) -> Result<bool, StoreError> {
213 let expected = self.key.expected_root(root)?;
214 if self.written != self.declared {
215 return Err(StoreError::Invalid("blob length does not match".into()));
216 }
217 if self.hasher.finalize() != expected {
218 return Err(StoreError::Invalid(
219 "blob hash does not match its key".into(),
220 ));
221 }
222 let file = self.file.take().ok_or_else(closed)?;
223 file.sync_all().map_err(io_error)?;
224 drop(file);
225 let tmp = self.tmp.as_ref().ok_or_else(closed)?;
226 let existed = self.dest.exists();
227 fs::rename(tmp, &self.dest).map_err(io_error)?;
230 self.tmp = None;
231 if let Some(dir) = self.dest.parent() {
232 sync_dir(dir).map_err(io_error)?;
233 }
234 Ok(existed)
235 }
236}
237
238fn outcome(existed: bool) -> CommitOutcome {
239 if existed {
240 CommitOutcome::AlreadyPresent
241 } else {
242 CommitOutcome::Created
243 }
244}
245
246fn closed() -> StoreError {
249 StoreError::unavailable(io::Error::other("blob upload already closed"))
250}
251
252impl Drop for FsPackSink {
253 fn drop(&mut self) {
254 self.file = None;
255 if let Some(tmp) = self.tmp.take() {
256 let _ = fs::remove_file(tmp);
257 }
258 }
259}
260
261impl PackSink for FsPackSink {
262 async fn write(&mut self, chunk: Bytes) -> Result<(), StoreError> {
263 let total = self.written.checked_add(chunk.len() as u64);
264 let Some(total) = total.filter(|t| *t <= self.declared) else {
265 return Err(StoreError::Invalid("blob is longer than declared".into()));
266 };
267 let file = self.file.as_mut().ok_or_else(closed)?;
268 file.write_all(&chunk).map_err(io_error)?;
269 self.hasher.update(&chunk);
270 self.written = total;
271 Ok(())
272 }
273
274 async fn commit(mut self) -> Result<CommitOutcome, StoreError> {
275 Ok(outcome(self.publish(None)?))
277 }
278
279 async fn commit_with_root(mut self, content_root: Hash) -> Result<CommitOutcome, StoreError> {
280 Ok(outcome(self.publish(Some(content_root))?))
281 }
282
283 async fn abort(self) {}
284}
285
286struct Blocks {
288 file: File,
289 remaining: u64,
290}
291
292impl Stream for Blocks {
293 type Item = Result<Bytes, StoreError>;
294
295 fn poll_next(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll<Option<Self::Item>> {
296 let this = self.get_mut();
297 if this.remaining == 0 {
298 return Poll::Ready(None);
299 }
300 let n = usize::try_from(this.remaining).map_or(READ_BLOCK, |r| r.min(READ_BLOCK));
301 let mut piece = BytesMut::zeroed(n);
302 match this.file.read_exact(&mut piece) {
303 Ok(()) => {
304 this.remaining -= n as u64;
305 Poll::Ready(Some(Ok(piece.freeze())))
306 }
307 Err(e) => {
308 this.remaining = 0;
309 Poll::Ready(Some(Err(io_error(e))))
310 }
311 }
312 }
313}
314
315const _: () = assert!(READ_BLOCK <= MAX_BLOB_PIECE_BYTES);
317
318impl BlobStore for FsBlobStore {
319 type Sink = FsPackSink;
320
321 async fn begin(&self, key: BlobKey, len: u64) -> Result<FsPackSink, StoreError> {
322 let dest = self.path(&key)?;
323 let dir = dest
324 .parent()
325 .ok_or_else(|| StoreError::Invalid("blob path has no directory".into()))?;
326 create_dir_all_durably(dir).map_err(io_error)?;
327 let tmp = temp_path(&dest).map_err(io_error)?;
328 let file = OpenOptions::new()
329 .write(true)
330 .create_new(true)
331 .open(&tmp)
332 .map_err(io_error)?;
333 Ok(FsPackSink {
334 file: Some(file),
335 tmp: Some(tmp),
336 dest,
337 key,
338 declared: len,
339 written: 0,
340 hasher: Hasher::new(),
341 })
342 }
343
344 async fn get(
345 &self,
346 key: &BlobKey,
347 range: Option<ByteRange>,
348 ) -> Result<Option<BlobBody>, StoreError> {
349 let mut file = match File::open(self.path(key)?) {
350 Ok(file) => file,
351 Err(e) if e.kind() == ErrorKind::NotFound => return Ok(None),
352 Err(e) => return Err(io_error(e)),
353 };
354 let len = file.metadata().map_err(io_error)?.len();
355 let span = match range {
356 Some(range) => range.resolve(len)?,
357 None => 0..len,
358 };
359 if span.start > 0 {
360 file.seek(SeekFrom::Start(span.start)).map_err(io_error)?;
361 }
362 let n = span.end - span.start;
363 if let Ok(whole) = usize::try_from(n)
364 && whole <= MAX_BLOB_PIECE_BYTES
365 {
366 let mut buf = vec![0; whole];
367 file.read_exact(&mut buf).map_err(io_error)?;
368 return Ok(Some(BlobBody::Bytes(Bytes::from(buf))));
369 }
370 Ok(Some(BlobBody::Stream {
371 len: n,
372 stream: Box::pin(Blocks { file, remaining: n }),
373 }))
374 }
375
376 async fn head(&self, key: &BlobKey) -> Result<Option<BlobMeta>, StoreError> {
377 match fs::metadata(self.path(key)?) {
378 Ok(meta) => Ok(Some(BlobMeta { len: meta.len() })),
379 Err(e) if e.kind() == ErrorKind::NotFound => Ok(None),
380 Err(e) => Err(io_error(e)),
381 }
382 }
383
384 async fn probe(&self) -> Result<(), StoreError> {
385 let meta = fs::metadata(&self.root).map_err(io_error)?;
386 if meta.is_dir() {
387 Ok(())
388 } else {
389 Err(unavailable(io::Error::other(
390 "blob root is not a directory",
391 )))
392 }
393 }
394
395 async fn delete(&self, key: &BlobKey) -> Result<bool, StoreError> {
396 let path = self.path(key)?;
397 match fs::remove_file(&path) {
398 Ok(()) => {}
399 Err(e) if e.kind() == ErrorKind::NotFound => return Ok(false),
400 Err(e) => return Err(io_error(e)),
401 }
402 if let Some(dir) = path.parent() {
403 sync_dir(dir).map_err(io_error)?;
404 }
405 Ok(true)
406 }
407}