diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index fca5282ad..8bf9ec340 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -17,9 +17,9 @@ use crate::data_usage::local_snapshot::ensure_data_usage_layout; use crate::disk::{ BUCKET_META_PREFIX, CHECK_PART_FILE_CORRUPT, CHECK_PART_FILE_NOT_FOUND, CHECK_PART_SUCCESS, CHECK_PART_UNKNOWN, CHECK_PART_VOLUME_NOT_FOUND, CheckPartsResp, DeleteOptions, DiskAPI, DiskInfo, DiskInfoOptions, DiskLocation, DiskMetrics, - FileInfoVersions, FileReader, FileWriter, RUSTFS_META_BUCKET, RUSTFS_META_TMP_DELETED_BUCKET, ReadMultipleReq, - ReadMultipleResp, ReadOptions, RenameDataResp, STORAGE_FORMAT_FILE, STORAGE_FORMAT_FILE_BACKUP, UpdateMetadataOpts, - VolumeInfo, WalkDirOptions, conv_part_err_to_int, + FileInfoVersions, FileReader, FileWriter, RUSTFS_META_BUCKET, RUSTFS_META_TMP_BUCKET, RUSTFS_META_TMP_DELETED_BUCKET, + ReadMultipleReq, ReadMultipleResp, ReadOptions, RenameDataResp, STORAGE_FORMAT_FILE, STORAGE_FORMAT_FILE_BACKUP, + UpdateMetadataOpts, VolumeInfo, WalkDirOptions, conv_part_err_to_int, endpoint::Endpoint, error::{DiskError, Error, FileAccessDeniedWithContext, Result}, error_conv::{to_access_error, to_file_error, to_unformatted_disk_error, to_volume_error}, @@ -61,6 +61,10 @@ use tokio::time::interval; use tracing::{debug, error, info, warn}; use uuid::Uuid; +const DELETED_OBJECTS_CLEANUP_INTERVAL: Duration = Duration::from_secs(60 * 5); +const STALE_TMP_OBJECT_EXPIRY: Duration = Duration::from_secs(24 * 60 * 60); +const RUSTFS_META_TMP_OLD_BUCKET: &str = ".rustfs.sys/tmp-old"; + #[derive(Debug, Clone)] pub struct FormatInfo { pub id: Option, @@ -133,8 +137,8 @@ impl LocalDisk { ensure_data_usage_layout(&root).await.map_err(DiskError::from)?; - if cleanup { - // TODO: remove temporary data + if cleanup && let Err(err) = Self::cleanup_tmp_on_startup(&root).await { + warn!(root = ?root, error = ?err, "failed to cleanup temporary data during disk startup"); } // Use optimized path resolution instead of absolutize_virtually @@ -251,13 +255,16 @@ impl LocalDisk { } async fn cleanup_deleted_objects_loop(root: PathBuf, mut exit_rx: tokio::sync::broadcast::Receiver<()>) { - let mut interval = interval(Duration::from_secs(60 * 5)); + let mut interval = interval(DELETED_OBJECTS_CLEANUP_INTERVAL); loop { tokio::select! { _ = interval.tick() => { if let Err(err) = Self::cleanup_deleted_objects(root.clone()).await { error!("cleanup_deleted_objects error: {:?}", err); } + if let Err(err) = Self::cleanup_stale_tmp_objects(root.clone()).await { + error!("cleanup_stale_tmp_objects error: {:?}", err); + } } _ = exit_rx.recv() => { info!("cleanup_deleted_objects_loop exit"); @@ -267,13 +274,83 @@ impl LocalDisk { } } - async fn cleanup_deleted_objects(root: PathBuf) -> Result<()> { + fn meta_path(root: &Path, meta_path: &str) -> PathBuf { #[cfg(windows)] - let trash_path = RUSTFS_META_TMP_DELETED_BUCKET.replace('/', "\\"); + let meta_path = meta_path.replace('/', "\\"); #[cfg(not(windows))] - let trash_path = RUSTFS_META_TMP_DELETED_BUCKET.to_string(); + let meta_path = meta_path.to_string(); - let trash = root.join(trash_path); + root.join(meta_path) + } + + async fn cleanup_tmp_on_startup(root: &Path) -> Result<()> { + let tmp_path = Self::meta_path(root, RUSTFS_META_TMP_BUCKET); + let tmp_old_path = Self::meta_path(root, RUSTFS_META_TMP_OLD_BUCKET).join(Uuid::new_v4().to_string()); + + rename_all(&tmp_path, &tmp_old_path, root).await?; + + let tmp_old_root = Self::meta_path(root, RUSTFS_META_TMP_OLD_BUCKET); + tokio::spawn(async move { + if let Err(err) = tokio::fs::remove_dir_all(&tmp_old_root).await + && err.kind() != ErrorKind::NotFound + { + warn!(path = ?tmp_old_root, error = ?err, "failed to remove old temporary data"); + } + }); + + tokio::fs::create_dir_all(Self::meta_path(root, RUSTFS_META_TMP_DELETED_BUCKET)).await?; + Ok(()) + } + + async fn cleanup_stale_tmp_objects(root: PathBuf) -> Result<()> { + Self::cleanup_stale_tmp_objects_with_expiry(root, STALE_TMP_OBJECT_EXPIRY).await + } + + async fn cleanup_stale_tmp_objects_with_expiry(root: PathBuf, expiry: Duration) -> Result<()> { + let tmp_path = Self::meta_path(&root, RUSTFS_META_TMP_BUCKET); + let mut entries = match fs::read_dir(&tmp_path).await { + Ok(entries) => entries, + Err(e) => { + if e.kind() == ErrorKind::NotFound { + return Ok(()); + } + return Err(e.into()); + } + }; + + while let Some(entry) = entries.next_entry().await? { + let name = entry.file_name().to_string_lossy().to_string(); + if name.is_empty() || name == "." || name == ".." || name == ".trash" { + continue; + } + + let file_type = entry.file_type().await?; + if !file_type.is_dir() { + continue; + } + + let Some(age) = entry + .metadata() + .await? + .modified() + .ok() + .and_then(|modified| modified.elapsed().ok()) + else { + continue; + }; + if age <= expiry { + continue; + } + + let target_path = Self::meta_path(&root, RUSTFS_META_TMP_DELETED_BUCKET).join(Uuid::new_v4().to_string()); + rename_all(entry.path(), target_path, Self::meta_path(&root, RUSTFS_META_BUCKET)).await?; + } + + Ok(()) + } + + async fn cleanup_deleted_objects(root: PathBuf) -> Result<()> { + let trash = Self::meta_path(&root, RUSTFS_META_TMP_DELETED_BUCKET); let mut entries = match fs::read_dir(&trash).await { Ok(entries) => entries, Err(e) => { @@ -2721,6 +2798,46 @@ mod test { } } + #[tokio::test] + async fn cleanup_tmp_on_startup_moves_existing_tmp_and_recreates_trash() { + use tempfile::tempdir; + + let dir = tempdir().unwrap(); + let tmp = LocalDisk::meta_path(dir.path(), RUSTFS_META_TMP_BUCKET); + let leftover = tmp.join("leftover").join("data"); + fs::create_dir_all(leftover.parent().unwrap()).await.unwrap(); + fs::write(&leftover, b"temporary").await.unwrap(); + + LocalDisk::cleanup_tmp_on_startup(dir.path()).await.unwrap(); + + assert!(!tmp.join("leftover").exists()); + assert!(LocalDisk::meta_path(dir.path(), RUSTFS_META_TMP_DELETED_BUCKET).exists()); + } + + #[tokio::test] + async fn cleanup_stale_tmp_objects_moves_expired_tmp_dirs_to_trash() { + use tempfile::tempdir; + + let dir = tempdir().unwrap(); + let tmp = LocalDisk::meta_path(dir.path(), RUSTFS_META_TMP_BUCKET); + let stale = tmp.join("stale").join("data"); + let trash = LocalDisk::meta_path(dir.path(), RUSTFS_META_TMP_DELETED_BUCKET); + fs::create_dir_all(stale.parent().unwrap()).await.unwrap(); + fs::create_dir_all(&trash).await.unwrap(); + fs::write(&stale, b"temporary").await.unwrap(); + + tokio::time::sleep(Duration::from_millis(2)).await; + LocalDisk::cleanup_stale_tmp_objects_with_expiry(dir.path().to_path_buf(), Duration::ZERO) + .await + .unwrap(); + + assert!(!tmp.join("stale").exists()); + assert!(trash.exists()); + + let mut entries = fs::read_dir(&trash).await.unwrap(); + assert!(entries.next_entry().await.unwrap().is_some()); + } + #[tokio::test] async fn test_scan_dir_includes_nested_object_dirs() { use rustfs_filemeta::MetacacheReader; diff --git a/crates/ecstore/src/set_disk.rs b/crates/ecstore/src/set_disk.rs index 023a4c726..1fcbe1e0d 100644 --- a/crates/ecstore/src/set_disk.rs +++ b/crates/ecstore/src/set_disk.rs @@ -809,198 +809,205 @@ impl ObjectIO for SetDisks { let tmp_object = format!("{}/{}/part.1", tmp_dir, fi.data_dir.unwrap()); - let erasure = erasure_coding::Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size); + let result: Result = async { + let erasure = erasure_coding::Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks, fi.erasure.block_size); - let is_inline_buffer = { - if let Some(sc) = GLOBAL_STORAGE_CLASS.get() { - sc.should_inline(erasure.shard_file_size(data.size()), opts.versioned) - } else { - false - } - }; + let is_inline_buffer = { + if let Some(sc) = GLOBAL_STORAGE_CLASS.get() { + sc.should_inline(erasure.shard_file_size(data.size()), opts.versioned) + } else { + false + } + }; - let mut writers = Vec::with_capacity(shuffle_disks.len()); - let mut errors = Vec::with_capacity(shuffle_disks.len()); - for disk_op in shuffle_disks.iter() { - if let Some(disk) = disk_op - && disk.is_online().await - { - let writer = match create_bitrot_writer( - is_inline_buffer, - Some(disk), - RUSTFS_META_TMP_BUCKET, - &tmp_object, - erasure.shard_file_size(data.size()), - erasure.shard_size(), - HashAlgorithm::HighwayHash256S, - ) - .await + let mut writers = Vec::with_capacity(shuffle_disks.len()); + let mut errors = Vec::with_capacity(shuffle_disks.len()); + for disk_op in shuffle_disks.iter() { + if let Some(disk) = disk_op + && disk.is_online().await { - Ok(writer) => writer, - Err(err) => { - warn!("create_bitrot_writer disk {}, err {:?}, skipping operation", disk.to_string(), err); - errors.push(Some(err)); - writers.push(None); - continue; - } - }; + let writer = match create_bitrot_writer( + is_inline_buffer, + Some(disk), + RUSTFS_META_TMP_BUCKET, + &tmp_object, + erasure.shard_file_size(data.size()), + erasure.shard_size(), + HashAlgorithm::HighwayHash256S, + ) + .await + { + Ok(writer) => writer, + Err(err) => { + warn!("create_bitrot_writer disk {}, err {:?}, skipping operation", disk.to_string(), err); + errors.push(Some(err)); + writers.push(None); + continue; + } + }; - writers.push(Some(writer)); - errors.push(None); - } else { - errors.push(Some(DiskError::DiskNotFound)); - writers.push(None); - } - } - - let nil_count = errors.iter().filter(|&e| e.is_none()).count(); - if nil_count < write_quorum { - error!("not enough disks to write: {:?}", errors); - if let Some(write_err) = reduce_write_quorum_errs(&errors, OBJECT_OP_IGNORED_ERRS, write_quorum) { - return Err(to_object_err(write_err.into(), vec![bucket, object])); + writers.push(Some(writer)); + errors.push(None); + } else { + errors.push(Some(DiskError::DiskNotFound)); + writers.push(None); + } } - return Err(Error::other(format!("not enough disks to write: {errors:?}"))); - } - - let stream = mem::replace( - &mut data.stream, - HashReader::from_stream(Cursor::new(Vec::new()), 0, 0, None, None, false)?, - ); - - let (reader, w_size) = match Arc::new(erasure).encode(stream, &mut writers, write_quorum).await { - Ok((r, w)) => (r, w), - Err(e) => { - error!("encode err {:?}", e); - return Err(e.into()); - } - }; // TODO: delete temporary directory on error - - let _ = mem::replace(&mut data.stream, reader); - // if let Err(err) = close_bitrot_writers(&mut writers).await { - // error!("close_bitrot_writers err {:?}", err); - // } - - if (w_size as i64) < data.size() { - warn!("put_object write size < data.size(), w_size={}, data.size={}", w_size, data.size()); - return Err(Error::other(format!( - "put_object write size < data.size(), w_size={}, data.size={}", - w_size, - data.size() - ))); - } - - if contains_key_str(&user_defined, SUFFIX_COMPRESSION) { - insert_str(&mut user_defined, SUFFIX_COMPRESSION_SIZE, w_size.to_string()); - } - - let index_op = data.stream.try_get_index().map(|v| v.clone().into_vec()); - - //TODO: userDefined - - let etag = data.stream.try_resolve_etag().unwrap_or_default(); - - user_defined.insert("etag".to_owned(), etag.clone()); - - if !user_defined.contains_key("content-type") { - // get content-type - } - - let mut actual_size = data.actual_size(); - if actual_size < 0 { - let is_compressed = fi.is_compressed(); - if !is_compressed { - actual_size = w_size as i64; - } - } - - if fi.checksum.is_none() - && let Some(content_hash) = data.as_hash_reader().content_hash() - { - fi.checksum = Some(content_hash.to_bytes(&[])); - } - - if let Some(sc) = user_defined.get(AMZ_STORAGE_CLASS) - && sc == storageclass::STANDARD - { - let _ = user_defined.remove(AMZ_STORAGE_CLASS); - } - - let mod_time = if let Some(mod_time) = opts.mod_time { - Some(mod_time) - } else { - Some(OffsetDateTime::now_utc()) - }; - - for (i, pfi) in parts_metadatas.iter_mut().enumerate() { - pfi.metadata = user_defined.clone(); - if is_inline_buffer { - if let Some(writer) = writers[i].take() { - pfi.data = Some(writer.into_inline_data().map(Bytes::from).unwrap_or_default()); + let nil_count = errors.iter().filter(|&e| e.is_none()).count(); + if nil_count < write_quorum { + error!("not enough disks to write: {:?}", errors); + if let Some(write_err) = reduce_write_quorum_errs(&errors, OBJECT_OP_IGNORED_ERRS, write_quorum) { + return Err(to_object_err(write_err.into(), vec![bucket, object])); } - pfi.set_inline_data(); + return Err(Error::other(format!("not enough disks to write: {errors:?}"))); } - pfi.mod_time = mod_time; - pfi.size = w_size as i64; - pfi.versioned = opts.versioned || opts.version_suspended; - pfi.add_object_part(1, etag.clone(), w_size, mod_time, actual_size, index_op.clone(), None); - pfi.checksum = fi.checksum.clone(); + let stream = mem::replace( + &mut data.stream, + HashReader::from_stream(Cursor::new(Vec::new()), 0, 0, None, None, false)?, + ); - if opts.data_movement { - pfi.set_data_moved(); + let (reader, w_size) = match Arc::new(erasure).encode(stream, &mut writers, write_quorum).await { + Ok((r, w)) => (r, w), + Err(e) => { + error!("encode err {:?}", e); + return Err(e.into()); + } + }; // TODO: delete temporary directory on error + + let _ = mem::replace(&mut data.stream, reader); + // if let Err(err) = close_bitrot_writers(&mut writers).await { + // error!("close_bitrot_writers err {:?}", err); + // } + + if (w_size as i64) < data.size() { + warn!("put_object write size < data.size(), w_size={}, data.size={}", w_size, data.size()); + return Err(Error::other(format!( + "put_object write size < data.size(), w_size={}, data.size={}", + w_size, + data.size() + ))); } - } - drop(writers); // drop writers to close all files, this is to prevent FileAccessDenied errors when renaming data + if contains_key_str(&user_defined, SUFFIX_COMPRESSION) { + insert_str(&mut user_defined, SUFFIX_COMPRESSION_SIZE, w_size.to_string()); + } - if !opts.no_lock && object_lock_guard.is_none() { - let ns_lock = self.new_ns_lock(bucket, object).await?; - object_lock_guard = Some(ns_lock.get_write_lock(get_lock_acquire_timeout()).await.map_err(|e| { - StorageError::other(format!( - "Failed to acquire write lock: {}", - self.format_lock_error_from_error(bucket, object, "write", &e) - )) - })?); - } + let index_op = data.stream.try_get_index().map(|v| v.clone().into_vec()); - let (online_disks, _, op_old_dir) = Self::rename_data( - &shuffle_disks, - RUSTFS_META_TMP_BUCKET, - tmp_dir.as_str(), - &parts_metadatas, - bucket, - object, - write_quorum, - ) - .await?; + //TODO: userDefined - if let Some(old_dir) = op_old_dir { - self.commit_rename_data_dir(&online_disks, bucket, object, &old_dir.to_string(), write_quorum) - .await?; - } + let etag = data.stream.try_resolve_etag().unwrap_or_default(); - drop(object_lock_guard); // drop object lock guard to release the lock + user_defined.insert("etag".to_owned(), etag.clone()); - self.delete_all(RUSTFS_META_TMP_BUCKET, &tmp_dir).await?; + if !user_defined.contains_key("content-type") { + // get content-type + } - for (i, op_disk) in online_disks.iter().enumerate() { - if let Some(disk) = op_disk - && disk.is_online().await + let mut actual_size = data.actual_size(); + if actual_size < 0 { + let is_compressed = fi.is_compressed(); + if !is_compressed { + actual_size = w_size as i64; + } + } + + if fi.checksum.is_none() + && let Some(content_hash) = data.as_hash_reader().content_hash() { - fi = parts_metadatas[i].clone(); - break; + fi.checksum = Some(content_hash.to_bytes(&[])); } + + if let Some(sc) = user_defined.get(AMZ_STORAGE_CLASS) + && sc == storageclass::STANDARD + { + let _ = user_defined.remove(AMZ_STORAGE_CLASS); + } + + let mod_time = if let Some(mod_time) = opts.mod_time { + Some(mod_time) + } else { + Some(OffsetDateTime::now_utc()) + }; + + for (i, pfi) in parts_metadatas.iter_mut().enumerate() { + pfi.metadata = user_defined.clone(); + if is_inline_buffer { + if let Some(writer) = writers[i].take() { + pfi.data = Some(writer.into_inline_data().map(Bytes::from).unwrap_or_default()); + } + + pfi.set_inline_data(); + } + + pfi.mod_time = mod_time; + pfi.size = w_size as i64; + pfi.versioned = opts.versioned || opts.version_suspended; + pfi.add_object_part(1, etag.clone(), w_size, mod_time, actual_size, index_op.clone(), None); + pfi.checksum = fi.checksum.clone(); + + if opts.data_movement { + pfi.set_data_moved(); + } + } + + drop(writers); // drop writers to close all files, this is to prevent FileAccessDenied errors when renaming data + + if !opts.no_lock && object_lock_guard.is_none() { + let ns_lock = self.new_ns_lock(bucket, object).await?; + object_lock_guard = Some(ns_lock.get_write_lock(get_lock_acquire_timeout()).await.map_err(|e| { + StorageError::other(format!( + "Failed to acquire write lock: {}", + self.format_lock_error_from_error(bucket, object, "write", &e) + )) + })?); + } + + let (online_disks, _, op_old_dir) = Self::rename_data( + &shuffle_disks, + RUSTFS_META_TMP_BUCKET, + tmp_dir.as_str(), + &parts_metadatas, + bucket, + object, + write_quorum, + ) + .await?; + + if let Some(old_dir) = op_old_dir { + self.commit_rename_data_dir(&online_disks, bucket, object, &old_dir.to_string(), write_quorum) + .await?; + } + + drop(object_lock_guard); // drop object lock guard to release the lock + + for (i, op_disk) in online_disks.iter().enumerate() { + if let Some(disk) = op_disk + && disk.is_online().await + { + fi = parts_metadatas[i].clone(); + break; + } + } + + record_capacity_scope_if_needed(opts.capacity_scope_token, &online_disks); + + fi.replication_state_internal = Some(opts.put_replication_state()); + + fi.is_latest = true; + + Ok(ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended)) + } + .await; + + if let Err(err) = self.delete_all(RUSTFS_META_TMP_BUCKET, &tmp_dir).await { + warn!(tmp_dir = %tmp_dir, error = ?err, "failed to cleanup put_object temporary data"); } - record_capacity_scope_if_needed(opts.capacity_scope_token, &online_disks); - - fi.replication_state_internal = Some(opts.put_replication_state()); - - fi.is_latest = true; - - Ok(ObjectInfo::from_file_info(&fi, bucket, object, opts.versioned || opts.version_suspended)) + result } }