Skip to main content

mkit_server/fs/
multipart.rs

1//! Durable multipart staging below `server-uploads/<ticket-id-hex>`.
2
3use std::fs::{self, File, OpenOptions};
4use std::io::{self, ErrorKind, Read as _, Write as _};
5use std::path::{Path, PathBuf};
6use std::sync::atomic::{AtomicU64, Ordering};
7use std::time::{Duration, SystemTime, UNIX_EPOCH};
8
9use bytes::{Bytes, BytesMut};
10use mkit_core::hash::{Hash, hash, to_hex_bytes};
11use mkit_core::upload_parts::{PartHasher, PartPlan, merge_to_root};
12use mkit_transport_file::{create_dir_all_durably, sync_dir, temp_path};
13
14use super::blob::{FsBlobStore, READ_BLOCK};
15use super::io_error;
16use crate::store::{
17    BlobKey, BlobNamespace, BlobStore, CommitOutcome, MultipartBlobStore, PackSink, PartRef,
18    PartSink, StoreError,
19};
20
21const UPLOADS_DIR: &str = "server-uploads";
22const META: &str = "meta";
23const META_MAGIC: &[u8; 5] = b"MKUP1";
24const SESSION_AGE: Duration = Duration::from_hours(169);
25static NEXT_SESSION: AtomicU64 = AtomicU64::new(0);
26
27fn session_lock(
28    store: &FsBlobStore,
29    session: &[u8],
30) -> Result<std::sync::Arc<tokio::sync::Mutex<()>>, StoreError> {
31    let id: [u8; 32] = session.try_into().map_err(|_| StoreError::SessionGone)?;
32    let mut locks = store
33        .multipart_locks
34        .lock()
35        .unwrap_or_else(std::sync::PoisonError::into_inner);
36    if locks.len() >= 1024 {
37        locks.retain(|_, weak| weak.strong_count() > 0);
38    }
39    if let Some(lock) = locks.get(&id).and_then(std::sync::Weak::upgrade) {
40        return Ok(lock);
41    }
42    let lock = std::sync::Arc::new(tokio::sync::Mutex::new(()));
43    locks.insert(id, std::sync::Arc::downgrade(&lock));
44    Ok(lock)
45}
46
47fn uploads(root: &Path) -> PathBuf {
48    root.join(UPLOADS_DIR)
49}
50
51fn session_dir(store: &FsBlobStore, session: &[u8]) -> Result<PathBuf, StoreError> {
52    if session.len() != 32 {
53        return Err(StoreError::SessionGone);
54    }
55    Ok(uploads(&store.root).join(to_hex_bytes(session)))
56}
57
58fn session_meta(key: BlobKey, len: u64, part_size: u64) -> Result<Vec<u8>, StoreError> {
59    if !matches!(key.namespace(), BlobNamespace::Pack | BlobNamespace::Object) {
60        return Err(StoreError::Invalid(
61            "multipart requires a pack or object key".into(),
62        ));
63    }
64    let mut value = Vec::with_capacity(53);
65    value.extend_from_slice(META_MAGIC);
66    value.extend_from_slice(key.hash());
67    value.extend_from_slice(&len.to_be_bytes());
68    value.extend_from_slice(&part_size.to_be_bytes());
69    Ok(value)
70}
71
72fn verify_session(dir: &Path, expected: &[u8]) -> Result<(), StoreError> {
73    match fs::read(dir.join(META)) {
74        Ok(actual) if actual == expected => Ok(()),
75        Ok(_) => Err(StoreError::SessionGone),
76        Err(e) if e.kind() == ErrorKind::NotFound => Err(StoreError::SessionGone),
77        Err(e) => Err(io_error(e)),
78    }
79}
80
81async fn create_session(
82    store: &FsBlobStore,
83    key: BlobKey,
84    len: u64,
85    part_size: u64,
86    session: &[u8],
87) -> Result<Vec<u8>, StoreError> {
88    PartPlan::new(len, part_size, u32::MAX)
89        .map_err(|e| StoreError::Invalid(e.to_string().into()))?;
90    let expected = session_meta(key, len, part_size)?;
91    let dir = session_dir(store, session)?;
92    let parent = uploads(&store.root);
93    let lock = session_lock(store, session)?;
94    let _guard = lock.lock().await;
95    create_dir_all_durably(&parent).map_err(io_error)?;
96    match fs::symlink_metadata(&dir) {
97        Ok(info) if !info.is_dir() => return Err(StoreError::SessionGone),
98        Ok(_) => match fs::read(dir.join(META)) {
99            Ok(actual) if actual == expected => return Ok(session.to_vec()),
100            Ok(_) => {}
101            Err(e) if e.kind() == ErrorKind::NotFound => {}
102            Err(e) => return Err(io_error(e)),
103        },
104        Err(e) if e.kind() == ErrorKind::NotFound => {}
105        Err(e) => return Err(io_error(e)),
106    }
107    if dir.exists() {
108        // An attempt that never issued a ticket may leave an incomplete or
109        // mismatched directory, including after a part-size configuration
110        // change. This open reservation replaces that orphan.
111        fs::remove_dir_all(&dir).map_err(io_error)?;
112        sync_dir(&parent).map_err(io_error)?;
113    }
114    fs::create_dir(&dir).map_err(io_error)?;
115    let dest = dir.join(META);
116    let tmp = temp_path(&dest).map_err(io_error)?;
117    let result = (|| {
118        sync_dir(&parent)?;
119        let mut file = OpenOptions::new().write(true).create_new(true).open(&tmp)?;
120        file.write_all(&expected)?;
121        file.sync_all()?;
122        fs::rename(&tmp, &dest)?;
123        sync_dir(&dir)?;
124        Ok::<(), io::Error>(())
125    })();
126    if let Err(e) = result {
127        let _ = fs::remove_file(&tmp);
128        let _ = fs::remove_dir_all(&dir);
129        let _ = sync_dir(&parent);
130        return Err(io_error(e));
131    }
132    Ok(session.to_vec())
133}
134
135fn fresh_session(key: BlobKey) -> [u8; 32] {
136    let now = SystemTime::now()
137        .duration_since(UNIX_EPOCH)
138        .unwrap_or_default();
139    let mut material = Vec::with_capacity(64);
140    material.extend_from_slice(key.hash());
141    material.extend_from_slice(&now.as_nanos().to_be_bytes());
142    material.extend_from_slice(&std::process::id().to_be_bytes());
143    material.extend_from_slice(&NEXT_SESSION.fetch_add(1, Ordering::Relaxed).to_be_bytes());
144    hash(&material)
145}
146
147fn part_name(index: u32, cv: &[u8; 32]) -> String {
148    format!("{index}-{}", to_hex_bytes(cv))
149}
150
151fn current_path(dir: &Path, index: u32) -> PathBuf {
152    dir.join(format!("{index}.current"))
153}
154
155fn publish_current(dir: &Path, index: u32, name: &str) -> Result<(), StoreError> {
156    let dest = current_path(dir, index);
157    let tmp = temp_path(&dest).map_err(io_error)?;
158    let result = (|| {
159        let mut file = OpenOptions::new().write(true).create_new(true).open(&tmp)?;
160        file.write_all(name.as_bytes())?;
161        file.sync_all()?;
162        fs::rename(&tmp, &dest)?;
163        sync_dir(dir)?;
164        Ok::<(), io::Error>(())
165    })();
166    if result.is_err() {
167        let _ = fs::remove_file(&tmp);
168    }
169    result.map_err(io_error)
170}
171
172fn is_part_name(name: &str, index: u32) -> bool {
173    name.strip_prefix(&format!("{index}-")).is_some_and(|cv| {
174        cv.len() == 64
175            && cv
176                .bytes()
177                .all(|b| b.is_ascii_hexdigit() && !b.is_ascii_uppercase())
178    })
179}
180
181/// A staged file and incremental subtree hasher. Dropping it removes only
182/// this attempt's temp file; the prior verified part remains intact.
183pub struct FsPartSink {
184    session_lock: std::sync::Arc<tokio::sync::Mutex<()>>,
185    file: Option<File>,
186    tmp: PathBuf,
187    dir: PathBuf,
188    name: String,
189    meta: Vec<u8>,
190    index: u32,
191    expected_cv: [u8; 32],
192    hasher: Option<PartHasher>,
193}
194
195impl std::fmt::Debug for FsPartSink {
196    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
197        f.debug_struct("FsPartSink")
198            .field("tmp", &self.tmp)
199            .field("index", &self.index)
200            .finish_non_exhaustive()
201    }
202}
203
204impl Drop for FsPartSink {
205    fn drop(&mut self) {
206        self.file = None;
207        let _ = fs::remove_file(&self.tmp);
208    }
209}
210
211impl PartSink for FsPartSink {
212    async fn write(&mut self, chunk: Bytes) -> Result<(), StoreError> {
213        if chunk.is_empty() {
214            return Err(StoreError::Invalid("empty part chunk".into()));
215        }
216        let hasher = self.hasher.as_mut().ok_or(StoreError::SessionGone)?;
217        hasher
218            .update(&chunk)
219            .map_err(|e| StoreError::Invalid(e.to_string().into()))?;
220        if let Err(e) = self
221            .file
222            .as_mut()
223            .ok_or(StoreError::SessionGone)?
224            .write_all(&chunk)
225        {
226            self.file = None;
227            return Err(io_error(e));
228        }
229        Ok(())
230    }
231
232    async fn commit(mut self) -> Result<Vec<u8>, StoreError> {
233        let cv = self
234            .hasher
235            .take()
236            .ok_or(StoreError::SessionGone)?
237            .finalize()
238            .map_err(|e| StoreError::Invalid(e.to_string().into()))?;
239        if cv != self.expected_cv {
240            return Err(StoreError::PartSubtreeMismatch);
241        }
242        let file = self.file.take().ok_or(StoreError::SessionGone)?;
243        file.sync_all().map_err(io_error)?;
244        drop(file);
245        let _guard = self.session_lock.lock().await;
246        verify_session(&self.dir, &self.meta)?;
247        // The new file is durable before the pointer changes. A crash before
248        // that change leaves the old receipt valid; a crash after it leaves
249        // the new receipt valid, even if old files remain for cleanup.
250        fs::rename(&self.tmp, self.dir.join(&self.name)).map_err(io_error)?;
251        sync_dir(&self.dir).map_err(io_error)?;
252        publish_current(&self.dir, self.index, &self.name)?;
253        for entry in fs::read_dir(&self.dir).map_err(io_error)? {
254            let entry = entry.map_err(io_error)?;
255            if entry.file_type().map_err(io_error)?.is_file()
256                && entry
257                    .file_name()
258                    .to_str()
259                    .is_some_and(|name| name != self.name && is_part_name(name, self.index))
260                && let Err(e) = fs::remove_file(entry.path())
261            {
262                tracing::warn!(error = %e, "old multipart part cleanup failed");
263            }
264        }
265        sync_dir(&self.dir).map_err(io_error)?;
266        Ok(self.name.as_bytes().to_vec())
267    }
268
269    async fn abort(self) {}
270}
271
272impl MultipartBlobStore for FsBlobStore {
273    type PartSink = FsPartSink;
274    const MAX_PARTS: u32 = u32::MAX;
275
276    fn supports_multipart(&self) -> bool {
277        true
278    }
279
280    async fn begin_multipart(
281        &self,
282        key: BlobKey,
283        len: u64,
284        part_size: u64,
285    ) -> Result<Vec<u8>, StoreError> {
286        for _ in 0..8 {
287            let session = fresh_session(key);
288            if !session_dir(self, &session)?.exists() {
289                return create_session(self, key, len, part_size, &session).await;
290            }
291        }
292        Err(StoreError::unavailable(io::Error::other(
293            "could not allocate a multipart session",
294        )))
295    }
296
297    async fn begin_multipart_for_ticket(
298        &self,
299        key: BlobKey,
300        len: u64,
301        part_size: u64,
302        ticket_id: [u8; 32],
303    ) -> Result<Vec<u8>, StoreError> {
304        create_session(self, key, len, part_size, &ticket_id).await
305    }
306
307    async fn begin_part(
308        &self,
309        key: BlobKey,
310        session: &[u8],
311        plan: &PartPlan,
312        index: u32,
313        expected_cv: [u8; 32],
314    ) -> Result<FsPartSink, StoreError> {
315        let meta = session_meta(key, plan.total(), plan.part_size())?;
316        let dir = session_dir(self, session)?;
317        let hasher =
318            PartHasher::new(plan, index).map_err(|e| StoreError::Invalid(e.to_string().into()))?;
319        let name = part_name(index, &expected_cv);
320        let session_lock = session_lock(self, session)?;
321        let _guard = session_lock.lock().await;
322        verify_session(&dir, &meta)?;
323        let tmp = temp_path(&dir.join(&name)).map_err(io_error)?;
324        let file = OpenOptions::new()
325            .write(true)
326            .create_new(true)
327            .open(&tmp)
328            .map_err(io_error)?;
329        Ok(FsPartSink {
330            session_lock: session_lock.clone(),
331            file: Some(file),
332            tmp,
333            dir,
334            name,
335            meta,
336            index,
337            expected_cv,
338            hasher: Some(hasher),
339        })
340    }
341
342    async fn complete(
343        &self,
344        key: BlobKey,
345        session: &[u8],
346        plan: &PartPlan,
347        parts: &[PartRef],
348    ) -> Result<CommitOutcome, StoreError> {
349        self.complete_with(key, session, plan, parts, None).await
350    }
351
352    async fn complete_with_root(
353        &self,
354        key: BlobKey,
355        session: &[u8],
356        plan: &PartPlan,
357        parts: &[PartRef],
358        content_root: Hash,
359    ) -> Result<CommitOutcome, StoreError> {
360        self.complete_with(key, session, plan, parts, Some(content_root))
361            .await
362    }
363
364    async fn abort(&self, key: BlobKey, session: &[u8]) -> Result<(), StoreError> {
365        self.abort_session(key, session).await
366    }
367}
368
369impl FsBlobStore {
370    /// Complete a multipart upload against the key's hash (`None`) or an
371    /// object's content root.
372    async fn complete_with(
373        &self,
374        key: BlobKey,
375        session: &[u8],
376        plan: &PartPlan,
377        parts: &[PartRef],
378        root: Option<Hash>,
379    ) -> Result<CommitOutcome, StoreError> {
380        let expected = key.expected_root(root)?;
381        let meta = session_meta(key, plan.total(), plan.part_size())?;
382        let dir = session_dir(self, session)?;
383        let session_lock = session_lock(self, session)?;
384        let _guard = session_lock.lock().await;
385        verify_session(&dir, &meta)?;
386        if parts.len() != plan.count() as usize {
387            return Err(StoreError::Invalid("wrong number of parts".into()));
388        }
389        let mut cvs = Vec::with_capacity(parts.len());
390        for (position, part) in parts.iter().enumerate() {
391            let index = u32::try_from(position)
392                .map_err(|_| StoreError::Invalid("part index overflow".into()))?;
393            let prefix = format!("{index}-");
394            let tag = std::str::from_utf8(&part.tag)
395                .map_err(|_| StoreError::Invalid("invalid part tag".into()))?;
396            let cv_hex = tag
397                .strip_prefix(&prefix)
398                .filter(|_| is_part_name(tag, index))
399                .ok_or_else(|| StoreError::Invalid("invalid part tag".into()))?;
400            let cv = mkit_core::hash::from_hex(cv_hex)
401                .map_err(|_| StoreError::Invalid("invalid part tag".into()))?;
402            if part.index != index
403                || part.len
404                    != plan
405                        .expected_len(index)
406                        .map_err(|e| StoreError::Invalid(e.to_string().into()))?
407            {
408                return Err(StoreError::Invalid("part geometry mismatch".into()));
409            }
410            cvs.push(cv);
411        }
412        if merge_to_root(plan, &cvs).map_err(|e| StoreError::Invalid(e.to_string().into()))?
413            != expected
414        {
415            return Err(StoreError::Invalid("merged part root mismatch".into()));
416        }
417        let mut sink = self.begin(key, plan.total()).await?;
418        let mut buf = BytesMut::with_capacity(READ_BLOCK);
419        for part in parts {
420            let tag = std::str::from_utf8(&part.tag)
421                .map_err(|_| StoreError::Invalid("invalid part tag".into()))?;
422            let active = match fs::read(current_path(&dir, part.index)) {
423                Ok(active) => active,
424                Err(e) if e.kind() == ErrorKind::NotFound => {
425                    verify_session(&dir, &meta)?;
426                    return Err(StoreError::Invalid("part tag is not current".into()));
427                }
428                Err(e) => return Err(io_error(e)),
429            };
430            if active != part.tag {
431                return Err(StoreError::Invalid("part tag is not current".into()));
432            }
433            let mut file = match File::open(dir.join(tag)) {
434                Ok(file) => file,
435                Err(e) if e.kind() == ErrorKind::NotFound => {
436                    verify_session(&dir, &meta)?;
437                    // The session exists, but this receipt's tag no longer
438                    // names its selected part (for example after replacement).
439                    return Err(StoreError::Invalid("part tag has no stored file".into()));
440                }
441                Err(e) => return Err(io_error(e)),
442            };
443            if file.metadata().map_err(io_error)?.len() != part.len {
444                return Err(StoreError::Invalid("stored part length mismatch".into()));
445            }
446            let mut remaining = part.len;
447            while remaining > 0 {
448                let n = usize::try_from(remaining).map_or(READ_BLOCK, |n| n.min(READ_BLOCK));
449                buf.resize(n, 0);
450                file.read_exact(&mut buf[..n]).map_err(io_error)?;
451                sink.write(buf.split().freeze()).await?;
452                remaining -= n as u64;
453            }
454        }
455        let outcome = match root {
456            Some(root) => sink.commit_with_root(root).await?,
457            None => sink.commit().await?,
458        };
459        // The pack is durable. Cleanup failure is harmless; startup sweep
460        // reclaims the directory after the ticket lifetime.
461        if let Err(e) = fs::remove_dir_all(&dir).and_then(|()| sync_dir(&uploads(&self.root))) {
462            tracing::warn!(error = %e, "completed multipart session cleanup failed");
463        }
464        Ok(outcome)
465    }
466
467    async fn abort_session(&self, key: BlobKey, session: &[u8]) -> Result<(), StoreError> {
468        let dir = session_dir(self, session)?;
469        let session_lock = session_lock(self, session)?;
470        let _guard = session_lock.lock().await;
471        let meta = match fs::read(dir.join(META)) {
472            Ok(meta) => meta,
473            Err(e) if e.kind() == ErrorKind::NotFound => return Ok(()),
474            Err(e) => return Err(io_error(e)),
475        };
476        if meta.get(5..37) != Some(key.hash().as_slice()) {
477            return Ok(());
478        }
479        fs::remove_dir_all(&dir).map_err(io_error)?;
480        sync_dir(&uploads(&self.root)).map_err(io_error)
481    }
482}
483
484pub(super) fn sweep_sessions(root: &Path, now: SystemTime) -> io::Result<usize> {
485    let parent = uploads(root);
486    let entries = match fs::read_dir(&parent) {
487        Ok(entries) => entries,
488        Err(e) if e.kind() == ErrorKind::NotFound => return Ok(0),
489        Err(e) => return Err(e),
490    };
491    let mut removed = 0;
492    for entry in entries {
493        let Ok(entry) = entry else { continue };
494        let name = entry.file_name();
495        let Some(name) = name.to_str() else { continue };
496        if name.len() != 64
497            || !name
498                .bytes()
499                .all(|b| b.is_ascii_hexdigit() && !b.is_ascii_uppercase())
500        {
501            continue;
502        }
503        let Ok(meta) = fs::symlink_metadata(entry.path()) else {
504            continue;
505        };
506        if !meta.is_dir() {
507            continue;
508        }
509        // The immutable metadata file records session creation. Directory
510        // mtime changes as parts arrive, but cannot extend a ticket's life.
511        let created = fs::symlink_metadata(entry.path().join(META))
512            .ok()
513            .filter(fs::Metadata::is_file)
514            .unwrap_or(meta);
515        if created
516            .modified()
517            .ok()
518            .and_then(|modified| now.duration_since(modified).ok())
519            .is_none_or(|age| age < SESSION_AGE)
520        {
521            continue;
522        }
523        if fs::remove_dir_all(entry.path()).is_ok() {
524            removed += 1;
525        }
526    }
527    if removed > 0 {
528        sync_dir(&parent)?;
529    }
530    Ok(removed)
531}