diff --git a/src/cluster_manager/CMakeLists.txt b/src/cluster_manager/CMakeLists.txt index 5090e49..6a729fc 100644 --- a/src/cluster_manager/CMakeLists.txt +++ b/src/cluster_manager/CMakeLists.txt @@ -16,7 +16,11 @@ list(REMOVE_ITEM MODULE_SRC_LIB "${CMAKE_CURRENT_SOURCE_DIR}/cm_main.cc" ) -add_executable(${MODULE_NAME} ${MODULE_SRC} ${MODULE_HEADERS}) +set(DS_SHARED_FLAG_SRC + "${PROJECT_SOURCE_DIR}/src/data_server/ds_flags.cc" +) + +add_executable(${MODULE_NAME} ${MODULE_SRC} ${MODULE_HEADERS} ${DS_SHARED_FLAG_SRC}) target_link_libraries(${MODULE_NAME} PRIVATE gflags folly @@ -32,7 +36,7 @@ set_target_properties(${MODULE_NAME} PROPERTIES INSTALL_RPATH "\$ORIGIN/../../../third_party/sict/lib" ) -add_library(${MODULE_NAME}_static STATIC ${MODULE_SRC_LIB}) +add_library(${MODULE_NAME}_static STATIC ${MODULE_SRC_LIB} ${DS_SHARED_FLAG_SRC}) # PRIVATE means that the libraries are only used by this target and # not propagated to targets that link against this one. # @@ -46,3 +50,10 @@ target_link_libraries(${MODULE_NAME}_static PRIVATE sict simm_common ) + +if(ENABLE_TESTS) + target_compile_definitions(${MODULE_NAME}_static PRIVATE SIMM_UNIT_TEST) + target_include_directories(${MODULE_NAME}_static PRIVATE + ${PROJECT_SOURCE_DIR}/third_party/gtest/googletest/include + ) +endif() diff --git a/src/cluster_manager/cm_node_manager.cc b/src/cluster_manager/cm_node_manager.cc index 4c9ada9..ced2568 100644 --- a/src/cluster_manager/cm_node_manager.cc +++ b/src/cluster_manager/cm_node_manager.cc @@ -1,3 +1,4 @@ +#include #include #include #include @@ -12,6 +13,7 @@ #include "rpc/connection.h" #include "rpc/rpc.h" +#include "cm_rpc_handler.h" #include "common/logging/logging.h" #include "common/utils/time_util.h" #include "cm_node_manager.h" @@ -19,6 +21,7 @@ DECLARE_string(dataserver_namespace); DECLARE_string(dataserver_svc_name); DECLARE_string(dataserver_port_name); +DECLARE_int32(mgt_service_port); DECLARE_uint32(rpc_timeout_inSecs); DECLARE_uint32(dataserver_resource_interval_inSecs); DECLARE_bool(cm_deferred_reshard_enabled); @@ -37,6 +40,7 @@ ClusterManagerNodeManager::ClusterManagerNodeManager() { ClusterManagerNodeManager::~ClusterManagerNodeManager() { MLOG_INFO("Start destruct Node Manager"); + Stop(); if (rpc_client_ != nullptr) { delete rpc_client_; rpc_client_ = nullptr; @@ -44,6 +48,9 @@ ClusterManagerNodeManager::~ClusterManagerNodeManager() { } void ClusterManagerNodeManager::Init() { + if (resource_thread_ != nullptr) { + return; + } start_timestamp_us_ = simm::utils::current_microseconds(); resource_thread_stop_.store(false); @@ -52,7 +59,8 @@ void ClusterManagerNodeManager::Init() { std::function resource_loop = [self]() { while (!self->resource_thread_stop_.load()) { self->updateAllNodeResource(); - self->resource_thread_baton_.timed_wait(std::chrono::milliseconds(FLAGS_dataserver_resource_interval_inSecs * 1000)); + self->resource_thread_baton_.timed_wait( + std::chrono::milliseconds(FLAGS_dataserver_resource_interval_inSecs * 1000)); self->resource_thread_baton_.reset(); } }; @@ -67,6 +75,12 @@ void ClusterManagerNodeManager::Stop() { if (resource_thread_ != nullptr && resource_thread_->joinable()) { resource_thread_->join(); } + { + std::lock_guard lock(resource_conn_mutex_); + resource_conn_map_.clear(); + } + delete resource_thread_; + resource_thread_ = nullptr; MLOG_INFO("Delete resource query thread in Node Manager succeed"); } @@ -88,11 +102,7 @@ std::vector> ClusterManagerNodeManage } error_code_t ClusterManagerNodeManager::AddNode(const std::string & addr_str) { - auto result = node_status_map_.insert_or_assign(addr_str, NodeStatus::RUNNING); - if (!result.second) { - MLOG_ERROR("Add node {} in Node Manager status map failed", addr_str); - return CmErr::NodeManagerAddNodeFailed; - } + node_status_map_.insert_or_assign(addr_str, NodeStatus::RUNNING); // FIXME(ytji): for test, just comment below codes // // get node resource info, only print log if error @@ -110,6 +120,10 @@ error_code_t ClusterManagerNodeManager::AddNode(const std::string & addr_str) { error_code_t ClusterManagerNodeManager::DelNode(const std::string & addr_str) { node_info_map_.erase(addr_str); node_status_map_.erase(addr_str); + { + std::lock_guard lock(resource_conn_mutex_); + resource_conn_map_.erase(addr_str); + } MLOG_DEBUG("Delete node {} in Node Manager succeed", addr_str); return CommonErr::OK; } @@ -122,8 +136,13 @@ error_code_t ClusterManagerNodeManager::UpdateNodeStatus(const std::string & add std::shared_ptr ClusterManagerNodeManager::getNodeResource( const std::string & addr_str) { +#if defined(SIMM_UNIT_TEST) + if (test_resource_query_hook_) { + return test_resource_query_hook_(addr_str); + } +#endif DataServerResourceRequestPB req; - auto resp = new DataServerResourceResponsePB; + auto resp = std::make_unique(); sicl::rpc::RpcContext *ctx_p; sicl::rpc::RpcContext::newInstance(ctx_p); std::shared_ptr ctx = std::shared_ptr(ctx_p); @@ -134,16 +153,54 @@ std::shared_ptr ClusterManagerNodeManager::getNodeRe return nullptr; } - rpc_client_->SendRequest(ds_addr->node_ip_, ds_addr->node_port_, static_cast(0), req, resp, ctx, nullptr); + std::shared_ptr conn; + { + std::lock_guard lock(resource_conn_mutex_); + auto it = resource_conn_map_.find(addr_str); + if (it != resource_conn_map_.end()) { + conn = it->second; + } + } + if (!conn) { + conn = rpc_client_->connect(ds_addr->node_ip_, FLAGS_mgt_service_port); + if (!conn) { + MLOG_ERROR("Get node {} resource failed, connect failed to management port {}", addr_str, FLAGS_mgt_service_port); + return nullptr; + } + std::lock_guard lock(resource_conn_mutex_); + resource_conn_map_[addr_str] = conn; + } + + ctx->set_timeout(sicl::transport::TimerTick::TIMER_3S); + rpc_client_->SendRequest(conn, + static_cast(cm::ClusterManagerRpcType::RPC_DATASERVER_RESOURCE_QUERY), + req, + resp.get(), + ctx); if (ctx->Failed()) { std::string errmsg = ctx->ErrorText(); + std::lock_guard lock(resource_conn_mutex_); + resource_conn_map_.erase(addr_str); MLOG_ERROR("Get node {} resource failed, err:{}", addr_str, errmsg); return nullptr; } + std::vector shard_mem_infos; + shard_mem_infos.reserve(resp->shard_mem_infos_size()); + for (const auto &shard_pb : resp->shard_mem_infos()) { + shard_mem_infos.push_back( + {static_cast(shard_pb.shard_id()), static_cast(shard_pb.shard_mem_used_bytes())}); + } + std::sort(shard_mem_infos.begin(), shard_mem_infos.end(), [](const auto &lhs, const auto &rhs) { + return lhs.shard_id_ < rhs.shard_id_; + }); + auto resource = std::make_shared(static_cast(resp->mem_total_bytes()), + static_cast(resp->mem_allocated_bytes()), + static_cast(resp->mem_used_bytes()), static_cast(resp->mem_free_bytes()), - static_cast(resp->mem_used_bytes())); + resp->last_report_timestamp_us(), + std::move(shard_mem_infos)); MLOG_DEBUG("Get node {} resource info in Node Manager succeed", addr_str); return resource; @@ -191,6 +248,15 @@ std::shared_ptr ClusterManagerNodeManager::GetNodeRe return it->second; } +std::shared_ptr ClusterManagerNodeManager::RefreshNodeResource( + const std::string &addr_str) { + auto resource_ret = getNodeResource(addr_str); + if (resource_ret != nullptr) { + node_info_map_.insert_or_assign(addr_str, resource_ret); + } + return resource_ret; +} + std::unordered_map ClusterManagerNodeManager::GetAllNodeStatus() { std::unordered_map status_map; for (auto &pair : node_status_map_) { @@ -225,7 +291,6 @@ HandshakeResult ClusterManagerNodeManager::ProcessHandshake( const std::string& logical_id, const std::string& new_ip_port, const std::vector& reported_shards) { - HandshakeResult result; auto it = logical_node_table_.find(logical_id); @@ -259,8 +324,7 @@ HandshakeResult ClusterManagerNodeManager::ProcessHandshake( entry.status = NodeStatus::RUNNING; entry.deferred_reshard_since = {}; logical_node_table_.assign(logical_id, entry); - MLOG_INFO("Node replacement: logical_id={} old_ip={} new_ip={}", - logical_id, result.old_ip_port, new_ip_port); + MLOG_INFO("Node replacement: logical_id={} old_ip={} new_ip={}", logical_id, result.old_ip_port, new_ip_port); return result; } @@ -276,8 +340,7 @@ HandshakeResult ClusterManagerNodeManager::ProcessHandshake( // Case 4: IP changed while still RUNNING (rare IP drift) migrateNodeIp(logical_id, entry, new_ip_port); logical_node_table_.assign(logical_id, entry); - MLOG_INFO("Node IP update: logical_id={} old_ip={} new_ip={}", - logical_id, result.old_ip_port, new_ip_port); + MLOG_INFO("Node IP update: logical_id={} old_ip={} new_ip={}", logical_id, result.old_ip_port, new_ip_port); } return result; } @@ -323,7 +386,6 @@ error_code_t ClusterManagerNodeManager::SetNodeStatus( NodeStatus status, std::chrono::steady_clock::time_point ts, std::optional expected_status) { - auto it = logical_node_table_.find(logical_id); if (it == logical_node_table_.end()) { MLOG_WARN("SetNodeStatus: logical_id={} not found", logical_id); diff --git a/src/cluster_manager/cm_node_manager.h b/src/cluster_manager/cm_node_manager.h index e38321f..7fdaace 100644 --- a/src/cluster_manager/cm_node_manager.h +++ b/src/cluster_manager/cm_node_manager.h @@ -6,6 +6,7 @@ #include #include #include +#include #include #include @@ -88,6 +89,9 @@ class ClusterManagerNodeManager : public std::enable_shared_from_this GetNodeResource(const std::string &addr_str); + // actively query one data node resource info and refresh cache + std::shared_ptr RefreshNodeResource(const std::string &addr_str); + // get status of all data nodes // Returns an unordered_map of address string to node status std::unordered_map GetAllNodeStatus(); @@ -149,13 +153,20 @@ class ClusterManagerNodeManager : public std::enable_shared_from_this addr_to_logical_; sicl::rpc::SiRPC *rpc_client_{nullptr}; + std::mutex resource_conn_mutex_; + std::unordered_map> resource_conn_map_; std::thread *resource_thread_{nullptr}; std::atomic resource_thread_stop_{false}; folly::Baton<> resource_thread_baton_; uint64_t start_timestamp_us_{0}; +#if defined(SIMM_UNIT_TEST) + std::function(const std::string &)> test_resource_query_hook_; +#endif + // only for UT test #if defined(SIMM_UNIT_TEST) + friend class ClusterManagerNodeManagerTestPeer; FRIEND_TEST(ClusterManagerHBMonitorTest, TestHBMonitor); #endif }; diff --git a/src/cluster_manager/cm_rpc_handler.cc b/src/cluster_manager/cm_rpc_handler.cc index 950b5d5..33124b1 100644 --- a/src/cluster_manager/cm_rpc_handler.cc +++ b/src/cluster_manager/cm_rpc_handler.cc @@ -35,6 +35,30 @@ DECLARE_LOG_MODULE("cluster_manager"); namespace simm { namespace cm { +static inline void FillNodeInfoHelper(const std::string &addr_str, + NodeStatus status, + const std::shared_ptr &resource, + NodeInfoPB *node_info) { + auto node_addr = simm::common::NodeAddress::ParseFromString(addr_str); + if (!node_addr) { + return; + } + + auto *addr_pb = node_info->mutable_node_address(); + addr_pb->set_ip(node_addr->node_ip_); + addr_pb->set_port(node_addr->node_port_); + node_info->set_node_status(static_cast(status)); + + if (resource) { + auto *res_pb = node_info->mutable_resource(); + res_pb->set_mem_free_bytes(resource->mem_free_bytes_); + res_pb->set_mem_total_bytes(resource->mem_total_bytes_); + res_pb->set_mem_used_bytes(resource->mem_used_bytes_); + res_pb->set_mem_allocated_bytes(resource->mem_allocated_bytes_); + res_pb->set_last_report_timestamp_us(resource->last_report_timestamp_us_); + } +} + template static inline void FillRespWithRoutingTableEntryHelper(const simm::cm::QueryResultMap &resmap, T *respb) { if (resmap.empty()) { @@ -79,7 +103,7 @@ void NewNodeHandshakeHandler::Work(const std::shared_ptr auto resp = std::make_shared(); simm::common::NodeAddress node_addr = {req->node().ip(), req->node().port()}; std::string addr_str = node_addr.toString(); - const std::string& logical_id = req->logical_node_id(); + const std::string &logical_id = req->logical_node_id(); error_code_t ret = CommonErr::OK; if (!logical_id.empty() && simm::common::ModuleServiceState::GetInstance().GracePeriodFinished()) { @@ -91,23 +115,25 @@ void NewNodeHandshakeHandler::Work(const std::shared_ptr case HandshakeResult::Action::DEFERRED_RESHARD_REPLACE: { // Replacement scenario: collect shards from old IP, assign to new IP auto shards = shard_manager_->GetShardsOwnedByNode(result.old_ip_port); - shard_manager_->BatchAssignRoutingTable( - shards, std::make_shared(node_addr)); + shard_manager_->BatchAssignRoutingTable(shards, std::make_shared(node_addr)); hb_monitor_->OnDeferredReshardResolved(logical_id); result.shards_to_assign = std::move(shards); MLOG_INFO("Node handshake replace complete: {} -> {} for {} shards", - result.old_ip_port, addr_str, result.shards_to_assign.size()); + result.old_ip_port, + addr_str, + result.shards_to_assign.size()); break; } case HandshakeResult::Action::IP_UPDATE: { // IP changed while RUNNING: update routing table entries to new IP auto shards = shard_manager_->GetShardsOwnedByNode(result.old_ip_port); - shard_manager_->BatchAssignRoutingTable( - shards, std::make_shared(node_addr)); + shard_manager_->BatchAssignRoutingTable(shards, std::make_shared(node_addr)); hb_monitor_->OnDeferredReshardResolved(logical_id); result.shards_to_assign = std::move(shards); MLOG_INFO("Node handshake IP update: {} -> {} for {} shards", - result.old_ip_port, addr_str, result.shards_to_assign.size()); + result.old_ip_port, + addr_str, + result.shards_to_assign.size()); break; } case HandshakeResult::Action::NEW_NODE: { @@ -115,7 +141,8 @@ void NewNodeHandshakeHandler::Work(const std::shared_ptr // Shard assignment for scale-out is not yet supported; no shards assigned here. // TODO: implement scale-out shard assignment MLOG_WARN("New node registered post-grace-period (scale-out not yet supported): logical_id={} ip={}", - logical_id, addr_str); + logical_id, + addr_str); break; } default: @@ -129,15 +156,13 @@ void NewNodeHandshakeHandler::Work(const std::shared_ptr } } else if (simm::common::ModuleServiceState::GetInstance().GracePeriodFinished()) { - MLOG_INFO("Grace period is already finished, new dataserver({}) will be waited for joining the cluster", - addr_str); + MLOG_INFO("Grace period is already finished, new dataserver({}) will be waited for joining the cluster", addr_str); // Post-grace-period without logical_node_id: legacy behavior } else { // Still in grace period: register node and assign shards MLOG_INFO("Still in Grace period new dataserver({}) will be added into cluster", addr_str); std::vector reported_shards(req->shard_ids().begin(), req->shard_ids().end()); - shard_manager_->BatchAssignRoutingTable(reported_shards, - std::make_shared(node_addr)); + shard_manager_->BatchAssignRoutingTable(reported_shards, std::make_shared(node_addr)); if (!logical_id.empty()) { // ProcessHandshake (Case 1) calls AddNode internally, so no separate AddNode needed. @@ -177,7 +202,7 @@ void NodeHeartBeatHandler::Work(const std::shared_ptr ctx auto resp = std::make_shared(); error_code_t ret = CommonErr::OK; std::string addr_str = req->node().ip() + ":" + std::to_string(req->node().port()); - const std::string& logical_id = req->logical_node_id(); + const std::string &logical_id = req->logical_node_id(); if (simm::utils::IsValidV4IPAddr(req->node().ip()) && simm::utils::IsValidPortNum(req->node().port())) { if (!logical_id.empty()) { @@ -406,31 +431,9 @@ void ListNodesHandler::Work(const std::shared_ptr ctx, // Build response with node info (address, status, resource) for (const auto &[addr_str, status] : node_stat_list) { - // Parse address string (ip:port format) - auto node_addr = simm::common::NodeAddress::ParseFromString(addr_str); - if (!node_addr) { - MLOG_ERROR("Failed to parse node address string({})", addr_str); - continue; - } - auto *node_info = resp->add_nodes(); - - // Set node address - auto *addr_pb = node_info->mutable_node_address(); - addr_pb->set_ip(node_addr->node_ip_); - addr_pb->set_port(node_addr->node_port_); - - // Set node status - node_info->set_node_status(static_cast(status)); - - // Set node resource if available auto res_it = res_map.find(addr_str); - if (res_it != res_map.end() && res_it->second) { - auto *res_pb = node_info->mutable_resource(); - res_pb->set_mem_free_bytes(res_it->second->mem_free_bytes_); - res_pb->set_mem_total_bytes(res_it->second->mem_total_bytes_); - res_pb->set_mem_used_bytes(res_it->second->mem_used_bytes_); - } + FillNodeInfoHelper(addr_str, status, res_it != res_map.end() ? res_it->second : nullptr, node_info); } resp->set_ret_code(CommonErr::OK); @@ -443,6 +446,74 @@ void ListNodesHandler::Work(const std::shared_ptr ctx, SEND_RESPONSE(ctx, resp); } +void GetNodeResourceHandler::Work(const std::shared_ptr ctx, + const std::shared_ptr conn, + const google::protobuf::Message *request) const { + auto req_begin_ts = std::chrono::steady_clock::now(); + auto req = dynamic_cast(request); + auto resp = std::make_shared(); + const std::string addr_str = req->node().ip() + ":" + std::to_string(req->node().port()); + + if (FOLLY_UNLIKELY(!simm::common::ModuleServiceState::GetInstance().IsServiceReady())) { + resp->set_ret_code(CmErr::InitInGracePeriod); + simm::common::Metrics::Instance("cluster_manager").IncErrorsTotal("get_node_resource"); + simm::common::Metrics::Instance("cluster_manager") + .ObserveRequestDuration("get_node_resource", + static_cast(std::chrono::duration_cast( + std::chrono::steady_clock::now() - req_begin_ts) + .count())); + simm::common::Metrics::Instance("cluster_manager").IncRequestsTotal("get_node_resource"); + SEND_RESPONSE(ctx, resp); + return; + } + + if (!node_manager_->QueryNodeExists(addr_str)) { + resp->set_ret_code(CommonErr::TargetNotFound); + simm::common::Metrics::Instance("cluster_manager").IncErrorsTotal("get_node_resource"); + simm::common::Metrics::Instance("cluster_manager") + .ObserveRequestDuration("get_node_resource", + static_cast(std::chrono::duration_cast( + std::chrono::steady_clock::now() - req_begin_ts) + .count())); + simm::common::Metrics::Instance("cluster_manager").IncRequestsTotal("get_node_resource"); + SEND_RESPONSE(ctx, resp); + return; + } + + auto resource = node_manager_->GetNodeResource(addr_str); + if (!resource) { + resource = node_manager_->RefreshNodeResource(addr_str); + } + if (!resource) { + resp->set_ret_code(CommonErr::TargetUnavailable); + simm::common::Metrics::Instance("cluster_manager").IncErrorsTotal("get_node_resource"); + simm::common::Metrics::Instance("cluster_manager") + .ObserveRequestDuration("get_node_resource", + static_cast(std::chrono::duration_cast( + std::chrono::steady_clock::now() - req_begin_ts) + .count())); + simm::common::Metrics::Instance("cluster_manager").IncRequestsTotal("get_node_resource"); + SEND_RESPONSE(ctx, resp); + return; + } + + FillNodeInfoHelper(addr_str, node_manager_->QueryNodeStatus(addr_str), resource, resp->mutable_node()); + for (const auto &shard_resource : resource->shard_mem_infos_) { + auto *shard_pb = resp->add_shard_resources(); + shard_pb->set_shard_id(shard_resource.shard_id_); + shard_pb->set_mem_used_bytes(shard_resource.mem_used_bytes_); + } + + resp->set_ret_code(CommonErr::OK); + simm::common::Metrics::Instance("cluster_manager") + .ObserveRequestDuration("get_node_resource", + static_cast(std::chrono::duration_cast( + std::chrono::steady_clock::now() - req_begin_ts) + .count())); + simm::common::Metrics::Instance("cluster_manager").IncRequestsTotal("get_node_resource"); + SEND_RESPONSE(ctx, resp); +} + void SetNodeStatusHandler::Work(const std::shared_ptr ctx, const std::shared_ptr conn, const google::protobuf::Message *request) const { diff --git a/src/cluster_manager/cm_rpc_handler.h b/src/cluster_manager/cm_rpc_handler.h index 2512e69..d3a6d93 100644 --- a/src/cluster_manager/cm_rpc_handler.h +++ b/src/cluster_manager/cm_rpc_handler.h @@ -152,6 +152,21 @@ class ListNodesHandler : public sicl::rpc::HandlerBase { std::shared_ptr node_manager_{nullptr}; }; +class GetNodeResourceHandler : public sicl::rpc::HandlerBase { + public: + explicit GetNodeResourceHandler(sicl::rpc::SiRPC *service, + google::protobuf::Message *request, + std::shared_ptr node_manager) + : HandlerBase(service, request), node_manager_(node_manager) {} + + virtual void Work(const std::shared_ptr ctx, + const std::shared_ptr conn, + const google::protobuf::Message *request) const override; + + private: + std::shared_ptr node_manager_{nullptr}; +}; + class SetNodeStatusHandler : public sicl::rpc::HandlerBase { public: explicit SetNodeStatusHandler( diff --git a/src/cluster_manager/cm_service.cc b/src/cluster_manager/cm_service.cc index bbae00e..6ab0daf 100644 --- a/src/cluster_manager/cm_service.cc +++ b/src/cluster_manager/cm_service.cc @@ -65,9 +65,7 @@ error_code_t ClusterManagerService::Start() { hb_monitor_->Start(); // start dataservers resource query background thread - // FIXME(ytji): for v0930 version, background thread(to query resource stats from ds) in - // node manager module is not actived yet, so just comment init action - // node_manager_->Init(); + node_manager_->Init(); MLOG_INFO("ClusterManager service starts successfully!"); } else { @@ -96,9 +94,7 @@ error_code_t ClusterManagerService::Stop() { // shard manager stop? // stop node manager background thread - // FIXME(ytji): for v0930 version, background thread(to query resource stats from ds) in - // node manager module is not actived yet, so just comment stop action - // node_manager_->Stop(); + node_manager_->Stop(); MLOG_INFO("Cluster Manager service stopped successfully..."); } else { @@ -150,14 +146,6 @@ error_code_t ClusterManagerService::StartRPCServices() { inter_rpc_service_->RegisterHandler( static_cast(simm::cm::ClusterManagerRpcType::RPC_NODE_HEARTBEAT), new NodeHeartBeatHandler(inter_rpc_service_.get(), new DataServerHeartBeatRequestPB, hb_monitor_)); - inter_rpc_service_->RegisterHandler( - static_cast(simm::cm::ClusterManagerRpcType::RPC_DATASERVER_RESOURCE_QUERY), - new DataServerResourceReportHandler(inter_rpc_service_.get(), new DataServerShardMemInfoPB)); - // TODO(ytji): will activate after v0930 version - // inter_rpc_service_->RegisterHandler(static_cast( - // simm::cm::ClusterManagerRpcType::RPC_ROUTING_TABLE_UPDATE), - // new RoutingTableUpdateHandler(inter_rpc_service_.get(), new - // RoutingTableUpdateRequestPB)); inter_rpc_service_->RegisterHandler( static_cast(simm::cm::ClusterManagerRpcType::RPC_ROUTING_TABLE_QUERY_SINGLE), new RoutingTableQuerySingleHandler( @@ -188,6 +176,9 @@ error_code_t ClusterManagerService::StartRPCServices() { admin_rpc_service_->RegisterHandler( static_cast(simm::common::CommonRpcType::RPC_LIST_NODE_REQ), new ListNodesHandler(admin_rpc_service_.get(), new ListNodesRequestPB, node_manager_)); + admin_rpc_service_->RegisterHandler( + static_cast(simm::common::CommonRpcType::RPC_GET_NODE_RESOURCE_REQ), + new GetNodeResourceHandler(admin_rpc_service_.get(), new GetNodeResourceRequestPB, node_manager_)); admin_rpc_service_->RegisterHandler( static_cast(simm::common::CommonRpcType::RPC_SET_NODE_STATUS_REQ), new SetNodeStatusHandler(admin_rpc_service_.get(), new SetNodeStatusRequestPB, node_manager_)); diff --git a/src/cluster_manager/cm_service.h b/src/cluster_manager/cm_service.h index 59baaf1..e1286c7 100644 --- a/src/cluster_manager/cm_service.h +++ b/src/cluster_manager/cm_service.h @@ -76,6 +76,7 @@ class ClusterManagerService { FRIEND_TEST(ClusterManagerServiceTest, TestQueryRoutingTableRejectsRequestsDuringGracePeriod); FRIEND_TEST(ClusterManagerServiceTest, TestQueryRoutingTableReturnsIncompleteWhenShardUnavailable); FRIEND_TEST(ClusterManagerServiceTest, TestNodeRejoinRPC); + FRIEND_TEST(ClusterManagerServiceTest, TestListNodesAndGetNodeResourceAdminRPCs); #endif }; diff --git a/src/common/base/common_types.h b/src/common/base/common_types.h index 5a32f2b..5f37d87 100644 --- a/src/common/base/common_types.h +++ b/src/common/base/common_types.h @@ -73,12 +73,36 @@ struct NodeAddress { struct NodeResource { NodeResource() {} - NodeResource(int64_t mem_total_bytes, int64_t mem_free_bytes, int64_t mem_used_bytes) - : mem_free_bytes_(mem_free_bytes), mem_total_bytes_(mem_total_bytes), mem_used_bytes_(mem_used_bytes) {} + struct ShardMemResource { + shard_id_t shard_id_{0}; + int64_t mem_used_bytes_{0}; + }; + + NodeResource(int64_t mem_total_bytes, int64_t mem_allocated_bytes, int64_t mem_used_bytes, int64_t mem_free_bytes) + : mem_free_bytes_(mem_free_bytes), + mem_total_bytes_(mem_total_bytes), + mem_allocated_bytes_(mem_allocated_bytes), + mem_used_bytes_(mem_used_bytes) {} + + NodeResource(int64_t mem_total_bytes, + int64_t mem_allocated_bytes, + int64_t mem_used_bytes, + int64_t mem_free_bytes, + uint64_t last_report_timestamp_us, + std::vector shard_mem_infos) + : mem_free_bytes_(mem_free_bytes), + mem_total_bytes_(mem_total_bytes), + mem_allocated_bytes_(mem_allocated_bytes), + mem_used_bytes_(mem_used_bytes), + last_report_timestamp_us_(last_report_timestamp_us), + shard_mem_infos_(std::move(shard_mem_infos)) {} int64_t mem_free_bytes_{0}; int64_t mem_total_bytes_{0}; + int64_t mem_allocated_bytes_{0}; int64_t mem_used_bytes_{0}; + uint64_t last_report_timestamp_us_{0}; + std::vector shard_mem_infos_; }; // timestmap info of node heatbeat @@ -134,4 +158,4 @@ class ModuleServiceState { }; } // namespace common -} // namespace simm \ No newline at end of file +} // namespace simm diff --git a/src/common/rpc_handlers/common_rpc_handlers.h b/src/common/rpc_handlers/common_rpc_handlers.h index 87a7133..5d9170d 100644 --- a/src/common/rpc_handlers/common_rpc_handlers.h +++ b/src/common/rpc_handlers/common_rpc_handlers.h @@ -18,6 +18,7 @@ enum class CommonRpcType { RPC_LIST_NODE_REQ, RPC_SET_NODE_STATUS_REQ, RPC_TRACE_TOGGLE_REQ, + RPC_GET_NODE_RESOURCE_REQ, //..... }; @@ -58,7 +59,6 @@ class ListGFlagsHandler : public sicl::rpc::HandlerBase { const google::protobuf::Message *request) const override; }; - #ifdef SIMM_ENABLE_TRACE class TraceToggleHandler : public sicl::rpc::HandlerBase { public: @@ -74,4 +74,4 @@ class TraceToggleHandler : public sicl::rpc::HandlerBase { #endif } // namespace common -} // namespace simm \ No newline at end of file +} // namespace simm diff --git a/src/data_server/kv_rpc_service.cc b/src/data_server/kv_rpc_service.cc index 26e3259..9083300 100644 --- a/src/data_server/kv_rpc_service.cc +++ b/src/data_server/kv_rpc_service.cc @@ -58,11 +58,8 @@ KVRpcService::~KVRpcService() { keepalive_thread_->join(); keepalive_thread_.reset(); } - (void)io_service_.reset(); - (void)mgmt_client_.reset(); - (void)mgt_service_.reset(); - (void)admin_rpc_service_.reset(); - (void)object_pool_.release(); // static variable cannot be deleted + // cache_pool_ relies on the mempool owned by io_service_, so release cache-related + // objects before destroying the RPC services. cache_pool_.reset(); cache_evictor_.reset(); for (auto [_, table] : all_tables_) { @@ -70,6 +67,11 @@ KVRpcService::~KVRpcService() { delete table; } all_tables_.clear(); + (void)mgmt_client_.reset(); + (void)mgt_service_.reset(); + (void)admin_rpc_service_.reset(); + (void)io_service_.reset(); + (void)object_pool_.release(); // static variable cannot be deleted } void KVRpcService::SetClusterDisconnectHandler(std::function handler) { @@ -631,8 +633,10 @@ error_code_t KVRpcService::KVPut(std::shared_ptr ctx, std::lock_guard table_lock(all_table_mtx_); table = all_tables_[shard_id]; if (!table) { - MLOG_ERROR("KVPut: shard {} table not initialized (not assigned by CM)", shard_id); - return DsErr::KeyNotFound; + MLOG_WARN("KVPut: shard {} table not initialized, create lazily on first write", shard_id); + table = new KVHashTable(); + table->Init(shard_id); + all_tables_[shard_id] = table; } } std::unique_ptr key_meta(object_pool_->AcquireKey(), @@ -861,10 +865,16 @@ void KVRpcService::GetResourceStats(const DataServerResourceRequestPB *req, Data rsp->mutable_node()->set_ip(ip); rsp->mutable_node()->set_port(port); size_t mem_total_bytes = cache_pool_->mem_total_bytes(); - size_t mem_used_bytes = cache_pool_->mem_used_bytes(); + size_t mem_allocated_bytes = cache_pool_->mem_used_bytes(); + size_t mem_used_bytes = 0; + for (uint32_t s = 0; s < FLAGS_shard_total_num; s++) { + mem_used_bytes += shard_used_bytes_[s].load(std::memory_order_relaxed); + } rsp->set_mem_total_bytes(mem_total_bytes); + rsp->set_mem_allocated_bytes(mem_allocated_bytes); rsp->set_mem_used_bytes(mem_used_bytes); rsp->set_mem_free_bytes(mem_total_bytes - mem_used_bytes); + rsp->set_last_report_timestamp_us(simm::utils::current_microseconds()); rsp->mutable_shard_mem_infos()->Reserve(FLAGS_shard_total_num); for (uint32_t s = 0; s < FLAGS_shard_total_num; s++) { diff --git a/src/proto/cm_clnt_rpcs.proto b/src/proto/cm_clnt_rpcs.proto index 9f8a627..bc9a37b 100644 --- a/src/proto/cm_clnt_rpcs.proto +++ b/src/proto/cm_clnt_rpcs.proto @@ -46,6 +46,13 @@ message NodeResourcePB { int64 mem_free_bytes = 1; int64 mem_total_bytes = 2; int64 mem_used_bytes = 3; + int64 mem_allocated_bytes = 4; + uint64 last_report_timestamp_us = 5; +} + +message ShardResourcePB { + uint32 shard_id = 1; + int64 mem_used_bytes = 2; } // Node Status Information @@ -66,6 +73,18 @@ message ListNodesResponsePB { repeated NodeInfoPB nodes = 2; } +// Clnt -> CM +message GetNodeResourceRequestPB { + proto.common.NodeAddressPB node = 1; +} + +// CM -> Clnt +message GetNodeResourceResponsePB { + sint32 ret_code = 1; + NodeInfoPB node = 2; + repeated ShardResourcePB shard_resources = 3; +} + // ********** Set Node Status RPCS ********** // Clnt -> CM message SetNodeStatusRequestPB { diff --git a/src/proto/ds_cm_rpcs.proto b/src/proto/ds_cm_rpcs.proto index 8fd0ce4..4fddff0 100644 --- a/src/proto/ds_cm_rpcs.proto +++ b/src/proto/ds_cm_rpcs.proto @@ -73,7 +73,9 @@ message DataServerResourceRequestPB { message DataServerResourceResponsePB { proto.common.NodeAddressPB node = 1; int64 mem_total_bytes = 2; - int64 mem_free_bytes = 3; + int64 mem_allocated_bytes = 3; int64 mem_used_bytes = 4; - repeated DataServerShardMemInfoPB shard_mem_infos = 5; + int64 mem_free_bytes = 5; + repeated DataServerShardMemInfoPB shard_mem_infos = 6; + uint64 last_report_timestamp_us = 7; } diff --git a/tests/cluster_manager/test_cm_node_manager.cc b/tests/cluster_manager/test_cm_node_manager.cc new file mode 100644 index 0000000..9000eca --- /dev/null +++ b/tests/cluster_manager/test_cm_node_manager.cc @@ -0,0 +1,209 @@ +#include +#include +#include +#include + +#include +#include + +#include + +#include "cluster_manager/cm_rpc_handler.h" +#include "cluster_manager/cm_node_manager.h" +#include "common/logging/logging.h" +#include "proto/ds_cm_rpcs.pb.h" +#include "rpc/connection.h" +#include "rpc/rpc.h" + +DECLARE_LOG_MODULE("cluster_manager_test"); + +DECLARE_uint32(dataserver_resource_interval_inSecs); +DECLARE_string(cm_log_file); +DECLARE_int32(mgt_service_port); + +namespace simm { +namespace cm { + +class ClusterManagerNodeManagerTestPeer { + public: + static void SetResourceQueryHook( + ClusterManagerNodeManager &node_manager, + std::function(const std::string &)> hook) { + node_manager.test_resource_query_hook_ = std::move(hook); + } +}; + +namespace { + +template +bool WaitUntil(Predicate pred, + std::chrono::milliseconds timeout, + std::chrono::milliseconds poll_interval = std::chrono::milliseconds(50)) { + const auto deadline = std::chrono::steady_clock::now() + timeout; + while (std::chrono::steady_clock::now() < deadline) { + if (pred()) { + return true; + } + std::this_thread::sleep_for(poll_interval); + } + return pred(); +} + +} // namespace + +class DelayedResourceQueryHandler : public sicl::rpc::HandlerBase { + public: + DelayedResourceQueryHandler(sicl::rpc::SiRPC *service, std::chrono::milliseconds delay) + : sicl::rpc::HandlerBase(service, new DataServerResourceRequestPB()), delay_(delay) {} + + void Work(const std::shared_ptr ctx, + const std::shared_ptr conn, + const google::protobuf::Message *request) const override { + (void)request; + auto rsp = std::make_shared(); + rsp->set_mem_total_bytes(4096); + rsp->set_mem_allocated_bytes(3584); + rsp->set_mem_free_bytes(1024); + rsp->set_mem_used_bytes(3072); + rsp->set_last_report_timestamp_us(123456); + auto *shard = rsp->add_shard_mem_infos(); + shard->set_shard_id(7); + shard->set_shard_mem_used_bytes(2048); + std::this_thread::sleep_for(delay_); + conn->SendResponse(*rsp, ctx, [rsp](std::shared_ptr) {}); + } + + private: + std::chrono::milliseconds delay_; +}; + +class ClusterManagerNodeManagerTest : public ::testing::Test { + protected: + void SetUp() override { + simm::logging::LogConfig cm_log_config = simm::logging::LogConfig{FLAGS_cm_log_file, "DEBUG"}; + simm::logging::LoggerManager::Instance().UpdateConfig("cluster_manager", cm_log_config); + old_interval_ = FLAGS_dataserver_resource_interval_inSecs; + FLAGS_dataserver_resource_interval_inSecs = 1; + } + + void TearDown() override { FLAGS_dataserver_resource_interval_inSecs = old_interval_; } + + private: + uint32_t old_interval_{0}; +}; + +TEST_F(ClusterManagerNodeManagerTest, TestResourcePollingRefreshesNodeAndShardStats) { + auto node_manager = std::make_shared(); + std::atomic generation{0}; + ClusterManagerNodeManagerTestPeer::SetResourceQueryHook(*node_manager, [&generation](const std::string &addr_str) { + EXPECT_EQ(addr_str, "127.0.0.1:44110"); + if (generation.load() == 0) { + return std::make_shared( + 1024, 896, 768, 256, 100, std::vector{{1, 128}, {8, 640}}); + } + return std::make_shared( + 2048, 1024, 512, 1536, 200, std::vector{{3, 512}}); + }); + + ASSERT_EQ(node_manager->AddNode("127.0.0.1:44110"), CommonErr::OK); + node_manager->Init(); + auto cleanup = folly::makeGuard([&]() { node_manager->Stop(); }); + + ASSERT_TRUE(WaitUntil( + [&]() { + auto resource = node_manager->GetNodeResource("127.0.0.1:44110"); + return resource != nullptr && resource->mem_used_bytes_ == 768 && resource->shard_mem_infos_.size() == 2; + }, + std::chrono::seconds(3))); + + auto resource = node_manager->GetNodeResource("127.0.0.1:44110"); + ASSERT_NE(resource, nullptr); + EXPECT_EQ(resource->mem_total_bytes_, 1024); + EXPECT_EQ(resource->mem_allocated_bytes_, 896); + EXPECT_EQ(resource->mem_free_bytes_, 256); + EXPECT_EQ(resource->mem_used_bytes_, 768); + EXPECT_EQ(resource->last_report_timestamp_us_, 100); + ASSERT_EQ(resource->shard_mem_infos_.size(), 2); + EXPECT_EQ(resource->shard_mem_infos_[0].shard_id_, 1); + EXPECT_EQ(resource->shard_mem_infos_[0].mem_used_bytes_, 128); + EXPECT_EQ(resource->shard_mem_infos_[1].shard_id_, 8); + EXPECT_EQ(resource->shard_mem_infos_[1].mem_used_bytes_, 640); + + generation.store(1); + ASSERT_TRUE(WaitUntil( + [&]() { + auto updated = node_manager->GetNodeResource("127.0.0.1:44110"); + return updated != nullptr && updated->mem_used_bytes_ == 512 && updated->last_report_timestamp_us_ == 200 && + updated->shard_mem_infos_.size() == 1 && updated->shard_mem_infos_[0].shard_id_ == 3; + }, + std::chrono::seconds(3))); +} + +TEST_F(ClusterManagerNodeManagerTest, TestRefreshNodeResourceWaitsForResponseCompletion) { + auto old_mgt_port = FLAGS_mgt_service_port; + auto restore_port = folly::makeGuard([&]() { FLAGS_mgt_service_port = old_mgt_port; }); + FLAGS_mgt_service_port = 44111; + + sicl::rpc::SiRPC *server_raw = nullptr; + ASSERT_EQ(sicl::rpc::SiRPC::newInstance(server_raw, true), sicl::transport::Result::SICL_SUCCESS); + std::unique_ptr server(server_raw); + ASSERT_TRUE(server->RegisterHandler( + static_cast(ClusterManagerRpcType::RPC_DATASERVER_RESOURCE_QUERY), + new DelayedResourceQueryHandler(server.get(), std::chrono::milliseconds(200)))); + ASSERT_EQ(server->Start(FLAGS_mgt_service_port), 0); + + std::jthread server_thread([&]() { server->RunUntilAskedToQuit(); }); + auto server_cleanup = folly::makeGuard([&]() { server->Stop(); }); + + auto node_manager = std::make_shared(); + auto resource = node_manager->RefreshNodeResource("127.0.0.1:44111"); + + ASSERT_NE(resource, nullptr); + EXPECT_EQ(resource->mem_total_bytes_, 4096); + EXPECT_EQ(resource->mem_allocated_bytes_, 3584); + EXPECT_EQ(resource->mem_free_bytes_, 1024); + EXPECT_EQ(resource->mem_used_bytes_, 3072); + EXPECT_EQ(resource->last_report_timestamp_us_, 123456); + ASSERT_EQ(resource->shard_mem_infos_.size(), 1); + EXPECT_EQ(resource->shard_mem_infos_[0].shard_id_, 7); + EXPECT_EQ(resource->shard_mem_infos_[0].mem_used_bytes_, 2048); +} + +TEST_F(ClusterManagerNodeManagerTest, TestRefreshNodeResourceUsesManagementPortForIoAddress) { + auto old_mgt_port = FLAGS_mgt_service_port; + auto restore_ports = folly::makeGuard([&]() { FLAGS_mgt_service_port = old_mgt_port; }); + FLAGS_mgt_service_port = 44121; + + sicl::rpc::SiRPC *server_raw = nullptr; + ASSERT_EQ(sicl::rpc::SiRPC::newInstance(server_raw, true), sicl::transport::Result::SICL_SUCCESS); + std::unique_ptr server(server_raw); + ASSERT_TRUE(server->RegisterHandler( + static_cast(ClusterManagerRpcType::RPC_DATASERVER_RESOURCE_QUERY), + new DelayedResourceQueryHandler(server.get(), std::chrono::milliseconds(10)))); + ASSERT_EQ(server->Start(FLAGS_mgt_service_port), 0); + + std::jthread server_thread([&]() { server->RunUntilAskedToQuit(); }); + auto server_cleanup = folly::makeGuard([&]() { server->Stop(); }); + + auto node_manager = std::make_shared(); + auto resource = node_manager->RefreshNodeResource("127.0.0.1:44120"); + + ASSERT_NE(resource, nullptr); + EXPECT_EQ(resource->mem_total_bytes_, 4096); + EXPECT_EQ(resource->mem_allocated_bytes_, 3584); + EXPECT_EQ(resource->mem_used_bytes_, 3072); + EXPECT_EQ(resource->last_report_timestamp_us_, 123456); +} + +TEST_F(ClusterManagerNodeManagerTest, TestAddNodeIsIdempotentForExistingAddress) { + auto node_manager = std::make_shared(); + ASSERT_EQ(node_manager->AddNode("127.0.0.1:44112"), CommonErr::OK); + ASSERT_EQ(node_manager->AddNode("127.0.0.1:44112"), CommonErr::OK); + + auto status_map = node_manager->GetAllNodeStatus(); + ASSERT_EQ(status_map.size(), 1); + EXPECT_EQ(status_map.at("127.0.0.1:44112"), NodeStatus::RUNNING); +} + +} // namespace cm +} // namespace simm diff --git a/tests/cluster_manager/test_cm_service.cc b/tests/cluster_manager/test_cm_service.cc index 6b60ac0..1086e09 100644 --- a/tests/cluster_manager/test_cm_service.cc +++ b/tests/cluster_manager/test_cm_service.cc @@ -1,9 +1,11 @@ #include #include #include +#include #include #include #include +#include #include #include @@ -37,6 +39,17 @@ namespace cm { using CmSrvPtr = std::unique_ptr; +class ClusterManagerNodeManagerTestPeer { + public: + static void SetResourceQueryHook( + ClusterManagerNodeManager &node_manager, + std::function(const std::string &)> hook) { + node_manager.test_resource_query_hook_ = std::move(hook); + } +}; + +namespace {} // namespace + class ClusterManagerServiceTest : public ::testing::Test { protected: void SetUp() override { @@ -779,6 +792,93 @@ TEST_F(ClusterManagerServiceTest, TestNodeRejoinRPC) { FLAGS_cm_cluster_init_grace_period_inSecs = old_flag_val; } +TEST_F(ClusterManagerServiceTest, TestListNodesAndGetNodeResourceAdminRPCs) { + ASSERT_EQ(cm_service_ptr_->Init(), CommonErr::OK); + simm::common::ModuleServiceState::GetInstance().MarkServiceReady(); + ASSERT_EQ(cm_service_ptr_->StartRPCServices(), CommonErr::OK); + auto stop_rpc = folly::makeGuard([&]() { EXPECT_EQ(cm_service_ptr_->StopRPCServices(), CommonErr::OK); }); + + ClusterManagerNodeManagerTestPeer::SetResourceQueryHook( + *cm_service_ptr_->node_manager_, [](const std::string &addr_str) { + EXPECT_EQ(addr_str, "127.0.0.1:44111"); + return std::make_shared( + 4096, 1024, 3072, 1024, 123456, std::vector{{1, 1024}, {9, 2048}}); + }); + ASSERT_EQ(cm_service_ptr_->node_manager_->AddNode("127.0.0.1:44111"), CommonErr::OK); + auto resource = cm_service_ptr_->node_manager_->RefreshNodeResource("127.0.0.1:44111"); + ASSERT_NE(resource, nullptr); + + sicl::rpc::SiRPC *sirpc_client = nullptr; + ASSERT_EQ(sicl::rpc::SiRPC::newInstance(sirpc_client, false), sicl::transport::Result::SICL_SUCCESS); + auto client = std::unique_ptr(sirpc_client); + + auto make_ctx = []() { + sicl::rpc::RpcContext *ctx_p = nullptr; + sicl::rpc::RpcContext::newInstance(ctx_p); + auto ctx = std::shared_ptr(ctx_p); + ctx->set_timeout(sicl::transport::TimerTick::TIMER_1S); + return ctx; + }; + + std::latch done_latch(2); + + ListNodesRequestPB list_req; + auto *list_rsp = new ListNodesResponsePB(); + client->SendRequest( + "127.0.0.1", + FLAGS_cm_rpc_admin_port, + static_cast(simm::common::CommonRpcType::RPC_LIST_NODE_REQ), + list_req, + list_rsp, + make_ctx(), + [&done_latch, list_rsp](const google::protobuf::Message *rsp, const std::shared_ptr ctx) { + EXPECT_EQ(ctx->ErrorCode(), sicl::transport::Result::SICL_SUCCESS); + auto response = dynamic_cast(rsp); + ASSERT_NE(response, nullptr); + ASSERT_EQ(response->ret_code(), CommonErr::OK); + ASSERT_EQ(response->nodes_size(), 1); + EXPECT_EQ(response->nodes(0).node_address().ip(), "127.0.0.1"); + EXPECT_EQ(response->nodes(0).node_address().port(), 44111); + EXPECT_EQ(response->nodes(0).resource().mem_total_bytes(), 4096); + EXPECT_EQ(response->nodes(0).resource().mem_allocated_bytes(), 1024); + EXPECT_EQ(response->nodes(0).resource().mem_used_bytes(), 3072); + EXPECT_EQ(response->nodes(0).resource().last_report_timestamp_us(), 123456); + delete list_rsp; + done_latch.count_down(); + }); + + GetNodeResourceRequestPB detail_req; + detail_req.mutable_node()->set_ip("127.0.0.1"); + detail_req.mutable_node()->set_port(44111); + auto *detail_rsp = new GetNodeResourceResponsePB(); + client->SendRequest("127.0.0.1", + FLAGS_cm_rpc_admin_port, + static_cast(simm::common::CommonRpcType::RPC_GET_NODE_RESOURCE_REQ), + detail_req, + detail_rsp, + make_ctx(), + [&done_latch, detail_rsp](const google::protobuf::Message *rsp, + const std::shared_ptr ctx) { + EXPECT_EQ(ctx->ErrorCode(), sicl::transport::Result::SICL_SUCCESS); + auto response = dynamic_cast(rsp); + ASSERT_NE(response, nullptr); + ASSERT_EQ(response->ret_code(), CommonErr::OK); + EXPECT_EQ(response->node().resource().mem_total_bytes(), 4096); + EXPECT_EQ(response->node().resource().mem_allocated_bytes(), 1024); + EXPECT_EQ(response->node().resource().mem_free_bytes(), 1024); + EXPECT_EQ(response->node().resource().mem_used_bytes(), 3072); + ASSERT_EQ(response->shard_resources_size(), 2); + EXPECT_EQ(response->shard_resources(0).shard_id(), 1); + EXPECT_EQ(response->shard_resources(0).mem_used_bytes(), 1024); + EXPECT_EQ(response->shard_resources(1).shard_id(), 9); + EXPECT_EQ(response->shard_resources(1).mem_used_bytes(), 2048); + delete detail_rsp; + done_latch.count_down(); + }); + + done_latch.wait(); +} + } // namespace cm } // namespace simm diff --git a/tests/data_server/test_ds_kv_service.cc b/tests/data_server/test_ds_kv_service.cc index c8bb463..c9ec28a 100644 --- a/tests/data_server/test_ds_kv_service.cc +++ b/tests/data_server/test_ds_kv_service.cc @@ -27,6 +27,7 @@ DECLARE_uint64(memory_limit_bytes); DECLARE_uint32(ds_free_memory_usable_ratio); DECLARE_uint32(cm_hb_tolerance_count); DECLARE_bool(ds_process_exit_cm_disconnection); +DECLARE_string(ds_logical_node_id); namespace simm { namespace ds { @@ -63,15 +64,21 @@ class KVServiceTest : public ::testing::Test { protected: void SetUp() override { rpcServicePtr = std::make_unique(); + old_logical_node_id_ = FLAGS_ds_logical_node_id; + FLAGS_ds_logical_node_id = "ut-ds-kv-service"; FLAGS_memory_limit_bytes = 1ULL << 31; FLAGS_ds_free_memory_usable_ratio = 100; auto ret = rpcServicePtr->Init(); ASSERT_EQ(ret, CommonErr::OK); } - void TearDown() override { rpcServicePtr.reset(); } + void TearDown() override { + rpcServicePtr.reset(); + FLAGS_ds_logical_node_id = old_logical_node_id_; + } std::unique_ptr rpcServicePtr; + std::string old_logical_node_id_; }; class KVServiceLightTest : public ::testing::Test { @@ -354,6 +361,29 @@ TEST_F(KVServiceTest, TestGetRejectsTooSmallClientBuffer) { clientPool->join(); } +TEST_F(KVServiceTest, TestKVPutLazilyInitializesMissingShardTable) { + ASSERT_EQ(rpcServicePtr->Start(), CommonErr::OK); + const std::string key = "test_lazy_init_put"; + const uint32_t shard_id = + hashkit::HashkitBase::Instance().generate_16bit_hash_value(key.c_str(), key.length()) % FLAGS_shard_total_num; + + ASSERT_EQ(rpcServicePtr->GetHashTable(shard_id), nullptr); + + auto req = std::make_unique(); + req->set_shard_id(shard_id); + req->set_key(key); + req->set_val_len(128); + + KVEntryIntrusivePtr entry; + auto ret = rpcServicePtr->KVPut(std::make_shared(), req.get(), entry); + + ASSERT_EQ(ret, CommonErr::OK); + ASSERT_NE(rpcServicePtr->GetHashTable(shard_id), nullptr); + ASSERT_TRUE(entry); + rpcServicePtr->KVPutSuccessHooks(entry); + ASSERT_EQ(rpcServicePtr->Stop(), CommonErr::OK); +} + TEST_F(KVServiceLightTest, TestHeartbeatFailureCountResetOnSuccess) { rpcServicePtr->heartbeat_failure_count_.store(4); rpcServicePtr->cm_ready_.store(true); diff --git a/tools/simm_ctl_admin.cc b/tools/simm_ctl_admin.cc index a0eb0b8..4895a08 100644 --- a/tools/simm_ctl_admin.cc +++ b/tools/simm_ctl_admin.cc @@ -19,11 +19,13 @@ #include #include -#include +#include #include +#include #include #include #include +#include #include #include #include @@ -47,6 +49,7 @@ #include "common/logging/logging.h" #include "common/trace/trace_server.h" #include "common/rpc_handlers/common_rpc_handlers.h" +#include "common/utils/time_util.h" #include "common/utils/tabulate.hpp" #include "proto/cm_clnt_rpcs.pb.h" #include "proto/common.pb.h" @@ -55,6 +58,18 @@ namespace po = boost::program_options; enum class ResourceType { Node, Shard, GFlag }; +static std::string FormatTimestampSeconds(uint64_t timestamp_us) { + if (timestamp_us == 0) { + return "-"; + } + + const auto secs = static_cast(timestamp_us / 1000000ULL); + std::tm local_tm = *std::localtime(&secs); + char time_buf[64]; + std::strftime(time_buf, sizeof(time_buf), "%Y-%m-%d_%H:%M:%S %Z", &local_tm); + return std::string(time_buf); +} + // Low-level channel interface class AdminChannel { public: @@ -74,17 +89,15 @@ class AdminChannel { // Message type for UDS admin channel enum class AdminMsgType : uint16_t { TRACE_TOGGLE = 1, - GFLAG_LIST = 2, - GFLAG_GET = 3, - GFLAG_SET = 4, + GFLAG_LIST = 2, + GFLAG_GET = 3, + GFLAG_SET = 4, }; // Unix domain socket implementation class UdsChannel : public AdminChannel { public: - explicit UdsChannel(std::string path, - AdminMsgType type) - : path_(std::move(path)), type_(type), fd_(-1) {} + explicit UdsChannel(std::string path, AdminMsgType type) : path_(std::move(path)), type_(type), fd_(-1) {} ~UdsChannel() override { if (fd_ >= 0) { ::close(fd_); @@ -105,8 +118,7 @@ class UdsChannel : public AdminChannel { std::strncpy(addr.sun_path, path_.c_str(), sizeof(addr.sun_path) - 1); addr.sun_path[sizeof(addr.sun_path) - 1] = '\0'; - socklen_t addr_len = static_cast( - offsetof(sockaddr_un, sun_path) + std::strlen(addr.sun_path)); + socklen_t addr_len = static_cast(offsetof(sockaddr_un, sun_path) + std::strlen(addr.sun_path)); if (::connect(fd_, reinterpret_cast(&addr), addr_len) < 0) { std::cerr << "trace uds: failed to connect to " << path_ << "\n"; @@ -184,8 +196,8 @@ class UdsChannel : public AdminChannel { public: bool Call(const google::protobuf::Message &req, google::protobuf::Message *resp, - std::function &)> done_cb) override { + std::function &)> + done_cb) override { // Extended framing: [uint32_t len][uint16_t type][payload] std::string buf; if (!req.SerializeToString(&buf)) { @@ -276,34 +288,29 @@ class RpcChannel : public AdminChannel { sicl::rpc::ReqType req_type, std::shared_ptr ctx, std::unique_ptr client) - : ip_(std::move(ip)), - port_(port), - req_type_(req_type), - ctx_(std::move(ctx)), - client_(std::move(client)) {} + : ip_(std::move(ip)), port_(port), req_type_(req_type), ctx_(std::move(ctx)), client_(std::move(client)) {} bool Init() override { return client_ != nullptr; } bool Call(const google::protobuf::Message &req, google::protobuf::Message *resp, - std::function &)> done_cb) override { + std::function &)> + done_cb) override { if (!client_) { return false; } google::protobuf::Message *resp_ptr = resp; client_->SendRequest(ip_, - port_, - req_type_, - req, - resp_ptr, - ctx_, - [done_cb](const google::protobuf::Message *rsp, - const std::shared_ptr ctx) { - if (done_cb) { - done_cb(rsp, ctx); - } - }); + port_, + req_type_, + req, + resp_ptr, + ctx_, + [done_cb](const google::protobuf::Message *rsp, const std::shared_ptr ctx) { + if (done_cb) { + done_cb(rsp, ctx); + } + }); return true; } @@ -338,19 +345,87 @@ static void CallbackNode(const std::string &operation, std::shared_ptr ctx_shared; InitRpcClientAndContext(rpc_client, ctx_shared); - if (operation == "list") { + if (operation == "list" || operation == "summary") { auto done_cb = [&](const google::protobuf::Message *rsp, const std::shared_ptr ctx) { if (ctx->Failed()) { std::cerr << "Error: RPC failed, err: " << ctx->ErrorText() << "\n"; } else { auto *response = dynamic_cast(rsp); if (response->ret_code() == CommonErr::OK) { + uint64_t cluster_total_bytes = 0; + uint64_t cluster_allocated_bytes = 0; + uint64_t cluster_free_bytes = 0; + uint64_t cluster_used_bytes = 0; + uint64_t running_nodes = 0; + uint64_t dead_nodes = 0; + + for (int i = 0; i < response->nodes_size(); ++i) { + const auto &node_info = response->nodes(i); + cluster_total_bytes += static_cast(node_info.resource().mem_total_bytes()); + cluster_allocated_bytes += static_cast(node_info.resource().mem_allocated_bytes()); + cluster_free_bytes += static_cast(node_info.resource().mem_free_bytes()); + cluster_used_bytes += static_cast(node_info.resource().mem_used_bytes()); + if (node_info.node_status() == static_cast(simm::common::NodeStatus::RUNNING)) { + running_nodes++; + } else { + dead_nodes++; + } + } + + if (operation == "summary") { + tabulate::Table cluster_tbl; + cluster_tbl.add_row( + {"Nodes", "Running", "Dead", "Total Memory (MB)", "Allocated Memory (MB)", "Used Memory (MB)", + "Free Memory (MB)"}); + cluster_tbl.add_row({std::to_string(response->nodes_size()), + std::to_string(running_nodes), + std::to_string(dead_nodes), + std::to_string(cluster_total_bytes / (1024 * 1024)), + std::to_string(cluster_allocated_bytes / (1024 * 1024)), + std::to_string(cluster_used_bytes / (1024 * 1024)), + std::to_string(cluster_free_bytes / (1024 * 1024))}); + cluster_tbl.row(0).format().font_style({tabulate::FontStyle::bold}); + std::cout << cluster_tbl << "\n"; + + tabulate::Table node_tbl; + node_tbl.add_row( + {"Node Address", "Status", "Total Memory (MB)", "Allocated Memory (MB)", "Used Memory (MB)", + "Free Memory (MB)"}) + .format() + .width(20); + for (int i = 0; i < response->nodes_size(); ++i) { + const auto &node_info = response->nodes(i); + const auto &node_addr = node_info.node_address(); + std::string addr_str = node_addr.ip() + ":" + std::to_string(node_addr.port()); + std::string_view status_str = + simm::common::NodeStatusToString(static_cast(node_info.node_status())); + node_tbl.add_row({addr_str, + std::string(status_str), + std::to_string(node_info.resource().mem_total_bytes() / (1024 * 1024)), + std::to_string(node_info.resource().mem_allocated_bytes() / (1024 * 1024)), + std::to_string(node_info.resource().mem_used_bytes() / (1024 * 1024)), + std::to_string(node_info.resource().mem_free_bytes() / (1024 * 1024))}); + } + node_tbl.column(0).format().width(20).font_style({tabulate::FontStyle::bold}); + node_tbl.column(1).format().width(12); + node_tbl.column(2).format().width(18); + node_tbl.column(3).format().width(18); + node_tbl.column(4).format().width(18); + node_tbl.column(5).format().width(18); + node_tbl.row(0).format().font_style({tabulate::FontStyle::bold}); + std::cout << node_tbl << std::endl; + delete static_cast(rsp); + done_latch.count_down(); + return; + } + tabulate::Table tbl; tbl.format().locale("C"); if (verbose) { // Verbose mode: show detailed information - tbl.add_row({"Node Address", "Status", "Total Memory (MB)", "Free Memory (MB)", "Used Memory (MB)"}) + tbl.add_row({"Node Address", "Status", "Total Memory (MB)", "Allocated Memory (MB)", "Used Memory (MB)", + "Free Memory (MB)"}) .format() .width(20); @@ -365,10 +440,11 @@ static void CallbackNode(const std::string &operation, // Convert memory bytes to MB std::string total_mem = std::to_string(node_info.resource().mem_total_bytes() / (1024 * 1024)); - std::string free_mem = std::to_string(node_info.resource().mem_free_bytes() / (1024 * 1024)); + std::string allocated_mem = std::to_string(node_info.resource().mem_allocated_bytes() / (1024 * 1024)); std::string used_mem = std::to_string(node_info.resource().mem_used_bytes() / (1024 * 1024)); + std::string free_mem = std::to_string(node_info.resource().mem_free_bytes() / (1024 * 1024)); - tbl.add_row({addr_str, status_str, total_mem, free_mem, used_mem}); + tbl.add_row({addr_str, status_str, total_mem, allocated_mem, used_mem, free_mem}); } tbl.column(0).format().width(20).font_style({tabulate::FontStyle::bold}); @@ -376,6 +452,7 @@ static void CallbackNode(const std::string &operation, tbl.column(2).format().width(18); tbl.column(3).format().width(18); tbl.column(4).format().width(18); + tbl.column(5).format().width(18); tbl.row(0).format().font_style({tabulate::FontStyle::bold}); } else { // Normal mode: show simple information @@ -402,6 +479,7 @@ static void CallbackNode(const std::string &operation, std::cerr << "Error: ListNodes RPC failed with ret_code: " << response->ret_code() << "\n"; } } + delete static_cast(rsp); done_latch.count_down(); }; @@ -414,6 +492,79 @@ static void CallbackNode(const std::string &operation, resp, ctx_shared, done_cb); + } else if (operation == "stat") { + size_t colon_pos = name.find(':'); + if (colon_pos == std::string::npos) { + std::cerr << "Invalid node address format. Expected IP:PORT but got: " << name << "\n"; + exit(1); + } + + GetNodeResourceRequestPB req; + req.mutable_node()->set_ip(name.substr(0, colon_pos)); + req.mutable_node()->set_port(std::stoi(name.substr(colon_pos + 1))); + + auto done_cb = [&](const google::protobuf::Message *rsp, const std::shared_ptr ctx) { + if (ctx->Failed()) { + std::cerr << "Error: RPC failed, err: " << ctx->ErrorText() << "\n"; + } else { + auto *response = dynamic_cast(rsp); + if (!response || response->ret_code() != CommonErr::OK) { + std::cerr << "Error: GetNodeResource RPC failed with ret_code: " << (response ? response->ret_code() : -1) + << "\n"; + } else { + tabulate::Table summary; + summary.add_row({"Node Address", + "Status", + "Total Memory (MB)", + "Allocated Memory (MB)", + "Used Memory (MB)", + "Free Memory (MB)", + "Last Report"}); + const auto &node_info = response->node(); + const auto &node_addr = node_info.node_address(); + std::string last_report = FormatTimestampSeconds(node_info.resource().last_report_timestamp_us()); + summary.add_row({node_addr.ip() + ":" + std::to_string(node_addr.port()), + std::string(simm::common::NodeStatusToString( + static_cast(node_info.node_status()))), + std::to_string(node_info.resource().mem_total_bytes() / (1024 * 1024)), + std::to_string(node_info.resource().mem_allocated_bytes() / (1024 * 1024)), + std::to_string(node_info.resource().mem_used_bytes() / (1024 * 1024)), + std::to_string(node_info.resource().mem_free_bytes() / (1024 * 1024)), + last_report}); + summary.row(0).format().font_style({tabulate::FontStyle::bold}); + std::cout << summary << "\n"; + + tabulate::Table shards; + shards.add_row({"Shard ID", "Used Memory (MB)"}); + int shown = 0; + for (int i = 0; i < response->shard_resources_size(); ++i) { + const auto &shard = response->shard_resources(i); + if (!verbose && shard.mem_used_bytes() == 0) { + continue; + } + shards.add_row({std::to_string(shard.shard_id()), std::to_string(shard.mem_used_bytes() / (1024 * 1024))}); + shown++; + } + shards.row(0).format().font_style({tabulate::FontStyle::bold}); + if (shown > 0) { + std::cout << shards << std::endl; + } else { + std::cout << "(no shard memory usage to display; use --verbose to show zero-usage shards)\n"; + } + } + } + delete static_cast(rsp); + done_latch.count_down(); + }; + + auto resp = new GetNodeResourceResponsePB(); + rpc_client->SendRequest(ip, + port, + static_cast(simm::common::CommonRpcType::RPC_GET_NODE_RESOURCE_REQ), + req, + resp, + ctx_shared, + done_cb); } else if (operation == "set") { // Parse node address from name (format: IP:PORT) size_t colon_pos = name.find(':'); @@ -588,10 +739,10 @@ static void CallbackShard(const std::string &operation, } static void CallbackGFlag(AdminChannel &channel, - const std::string &operation, - const std::string &name, - const std::string &value, - bool verbose) { + const std::string &operation, + const std::string &name, + const std::string &value, + bool verbose) { std::latch done_latch(1); if (operation == "list") { @@ -599,10 +750,7 @@ static void CallbackGFlag(AdminChannel &channel, proto::common::ListAllGFlagsRequestPB req; if (!channel.Call( - req, - resp, - [&](const google::protobuf::Message *rsp, - const std::shared_ptr &ctx) { + req, resp, [&](const google::protobuf::Message *rsp, const std::shared_ptr &ctx) { const auto *response = dynamic_cast(rsp); if (ctx && ctx->Failed()) { LOG_ERROR("RPC failed, err:{}", ctx->ErrorText()); @@ -613,11 +761,8 @@ static void CallbackGFlag(AdminChannel &channel, tbl.add_row({"Flag Name", "VALUE", "Default Value", "TYPE", "Description"}).format().width(15); for (int i = 0; i < response->flags_size(); ++i) { const auto &f = response->flags(i); - tbl.add_row({f.flag_name(), - f.flag_value(), - f.flag_default_value(), - f.flag_type(), - f.flag_description()}); + tbl.add_row( + {f.flag_name(), f.flag_value(), f.flag_default_value(), f.flag_type(), f.flag_description()}); } tbl.column(0).format().width(30).font_style({tabulate::FontStyle::bold}); tbl.column(4).format().width(30); @@ -631,8 +776,7 @@ static void CallbackGFlag(AdminChannel &channel, } std::cout << tbl << std::endl; } else { - LOG_ERROR("ListAllGFlagsResponsePB not ok, ret_code:{}", - response ? response->ret_code() : -1); + LOG_ERROR("ListAllGFlagsResponsePB not ok, ret_code:{}", response ? response->ret_code() : -1); std::cerr << "Error: ListGFlags RPC failed" << std::endl; } done_latch.count_down(); @@ -648,10 +792,7 @@ static void CallbackGFlag(AdminChannel &channel, req.set_flag_name(name); if (!channel.Call( - req, - resp, - [&](const google::protobuf::Message *rsp, - const std::shared_ptr &ctx) { + req, resp, [&](const google::protobuf::Message *rsp, const std::shared_ptr &ctx) { const auto *response = dynamic_cast(rsp); if (ctx && ctx->Failed()) { LOG_ERROR("RPC failed, err:{}", ctx->ErrorText()); @@ -670,8 +811,7 @@ static void CallbackGFlag(AdminChannel &channel, tbl.column(1).format().width(50); std::cout << tbl << std::endl; } else { - LOG_ERROR("GetGFlagValueResponsePB not ok, ret_code:{}", - response ? response->ret_code() : -1); + LOG_ERROR("GetGFlagValueResponsePB not ok, ret_code:{}", response ? response->ret_code() : -1); std::cerr << "Error: GetGFlagValue RPC failed" << std::endl; } done_latch.count_down(); @@ -688,10 +828,7 @@ static void CallbackGFlag(AdminChannel &channel, req.set_flag_value(value); if (!channel.Call( - req, - resp, - [&](const google::protobuf::Message *rsp, - const std::shared_ptr &ctx) { + req, resp, [&](const google::protobuf::Message *rsp, const std::shared_ptr &ctx) { const auto *response = dynamic_cast(rsp); if (ctx && ctx->Failed()) { LOG_ERROR("RPC failed, err:{}", ctx->ErrorText()); @@ -704,8 +841,7 @@ static void CallbackGFlag(AdminChannel &channel, tbl.column(1).format().width(50); std::cout << tbl << std::endl; } else { - LOG_ERROR("SetGFlagValueResponsePB not ok, ret_code:{}", - response ? response->ret_code() : -1); + LOG_ERROR("SetGFlagValueResponsePB not ok, ret_code:{}", response ? response->ret_code() : -1); std::cerr << "Error: SetGFlagValue RPC failed" << std::endl; } done_latch.count_down(); @@ -723,9 +859,7 @@ static void CallbackGFlag(AdminChannel &channel, done_latch.wait(); } -static void CallbackTrace(AdminChannel &channel, - const std::string &value, - bool verbose) { +static void CallbackTrace(AdminChannel &channel, const std::string &value, bool verbose) { (void)verbose; proto::common::TraceToggleRequestPB req; @@ -733,21 +867,18 @@ static void CallbackTrace(AdminChannel &channel, req.set_enable_trace(value == "1"); std::latch done_latch(1); - if (!channel.Call(req, - resp, - [&](const google::protobuf::Message *rsp, - const std::shared_ptr &ctx) { - const auto *response = dynamic_cast(rsp); - if (ctx && ctx->Failed()) { - LOG_ERROR("RPC failed, err:{}", ctx->ErrorText()); - } else if (!response || response->ret_code() != CommonErr::OK) { - LOG_ERROR("SetGFlagValueResponsePB not ok, ret_code:{}", - response ? response->ret_code() : -1); - std::cerr << "Error: SetGFlagValue RPC failed" << std::endl; - } - done_latch.count_down(); - delete resp; - })) { + if (!channel.Call( + req, resp, [&](const google::protobuf::Message *rsp, const std::shared_ptr &ctx) { + const auto *response = dynamic_cast(rsp); + if (ctx && ctx->Failed()) { + LOG_ERROR("RPC failed, err:{}", ctx->ErrorText()); + } else if (!response || response->ret_code() != CommonErr::OK) { + LOG_ERROR("SetGFlagValueResponsePB not ok, ret_code:{}", response ? response->ret_code() : -1); + std::cerr << "Error: SetGFlagValue RPC failed" << std::endl; + } + done_latch.count_down(); + delete resp; + })) { std::cerr << "trace: channel.Call() failed" << std::endl; return; } @@ -790,6 +921,8 @@ int main(int argc, char *argv[]) { std::cout << "OPTIONS:\n" << desc << "\n"; std::cout << "SUBCOMMANDS:\n" << " node list [OPTIONS] List all nodes\n" + << " node summary [OPTIONS] Show cluster-wide node resource summary\n" + << " node stat Show detailed resource stats for one node\n" << " node set Set node status (0=DEAD, 1=RUNNING)\n" << " shard list [OPTIONS] List all shards\n" << " gflag list [OPTIONS] List all gflags\n" @@ -810,12 +943,20 @@ int main(int argc, char *argv[]) { // Parse subcommand format: "resource_type operation" if (subcommand == "node") { if (args.empty()) { - std::cerr << "Error: node subcommand requires an operation (list or set)\n"; + std::cerr << "Error: node subcommand requires an operation (list, summary, stat or set)\n"; return 1; } operation = args[0]; if (operation == "list") { CallbackNode("list", "", "", ip, port, verbose); + } else if (operation == "summary") { + CallbackNode("summary", "", "", ip, port, verbose); + } else if (operation == "stat") { + if (args.size() < 2) { + std::cerr << "Error: node stat requires argument: \n"; + return 1; + } + CallbackNode("stat", args[1], "", ip, port, verbose); } else if (operation == "set") { if (args.size() < 3) { std::cerr << "Error: node set requires arguments: \n"; @@ -840,7 +981,7 @@ int main(int argc, char *argv[]) { std::cerr << "Error: Unknown shard operation: " << operation << "\n"; return 1; } - } else if(subcommand == "gflag" || subcommand == "trace"){ + } else if (subcommand == "gflag" || subcommand == "trace") { std::unique_ptr channel_ptr{nullptr}; operation = args[0]; if (pid == -1) { @@ -910,7 +1051,7 @@ int main(int argc, char *argv[]) { } std::string flag_name = args[1]; CallbackGFlag(*channel_ptr, "get", flag_name, "", verbose); - } else if (operation == "set"){ // set + } else if (operation == "set") { // set if (args.size() < 3) { std::cerr << "Error: gflag set requires arguments: \n"; return 1; @@ -950,4 +1091,4 @@ int main(int argc, char *argv[]) { } return 0; -} \ No newline at end of file +} diff --git a/tools/simm_stable_test.cc b/tools/simm_stable_test.cc index 7e4cf73..6c43d39 100644 --- a/tools/simm_stable_test.cc +++ b/tools/simm_stable_test.cc @@ -27,6 +27,7 @@ #include #include #include +#include #include #include #include @@ -105,8 +106,38 @@ std::string convert_to_readable_size(uint64_t bytes_num) { } struct LatencyStats { - static constexpr std::array kBucketUpperUs = - {50, 100, 200, 500, 1000, 2000, 5000, 10000, 20000, 50000, 100000, 200000, 500000, 1000000, UINT64_MAX}; + static constexpr std::array kBucketUpperUs = {25, + 50, + 75, + 100, + 150, + 200, + 300, + 400, + 500, + 750, + 1000, + 1500, + 2000, + 3000, + 5000, + 7500, + 10000, + 15000, + 20000, + 30000, + 50000, + 75000, + 100000, + 150000, + 200000, + 300000, + 500000, + 750000, + 1000000, + 2000000, + 5000000, + UINT64_MAX}; void add(micro_ts d) { const auto us = static_cast(std::max(0, d.count())); @@ -132,6 +163,26 @@ struct LatencyStats { } } + LatencyStats delta_from(const LatencyStats &base) const { + LatencyStats delta; + delta.count_ = count_ >= base.count_ ? count_ - base.count_ : 0; + delta.total_us_ = total_us_ >= base.total_us_ ? total_us_ - base.total_us_ : 0; + delta.min_us_ = delta.count_ == 0 ? 0 : std::numeric_limits::max(); + delta.max_us_ = 0; + for (size_t i = 0; i < bucket_counts_.size(); ++i) { + delta.bucket_counts_[i] = + bucket_counts_[i] >= base.bucket_counts_[i] ? bucket_counts_[i] - base.bucket_counts_[i] : 0; + if (delta.bucket_counts_[i] > 0) { + delta.min_us_ = std::min(delta.min_us_, kBucketUpperUs[i]); + delta.max_us_ = kBucketUpperUs[i]; + } + } + if (delta.min_us_ == std::numeric_limits::max()) { + delta.min_us_ = 0; + } + return delta; + } + double avg_us() const { return count_ == 0 ? 0.0 : static_cast(total_us_) / static_cast(count_); } uint64_t percentile(double pct) const { @@ -156,6 +207,46 @@ struct LatencyStats { std::array bucket_counts_{}; }; +struct FailureBreakdown { + void merge(const FailureBreakdown &o) { + put_error_rc_ += o.put_error_rc_; + overwrite_error_rc_ += o.overwrite_error_rc_; + get_error_rc_ += o.get_error_rc_; + get_size_mismatch_ += o.get_size_mismatch_; + get_data_mismatch_ += o.get_data_mismatch_; + exists_error_rc_ += o.exists_error_rc_; + exists_unexpected_hit_ += o.exists_unexpected_hit_; + delete_error_rc_ += o.delete_error_rc_; + submit_error_rc_ += o.submit_error_rc_; + for (const auto &[rc, cnt] : o.error_code_counts_) { + error_code_counts_[rc] += cnt; + } + } + + void add_error_code(int32_t rc) { + if (rc != CommonErr::OK) { + ++error_code_counts_[rc]; + } + } + + bool empty() const { + return put_error_rc_ == 0 && overwrite_error_rc_ == 0 && get_error_rc_ == 0 && get_size_mismatch_ == 0 && + get_data_mismatch_ == 0 && exists_error_rc_ == 0 && exists_unexpected_hit_ == 0 && delete_error_rc_ == 0 && + submit_error_rc_ == 0; + } + + uint64_t put_error_rc_{0}; + uint64_t overwrite_error_rc_{0}; + uint64_t get_error_rc_{0}; + uint64_t get_size_mismatch_{0}; + uint64_t get_data_mismatch_{0}; + uint64_t exists_error_rc_{0}; + uint64_t exists_unexpected_hit_{0}; + uint64_t delete_error_rc_{0}; + uint64_t submit_error_rc_{0}; + std::map error_code_counts_; +}; + struct ThreadStats { void merge(const ThreadStats &o) { put_ += o.put_; @@ -186,13 +277,14 @@ struct ThreadStats { put_size_bytes += o.put_size_bytes; get_size_bytes += o.get_size_bytes; latency_.merge(o.latency_); + failure_breakdown_.merge(o.failure_breakdown_); } uint64_t total_ops() const { return put_ + overwrite_put_ + get_ + exists_ + del_ + mput_ + mget_; } uint64_t total_failures() const { return put_fails_ + overwrite_put_fails_ + get_fails_ + exists_fails_ + del_fails_ + mput_fails_ + mget_fails_ + - submit_fails_ + data_mismatch_; + submit_fails_; } uint64_t put_{0}; @@ -223,6 +315,7 @@ struct ThreadStats { uint64_t put_size_bytes{0}; uint64_t get_size_bytes{0}; LatencyStats latency_{}; + FailureBreakdown failure_breakdown_{}; }; struct ValueSpec { @@ -240,8 +333,9 @@ enum class OpType { struct WorkerRuntime { explicit WorkerRuntime(uint32_t tid, size_t keyspace, size_t key_len_limit) - : slots(keyspace), keys(simm::tools::stable_test::BuildWorkerKeyspace(tid, keyspace, key_len_limit)) {} + : tid(tid), slots(keyspace), keys(simm::tools::stable_test::BuildWorkerKeyspace(tid, keyspace, key_len_limit)) {} + uint32_t tid; mutable std::mutex mutex; std::condition_variable cv; ThreadStats stats; @@ -262,6 +356,116 @@ struct PendingSingleOp { steady_clock_t::time_point start_ts; }; +ThreadStats DeltaStats(const ThreadStats &snapshot, const ThreadStats &base) { + ThreadStats delta = snapshot; + delta.put_ -= base.put_; + delta.put_fails_ -= base.put_fails_; + delta.put_succs_ -= base.put_succs_; + delta.overwrite_put_ -= base.overwrite_put_; + delta.overwrite_put_fails_ -= base.overwrite_put_fails_; + delta.overwrite_put_succs_ -= base.overwrite_put_succs_; + delta.get_ -= base.get_; + delta.get_fails_ -= base.get_fails_; + delta.get_succs_ -= base.get_succs_; + delta.exists_ -= base.exists_; + delta.exists_fails_ -= base.exists_fails_; + delta.exists_succs_ -= base.exists_succs_; + delta.del_ -= base.del_; + delta.del_fails_ -= base.del_fails_; + delta.del_succs_ -= base.del_succs_; + delta.mput_ -= base.mput_; + delta.mput_fails_ -= base.mput_fails_; + delta.mput_succs_ -= base.mput_succs_; + delta.mget_ -= base.mget_; + delta.mget_fails_ -= base.mget_fails_; + delta.mget_succs_ -= base.mget_succs_; + delta.data_match_ -= base.data_match_; + delta.data_mismatch_ -= base.data_mismatch_; + delta.expected_miss_ -= base.expected_miss_; + delta.submit_fails_ -= base.submit_fails_; + delta.put_size_bytes -= base.put_size_bytes; + delta.get_size_bytes -= base.get_size_bytes; + delta.latency_ = snapshot.latency_.delta_from(base.latency_); + return delta; +} + +std::string FormatTopErrorCodes(const std::map &error_code_counts, size_t limit = 6) { + if (error_code_counts.empty()) { + return "-"; + } + std::vector> entries(error_code_counts.begin(), error_code_counts.end()); + std::sort(entries.begin(), entries.end(), [](const auto &lhs, const auto &rhs) { + if (lhs.second != rhs.second) { + return lhs.second > rhs.second; + } + return lhs.first < rhs.first; + }); + + std::ostringstream oss; + for (size_t i = 0; i < std::min(limit, entries.size()); ++i) { + if (i > 0) { + oss << ", "; + } + oss << entries[i].first << ":" << entries[i].second; + } + return oss.str(); +} + +void PrintFailureBreakdown(const FailureBreakdown &b) { + if (b.empty()) { + std::cout << "FailureBreakdown: -\n"; + return; + } + + std::cout << "FailureBreakdown: " + << "PutErrRc=" << b.put_error_rc_ << ", " + << "OverwriteErrRc=" << b.overwrite_error_rc_ << ", " + << "GetErrRc=" << b.get_error_rc_ << ", " + << "GetSizeMismatch=" << b.get_size_mismatch_ << ", " + << "GetDataMismatch=" << b.get_data_mismatch_ << ", " + << "ExistsErrRc=" << b.exists_error_rc_ << ", " + << "ExistsUnexpectedHit=" << b.exists_unexpected_hit_ << ", " + << "DeleteErrRc=" << b.delete_error_rc_ << ", " + << "SubmitErrRc=" << b.submit_error_rc_ << "\n" + << "TopErrorCodes : " << FormatTopErrorCodes(b.error_code_counts_) << "\n"; +} + +void PrintTopThreadStats(const std::vector> &snapshots, uint64_t elapsed_secs, bool is_delta) { + struct Entry { + uint32_t tid{0}; + ThreadStats stats; + }; + + std::vector entries; + entries.reserve(snapshots.size()); + for (const auto &[tid, stats] : snapshots) { + entries.push_back(Entry{tid, stats}); + } + + std::sort(entries.begin(), entries.end(), [](const Entry &lhs, const Entry &rhs) { + if (lhs.stats.total_failures() != rhs.stats.total_failures()) { + return lhs.stats.total_failures() > rhs.stats.total_failures(); + } + if (lhs.stats.total_ops() != rhs.stats.total_ops()) { + return lhs.stats.total_ops() > rhs.stats.total_ops(); + } + return lhs.tid < rhs.tid; + }); + + const size_t limit = entries.size() <= 8 ? entries.size() : 8; + std::cout << "ThreadBreakdown : " << (is_delta ? "delta" : "total") << ", top " << limit << " thread(s)\n"; + for (size_t i = 0; i < limit && i < entries.size(); ++i) { + const auto &e = entries[i]; + const double qps = + elapsed_secs == 0 ? 0.0 : static_cast(e.stats.total_ops()) / static_cast(elapsed_secs); + std::cout << " tid=" << e.tid << " ops=" << e.stats.total_ops() << " fails=" << e.stats.total_failures() + << " qps=" << std::fixed << std::setprecision(2) << qps << " get_fails=" << e.stats.get_fails_ + << " exists_fails=" << e.stats.exists_fails_ << " del_fails=" << e.stats.del_fails_ + << " data_mismatch=" << e.stats.data_mismatch_ << " avg_us=" << e.stats.latency_.avg_us() + << " p99_us=" << e.stats.latency_.percentile(0.99) << "\n"; + } +} + uint32_t choose_size(uint32_t limit) { return FLAGS_fixed_kvsize ? limit : folly::Random::rand32(1, limit + 1); } @@ -294,6 +498,12 @@ void RecordLatency(ThreadStats &stats, steady_clock_t::time_point start_ts) { } void print_stats(const ThreadStats &s, uint32_t threads, uint64_t elapsed_secs, bool is_delta) { + static constexpr const char *kRed = "\033[31m"; + static constexpr const char *kGreen = "\033[32m"; + static constexpr const char *kReset = "\033[0m"; + auto red_label = [&](const char *label) { return std::string(kRed) + label + kReset; }; + auto green_label = [&](const char *label) { return std::string(kGreen) + label + kReset; }; + const double qps = elapsed_secs == 0 ? 0.0 : static_cast(s.total_ops()) / static_cast(elapsed_secs); const double throughput_mb = elapsed_secs == 0 ? 0.0 : static_cast(s.put_size_bytes + s.get_size_bytes) / @@ -304,31 +514,31 @@ void print_stats(const ThreadStats &s, uint32_t threads, uint64_t elapsed_secs, << "Threads : " << threads << "\n" << "ElapsedSecs : " << elapsed_secs << "\n" << "TotalOps : " << s.total_ops() << "\n" - << "Failures : " << s.total_failures() << "\n" - << "SubmitFails : " << s.submit_fails_ << "\n" + << red_label("Failures ") << ": " << s.total_failures() << "\n" + << red_label("SubmitFails ") << ": " << s.submit_fails_ << "\n" << "PutCnt : " << s.put_ << "\n" - << "PutFails : " << s.put_fails_ << "\n" - << "PutSuccs : " << s.put_succs_ << "\n" + << red_label("PutFails ") << ": " << s.put_fails_ << "\n" + << green_label("PutSuccs ") << ": " << s.put_succs_ << "\n" << "OverwriteCnt : " << s.overwrite_put_ << "\n" - << "OverwriteFails : " << s.overwrite_put_fails_ << "\n" - << "OverwriteSuccs : " << s.overwrite_put_succs_ << "\n" + << red_label("OverwriteFails ") << ": " << s.overwrite_put_fails_ << "\n" + << green_label("OverwriteSuccs ") << ": " << s.overwrite_put_succs_ << "\n" << "GetCnt : " << s.get_ << "\n" - << "GetFails : " << s.get_fails_ << "\n" - << "GetSuccs : " << s.get_succs_ << "\n" + << red_label("GetFails ") << ": " << s.get_fails_ << "\n" + << green_label("GetSuccs ") << ": " << s.get_succs_ << "\n" << "ExistsCnt : " << s.exists_ << "\n" - << "ExistsFails : " << s.exists_fails_ << "\n" - << "ExistsSuccs : " << s.exists_succs_ << "\n" + << red_label("ExistsFails ") << ": " << s.exists_fails_ << "\n" + << green_label("ExistsSuccs ") << ": " << s.exists_succs_ << "\n" << "DeleteCnt : " << s.del_ << "\n" - << "DeleteFails : " << s.del_fails_ << "\n" - << "DeleteSuccs : " << s.del_succs_ << "\n" + << red_label("DeleteFails ") << ": " << s.del_fails_ << "\n" + << green_label("DeleteSuccs ") << ": " << s.del_succs_ << "\n" << "MPutCnt : " << s.mput_ << "\n" - << "MPutFails : " << s.mput_fails_ << "\n" - << "MPutSuccs : " << s.mput_succs_ << "\n" + << red_label("MPutFails ") << ": " << s.mput_fails_ << "\n" + << green_label("MPutSuccs ") << ": " << s.mput_succs_ << "\n" << "MGetCnt : " << s.mget_ << "\n" - << "MGetFails : " << s.mget_fails_ << "\n" - << "MGetSuccs : " << s.mget_succs_ << "\n" - << "DataMatch : " << s.data_match_ << "\n" - << "DataMismatch : " << s.data_mismatch_ << "\n" + << red_label("MGetFails ") << ": " << s.mget_fails_ << "\n" + << green_label("MGetSuccs ") << ": " << s.mget_succs_ << "\n" + << green_label("DataMatch ") << ": " << s.data_match_ << "\n" + << red_label("DataMismatch ") << ": " << s.data_mismatch_ << "\n" << "ExpectedMiss : " << s.expected_miss_ << "\n" << "PutSize : " << convert_to_readable_size(s.put_size_bytes) << "\n" << "GetSize : " << convert_to_readable_size(s.get_size_bytes) << "\n" @@ -338,8 +548,9 @@ void print_stats(const ThreadStats &s, uint32_t threads, uint64_t elapsed_secs, << "P50LatencyUs : " << s.latency_.percentile(0.50) << "\n" << "P95LatencyUs : " << s.latency_.percentile(0.95) << "\n" << "P99LatencyUs : " << s.latency_.percentile(0.99) << "\n" - << "MaxLatencyUs : " << s.latency_.max_us_ << "\n" - << std::endl; + << "MaxLatencyUs : " << s.latency_.max_us_ << "\n"; + PrintFailureBreakdown(s.failure_breakdown_); + std::cout << std::endl; } void usage() { @@ -530,6 +741,31 @@ ThreadStats SnapshotStats(const std::vector> &wor return aggregated; } +std::vector> SnapshotPerWorkerStats(const std::vector> &workers) { + std::vector> snapshots; + snapshots.reserve(workers.size()); + for (const auto &worker : workers) { + std::lock_guard lock(worker->mutex); + snapshots.emplace_back(worker->tid, worker->stats); + } + return snapshots; +} + +std::vector> DeltaPerWorkerStats( + const std::vector> &snapshot, + const std::vector> &base) { + std::vector> deltas; + deltas.reserve(snapshot.size()); + for (size_t i = 0; i < snapshot.size(); ++i) { + if (i < base.size() && snapshot[i].first == base[i].first) { + deltas.emplace_back(snapshot[i].first, DeltaStats(snapshot[i].second, base[i].second)); + } else { + deltas.emplace_back(snapshot[i].first, snapshot[i].second); + } + } + return deltas; +} + void RecordExpectedExistsResult(ThreadStats &stats, bool expected_exists, int rc) { ++stats.exists_; if (expected_exists) { @@ -537,6 +773,8 @@ void RecordExpectedExistsResult(ThreadStats &stats, bool expected_exists, int rc ++stats.exists_succs_; } else { ++stats.exists_fails_; + ++stats.failure_breakdown_.exists_error_rc_; + stats.failure_breakdown_.add_error_code(rc); } return; } @@ -547,6 +785,7 @@ void RecordExpectedExistsResult(ThreadStats &stats, bool expected_exists, int rc ++stats.expected_miss_; } else { ++stats.exists_fails_; + ++stats.failure_breakdown_.exists_unexpected_hit_; } } else { ++stats.exists_succs_; @@ -613,9 +852,12 @@ void RunSyncSingleOp(uint32_t tid, simm::clnt::KVStore &kvstore, WorkerRuntime & } else { if (op.overwrite) { ++runtime.stats.overwrite_put_fails_; + ++runtime.stats.failure_breakdown_.overwrite_error_rc_; } else { ++runtime.stats.put_fails_; + ++runtime.stats.failure_breakdown_.put_error_rc_; } + runtime.stats.failure_breakdown_.add_error_code(rc); } return; } @@ -633,9 +875,18 @@ void RunSyncSingleOp(uint32_t tid, simm::clnt::KVStore &kvstore, WorkerRuntime & ++runtime.stats.get_succs_; ++runtime.stats.data_match_; runtime.stats.get_size_bytes += op.previous_spec.size; - } else { + } else if (rc == static_cast(op.previous_spec.size)) { ++runtime.stats.get_fails_; ++runtime.stats.data_mismatch_; + ++runtime.stats.failure_breakdown_.get_data_mismatch_; + } else { + ++runtime.stats.get_fails_; + if (rc >= 0) { + ++runtime.stats.failure_breakdown_.get_size_mismatch_; + } else { + ++runtime.stats.failure_breakdown_.get_error_rc_; + runtime.stats.failure_breakdown_.add_error_code(rc); + } } return; } @@ -664,6 +915,8 @@ void RunSyncSingleOp(uint32_t tid, simm::clnt::KVStore &kvstore, WorkerRuntime & } } else { ++runtime.stats.del_fails_; + ++runtime.stats.failure_breakdown_.delete_error_rc_; + runtime.stats.failure_breakdown_.add_error_code(rc); } } @@ -727,6 +980,8 @@ void RunSyncBatchOp(uint32_t tid, simm::clnt::KVStore &kvstore, WorkerRuntime &r runtime.slots[slots[i]] = specs[i]; } else { ++runtime.stats.put_fails_; + ++runtime.stats.failure_breakdown_.put_error_rc_; + runtime.stats.failure_breakdown_.add_error_code(put_rets[i]); all_ok = false; } } @@ -761,7 +1016,15 @@ void RunSyncBatchOp(uint32_t tid, simm::clnt::KVStore &kvstore, WorkerRuntime &r runtime.stats.get_size_bytes += specs[i].size; } else { ++runtime.stats.get_fails_; - ++runtime.stats.data_mismatch_; + if (get_rets[i] == static_cast(specs[i].size)) { + ++runtime.stats.data_mismatch_; + ++runtime.stats.failure_breakdown_.get_data_mismatch_; + } else if (get_rets[i] >= 0) { + ++runtime.stats.failure_breakdown_.get_size_mismatch_; + } else { + ++runtime.stats.failure_breakdown_.get_error_rc_; + runtime.stats.failure_breakdown_.add_error_code(get_rets[i]); + } all_ok = false; } } @@ -840,9 +1103,12 @@ void SubmitAsyncSingleOp(uint32_t tid, } else { if (op.overwrite) { ++runtime->stats.overwrite_put_fails_; + ++runtime->stats.failure_breakdown_.overwrite_error_rc_; } else { ++runtime->stats.put_fails_; + ++runtime->stats.failure_breakdown_.put_error_rc_; } + runtime->stats.failure_breakdown_.add_error_code(result); } if (runtime->inflight > 0) { --runtime->inflight; @@ -852,6 +1118,8 @@ void SubmitAsyncSingleOp(uint32_t tid, if (rc != CommonErr::OK) { std::lock_guard lock(runtime->mutex); ++runtime->stats.submit_fails_; + ++runtime->stats.failure_breakdown_.submit_error_rc_; + runtime->stats.failure_breakdown_.add_error_code(rc); if (runtime->inflight > 0) { --runtime->inflight; } @@ -873,9 +1141,18 @@ void SubmitAsyncSingleOp(uint32_t tid, ++runtime->stats.get_succs_; ++runtime->stats.data_match_; runtime->stats.get_size_bytes += op.previous_spec.size; - } else { + } else if (result == static_cast(op.previous_spec.size)) { ++runtime->stats.get_fails_; ++runtime->stats.data_mismatch_; + ++runtime->stats.failure_breakdown_.get_data_mismatch_; + } else { + ++runtime->stats.get_fails_; + if (result >= 0) { + ++runtime->stats.failure_breakdown_.get_size_mismatch_; + } else { + ++runtime->stats.failure_breakdown_.get_error_rc_; + runtime->stats.failure_breakdown_.add_error_code(result); + } } if (runtime->inflight > 0) { --runtime->inflight; @@ -885,6 +1162,8 @@ void SubmitAsyncSingleOp(uint32_t tid, if (rc != CommonErr::OK) { std::lock_guard lock(runtime->mutex); ++runtime->stats.submit_fails_; + ++runtime->stats.failure_breakdown_.submit_error_rc_; + runtime->stats.failure_breakdown_.add_error_code(rc); if (runtime->inflight > 0) { --runtime->inflight; } @@ -906,6 +1185,8 @@ void SubmitAsyncSingleOp(uint32_t tid, if (rc != CommonErr::OK) { std::lock_guard lock(runtime->mutex); ++runtime->stats.submit_fails_; + ++runtime->stats.failure_breakdown_.submit_error_rc_; + runtime->stats.failure_breakdown_.add_error_code(rc); if (runtime->inflight > 0) { --runtime->inflight; } @@ -930,6 +1211,8 @@ void SubmitAsyncSingleOp(uint32_t tid, } } else { ++runtime->stats.del_fails_; + ++runtime->stats.failure_breakdown_.delete_error_rc_; + runtime->stats.failure_breakdown_.add_error_code(result); } if (runtime->inflight > 0) { --runtime->inflight; @@ -939,6 +1222,8 @@ void SubmitAsyncSingleOp(uint32_t tid, if (rc != CommonErr::OK) { std::lock_guard lock(runtime->mutex); ++runtime->stats.submit_fails_; + ++runtime->stats.failure_breakdown_.submit_error_rc_; + runtime->stats.failure_breakdown_.add_error_code(rc); if (runtime->inflight > 0) { --runtime->inflight; } @@ -1000,6 +1285,7 @@ int main(int argc, char **argv) { auto test_start_ts = steady_clock_t::now(); auto test_end_ts = test_start_ts + std::chrono::seconds(FLAGS_time); ThreadStats last_snapshot; + std::vector> last_worker_snapshots; auto reporter = std::thread([&]() { while (!stop_threads.load(std::memory_order_acquire)) { std::this_thread::sleep_for(std::chrono::seconds(FLAGS_report_interval_inSecs)); @@ -1007,36 +1293,14 @@ int main(int argc, char **argv) { break; } auto snapshot = SnapshotStats(workers); - ThreadStats delta = snapshot; - delta.put_ -= last_snapshot.put_; - delta.put_fails_ -= last_snapshot.put_fails_; - delta.put_succs_ -= last_snapshot.put_succs_; - delta.overwrite_put_ -= last_snapshot.overwrite_put_; - delta.overwrite_put_fails_ -= last_snapshot.overwrite_put_fails_; - delta.overwrite_put_succs_ -= last_snapshot.overwrite_put_succs_; - delta.get_ -= last_snapshot.get_; - delta.get_fails_ -= last_snapshot.get_fails_; - delta.get_succs_ -= last_snapshot.get_succs_; - delta.exists_ -= last_snapshot.exists_; - delta.exists_fails_ -= last_snapshot.exists_fails_; - delta.exists_succs_ -= last_snapshot.exists_succs_; - delta.del_ -= last_snapshot.del_; - delta.del_fails_ -= last_snapshot.del_fails_; - delta.del_succs_ -= last_snapshot.del_succs_; - delta.mput_ -= last_snapshot.mput_; - delta.mput_fails_ -= last_snapshot.mput_fails_; - delta.mput_succs_ -= last_snapshot.mput_succs_; - delta.mget_ -= last_snapshot.mget_; - delta.mget_fails_ -= last_snapshot.mget_fails_; - delta.mget_succs_ -= last_snapshot.mget_succs_; - delta.data_match_ -= last_snapshot.data_match_; - delta.data_mismatch_ -= last_snapshot.data_mismatch_; - delta.expected_miss_ -= last_snapshot.expected_miss_; - delta.submit_fails_ -= last_snapshot.submit_fails_; - delta.put_size_bytes -= last_snapshot.put_size_bytes; - delta.get_size_bytes -= last_snapshot.get_size_bytes; + auto worker_snapshot = SnapshotPerWorkerStats(workers); + ThreadStats delta = DeltaStats(snapshot, last_snapshot); + auto worker_delta = DeltaPerWorkerStats(worker_snapshot, last_worker_snapshots); last_snapshot = snapshot; + last_worker_snapshots = worker_snapshot; print_stats(delta, FLAGS_threads, FLAGS_report_interval_inSecs, true); + PrintTopThreadStats(worker_delta, FLAGS_report_interval_inSecs, true); + std::cout << std::endl; } }); @@ -1085,7 +1349,9 @@ int main(int argc, char **argv) { auto elapsed = std::chrono::duration_cast(steady_clock_t::now() - test_start_ts).count(); auto final_stats = SnapshotStats(workers); + auto final_worker_stats = SnapshotPerWorkerStats(workers); print_stats(final_stats, FLAGS_threads, static_cast(elapsed), false); + PrintTopThreadStats(final_worker_stats, static_cast(elapsed), false); if (final_stats.total_ops() == 0) { std::cerr << "[simm_stable_test] no operations completed" << std::endl; return EIO;