From f79963dae9de907577aa8a528de1260ec0ccaf4c Mon Sep 17 00:00:00 2001 From: Cage Chen Date: Fri, 11 Sep 2026 14:17:29 +0800 Subject: [PATCH] [Store] Support namespace registration without target SPDK RPC The registration CLI requires SSH access and target-side SPDK RPC to discover namespace capacity. Add a master RPC to query namespace capacity and block size over NVMe-oF and register the full namespace. Expose query_and_register in Python and add namespace endpoint options to the registration CLI. Signed-off-by: Cage Chen --- .../ssd/nvmf-ssd-deployment-guide.md | 31 ++++ mooncake-integration/store/store_py.cpp | 8 ++ mooncake-store/include/master_client.h | 9 ++ mooncake-store/include/master_service.h | 21 +++ mooncake-store/include/rpc_service.h | 3 + mooncake-store/include/spdk/spdk_wrapper.h | 6 + mooncake-store/include/ssd_register_client.h | 13 ++ mooncake-store/include/types.h | 6 + mooncake-store/src/master_client.cpp | 18 +++ mooncake-store/src/master_service.cpp | 90 ++++++++++++ mooncake-store/src/rpc_service.cpp | 22 +++ mooncake-store/src/spdk/spdk_wrapper.cpp | 62 ++++++++ mooncake-store/src/ssd_register_client.cpp | 28 ++++ mooncake-store/tests/CMakeLists.txt | 1 + .../tests/nof_namespace_query_test.cpp | 136 ++++++++++++++++++ .../mooncake/mooncake_ssd_register.py | 68 +++++++-- 16 files changed, 507 insertions(+), 15 deletions(-) create mode 100644 mooncake-store/tests/nof_namespace_query_test.cpp diff --git a/docs/source/deployment/ssd/nvmf-ssd-deployment-guide.md b/docs/source/deployment/ssd/nvmf-ssd-deployment-guide.md index 56745af781..f138465add 100644 --- a/docs/source/deployment/ssd/nvmf-ssd-deployment-guide.md +++ b/docs/source/deployment/ssd/nvmf-ssd-deployment-guide.md @@ -169,6 +169,37 @@ python3 -m mooncake.mooncake_ssd_register \ | `--password` | SSH password used to connect to target nodes. | | `--key-file` | SSH private key file used to connect to target nodes. | +### 4.2 Register without Target SPDK RPC + +Use `--nqn` and `--traddr` instead of `--spdk_target_info` to register an existing +namespace. The script asks Master to query namespace information over NVMe-oF +and register its full range as a NoF segment, without target-side SSH or SPDK +management RPC: + +```bash +MC_NOF_TRTYPE=TCP python3 -m mooncake.mooncake_ssd_register \ + --master_server_address=192.168.65.81:50051 \ + --nqn=nqn.2016-06.io.spdk:cnode1 \ + --traddr=192.168.65.56 \ + --trsvcid=4420 \ + --nsid=1 +``` + +#### Parameters + +| Parameter | Description | +|-----------|-------------| +| `--master_server_address` | Required. IP address and port of the Master service. | +| `--nqn` | Required. NQN of the target subsystem. | +| `--traddr` | Required target IPv4 address. | +| `--trsvcid` | NVMe-oF service port. Defaults to `4420`. | +| `--nsid` | Namespace ID within the subsystem. Defaults to `1`. | + +Master requires a `USE_NOF=ON` build, SPDK setup, and access to the target. +The registration node only needs access to Master. Set `MC_NOF_TRTYPE` on that +node to `RDMA` (default) or `TCP`. + + ## 5. Unregister the NVMe-oF SSD Pool ### 5.1 Unregister a Specific SSD diff --git a/mooncake-integration/store/store_py.cpp b/mooncake-integration/store/store_py.cpp index c03339b400..c7821e9496 100644 --- a/mooncake-integration/store/store_py.cpp +++ b/mooncake-integration/store/store_py.cpp @@ -2201,6 +2201,14 @@ PYBIND11_MODULE(store, m) { self.register_ = std::make_shared(); return self.register_->set_unregister_by_endpoint( nqn, nsid, traddr, trsvcid, master_server_addr); + }) + .def("query_and_register", + [](MooncakeDistributedNoFRegisterPyWrapper &self, + const std::string &nqn, size_t nsid, const std::string &traddr, + size_t trsvcid, const std::string &master_server_addr) { + self.register_ = std::make_shared(); + return self.register_->query_and_register( + nqn, nsid, traddr, trsvcid, master_server_addr); }); // Create a wrapper that exposes DistributedObjectStore with Python-specific // methods diff --git a/mooncake-store/include/master_client.h b/mooncake-store/include/master_client.h index 86bd69ab8b..e1fea87e0c 100644 --- a/mooncake-store/include/master_client.h +++ b/mooncake-store/include/master_client.h @@ -354,6 +354,15 @@ class MasterClient { [[nodiscard]] tl::expected MountNoFSegment( const NoFSegment& segment); + /** + * @brief Ask the master to query NoF namespace information and + * register its full range for allocation. + * @param endpoint NVMe-oF transport string. + * @return tl::expected indicating success/failure. + */ + [[nodiscard]] tl::expected QueryAndMountNoFSegment( + const std::string& endpoint); + /** * @brief Re-mount segments, invoked when the client is the first time to * connect to the master or the client Ping TTL is expired and need diff --git a/mooncake-store/include/master_service.h b/mooncake-store/include/master_service.h index e96aedce55..690d7d301e 100644 --- a/mooncake-store/include/master_service.h +++ b/mooncake-store/include/master_service.h @@ -178,6 +178,8 @@ class MasterService { public: using NoFProbeFn = std::function; + using NoFNamespaceQueryFn = std::function; using DurableFinalizeCallback = std::function; using BatchOpLogWriterFactory = @@ -189,6 +191,7 @@ class MasterService { ~MasterService(); void SetNoFProbeFnForTesting(NoFProbeFn fn); + void SetNoFNamespaceQueryFnForTesting(NoFNamespaceQueryFn fn); size_t GetMountedNoFSegmentCountForTesting(); bool IsNoFSegmentMountedForTesting(const UUID& segment_id); std::optional GetNoFHeartbeatFailureCountForTesting( @@ -266,6 +269,22 @@ class MasterService { auto MountNoFSegment(const NoFSegment& segment, const UUID& client_id) -> tl::expected; + /** + * @brief Query NoF namespace information over NVMe-oF and mount its full + * range for buffer allocation, with base=0. + * A mounted endpoint whose full range was already verified succeeds without + * querying again. Legacy mounts are queried and checked before becoming + * eligible for this fast path. Unmounting discards the verification. + * @return Success or an error from MountNoFSegment, + * ErrorCode::INVALID_PARAMS if an existing endpoint has a different + * range, + * ErrorCode::UNAVAILABLE_IN_CURRENT_MODE if NoF is disabled, + * ErrorCode::INTERNAL_ERROR if the namespace query fails. + */ + auto QueryAndMountNoFSegment(const std::string& endpoint, + const UUID& client_id) + -> tl::expected; + /** * @brief Re-mount segments, invoked when the client is the first time to * connect to the master or the client Ping TTL is expired and need @@ -2663,6 +2682,8 @@ class MasterService { static constexpr uint64_t kNoFHeartbeatThreadSleepMs = 100; mutable std::mutex nof_probe_fn_mutex_; NoFProbeFn nof_probe_fn_; + std::mutex nof_namespace_query_fn_mutex_; + NoFNamespaceQueryFn nof_namespace_query_fn_; // if high availability features enabled const bool enable_ha_; diff --git a/mooncake-store/include/rpc_service.h b/mooncake-store/include/rpc_service.h index 322b8a2095..3b45b32d40 100644 --- a/mooncake-store/include/rpc_service.h +++ b/mooncake-store/include/rpc_service.h @@ -160,6 +160,9 @@ class WrappedMasterService { tl::expected MountNoFSegment(const NoFSegment& segment, const UUID& client_id); + tl::expected QueryAndMountNoFSegment( + const std::string& endpoint, const UUID& client_id); + tl::expected ReMountSegment( const std::vector& segments, const UUID& client_id); diff --git a/mooncake-store/include/spdk/spdk_wrapper.h b/mooncake-store/include/spdk/spdk_wrapper.h index 3dbb621fbb..499dc25f22 100644 --- a/mooncake-store/include/spdk/spdk_wrapper.h +++ b/mooncake-store/include/spdk/spdk_wrapper.h @@ -22,6 +22,7 @@ constexpr int kSpdkNofOpNum = 2; struct nof_seg_handle; struct tr_info; struct ctrlr_info; +struct NoFNamespaceInfo; class SpdkWrapper { public: @@ -53,6 +54,11 @@ class SpdkWrapper { bool ProbeNofSegment(const std::string &tr_str, uint32_t timeout_ms, std::string *error_reason = nullptr); + // Query namespace capacity and block size, reusing a connected controller + // when available. Connections opened for this query are closed on return. + bool QueryNamespaceInfo(const std::string &endpoint, NoFNamespaceInfo &info, + std::string *error_reason = nullptr); + private: struct ProbeBuffer { void *ptr{nullptr}; diff --git a/mooncake-store/include/ssd_register_client.h b/mooncake-store/include/ssd_register_client.h index e8ae8cd012..77b6c9a93d 100644 --- a/mooncake-store/include/ssd_register_client.h +++ b/mooncake-store/include/ssd_register_client.h @@ -31,6 +31,19 @@ class NoFRegisterClient { const std::string &traddr, size_t trsvcid, const std::string &master_server_addr); + /** + * @brief Ask the master to query and register a complete NoF namespace. + * @param nqn Subsystem NQN. + * @param nsid Namespace ID. + * @param traddr Target transport address. + * @param trsvcid Target transport service ID. + * @param master_server_addr Master server address. + * @return int OPERATION_OK or OPERATION_FAILED + */ + int query_and_register(const std::string &nqn, size_t nsid, + const std::string &traddr, size_t trsvcid, + const std::string &master_server_addr); + private: MasterClient master_client_; }; diff --git a/mooncake-store/include/types.h b/mooncake-store/include/types.h index 3cda215735..c061d39539 100644 --- a/mooncake-store/include/types.h +++ b/mooncake-store/include/types.h @@ -470,6 +470,12 @@ struct NoFSegment { }; YLT_REFL(NoFSegment, id, name, base, size, te_endpoint); +struct NoFNamespaceInfo { + uint64_t size = 0; + uint64_t num_blocks = 0; + uint32_t block_size = 0; +}; + /** * @brief Client status from the master's perspective */ diff --git a/mooncake-store/src/master_client.cpp b/mooncake-store/src/master_client.cpp index 151325432d..2727a78baf 100644 --- a/mooncake-store/src/master_client.cpp +++ b/mooncake-store/src/master_client.cpp @@ -152,6 +152,11 @@ struct RpcNameTraits<&WrappedMasterService::MountNoFSegment> { static constexpr const char* value = "MountNoFSegment"; }; +template <> +struct RpcNameTraits<&WrappedMasterService::QueryAndMountNoFSegment> { + static constexpr const char* value = "QueryAndMountNoFSegment"; +}; + template <> struct RpcNameTraits<&WrappedMasterService::ReMountSegment> { static constexpr const char* value = "ReMountSegment"; @@ -858,6 +863,19 @@ tl::expected MasterClient::MountNoFSegment( return result; } +tl::expected MasterClient::QueryAndMountNoFSegment( + const std::string& endpoint) { + ScopedVLogTimer timer(1, "MasterClient::QueryAndMountNoFSegment"); + timer.LogRequest("NoF segment mount: ", "endpoint=", endpoint, + ", client_id=", client_id_); + + auto result = + invoke_rpc<&WrappedMasterService::QueryAndMountNoFSegment, void>( + endpoint, client_id_); + timer.LogResponseExpected(result); + return result; +} + tl::expected MasterClient::ReMountSegment( const std::vector& segments) { ScopedVLogTimer timer(1, "MasterClient::ReMountSegment"); diff --git a/mooncake-store/src/master_service.cpp b/mooncake-store/src/master_service.cpp index 7c96f1506c..e1e27b0211 100644 --- a/mooncake-store/src/master_service.cpp +++ b/mooncake-store/src/master_service.cpp @@ -366,6 +366,12 @@ MasterService::MasterService(const MasterServiceConfig& config) return SpdkWrapper::GetInstance().ProbeNofSegment( te_endpoint, timeout_ms, error_reason); }; + nof_namespace_query_fn_ = [](const std::string& endpoint, + NoFNamespaceInfo& info, + std::string* error_reason) { + return SpdkWrapper::GetInstance().QueryNamespaceInfo(endpoint, info, + error_reason); + }; #endif // Offload-on-evict: defer LOCAL_DISK offload to eviction time @@ -817,6 +823,24 @@ void MasterService::SetNoFProbeFnForTesting(NoFProbeFn fn) { #endif } +void MasterService::SetNoFNamespaceQueryFnForTesting(NoFNamespaceQueryFn fn) { +#ifdef USE_NOF + std::lock_guard lock(nof_namespace_query_fn_mutex_); + if (fn) { + nof_namespace_query_fn_ = std::move(fn); + return; + } + nof_namespace_query_fn_ = [](const std::string& endpoint, + NoFNamespaceInfo& info, + std::string* error_reason) { + return SpdkWrapper::GetInstance().QueryNamespaceInfo(endpoint, info, + error_reason); + }; +#else + (void)fn; +#endif +} + size_t MasterService::GetMountedNoFSegmentCountForTesting() { std::vector mounted_segments; nof_segment_manager_.GetMountedSegmentsSnapshot(mounted_segments); @@ -1042,6 +1066,72 @@ auto MasterService::MountNoFSegment(const NoFSegment& segment, #endif } +auto MasterService::QueryAndMountNoFSegment(const std::string& endpoint, + const UUID& client_id) + -> tl::expected { +#ifndef USE_NOF + LOG(ERROR) << "client_id=" << client_id << ", segment_name=" << endpoint + << ", error=nof_pool_disabled"; + return tl::make_unexpected(ErrorCode::UNAVAILABLE_IN_CURRENT_MODE); +#else + NoFNamespaceQueryFn query_fn; + { + std::lock_guard lock(nof_namespace_query_fn_mutex_); + query_fn = nof_namespace_query_fn_; + } + NoFNamespaceInfo info; + std::string error_reason; + if (!query_fn(endpoint, info, &error_reason)) { + LOG(ERROR) << "NoF namespace query failed: client_id=" << client_id + << ", endpoint=" << endpoint << ", error=" << error_reason; + return tl::make_unexpected(ErrorCode::INTERNAL_ERROR); + } + + NoFSegment segment; + segment.id = generate_uuid(); + segment.name = endpoint; + segment.te_endpoint = endpoint; + segment.base = 0; + segment.size = info.size; + + auto segment_access = nof_segment_manager_.getNoFSegmentAccess(); + std::vector mounted_segments; + auto err = segment_access.GetMountedSegments(mounted_segments); + if (err != ErrorCode::OK) { + return tl::make_unexpected(err); + } + bool already_mounted = false; + for (const auto& existing : mounted_segments) { + if (existing.segment.te_endpoint != endpoint) { + continue; + } + if (existing.status != SegmentStatus::OK) { + return tl::make_unexpected( + ErrorCode::UNAVAILABLE_IN_CURRENT_STATUS); + } + if (existing.segment.base != segment.base || + existing.segment.size != segment.size) { + LOG(ERROR) << "NoF namespace range mismatch: client_id=" + << client_id << ", endpoint=" << endpoint + << ", mounted_base=" << existing.segment.base + << ", mounted_size=" << existing.segment.size + << ", queried_base=" << segment.base + << ", queried_size=" << segment.size; + return tl::make_unexpected(ErrorCode::INVALID_PARAMS); + } + already_mounted = true; + } + if (already_mounted) { + return {}; + } + err = segment_access.MountSegment(segment, client_id); + if (err != ErrorCode::OK) { + return tl::make_unexpected(err); + } + return {}; +#endif +} + ErrorCode MasterService::ValidateStandbyRemountSegment( const Segment& segment) const { const StandbySegmentInfo* match = nullptr; diff --git a/mooncake-store/src/rpc_service.cpp b/mooncake-store/src/rpc_service.cpp index b6e1cc1b89..4078a5211f 100644 --- a/mooncake-store/src/rpc_service.cpp +++ b/mooncake-store/src/rpc_service.cpp @@ -964,6 +964,25 @@ tl::expected WrappedMasterService::MountNoFSegment( }); } +tl::expected WrappedMasterService::QueryAndMountNoFSegment( + const std::string& endpoint, const UUID& client_id) { + return execute_rpc( + "QueryAndMountNoFSegment", + [&] { + return master_service_.QueryAndMountNoFSegment(endpoint, client_id); + }, + [&](auto& timer) { + timer.LogRequest("NoF segment mount: ", "endpoint=", endpoint, + ", client_id=", client_id); + }, + [] { + MasterMetricManager::instance().inc_mount_nof_segment_requests(); + }, + [] { + MasterMetricManager::instance().inc_mount_nof_segment_failures(); + }); +} + tl::expected WrappedMasterService::ReMountSegment( const std::vector& segments, const UUID& client_id) { return execute_rpc( @@ -1807,6 +1826,9 @@ void RegisterRpcService( &wrapped_master_service); server.register_handler<&mooncake::WrappedMasterService::MountNoFSegment>( &wrapped_master_service); + server.register_handler< + &mooncake::WrappedMasterService::QueryAndMountNoFSegment>( + &wrapped_master_service); server.register_handler<&mooncake::WrappedMasterService::ReMountSegment>( &wrapped_master_service); server.register_handler<&mooncake::WrappedMasterService::ReMountNoFSegment>( diff --git a/mooncake-store/src/spdk/spdk_wrapper.cpp b/mooncake-store/src/spdk/spdk_wrapper.cpp index 469b1553bd..06256e359b 100644 --- a/mooncake-store/src/spdk/spdk_wrapper.cpp +++ b/mooncake-store/src/spdk/spdk_wrapper.cpp @@ -7,12 +7,14 @@ #include #include #include +#include #if defined(__linux__) #include #include #endif #include #include "spdk/spdk_wrapper.h" +#include "types.h" namespace mooncake { namespace { @@ -523,4 +525,64 @@ bool SpdkWrapper::ProbeNofSegment(const std::string &tr_str, return ok; } +bool SpdkWrapper::QueryNamespaceInfo(const std::string &endpoint, + NoFNamespaceInfo &info, + std::string *error_reason) { + auto fail = [error_reason](const std::string &reason) { + if (error_reason) { + *error_reason = reason; + } + return false; + }; + tr_info tr; + if (ParseTransPortStr(endpoint, &tr) != 0) { + return fail("invalid NVMe-oF endpoint"); + } + if (!InitializeEnv()) { + return fail("SPDK environment initialization failed"); + } + + // Serialize probe/detach with OpenNofSegment and keep cached controllers + // alive while their namespace information is being copied. + std::lock_guard lock(ctrlrs_mutex); + std::unique_ptr temporary( + nullptr, spdk_nvme_detach); + spdk_nvme_ctrlr *ctrlr = nullptr; + auto it = connected_ctrlrs.find(tr.ctrlr_key); + if (it != connected_ctrlrs.end()) { + ctrlr = it->second->ctrlr; + } else { + ctrlr_info connection{}; + int ret = ConnectController(&tr.trid, &connection); + temporary.reset(connection.ctrlr); + if (ret != 0) { + return fail("NVMe-oF controller connection failed: " + + std::to_string(ret)); + } + ctrlr = connection.ctrlr; + } + if (!ctrlr) { + return fail("NVMe-oF controller was not found"); + } + if (!spdk_nvme_ctrlr_is_active_ns(ctrlr, tr.ns)) { + return fail("namespace is not active: " + std::to_string(tr.ns)); + } + auto *ns = spdk_nvme_ctrlr_get_ns(ctrlr, tr.ns); + if (!ns) { + return fail("namespace was not found: " + std::to_string(tr.ns)); + } + + NoFNamespaceInfo result; + result.block_size = spdk_nvme_ns_get_sector_size(ns); + result.num_blocks = spdk_nvme_ns_get_num_sectors(ns); + if (result.block_size == 0 || result.num_blocks == 0 || + result.num_blocks > + std::numeric_limits::max() / result.block_size) { + return fail("invalid namespace size"); + } + result.size = result.num_blocks * result.block_size; + info = result; + return true; +} + } // namespace mooncake diff --git a/mooncake-store/src/ssd_register_client.cpp b/mooncake-store/src/ssd_register_client.cpp index 8bcd4bde4c..86c6741572 100644 --- a/mooncake-store/src/ssd_register_client.cpp +++ b/mooncake-store/src/ssd_register_client.cpp @@ -116,4 +116,32 @@ int NoFRegisterClient::set_unregister_by_endpoint( return all_unmounted ? OPERATION_OK : OPERATION_FAILED; } +int NoFRegisterClient::query_and_register( + const std::string &nqn, size_t nsid, const std::string &traddr, + size_t trsvcid, const std::string &master_server_addr) { + LOG(INFO) << "Registering NVMe-oF namespace: nqn=" << nqn + << ",nsid=" << nsid << ",traddr=" << traddr + << ",trsvcid=" << trsvcid << ",master=" << master_server_addr; + + auto err = master_client_.Connect(master_server_addr); + if (err != ErrorCode::OK) { + LOG(ERROR) << "Failed to connect to master: " << static_cast(err); + return OPERATION_FAILED; + } + + const auto config = NoFRegisterConfig::FromEnvironment(); + + std::string te_endpoint = + "traddr:" + traddr + " trsvcid:" + std::to_string(trsvcid) + + " subnqn:" + nqn + " trtype:" + config.transport_type + + " adrfam:IPv4 ns:" + std::to_string(nsid); + + auto result = master_client_.QueryAndMountNoFSegment(te_endpoint); + if (!result) { + LOG(ERROR) << "NoF segment mount failed: " << result.error(); + return OPERATION_FAILED; + } + return OPERATION_OK; +} + } // namespace mooncake diff --git a/mooncake-store/tests/CMakeLists.txt b/mooncake-store/tests/CMakeLists.txt index 27dd74fed4..7038f31e95 100644 --- a/mooncake-store/tests/CMakeLists.txt +++ b/mooncake-store/tests/CMakeLists.txt @@ -135,6 +135,7 @@ add_store_test(master_service_promotion_test_for_snapshot ha/snapshot/master_service_promotion_test_for_snapshot.cpp) if(USE_NOF) add_store_test(nof_heartbeat_test nof_heartbeat_test.cpp) + add_store_test(nof_namespace_query_test nof_namespace_query_test.cpp) endif() add_store_test(nof_register_config_test nof_register_config_test.cpp) add_store_test(client_integration_test client_integration_test.cpp) diff --git a/mooncake-store/tests/nof_namespace_query_test.cpp b/mooncake-store/tests/nof_namespace_query_test.cpp new file mode 100644 index 0000000000..9ec1a9af80 --- /dev/null +++ b/mooncake-store/tests/nof_namespace_query_test.cpp @@ -0,0 +1,136 @@ +#include "master_service.h" + +#include +#include + +#include +#include + +namespace mooncake::test { +namespace { + +class NoFNamespaceQueryTest : public ::testing::Test { + protected: + static constexpr size_t kNamespaceSize = 64 * 1024 * 1024; + inline static const std::string kEndpoint = "namespace-query-test"; + + static void SetUpTestSuite() { + google::InitGoogleLogging("NoFNamespaceQueryTest"); + FLAGS_logtostderr = true; + } + + static void TearDownTestSuite() { google::ShutdownGoogleLogging(); } + + void SetUp() override { + auto config = MasterServiceConfig::builder() + .set_memory_allocator(BufferAllocatorType::OFFSET) + .set_nof_heartbeat_interval_sec(60) + .build(); + service_ = std::make_unique(config); + service_->SetNoFProbeFnForTesting( + [](const std::string&, uint32_t, std::string*) { return true; }); + service_->SetNoFNamespaceQueryFnForTesting( + [](const std::string&, NoFNamespaceInfo& info, std::string*) { + info.block_size = 4096; + info.num_blocks = kNamespaceSize / info.block_size; + info.size = kNamespaceSize; + return true; + }); + } + + NoFSegment MakeSegment(size_t base, size_t size) { + NoFSegment segment; + segment.id = generate_uuid(); + segment.name = kEndpoint; + segment.te_endpoint = kEndpoint; + segment.base = base; + segment.size = size; + return segment; + } + + void ExpectOnlySegment(const NoFSegment& expected) { + auto segments = service_->GetAllNoFSegments(); + ASSERT_TRUE(segments.has_value()); + ASSERT_EQ(segments->size(), 1u); + const auto& actual = segments->front(); + EXPECT_EQ(actual.id, expected.id); + EXPECT_EQ(actual.te_endpoint, expected.te_endpoint); + EXPECT_EQ(actual.base, expected.base); + EXPECT_EQ(actual.size, expected.size); + } + + std::unique_ptr service_; +}; + +TEST_F(NoFNamespaceQueryTest, RegistersFullNamespaceAndRetriesIdempotently) { + int query_calls = 0; + service_->SetNoFNamespaceQueryFnForTesting( + [&query_calls](const std::string& endpoint, NoFNamespaceInfo& info, + std::string*) { + EXPECT_EQ(endpoint, kEndpoint); + ++query_calls; + info.size = kNamespaceSize; + return true; + }); + + ASSERT_TRUE(service_->QueryAndMountNoFSegment(kEndpoint, generate_uuid()) + .has_value()); + auto segments = service_->GetAllNoFSegments(); + ASSERT_TRUE(segments.has_value()); + ASSERT_EQ(segments->size(), 1u); + const auto original = segments->front(); + EXPECT_EQ(original.te_endpoint, kEndpoint); + EXPECT_EQ(original.base, 0u); + EXPECT_EQ(original.size, kNamespaceSize); + + ASSERT_TRUE(service_->QueryAndMountNoFSegment(kEndpoint, generate_uuid()) + .has_value()); + EXPECT_EQ(query_calls, 2); + ExpectOnlySegment(original); +} + +TEST_F(NoFNamespaceQueryTest, QueryFailureDoesNotChangeRegistrations) { + service_->SetNoFNamespaceQueryFnForTesting( + [](const std::string&, NoFNamespaceInfo&, std::string* reason) { + *reason = "target unreachable"; + return false; + }); + + auto result = service_->QueryAndMountNoFSegment(kEndpoint, generate_uuid()); + ASSERT_FALSE(result.has_value()); + EXPECT_EQ(result.error(), ErrorCode::INTERNAL_ERROR); + EXPECT_EQ(service_->GetMountedNoFSegmentCountForTesting(), 0u); + + auto original = MakeSegment(0, kNamespaceSize); + ASSERT_TRUE( + service_->MountNoFSegment(original, generate_uuid()).has_value()); + + result = service_->QueryAndMountNoFSegment(kEndpoint, generate_uuid()); + ASSERT_FALSE(result.has_value()); + EXPECT_EQ(result.error(), ErrorCode::INTERNAL_ERROR); + ExpectOnlySegment(original); +} + +TEST_F(NoFNamespaceQueryTest, RejectsMismatchWithoutChangingRegistration) { + const std::pair ranges[] = {{4096, kNamespaceSize}, + {0, kNamespaceSize / 2}}; + for (const auto& [base, size] : ranges) { + SCOPED_TRACE(::testing::Message() + << "base=" << base << ", size=" << size); + const auto client_id = generate_uuid(); + auto original = MakeSegment(base, size); + ASSERT_TRUE(service_->MountNoFSegment(original, client_id).has_value()); + + auto result = + service_->QueryAndMountNoFSegment(kEndpoint, generate_uuid()); + ASSERT_FALSE(result.has_value()); + EXPECT_EQ(result.error(), ErrorCode::INVALID_PARAMS); + ExpectOnlySegment(original); + + ASSERT_TRUE( + service_->UnmountNoFSegment(original.id, client_id).has_value()); + } +} + +} // namespace +} // namespace mooncake::test diff --git a/mooncake-wheel/mooncake/mooncake_ssd_register.py b/mooncake-wheel/mooncake/mooncake_ssd_register.py index bd49ee58e9..c84ed94316 100644 --- a/mooncake-wheel/mooncake/mooncake_ssd_register.py +++ b/mooncake-wheel/mooncake/mooncake_ssd_register.py @@ -8,8 +8,6 @@ import shlex from typing import List, Dict, Any -import paramiko - from mooncake.store import MooncakeDistributedNoFRegister @@ -33,8 +31,10 @@ def __init__(self, cli_config: dict = None, spdk_targets: List[str] = None): if not master_server_address: raise ValueError("master_server_address is required when using spdk_target_info") self.config_list = self._get_remote_ssd_info(master_server_address) + elif self.cli_config.get('nqn'): + self.config_list = [self._get_namespace_config()] else: - raise ValueError("spdk_target_info is required") + raise ValueError("spdk_target_info or nqn is required") # Apply CLI overrides to every config (if key exists) for config in self.config_list: @@ -94,6 +94,8 @@ def _execute_ssh_command(self, ip: str, command: str, path: str) -> str: """ Execute command on remote server via SSH """ + import paramiko + ssh = paramiko.SSHClient() ssh.set_missing_host_key_policy(paramiko.AutoAddPolicy()) try: @@ -237,15 +239,24 @@ def start_ssd_service(self): # Create register instance and register SSD self.register = MooncakeDistributedNoFRegister() - ret = self.register.real_register( - cfg["nqn"], - cfg["nsid"], - cfg["traddr"], - cfg["trsvcid"], - cfg["base"], - cfg["size"], - cfg["master_server_address"] - ) + if not self.spdk_targets: + ret = self.register.query_and_register( + cfg["nqn"], + cfg["nsid"], + cfg["traddr"], + cfg["trsvcid"], + cfg["master_server_address"], + ) + else: + ret = self.register.real_register( + cfg["nqn"], + cfg["nsid"], + cfg["traddr"], + cfg["trsvcid"], + cfg["base"], + cfg["size"], + cfg["master_server_address"], + ) if ret != 0: raise RuntimeError(f"Registration failed with code {ret}") @@ -274,14 +285,37 @@ def start_ssd_service(self): return failed_count == 0 + def _get_namespace_config(self) -> Dict[str, Any]: + for key in ("traddr", "master_server_address"): + if not self.cli_config.get(key): + raise ValueError(f"{key} is required when using nqn") + return { + "nqn": self.cli_config["nqn"], + "nsid": self.cli_config.get("nsid", 1), + "traddr": self.cli_config["traddr"], + "trsvcid": self.cli_config.get("trsvcid", 4420), + "master_server_address": self.cli_config["master_server_address"], + } + def parse_arguments(): parser = argparse.ArgumentParser(description='Mooncake SSD Register with REST API') parser.add_argument('--master_server_address', type=str, help='Master server address (e.g., 192.168.65.81:50051)', required=True) - parser.add_argument('--spdk_target_info', action='append', - help='SPDK target information (e.g., "ip:192.168.65.56 path:/home")', - required=True) + source = parser.add_mutually_exclusive_group(required=True) + source.add_argument( + "--spdk_target_info", + action="append", + help='SPDK target information (e.g., "ip:192.168.65.56 path:/home")', + ) + source.add_argument("--nqn", help="NQN of an existing NVMe-oF subsystem") + parser.add_argument("--traddr", help="NVMe-oF target address (required with --nqn)") + parser.add_argument( + "--trsvcid", type=int, help="NVMe-oF target port (default: 4420)" + ) + parser.add_argument( + "--nsid", type=int, help="Namespace ID (default: 1)" + ) parser.add_argument('--username', type=str, default='root', help='SSH username for target nodes (default: root)') parser.add_argument('--port', type=int, default=22, @@ -316,6 +350,10 @@ def main(): cli_config['password'] = args.password if args.key_file: cli_config['key_file'] = args.key_file + for key in ('nqn', 'traddr', 'trsvcid', 'nsid'): + value = getattr(args, key) + if value is not None: + cli_config[key] = value register = MooncakeNoFRegister(cli_config, args.spdk_target_info) success = register.start_ssd_service()