1use 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 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
181pub 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 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 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 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 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 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}