diff --git a/protocols/pain/proto/asura.proto b/protocols/pain/proto/asura.proto index 3425389..34b9c43 100644 --- a/protocols/pain/proto/asura.proto +++ b/protocols/pain/proto/asura.proto @@ -8,9 +8,6 @@ option cc_generic_services = true; service AsuraService { rpc RegisterDeva(RegisterDevaRequest) returns (RegisterDevaResponse); rpc ListDeva(ListDevaRequest) returns (ListDevaResponse); - rpc RegisterManusya(RegisterManusyaRequest) - returns (RegisterManusyaResponse); - rpc ListManusya(ListManusyaRequest) returns (ListManusyaResponse); } message DevaServer { @@ -40,32 +37,3 @@ message ListDevaResponse { Header header = 1; repeated DevaServer deva_servers = 2; } - -message ManusyaServer { - UUID id = 1; - string ip = 2; - uint32 port = 3; - UUID pool_id = 4; -} - -message RegisterManusyaRequest { - repeated ManusyaServer manusya_servers = 1; -} - -message RegisterManusyaResponse { - message Result { - UUID id = 1; - int32 code = 2; - string message = 3; - } - - Header header = 1; - repeated Result results = 2; -} - -message ListManusyaRequest {} - -message ListManusyaResponse { - Header header = 1; - repeated ManusyaServer manusya_servers = 2; -} diff --git a/protocols/pain/proto/common.proto b/protocols/pain/proto/common.proto index 7e4a56a..2650d00 100644 --- a/protocols/pain/proto/common.proto +++ b/protocols/pain/proto/common.proto @@ -85,3 +85,26 @@ message DirEntry { message DirEntries { repeated DirEntry entries = 1; } + +message ManusyaID { + string ip = 1; + int32 port = 2; + UUID uuid = 3; +} + +message StorageInfo { + string cluster_id = 1; + int64 ctime = 2; +} + +message ManusyaRegistration { + ManusyaID manusya_id = 1; + StorageInfo storage_info = 2; +} + +message ManusyaDescriptor { + ManusyaID manusya_id = 1; + StorageInfo storage_info = 2; + bool is_alive = 3; + uint64 last_heartbeat_time = 4; +} diff --git a/protocols/pain/proto/deva.proto b/protocols/pain/proto/deva.proto index a830ff2..929905b 100644 --- a/protocols/pain/proto/deva.proto +++ b/protocols/pain/proto/deva.proto @@ -18,6 +18,10 @@ service DevaService { rpc SealChunk(SealChunkRequest) returns (SealChunkResponse); rpc SealAndNewChunk(SealAndNewChunkRequest) returns (SealAndNewChunkResponse); + + rpc ManusyaHeartbeat(ManusyaHeartbeatRequest) + returns (ManusyaHeartbeatResponse); + rpc ListManusya(ListManusyaRequest) returns (ListManusyaResponse); } enum OpenFlag { @@ -114,3 +118,18 @@ message SealChunkRequest { message SealChunkResponse { Header header = 1; } + +message ManusyaHeartbeatRequest { + ManusyaRegistration manusya_registration = 1; +} + +message ManusyaHeartbeatResponse { + Header header = 1; +} + +message ListManusyaRequest {} + +message ListManusyaResponse { + Header header = 1; + repeated ManusyaDescriptor manusya_descriptors = 2; +} diff --git a/protocols/pain/proto/deva_store.proto b/protocols/pain/proto/deva_store.proto index f6fbd79..477e009 100644 --- a/protocols/pain/proto/deva_store.proto +++ b/protocols/pain/proto/deva_store.proto @@ -18,6 +18,14 @@ message CreateFileResponse { FileInfo file_info = 1; } +message GetFileInfoRequest { + string path = 1; +} + +message GetFileInfoResponse { + FileInfo file_info = 1; +} + message CreateDirRequest { string path = 1; UUID dir_id = 2; @@ -64,3 +72,15 @@ message SealChunkResponse {} message SealAndNewChunkRequest {} message SealAndNewChunkResponse {} + +message ManusyaHeartbeatRequest { + ManusyaRegistration manusya_registration = 1; +} + +message ManusyaHeartbeatResponse {} + +message ListManusyaRequest {} + +message ListManusyaResponse { + repeated ManusyaDescriptor manusya_descriptors = 1; +} diff --git a/src/asura/asura_service_impl.cc b/src/asura/asura_service_impl.cc index dfcf019..6066f0c 100644 --- a/src/asura/asura_service_impl.cc +++ b/src/asura/asura_service_impl.cc @@ -39,34 +39,6 @@ void AsuraServiceImpl::RegisterDeva(::google::protobuf::RpcController* controlle } } -void AsuraServiceImpl::RegisterManusya(::google::protobuf::RpcController* controller, - [[maybe_unused]] const pain::proto::asura::RegisterManusyaRequest* request, - [[maybe_unused]] pain::proto::asura::RegisterManusyaResponse* response, - ::google::protobuf::Closure* done) { // NOLINT(readability-non-const-parameter) - ASURA_SPAN(span, controller); - brpc::ClosureGuard done_guard(done); - for (const auto& manusya_server : request->manusya_servers()) { - auto uuid = UUID(manusya_server.id().high(), manusya_server.id().low()); - auto result = response->add_results(); - result->mutable_id()->CopyFrom(manusya_server.id()); - result->set_code(0); - result->set_message("ok"); - - if (!_store->hexists(ASURA_MANUSYA, uuid.str())) { - result->set_code(EEXIST); - result->set_message( - fmt::format("{} existed. ip:{}, port:{}", uuid.str(), manusya_server.ip(), manusya_server.port())); - continue; - } - - auto status = _store->hset(ASURA_MANUSYA, uuid.str(), manusya_server.SerializeAsString()); - if (!status.ok()) { - result->set_code(status.error_code()); - result->set_message(status.error_cstr()); - } - } -} - void AsuraServiceImpl::ListDeva(::google::protobuf::RpcController* controller, [[maybe_unused]] const pain::proto::asura::ListDevaRequest* request, [[maybe_unused]] pain::proto::asura::ListDevaResponse* response, @@ -82,19 +54,4 @@ void AsuraServiceImpl::ListDeva(::google::protobuf::RpcController* controller, } } -void AsuraServiceImpl::ListManusya(::google::protobuf::RpcController* controller, - [[maybe_unused]] const pain::proto::asura::ListManusyaRequest* request, - [[maybe_unused]] pain::proto::asura::ListManusyaResponse* response, - ::google::protobuf::Closure* done) { // NOLINT(readability-non-const-parameter) - ASURA_SPAN(span, controller); - brpc::ClosureGuard done_guard(done); - auto it = _store->hgetall(ASURA_MANUSYA); - while (it->valid()) { - auto manusya = response->add_manusya_servers(); - auto value = it->value(); - manusya->ParseFromArray(value.data(), value.size()); - it->next(); - } -} - } // namespace pain::asura diff --git a/src/asura/asura_service_impl.h b/src/asura/asura_service_impl.h index e7056f8..3af165d 100644 --- a/src/asura/asura_service_impl.h +++ b/src/asura/asura_service_impl.h @@ -10,9 +10,7 @@ class AsuraServiceImpl : public pain::proto::asura::AsuraService { public: AsuraServiceImpl(common::StorePtr store) : _store(store) {} ASURA_RPC_ENTRY(RegisterDeva); - ASURA_RPC_ENTRY(RegisterManusya); ASURA_RPC_ENTRY(ListDeva); - ASURA_RPC_ENTRY(ListManusya); private: common::StorePtr _store; diff --git a/src/deva/deva.cc b/src/deva/deva.cc index 34e70bb..41c93c8 100644 --- a/src/deva/deva.cc +++ b/src/deva/deva.cc @@ -60,7 +60,10 @@ DEVA_METHOD(CreateFile) { file_info->set_mode(request->mode()); file_info->set_uid(request->uid()); file_info->set_gid(request->gid()); - _file_infos[file_uuid] = *file_info; + status = update_file_info(file_uuid, *file_info); + if (!status.ok()) { + return status; + } return Status::OK(); } @@ -85,7 +88,10 @@ DEVA_METHOD(CreateDir) { file_info->set_mode(request->mode()); file_info->set_uid(request->uid()); file_info->set_gid(request->gid()); - _file_infos[dir_uuid] = *file_info; + status = update_file_info(dir_uuid, *file_info); + if (!status.ok()) { + return status; + } return Status::OK(); } @@ -150,6 +156,76 @@ DEVA_METHOD(SealAndNewChunk) { return Status::OK(); } +DEVA_METHOD(GetFileInfo) { + SPAN(span); + PLOG_INFO(("desc", "get_file_info")("version", version)("index", index)); + auto& path = request->path(); + UUID file_uuid; + FileType file_type = FileType::kNone; + auto status = _namespace.lookup(path.c_str(), &file_uuid, &file_type); + if (!status.ok()) { + return status; + } + if (file_type != FileType::kFile) { + return Status(EINVAL, fmt::format("{} is not a file", path.c_str())); + } + proto::FileInfo file_info; + status = get_file_info(file_uuid, &file_info); + if (!status.ok()) { + return status; + } + response->mutable_file_info()->Swap(&file_info); + return Status::OK(); +} + +DEVA_METHOD(ManusyaHeartbeat) { + SPAN(span); + PLOG_INFO(("desc", "manusya_heartbeat")("version", version)("index", index)); + auto& manusya_id = request->manusya_registration().manusya_id(); + // auto& storage_info = request->manusya_registration().storage_info(); + // TODO: check cluster id + + UUID manusya_uuid(manusya_id.uuid().high(), manusya_id.uuid().low()); + auto it = _manusya_descriptors.find(manusya_uuid); + if (it == _manusya_descriptors.end()) { + // new manusya + PLOG_INFO(("desc", "new manusya") // + ("manusya_uuid", manusya_uuid.str()) // + ("ip", manusya_id.ip()) // + ("port", manusya_id.port())); + ManusyaDescriptor manusya_descriptor; + manusya_descriptor.ip = manusya_id.ip(); + manusya_descriptor.port = manusya_id.port(); + manusya_descriptor.uuid = manusya_uuid; + manusya_descriptor.is_alive = true; + manusya_descriptor.update_heartbeat(); + _manusya_descriptors[manusya_uuid] = manusya_descriptor; + } else { + auto& manusya_descriptor = it->second; + manusya_descriptor.ip = manusya_id.ip(); + manusya_descriptor.port = manusya_id.port(); + manusya_descriptor.update_heartbeat(); + } + return Status::OK(); +} + +DEVA_METHOD(ListManusya) { + SPAN(span); + PLOG_INFO(("desc", "list_manusya")("version", version)("index", index)); + auto manusya_descriptors = response->mutable_manusya_descriptors(); + for (auto& [_, manusya_descriptor] : _manusya_descriptors) { + auto manusya_descriptor_proto = manusya_descriptors->Add(); + manusya_descriptor_proto->mutable_manusya_id()->set_ip(manusya_descriptor.ip); + manusya_descriptor_proto->mutable_manusya_id()->set_port(manusya_descriptor.port); + manusya_descriptor_proto->mutable_manusya_id()->mutable_uuid()->set_high(manusya_descriptor.uuid.high()); + manusya_descriptor_proto->mutable_manusya_id()->mutable_uuid()->set_low(manusya_descriptor.uuid.low()); + manusya_descriptor_proto->mutable_storage_info()->set_cluster_id(manusya_descriptor.cluster_id); + manusya_descriptor_proto->set_is_alive(manusya_descriptor.is_alive); + manusya_descriptor_proto->set_last_heartbeat_time(manusya_descriptor.last_heartbeat_time); + } + return Status::OK(); +} + Status Deva::set_applied_index(int64_t index) { if (index <= _applied_index) { PLOG_WARN(("desc", "index is applied already")("index", index)); @@ -171,6 +247,48 @@ Status Deva::set_applied_index(int64_t index) { return status; } +Status Deva::update_file_info(const UUID& id, const proto::FileInfo& file_info) { + auto in_txn = common::TxnManager::instance().in_txn(); + auto this_txn = _store->begin_txn(); + auto txn = in_txn ? common::TxnManager::instance().get_txn_store() : this_txn.get(); + if (txn == nullptr) { + return Status(EIO, "Failed to begin transaction"); + } + auto status = txn->hset(_file_info_key, id.str(), file_info.SerializeAsString()); + if (!status.ok()) { + return status; + } + if (in_txn) { + return Status::OK(); + } + status = txn->commit(); + if (!status.ok()) { + return status; + } + return status; +} + +Status Deva::get_file_info(const UUID& id, proto::FileInfo* file_info) { + std::string file_info_str; + auto status = _store->hget(_file_info_key, id.str(), &file_info_str); + if (!status.ok()) { + return status; + } + auto ret = file_info->ParseFromString(file_info_str); + if (!ret) { + return Status(EIO, "Failed to parse file info"); + } + return Status::OK(); +} + +Status Deva::remove_file_info(const UUID& id) { + auto status = _store->hdel(_file_info_key, id.str()); + if (!status.ok()) { + return status; + } + return Status::OK(); +} + Status Deva::save_snapshot(std::string_view path, std::vector* files) { return _store->check_point(path.data(), files); } diff --git a/src/deva/deva.h b/src/deva/deva.h index 26d3de9..92b2e8d 100644 --- a/src/deva/deva.h +++ b/src/deva/deva.h @@ -8,6 +8,7 @@ #include "common/txn_manager.h" #include "common/txn_store.h" #include "deva/container.h" +#include "deva/manusya_descriptor.h" #include "deva/namespace.h" #include "deva/op.h" @@ -58,6 +59,9 @@ class Deva : public Container { DEVA_ENTRY(CheckInChunk); DEVA_ENTRY(SealChunk); DEVA_ENTRY(SealAndNewChunk); + DEVA_ENTRY(GetFileInfo); + DEVA_ENTRY(ManusyaHeartbeat); + DEVA_ENTRY(ListManusya); Status save_snapshot(std::string_view path, std::vector* files) override; Status load_snapshot(std::string_view path) override; @@ -69,15 +73,22 @@ class Deva : public Container { return index != 0 && index <= _applied_index; } + Status update_file_info(const UUID& id, const proto::FileInfo& file_info); + Status get_file_info(const UUID& id, proto::FileInfo* file_info); + Status remove_file_info(const UUID& id); + private: std::atomic _use_count; common::StorePtr _store; Namespace _namespace; - std::unordered_map _file_infos; + const char* _file_info_key = "file_info"; const char* _meta_key = "meta"; const char* _applied_index_key = "applied_index"; int64_t _applied_index = 0; + // don't need persist + std::unordered_map _manusya_descriptors; + friend void intrusive_ptr_add_ref(Deva* deva) { ++deva->_use_count; } diff --git a/src/deva/deva_service_impl.cc b/src/deva/deva_service_impl.cc index 49aac81..23ae655 100644 --- a/src/deva/deva_service_impl.cc +++ b/src/deva/deva_service_impl.cc @@ -48,7 +48,18 @@ DEVA_SERVICE_METHOD(OpenFile) { response->mutable_header()->set_status(0); response->mutable_header()->set_message("ok"); } else { - // TODO: open file + // readdir + // get file info + pain::proto::deva::store::GetFileInfoRequest get_file_info_request; + pain::proto::deva::store::GetFileInfoResponse get_file_info_response; + get_file_info_request.set_path(path); + auto status = bridge(1, _rsm, get_file_info_request, &get_file_info_response).get(); + if (!status.ok()) { + PLOG_ERROR(("desc", "failed to get file info")("error", status.error_str())); + } + response->mutable_file_info()->Swap(get_file_info_response.mutable_file_info()); + response->mutable_header()->set_status(0); + response->mutable_header()->set_message("ok"); } response->mutable_header()->set_status(0); @@ -137,6 +148,44 @@ DEVA_SERVICE_METHOD(SealAndNewChunk) { DEFINE_SPAN(span, controller); } +DEVA_SERVICE_METHOD(ManusyaHeartbeat) { + brpc::ClosureGuard done_guard(done); + DEFINE_SPAN(span, controller); + auto& manusya_registration = request->manusya_registration(); + pain::proto::deva::store::ManusyaHeartbeatRequest manusya_heartbeat_request; + pain::proto::deva::store::ManusyaHeartbeatResponse manusya_heartbeat_response; + + manusya_heartbeat_request.mutable_manusya_registration()->CopyFrom(manusya_registration); + + auto status = + bridge(1, _rsm, manusya_heartbeat_request, &manusya_heartbeat_response).get(); + if (!status.ok()) { + PLOG_ERROR(("desc", "failed to handle manusya heartbeat")("error", status.error_str())); + response->mutable_header()->set_status(status.error_code()); + response->mutable_header()->set_message(status.error_str()); + return; + } + response->mutable_header()->set_status(0); + response->mutable_header()->set_message("ok"); +} + +DEVA_SERVICE_METHOD(ListManusya) { + brpc::ClosureGuard done_guard(done); + DEFINE_SPAN(span, controller); + pain::proto::deva::store::ListManusyaRequest list_manusya_request; + pain::proto::deva::store::ListManusyaResponse list_manusya_response; + auto status = bridge(1, _rsm, list_manusya_request, &list_manusya_response).get(); + if (!status.ok()) { + PLOG_ERROR(("desc", "failed to list manusya")("error", status.error_str())); + response->mutable_header()->set_status(status.error_code()); + response->mutable_header()->set_message(status.error_str()); + return; + } + response->mutable_manusya_descriptors()->Swap(list_manusya_response.mutable_manusya_descriptors()); + response->mutable_header()->set_status(0); + response->mutable_header()->set_message("ok"); +} + } // namespace pain::deva #undef DEVA_SERVICE_METHOD diff --git a/src/deva/deva_service_impl.h b/src/deva/deva_service_impl.h index 1fa9cef..6d29819 100644 --- a/src/deva/deva_service_impl.h +++ b/src/deva/deva_service_impl.h @@ -25,6 +25,8 @@ class DevaServiceImpl : public pain::proto::deva::DevaService { DEVA_SERVICE_METHOD(CheckInChunk); DEVA_SERVICE_METHOD(SealChunk); DEVA_SERVICE_METHOD(SealAndNewChunk); + DEVA_SERVICE_METHOD(ManusyaHeartbeat); + DEVA_SERVICE_METHOD(ListManusya); private: RsmPtr _rsm; diff --git a/src/deva/manusya_descriptor.h b/src/deva/manusya_descriptor.h new file mode 100644 index 0000000..9124436 --- /dev/null +++ b/src/deva/manusya_descriptor.h @@ -0,0 +1,41 @@ +#pragma once + +#include + +namespace pain::deva { + +struct ManusyaDescriptor { + enum class AdminState { + kNormal = 0, + kDecommissioned = 1, + kMaintenance = 2, + }; + + std::string ip; + int32_t port; + UUID uuid; + std::string cluster_id; + std::string network_location; // rack info, like "/rack1" or "/dc1/az1/rack1" + + uint64_t total_capacity; // total capacity (bytes) + uint64_t used_space; // used space (bytes) + uint64_t remaining_space; // remaining space (bytes) + + time_t last_heartbeat_time; // last heartbeat time + bool is_alive; // is alive + + std::vector chunks; + AdminState admin_state; // admin state (normal, decommissioned, maintenance, etc.) + + void update_heartbeat() { + last_heartbeat_time = time(nullptr); + is_alive = true; + } + + // 示例方法:标记为不活跃 + void mark_dead() { + is_alive = false; + } +}; + +} // namespace pain::deva diff --git a/src/deva/op.h b/src/deva/op.h index 05f1c6f..0eb652d 100644 --- a/src/deva/op.h +++ b/src/deva/op.h @@ -22,6 +22,9 @@ enum class OpType : uint32_t { DEFINE_DEVA_OP(7, SealChunk, true), DEFINE_DEVA_OP(8, SealAndNewChunk, true), DEFINE_DEVA_OP(9, ReadDir, false), + DEFINE_DEVA_OP(10, GetFileInfo, false), + DEFINE_DEVA_OP(20, ManusyaHeartbeat, false), + DEFINE_DEVA_OP(21, ListManusya, false), DEFINE_DEVA_OP(100, MaxDevaOp, true), }; diff --git a/src/deva/test/test_deva.cc b/src/deva/test/test_deva.cc index a381565..993809b 100644 --- a/src/deva/test/test_deva.cc +++ b/src/deva/test/test_deva.cc @@ -15,9 +15,8 @@ class TestDeva : public ::testing::Test { _mock_deva.group().c_str(), &pain::proto::deva::DevaService::Mkdir, &request, response); } - pain::Status open(const std::string& path, - const pain::proto::deva::OpenFlag& flags, - pain::proto::deva::OpenFileResponse* response) { + pain::Status + open(const std::string& path, pain::proto::deva::OpenFlag flags, pain::proto::deva::OpenFileResponse* response) { pain::proto::deva::OpenFileRequest request; request.set_path(path); request.set_flags(flags); @@ -32,6 +31,26 @@ class TestDeva : public ::testing::Test { _mock_deva.group().c_str(), &pain::proto::deva::DevaService::ReadDir, &request, response); } + pain::Status manusya_heartbeat(const pain::UUID& uuid, + const char* ip, + int32_t port, + pain::proto::deva::ManusyaHeartbeatResponse* response) { + pain::proto::deva::ManusyaHeartbeatRequest request; + auto manusya_id = request.mutable_manusya_registration()->mutable_manusya_id(); + manusya_id->mutable_uuid()->set_high(uuid.high()); + manusya_id->mutable_uuid()->set_low(uuid.low()); + manusya_id->set_ip(ip); + manusya_id->set_port(port); + return pain::deva::call_rpc( + _mock_deva.group().c_str(), &pain::proto::deva::DevaService::ManusyaHeartbeat, &request, response); + } + + pain::Status list_manusya(pain::proto::deva::ListManusyaResponse* response) { + pain::proto::deva::ListManusyaRequest request; + return pain::deva::call_rpc( + _mock_deva.group().c_str(), &pain::proto::deva::DevaService::ListManusya, &request, response); + } + void TearDown() override { if (::testing::Test::HasFailure()) { _mock_deva.do_not_remove_data_path(); @@ -104,15 +123,30 @@ TEST_F(TestDeva, CreateDirectoryAndFile) { EXPECT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; std::cout << "leader: " << leader << std::endl; - pain::proto::deva::MkdirResponse mkdir_response; - status = mkdir("/test", &mkdir_response); - ASSERT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; - std::cout << "mkdir response: " << mkdir_response.DebugString() << std::endl; + { + pain::proto::deva::MkdirResponse mkdir_response; + status = mkdir("/test", &mkdir_response); + ASSERT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; + std::cout << "mkdir response: " << mkdir_response.DebugString() << std::endl; + } - pain::proto::deva::OpenFileResponse response; - status = open("/test/test.txt", pain::proto::deva::OpenFlag::OPEN_CREATE, &response); - EXPECT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; - std::cout << "response: " << response.DebugString() << std::endl; + pain::proto::FileInfo file_info; + { + pain::proto::deva::OpenFileResponse response; + status = open("/test/test.txt", pain::proto::deva::OpenFlag::OPEN_CREATE, &response); + EXPECT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; + std::cout << "response: " << response.DebugString() << std::endl; + file_info.Swap(response.mutable_file_info()); + } + + { + pain::proto::deva::OpenFileResponse response; + status = open("/test/test.txt", pain::proto::deva::OpenFlag::OPEN_READ, &response); + EXPECT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; + std::cout << "response: " << response.DebugString() << std::endl; + EXPECT_EQ(response.file_info().file_id().high(), file_info.file_id().high()); + EXPECT_EQ(response.file_info().file_id().low(), file_info.file_id().low()); + } } TEST_F(TestDeva, ReadDir) { @@ -202,4 +236,45 @@ TEST_F(TestDeva, Snapshot) { EXPECT_EQ(readdir_response.entries(0).type(), pain::proto::FileType::FILE_TYPE_FILE); } +TEST_F(TestDeva, ListManusya) { + _mock_deva.start(); + SCOPE_EXIT { + _mock_deva.stop(); + }; + std::string leader; + auto status = _mock_deva.wait_for_leader(&leader); + ASSERT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; + std::cout << "leader: " << leader << std::endl; + + { + pain::proto::deva::ListManusyaResponse response; + status = list_manusya(&response); + ASSERT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; + std::cout << "list_manusya response: " << response.DebugString() << std::endl; + } + + { + // heartbeat + pain::proto::deva::ManusyaHeartbeatResponse response; + pain::UUID uuid(1, 2); + status = manusya_heartbeat(uuid, "127.0.0.1", 12345, &response); // NOLINT(readability-magic-numbers) + ASSERT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; + std::cout << "manusya_heartbeat response: " << response.DebugString() << std::endl; + } + + { + pain::proto::deva::ListManusyaResponse response; + status = list_manusya(&response); + ASSERT_TRUE(status.ok()) << status.error_str() << "(" << status.error_code() << ")"; + std::cout << "list_manusya response: " << response.DebugString() << std::endl; + pain::UUID uuid(1, 2); + ASSERT_EQ(response.manusya_descriptors_size(), 1); + EXPECT_EQ(response.manusya_descriptors(0).manusya_id().uuid().high(), uuid.high()); + EXPECT_EQ(response.manusya_descriptors(0).manusya_id().uuid().low(), uuid.low()); + EXPECT_EQ(response.manusya_descriptors(0).manusya_id().ip(), "127.0.0.1"); + EXPECT_EQ(response.manusya_descriptors(0).manusya_id().port(), 12345); + EXPECT_TRUE(response.manusya_descriptors(0).is_alive()); + } +} + } // namespace diff --git a/src/sad/asura.cc b/src/sad/asura.cc index d793a8a..9a4c504 100644 --- a/src/sad/asura.cc +++ b/src/sad/asura.cc @@ -125,85 +125,4 @@ COMMAND(list_deva) { return Status::OK(); } -REGISTER_ASURA_CMD(register_manusya, [](argparse::ArgumentParser& parser) { - parser.add_description("add manusya"); - parser.add_argument("--ip").required(); - // NOLINTNEXTLINE - parser.add_argument("--port").default_value(0U).scan<'i', uint32_t>().required(); - parser.add_argument("--pool").required(); -}); -COMMAND(register_manusya) { - SPAN(span); - auto host = args.get("--host"); - auto ip = args.get("--ip"); - auto port = args.get("--port"); - auto pool = args.get("--pool"); - if (!UUID::valid(pool)) { - return Status(EINVAL, "Invalid pool id"); - } - auto pool_id = UUID::from_str_or_die(pool); - - brpc::Channel channel; - brpc::ChannelOptions options; - options.connect_timeout_ms = 2000; // NOLINT(readability-magic-numbers) - options.timeout_ms = 10000; // NOLINT(readability-magic-numbers) - options.max_retry = 0; // NOLINT(readability-magic-numbers) - if (channel.Init(host.c_str(), &options) != 0) { - return Status(EAGAIN, "Fail to initialize channel"); - } - - brpc::Controller cntl; - pain::proto::asura::RegisterManusyaRequest request; - pain::proto::asura::RegisterManusyaResponse response; - pain::proto::asura::AsuraService::Stub stub(&channel); - pain::inject_tracer(&cntl); - - auto id = pain::UUID::generate(); - auto manusya_server = request.add_manusya_servers(); - manusya_server->mutable_id()->set_low(id.low()); - manusya_server->mutable_id()->set_high(id.high()); - manusya_server->set_ip(ip); - manusya_server->set_port(port); - manusya_server->mutable_pool_id()->set_low(pool_id.low()); - manusya_server->mutable_pool_id()->set_high(pool_id.high()); - stub.RegisterManusya(&cntl, &request, &response, nullptr); - if (cntl.Failed()) { - return Status(cntl.ErrorCode(), cntl.ErrorText()); - } - - print(cntl, &response); - return Status::OK(); -} - -REGISTER_ASURA_CMD(list_manusya, [](argparse::ArgumentParser& parser) { - parser.add_description("list manusya"); -}); -COMMAND(list_manusya) { - SPAN(span); - auto host = args.get("--host"); - - brpc::Channel channel; - brpc::ChannelOptions options; - options.connect_timeout_ms = 2000; // NOLINT(readability-magic-numbers) - options.timeout_ms = 10000; // NOLINT(readability-magic-numbers) - options.max_retry = 0; // NOLINT(readability-magic-numbers) - if (channel.Init(host.c_str(), &options) != 0) { - return Status(EAGAIN, "Fail to initialize channel"); - } - - brpc::Controller cntl; - pain::proto::asura::ListManusyaRequest request; - pain::proto::asura::ListManusyaResponse response; - pain::proto::asura::AsuraService::Stub stub(&channel); - pain::inject_tracer(&cntl); - - stub.ListManusya(&cntl, &request, &response, nullptr); - if (cntl.Failed()) { - return Status(cntl.ErrorCode(), cntl.ErrorText()); - } - - print(cntl, &response); - return Status::OK(); -} - } // namespace pain::sad::asura