diff --git a/include/paimon/fs/file_system.h b/include/paimon/fs/file_system.h index 99081b15..822f6ccc 100644 --- a/include/paimon/fs/file_system.h +++ b/include/paimon/fs/file_system.h @@ -135,57 +135,75 @@ class PAIMON_EXPORT OutputStream : public Stream { virtual Result GetUri() const = 0; }; -/// Basic file status information interface. +/// Basic file status information. /// /// This class provides fundamental file system metadata for files and directories. It serves as a -/// lightweight interface for basic file operations that only require path information and directory -/// status. +/// lightweight value type for basic file operations that only require path information and +/// directory status. It can be copied and stored cheaply. class PAIMON_EXPORT BasicFileStatus { public: BasicFileStatus() = default; - virtual ~BasicFileStatus() = default; + + /// Create a basic file status from caller-supplied metadata. + /// @param path The path of the file or directory. + /// @param is_dir Whether the path represents a directory. + BasicFileStatus(std::string path, bool is_dir) : path_(std::move(path)), is_dir_(is_dir) {} /// Check if this entry represents a directory. - virtual bool IsDir() const = 0; + bool IsDir() const { + return is_dir_; + } /// Get the path of this file or directory. - virtual std::string GetPath() const = 0; + std::string GetPath() const { + return path_; + } + + private: + std::string path_; + bool is_dir_ = false; }; -/// Extended file status information. +/// File status information. /// -/// This class extends BasicFileStatus to provide comprehensive file system metadata including file -/// size, modification time, and other attributes. It's used for operations that require detailed -/// file information. +/// This class provides comprehensive file system metadata including file path, size, directory +/// flag, and modification time. It is a concrete value type that can be copied and stored cheaply. class PAIMON_EXPORT FileStatus { public: - FileStatus() = default; - virtual ~FileStatus() = default; - /// Sentinel returned by `GetModificationTime()` when the modification time is not known. static constexpr int64_t kUnknownModificationTime = -1; + static constexpr int64_t kNoSize = -1; + + FileStatus() = default; + /// Create a file status from caller-supplied metadata. /// @param path The path of the file or directory. /// @param length The size of the file in bytes. It may be negative only when the size is /// unknown. /// @param is_dir Whether the path represents a directory. Defaults to false. - FileStatus(std::string path, int64_t length, bool is_dir = false) - : path_(std::move(path)), length_(length), is_dir_(is_dir) {} + /// @param modification_time The last modification time in milliseconds since the epoch + /// (UTC January 1, 1970). Defaults to `kUnknownModificationTime`. + FileStatus(std::string path, int64_t length, bool is_dir = false, + int64_t modification_time = kUnknownModificationTime) + : path_(std::move(path)), + length_(length), + is_dir_(is_dir), + modification_time_(modification_time) {} /// Get the size of the file in bytes. /// @note For directories, this method is undefined behavior. - virtual int64_t GetLen() const { + int64_t GetLen() const { return length_; } /// Check if this entry represents a directory. - virtual bool IsDir() const { + bool IsDir() const { return is_dir_; } /// Get the path of this file or directory. - virtual std::string GetPath() const { + std::string GetPath() const { return path_; } @@ -193,14 +211,15 @@ class PAIMON_EXPORT FileStatus { /// /// @return A long value representing the time the file was last modified, measured in /// milliseconds since the epoch (UTC January 1, 1970). - virtual int64_t GetModificationTime() const { - return kUnknownModificationTime; + int64_t GetModificationTime() const { + return modification_time_; } private: std::string path_; - int64_t length_ = -1; + int64_t length_ = kNoSize; bool is_dir_ = false; + int64_t modification_time_ = kUnknownModificationTime; }; /// Abstract file system interface. @@ -263,25 +282,23 @@ class PAIMON_EXPORT FileSystem { virtual Status Delete(const std::string& path, bool recursive = true) const = 0; /// Get detailed status information for a file or directory. /// @param path The file or directory path to query. - /// @return Result containing a unique pointer to `FileStatus` on success, or error status on + /// @return Result containing a `FileStatus` on success, or error status on /// failure (e.g., path not found, permission denied). - virtual Result> GetFileStatus(const std::string& path) const = 0; + virtual Result GetFileStatus(const std::string& path) const = 0; /// List files of a directory (basic information only). /// @param directory The directory path to list. /// @param[out] file_status_list Output vector to store `BasicFileStatus` objects. /// @return Status indicating success (OK) or failure with error information. - virtual Status ListDir( - const std::string& directory, - std::vector>* file_status_list) const = 0; + virtual Status ListDir(const std::string& directory, + std::vector* file_status_list) const = 0; /// List file status with detailed information. /// @param path The file or directory path to list. /// @param[out] file_status_list Output vector to store `FileStatus` objects. /// @return Status indicating success (OK) or failure with error information. - virtual Status ListFileStatus( - const std::string& path, - std::vector>* file_status_list) const = 0; + virtual Status ListFileStatus(const std::string& path, + std::vector* file_status_list) const = 0; /// Check if a file or directory exists. /// @param path The file or directory path to check. diff --git a/src/paimon/common/data/decimal_test.cpp b/src/paimon/common/data/decimal_test.cpp index f8931b10..b5674261 100644 --- a/src/paimon/common/data/decimal_test.cpp +++ b/src/paimon/common/data/decimal_test.cpp @@ -100,7 +100,7 @@ TEST(DecimalTest, TestCompatibleWithJava) { auto pool = GetDefaultPool(); auto file_system = std::make_unique(); auto file_name = paimon::test::GetDataDir() + "/decimal_bytes.data"; - int64_t file_length = file_system->GetFileStatus(file_name).value()->GetLen(); + int64_t file_length = file_system->GetFileStatus(file_name).value().GetLen(); ASSERT_GT(file_length, 0); ASSERT_OK_AND_ASSIGN(auto input_stream, file_system->Open(file_name)); auto data_bytes = Bytes::AllocateBytes(file_length, pool.get()); diff --git a/src/paimon/common/fs/file_system_test.cpp b/src/paimon/common/fs/file_system_test.cpp index 22cbf3f0..b4527dd2 100644 --- a/src/paimon/common/fs/file_system_test.cpp +++ b/src/paimon/common/fs/file_system_test.cpp @@ -113,19 +113,19 @@ class FileSystemTest : public ::testing::Test, public ::testing::WithParamInterf return new_path; } - void CheckFileStatus(const std::vector>& actual_statuses, + void CheckFileStatus(const std::vector& actual_statuses, const std::set& expected_files, const std::set& expected_dirs) const { ASSERT_EQ(actual_statuses.size(), expected_files.size() + expected_dirs.size()); std::set actual_files; std::set actual_dirs; for (const auto& file_status : actual_statuses) { - if (file_status->IsDir()) { - actual_dirs.insert(RemoveLastSlashInPath(file_status->GetPath())); + if (file_status.IsDir()) { + actual_dirs.insert(RemoveLastSlashInPath(file_status.GetPath())); } else { - actual_files.insert(file_status->GetPath()); - ASSERT_GT(file_status->GetLen(), 0); - int64_t modification_time = file_status->GetModificationTime(); + actual_files.insert(file_status.GetPath()); + ASSERT_GT(file_status.GetLen(), 0); + int64_t modification_time = file_status.GetModificationTime(); ASSERT_GT(modification_time, 10000000000L); // MIN_VALID_FILE_MODIFICATION_MS ASSERT_LT(modification_time, 10000000000000L); // MAX_VALID_FILE_MODIFICATION_MS } @@ -138,17 +138,17 @@ class FileSystemTest : public ::testing::Test, public ::testing::WithParamInterf ASSERT_EQ(actual_dirs, normalized_expected_dirs); } - void CheckBasicFileStatus(const std::vector>& actual_statuses, + void CheckBasicFileStatus(const std::vector& actual_statuses, const std::set& expected_files, const std::set& expected_dirs) const { ASSERT_EQ(actual_statuses.size(), expected_files.size() + expected_dirs.size()); std::set actual_files; std::set actual_dirs; for (const auto& file_status : actual_statuses) { - if (file_status->IsDir()) { - actual_dirs.insert(RemoveLastSlashInPath(file_status->GetPath())); + if (file_status.IsDir()) { + actual_dirs.insert(RemoveLastSlashInPath(file_status.GetPath())); } else { - actual_files.insert(file_status->GetPath()); + actual_files.insert(file_status.GetPath()); } } std::set normalized_expected_dirs; @@ -354,11 +354,11 @@ TEST_P(FileSystemTest, TestWriteEmptyFile) { ASSERT_OK(out_stream->Close()); // get file status - ASSERT_OK_AND_ASSIGN(auto st, fs_->GetFileStatus(file_path)); - ASSERT_EQ(st->GetPath(), file_path); - ASSERT_FALSE(st->IsDir()); - ASSERT_EQ(st->GetLen(), 0); - auto modification_time = st->GetModificationTime(); + ASSERT_OK_AND_ASSIGN(FileStatus st, fs_->GetFileStatus(file_path)); + ASSERT_EQ(st.GetPath(), file_path); + ASSERT_FALSE(st.IsDir()); + ASSERT_EQ(st.GetLen(), 0); + int64_t modification_time = st.GetModificationTime(); ASSERT_GT(modification_time, 10000000000L); ASSERT_LT(modification_time, 10000000000000L); @@ -1202,10 +1202,10 @@ TEST_P(FileSystemTest, TestGetFileStatus1) { ASSERT_OK(fs_->Mkdirs(dir_path)); ASSERT_OK_AND_ASSIGN(bool is_exist, fs_->Exists(dir_path)); ASSERT_TRUE(is_exist); - ASSERT_OK_AND_ASSIGN(std::unique_ptr st, fs_->GetFileStatus(dir_path)); - ASSERT_EQ(RemoveLastSlashInPath(st->GetPath()), RemoveLastSlashInPath(dir_path)); - ASSERT_TRUE(st->IsDir()); - auto modification_time = st->GetModificationTime(); + ASSERT_OK_AND_ASSIGN(FileStatus st, fs_->GetFileStatus(dir_path)); + ASSERT_EQ(RemoveLastSlashInPath(st.GetPath()), RemoveLastSlashInPath(dir_path)); + ASSERT_TRUE(st.IsDir()); + auto modification_time = st.GetModificationTime(); ASSERT_GT(modification_time, 10000000000L); ASSERT_LT(modification_time, 10000000000000L); @@ -1216,9 +1216,9 @@ TEST_P(FileSystemTest, TestGetFileStatus1) { // check meta in dir ASSERT_OK_AND_ASSIGN(st, fs_->GetFileStatus(dir_path)); - ASSERT_EQ(RemoveLastSlashInPath(st->GetPath()), RemoveLastSlashInPath(dir_path)); - ASSERT_TRUE(st->IsDir()); - modification_time = st->GetModificationTime(); + ASSERT_EQ(RemoveLastSlashInPath(st.GetPath()), RemoveLastSlashInPath(dir_path)); + ASSERT_TRUE(st.IsDir()); + modification_time = st.GetModificationTime(); ASSERT_GT(modification_time, 10000000000L); ASSERT_LT(modification_time, 10000000000000L); } @@ -1237,11 +1237,11 @@ TEST_P(FileSystemTest, TestGetFileStatus1) { ASSERT_OK_AND_ASSIGN(bool is_exist, fs_->Exists(file_path)); ASSERT_TRUE(is_exist); - ASSERT_OK_AND_ASSIGN(std::unique_ptr st, fs_->GetFileStatus(file_path)); - ASSERT_EQ(st->GetPath(), file_path); - ASSERT_FALSE(st->IsDir()); - ASSERT_EQ(st->GetLen(), content.size()); - auto modification_time = st->GetModificationTime(); + ASSERT_OK_AND_ASSIGN(FileStatus st, fs_->GetFileStatus(file_path)); + ASSERT_EQ(st.GetPath(), file_path); + ASSERT_FALSE(st.IsDir()); + ASSERT_EQ(st.GetLen(), content.size()); + auto modification_time = st.GetModificationTime(); ASSERT_GT(modification_time, 10000000000L); ASSERT_LT(modification_time, 10000000000000L); } @@ -1259,31 +1259,29 @@ TEST_P(FileSystemTest, TestGetFileStatus2) { { // input is a dir, with a trailing '/' std::string dir_name = test_path + "/"; - std::vector> status_list; - ASSERT_OK_AND_ASSIGN(std::unique_ptr file_status, fs_->GetFileStatus(dir_name)); + std::vector status_list; + ASSERT_OK_AND_ASSIGN(FileStatus file_status, fs_->GetFileStatus(dir_name)); status_list.emplace_back(std::move(file_status)); CheckFileStatus(status_list, /*expected_files=*/{}, /*expected_dirs=*/{dir_name}); } { // input is a dir, without a trailing '/' - std::vector> status_list; - ASSERT_OK_AND_ASSIGN(std::unique_ptr file_status, - fs_->GetFileStatus(test_path)); + std::vector status_list; + ASSERT_OK_AND_ASSIGN(FileStatus file_status, fs_->GetFileStatus(test_path)); status_list.emplace_back(std::move(file_status)); CheckFileStatus(status_list, /*expected_files=*/{}, /*expected_dirs=*/{test_path}); } { // input is a file - std::vector> status_list; + std::vector status_list; std::string file_name = PathUtil::JoinPath(test_path, "README"); - ASSERT_OK_AND_ASSIGN(std::unique_ptr file_status, - fs_->GetFileStatus(file_name)); + ASSERT_OK_AND_ASSIGN(FileStatus file_status, fs_->GetFileStatus(file_name)); status_list.emplace_back(std::move(file_status)); CheckFileStatus(status_list, /*expected_files=*/{file_name}, /*expected_dirs=*/{}); } { // input is not exist - std::vector> status_list; + std::vector status_list; ASSERT_NOK(fs_->GetFileStatus(PathUtil::JoinPath(test_path, "NOT_EXIST"))); } } @@ -1291,7 +1289,7 @@ TEST_P(FileSystemTest, TestGetFileStatus2) { TEST_P(FileSystemTest, TestInvalidListFileStatus) { { // list non exist dir will return ok - std::vector> file_status_list; + std::vector file_status_list; ASSERT_OK(fs_->ListFileStatus(test_root_ + "/non-exist/", &file_status_list)); ASSERT_TRUE(file_status_list.empty()); } @@ -1301,11 +1299,11 @@ TEST_P(FileSystemTest, TestInvalidListFileStatus) { ASSERT_OK(fs_->WriteFile(file_path, "hello", /*overwrite=*/false)); ASSERT_OK_AND_ASSIGN(bool is_exist, fs_->Exists(file_path)); ASSERT_TRUE(is_exist); - std::vector> file_status_list; + std::vector file_status_list; ASSERT_OK(fs_->ListFileStatus(file_path, &file_status_list)); - ASSERT_EQ(file_status_list[0]->GetPath(), file_path); - ASSERT_EQ(file_status_list[0]->GetLen(), 5); - ASSERT_FALSE(file_status_list[0]->IsDir()); + ASSERT_EQ(file_status_list[0].GetPath(), file_path); + ASSERT_EQ(file_status_list[0].GetLen(), 5); + ASSERT_FALSE(file_status_list[0].IsDir()); } } @@ -1315,7 +1313,7 @@ TEST_P(FileSystemTest, TestListFileStatus1) { ASSERT_OK_AND_ASSIGN(bool is_exist, fs_->Exists(file_path)); ASSERT_TRUE(is_exist); - std::vector> file_status_list; + std::vector file_status_list; ASSERT_OK(fs_->ListFileStatus(test_root_, &file_status_list)); ASSERT_EQ(file_status_list.size(), 1); std::set expected_dirs = {test_root_ + "/file_dir1/"}; @@ -1336,13 +1334,13 @@ TEST_P(FileSystemTest, TestListFileStatus1) { ASSERT_OK_AND_ASSIGN(is_exist, fs_->Exists(file_path2)); ASSERT_TRUE(is_exist); - std::vector> file_status_list2; + std::vector file_status_list2; ASSERT_OK(fs_->ListFileStatus(test_root_, &file_status_list2)); ASSERT_EQ(file_status_list2.size(), 2); expected_dirs = {test_root_ + "/file_dir1/", dir_path2}; CheckFileStatus(file_status_list2, /*expected_files=*/{}, expected_dirs); - std::vector> file_status_list3; + std::vector file_status_list3; ASSERT_OK(fs_->ListFileStatus(test_root_ + "/file_dir1/", &file_status_list3)); ASSERT_EQ(file_status_list3.size(), 3); expected_dirs = {dir_path3}; @@ -1360,25 +1358,25 @@ TEST_P(FileSystemTest, TestListFileStatus2) { PathUtil::JoinPath(test_path, "snapshot")}; { // input is a dir, with a trailing '/' - std::vector> status_list; + std::vector status_list; ASSERT_OK(fs_->ListFileStatus(test_path + "/", &status_list)); CheckFileStatus(status_list, expected_files, expected_dirs); } { // input is a dir, without a trailing '/' - std::vector> status_list; + std::vector status_list; ASSERT_OK(fs_->ListFileStatus(test_path, &status_list)); CheckFileStatus(status_list, expected_files, expected_dirs); } { // input is a file - std::vector> status_list; + std::vector status_list; ASSERT_OK(fs_->ListFileStatus(PathUtil::JoinPath(test_path, "README"), &status_list)); CheckFileStatus(status_list, expected_files, /*expected_dirs=*/{}); } { // input is not exist - std::vector> status_list; + std::vector status_list; ASSERT_OK(fs_->ListFileStatus(PathUtil::JoinPath(test_path, "NOT_EXIST"), &status_list)); CheckFileStatus(status_list, /*expected_files=*/{}, /*expected_dirs=*/{}); } @@ -1391,7 +1389,7 @@ TEST_P(FileSystemTest, TestListDir1) { ASSERT_TRUE(is_exist); auto dir_path1 = test_root_ + "/file_dir1/"; - std::vector> file_status_list; + std::vector file_status_list; ASSERT_OK(fs_->ListDir(test_root_, &file_status_list)); ASSERT_EQ(file_status_list.size(), 1); std::set expected_dirs = {dir_path1}; @@ -1407,13 +1405,13 @@ TEST_P(FileSystemTest, TestListDir1) { ASSERT_OK_AND_ASSIGN(is_exist, fs_->Exists(dir_path3)); ASSERT_TRUE(is_exist); - std::vector> file_status_list2; + std::vector file_status_list2; ASSERT_OK(fs_->ListDir(test_root_, &file_status_list2)); ASSERT_EQ(file_status_list2.size(), 2); expected_dirs = {dir_path1, dir_path2}; CheckBasicFileStatus(file_status_list2, std::set(), expected_dirs); - std::vector> file_status_list3; + std::vector file_status_list3; ASSERT_OK(fs_->ListDir(dir_path1, &file_status_list3)); ASSERT_EQ(file_status_list3.size(), 2); std::set expected_files = {file_path}; @@ -1421,12 +1419,12 @@ TEST_P(FileSystemTest, TestListDir1) { CheckBasicFileStatus(file_status_list3, expected_files, expected_dirs); // list non exist dir will return ok - std::vector> file_status_list4; + std::vector file_status_list4; ASSERT_OK(fs_->ListDir(test_root_ + "/non-exist/", &file_status_list4)); ASSERT_TRUE(file_status_list4.empty()); // list invalid path, a data file path - std::vector> file_status_list5; + std::vector file_status_list5; ASSERT_NOK_WITH_MSG(fs_->ListDir(file_path, &file_status_list5), "is not a directory"); } @@ -1440,26 +1438,26 @@ TEST_P(FileSystemTest, TestListDir2) { PathUtil::JoinPath(test_path, "snapshot")}; { // input is a dir, with a trailing '/' - std::vector> status_list; + std::vector status_list; ASSERT_OK(fs_->ListDir(test_path + "/", &status_list)); CheckBasicFileStatus(status_list, expected_files, expected_dirs); } { // input is a dir, without a trailing '/' - std::vector> status_list; + std::vector status_list; ASSERT_OK(fs_->ListDir(test_path, &status_list)); CheckBasicFileStatus(status_list, expected_files, expected_dirs); } { // input is a file - std::vector> status_list; + std::vector status_list; ASSERT_NOK_WITH_MSG(fs_->ListDir(PathUtil::JoinPath(test_path, "README"), &status_list), "file " + PathUtil::JoinPath(test_path, "README") + " already exists and is not a directory"); } { // input is not exist - std::vector> status_list; + std::vector status_list; ASSERT_OK(fs_->ListDir(PathUtil::JoinPath(test_path, "NOT_EXIST"), &status_list)); CheckBasicFileStatus(status_list, /*expected_files=*/{}, /*expected_dirs=*/{}); diff --git a/src/paimon/common/fs/object_store_file_system.cpp b/src/paimon/common/fs/object_store_file_system.cpp index 77193d9b..ae44835e 100644 --- a/src/paimon/common/fs/object_store_file_system.cpp +++ b/src/paimon/common/fs/object_store_file_system.cpp @@ -43,49 +43,6 @@ std::string NormalizeDirectoryPrefix(const std::string& key) { return key + "/"; } -class ObjectStoreBasicFileStatus : public BasicFileStatus { - public: - ObjectStoreBasicFileStatus(std::string path, bool is_dir) - : path_(std::move(path)), is_dir_(is_dir) {} - bool IsDir() const override { - return is_dir_; - } - std::string GetPath() const override { - return path_; - } - - private: - std::string path_; - bool is_dir_; -}; - -class ObjectStoreFileStatus : public FileStatus { - public: - ObjectStoreFileStatus(std::string path, int64_t size, int64_t modification_time, bool is_dir) - : path_(std::move(path)), - size_(size), - modification_time_(modification_time), - is_dir_(is_dir) {} - int64_t GetLen() const override { - return size_; - } - bool IsDir() const override { - return is_dir_; - } - std::string GetPath() const override { - return path_; - } - int64_t GetModificationTime() const override { - return modification_time_; - } - - private: - std::string path_; - int64_t size_; - int64_t modification_time_; - bool is_dir_; -}; - class ObjectStoreInputStream : public InputStream { public: ObjectStoreInputStream(std::shared_ptr client, @@ -401,15 +358,13 @@ Result> ObjectStoreFileSystem::Open( ToUri(object_path), file_size); } -Result> ObjectStoreFileSystem::GetFileStatus( - const std::string& path) const { +Result ObjectStoreFileSystem::GetFileStatus(const std::string& path) const { PAIMON_ASSIGN_OR_RAISE(ObjectStorePath object_path, ParsePath(path)); if (!object_path.key.empty()) { Result metadata = client_->HeadObject(object_path); if (metadata.ok()) { - return std::make_unique( - ToUri(object_path), metadata.value().size, metadata.value().modification_time, - false); + return FileStatus(ToUri(object_path), metadata.value().size, /*is_dir=*/false, + metadata.value().modification_time); } if (!metadata.status().IsNotExist()) { return metadata.status(); @@ -419,12 +374,12 @@ Result> ObjectStoreFileSystem::GetFileStatus( if (!exists) { return Status::NotExist(fmt::format("{} does not exist", path)); } - return std::make_unique(ToUri(object_path, true), 0, 0, true); + return FileStatus(ToUri(object_path, true), 0, /*is_dir=*/true); } -Status ObjectStoreFileSystem::ListDirectory( - const ObjectStorePath& path, std::vector>* basic_statuses, - std::vector>* statuses) const { +Status ObjectStoreFileSystem::ListDirectory(const ObjectStorePath& path, + std::vector* basic_statuses, + std::vector* statuses) const { if (basic_statuses == nullptr && statuses == nullptr) { return Status::Invalid("a destination status list is required"); } @@ -443,21 +398,18 @@ Status ObjectStoreFileSystem::ListDirectory( } ObjectStorePath child{path.bucket, object.key}; if (basic_statuses) { - basic_statuses->push_back( - std::make_unique(ToUri(child), false)); + basic_statuses->emplace_back(ToUri(child), /*is_dir=*/false); } else { - statuses->push_back(std::make_unique( - ToUri(child), object.size, object.modification_time, false)); + statuses->emplace_back(ToUri(child), object.size, /*is_dir=*/false, + object.modification_time); } } for (const auto& prefix : result.common_prefixes) { ObjectStorePath child{path.bucket, prefix}; if (basic_statuses) { - basic_statuses->push_back( - std::make_unique(ToUri(child, true), true)); + basic_statuses->emplace_back(ToUri(child, true), /*is_dir=*/true); } else { - statuses->push_back( - std::make_unique(ToUri(child, true), 0, 0, true)); + statuses->emplace_back(ToUri(child, true), 0, /*is_dir=*/true); } } token = result.continuation_token; @@ -468,9 +420,8 @@ Status ObjectStoreFileSystem::ListDirectory( return Status::OK(); } -Status ObjectStoreFileSystem::ListDir( - const std::string& directory, - std::vector>* file_status_list) const { +Status ObjectStoreFileSystem::ListDir(const std::string& directory, + std::vector* file_status_list) const { PAIMON_ASSIGN_OR_RAISE(ObjectStorePath path, ParsePath(directory)); if (!path.key.empty() && path.key.back() != '/') { Result metadata = client_->HeadObject(path); @@ -484,15 +435,14 @@ Status ObjectStoreFileSystem::ListDir( return ListDirectory(path, file_status_list, nullptr); } -Status ObjectStoreFileSystem::ListFileStatus( - const std::string& path, std::vector>* file_status_list) const { +Status ObjectStoreFileSystem::ListFileStatus(const std::string& path, + std::vector* file_status_list) const { PAIMON_ASSIGN_OR_RAISE(ObjectStorePath object_path, ParsePath(path)); if (!object_path.key.empty() && object_path.key.back() != '/') { Result metadata = client_->HeadObject(object_path); if (metadata.ok()) { - file_status_list->push_back( - std::make_unique(ToUri(object_path), metadata.value().size, - metadata.value().modification_time, false)); + file_status_list->emplace_back(ToUri(object_path), metadata.value().size, + /*is_dir=*/false, metadata.value().modification_time); return Status::OK(); } if (!metadata.status().IsNotExist()) { diff --git a/src/paimon/common/fs/object_store_file_system.h b/src/paimon/common/fs/object_store_file_system.h index 5bc54a43..a1313b85 100644 --- a/src/paimon/common/fs/object_store_file_system.h +++ b/src/paimon/common/fs/object_store_file_system.h @@ -86,12 +86,11 @@ class PAIMON_EXPORT ObjectStoreFileSystem : public FileSystem { Result> Open(const std::string& path) const override; Result> Open(const FileStatus& file_status) const override; - Result> GetFileStatus(const std::string& path) const override; + Result GetFileStatus(const std::string& path) const override; Status ListDir(const std::string& directory, - std::vector>* file_status_list) const override; - Status ListFileStatus( - const std::string& path, - std::vector>* file_status_list) const override; + std::vector* file_status_list) const override; + Status ListFileStatus(const std::string& path, + std::vector* file_status_list) const override; Result Exists(const std::string& path) const override; Result> Create(const std::string& path, @@ -108,9 +107,8 @@ class PAIMON_EXPORT ObjectStoreFileSystem : public FileSystem { private: Result DirectoryExists(const ObjectStorePath& path) const; - Status ListDirectory(const ObjectStorePath& path, - std::vector>* basic_statuses, - std::vector>* statuses) const; + Status ListDirectory(const ObjectStorePath& path, std::vector* basic_statuses, + std::vector* statuses) const; Status ReadOnlyStatus() const; std::string scheme_; diff --git a/src/paimon/common/fs/object_store_file_system_test.cpp b/src/paimon/common/fs/object_store_file_system_test.cpp index 131bfd17..824060f8 100644 --- a/src/paimon/common/fs/object_store_file_system_test.cpp +++ b/src/paimon/common/fs/object_store_file_system_test.cpp @@ -103,11 +103,11 @@ TEST(ObjectStoreFileSystemTest, TestObjectWinsOverPrefix) { client->objects_["foo/bar"] = "child"; ObjectStoreFileSystem fs("s3", client); ASSERT_OK_AND_ASSIGN(auto status, fs.GetFileStatus("s3://bucket/foo")); - ASSERT_FALSE(status->IsDir()); - std::vector> statuses; + ASSERT_FALSE(status.IsDir()); + std::vector statuses; ASSERT_OK(fs.ListFileStatus("s3://bucket/foo", &statuses)); ASSERT_EQ(statuses.size(), 1); - ASSERT_FALSE(statuses[0]->IsDir()); + ASSERT_FALSE(statuses[0].IsDir()); } TEST(ObjectStoreFileSystemTest, TestHeadErrorIsNotMasked) { @@ -132,11 +132,11 @@ TEST(ObjectStoreFileSystemTest, TestPaginationAndDirectoryMarker) { client->pages_[""] = first; client->pages_["next"] = second; ObjectStoreFileSystem fs("s3", client); - std::vector> statuses; + std::vector statuses; ASSERT_OK(fs.ListFileStatus("s3://bucket/dir/", &statuses)); ASSERT_EQ(statuses.size(), 2); - ASSERT_EQ(statuses[0]->GetPath(), "s3://bucket/dir/a"); - ASSERT_TRUE(statuses[1]->IsDir()); + ASSERT_EQ(statuses[0].GetPath(), "s3://bucket/dir/a"); + ASSERT_TRUE(statuses[1].IsDir()); } TEST(ObjectStoreFileSystemTest, TestTruncatedPageRequiresContinuationToken) { @@ -146,7 +146,7 @@ TEST(ObjectStoreFileSystemTest, TestTruncatedPageRequiresContinuationToken) { page.is_truncated = true; client->pages_[""] = page; ObjectStoreFileSystem fs("s3", client); - std::vector> statuses; + std::vector statuses; ASSERT_TRUE(fs.ListFileStatus("s3://bucket/dir/", &statuses).IsIOError()); ASSERT_TRUE(statuses.empty()); } @@ -172,8 +172,8 @@ TEST(ObjectStoreFileSystemTest, TestPathWithLeadingSlashes) { auto client = std::make_shared(); client->objects_["file"] = "data"; ObjectStoreFileSystem fs("s3", client); - ASSERT_OK_AND_ASSIGN(auto status, fs.GetFileStatus("s3://bucket///file")); - ASSERT_EQ(status->GetPath(), "s3://bucket/file"); + ASSERT_OK_AND_ASSIGN(FileStatus status, fs.GetFileStatus("s3://bucket///file")); + ASSERT_EQ(status.GetPath(), "s3://bucket/file"); } TEST(ObjectStoreInputStreamTest, TestBoundsCloseAndSeekOverflow) { diff --git a/src/paimon/common/fs/resolving_file_system.cpp b/src/paimon/common/fs/resolving_file_system.cpp index a9d6aec5..ce61b9b0 100644 --- a/src/paimon/common/fs/resolving_file_system.cpp +++ b/src/paimon/common/fs/resolving_file_system.cpp @@ -107,21 +107,19 @@ Status ResolvingFileSystem::Delete(const std::string& path, bool recursive) cons return fs->Delete(path, recursive); } -Result> ResolvingFileSystem::GetFileStatus( - const std::string& path) const { +Result ResolvingFileSystem::GetFileStatus(const std::string& path) const { PAIMON_ASSIGN_OR_RAISE(std::shared_ptr fs, GetRealFileSystem(path)); return fs->GetFileStatus(path); } -Status ResolvingFileSystem::ListDir( - const std::string& directory, - std::vector>* file_status_list) const { +Status ResolvingFileSystem::ListDir(const std::string& directory, + std::vector* file_status_list) const { PAIMON_ASSIGN_OR_RAISE(std::shared_ptr fs, GetRealFileSystem(directory)); return fs->ListDir(directory, file_status_list); } -Status ResolvingFileSystem::ListFileStatus( - const std::string& path, std::vector>* file_status_list) const { +Status ResolvingFileSystem::ListFileStatus(const std::string& path, + std::vector* file_status_list) const { PAIMON_ASSIGN_OR_RAISE(std::shared_ptr fs, GetRealFileSystem(path)); return fs->ListFileStatus(path, file_status_list); } diff --git a/src/paimon/common/fs/resolving_file_system.h b/src/paimon/common/fs/resolving_file_system.h index c3625caa..8df32512 100644 --- a/src/paimon/common/fs/resolving_file_system.h +++ b/src/paimon/common/fs/resolving_file_system.h @@ -48,12 +48,11 @@ class ResolvingFileSystem : public FileSystem { Status Mkdirs(const std::string& path) const override; Status Rename(const std::string& src, const std::string& dst) const override; Status Delete(const std::string& path, bool recursive = true) const override; - Result> GetFileStatus(const std::string& path) const override; + Result GetFileStatus(const std::string& path) const override; Status ListDir(const std::string& directory, - std::vector>* file_status_list) const override; - Status ListFileStatus( - const std::string& path, - std::vector>* file_status_list) const override; + std::vector* file_status_list) const override; + Status ListFileStatus(const std::string& path, + std::vector* file_status_list) const override; Result Exists(const std::string& path) const override; private: diff --git a/src/paimon/common/fs/resolving_file_system_test.cpp b/src/paimon/common/fs/resolving_file_system_test.cpp index c459c004..ed1bb781 100644 --- a/src/paimon/common/fs/resolving_file_system_test.cpp +++ b/src/paimon/common/fs/resolving_file_system_test.cpp @@ -51,15 +51,12 @@ class GmockFileSystem : public FileSystem { MOCK_METHOD(Status, Rename, (const std::string& src, const std::string& dst), (const, override)); MOCK_METHOD(Status, Delete, (const std::string& path, bool recursive), (const, override)); - MOCK_METHOD(Result>, GetFileStatus, (const std::string& path), - (const, override)); + MOCK_METHOD(Result, GetFileStatus, (const std::string& path), (const, override)); MOCK_METHOD(Status, ListDir, - (const std::string& directory, - std::vector>* file_status_list), + (const std::string& directory, std::vector* file_status_list), (const, override)); MOCK_METHOD(Status, ListFileStatus, - (const std::string& path, - std::vector>* file_status_list), + (const std::string& path, std::vector* file_status_list), (const, override)); MOCK_METHOD(Result, Exists, (const std::string& path), (const, override)); diff --git a/src/paimon/common/global_index/btree/btree_compatibility_test.cpp b/src/paimon/common/global_index/btree/btree_compatibility_test.cpp index 0adde6c0..c76b3ea8 100644 --- a/src/paimon/common/global_index/btree/btree_compatibility_test.cpp +++ b/src/paimon/common/global_index/btree/btree_compatibility_test.cpp @@ -115,8 +115,8 @@ class BTreeCompatibilityTest : public ::testing::Test { const std::shared_ptr& arrow_type) const { auto meta_str = ReadFileAsString(meta_path); std::shared_ptr meta_bytes = Bytes::AllocateBytes(meta_str, pool_.get()); - PAIMON_ASSIGN_OR_RAISE(auto file_status, fs_->GetFileStatus(bin_path)); - auto file_size = file_status->GetLen(); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, fs_->GetFileStatus(bin_path)); + auto file_size = file_status.GetLen(); GlobalIndexIOMeta io_meta(bin_path, file_size, meta_bytes); std::vector metas = {io_meta}; diff --git a/src/paimon/common/global_index/btree/btree_global_index_integration_test.cpp b/src/paimon/common/global_index/btree/btree_global_index_integration_test.cpp index e5aeaac2..29ea456d 100644 --- a/src/paimon/common/global_index/btree/btree_global_index_integration_test.cpp +++ b/src/paimon/common/global_index/btree/btree_global_index_integration_test.cpp @@ -55,8 +55,9 @@ class FakeGlobalIndexFileWriter : public GlobalIndexFileWriter { } Result GetFileSize(const std::string& file_name) const override { - PAIMON_ASSIGN_OR_RAISE(auto file_status, fs_->GetFileStatus(base_path_ + "/" + file_name)); - return file_status->GetLen(); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, + fs_->GetFileStatus(base_path_ + "/" + file_name)); + return file_status.GetLen(); } std::string ToPath(const std::string& file_name) const override { diff --git a/src/paimon/common/global_index/btree/lazy_filtered_btree_reader_test.cpp b/src/paimon/common/global_index/btree/lazy_filtered_btree_reader_test.cpp index 6a45aee3..bf5503fa 100644 --- a/src/paimon/common/global_index/btree/lazy_filtered_btree_reader_test.cpp +++ b/src/paimon/common/global_index/btree/lazy_filtered_btree_reader_test.cpp @@ -54,8 +54,9 @@ class FakeLazyFileWriter : public GlobalIndexFileWriter { } Result GetFileSize(const std::string& file_name) const override { - PAIMON_ASSIGN_OR_RAISE(auto file_status, fs_->GetFileStatus(base_path_ + "/" + file_name)); - return file_status->GetLen(); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, + fs_->GetFileStatus(base_path_ + "/" + file_name)); + return file_status.GetLen(); } std::string ToPath(const std::string& file_name) const override { diff --git a/src/paimon/common/logging/logging_test.cpp b/src/paimon/common/logging/logging_test.cpp index 6ffec1f8..d7fe9832 100644 --- a/src/paimon/common/logging/logging_test.cpp +++ b/src/paimon/common/logging/logging_test.cpp @@ -135,16 +135,16 @@ TEST(LoggerTest, TestGlogWritesLogFileToDisk) { // Collect the content of whatever file glog created in our directory. std::string on_disk_path; std::string content; - std::vector> entries; + std::vector entries; ASSERT_OK(fs->ListDir(tmp_dir->Str(), &entries)); for (const auto& entry : entries) { - if (entry->IsDir()) { + if (entry.IsDir()) { continue; } std::string file_content; - ASSERT_OK(fs->ReadFile(entry->GetPath(), &file_content)); + ASSERT_OK(fs->ReadFile(entry.GetPath(), &file_content)); if (file_content.find(token) != std::string::npos) { - on_disk_path = entry->GetPath(); + on_disk_path = entry.GetPath(); content = std::move(file_content); break; } diff --git a/src/paimon/common/reader/prefetch_file_batch_reader_impl_test.cpp b/src/paimon/common/reader/prefetch_file_batch_reader_impl_test.cpp index ed840a24..b3828f7e 100644 --- a/src/paimon/common/reader/prefetch_file_batch_reader_impl_test.cpp +++ b/src/paimon/common/reader/prefetch_file_batch_reader_impl_test.cpp @@ -206,7 +206,7 @@ class PrefetchFileBatchReaderImplTest : public ::testing::Test, EXPECT_OK_AND_ASSIGN( std::unique_ptr reader, PrefetchFileBatchReaderImpl::Create( - data_file_path, data_file_status->GetLen(), reader_builder.get(), local_fs_, + data_file_path, data_file_status.GetLen(), reader_builder.get(), local_fs_, prefetch_max_parallel_num, batch_size, prefetch_max_parallel_num * 2, /*enable_adaptive_prefetch_strategy=*/false, executor, /*initialize_read_ranges=*/false, cache_mode, CacheConfig(), GetDefaultPool())); diff --git a/src/paimon/core/append/append_compact_coordinator_test.cpp b/src/paimon/core/append/append_compact_coordinator_test.cpp index 2d1ff12c..1c8c7aff 100644 --- a/src/paimon/core/append/append_compact_coordinator_test.cpp +++ b/src/paimon/core/append/append_compact_coordinator_test.cpp @@ -580,7 +580,7 @@ TEST_F(AppendCompactCoordinatorTest, TestCompactWithExternalPath) { { auto filesystem = external_dir->GetFileSystem(); auto bucket_dir = external_path + "/bucket-0/"; - std::vector> file_statuses; + std::vector file_statuses; ASSERT_OK(filesystem->ListDir(bucket_dir, &file_statuses)); ASSERT_EQ(file_statuses.size(), 3) << "External path directory should contain compact output files"; diff --git a/src/paimon/core/append/append_only_writer_test.cpp b/src/paimon/core/append/append_only_writer_test.cpp index fdde9a9c..8ef6e135 100644 --- a/src/paimon/core/append/append_only_writer_test.cpp +++ b/src/paimon/core/append/append_only_writer_test.cpp @@ -444,7 +444,7 @@ TEST_F(AppendOnlyWriterTest, TestWriteAndClose) { ASSERT_OK(writer->Close()); auto file_system = std::make_shared(); - std::vector> file_status_list; + std::vector file_status_list; ASSERT_OK(file_system->ListDir(dir->Str(), &file_status_list)); ASSERT_TRUE(file_status_list.empty()); } @@ -489,7 +489,7 @@ TEST_F(AppendOnlyWriterTest, TestInvalidRowKind) { ASSERT_OK(writer->Close()); auto file_system = std::make_shared(); - std::vector> file_status_list; + std::vector file_status_list; ASSERT_OK(file_system->ListDir(dir->Str(), &file_status_list)); ASSERT_TRUE(file_status_list.empty()); } diff --git a/src/paimon/core/catalog/file_system_catalog.cpp b/src/paimon/core/catalog/file_system_catalog.cpp index 3c587709..85907c12 100644 --- a/src/paimon/core/catalog/file_system_catalog.cpp +++ b/src/paimon/core/catalog/file_system_catalog.cpp @@ -217,12 +217,12 @@ Result FileSystemCatalog::NewDataTablePath(const std::string& wareh } Result> FileSystemCatalog::ListDatabases() const { - std::vector> file_status_list; + std::vector file_status_list; PAIMON_RETURN_NOT_OK(fs_->ListDir(warehouse_, &file_status_list)); std::vector db_names; for (const auto& file_status : file_status_list) { - if (file_status->IsDir()) { - std::string name = PathUtil::GetName(file_status->GetPath()); + if (file_status.IsDir()) { + std::string name = PathUtil::GetName(file_status.GetPath()); if (StringUtils::EndsWith(name, DB_SUFFIX)) { db_names.push_back(name.substr(0, name.length() - std::strlen(DB_SUFFIX))); } @@ -236,12 +236,12 @@ Result> FileSystemCatalog::ListTables(const std::string return GlobalSystemTableLoader::GetSupportedTableNames(catalog_options_); } std::string database_path = NewDatabasePath(warehouse_, db_name); - std::vector> file_status_list; + std::vector file_status_list; PAIMON_RETURN_NOT_OK(fs_->ListDir(database_path, &file_status_list)); std::vector table_names; for (const auto& file_status : file_status_list) { - if (file_status->IsDir()) { - std::string table_path = file_status->GetPath(); + if (file_status.IsDir()) { + std::string table_path = file_status.GetPath(); PAIMON_ASSIGN_OR_RAISE(bool table_exist, TableExistsInFileSystem(table_path)); if (table_exist) { table_names.push_back(PathUtil::GetName(table_path)); diff --git a/src/paimon/core/global_index/global_index_file_manager.h b/src/paimon/core/global_index/global_index_file_manager.h index 584be053..3db9965b 100644 --- a/src/paimon/core/global_index/global_index_file_manager.h +++ b/src/paimon/core/global_index/global_index_file_manager.h @@ -63,9 +63,8 @@ class GlobalIndexFileManager : public GlobalIndexFileReader, public GlobalIndexF } Result GetFileSize(const std::string& file_name) const override { - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file_status, - fs_->GetFileStatus(ToPath(file_name))); - return file_status->GetLen(); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, fs_->GetFileStatus(ToPath(file_name))); + return file_status.GetLen(); } bool IsExternalPath() const { diff --git a/src/paimon/core/index/index_file.h b/src/paimon/core/index/index_file.h index f3829b4e..c3553367 100644 --- a/src/paimon/core/index/index_file.h +++ b/src/paimon/core/index/index_file.h @@ -46,8 +46,8 @@ class IndexFile { } virtual Result FileSize(const std::string& file) const { - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file_status, fs_->GetFileStatus(file)); - return file_status->GetLen(); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, fs_->GetFileStatus(file)); + return file_status.GetLen(); } virtual void Delete(const std::shared_ptr& file) const { diff --git a/src/paimon/core/io/single_file_writer.h b/src/paimon/core/io/single_file_writer.h index 224e8d9e..99507b57 100644 --- a/src/paimon/core/io/single_file_writer.h +++ b/src/paimon/core/io/single_file_writer.h @@ -232,8 +232,8 @@ Status SingleFileWriter::Close() { std::shared_ptr out = std::move(out_); PAIMON_RETURN_NOT_OK(out->Close()); } else { - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file_status, fs_->GetFileStatus(path_)); - output_bytes_ = file_status->GetLen(); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, fs_->GetFileStatus(path_)); + output_bytes_ = file_status.GetLen(); } // Completing the format writer and stream is terminal even if publication fails. The scope // guard still removes the file on a callback error, while a repeated Close() does not publish diff --git a/src/paimon/core/io/single_file_writer_test.cpp b/src/paimon/core/io/single_file_writer_test.cpp index 8ed940e2..78fce54c 100644 --- a/src/paimon/core/io/single_file_writer_test.cpp +++ b/src/paimon/core/io/single_file_writer_test.cpp @@ -74,7 +74,7 @@ TEST(SingleFileWriterTest, TestSimple) { ASSERT_NOK_WITH_MSG(writer.GetAbortExecutor(), "Writer should be closed"); ASSERT_OK(writer.Close()); ASSERT_OK_AND_ASSIGN(auto file_status, file_system->GetFileStatus(file_path)); - ASSERT_GT(file_status->GetLen(), 0); + ASSERT_GT(file_status.GetLen(), 0); ASSERT_OK_AND_ASSIGN(auto abort_executor, writer.GetAbortExecutor()); abort_executor.Abort(); ASSERT_OK_AND_ASSIGN(auto exist, file_system->Exists(file_path)); diff --git a/src/paimon/core/manifest/manifest_committable_test.cpp b/src/paimon/core/manifest/manifest_committable_test.cpp index 00a8a30d..4ed62827 100644 --- a/src/paimon/core/manifest/manifest_committable_test.cpp +++ b/src/paimon/core/manifest/manifest_committable_test.cpp @@ -34,7 +34,7 @@ class ManifestCommittableTest : public testing::Test { std::vector> GetCommitMessages(const std::string& path, int32_t version) const { auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(path).value().GetLen(); std::vector buffer(buffer_length, 0); EXPECT_OK_AND_ASSIGN(auto in_stream, file_system->Open(path)); diff --git a/src/paimon/core/manifest/manifest_file_test.cpp b/src/paimon/core/manifest/manifest_file_test.cpp index 8bb2c22e..8f6e0b2e 100644 --- a/src/paimon/core/manifest/manifest_file_test.cpp +++ b/src/paimon/core/manifest/manifest_file_test.cpp @@ -72,19 +72,18 @@ class CountingFileSystem : public FileSystem { return local_.Delete(path, recursive); } - Result> GetFileStatus(const std::string& path) const override { + Result GetFileStatus(const std::string& path) const override { ++get_file_status_count; return local_.GetFileStatus(path); } Status ListDir(const std::string& directory, - std::vector>* file_status_list) const override { + std::vector* file_status_list) const override { return local_.ListDir(directory, file_status_list); } - Status ListFileStatus( - const std::string& path, - std::vector>* file_status_list) const override { + Status ListFileStatus(const std::string& path, + std::vector* file_status_list) const override { return local_.ListFileStatus(path, file_status_list); } diff --git a/src/paimon/core/mergetree/lookup/remote_lookup_file_manager.cpp b/src/paimon/core/mergetree/lookup/remote_lookup_file_manager.cpp index 6d962572..28a5ed17 100644 --- a/src/paimon/core/mergetree/lookup/remote_lookup_file_manager.cpp +++ b/src/paimon/core/mergetree/lookup/remote_lookup_file_manager.cpp @@ -53,9 +53,9 @@ Result> RemoteLookupFileManager::GenRemoteLookupFi std::string local_file_path = lookup_file->LocalFile(); // Get the file size from the local file system - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr local_file_status, + PAIMON_ASSIGN_OR_RAISE(FileStatus local_file_status, file_system_->GetFileStatus(local_file_path)); - int64_t length = local_file_status->GetLen(); + int64_t length = local_file_status.GetLen(); std::string remote_sst_name = lookup_levels->NewRemoteSst(file, length); std::string remote_sst_path = RemoteSstPath(file, remote_sst_name); diff --git a/src/paimon/core/mergetree/lookup/remote_lookup_file_manager_test.cpp b/src/paimon/core/mergetree/lookup/remote_lookup_file_manager_test.cpp index cd8f3e18..2d146086 100644 --- a/src/paimon/core/mergetree/lookup/remote_lookup_file_manager_test.cpp +++ b/src/paimon/core/mergetree/lookup/remote_lookup_file_manager_test.cpp @@ -291,7 +291,7 @@ TEST_F(RemoteLookupFileManagerTest, TryToDownloadLargeFileAcrossMultipleBuffers) // Verify the downloaded file has the correct size ASSERT_OK_AND_ASSIGN(auto local_status, fs_->GetFileStatus(local_file_path)); - ASSERT_EQ(local_status->GetLen(), file_size); + ASSERT_EQ(local_status.GetLen(), file_size); // Verify the downloaded file content matches the original { diff --git a/src/paimon/core/mergetree/lookup_file_test.cpp b/src/paimon/core/mergetree/lookup_file_test.cpp index a9a9051b..f6f814ef 100644 --- a/src/paimon/core/mergetree/lookup_file_test.cpp +++ b/src/paimon/core/mergetree/lookup_file_test.cpp @@ -216,7 +216,7 @@ TEST(LookupFileTest, TestLookupFileCacheLifecycle) { ASSERT_EQ(cache->GetCurrentWeight(), 0); // d.sst should be deleted ASSERT_FALSE(fs->Exists(path_d).value()); - std::vector> file_status_list; + std::vector file_status_list; ASSERT_OK(fs->ListDir(tmp_dir->Str(), &file_status_list)); ASSERT_TRUE(file_status_list.empty()); ASSERT_EQ(call_back_files, diff --git a/src/paimon/core/mergetree/lookup_levels.cpp b/src/paimon/core/mergetree/lookup_levels.cpp index 51bb5651..8b5f69e7 100644 --- a/src/paimon/core/mergetree/lookup_levels.cpp +++ b/src/paimon/core/mergetree/lookup_levels.cpp @@ -236,8 +236,8 @@ Result> LookupLevels::CreateLookupFile( } // Get file size for cache weight calculation - PAIMON_ASSIGN_OR_RAISE(auto file_status, fs_->GetFileStatus(kv_file_path)); - int64_t file_size = file_status->GetLen(); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, fs_->GetFileStatus(kv_file_path)); + int64_t file_size = file_status.GetLen(); PAIMON_ASSIGN_OR_RAISE(std::unique_ptr reader, lookup_store_factory_->CreateReader(fs_, kv_file_path, pool_)); diff --git a/src/paimon/core/mergetree/lookup_levels_test.cpp b/src/paimon/core/mergetree/lookup_levels_test.cpp index a136b625..09565692 100644 --- a/src/paimon/core/mergetree/lookup_levels_test.cpp +++ b/src/paimon/core/mergetree/lookup_levels_test.cpp @@ -238,10 +238,10 @@ TEST_F(LookupLevelsTest, TestMultiLevels) { ASSERT_EQ(lookup_levels->GetLevels()->NonEmptyHighestLevel(), 2); // test lookup file in tmp dir - std::vector> file_status_list; + std::vector file_status_list; ASSERT_OK(fs_->ListDir(tmp_dir_->Str(), &file_status_list)); ASSERT_EQ(file_status_list.size(), 1); - auto channel_dir = file_status_list[0]->GetPath(); + auto channel_dir = file_status_list[0].GetPath(); file_status_list.clear(); ASSERT_OK(fs_->ListDir(channel_dir, &file_status_list)); ASSERT_EQ(file_status_list.size(), 2); @@ -740,10 +740,10 @@ TEST_F(LookupLevelsTest, TestLookupFileCacheIntegration) { ASSERT_EQ(shared_cache->Size(), 3); // file0, file1, file2 // Collect local file paths for later verification - std::vector> tmp_files; + std::vector tmp_files; ASSERT_OK(fs_->ListDir(tmp_dir_->Str(), &tmp_files)); ASSERT_EQ(tmp_files.size(), 1); - auto channel_dir = tmp_files[0]->GetPath(); + auto channel_dir = tmp_files[0].GetPath(); tmp_files.clear(); ASSERT_OK(fs_->ListDir(channel_dir, &tmp_files)); ASSERT_EQ(tmp_files.size(), 3); diff --git a/src/paimon/core/mergetree/merge_tree_writer_test.cpp b/src/paimon/core/mergetree/merge_tree_writer_test.cpp index 57a62654..2155647a 100644 --- a/src/paimon/core/mergetree/merge_tree_writer_test.cpp +++ b/src/paimon/core/mergetree/merge_tree_writer_test.cpp @@ -256,7 +256,7 @@ TEST_P(MergeTreeWriterTest, TestSimple) { // check data file exist and read ok std::string expected_data_file_name = "data-" + uuid + "-0.orc"; std::string expected_data_file_path = dir->Str() + "/" + expected_data_file_name; - ASSERT_OK_AND_ASSIGN(std::unique_ptr data_file_status, + ASSERT_OK_AND_ASSIGN(FileStatus data_file_status, options.GetFileSystem()->GetFileStatus(expected_data_file_path)); std::shared_ptr expected_array; @@ -273,7 +273,7 @@ TEST_P(MergeTreeWriterTest, TestSimple) { ASSERT_TRUE(commit_increment.GetCompactIncrement().IsEmpty()); ASSERT_EQ(1, commit_increment.GetNewFilesIncrement().NewFiles().size()); auto expected_data_file_meta = std::make_shared( - expected_data_file_name, /*file_size=*/data_file_status->GetLen(), /*row_count=*/3, + expected_data_file_name, /*file_size=*/data_file_status.GetLen(), /*row_count=*/3, /*min_key=*/BinaryRowGenerator::GenerateRow({std::string("Alice")}, pool_.get()), /*max_key=*/BinaryRowGenerator::GenerateRow({std::string("Paul")}, pool_.get()), /*key_stats=*/ @@ -336,7 +336,7 @@ TEST_P(MergeTreeWriterTest, TestWriteMultiBatch) { // check data file exist and read ok std::string expected_data_file_name = "data-" + uuid + "-0.orc"; std::string expected_data_file_path = dir->Str() + "/" + expected_data_file_name; - ASSERT_OK_AND_ASSIGN(std::unique_ptr data_file_status, + ASSERT_OK_AND_ASSIGN(FileStatus data_file_status, options.GetFileSystem()->GetFileStatus(expected_data_file_path)); std::shared_ptr expected_array; @@ -354,7 +354,7 @@ TEST_P(MergeTreeWriterTest, TestWriteMultiBatch) { ASSERT_TRUE(commit_increment.GetCompactIncrement().IsEmpty()); ASSERT_EQ(1, commit_increment.GetNewFilesIncrement().NewFiles().size()); auto expected_data_file_meta = std::make_shared( - expected_data_file_name, /*file_size=*/data_file_status->GetLen(), /*row_count=*/4, + expected_data_file_name, /*file_size=*/data_file_status.GetLen(), /*row_count=*/4, /*min_key=*/BinaryRowGenerator::GenerateRow({std::string("Alice")}, pool_.get()), /*max_key=*/BinaryRowGenerator::GenerateRow({std::string("Skye")}, pool_.get()), /*key_stats=*/ @@ -440,7 +440,7 @@ TEST_P(MergeTreeWriterTest, TestSharedShreddingMapDataFileMetaInfo) { std::string expected_data_file_name = "data-" + uuid + "-0.orc"; std::string expected_data_file_path = dir->Str() + "/" + expected_data_file_name; - ASSERT_OK_AND_ASSIGN(std::unique_ptr data_file_status, + ASSERT_OK_AND_ASSIGN(FileStatus data_file_status, options.GetFileSystem()->GetFileStatus(expected_data_file_path)); std::map column_to_k = {{"tags", 3}}; @@ -466,7 +466,7 @@ TEST_P(MergeTreeWriterTest, TestSharedShreddingMapDataFileMetaInfo) { expected_shredding_meta); auto expected_data_file_meta = std::make_shared( - expected_data_file_name, /*file_size=*/data_file_status->GetLen(), /*row_count=*/2, + expected_data_file_name, /*file_size=*/data_file_status.GetLen(), /*row_count=*/2, /*min_key=*/BinaryRowGenerator::GenerateRow({1}, pool_.get()), /*max_key=*/BinaryRowGenerator::GenerateRow({2}, pool_.get()), /*key_stats=*/ @@ -625,7 +625,7 @@ TEST_P(MergeTreeWriterTest, TestWriteWithDeleteRow) { // check data file exist and read ok std::string expected_data_file_name = "data-" + uuid + "-0.orc"; std::string expected_data_file_path = dir->Str() + "/" + expected_data_file_name; - ASSERT_OK_AND_ASSIGN(std::unique_ptr data_file_status, + ASSERT_OK_AND_ASSIGN(FileStatus data_file_status, options.GetFileSystem()->GetFileStatus(expected_data_file_path)); std::shared_ptr expected_array; @@ -642,7 +642,7 @@ TEST_P(MergeTreeWriterTest, TestWriteWithDeleteRow) { ASSERT_TRUE(commit_increment.GetCompactIncrement().IsEmpty()); ASSERT_EQ(1, commit_increment.GetNewFilesIncrement().NewFiles().size()); auto expected_data_file_meta = std::make_shared( - expected_data_file_name, /*file_size=*/data_file_status->GetLen(), /*row_count=*/3, + expected_data_file_name, /*file_size=*/data_file_status.GetLen(), /*row_count=*/3, /*min_key=*/BinaryRowGenerator::GenerateRow({std::string("Alice")}, pool_.get()), /*max_key=*/BinaryRowGenerator::GenerateRow({std::string("Paul")}, pool_.get()), /*key_stats=*/ @@ -721,10 +721,10 @@ TEST_P(MergeTreeWriterTest, TestMultiplePrepareCommit) { std::string expected_data_file_dir = dir->Str() + "/"; ASSERT_OK_AND_ASSIGN( - std::unique_ptr data_file_status1, + FileStatus data_file_status1, options.GetFileSystem()->GetFileStatus(expected_data_file_dir + expected_data_file_name1)); ASSERT_OK_AND_ASSIGN( - std::unique_ptr data_file_status2, + FileStatus data_file_status2, options.GetFileSystem()->GetFileStatus(expected_data_file_dir + expected_data_file_name2)); std::shared_ptr expected_array1; @@ -753,7 +753,7 @@ TEST_P(MergeTreeWriterTest, TestMultiplePrepareCommit) { ASSERT_EQ(1, commit_increment1.GetNewFilesIncrement().NewFiles().size()); ASSERT_EQ(1, commit_increment2.GetNewFilesIncrement().NewFiles().size()); auto expected_data_file_meta1 = std::make_shared( - expected_data_file_name1, /*file_size=*/data_file_status1->GetLen(), /*row_count=*/3, + expected_data_file_name1, /*file_size=*/data_file_status1.GetLen(), /*row_count=*/3, /*min_key=*/BinaryRowGenerator::GenerateRow({std::string("Alice")}, pool_.get()), /*max_key=*/BinaryRowGenerator::GenerateRow({std::string("Paul")}, pool_.get()), /*key_stats=*/ @@ -772,7 +772,7 @@ TEST_P(MergeTreeWriterTest, TestMultiplePrepareCommit) { /*write_cols=*/std::nullopt); auto expected_data_file_meta2 = std::make_shared( - expected_data_file_name2, /*file_size=*/data_file_status2->GetLen(), /*row_count=*/3, + expected_data_file_name2, /*file_size=*/data_file_status2.GetLen(), /*row_count=*/3, /*min_key=*/BinaryRowGenerator::GenerateRow({std::string("Alice")}, pool_.get()), /*max_key=*/BinaryRowGenerator::GenerateRow({std::string("Skye")}, pool_.get()), /*key_stats=*/ @@ -945,10 +945,10 @@ TEST_P(MergeTreeWriterTest, TestAutoFlush) { std::string expected_data_file_dir = dir->Str() + "/"; ASSERT_OK_AND_ASSIGN( - std::unique_ptr data_file_status1, + FileStatus data_file_status1, options.GetFileSystem()->GetFileStatus(expected_data_file_dir + expected_data_file_name1)); ASSERT_OK_AND_ASSIGN( - std::unique_ptr data_file_status2, + FileStatus data_file_status2, options.GetFileSystem()->GetFileStatus(expected_data_file_dir + expected_data_file_name2)); std::shared_ptr expected_array1; @@ -975,7 +975,7 @@ TEST_P(MergeTreeWriterTest, TestAutoFlush) { ASSERT_TRUE(commit_increment.GetCompactIncrement().IsEmpty()); ASSERT_EQ(2, commit_increment.GetNewFilesIncrement().NewFiles().size()); auto expected_data_file_meta1 = std::make_shared( - expected_data_file_name1, /*file_size=*/data_file_status1->GetLen(), /*row_count=*/3, + expected_data_file_name1, /*file_size=*/data_file_status1.GetLen(), /*row_count=*/3, /*min_key=*/BinaryRowGenerator::GenerateRow({std::string("Alice")}, pool_.get()), /*max_key=*/BinaryRowGenerator::GenerateRow({std::string("Paul")}, pool_.get()), /*key_stats=*/ @@ -994,7 +994,7 @@ TEST_P(MergeTreeWriterTest, TestAutoFlush) { /*write_cols=*/std::nullopt); auto expected_data_file_meta2 = std::make_shared( - expected_data_file_name2, /*file_size=*/data_file_status2->GetLen(), /*row_count=*/3, + expected_data_file_name2, /*file_size=*/data_file_status2.GetLen(), /*row_count=*/3, /*min_key=*/BinaryRowGenerator::GenerateRow({std::string("Alice")}, pool_.get()), /*max_key=*/BinaryRowGenerator::GenerateRow({std::string("Skye")}, pool_.get()), /*key_stats=*/ @@ -1103,12 +1103,12 @@ TEST_P(MergeTreeWriterTest, TestBulkData) { for (size_t i = 0; i < batch_size; ++i) { std::string expected_data_file_name = "data-" + uuid + "-" + std::to_string(i) + ".orc"; // check data file exist and read ok - ASSERT_OK_AND_ASSIGN(std::unique_ptr data_file_status, + ASSERT_OK_AND_ASSIGN(FileStatus data_file_status, options.GetFileSystem()->GetFileStatus(expected_data_file_dir + expected_data_file_name)); // check data file meta auto expected_data_file_meta = std::make_shared( - expected_data_file_name, /*file_size=*/data_file_status->GetLen(), /*row_count=*/3, + expected_data_file_name, /*file_size=*/data_file_status.GetLen(), /*row_count=*/3, /*min_key=*/BinaryRowGenerator::GenerateRow({std::string("Alice")}, pool_.get()), /*max_key=*/BinaryRowGenerator::GenerateRow({std::string("Paul")}, pool_.get()), /*key_stats=*/ diff --git a/src/paimon/core/mergetree/spill_reader.cpp b/src/paimon/core/mergetree/spill_reader.cpp index cb2598e5..c7916b22 100644 --- a/src/paimon/core/mergetree/spill_reader.cpp +++ b/src/paimon/core/mergetree/spill_reader.cpp @@ -54,8 +54,8 @@ Result> SpillReader::Create( Status SpillReader::Open(const FileIOChannel::ID& channel_id) { const std::string& file_path = channel_id.GetPath(); PAIMON_ASSIGN_OR_RAISE(in_stream_, fs_->Open(file_path)); - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file_status, fs_->GetFileStatus(file_path)); - int64_t file_len = file_status->GetLen(); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, fs_->GetFileStatus(file_path)); + int64_t file_len = file_status.GetLen(); arrow_input_stream_adapter_ = std::make_shared(in_stream_, file_len, arrow_pool_); auto ipc_read_options = arrow::ipc::IpcReadOptions::Defaults(); diff --git a/src/paimon/core/mergetree/spill_writer.cpp b/src/paimon/core/mergetree/spill_writer.cpp index 170057fc..1e0e1a33 100644 --- a/src/paimon/core/mergetree/spill_writer.cpp +++ b/src/paimon/core/mergetree/spill_writer.cpp @@ -115,9 +115,8 @@ Result SpillWriter::GetFileSize() const { PAIMON_ASSIGN_OR_RAISE_FROM_ARROW(int64_t file_size, arrow_output_stream_adapter_->Tell()); return file_size; } - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file_status, - fs_->GetFileStatus(channel_id_.GetPath())); - return file_status->GetLen(); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, fs_->GetFileStatus(channel_id_.GetPath())); + return file_status.GetLen(); } const FileIOChannel::ID& SpillWriter::GetChannelId() const { diff --git a/src/paimon/core/mergetree/write_buffer_test.cpp b/src/paimon/core/mergetree/write_buffer_test.cpp index 26f02933..60ca375d 100644 --- a/src/paimon/core/mergetree/write_buffer_test.cpp +++ b/src/paimon/core/mergetree/write_buffer_test.cpp @@ -97,12 +97,12 @@ class WriteBufferTest : public ::testing::Test { Result GetOnlySpillFileSize() const { PAIMON_ASSIGN_OR_RAISE(std::string spill_dir, io_manager_->GetSpillDir()); - std::vector> spill_files; + std::vector spill_files; PAIMON_RETURN_NOT_OK(tmp_dir_->GetFileSystem()->ListFileStatus(spill_dir, &spill_files)); - if (spill_files.size() != 1 || spill_files[0]->IsDir()) { + if (spill_files.size() != 1 || spill_files[0].IsDir()) { return Status::Invalid("expected exactly one spill file"); } - return spill_files[0]->GetLen(); + return spill_files[0].GetLen(); } Result ReadReaderResult(KeyValueRecordReader* reader) const { diff --git a/src/paimon/core/migrate/file_meta_utils.cpp b/src/paimon/core/migrate/file_meta_utils.cpp index cfa30adc..ee1f6599 100644 --- a/src/paimon/core/migrate/file_meta_utils.cpp +++ b/src/paimon/core/migrate/file_meta_utils.cpp @@ -83,11 +83,9 @@ Result> ConstructFileMeta( stats_extractor->ExtractWithFileInfo(fs, dst_file_path, memory_pool)); PAIMON_ASSIGN_OR_RAISE(SimpleStats simple_stats, SimpleStatsConverter::ToBinary(stats.first, memory_pool.get())); - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file_status, - fs->GetFileStatus(dst_file_path)); - assert(file_status); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, fs->GetFileStatus(dst_file_path)); return DataFileMeta::ForAppend( - new_file_name, file_status->GetLen(), stats.second.GetRowCount(), simple_stats, + new_file_name, file_status.GetLen(), stats.second.GetRowCount(), simple_stats, /*min_sequence_number=*/0, /*max_sequence_number=*/0, schema_id, /*extra_files=*/{}, /*embedded_index=*/nullptr, FileSource::Append(), /*value_stats_cols=*/std::nullopt, /*external_path=*/std::nullopt, /*first_row_id=*/std::nullopt, /*write_cols=*/std::nullopt); diff --git a/src/paimon/core/operation/file_store_commit_impl_test.cpp b/src/paimon/core/operation/file_store_commit_impl_test.cpp index 3590c33e..f32d86fa 100644 --- a/src/paimon/core/operation/file_store_commit_impl_test.cpp +++ b/src/paimon/core/operation/file_store_commit_impl_test.cpp @@ -90,8 +90,7 @@ class GmockFileSystem : public LocalFileSystem { public: MOCK_METHOD(Status, ReadFile, (const std::string& path, std::string* content), (override)); MOCK_METHOD(Status, ListDir, - (const std::string& directory, - std::vector>* file_status_list), + (const std::string& directory, std::vector* file_status_list), (const, override)); MOCK_METHOD(Status, AtomicStore, (const std::string& path, const std::string& content), (override)); @@ -111,13 +110,11 @@ class GmockFileSystemFactory : public LocalFileSystemFactory { using ::testing::A; using ::testing::Invoke; - ON_CALL(*fs, ListDir(A(), - A>*>())) - .WillByDefault( - Invoke([fs_ptr](const std::string& directory, - std::vector>* file_status_list) { - return fs_ptr->LocalFileSystem::ListDir(directory, file_status_list); - })); + ON_CALL(*fs, ListDir(A(), A*>())) + .WillByDefault(Invoke([fs_ptr](const std::string& directory, + std::vector* file_status_list) { + return fs_ptr->LocalFileSystem::ListDir(directory, file_status_list); + })); ON_CALL(*fs, ReadFile(A(), A())) .WillByDefault(Invoke([fs_ptr](const std::string& path, std::string* content) { @@ -350,8 +347,8 @@ class FileStoreCommitImplTest : public testing::Test { std::vector> GetCommitMessages(const std::string& path, int32_t version) const { auto file_system = GetFileSystem(); - EXPECT_OK_AND_ASSIGN(std::unique_ptr file, file_system->GetFileStatus(path)); - std::vector buffer(file->GetLen(), 0); + EXPECT_OK_AND_ASSIGN(FileStatus file, file_system->GetFileStatus(path)); + std::vector buffer(file.GetLen(), 0); EXPECT_OK_AND_ASSIGN(std::unique_ptr in_stream, file_system->Open(path)); EXPECT_TRUE(in_stream); @@ -537,9 +534,8 @@ TEST_F(FileStoreCommitImplTest, TestCommitWithConflictSnapshotAndRetryTenTimes) EXPECT_CALL(*mock_fs, ListDir(testing::_, testing::_)).Times(testing::AnyNumber()); EXPECT_CALL(*mock_fs, ListDir(testing::StrEq(PathUtil::JoinPath(table_path, "snapshot")), testing::_)) - .WillRepeatedly( - testing::Invoke([](const std::string& directory, - std::vector>* file_status_list) { + .WillRepeatedly(testing::Invoke( + [](const std::string& directory, std::vector* file_status_list) { return Status::OK(); })); @@ -584,14 +580,12 @@ TEST_F(FileStoreCommitImplTest, TestCommitWithConflictSnapshotAndRetryOnce) { EXPECT_CALL(*mock_fs, ListDir(testing::_, testing::_)).Times(testing::AnyNumber()); EXPECT_CALL(*mock_fs, ListDir(testing::StrEq(PathUtil::JoinPath(table_path, "snapshot")), testing::_)) - .WillOnce( - testing::Invoke([](const std::string& directory, - std::vector>* file_status_list) { + .WillOnce(testing::Invoke( + [](const std::string& directory, std::vector* file_status_list) { return Status::OK(); })) - .WillRepeatedly( - testing::Invoke([&](const std::string& directory, - std::vector>* file_status_list) { + .WillRepeatedly(testing::Invoke( + [&](const std::string& directory, std::vector* file_status_list) { return mock_fs->LocalFileSystem::ListDir(directory, file_status_list); })); diff --git a/src/paimon/core/operation/manifest_file_merger_test.cpp b/src/paimon/core/operation/manifest_file_merger_test.cpp index 1ff1ffd7..bff1659e 100644 --- a/src/paimon/core/operation/manifest_file_merger_test.cpp +++ b/src/paimon/core/operation/manifest_file_merger_test.cpp @@ -113,13 +113,13 @@ class ManifestFileMergerTest : public testing::Test { } std::set ListManifestFiles() const { - std::vector> file_statuses; + std::vector file_statuses; EXPECT_OK(file_system_->ListFileStatus( FileStorePathFactory::ManifestPath(path_factory_->RootPath()), &file_statuses)); std::set files; for (const auto& status : file_statuses) { - if (!status->IsDir()) { - files.insert(status->GetPath()); + if (!status.IsDir()) { + files.insert(status.GetPath()); } } return files; diff --git a/src/paimon/core/operation/orphan_files_cleaner_impl.cpp b/src/paimon/core/operation/orphan_files_cleaner_impl.cpp index e582ae45..6b8e5ba0 100644 --- a/src/paimon/core/operation/orphan_files_cleaner_impl.cpp +++ b/src/paimon/core/operation/orphan_files_cleaner_impl.cpp @@ -98,7 +98,7 @@ Result> OrphanFilesCleanerImpl::Clean() { "OrphanFilesCleaner do not support cleaning table with branch"); } PAIMON_ASSIGN_OR_RAISE(std::set all_dirs, ListPaimonFileDirs()); - std::vector>>> file_statuses_futures; + std::vector>> file_statuses_futures; ScopeGuard file_statuses_guard( [&file_statuses_futures]() { CollectAll(file_statuses_futures); }); for (const auto& dir : all_dirs) { @@ -115,23 +115,23 @@ Result> OrphanFilesCleanerImpl::Clean() { uint64_t file_statuses_duration = duration.Reset(); for (const auto& file_statuses : CollectAll(file_statuses_futures)) { for (const auto& file_status : file_statuses) { - if (file_status->IsDir()) { + if (file_status.IsDir()) { continue; } - std::string path = file_status->GetPath(); + std::string path = file_status.GetPath(); std::string file_name = PathUtil::GetName(path); if (!SupportToClean(file_name)) { continue; } - if (file_status->GetModificationTime() < older_than_ms_ && + if (file_status.GetModificationTime() < older_than_ms_ && !used_file_names.count(file_name)) { if (should_be_retained_ && should_be_retained_(file_name)) { continue; } - if (file_status->GetModificationTime() <= MIN_VALID_FILE_MODIFICATION_MS) { + if (file_status.GetModificationTime() <= MIN_VALID_FILE_MODIFICATION_MS) { return Status::Invalid( fmt::format("file '{}' modification '{}' is not in millisecond", path, - file_status->GetModificationTime())); + file_status.GetModificationTime())); } need_to_deletes.insert(path); futures.push_back(Via(executor_.get(), [this, path]() { @@ -188,7 +188,7 @@ std::set OrphanFilesCleanerImpl::ListFileDirs(const std::string& pa queue.push(path); std::set results; for (int32_t current_level = 0; current_level <= max_level; current_level++) { - std::vector>>> futures; + std::vector>> futures; while (!queue.empty()) { auto current_path = queue.front(); futures.push_back(Via(executor_.get(), [this, current_path] { @@ -198,15 +198,15 @@ std::set OrphanFilesCleanerImpl::ListFileDirs(const std::string& pa } for (const auto& dirs : CollectAll(futures)) { for (const auto& dir : dirs) { - const auto& dir_name = PathUtil::GetName(dir->GetPath()); + const auto& dir_name = PathUtil::GetName(dir.GetPath()); if (current_level == max_level) { if (StringUtils::StartsWith( dir_name, std::string(FileStorePathFactory::BUCKET_PATH_PREFIX))) { - results.insert(dir->GetPath()); + results.insert(dir.GetPath()); } } else { if (dir_name.find("=") != std::string::npos) { - queue.push(dir->GetPath()); + queue.push(dir.GetPath()); } } } @@ -215,13 +215,12 @@ std::set OrphanFilesCleanerImpl::ListFileDirs(const std::string& pa return results; } -std::vector> OrphanFilesCleanerImpl::TryBestListingDirs( - const std::string& path) const { +std::vector OrphanFilesCleanerImpl::TryBestListingDirs(const std::string& path) const { Result is_exist = fs_->Exists(path); if (!is_exist.ok()) { return {}; } - std::vector> file_statuses; + std::vector file_statuses; auto status = fs_->ListFileStatus(path, &file_statuses); if (!status.ok()) { return {}; @@ -229,13 +228,13 @@ std::vector> OrphanFilesCleanerImpl::TryBestListingD return file_statuses; } -std::vector> OrphanFilesCleanerImpl::MinimalTryBestListingDirs( +std::vector OrphanFilesCleanerImpl::MinimalTryBestListingDirs( const std::string& path) const { Result is_exist = fs_->Exists(path); if (!is_exist.ok()) { return {}; } - std::vector> file_statuses; + std::vector file_statuses; auto status = fs_->ListDir(path, &file_statuses); if (!status.ok()) { return {}; diff --git a/src/paimon/core/operation/orphan_files_cleaner_impl.h b/src/paimon/core/operation/orphan_files_cleaner_impl.h index be96b3b5..88cde44e 100644 --- a/src/paimon/core/operation/orphan_files_cleaner_impl.h +++ b/src/paimon/core/operation/orphan_files_cleaner_impl.h @@ -71,9 +71,8 @@ class OrphanFilesCleanerImpl : public OrphanFilesCleaner { private: Result> ListPaimonFileDirs() const; - std::vector> TryBestListingDirs(const std::string& path) const; - std::vector> MinimalTryBestListingDirs( - const std::string& path) const; + std::vector TryBestListingDirs(const std::string& path) const; + std::vector MinimalTryBestListingDirs(const std::string& path) const; std::set ListFileDirs(const std::string& path, int32_t max_level) const; Result> GetUsedFiles() const; Result> GetUsedFilesBySnapshot(const Snapshot& snapshot) const; diff --git a/src/paimon/core/postpone/postpone_bucket_writer_test.cpp b/src/paimon/core/postpone/postpone_bucket_writer_test.cpp index a629a6bd..e932d40d 100644 --- a/src/paimon/core/postpone/postpone_bucket_writer_test.cpp +++ b/src/paimon/core/postpone/postpone_bucket_writer_test.cpp @@ -196,7 +196,7 @@ TEST_P(PostponeBucketWriterTest, TestSimple) { // check data file exist and read ok std::string expected_data_file_name = "data-" + uuid + "-0." + file_format; std::string expected_data_file_path = dir->Str() + "/" + expected_data_file_name; - ASSERT_OK_AND_ASSIGN(std::unique_ptr data_file_status, + ASSERT_OK_AND_ASSIGN(FileStatus data_file_status, options.GetFileSystem()->GetFileStatus(expected_data_file_path)); std::shared_ptr expected_array; @@ -214,7 +214,7 @@ TEST_P(PostponeBucketWriterTest, TestSimple) { ASSERT_TRUE(commit_increment.GetCompactIncrement().IsEmpty()); ASSERT_EQ(1, commit_increment.GetNewFilesIncrement().NewFiles().size()); auto expected_data_file_meta = std::make_shared( - expected_data_file_name, /*file_size=*/data_file_status->GetLen(), /*row_count=*/4, + expected_data_file_name, /*file_size=*/data_file_status.GetLen(), /*row_count=*/4, /*min_key=*/BinaryRowGenerator::GenerateRow({std::string("Lucy")}, pool_.get()), /*max_key=*/BinaryRowGenerator::GenerateRow({std::string("Alice")}, pool_.get()), /*key_stats=*/ @@ -274,7 +274,7 @@ TEST_P(PostponeBucketWriterTest, TestNestedType) { // check data file exist and read ok std::string expected_data_file_name = "data-" + uuid + "-0." + file_format; std::string expected_data_file_path = dir->Str() + "/" + expected_data_file_name; - ASSERT_OK_AND_ASSIGN(std::unique_ptr data_file_status, + ASSERT_OK_AND_ASSIGN(FileStatus data_file_status, options.GetFileSystem()->GetFileStatus(expected_data_file_path)); arrow::FieldVector write_fields = {arrow::field("_SEQUENCE_NUMBER", arrow::int64()), @@ -296,7 +296,7 @@ TEST_P(PostponeBucketWriterTest, TestNestedType) { ASSERT_TRUE(commit_increment.GetCompactIncrement().IsEmpty()); ASSERT_EQ(1, commit_increment.GetNewFilesIncrement().NewFiles().size()); auto expected_data_file_meta = std::make_shared( - expected_data_file_name, /*file_size=*/data_file_status->GetLen(), /*row_count=*/4, + expected_data_file_name, /*file_size=*/data_file_status.GetLen(), /*row_count=*/4, /*min_key=*/BinaryRowGenerator::GenerateRow({std::string("Lucy")}, pool_.get()), /*max_key=*/BinaryRowGenerator::GenerateRow({std::string("Alice")}, pool_.get()), /*key_stats=*/ @@ -353,9 +353,9 @@ TEST_F(PostponeBucketWriterTest, TestSharedShreddingMap) { std::string data_file_name = "data-" + uuid + "-0." + file_format; std::string data_file_path = dir->Str() + "/" + data_file_name; - ASSERT_OK_AND_ASSIGN(std::unique_ptr data_file_status, + ASSERT_OK_AND_ASSIGN(FileStatus data_file_status, options.GetFileSystem()->GetFileStatus(data_file_path)); - ASSERT_GT(data_file_status->GetLen(), 0); + ASSERT_GT(data_file_status.GetLen(), 0); arrow::FieldVector write_fields = {arrow::field("_SEQUENCE_NUMBER", arrow::int64()), arrow::field("_VALUE_KIND", arrow::int8())}; @@ -439,7 +439,7 @@ TEST_P(PostponeBucketWriterTest, TestWriteMultiBatch) { // check data file exist and read ok std::string expected_data_file_name = "data-" + uuid + "-0." + file_format; std::string expected_data_file_path = dir->Str() + "/" + expected_data_file_name; - ASSERT_OK_AND_ASSIGN(std::unique_ptr data_file_status, + ASSERT_OK_AND_ASSIGN(FileStatus data_file_status, options.GetFileSystem()->GetFileStatus(expected_data_file_path)); std::shared_ptr expected_array; @@ -462,7 +462,7 @@ TEST_P(PostponeBucketWriterTest, TestWriteMultiBatch) { ASSERT_TRUE(commit_increment.GetCompactIncrement().IsEmpty()); ASSERT_EQ(1, commit_increment.GetNewFilesIncrement().NewFiles().size()); auto expected_data_file_meta = std::make_shared( - expected_data_file_name, /*file_size=*/data_file_status->GetLen(), /*row_count=*/9, + expected_data_file_name, /*file_size=*/data_file_status.GetLen(), /*row_count=*/9, /*min_key=*/BinaryRowGenerator::GenerateRow({std::string("David")}, pool_.get()), /*max_key=*/BinaryRowGenerator::GenerateRow({std::string("Tom")}, pool_.get()), /*key_stats=*/ @@ -581,10 +581,10 @@ TEST_P(PostponeBucketWriterTest, TestMultiplePrepareCommit) { std::string expected_data_file_dir = dir->Str() + "/"; ASSERT_OK_AND_ASSIGN( - std::unique_ptr data_file_status1, + FileStatus data_file_status1, options.GetFileSystem()->GetFileStatus(expected_data_file_dir + expected_data_file_name1)); ASSERT_OK_AND_ASSIGN( - std::unique_ptr data_file_status2, + FileStatus data_file_status2, options.GetFileSystem()->GetFileStatus(expected_data_file_dir + expected_data_file_name2)); std::shared_ptr expected_array1; @@ -612,7 +612,7 @@ TEST_P(PostponeBucketWriterTest, TestMultiplePrepareCommit) { ASSERT_TRUE(commit_increment1.GetCompactIncrement().IsEmpty()); ASSERT_EQ(1, commit_increment1.GetNewFilesIncrement().NewFiles().size()); auto expected_data_file_meta1 = std::make_shared( - expected_data_file_name1, /*file_size=*/data_file_status1->GetLen(), /*row_count=*/3, + expected_data_file_name1, /*file_size=*/data_file_status1.GetLen(), /*row_count=*/3, /*min_key=*/BinaryRowGenerator::GenerateRow({std::string("David")}, pool_.get()), /*max_key=*/BinaryRowGenerator::GenerateRow({std::string("Alex")}, pool_.get()), /*key_stats=*/ @@ -632,7 +632,7 @@ TEST_P(PostponeBucketWriterTest, TestMultiplePrepareCommit) { ASSERT_TRUE(commit_increment2.GetCompactIncrement().IsEmpty()); ASSERT_EQ(1, commit_increment2.GetNewFilesIncrement().NewFiles().size()); auto expected_data_file_meta2 = std::make_shared( - expected_data_file_name2, /*file_size=*/data_file_status2->GetLen(), /*row_count=*/2, + expected_data_file_name2, /*file_size=*/data_file_status2.GetLen(), /*row_count=*/2, /*min_key=*/BinaryRowGenerator::GenerateRow({std::string("Judy")}, pool_.get()), /*max_key=*/BinaryRowGenerator::GenerateRow({std::string("Tom")}, pool_.get()), /*key_stats=*/ diff --git a/src/paimon/core/table/sink/commit_message_impl_test.cpp b/src/paimon/core/table/sink/commit_message_impl_test.cpp index fb6daa25..b5d4af28 100644 --- a/src/paimon/core/table/sink/commit_message_impl_test.cpp +++ b/src/paimon/core/table/sink/commit_message_impl_test.cpp @@ -38,9 +38,8 @@ TEST(CommitMessageImplTest, TestToString) { "/orc/append_10_external_path.db/append_10_external_path/" "commit_messages/commit_messages-01"; auto file_system = std::make_shared(); - ASSERT_OK_AND_ASSIGN(std::unique_ptr file_status, - file_system->GetFileStatus(data_path)); - auto buffer_length = file_status->GetLen(); + ASSERT_OK_AND_ASSIGN(FileStatus file_status, file_system->GetFileStatus(data_path)); + auto buffer_length = file_status.GetLen(); std::vector buffer(buffer_length, 0); auto in_stream = file_system->Open(data_path).value_or(nullptr); diff --git a/src/paimon/core/table/sink/commit_message_test.cpp b/src/paimon/core/table/sink/commit_message_test.cpp index 935e45d2..d602da22 100644 --- a/src/paimon/core/table/sink/commit_message_test.cpp +++ b/src/paimon/core/table/sink/commit_message_test.cpp @@ -72,7 +72,7 @@ TEST(CommitMessageTest, TestCompatibleWithVersion12) { "orc/pk_btree_source_meta.db/pk_btree_source_meta/" "commit_messages/commit_messages-01"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); @@ -106,7 +106,7 @@ TEST(CommitMessageTest, TestCompatibleWithVersion11) { "orc/append_with_global_index_with_partition.db/append_with_global_index_with_partition/" "commit_messages/commit_messages-01"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); @@ -150,7 +150,7 @@ TEST(CommitMessageTest, TestCompatibleWithVersion10) { "pk_dv_index_with_commit_message_version10/" "commit_messages/commit_messages-01"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); @@ -214,7 +214,7 @@ TEST(CommitMessageTest, TestCompatibleWithVersion9) { "orc/pk_dv_index_not_in_data_no_external.db/pk_dv_index_not_in_data_no_external/" "commit_messages/commit_messages-01"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); @@ -279,7 +279,7 @@ TEST(CommitMessageTest, TestCompatibleWithVersion9WithExternalPathForIndex) { "orc/pk_dv_index_in_data_with_external.db/pk_dv_index_in_data_with_external/" "commit_messages/commit_messages-01"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); @@ -347,7 +347,7 @@ TEST(CommitMessageTest, TestCompatibleWithVersion8) { "/orc/append_table_with_first_row_id.db/append_table_with_first_row_id/" "commit_messages/commit_messages-01"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); @@ -415,7 +415,7 @@ TEST(CommitMessageTest, TestCompatibleWithVersion7) { "/orc/pk_table_with_total_buckets.db/pk_table_with_total_buckets/" "commit_messages/commit_messages-01"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); ASSERT_OK(in_stream->Read(reinterpret_cast(buffer.data()), buffer.size())); @@ -487,7 +487,7 @@ TEST(CommitMessageTest, TestCompatibleWithVersion6) { "/orc/append_10_external_path.db/append_10_external_path/" "commit_messages/commit_messages-01"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); ASSERT_OK(in_stream->Read(reinterpret_cast(buffer.data()), buffer.size())); @@ -532,7 +532,7 @@ TEST(CommitMessageTest, TestCompatibleWithVersion5) { "/orc/pk_table_with_dv_cardinality.db/pk_table_with_dv_cardinality/" "commit_messages/commit_messages-01"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); ASSERT_OK(in_stream->Read(reinterpret_cast(buffer.data()), buffer.size())); @@ -632,7 +632,7 @@ TEST(CommitMessageTest, TestCompatibleWithVersion4) { std::string data_path = paimon::test::GetDataDir() + "/orc/append_10.db/append_10/commit_messages/commit_messages-01"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); @@ -722,7 +722,7 @@ TEST(CommitMessageTest, TestCompatibleWithJavaPaimon10WithStatsDenseStore) { "/orc/append_10_stats_dense_store.db/append_10_stats_dense_store/commit_messages/" "commit_messages-01"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); @@ -804,7 +804,7 @@ TEST(CommitMessageTest, TestCompatibleWith09JavaPaimon1) { std::string data_path = paimon::test::GetDataDir() + "/orc/append_09.db/append_09/commit_messages/commit_messages-01"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); @@ -892,7 +892,7 @@ TEST(CommitMessageTest, TestCompatibleWith09JavaPaimon2) { std::string data_path = paimon::test::GetDataDir() + "/orc/append_09.db/append_09/commit_messages/commit_messages-02"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); @@ -962,7 +962,7 @@ TEST(CommitMessageTest, TestCompatibleWith09JavaPaimon3) { std::string data_path = paimon::test::GetDataDir() + "/orc/append_09.db/append_09/commit_messages/commit_messages-03"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); @@ -1013,7 +1013,7 @@ TEST(CommitMessageTest, TestPkTableCompatibleWithJavaPaimon09) { paimon::test::GetDataDir() + "/orc/pk_09_with_dv.db/pk_09_with_dv/commit_messages/commit_messages-01"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); @@ -1089,7 +1089,7 @@ TEST(CommitMessageTest, TestInvalidMessages) { std::string data_path = paimon::test::GetDataDir() + "/orc/append_09.db/append_09/commit_messages/commit_messages-01"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); @@ -1107,7 +1107,7 @@ TEST(CommitMessageTest, TestCompatibleWithComplexDataType) { std::string data_path = paimon::test::GetDataDir() + "/orc/append_09.db/append_09/commit_messages/commit_messages_complex"; auto file_system = std::make_shared(); - auto buffer_length = file_system->GetFileStatus(data_path).value()->GetLen(); + auto buffer_length = file_system->GetFileStatus(data_path).value().GetLen(); std::vector buffer(buffer_length, 0); ASSERT_OK_AND_ASSIGN(auto in_stream, file_system->Open(data_path)); @@ -1235,11 +1235,11 @@ TEST(CommitMessageTest, TestSerialize) { ASSERT_EQ(serialized_commit_message, reserialized_commit_message); auto fs = std::make_shared(); - std::vector> status_list; + std::vector status_list; ASSERT_OK(fs->ListDir(root_path + "/f0=true/f3=1/bucket-0/", &status_list)); int32_t file_nums = 0; for (const auto& file_status : status_list) { - if (!file_status->IsDir()) { + if (!file_status.IsDir()) { file_nums++; } } diff --git a/src/paimon/core/table/system/metadata_system_tables.cpp b/src/paimon/core/table/system/metadata_system_tables.cpp index 4467b7b9..09399431 100644 --- a/src/paimon/core/table/system/metadata_system_tables.cpp +++ b/src/paimon/core/table/system/metadata_system_tables.cpp @@ -724,12 +724,12 @@ Result> BranchesSystemTable::BuildRows() const { for (const auto& name : branches) { PAIMON_ASSIGN_OR_RAISE( - std::unique_ptr branch_status, + FileStatus branch_status, context_.fs->GetFileStatus(BranchManager::BranchPath(context_.table_path, name))); GenericRow row(schema->num_fields()); row.SetField(0, StringValue(name)); PAIMON_ASSIGN_OR_RAISE(VariantType create_time, - LocalTimestampMillisValue(branch_status->GetModificationTime())); + LocalTimestampMillisValue(branch_status.GetModificationTime())); row.SetField(1, create_time); rows.push_back(std::move(row)); } diff --git a/src/paimon/core/utils/branch_manager.cpp b/src/paimon/core/utils/branch_manager.cpp index 3ca52cf4..11818900 100644 --- a/src/paimon/core/utils/branch_manager.cpp +++ b/src/paimon/core/utils/branch_manager.cpp @@ -37,14 +37,14 @@ Result> BranchManager::ListBranches(const std::shared_p return branches; } - std::vector> file_status_list; + std::vector file_status_list; PAIMON_RETURN_NOT_OK(fs->ListDir(branch_dir, &file_status_list)); std::string branch_prefix = BRANCH_PREFIX; for (const auto& file_status : file_status_list) { - if (!file_status->IsDir()) { + if (!file_status.IsDir()) { continue; } - std::string dir_name = PathUtil::GetName(file_status->GetPath()); + std::string dir_name = PathUtil::GetName(file_status.GetPath()); if (StringUtils::StartsWith(dir_name, branch_prefix, /*start_pos=*/0)) { branches.push_back(dir_name.substr(branch_prefix.length())); } diff --git a/src/paimon/core/utils/consumer_manager.cpp b/src/paimon/core/utils/consumer_manager.cpp index 4bde69b1..131b8bb4 100644 --- a/src/paimon/core/utils/consumer_manager.cpp +++ b/src/paimon/core/utils/consumer_manager.cpp @@ -59,14 +59,14 @@ Result> ConsumerManager::ListConsumers() const { return consumers; } - std::vector> file_status_list; + std::vector file_status_list; PAIMON_RETURN_NOT_OK(fs_->ListDir(consumer_dir, &file_status_list)); std::string prefix = kConsumerPrefix; for (const auto& file_status : file_status_list) { - if (file_status->IsDir()) { + if (file_status.IsDir()) { continue; } - std::string file_name = PathUtil::GetName(file_status->GetPath()); + std::string file_name = PathUtil::GetName(file_status.GetPath()); if (StringUtils::StartsWith(file_name, prefix, /*start_pos=*/0)) { consumers.push_back(file_name.substr(prefix.length())); } diff --git a/src/paimon/core/utils/file_utils.cpp b/src/paimon/core/utils/file_utils.cpp index 1a803e32..6b087730 100644 --- a/src/paimon/core/utils/file_utils.cpp +++ b/src/paimon/core/utils/file_utils.cpp @@ -48,24 +48,24 @@ Status FileUtils::ListVersionedFiles(const std::shared_ptr& fs, cons Status FileUtils::ListOriginalVersionedFiles(const std::shared_ptr& fs, const std::string& dir, const std::string& prefix, std::vector* files) { - std::vector> file_status_list; + std::vector file_status_list; PAIMON_RETURN_NOT_OK(ListVersionedFileStatus(fs, dir, prefix, &file_status_list)); for (auto& file_status : file_status_list) { - std::string file_name = PathUtil::GetName(file_status->GetPath()); + std::string file_name = PathUtil::GetName(file_status.GetPath()); files->emplace_back(file_name.substr(prefix.size())); } return Status::OK(); } -Status FileUtils::ListVersionedFileStatus( - const std::shared_ptr& fs, const std::string& dir, const std::string& prefix, - std::vector>* file_status_list) { +Status FileUtils::ListVersionedFileStatus(const std::shared_ptr& fs, + const std::string& dir, const std::string& prefix, + std::vector* file_status_list) { PAIMON_ASSIGN_OR_RAISE(bool exist, fs->Exists(dir)); if (exist) { - std::vector> file_statuses; + std::vector file_statuses; PAIMON_RETURN_NOT_OK(fs->ListDir(dir, &file_statuses)); for (auto& file_status : file_statuses) { - std::string file_name = PathUtil::GetName(file_status->GetPath()); + std::string file_name = PathUtil::GetName(file_status.GetPath()); if (StringUtils::StartsWith(file_name, prefix)) { file_status_list->emplace_back(std::move(file_status)); } diff --git a/src/paimon/core/utils/file_utils.h b/src/paimon/core/utils/file_utils.h index cbb5f766..16901ae8 100644 --- a/src/paimon/core/utils/file_utils.h +++ b/src/paimon/core/utils/file_utils.h @@ -24,11 +24,10 @@ #include #include +#include "paimon/fs/file_system.h" #include "paimon/status.h" -#include "paimon/type_fwd.h" namespace paimon { -class BasicFileStatus; class FileSystem; /// Utils for file reading and writing. @@ -50,9 +49,9 @@ class FileUtils { /// List versioned file status for the directory. /// /// @return Status - static Status ListVersionedFileStatus( - const std::shared_ptr& fs, const std::string& dir, const std::string& prefix, - std::vector>* file_status_list); + static Status ListVersionedFileStatus(const std::shared_ptr& fs, + const std::string& dir, const std::string& prefix, + std::vector* file_status_list); }; } // namespace paimon diff --git a/src/paimon/core/utils/snapshot_manager.cpp b/src/paimon/core/utils/snapshot_manager.cpp index 398272a5..3d38a695 100644 --- a/src/paimon/core/utils/snapshot_manager.cpp +++ b/src/paimon/core/utils/snapshot_manager.cpp @@ -177,14 +177,14 @@ Result> SnapshotManager::FindByListFiles( Result> SnapshotManager::TryGetNonSnapshotFiles(int64_t older_than_ms) const { std::set non_snapshot_files; - std::vector> file_status_list; + std::vector file_status_list; PAIMON_RETURN_NOT_OK(fs_->ListFileStatus(SnapshotDirectory(), &file_status_list)); for (const auto& file_status : file_status_list) { - std::string file_name = PathUtil::GetName(file_status->GetPath()); + std::string file_name = PathUtil::GetName(file_status.GetPath()); if (!StringUtils::StartsWith(file_name, std::string(SNAPSHOT_PREFIX)) && file_name != std::string(EARLIEST) && file_name != std::string(LATEST)) { - if (file_status->GetModificationTime() < older_than_ms) { - non_snapshot_files.insert(file_status->GetPath()); + if (file_status.GetModificationTime() < older_than_ms) { + non_snapshot_files.insert(file_status.GetPath()); } } } @@ -239,11 +239,11 @@ Status SnapshotManager::CommitHint(int64_t snapshot_id, const std::string& file_ Result> SnapshotManager::GetAllSnapshots() const { std::vector snapshots; - std::vector> file_statuses; + std::vector file_statuses; PAIMON_RETURN_NOT_OK(FileUtils::ListVersionedFileStatus(fs_, SnapshotDirectory(), SNAPSHOT_PREFIX, &file_statuses)); for (const auto& file_status : file_statuses) { - auto snapshot_path = file_status->GetPath(); + auto snapshot_path = file_status.GetPath(); PAIMON_ASSIGN_OR_RAISE(Snapshot snapshot, Snapshot::FromPath(fs_, snapshot_path)); snapshots.push_back(snapshot); } diff --git a/src/paimon/core/utils/tag_manager.cpp b/src/paimon/core/utils/tag_manager.cpp index 25cdba2e..080a072b 100644 --- a/src/paimon/core/utils/tag_manager.cpp +++ b/src/paimon/core/utils/tag_manager.cpp @@ -65,14 +65,14 @@ Result> TagManager::ListTagNames() const { return tag_names; } - std::vector> file_status_list; + std::vector file_status_list; PAIMON_RETURN_NOT_OK(fs_->ListDir(tag_dir, &file_status_list)); std::string tag_prefix = TAG_PREFIX; for (const auto& file_status : file_status_list) { - if (file_status->IsDir()) { + if (file_status.IsDir()) { continue; } - std::string file_name = PathUtil::GetName(file_status->GetPath()); + std::string file_name = PathUtil::GetName(file_status.GetPath()); if (StringUtils::StartsWith(file_name, tag_prefix, /*start_pos=*/0)) { tag_names.push_back(file_name.substr(tag_prefix.length())); } diff --git a/src/paimon/format/avro/avro_stats_extractor_test.cpp b/src/paimon/format/avro/avro_stats_extractor_test.cpp index 081c9550..c4d554d6 100644 --- a/src/paimon/format/avro/avro_stats_extractor_test.cpp +++ b/src/paimon/format/avro/avro_stats_extractor_test.cpp @@ -70,7 +70,7 @@ class AvroStatsExtractorTest : public ::testing::Test { ASSERT_OK(writer->Finish()); ASSERT_OK_AND_ASSIGN(auto file_status, fs->GetFileStatus(file_path)); - ASSERT_GT(file_status->GetLen(), 0); + ASSERT_GT(file_status.GetLen(), 0); } private: diff --git a/src/paimon/format/orc/orc_input_output_stream_test.cpp b/src/paimon/format/orc/orc_input_output_stream_test.cpp index 7db0e1bb..fcbb1304 100644 --- a/src/paimon/format/orc/orc_input_output_stream_test.cpp +++ b/src/paimon/format/orc/orc_input_output_stream_test.cpp @@ -135,7 +135,7 @@ TEST(OrcInputOutputStreamTest, TestSimple) { ASSERT_OK_AND_ASSIGN(std::shared_ptr input_stream, file_system->Open(file_name)); ASSERT_OK_AND_ASSIGN(auto in_stream, OrcInputStreamImpl::Create(input_stream, DEFAULT_NATURAL_READ_SIZE)); - auto length = file_system->GetFileStatus(file_name).value()->GetLen(); + auto length = file_system->GetFileStatus(file_name).value().GetLen(); ASSERT_EQ(in_stream->getName(), normalized_file_name); ASSERT_EQ(in_stream->getLength(), length); ASSERT_EQ(in_stream->getNaturalReadSize(), 1024 * 1024); diff --git a/src/paimon/format/parquet/parquet_file_batch_reader_test.cpp b/src/paimon/format/parquet/parquet_file_batch_reader_test.cpp index 3810af8a..bcd4af21 100644 --- a/src/paimon/format/parquet/parquet_file_batch_reader_test.cpp +++ b/src/paimon/format/parquet/parquet_file_batch_reader_test.cpp @@ -203,7 +203,7 @@ class ParquetFileBatchReaderTest : public ::testing::Test, const std::optional& selection_bitmap, int32_t batch_size, bool enable_page_level_filter = false) const { EXPECT_OK_AND_ASSIGN(auto input_stream, fs_->Open(file_name)); - auto length = fs_->GetFileStatus(file_name).value()->GetLen(); + auto length = fs_->GetFileStatus(file_name).value().GetLen(); auto in_stream = std::make_unique(std::move(input_stream), length, pool_); auto storage_read_bytes = in_stream->StorageReadBytes(); @@ -385,7 +385,7 @@ TEST_F(ParquetFileBatchReaderTest, TestSetReadSchema) { "parquet/parquet_append_table.db/parquet_append_table/bucket-0/" "data-9ea62f34-1dca-49c1-bf7a-d37303d8fb76-0.parquet"; ASSERT_OK_AND_ASSIGN(std::shared_ptr input_stream, fs_->Open(file_name)); - auto length = fs_->GetFileStatus(file_name).value()->GetLen(); + auto length = fs_->GetFileStatus(file_name).value().GetLen(); auto in_stream = std::make_unique(std::move(input_stream), length, pool_); std::map options; diff --git a/src/paimon/format/parquet/variant_parquet_test.cpp b/src/paimon/format/parquet/variant_parquet_test.cpp index c17b7275..fd60e7b2 100644 --- a/src/paimon/format/parquet/variant_parquet_test.cpp +++ b/src/paimon/format/parquet/variant_parquet_test.cpp @@ -378,7 +378,7 @@ class VariantParquetTest : public ::testing::Test { void OpenFile(std::unique_ptr* file_reader, std::shared_ptr* file_schema) { ASSERT_OK_AND_ASSIGN(auto input_stream, fs_->Open(file_path_)); - auto length = fs_->GetFileStatus(file_path_).value()->GetLen(); + auto length = fs_->GetFileStatus(file_path_).value().GetLen(); auto in_stream = std::make_unique(std::move(input_stream), length, arrow_pool_); std::map options = {}; @@ -565,7 +565,7 @@ TEST_F(VariantParquetTest, WriteAndReadRoundTrip) { } ASSERT_OK_AND_ASSIGN(auto input_stream, fs_->Open(file_path_)); - auto length = fs_->GetFileStatus(file_path_).value()->GetLen(); + auto length = fs_->GetFileStatus(file_path_).value().GetLen(); auto in_stream = std::make_unique(std::move(input_stream), length, arrow_pool_); std::map options = {}; diff --git a/src/paimon/fs/jindo/jindo_file_status.h b/src/paimon/fs/jindo/jindo_file_status.h deleted file mode 100644 index bc291eb0..00000000 --- a/src/paimon/fs/jindo/jindo_file_status.h +++ /dev/null @@ -1,67 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -#pragma once - -#include -#include - -#include "JdoFileInfo.hpp" // NOLINT(build/include_subdir) -#include "paimon/fs/file_system.h" - -namespace paimon::jindo { -class JindoBasicFileStatus : public BasicFileStatus { - public: - explicit JindoBasicFileStatus(JdoFileInfo&& file_info) : file_info_(std::move(file_info)) {} - - std::string GetPath() const override { - return file_info_.getPath(); - } - - bool IsDir() const override { - return file_info_.isDir(); - } - - private: - JdoFileInfo file_info_; -}; - -class JindoFileStatus : public FileStatus { - public: - explicit JindoFileStatus(JdoFileInfo&& file_info) : file_info_(std::move(file_info)) {} - - std::string GetPath() const override { - return file_info_.getPath(); - } - - int64_t GetLen() const override { - return file_info_.getLength(); - } - - int64_t GetModificationTime() const override { - return file_info_.getMtime(); - } - - bool IsDir() const override { - return file_info_.isDir(); - } - - private: - JdoFileInfo file_info_; -}; -} // namespace paimon::jindo diff --git a/src/paimon/fs/jindo/jindo_file_system.cpp b/src/paimon/fs/jindo/jindo_file_system.cpp index 0bfb7c8c..3312f8be 100644 --- a/src/paimon/fs/jindo/jindo_file_system.cpp +++ b/src/paimon/fs/jindo/jindo_file_system.cpp @@ -31,7 +31,6 @@ #include "jdo_error.h" // NOLINT(build/include_subdir) #include "paimon/common/utils/math.h" #include "paimon/common/utils/path_util.h" -#include "paimon/fs/jindo/jindo_file_status.h" #include "paimon/fs/jindo/jindo_utils.h" namespace paimon::jindo { @@ -129,8 +128,8 @@ Status JindoFileSystem::Rename(const std::string& src, const std::string& dst) c return Status::Invalid( fmt::format("rename {} to {} failed, because: dst file already exist", src, dst)); } - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file_status, GetFileStatus(src)); - if (!file_status->IsDir() && dst.back() == '/') { + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, GetFileStatus(src)); + if (!file_status.IsDir() && dst.back() == '/') { return Status::Invalid( fmt::format("rename {} to {} failed, because: src file is not a dir", src, dst)); } @@ -153,23 +152,23 @@ Status JindoFileSystem::Delete(const std::string& path, bool recursive) const { return Status::OK(); } -Result> JindoFileSystem::GetFileStatus(const std::string& path) const { +Result JindoFileSystem::GetFileStatus(const std::string& path) const { JdoFileInfo file_info; PAIMON_RETURN_NOT_OK_FROM_JINDO(impl_->GetFileSystem()->getFileInfo(path, &file_info)); - return std::make_unique(std::move(file_info)); + return FileStatus(file_info.getPath(), file_info.getLength(), file_info.isDir(), + file_info.getMtime()); } -Status JindoFileSystem::ListDir( - const std::string& directory, - std::vector>* file_status_list) const { +Status JindoFileSystem::ListDir(const std::string& directory, + std::vector* file_status_list) const { PAIMON_ASSIGN_OR_RAISE(bool exist, Exists(directory)); if (!exist) { return Status::OK(); } - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file_status, GetFileStatus(directory)); - if (!file_status->IsDir()) { + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, GetFileStatus(directory)); + if (!file_status.IsDir()) { return Status::Invalid( - fmt::format("file {} already exists and is not a directory", file_status->GetPath())); + fmt::format("file {} already exists and is not a directory", file_status.GetPath())); } JdoListResult list_result; while (true) { @@ -178,8 +177,7 @@ Status JindoFileSystem::ListDir( auto file_infos = list_result.getFileInfos(); file_status_list->reserve(file_status_list->size() + file_infos.size()); for (auto& file_info : file_infos) { - file_status_list->push_back( - std::make_unique(std::move(file_info))); + file_status_list->emplace_back(file_info.getPath(), file_info.isDir()); } if (!list_result.isTruncated()) { break; @@ -191,8 +189,8 @@ Status JindoFileSystem::ListDir( return Status::OK(); } -Status JindoFileSystem::ListFileStatus( - const std::string& path, std::vector>* file_status_list) const { +Status JindoFileSystem::ListFileStatus(const std::string& path, + std::vector* file_status_list) const { PAIMON_ASSIGN_OR_RAISE(bool exist, Exists(path)); if (!exist) { return Status::OK(); @@ -207,8 +205,7 @@ Status JindoFileSystem::ListFileStatus( for (auto& file_info : file_infos) { // oss not support list FileStatus, only return BasicFileStatus // call GetFileStatus to return FileStatus - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file_status, - GetFileStatus(file_info.getPath())); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, GetFileStatus(file_info.getPath())); file_status_list->push_back(std::move(file_status)); } if (!list_result.isTruncated()) { diff --git a/src/paimon/fs/jindo/jindo_file_system.h b/src/paimon/fs/jindo/jindo_file_system.h index 89b3081e..fed08bcb 100644 --- a/src/paimon/fs/jindo/jindo_file_system.h +++ b/src/paimon/fs/jindo/jindo_file_system.h @@ -49,14 +49,13 @@ class JindoFileSystem : public FileSystem { Status Rename(const std::string& src, const std::string& dst) const override; Status Delete(const std::string& path, bool recursive = true) const override; - Result> GetFileStatus(const std::string& path) const override; + Result GetFileStatus(const std::string& path) const override; Status ListDir(const std::string& directory, - std::vector>* file_status_list) const override; + std::vector* file_status_list) const override; - Status ListFileStatus( - const std::string& path, - std::vector>* file_status_list) const override; + Status ListFileStatus(const std::string& path, + std::vector* file_status_list) const override; Result Exists(const std::string& path) const override; diff --git a/src/paimon/fs/jindo/jindo_file_system_test.cpp b/src/paimon/fs/jindo/jindo_file_system_test.cpp index efaf2a72..54c34a19 100644 --- a/src/paimon/fs/jindo/jindo_file_system_test.cpp +++ b/src/paimon/fs/jindo/jindo_file_system_test.cpp @@ -143,14 +143,14 @@ TEST(JindoFileSystemPaginationTest, TestListDirAcrossOssPageBoundary) { auto fs_factory = std::make_shared(); ASSERT_OK_AND_ASSIGN(std::unique_ptr fs, fs_factory->Create(test_dir, options)); - std::vector> file_statuses; + std::vector file_statuses; ASSERT_OK(fs->ListDir(test_dir, &file_statuses)); ASSERT_EQ(file_statuses.size(), kFileCount); std::unordered_set actual_paths; - for (const std::unique_ptr& file_status : file_statuses) { - ASSERT_TRUE(actual_paths.insert(file_status->GetPath()).second) - << "duplicate path: " << file_status->GetPath(); + for (const BasicFileStatus& file_status : file_statuses) { + ASSERT_TRUE(actual_paths.insert(file_status.GetPath()).second) + << "duplicate path: " << file_status.GetPath(); } for (int32_t i = 0; i < kFileCount; ++i) { std::string index = std::to_string(i); diff --git a/src/paimon/fs/local/local_file.cpp b/src/paimon/fs/local/local_file.cpp index c603154c..64530667 100644 --- a/src/paimon/fs/local/local_file.cpp +++ b/src/paimon/fs/local/local_file.cpp @@ -31,7 +31,6 @@ #include "paimon/common/utils/math.h" #include "paimon/common/utils/path_util.h" #include "paimon/common/utils/string_utils.h" -#include "paimon/fs/local/local_file_status.h" namespace paimon { @@ -97,8 +96,8 @@ Result LocalFile::IsFile() const { Result LocalFile::IsDir() const { CHECK_HOOK(); - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file_status, GetFileStatus()); - return file_status->IsDir(); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, GetFileStatus()); + return file_status.IsDir(); } Status LocalFile::List(std::vector* file_list) const { @@ -160,7 +159,7 @@ Result LocalFile::Mkdir() const { return mkdir(path_.c_str(), 0755) == 0; } -Result> LocalFile::GetFileStatus() const { +Result LocalFile::GetFileStatus() const { CHECK_HOOK(); struct stat buf; if (stat(path_.c_str(), &buf) < 0) { @@ -168,20 +167,19 @@ Result> LocalFile::GetFileStatus() const { return Status::IOError( fmt::format("get file '{}' status failed, ec: {}", path_, std::strerror(cur_errno))); } - return std::make_unique(path_, buf.st_size, buf.st_mtime * 1000, - S_ISDIR(buf.st_mode)); + return FileStatus(path_, buf.st_size, S_ISDIR(buf.st_mode), buf.st_mtime * 1000); } Result LocalFile::Length() const { CHECK_HOOK(); - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file_status, GetFileStatus()); - return file_status->GetLen(); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, GetFileStatus()); + return file_status.GetLen(); } Result LocalFile::LastModifiedTimeMs() const { CHECK_HOOK(); - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file_status, GetFileStatus()); - return file_status->GetModificationTime(); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, GetFileStatus()); + return file_status.GetModificationTime(); } std::unique_ptr LocalFile::GetParentFile() const { diff --git a/src/paimon/fs/local/local_file.h b/src/paimon/fs/local/local_file.h index 7ebfd13f..6a3c7ed8 100644 --- a/src/paimon/fs/local/local_file.h +++ b/src/paimon/fs/local/local_file.h @@ -28,6 +28,7 @@ #include #include +#include "paimon/fs/file_system.h" #include "paimon/macros.h" #include "paimon/result.h" #include "paimon/status.h" @@ -35,7 +36,6 @@ namespace paimon { class IOHook; -class LocalFileStatus; class LocalFile { public: @@ -52,7 +52,7 @@ class LocalFile { const std::string& GetPath() const; std::unique_ptr GetParentFile() const; Result Mkdir() const; - Result> GetFileStatus() const; + Result GetFileStatus() const; Result Length() const; Result LastModifiedTimeMs() const; Status OpenFile(bool is_read_file); diff --git a/src/paimon/fs/local/local_file_status.h b/src/paimon/fs/local/local_file_status.h deleted file mode 100644 index 8987d0be..00000000 --- a/src/paimon/fs/local/local_file_status.h +++ /dev/null @@ -1,76 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -#pragma once - -#include - -#include "paimon/fs/file_system.h" - -namespace paimon { - -class LocalBasicFileStatus : public BasicFileStatus { - public: - LocalBasicFileStatus(const std::string& path, bool is_dir) : path_(path), is_dir_(is_dir) {} - - std::string GetPath() const override { - return path_; - } - - bool IsDir() const override { - return is_dir_; - } - - private: - std::string path_; - bool is_dir_; -}; - -class LocalFileStatus : public FileStatus { - public: - LocalFileStatus(const std::string& path, int64_t length, int64_t last_modification_time, - bool is_dir) - : path_(path), - length_(length), - last_modification_time_(last_modification_time), - is_dir_(is_dir) {} - - std::string GetPath() const override { - return path_; - } - - int64_t GetLen() const override { - return length_; - } - - int64_t GetModificationTime() const override { - return last_modification_time_; - } - - bool IsDir() const override { - return is_dir_; - } - - private: - std::string path_; - int64_t length_; - int64_t last_modification_time_; - bool is_dir_; -}; - -} // namespace paimon diff --git a/src/paimon/fs/local/local_file_system.cpp b/src/paimon/fs/local/local_file_system.cpp index 917506be..6adc0309 100644 --- a/src/paimon/fs/local/local_file_system.cpp +++ b/src/paimon/fs/local/local_file_system.cpp @@ -28,7 +28,6 @@ #include "fmt/format.h" #include "paimon/common/utils/path_util.h" -#include "paimon/fs/local/local_file_status.h" namespace paimon { @@ -99,7 +98,7 @@ Status LocalFileSystem::MkdirsInternal(std::unique_ptr&& file) const return Status::OK(); } -Result> LocalFileSystem::GetFileStatus(const std::string& path) const { +Result LocalFileSystem::GetFileStatus(const std::string& path) const { PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file, LocalFile::Create(path)); PAIMON_ASSIGN_OR_RAISE(bool is_exist, file->Exists()); if (is_exist) { @@ -112,9 +111,8 @@ Result> LocalFileSystem::GetFileStatus(const std::st } } -Status LocalFileSystem::ListDir( - const std::string& directory, - std::vector>* file_status_list) const { +Status LocalFileSystem::ListDir(const std::string& directory, + std::vector* file_status_list) const { PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file, LocalFile::Create(directory)); PAIMON_ASSIGN_OR_RAISE(bool is_exist, file->Exists()); if (!is_exist) { @@ -129,21 +127,20 @@ Status LocalFileSystem::ListDir( PAIMON_RETURN_NOT_OK(file->List(&file_list)); file_status_list->reserve(file_status_list->size() + file_list.size()); for (const auto& f : file_list) { - Result> file_status = - GetFileStatus(PathUtil::JoinPath(directory, f)); + Result file_status = GetFileStatus(PathUtil::JoinPath(directory, f)); if (!file_status.ok() && !file_status.status().IsNotExist()) { return file_status.status(); } else if (file_status.ok()) { - file_status_list->emplace_back(std::make_unique( - file_status.value()->GetPath(), file_status.value()->IsDir())); + file_status_list->emplace_back(file_status.value().GetPath(), + file_status.value().IsDir()); } } return Status::OK(); } } -Status LocalFileSystem::ListFileStatus( - const std::string& path, std::vector>* file_status_list) const { +Status LocalFileSystem::ListFileStatus(const std::string& path, + std::vector* file_status_list) const { PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file, LocalFile::Create(path)); PAIMON_ASSIGN_OR_RAISE(bool is_exist, file->Exists()); if (!is_exist) { @@ -151,15 +148,14 @@ Status LocalFileSystem::ListFileStatus( } PAIMON_ASSIGN_OR_RAISE(bool is_file, file->IsFile()); if (is_file) { - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file_status, file->GetFileStatus()); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, file->GetFileStatus()); file_status_list->emplace_back(std::move(file_status)); } else { std::vector file_list; PAIMON_RETURN_NOT_OK(file->List(&file_list)); file_status_list->reserve(file_status_list->size() + file_list.size()); for (const auto& f : file_list) { - Result> file_status = - GetFileStatus(PathUtil::JoinPath(path, f)); + Result file_status = GetFileStatus(PathUtil::JoinPath(path, f)); if (!file_status.ok() && !file_status.status().IsNotExist()) { return file_status.status(); } else if (file_status.ok()) { diff --git a/src/paimon/fs/local/local_file_system.h b/src/paimon/fs/local/local_file_system.h index 12ee5646..e9d051c0 100644 --- a/src/paimon/fs/local/local_file_system.h +++ b/src/paimon/fs/local/local_file_system.h @@ -48,12 +48,11 @@ class LocalFileSystem : public FileSystem { Status Mkdirs(const std::string& path) const override; Status Rename(const std::string& src, const std::string& dst) const override; Status Delete(const std::string& path, bool recursive = true) const override; - Result> GetFileStatus(const std::string& path) const override; + Result GetFileStatus(const std::string& path) const override; Status ListDir(const std::string& directory, - std::vector>* file_status_list) const override; - Status ListFileStatus( - const std::string& path, - std::vector>* file_status_list) const override; + std::vector* file_status_list) const override; + Status ListFileStatus(const std::string& path, + std::vector* file_status_list) const override; Result Exists(const std::string& path) const override; private: diff --git a/src/paimon/global_index/lucene/lucene_global_index_test.cpp b/src/paimon/global_index/lucene/lucene_global_index_test.cpp index 4fac71a3..1e30b366 100644 --- a/src/paimon/global_index/lucene/lucene_global_index_test.cpp +++ b/src/paimon/global_index/lucene/lucene_global_index_test.cpp @@ -89,7 +89,7 @@ class LuceneGlobalIndexTest : public ::testing::Test, PAIMON_ASSIGN_OR_RAISE(auto result_metas, global_writer->Finish()); // check tmp dir - std::vector> file_status_list; + std::vector file_status_list; EXPECT_OK(fs_->ListDir(tmp_dir, &file_status_list)); EXPECT_EQ(file_status_list.size(), 1); diff --git a/src/paimon/global_index/lumina/lumina_file_io_test.cpp b/src/paimon/global_index/lumina/lumina_file_io_test.cpp index 39dd1b2b..7ac7327f 100644 --- a/src/paimon/global_index/lumina/lumina_file_io_test.cpp +++ b/src/paimon/global_index/lumina/lumina_file_io_test.cpp @@ -43,9 +43,9 @@ TEST_F(LuminaFileIOTest, TestSimple) { // check file exist ASSERT_OK_AND_ASSIGN(bool exist, fs->Exists(index_path)); ASSERT_TRUE(exist); - ASSERT_OK_AND_ASSIGN(std::unique_ptr file_status, fs->GetFileStatus(index_path)); - ASSERT_FALSE(file_status->IsDir()); - ASSERT_EQ(file_status->GetLen(), content.length()); + ASSERT_OK_AND_ASSIGN(FileStatus file_status, fs->GetFileStatus(index_path)); + ASSERT_FALSE(file_status.IsDir()); + ASSERT_EQ(file_status.GetLen(), content.length()); // read content ASSERT_OK_AND_ASSIGN(std::shared_ptr in, fs->Open(index_path)); diff --git a/src/paimon/global_index/tantivy/tantivy_java_compat_test.cpp b/src/paimon/global_index/tantivy/tantivy_java_compat_test.cpp index bd9339e3..b998402c 100644 --- a/src/paimon/global_index/tantivy/tantivy_java_compat_test.cpp +++ b/src/paimon/global_index/tantivy/tantivy_java_compat_test.cpp @@ -84,7 +84,7 @@ class JavaCompatTest : public ::testing::Test { std::string archive_path = PathUtil::JoinPath(fixture_dir, fixture_name); EXPECT_OK_AND_ASSIGN(auto file_status, fs_->GetFileStatus(archive_path)); - int64_t file_size = file_status->GetLen(); + int64_t file_size = file_status.GetLen(); EXPECT_GT(file_size, 4) << "fixture archive must exist and be > 4 bytes"; // Empty metadata (options not needed for cross-read — we use defaults) @@ -312,9 +312,8 @@ class FixedNameGlobalIndexFileWriter : public GlobalIndexFileWriter { return fs_->Create(ToPath(file_name), /*overwrite=*/true); } Result GetFileSize(const std::string& file_name) const override { - PAIMON_ASSIGN_OR_RAISE(std::unique_ptr file_status, - fs_->GetFileStatus(ToPath(file_name))); - return file_status->GetLen(); + PAIMON_ASSIGN_OR_RAISE(FileStatus file_status, fs_->GetFileStatus(ToPath(file_name))); + return file_status.GetLen(); } private: @@ -412,7 +411,7 @@ TEST_F(JavaCompatTest, CppWriteDefaultTokenizerForJavaCrossRead) { // Build a reader directly off the archive path (mirrors OpenFixture // but rooted at the cpp fixtures dir). ASSERT_OK_AND_ASSIGN(auto file_status, fs_->GetFileStatus(archive_path)); - int64_t file_size = file_status->GetLen(); + int64_t file_size = file_status.GetLen(); auto meta_bytes = std::make_shared(std::string("{}"), pool_.get()); GlobalIndexIOMeta io_meta(archive_path, file_size, meta_bytes); auto reader_factory = diff --git a/src/paimon/testing/mock/mock_file_system.h b/src/paimon/testing/mock/mock_file_system.h index a691c3de..8df35827 100644 --- a/src/paimon/testing/mock/mock_file_system.h +++ b/src/paimon/testing/mock/mock_file_system.h @@ -79,25 +79,6 @@ class MockOutputStream : public OutputStream { } }; -class MockFileStatus : public FileStatus { - public: - MockFileStatus() = default; - ~MockFileStatus() override = default; - - std::string GetPath() const override { - return ""; - } - int64_t GetLen() const override { - return 0; - } - int64_t GetModificationTime() const override { - return 0; - } - bool IsDir() const override { - return false; - } -}; - class MockFileSystem : public FileSystem { public: MockFileSystem() = default; @@ -121,15 +102,15 @@ class MockFileSystem : public FileSystem { Status Delete(const std::string& path, bool recursive = true) const override { return Status::OK(); } - Result> GetFileStatus(const std::string& path) const override { - return std::make_unique(); + Result GetFileStatus(const std::string& path) const override { + return FileStatus(/*path=*/"", /*length=*/0, /*is_dir=*/false, /*modification_time=*/0); } Status ListDir(const std::string& directory, - std::vector>* status_list) const override { + std::vector* status_list) const override { return Status::OK(); } Status ListFileStatus(const std::string& path, - std::vector>* status_list) const override { + std::vector* status_list) const override { return Status::OK(); } Result Exists(const std::string& path) const override { diff --git a/src/paimon/testing/utils/test_helper.h b/src/paimon/testing/utils/test_helper.h index dfca5f49..925a59ad 100644 --- a/src/paimon/testing/utils/test_helper.h +++ b/src/paimon/testing/utils/test_helper.h @@ -395,18 +395,18 @@ class TestHelper { static int64_t CountChannelFiles(const std::shared_ptr& file_system, const std::string& tmp_path) { - std::vector> dir_statuses; + std::vector dir_statuses; EXPECT_OK(file_system->ListDir(tmp_path, &dir_statuses)); int64_t channel_file_count = 0; for (const auto& dir_status : dir_statuses) { - const std::string dir_path = dir_status->GetPath(); - if (dir_status->IsDir() && dir_path.find("paimon-io-") != std::string::npos) { - std::vector> file_statuses; + const std::string dir_path = dir_status.GetPath(); + if (dir_status.IsDir() && dir_path.find("paimon-io-") != std::string::npos) { + std::vector file_statuses; EXPECT_OK(file_system->ListDir(dir_path, &file_statuses)); for (const auto& file_status : file_statuses) { - if (StringUtils::EndsWith(file_status->GetPath(), ".channel")) { + if (StringUtils::EndsWith(file_status.GetPath(), ".channel")) { ++channel_file_count; } } diff --git a/test/inte/pk_compaction_inte_test.cpp b/test/inte/pk_compaction_inte_test.cpp index a7441aca..18e77f9d 100644 --- a/test/inte/pk_compaction_inte_test.cpp +++ b/test/inte/pk_compaction_inte_test.cpp @@ -830,7 +830,7 @@ TEST_F(PkCompactionInteTest, CompactWithExternalPath) { { auto filesystem = external_dir->GetFileSystem(); auto bucket_dir = external_path + "/f1=10/bucket-0/"; - std::vector> file_statuses; + std::vector file_statuses; ASSERT_OK(filesystem->ListDir(bucket_dir, &file_statuses)); ASSERT_FALSE(file_statuses.empty()) << "External path directory should contain compact output files"; diff --git a/test/inte/read_inte_test.cpp b/test/inte/read_inte_test.cpp index 4ad82d0c..6d7e5562 100644 --- a/test/inte/read_inte_test.cpp +++ b/test/inte/read_inte_test.cpp @@ -3789,16 +3789,15 @@ TEST_P(ReadInteTest, TestSpecificFs) { Status Delete(const std::string& path, bool recursive = true) const override { return fs_->Delete(path, recursive); } - Result> GetFileStatus(const std::string& path) const override { + Result GetFileStatus(const std::string& path) const override { return fs_->GetFileStatus(path); } Status ListDir(const std::string& directory, - std::vector>* status_list) const override { + std::vector* status_list) const override { return fs_->ListDir(directory, status_list); } - Status ListFileStatus( - const std::string& path, - std::vector>* status_list) const override { + Status ListFileStatus(const std::string& path, + std::vector* status_list) const override { return fs_->ListFileStatus(path, status_list); } Result Exists(const std::string& path) const override { diff --git a/test/inte/scan_and_read_inte_test.cpp b/test/inte/scan_and_read_inte_test.cpp index b50d6dcc..96a7c9a1 100644 --- a/test/inte/scan_and_read_inte_test.cpp +++ b/test/inte/scan_and_read_inte_test.cpp @@ -119,14 +119,14 @@ class ScanAndReadInteTest : public testing::Test, void CheckPostponeFile(const std::string& root_path, const std::vector& subdirs) const { - std::vector> status_list; + std::vector status_list; auto file_system = std::make_shared(); for (const auto& dir : subdirs) { ASSERT_OK(file_system->ListDir(PathUtil::JoinPath(root_path, dir), &status_list)); } ASSERT_FALSE(status_list.empty()); for (const auto& file_status : status_list) { - std::string path = file_status->GetPath(); + std::string path = file_status.GetPath(); ASSERT_TRUE(path.find("-u-") != std::string::npos); ASSERT_TRUE(path.find("-s-") != std::string::npos); ASSERT_TRUE(path.find("-w-") != std::string::npos); diff --git a/test/inte/write_and_read_inte_test.cpp b/test/inte/write_and_read_inte_test.cpp index a14461b8..d41b1a4a 100644 --- a/test/inte/write_and_read_inte_test.cpp +++ b/test/inte/write_and_read_inte_test.cpp @@ -582,11 +582,11 @@ TEST_P(WriteAndReadInteTest, TestAppendExternalPath) { auto get_file_list_in_external_path = [&](const UniqueTestDirectory* external_dir) { auto fs = external_dir->GetFileSystem(); auto bucket_dir = external_dir->Str() + "/f1=10/bucket-0/"; - std::vector> all_file_status; + std::vector all_file_status; ASSERT_OK(fs->ListDir(bucket_dir, &all_file_status)); ASSERT_FALSE(all_file_status.empty()); for (const auto& file_status : all_file_status) { - ASSERT_TRUE(PathUtil::GetName(file_status->GetPath()).find("test-data-") != + ASSERT_TRUE(PathUtil::GetName(file_status.GetPath()).find("test-data-") != std::string::npos); } }; @@ -676,7 +676,7 @@ TEST_P(WriteAndReadInteTest, TestAppendExternalPathAndNoneExternalPathStrategy) // check external path does not have any data file { auto fs = external_dir->GetFileSystem(); - std::vector> file_status_list; + std::vector file_status_list; ASSERT_OK(fs->ListDir(external_test_dir, &file_status_list)); ASSERT_TRUE(file_status_list.empty()); } diff --git a/test/inte/write_inte_test.cpp b/test/inte/write_inte_test.cpp index ffa02140..451c3ce2 100644 --- a/test/inte/write_inte_test.cpp +++ b/test/inte/write_inte_test.cpp @@ -147,13 +147,13 @@ class WriteInteTest : public testing::Test, public ::testing::WithParamInterface void CheckFileCount(const std::string& root_path, const std::vector& subdirs, int32_t expect_file_count) const { - std::vector> status_list; + std::vector status_list; for (const auto& dir : subdirs) { ASSERT_OK(file_system_->ListDir(PathUtil::JoinPath(root_path, dir), &status_list)); } int32_t file_count = 0; for (const auto& file_status : status_list) { - if (!file_status->IsDir()) { + if (!file_status.IsDir()) { file_count++; } } @@ -2537,12 +2537,12 @@ TEST_P(WriteInteTest, TestWriteWithFieldId) { ASSERT_OK(CommitMessages(table_path, commit_messages, /*ignore_empty_commit=*/false)); // check data file has field id meta - std::vector> status_list; + std::vector status_list; ASSERT_OK(file_system_->ListDir(table_path + "/bucket-0/", &status_list)); std::vector files; for (const auto& file_status : status_list) { - if (!file_status->IsDir()) { - files.emplace_back(file_status->GetPath()); + if (!file_status.IsDir()) { + files.emplace_back(file_status.GetPath()); } } ASSERT_EQ(files.size(), 1);