diff --git a/kv_cache_manager/data_storage/nfs_backend.cc b/kv_cache_manager/data_storage/nfs_backend.cc index 405fc404f..f2dd1f1b0 100644 --- a/kv_cache_manager/data_storage/nfs_backend.cc +++ b/kv_cache_manager/data_storage/nfs_backend.cc @@ -1,6 +1,10 @@ #include "nfs_backend.h" +#include +#include +#include #include +#include #include #include "kv_cache_manager/common/hash/hash.h" @@ -10,6 +14,29 @@ namespace kv_cache_manager { +namespace { + +std::pair ShouldDeleteFile(const DataStorageUri &storage_uri) { + const std::string blkid = storage_uri.GetParam("blkid"); + if (blkid.empty()) { + return {EC_OK, true}; + } + + std::uint64_t blkid_value = 0; + const auto *begin = blkid.data(); + const auto *end = blkid.data() + blkid.size(); + const auto [ptr, ec] = std::from_chars(begin, end, blkid_value); + if (ec != std::errc{} || ptr != end) { + KVCM_LOG_WARN("Skip deleting nfs file for invalid blkid [%s], uri: [%s]", + blkid.c_str(), + storage_uri.ToUriString().c_str()); + return {EC_ERROR, false}; + } + return {EC_OK, blkid_value == 0}; +} + +} // namespace + NfsBackend::NfsBackend(std::shared_ptr metrics_registry) : DataStorageBackend(std::move(metrics_registry)) {} @@ -79,15 +106,54 @@ std::vector> NfsBackend::Create(const std:: std::vector NfsBackend::Delete(const std::vector &storage_uris, const std::string &trace_id, std::function cb) { - std::vector result(storage_uris.size(), EC_OK); - // not supported yet + std::vector result; + result.reserve(storage_uris.size()); + for (const auto &storage_uri : storage_uris) { + const auto [decision, should_delete] = ShouldDeleteFile(storage_uri); + if (decision != EC_OK) { + result.push_back(decision); + continue; + } + if (!should_delete) { + result.push_back(EC_OK); + continue; + } + std::filesystem::path file_path = storage_uri.GetPath(); + std::error_code ec; + const bool removed = std::filesystem::remove(file_path, ec); + if (ec) { + KVCM_LOG_ERROR("Failed to delete nfs file [%s]: [%s]", file_path.string().c_str(), ec.message().c_str()); + result.push_back(EC_ERROR); + continue; + } + if (!removed) { + KVCM_LOG_WARN("Try delete nfs file not exist, file: [%s]", file_path.string().c_str()); + } + result.push_back(EC_OK); + } + if (cb) { + cb(); + } return result; } std::vector NfsBackend::Exist(const std::vector &storage_uris) { - std::vector result(storage_uris.size(), true); - // not supported yet + std::vector result; + result.reserve(storage_uris.size()); + for (const auto &storage_uri : storage_uris) { + std::error_code ec; + const bool exists = std::filesystem::exists(storage_uri.GetPath(), ec); + if (ec) { + KVCM_LOG_ERROR("Failed to check nfs file [%s]: [%s]", storage_uri.GetPath().c_str(), ec.message().c_str()); + result.push_back(false); + continue; + } + result.push_back(exists); + } return result; } +std::vector NfsBackend::MightExist(const std::vector &storage_uris) { + return std::vector(storage_uris.size(), true); +} std::vector NfsBackend::Lock(const std::vector &storage_uris) { std::vector result(storage_uris.size(), EC_OK); // not supported yet diff --git a/kv_cache_manager/data_storage/nfs_backend.h b/kv_cache_manager/data_storage/nfs_backend.h index 2c61cb80d..25fe2cfeb 100644 --- a/kv_cache_manager/data_storage/nfs_backend.h +++ b/kv_cache_manager/data_storage/nfs_backend.h @@ -29,6 +29,7 @@ class NfsBackend : public DataStorageBackend { const std::string &trace_id, std::function cb) override; std::vector Exist(const std::vector &storage_uris) override; + std::vector MightExist(const std::vector &storage_uris) override; std::vector Lock(const std::vector &storage_uris) override; std::vector UnLock(const std::vector &storage_uris) override; diff --git a/kv_cache_manager/data_storage/test/nfs_backend_test.cc b/kv_cache_manager/data_storage/test/nfs_backend_test.cc index e12d066eb..539964374 100644 --- a/kv_cache_manager/data_storage/test/nfs_backend_test.cc +++ b/kv_cache_manager/data_storage/test/nfs_backend_test.cc @@ -1,3 +1,5 @@ +#include +#include #include #include @@ -114,31 +116,105 @@ TEST_F(NfsBackendTest, TestDeleteReturnsOkAndSameSize) { NfsBackend backend(metrics_registry_); std::shared_ptr spec(new NfsStorageSpec); spec->set_key_count_per_file(1); - spec->set_root_path("/data/"); + spec->set_root_path(GetPrivateTestRuntimeDataPath()); StorageConfig storage_config(DataStorageType::DATA_STORAGE_TYPE_NFS, "test", spec); ASSERT_EQ(EC_OK, backend.Open(storage_config, "fake_trace_id_1")); - std::vector uris(3); - auto res = backend.Delete(uris, "fake_trace_id_2", []() {}); + + const auto file1 = GetPrivateTestRuntimeDataPath() + "delete-1"; + const auto file2 = GetPrivateTestRuntimeDataPath() + "delete-2"; + { + std::ofstream(file1) << "1"; + std::ofstream(file2) << "2"; + } + ASSERT_TRUE(std::filesystem::exists(file1)); + ASSERT_TRUE(std::filesystem::exists(file2)); + std::vector uris; + uris.emplace_back("file://nfs_test" + file1); + uris.emplace_back("file://nfs_test" + file2); + uris.emplace_back("file://nfs_test" + GetPrivateTestRuntimeDataPath() + "missing"); + + bool callback_called = false; + auto res = backend.Delete(uris, "fake_trace_id_2", [&callback_called]() { callback_called = true; }); + ASSERT_TRUE(callback_called); ASSERT_EQ(res.size(), uris.size()); for (auto code : res) { ASSERT_EQ(code, EC_OK); } + ASSERT_FALSE(std::filesystem::exists(file1)); + ASSERT_FALSE(std::filesystem::exists(file2)); +} + +TEST_F(NfsBackendTest, TestDeleteOnlyRemovesMultiBlockFileForBlockZero) { + NfsBackend backend(metrics_registry_); + std::shared_ptr spec(new NfsStorageSpec); + spec->set_key_count_per_file(2); + spec->set_root_path(GetPrivateTestRuntimeDataPath()); + StorageConfig storage_config(DataStorageType::DATA_STORAGE_TYPE_NFS, "test", spec); + ASSERT_EQ(EC_OK, backend.Open(storage_config, "fake_trace_id_1")); + + const auto file = GetPrivateTestRuntimeDataPath() + "multi-block-delete"; + std::ofstream(file) << "shared"; + ASSERT_TRUE(std::filesystem::exists(file)); + + DataStorageUri block_zero("file://nfs_test" + file + "?blkid=0&size=6"); + DataStorageUri block_one("file://nfs_test" + file + "?blkid=1&size=6"); + DataStorageUri invalid_block("file://nfs_test" + file + "?blkid=bad&size=6"); + + auto skip_res = backend.Delete({block_one}, "fake_trace_id_2", []() {}); + ASSERT_EQ(skip_res.size(), 1u); + ASSERT_EQ(skip_res[0], EC_OK); + ASSERT_TRUE(std::filesystem::exists(file)); + + auto invalid_res = backend.Delete({invalid_block}, "fake_trace_id_3", []() {}); + ASSERT_EQ(invalid_res.size(), 1u); + ASSERT_EQ(invalid_res[0], EC_ERROR); + ASSERT_TRUE(std::filesystem::exists(file)); + + DataStorageUri missing_blkid("file://nfs_test" + file + "?size=6"); + auto missing_blkid_res = backend.Delete({missing_blkid}, "fake_trace_id_4", []() {}); + ASSERT_EQ(missing_blkid_res.size(), 1u); + ASSERT_EQ(missing_blkid_res[0], EC_OK); + ASSERT_FALSE(std::filesystem::exists(file)); + + std::ofstream(file) << "shared"; + ASSERT_TRUE(std::filesystem::exists(file)); + DataStorageUri empty_blkid("file://nfs_test" + file + "?blkid=&size=6"); + auto empty_blkid_res = backend.Delete({empty_blkid}, "fake_trace_id_5", []() {}); + ASSERT_EQ(empty_blkid_res.size(), 1u); + ASSERT_EQ(empty_blkid_res[0], EC_OK); + ASSERT_FALSE(std::filesystem::exists(file)); + + std::ofstream(file) << "shared"; + ASSERT_TRUE(std::filesystem::exists(file)); + auto delete_res = backend.Delete({block_zero}, "fake_trace_id_6", []() {}); + ASSERT_EQ(delete_res.size(), 1u); + ASSERT_EQ(delete_res[0], EC_OK); + ASSERT_FALSE(std::filesystem::exists(file)); } -TEST_F(NfsBackendTest, TestExistReturnsTrues) { +TEST_F(NfsBackendTest, TestExistReturnsActualFileState) { NfsBackend backend(metrics_registry_); std::shared_ptr spec(new NfsStorageSpec); spec->set_key_count_per_file(1); - spec->set_root_path("/data/"); + spec->set_root_path(GetPrivateTestRuntimeDataPath()); StorageConfig storage_config(DataStorageType::DATA_STORAGE_TYPE_NFS, "test", spec); ASSERT_EQ(EC_OK, backend.Open(storage_config, "fake_trace_id_1")); - // TODO(qisa.cb) 没实现 - std::vector uris(5); + + const auto file1 = GetPrivateTestRuntimeDataPath() + "exists-1"; + std::ofstream(file1) << "1"; + std::vector uris; + uris.emplace_back("file://nfs_test" + file1); + uris.emplace_back("file://nfs_test" + GetPrivateTestRuntimeDataPath() + "missing"); + auto res = backend.Exist(uris); ASSERT_EQ(res.size(), uris.size()); - for (bool flag : res) { - ASSERT_TRUE(flag); - } + ASSERT_TRUE(res[0]); + ASSERT_FALSE(res[1]); + + auto might_res = backend.MightExist(uris); + ASSERT_EQ(might_res.size(), uris.size()); + ASSERT_TRUE(might_res[0]); + ASSERT_TRUE(might_res[1]); } TEST_F(NfsBackendTest, TestLockAndUnLockReturnOk) {