diff --git a/ecstore/src/disk/mod.rs b/ecstore/src/disk/mod.rs index 9a961c112..000d0d0dd 100644 --- a/ecstore/src/disk/mod.rs +++ b/ecstore/src/disk/mod.rs @@ -18,7 +18,7 @@ use crate::{ bucket::{metadata_sys::get_versioning_config, versioning::VersioningApi}, erasure::Writer, error::{Error, Result}, - file_meta::{merge_file_meta_versions, FileMeta, FileMetaShallowVersion}, + file_meta::{merge_file_meta_versions, FileMeta, FileMetaShallowVersion, VersionType}, heal::{ data_scanner::ShouldSleepFn, data_usage_cache::{DataUsageCache, DataUsageEntry}, @@ -629,6 +629,39 @@ impl MetaCacheEntry { !self.metadata.is_empty() } + pub fn is_latest_deletemarker(&mut self) -> bool { + if let Some(cached) = &self.cached { + if cached.versions.is_empty() { + return true; + } + + return cached.versions[0].header.version_type == VersionType::Delete; + } + + if !FileMeta::is_xl2_v1_format(&self.metadata) { + return false; + } + + match FileMeta::check_xl2_v1(&self.metadata) { + Ok((meta, _, _)) => { + if !meta.is_empty() { + // TODO: IsLatestDeleteMarker + } + } + Err(_) => return true, + } + + match self.xl_meta() { + Ok(res) => { + if res.versions.is_empty() { + return true; + } + res.versions[0].header.version_type == VersionType::Delete + } + Err(_) => true, + } + } + #[tracing::instrument(level = "debug", skip(self))] pub fn to_fileinfo(&self, bucket: &str) -> Result> { if self.is_dir() { diff --git a/ecstore/src/store_list_objects.rs b/ecstore/src/store_list_objects.rs index 2102eefdf..6db2c5cb9 100644 --- a/ecstore/src/store_list_objects.rs +++ b/ecstore/src/store_list_objects.rs @@ -17,7 +17,7 @@ use std::io::ErrorKind; use std::sync::Arc; use tokio::sync::broadcast::{self, Receiver as B_Receiver}; use tokio::sync::mpsc::{self, Receiver, Sender}; -use tracing::error; +use tracing::{error, warn}; const MAX_OBJECT_LIST: i32 = 1000; // const MAX_DELETE_LIST: i32 = 1000; @@ -318,6 +318,8 @@ impl ECStore { opts: ListPathOptions, sender: Sender, ) -> Result> { + // warn!("list_merged ops {:?}", &opts); + let mut futures = Vec::new(); let mut inputs = Vec::new(); @@ -396,7 +398,8 @@ async fn gather_results( ) -> Result> { let mut returned = false; let mut results = Vec::new(); - while let Some(entry) = recv.recv().await { + while let Some(mut entry) = recv.recv().await { + // warn!("gather_results entry {}", &entry.name); if returned { continue; } @@ -404,7 +407,9 @@ async fn gather_results( // TODO: rx.recv() // TODO: isLatestDeletemarker - if !opts.include_directories && (entry.is_dir() || (!opts.versioned && entry.is_object())) { + if !opts.include_directories + && (entry.is_dir() || (!opts.versioned && entry.is_object() && entry.is_latest_deletemarker())) + { continue; } @@ -421,6 +426,7 @@ async fn gather_results( results.push(entry); } + // warn!("gather_results results {:?}", &results); Ok(results) } @@ -456,6 +462,7 @@ async fn merge_entry_channels( tokio::select! { has_entry = in_channels[0].recv()=>{ if let Some(entry) = has_entry{ + // warn!("merge_entry_channels entry {}", &entry.name); out_channel.send(entry).await?; } else { return Ok(())