From 1519ce60ade4edb59506447acf450ddadf2a05ce Mon Sep 17 00:00:00 2001 From: Teng Ma Date: Fri, 25 Apr 2025 17:35:16 +0800 Subject: [PATCH] [TransferEngine] implement closeLocalSegment by segment name (#286) --- mooncake-transfer-engine/include/transfer_engine.h | 2 ++ mooncake-transfer-engine/include/transfer_engine_c.h | 2 ++ mooncake-transfer-engine/include/transfer_metadata.h | 2 ++ mooncake-transfer-engine/src/transfer_engine.cpp | 9 +++++++++ mooncake-transfer-engine/src/transfer_engine_c.cpp | 5 +++++ mooncake-transfer-engine/src/transfer_metadata.cpp | 10 ++++++++++ .../tests/transfer_metadata_test.cpp | 4 ++++ 7 files changed, 34 insertions(+) diff --git a/mooncake-transfer-engine/include/transfer_engine.h b/mooncake-transfer-engine/include/transfer_engine.h index d9e475e5..4acf43be 100644 --- a/mooncake-transfer-engine/include/transfer_engine.h +++ b/mooncake-transfer-engine/include/transfer_engine.h @@ -89,6 +89,8 @@ class TransferEngine { int closeSegment(SegmentHandle handle); + int removeLocalSegment(const std::string &segment_name); + int registerLocalMemory(void *addr, size_t length, const std::string &location = kWildcardLocation, bool remote_accessible = true, diff --git a/mooncake-transfer-engine/include/transfer_engine_c.h b/mooncake-transfer-engine/include/transfer_engine_c.h index 5faa29da..2c39a23a 100644 --- a/mooncake-transfer-engine/include/transfer_engine_c.h +++ b/mooncake-transfer-engine/include/transfer_engine_c.h @@ -112,6 +112,8 @@ segment_id_t openSegmentNoCache(transfer_engine_t engine, const char *segment_na int closeSegment(transfer_engine_t engine, segment_id_t segment_id); +int removeLocalSegment(transfer_engine_t engine, const char *segment_name); + void destroyTransferEngine(transfer_engine_t engine); int registerLocalMemory(transfer_engine_t engine, void *addr, size_t length, diff --git a/mooncake-transfer-engine/include/transfer_metadata.h b/mooncake-transfer-engine/include/transfer_metadata.h index ba560224..5e880553 100644 --- a/mooncake-transfer-engine/include/transfer_metadata.h +++ b/mooncake-transfer-engine/include/transfer_metadata.h @@ -120,6 +120,8 @@ class TransferMetadata { int addLocalSegment(SegmentID segment_id, const std::string &segment_name, std::shared_ptr &&desc); + + int removeLocalSegment(const std::string &segment_name); int addRpcMetaEntry(const std::string &server_name, RpcMetaDesc &desc); diff --git a/mooncake-transfer-engine/src/transfer_engine.cpp b/mooncake-transfer-engine/src/transfer_engine.cpp index 431ef897..51bbaca0 100644 --- a/mooncake-transfer-engine/src/transfer_engine.cpp +++ b/mooncake-transfer-engine/src/transfer_engine.cpp @@ -199,6 +199,15 @@ Transport::SegmentHandle TransferEngine::openSegment( int TransferEngine::closeSegment(Transport::SegmentHandle handle) { return 0; } +int TransferEngine::removeLocalSegment(const std::string &segment_name) { + if (segment_name.empty()) return ERR_INVALID_ARGUMENT; + std::string trimmed_segment_name = segment_name; + while (!trimmed_segment_name.empty() && trimmed_segment_name[0] == '/') + trimmed_segment_name.erase(0, 1); + if (trimmed_segment_name.empty()) return ERR_INVALID_ARGUMENT; + return metadata_->removeLocalSegment(trimmed_segment_name); +} + bool TransferEngine::checkOverlap(void *addr, uint64_t length) { std::shared_lock lock(mutex_); for (auto &local_memory_region : local_memory_regions_) { diff --git a/mooncake-transfer-engine/src/transfer_engine_c.cpp b/mooncake-transfer-engine/src/transfer_engine_c.cpp index 559d409f..6ce637d3 100644 --- a/mooncake-transfer-engine/src/transfer_engine_c.cpp +++ b/mooncake-transfer-engine/src/transfer_engine_c.cpp @@ -77,6 +77,11 @@ int closeSegment(transfer_engine_t engine, segment_id_t segment_id) { return native->closeSegment(segment_id); } +int removeLocalSegment(transfer_engine_t engine, const char *segment_name) { + TransferEngine *native = (TransferEngine *)engine; + return native->removeLocalSegment(segment_name); +} + int registerLocalMemory(transfer_engine_t engine, void *addr, size_t length, const char *location, int remote_accessible) { TransferEngine *native = (TransferEngine *)engine; diff --git a/mooncake-transfer-engine/src/transfer_metadata.cpp b/mooncake-transfer-engine/src/transfer_metadata.cpp index a5343981..5637f63d 100644 --- a/mooncake-transfer-engine/src/transfer_metadata.cpp +++ b/mooncake-transfer-engine/src/transfer_metadata.cpp @@ -380,6 +380,16 @@ int TransferMetadata::addLocalSegment(SegmentID segment_id, return 0; } +int TransferMetadata::removeLocalSegment(const std::string &segment_name) { + RWSpinlock::WriteGuard guard(segment_lock_); + if (segment_name_to_id_map_.count(segment_name)) { + int segment_id = segment_name_to_id_map_[segment_name]; + segment_name_to_id_map_.erase(segment_name); + segment_id_to_desc_map_.erase(segment_id); + } + return 0; +} + int TransferMetadata::addLocalMemoryBuffer(const BufferDesc &buffer_desc, bool update_metadata) { { diff --git a/mooncake-transfer-engine/tests/transfer_metadata_test.cpp b/mooncake-transfer-engine/tests/transfer_metadata_test.cpp index c5025724..d8fd0d88 100644 --- a/mooncake-transfer-engine/tests/transfer_metadata_test.cpp +++ b/mooncake-transfer-engine/tests/transfer_metadata_test.cpp @@ -75,6 +75,8 @@ TEST_F(TransferMetadataTest, LocalSegmentTest) { ASSERT_EQ(des, segment_des); auto id = metadata_client->getSegmentID(segment_name); ASSERT_EQ(id, segment_id); + re = metadata_client->removeLocalSegment(segment_name); + ASSERT_EQ(re, 0); } // add and remove LocalMemoryBufferMeta @@ -101,6 +103,8 @@ TEST_F(TransferMetadataTest, LocalMemoryBufferTest) { re = metadata_client->removeLocalMemoryBuffer((void*)addr, false); ASSERT_EQ(re, 0); } + re = metadata_client->removeLocalSegment("test_local_segment"); + ASSERT_EQ(re, 0); } // add, get and remove RPCMetaEntryMeta