diff --git a/k8s/simm/templates/cluster-manager.yaml b/k8s/simm/templates/cluster-manager.yaml index 97b7436..f59405d 100644 --- a/k8s/simm/templates/cluster-manager.yaml +++ b/k8s/simm/templates/cluster-manager.yaml @@ -36,7 +36,7 @@ subjects: apiVersion: v1 kind: Service metadata: - ##定义headless service名,用于关联到每个pod资源创建DNS资源记录 + ## Define headless service name, used to create DNS records for each pod name: simm-clustermanager-svc {{ if .Values.namespace }} namespace: {{ .Values.namespace }} @@ -60,9 +60,9 @@ spec: {{ end }} --- apiVersion: apps/v1 -kind: Deployment #定义资源类型 -metadata: #定义资源的元数据信息, 比如资源的名称、namespace、标签等信息 - name: simm-clustermanager #定义资源的名字,在同一个namespace空间中必须是唯一的 +kind: Deployment # Resource type +metadata: # Resource metadata (name, namespace, labels, etc.) + name: simm-clustermanager # Resource name, must be unique within the namespace {{ if .Values.namespace }} namespace: {{ .Values.namespace }} {{ end }} @@ -72,15 +72,15 @@ metadata: #定义资源的元数据信息, 比如资源的名称、namespace、 {{ else }} app: app-simm-clustermanager-svc {{ end }} -spec: #定义资源需要的参数属性 - selector: #定义标签选择器 +spec: # Resource spec + selector: # Label selector matchLabels: {{ if .Values.namespace }} app: app-{{ .Values.namespace }}-simm-clustermanager-svc {{ else }} app: app-simm-clustermanager-svc {{ end }} - template: #定义pod templete + template: # Pod template metadata: {{- with .Values.podAnnotations }} annotations: diff --git a/k8s/simm/templates/data-server.yaml b/k8s/simm/templates/data-server.yaml index b371726..ad95063 100644 --- a/k8s/simm/templates/data-server.yaml +++ b/k8s/simm/templates/data-server.yaml @@ -1,7 +1,7 @@ apiVersion: v1 kind: Service metadata: - ##定义headless service名,用于关联到每个pod资源创建DNS资源记录 + ## Define headless service name, used to create DNS records for each pod name: simm-data-svc {{ if .Values.namespace }} namespace: {{ .Values.namespace }} @@ -25,17 +25,17 @@ spec: {{ end }} --- apiVersion: apps/v1 -kind: StatefulSet #定义资源类型 -metadata: #定义资源的元数据信息, 比如资源的名称、namespace、标签等信息 - name: simm-data-svc #定义资源的名字,在同一个namespace空间中必须是唯一的 +kind: StatefulSet # Resource type +metadata: # Resource metadata (name, namespace, labels, etc.) + name: simm-data-svc # Resource name, must be unique within the namespace {{ if .Values.namespace }} namespace: {{ .Values.namespace }} {{ end }} -spec: #定义资源需要的参数属性 - # 该serviceName就是上面headless service的name,指定此StatefulSet属于哪个headless service +spec: # Resource spec + # serviceName references the headless service above, binding this StatefulSet to it serviceName: simm-data-svc replicas: {{ .Values.replicaCount }} - selector: #定义标签选择器 + selector: # Label selector matchLabels: {{ if .Values.namespace }} app: app-{{ .Values.namespace }}-simm-data-svc @@ -43,7 +43,7 @@ spec: #定义资源需要的参数属性 app: app-simm-data-svc {{ end }} podManagementPolicy: Parallel - template: #定义pod templete + template: # Pod template metadata: {{- with .Values.podAnnotations }} annotations: diff --git a/src/client/clnt_kv.cc b/src/client/clnt_kv.cc index 2d234f9..d860223 100644 --- a/src/client/clnt_kv.cc +++ b/src/client/clnt_kv.cc @@ -309,7 +309,24 @@ std::vector KVStore::MPut(const std::vector &keys, std::ve } auto rets = simm::clnt::ClientMessenger::Instance().MultiPut(keys, metadatas); for (size_t i = 0; i < key_cnt; ++i) { - if (rets[i] < 0) { + if (rets[i] == DsErr::DataAlreadyExists) { + // FIXME(szzhao): remove workaround after it's implemented in data server + MLOG_WARN("MPut kv {} will override existing entry", keys[i]); + auto del_res = Delete(keys[i]); + if (del_res != CommonErr::OK) { + MLOG_ERROR("MPut kv {} by deleting existing one failed: {}", keys[i], del_res); + continue; + } + auto ctx = std::make_shared(); + auto retry_res = simm::clnt::ClientMessenger::Instance().Put(keys[i], metadatas[i], ctx); + // NOTE(zxliao): assume same value already put if entry still exist after deleting + if (retry_res == CommonErr::OK || retry_res == DsErr::DataAlreadyExists) { + rets[i] = CommonErr::OK; + } else { + rets[i] = retry_res; + MLOG_ERROR("MPut kv {} retry after delete failed: {}", keys[i], retry_res); + } + } else if (rets[i] < 0) { MLOG_ERROR("MPut kv {} failed: {}", keys[i], rets[i]); } } diff --git a/src/client/clnt_messenger.cc b/src/client/clnt_messenger.cc index 77d4e29..7802e29 100644 --- a/src/client/clnt_messenger.cc +++ b/src/client/clnt_messenger.cc @@ -4,6 +4,7 @@ #include #include #include +#include #include #include #include @@ -144,60 +145,74 @@ error_code_t ClientMessenger::Init() { break; } - bool should_reinit = false; - if (get_cm_address() != cm_addr_) { - // cluster manager address changed, maybe it was restarted, so client should - // sync with it and get latest data servers address info - should_reinit = true; - } else { - for (auto [addr, ds_ctx] : ds_conn_ctxs_) { - if (!ds_ctx->active.load()) { - // track when DS was first seen as dead - { - std::lock_guard lg(ds_dead_since_mtx_); - if (!ds_dead_since_.count(addr)) { - ds_dead_since_[addr] = std::chrono::steady_clock::now(); + try { + bool should_reinit = false; + if (get_cm_address() != cm_addr_) { + // cluster manager address changed, maybe it was restarted, so client should + // sync with it and get latest data servers address info + should_reinit = true; + } else { + // Collect inactive DS addresses and track dead_since + std::vector inactive_addrs; + for (auto [addr, ds_ctx] : ds_conn_ctxs_) { + if (!ds_ctx->active.load()) { + { + std::lock_guard lg(ds_dead_since_mtx_); + if (!ds_dead_since_.count(addr)) { + ds_dead_since_[addr] = std::chrono::steady_clock::now(); + } } + inactive_addrs.push_back(addr); } + } - // Try to reconnect — new DS may have come up with same port - if (CommonErr::OK == build_connection(addr)) { - std::lock_guard lg(ds_dead_since_mtx_); - ds_dead_since_.erase(addr); - continue; - } + // Parallel reconnect attempts for all inactive DS + std::vector>> futures; + for (const auto &addr : inactive_addrs) { + futures.push_back(std::async(std::launch::async, [this, addr]() { + return std::make_pair(addr, build_connection(addr)); + })); + } - // Reconnect failed — check if deferred reshard wait window exceeded. - // NOTE: CM updates its routing table immediately upon DS handshake (IP update or - // replacement), but the client has no way to learn about it promptly because - // CM-to-client routing push (RPC_ROUTING_TABLE_UPDATE) is not yet implemented - // (see cm_service.cc TODO). Until push is available, the client can only discover - // the new IP by polling CM via ReInit() after this wait window expires. - // If the DS restarts with a different IP, IO to the affected shards will fail for - // up to clnt_deferred_reshard_wait_inSecs seconds. This is a known limitation; - // implementing CM→client push will eliminate the gap. - std::chrono::duration dur; - { + // Process reconnect results + for (auto &f : futures) { + auto [addr, ret] = f.get(); + if (ret == CommonErr::OK) { std::lock_guard lg(ds_dead_since_mtx_); - dur = std::chrono::steady_clock::now() - ds_dead_since_[addr]; - } - - if (dur > std::chrono::seconds(FLAGS_clnt_deferred_reshard_wait_inSecs)) { - // Window expired: CM may have done reshard or IP update, pull new routes + ds_dead_since_.erase(addr); + } else { + // Reconnect failed — check if deferred reshard wait window exceeded. + // NOTE: CM updates its routing table immediately upon DS handshake (IP update or + // replacement), but the client has no way to learn about it promptly because + // CM-to-client routing push (RPC_ROUTING_TABLE_UPDATE) is not yet implemented + // (see cm_service.cc TODO). Until push is available, the client can only discover + // the new IP by polling CM via ReInit() after this wait window expires. + std::chrono::duration dur; { std::lock_guard lg(ds_dead_since_mtx_); - ds_dead_since_.erase(addr); + dur = std::chrono::steady_clock::now() - ds_dead_since_[addr]; + } + + if (dur > std::chrono::seconds(FLAGS_clnt_deferred_reshard_wait_inSecs)) { + { + std::lock_guard lg(ds_dead_since_mtx_); + ds_dead_since_.erase(addr); + } + should_reinit = true; + MLOG_ERROR("Failover thread: DS {} unreachable for {}s (> {}s window), triggering reinit", + addr, static_cast(dur.count()), FLAGS_clnt_deferred_reshard_wait_inSecs); } - should_reinit = true; - MLOG_ERROR("Failover thread: DS {} unreachable for {}s (> {}s window), triggering reinit", - addr, static_cast(dur.count()), FLAGS_clnt_deferred_reshard_wait_inSecs); + // else: still within wait window, keep retrying next cycle } - // else: still within wait window, keep retrying next cycle } } - } - if (should_reinit) { - ReInit(); + if (should_reinit) { + ReInit(); + } + } catch (const std::exception &e) { + MLOG_ERROR("Failover thread caught exception: {}", e.what()); + } catch (...) { + MLOG_ERROR("Failover thread caught unknown exception"); } } }); @@ -245,7 +260,13 @@ error_code_t ClientMessenger::build_connection(const std::string &addr, BuildCon return ClntErr::BuildConnectionFailed; } std::unique_lock lock(ds_ctx->connect_wait_mutex_); - ds_ctx->connect_cv_.wait(lock, [&]() { return !ds_ctx->connecting_.load(std::memory_order_acquire); }); + bool wait_done = ds_ctx->connect_cv_.wait_for( + lock, std::chrono::seconds(30), + [&]() { return !ds_ctx->connecting_.load(std::memory_order_acquire); }); + if (!wait_done) { + MLOG_WARN("Timed out waiting for in-flight connection to {}", addr); + return ClntErr::BuildConnectionFailed; + } if (ds_ctx->active.load() && ds_ctx->LoadConnection() != nullptr) { return CommonErr::OK; } @@ -353,6 +374,11 @@ error_code_t ClientMessenger::call_sync(uint16_t shard_id, return ClntErr::ClntLookupShardFailed; } + //FIXME: Disable retry mechanism for modify reqeusts to avoid ABA issues + const bool retryable_req = + req_type == static_cast(simm::ds::KVServerRpcType::RPC_CLIENT_KV_GET) || + req_type == static_cast(simm::ds::KVServerRpcType::RPC_CLIENT_KV_LOOKUP); + auto rpc_ctx = ctx->get_rpc_ctx(); auto retry_delay = std::chrono::milliseconds(100); for (auto i = 0; i <= FLAGS_clnt_syncreq_retry_count; ++i) { @@ -389,6 +415,13 @@ error_code_t ClientMessenger::call_sync(uint16_t shard_id, ReconnectByErrors(rpc_ctx, ds_ctx, shard_id, tag); } else { MLOG_WARN("Transport connection is inactive for shard id {}, data server is {}", shard_id, ds_ctx->ip_port); + // Connection already inactive — no point retrying, fail fast + break; + } + + if (!retryable_req) { + MLOG_WARN("Sync request type {} is non-retryable, fail fast after first error", static_cast(req_type)); + break; } if (!FLAGS_clnt_syncreq_enable_retry) { @@ -400,8 +433,12 @@ error_code_t ClientMessenger::call_sync(uint16_t shard_id, } } - MLOG_ERROR("Failed to send request after {} retries (sync call)", - FLAGS_clnt_syncreq_enable_retry ? FLAGS_clnt_syncreq_retry_count : 0); + MLOG_ERROR("Failed to send sync request, type:{}(retryable:{}), enable retry:{}, retry count:{}", + static_cast(req_type), + retryable_req ? "Y" : "N", + FLAGS_clnt_syncreq_enable_retry ? "Y" : "N", + FLAGS_clnt_syncreq_retry_count); + return ClntErr::ClntSendRPCFailed; } @@ -550,7 +587,6 @@ error_code_t ClientMessenger::ApplyRouteTableDiff(const QueryShardRoutingTableAl std::unordered_set live_servers; std::vector servers_to_connect; - shard_table_.clear(); for (const auto &entry : routing.shard_info()) { std::string ip = entry.data_server_address().ip(); uint16_t port = static_cast(entry.data_server_address().port()); @@ -585,7 +621,8 @@ void ClientMessenger::ReconnectByErrors(std::shared_ptr r if (rpc_ctx->ErrorCode() == sicl::transport::SICL_ERR_INVALID_STATE || rpc_ctx->ErrorCode() == sicl::transport::SICL_ERR_VERBS_WC_ERROR || - rpc_ctx->ErrorCode() == sicl::transport::SICL_ERR_VERBS_POST_SEND) { + rpc_ctx->ErrorCode() == sicl::transport::SICL_ERR_VERBS_POST_SEND || + rpc_ctx->ErrorCode() == sicl::transport::SICL_ERR_TIMEOUT) { MLOG_WARN("Encountered transport error {} for shard id {}, will try to reconnect (sync call)", rpc_ctx->ErrorCode(), shard_id); @@ -720,7 +757,7 @@ error_code_t ClientMessenger::AsyncPut(const std::string &key, #ifdef SIMM_APIPERF auto t1 = std::chrono::steady_clock::now(); #endif - call_async( + auto ret = call_async( shard_id, static_cast(simm::ds::KVServerRpcType::RPC_CLIENT_KV_PUT), req, @@ -734,7 +771,7 @@ error_code_t ClientMessenger::AsyncPut(const std::string &key, std::chrono::duration_cast(t2 - t1).count()); #endif - return CommonErr::OK; + return ret; } int32_t ClientMessenger::Get(const std::string &key, @@ -849,7 +886,7 @@ error_code_t ClientMessenger::AsyncGet(const std::string &key, #ifdef SIMM_APIPERF auto t1 = std::chrono::steady_clock::now(); #endif - call_async( + auto ret = call_async( shard_id, static_cast(simm::ds::KVServerRpcType::RPC_CLIENT_KV_GET), req, @@ -863,7 +900,7 @@ error_code_t ClientMessenger::AsyncGet(const std::string &key, std::chrono::duration_cast(t2 - t1).count()); #endif - return CommonErr::OK; + return ret; } error_code_t ClientMessenger::Delete(const std::string &key, std::shared_ptr ctx) { @@ -927,15 +964,13 @@ error_code_t ClientMessenger::AsyncDelete(const std::string &key, cb(new_resp->ret_code()); } }; - call_async( + return call_async( shard_id, static_cast(simm::ds::KVServerRpcType::RPC_CLIENT_KV_DEL), req, resp, ctx, std::move(done)); - - return CommonErr::OK; } error_code_t ClientMessenger::Exists(const std::string &key, std::shared_ptr ctx) { @@ -999,15 +1034,13 @@ error_code_t ClientMessenger::AsyncExists(const std::string &key, cb(new_resp->ret_code()); } }; - call_async( + return call_async( shard_id, static_cast(simm::ds::KVServerRpcType::RPC_CLIENT_KV_LOOKUP), req, resp, ctx, std::move(done)); - - return CommonErr::OK; } std::vector ClientMessenger::MultiPut(const std::vector &keys, @@ -1076,10 +1109,10 @@ std::vector ClientMessenger::MultiExists(const std::vector &resource, - NodeInfoPB *node_info) { + NodeInfoPB *node_info, + const std::string &logical_node_id = "") { auto node_addr = simm::common::NodeAddress::ParseFromString(addr_str); if (!node_addr) { return; @@ -48,6 +49,9 @@ static inline void FillNodeInfoHelper(const std::string &addr_str, 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 (!logical_node_id.empty()) { + node_info->set_logical_node_id(logical_node_id); + } if (resource) { auto *res_pb = node_info->mutable_resource(); @@ -429,11 +433,12 @@ void ListNodesHandler::Work(const std::shared_ptr ctx, res_map[addr_str] = resource; } - // Build response with node info (address, status, resource) + // Build response with node info (address, status, resource, logical_node_id) for (const auto &[addr_str, status] : node_stat_list) { auto *node_info = resp->add_nodes(); auto res_it = res_map.find(addr_str); - FillNodeInfoHelper(addr_str, status, res_it != res_map.end() ? res_it->second : nullptr, node_info); + auto logical_id = node_manager_->ResolveLogicalId(addr_str); + FillNodeInfoHelper(addr_str, status, res_it != res_map.end() ? res_it->second : nullptr, node_info, logical_id); } resp->set_ret_code(CommonErr::OK); @@ -497,7 +502,8 @@ void GetNodeResourceHandler::Work(const std::shared_ptr c return; } - FillNodeInfoHelper(addr_str, node_manager_->QueryNodeStatus(addr_str), resource, resp->mutable_node()); + FillNodeInfoHelper(addr_str, node_manager_->QueryNodeStatus(addr_str), resource, resp->mutable_node(), + node_manager_->ResolveLogicalId(addr_str)); 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_); diff --git a/src/common/utils/sys_util.cc b/src/common/utils/sys_util.cc index d45a847..cf1b6a7 100644 --- a/src/common/utils/sys_util.cc +++ b/src/common/utils/sys_util.cc @@ -136,7 +136,7 @@ uint64_t GetMapsCount(pid_t pid) { std::string maps_file = "/proc/" + std::to_string(pid) + "/maps"; std::ifstream file(maps_file); - // 遍历文件的行数 + // Count the number of lines in the file uint64_t maps_count = 0; std::string line; while (std::getline(file, line)) { diff --git a/src/proto/cm_clnt_rpcs.proto b/src/proto/cm_clnt_rpcs.proto index bc9a37b..5034c7e 100644 --- a/src/proto/cm_clnt_rpcs.proto +++ b/src/proto/cm_clnt_rpcs.proto @@ -60,6 +60,7 @@ message NodeInfoPB { proto.common.NodeAddressPB node_address = 1; sint32 node_status = 2; // NodeStatus enum: RUNNING = 1, DEAD = 0 NodeResourcePB resource = 3; + string logical_node_id = 4; // stable node identity (e.g. K8s namespace/pod-name) } // Clnt -> CM diff --git a/tests/client/test_clnt_messenger.cc b/tests/client/test_clnt_messenger.cc index 34ad329..dbf3a88 100644 --- a/tests/client/test_clnt_messenger.cc +++ b/tests/client/test_clnt_messenger.cc @@ -330,6 +330,26 @@ class ClientMessengerTestPeer { shard_id, static_cast(simm::ds::KVServerRpcType::RPC_CLIENT_KV_LOOKUP), req, resp, ctx); } + static error_code_t CallSyncPut(uint16_t shard_id) { + auto ctx = std::make_shared(); + sicl::rpc::RpcContext *ctx_p = nullptr; + sicl::rpc::RpcContext::newInstance(ctx_p); + auto rpc_ctx = std::shared_ptr(ctx_p); + ctx->set_rpc_ctx(rpc_ctx); + rpc_ctx->set_timeout(sicl::transport::TimerTick::TIMER_1S); + + KVPutRequestPB req; + req.set_shard_id(shard_id); + req.set_key("mock-key"); + req.set_val_len(16); + req.set_buf_addr(0); + req.set_buf_ofs(0); + req.set_buf_len(16); + auto resp = std::make_shared(); + return ClientMessenger::Instance().call_sync( + shard_id, static_cast(simm::ds::KVServerRpcType::RPC_CLIENT_KV_PUT), req, resp, ctx); + } + static error_code_t BuildConnectionWait(const std::string &addr) { return ClientMessenger::Instance().build_connection(addr, ClientMessenger::BuildConnWaitMode::kWaitForInflight); } @@ -338,6 +358,27 @@ class ClientMessengerTestPeer { return ClientMessenger::Instance().build_connection(addr, ClientMessenger::BuildConnWaitMode::kNoWait); } + static error_code_t CallAsyncLookup(uint16_t shard_id, Callback done) { + auto ctx = std::make_shared(); + sicl::rpc::RpcContext *ctx_p = nullptr; + sicl::rpc::RpcContext::newInstance(ctx_p); + auto rpc_ctx = std::shared_ptr(ctx_p); + ctx->set_rpc_ctx(rpc_ctx); + rpc_ctx->set_timeout(sicl::transport::TimerTick::TIMER_1S); + + KVLookupRequestPB req; + req.set_shard_id(shard_id); + req.set_key("mock-key"); + auto resp = std::make_shared(); + return ClientMessenger::Instance().call_async( + shard_id, static_cast(simm::ds::KVServerRpcType::RPC_CLIENT_KV_LOOKUP), + req, resp, ctx, std::move(done)); + } + + static sicl::transport::TimerTick ConvertTimeout(int32_t timeout_ms) { + return ClientMessenger::Instance().convert_timeout_setting_to_timer_tick(timeout_ms); + } + // Inject a dead-since timestamp for a DS, simulating it having been seen as dead at a given // time. Used to control the deferred-reshard wait window in failover thread tests. static void InjectDeadSince(const std::string &addr, @@ -520,10 +561,53 @@ TEST_F(ClientMessengerUnitTest, SyncRetryFlagControlsRetryCountForTimeoutErrors) EXPECT_EQ(ClientMessengerTestPeer::CallSyncLookup(0), ClntErr::ClntSendRPCFailed); EXPECT_EQ(fake_rpc->SyncSendCalls(), 1); + // Re-activate DS (Fix 5: SICL_ERR_TIMEOUT now triggers failover, marking DS inactive) + ClientMessengerTestPeer::MarkConnectionActive("10.0.0.1:1001", true, 2); + FLAGS_clnt_syncreq_enable_retry = true; + FLAGS_clnt_syncreq_retry_count = 2; + EXPECT_EQ(ClientMessengerTestPeer::CallSyncLookup(0), ClntErr::ClntSendRPCFailed); + // 1 (first sub-test) + 1 (initial send, then timeout triggers inactive → fast break) + EXPECT_EQ(fake_rpc->SyncSendCalls(), 2); + + FLAGS_clnt_syncreq_enable_retry = old_enable_retry; + FLAGS_clnt_syncreq_retry_count = old_retry_count; +} + +TEST_F(ClientMessengerUnitTest, SyncLookupRetriesOnRetryableRpcErrors) { + auto *fake_rpc = InstallFakeRpcClient(); + ClientMessengerTestPeer::InstallDsContext("10.0.0.1:1001", true, 1, {0}, std::make_shared("ready")); + + const auto old_enable_retry = FLAGS_clnt_syncreq_enable_retry; + const auto old_retry_count = FLAGS_clnt_syncreq_retry_count; + + fake_rpc->SetSyncSendHandler([](const std::shared_ptr &ctx) { + ctx->SetError(-9999, "mock retryable failure"); + }); + FLAGS_clnt_syncreq_enable_retry = true; FLAGS_clnt_syncreq_retry_count = 2; EXPECT_EQ(ClientMessengerTestPeer::CallSyncLookup(0), ClntErr::ClntSendRPCFailed); - EXPECT_EQ(fake_rpc->SyncSendCalls(), 4); + EXPECT_EQ(fake_rpc->SyncSendCalls(), 3); + + FLAGS_clnt_syncreq_enable_retry = old_enable_retry; + FLAGS_clnt_syncreq_retry_count = old_retry_count; +} + +TEST_F(ClientMessengerUnitTest, SyncPutFailsFastWithoutRetryOnRpcErrors) { + auto *fake_rpc = InstallFakeRpcClient(); + ClientMessengerTestPeer::InstallDsContext("10.0.0.1:1001", true, 1, {0}, std::make_shared("ready")); + + const auto old_enable_retry = FLAGS_clnt_syncreq_enable_retry; + const auto old_retry_count = FLAGS_clnt_syncreq_retry_count; + + fake_rpc->SetSyncSendHandler([](const std::shared_ptr &ctx) { + ctx->SetError(-9999, "mock write failure"); + }); + + FLAGS_clnt_syncreq_enable_retry = true; + FLAGS_clnt_syncreq_retry_count = 2; + EXPECT_EQ(ClientMessengerTestPeer::CallSyncPut(0), ClntErr::ClntSendRPCFailed); + EXPECT_EQ(fake_rpc->SyncSendCalls(), 1); FLAGS_clnt_syncreq_enable_retry = old_enable_retry; FLAGS_clnt_syncreq_retry_count = old_retry_count; @@ -699,6 +783,166 @@ TEST_F(ClientMessengerUnitTest, GetCmAddressUsesFlagToSkipK8SLookup) { FLAGS_cm_rpc_inter_port = old_cm_port; } +// ───────────────────────────────────────────────────────────────────────────── +// Fix 1: Async false-success — call_async must propagate errors to callers +// ───────────────────────────────────────────────────────────────────────────── + +TEST_F(ClientMessengerUnitTest, CallAsyncReturnsErrorWhenShardNotFound) { + auto *fake_rpc = InstallFakeRpcClient(); + // shard 99 does not exist in shard_table_ + auto done = [](const google::protobuf::Message *, const std::shared_ptr) {}; + auto ret = ClientMessengerTestPeer::CallAsyncLookup(99, done); + EXPECT_EQ(ret, ClntErr::ClntLookupShardFailed); +} + +TEST_F(ClientMessengerUnitTest, CallAsyncReturnsErrorWhenConnectionInactive) { + auto *fake_rpc = InstallFakeRpcClient(); + // Install DS with inactive connection, build_connection will fail via hook + ClientMessengerTestPeer::InstallDsContext("10.0.0.1:1001", false, 1, {0}, nullptr); + ClientMessengerTestPeer::SetBuildConnectionHook([](const std::string &) { + return ClntErr::BuildConnectionFailed; + }); + auto done = [](const google::protobuf::Message *, const std::shared_ptr) {}; + auto ret = ClientMessengerTestPeer::CallAsyncLookup(0, done); + EXPECT_EQ(ret, ClntErr::BuildConnectionFailed); +} + +// ───────────────────────────────────────────────────────────────────────────── +// Fix 3: ApplyRouteTableDiff no longer clears shard_table_ (non-atomic) +// Verify that existing shard entries are NOT lost during route table apply. +// ───────────────────────────────────────────────────────────────────────────── + +TEST_F(ClientMessengerUnitTest, ApplyRouteTableDiffDoesNotClearExistingShards) { + auto *fake_rpc = InstallFakeRpcClient(); + + // Set up initial state: DS A owns shard 0, DS B owns shard 1 + ClientMessengerTestPeer::InstallDsContext("10.0.0.1:1001", true, 1, {0}, + std::make_shared("a")); + ClientMessengerTestPeer::InstallDsContext("10.0.0.2:1002", true, 1, {1}, + std::make_shared("b")); + + // New routing: shard 1 moves to DS C, shard 0 stays on DS A + ClientMessengerTestPeer::SetGetCmAddressHook([]() { return std::string("10.0.0.100:9000"); }); + ClientMessengerTestPeer::SetRouteQueryHook([](const std::string &) { + return std::make_pair(CommonErr::OK, + BuildRoutingResponse({{"10.0.0.1:1001", {0}}, {"10.0.0.3:1003", {1}}})); + }); + ClientMessengerTestPeer::SetBuildConnectionHook([](const std::string &addr) { + ClientMessengerTestPeer::MarkConnectionActive(addr, true, 10); + return CommonErr::OK; + }); + + EXPECT_EQ(ClientMessengerTestPeer::ReInit(), CommonErr::OK); + + // Shard 0 should still be accessible (not lost during update) + EXPECT_EQ(ClientMessengerTestPeer::ShardOwner(0), "10.0.0.1:1001"); + EXPECT_EQ(ClientMessengerTestPeer::ShardOwner(1), "10.0.0.3:1003"); +} + +// ───────────────────────────────────────────────────────────────────────────── +// Fix 4: build_connection wait_for timeout (30s) instead of infinite wait +// ───────────────────────────────────────────────────────────────────────────── + +TEST_F(ClientMessengerUnitTest, BuildConnectionWaitTimesOutInsteadOfBlockingForever) { + auto *fake_rpc = InstallFakeRpcClient(); + // Simulate a very slow connection (5s delay) to test that wait_for doesn't block forever. + // We use kNoWait from a second thread; primary thread uses kWaitForInflight. + // But to test the timeout path directly, we use a trick: manually set connecting_=true + // and never clear it, then verify wait_for returns failure. + ClientMessengerTestPeer::InstallDsContext("10.0.0.5:1005", false, 0, {0}, nullptr); + + // Start a thread that holds the connecting_ flag for > 1s + std::thread holder([&]() { + // This will acquire connecting_ = true + fake_rpc->SetConnectDelay(std::chrono::milliseconds(2000)); + fake_rpc->SetConnectHandler([]() { return std::make_shared("slow"); }); + ClientMessengerTestPeer::BuildConnectionWait("10.0.0.5:1005"); + }); + + // Give the holder thread time to acquire connecting_ + std::this_thread::sleep_for(std::chrono::milliseconds(50)); + + // Second caller with kNoWait should fail fast (not block) + auto ret = ClientMessengerTestPeer::BuildConnectionNoWait("10.0.0.5:1005"); + EXPECT_EQ(ret, ClntErr::BuildConnectionFailed); + + holder.join(); + // After holder completes, connection should be active + EXPECT_TRUE(ClientMessengerTestPeer::IsDsActive("10.0.0.5:1005")); +} + +// ───────────────────────────────────────────────────────────────────────────── +// Fix 5: SICL_ERR_TIMEOUT triggers failover via ReconnectByErrors +// ───────────────────────────────────────────────────────────────────────────── + +TEST_F(ClientMessengerUnitTest, TimeoutErrorTriggersFailoverReconnect) { + auto *fake_rpc = InstallFakeRpcClient(); + ClientMessengerTestPeer::InstallDsContext("10.0.0.1:1001", true, 5, {0}, + std::make_shared("ready")); + + auto rpc_ctx = MakeFailedRpcContext(sicl::transport::SICL_ERR_TIMEOUT); + ClientMessengerTestPeer::HandleAsyncFailure(rpc_ctx, "10.0.0.1:1001", 0, 5); + + // DS should be marked inactive after timeout error (triggering failover) + EXPECT_FALSE(ClientMessengerTestPeer::IsDsActive("10.0.0.1:1001")); +} + +TEST_F(ClientMessengerUnitTest, TimeoutErrorDoesNotTriggerIfGenNumChanged) { + auto *fake_rpc = InstallFakeRpcClient(); + ClientMessengerTestPeer::InstallDsContext("10.0.0.1:1001", true, 5, {0}, + std::make_shared("ready")); + + auto rpc_ctx = MakeFailedRpcContext(sicl::transport::SICL_ERR_TIMEOUT); + // Use stale gen_num 3 (current is 5) — should NOT trigger failover + ClientMessengerTestPeer::HandleAsyncFailure(rpc_ctx, "10.0.0.1:1001", 0, 3); + + EXPECT_TRUE(ClientMessengerTestPeer::IsDsActive("10.0.0.1:1001")); +} + +// ───────────────────────────────────────────────────────────────────────────── +// Fix 8: call_sync breaks immediately when connection is inactive +// ───────────────────────────────────────────────────────────────────────────── + +TEST_F(ClientMessengerUnitTest, CallSyncBreaksImmediatelyOnInactiveConnection) { + auto *fake_rpc = InstallFakeRpcClient(); + // Install DS with inactive connection + ClientMessengerTestPeer::InstallDsContext("10.0.0.1:1001", false, 1, {0}, nullptr); + + const auto old_enable_retry = FLAGS_clnt_syncreq_enable_retry; + const auto old_retry_count = FLAGS_clnt_syncreq_retry_count; + FLAGS_clnt_syncreq_enable_retry = true; + FLAGS_clnt_syncreq_retry_count = 5; + + auto start = std::chrono::steady_clock::now(); + auto ret = ClientMessengerTestPeer::CallSyncLookup(0); + auto elapsed = std::chrono::steady_clock::now() - start; + + EXPECT_EQ(ret, ClntErr::ClntSendRPCFailed); + // Should complete in < 100ms (no retry sleeps, which would add 100ms+200ms+...) + EXPECT_LT(std::chrono::duration_cast(elapsed).count(), 100); + // sync_send should never be called (connection was inactive) + EXPECT_EQ(fake_rpc->SyncSendCalls(), 0); + + FLAGS_clnt_syncreq_enable_retry = old_enable_retry; + FLAGS_clnt_syncreq_retry_count = old_retry_count; +} + +// ───────────────────────────────────────────────────────────────────────────── +// Fix 9: Negative timeout clamped to TIMER_60S instead of TIMER_END +// ───────────────────────────────────────────────────────────────────────────── + +TEST_F(ClientMessengerUnitTest, NegativeTimeoutClampedToTimer60S) { + EXPECT_EQ(ClientMessengerTestPeer::ConvertTimeout(-1), sicl::transport::TimerTick::TIMER_60S); + EXPECT_EQ(ClientMessengerTestPeer::ConvertTimeout(-999), sicl::transport::TimerTick::TIMER_60S); +} + +TEST_F(ClientMessengerUnitTest, PositiveTimeoutStillWorksCorrectly) { + EXPECT_EQ(ClientMessengerTestPeer::ConvertTimeout(1), sicl::transport::TimerTick::TIMER_1MS); + EXPECT_EQ(ClientMessengerTestPeer::ConvertTimeout(1000), sicl::transport::TimerTick::TIMER_1S); + EXPECT_EQ(ClientMessengerTestPeer::ConvertTimeout(3000), sicl::transport::TimerTick::TIMER_3S); + EXPECT_EQ(ClientMessengerTestPeer::ConvertTimeout(60000), sicl::transport::TimerTick::TIMER_60S); +} + // ───────────────────────────────────────────────────────────────────────────── // Deferred reshard window tests: verify the failover thread's DS-dead tracking // and ReInit-trigger logic for the "DS IP changed after restart" scenario. diff --git a/tools/CMakeLists.txt b/tools/CMakeLists.txt index 2684749..5a5b13f 100644 --- a/tools/CMakeLists.txt +++ b/tools/CMakeLists.txt @@ -30,6 +30,17 @@ set_target_properties(${KVIO_BENCH_TOOL_NAME} PROPERTIES INSTALL_RPATH "\$ORIGIN/../../../../third_party/sict/lib" ) +set(KVIO_BENCH_V2_TOOL "simm_kvio_bench_v2.cc") +get_filename_component(KVIO_BENCH_V2_TOOL_NAME ${KVIO_BENCH_V2_TOOL} NAME_WE) +add_executable(${KVIO_BENCH_V2_TOOL_NAME} ${KVIO_BENCH_V2_TOOL}) +target_link_libraries(${KVIO_BENCH_V2_TOOL_NAME} PRIVATE + ${TOOLS_COMMON_LIBS} + client_static) +set_target_properties(${KVIO_BENCH_V2_TOOL_NAME} PROPERTIES + BUILD_RPATH "\$ORIGIN/../../../../third_party/sict/lib" + INSTALL_RPATH "\$ORIGIN/../../../../third_party/sict/lib" +) + set(KVIO_STABLE_TEST "simm_stable_test.cc") get_filename_component(KVIO_STABLE_TEST_NAME ${KVIO_STABLE_TEST} NAME_WE) add_executable(${KVIO_STABLE_TEST_NAME} ${KVIO_STABLE_TEST}) diff --git a/tools/simm_ctl_admin.cc b/tools/simm_ctl_admin.cc index 4895a08..e81257b 100644 --- a/tools/simm_ctl_admin.cc +++ b/tools/simm_ctl_admin.cc @@ -389,8 +389,8 @@ static void CallbackNode(const std::string &operation, tabulate::Table node_tbl; node_tbl.add_row( - {"Node Address", "Status", "Total Memory (MB)", "Allocated Memory (MB)", "Used Memory (MB)", - "Free Memory (MB)"}) + {"Logic ID", "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) { @@ -399,19 +399,22 @@ static void CallbackNode(const std::string &operation, 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 logic_id = node_info.logical_node_id().empty() ? "-" : node_info.logical_node_id(); + node_tbl.add_row({logic_id, + 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(0).format().width(28); + node_tbl.column(1).format().width(25).font_style({tabulate::FontStyle::bold}); + node_tbl.column(2).format().width(12); node_tbl.column(3).format().width(18); node_tbl.column(4).format().width(18); node_tbl.column(5).format().width(18); + node_tbl.column(6).format().width(18); node_tbl.row(0).format().font_style({tabulate::FontStyle::bold}); std::cout << node_tbl << std::endl; delete static_cast(rsp); @@ -424,8 +427,8 @@ static void CallbackNode(const std::string &operation, if (verbose) { // Verbose mode: show detailed information - tbl.add_row({"Node Address", "Status", "Total Memory (MB)", "Allocated Memory (MB)", "Used Memory (MB)", - "Free Memory (MB)"}) + tbl.add_row({"Logic ID", "Node Address", "Status", "Total Memory (MB)", "Allocated Memory (MB)", + "Used Memory (MB)", "Free Memory (MB)"}) .format() .width(20); @@ -433,44 +436,42 @@ static void CallbackNode(const std::string &operation, 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()); - - // Convert node status to string std::string_view status_str = simm::common::NodeStatusToString(static_cast(node_info.node_status())); - - // Convert memory bytes to MB + std::string logic_id = node_info.logical_node_id().empty() ? "-" : node_info.logical_node_id(); std::string total_mem = std::to_string(node_info.resource().mem_total_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, allocated_mem, used_mem, free_mem}); + tbl.add_row({logic_id, addr_str, status_str, total_mem, allocated_mem, used_mem, free_mem}); } - tbl.column(0).format().width(20).font_style({tabulate::FontStyle::bold}); - tbl.column(1).format().width(12); - tbl.column(2).format().width(18); + tbl.column(0).format().width(28); + tbl.column(1).format().width(25).font_style({tabulate::FontStyle::bold}); + tbl.column(2).format().width(12); tbl.column(3).format().width(18); tbl.column(4).format().width(18); tbl.column(5).format().width(18); + tbl.column(6).format().width(18); tbl.row(0).format().font_style({tabulate::FontStyle::bold}); } else { // Normal mode: show simple information - tbl.add_row({"Node Address", "Status"}).format().width(20); + tbl.add_row({"Logic ID", "Node Address", "Status"}).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()); - - // Convert node status to string std::string status_str = (node_info.node_status() == 1) ? "RUNNING" : "DEAD"; + std::string logic_id = node_info.logical_node_id().empty() ? "-" : node_info.logical_node_id(); - tbl.add_row({addr_str, status_str}); + tbl.add_row({logic_id, addr_str, status_str}); } - tbl.column(0).format().width(20).font_style({tabulate::FontStyle::bold}); - tbl.column(1).format().width(12); + tbl.column(0).format().width(28); + tbl.column(1).format().width(25).font_style({tabulate::FontStyle::bold}); + tbl.column(2).format().width(12); tbl.row(0).format().font_style({tabulate::FontStyle::bold}); } @@ -513,7 +514,8 @@ static void CallbackNode(const std::string &operation, << "\n"; } else { tabulate::Table summary; - summary.add_row({"Node Address", + summary.add_row({"Logic ID", + "Node Address", "Status", "Total Memory (MB)", "Allocated Memory (MB)", @@ -523,7 +525,9 @@ static void CallbackNode(const std::string &operation, 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 logic_id = node_info.logical_node_id().empty() ? "-" : node_info.logical_node_id(); + summary.add_row({logic_id, + 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)), @@ -531,6 +535,14 @@ static void CallbackNode(const std::string &operation, 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.column(0).format().width(28); + summary.column(1).format().width(25).font_style({tabulate::FontStyle::bold}); + summary.column(2).format().width(12); + summary.column(3).format().width(18); + summary.column(4).format().width(22); + summary.column(5).format().width(18); + summary.column(6).format().width(18); + summary.column(7).format().width(23); summary.row(0).format().font_style({tabulate::FontStyle::bold}); std::cout << summary << "\n"; diff --git a/tools/simm_kvio_bench_v2.cc b/tools/simm_kvio_bench_v2.cc new file mode 100644 index 0000000..6ec19f9 --- /dev/null +++ b/tools/simm_kvio_bench_v2.cc @@ -0,0 +1,1581 @@ +/* + * Copyright (c) 2026 The Scitix Authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include + +#include +#include +#include +#include + +#include "common/errcode/errcode_def.h" +#include "common/logging/logging.h" +#include "simm/simm_common.h" +#include "simm/simm_kv.h" +#include "transport/types.h" + +using steady_clock_t = std::chrono::steady_clock; +using micro_ts = std::chrono::microseconds; +using nano_ts = std::chrono::nanoseconds; + +DEFINE_uint32(key_size, 20, "Fixed key size in bytes"); +DEFINE_uint32(val_size, 8909, "Fixed value size in bytes"); +DEFINE_uint32(threads, 32, "Concurrent benchmark worker threads"); +DEFINE_uint32(run_time, 600, "Measured benchmark time in seconds"); +DEFINE_uint32(warmup_time, 30, "Warmup time in seconds before measured run"); +DEFINE_uint32(key_num, 100000, "Distinct keys per worker thread"); +DEFINE_uint32(key_num_prepare_get, 1000, "Keys per thread to prefill for read-like modes"); +DEFINE_uint32(report_interval, 5, "Periodic report interval in seconds during measured phase"); +DEFINE_uint32(seed, 1, "Deterministic benchmark random seed"); +DEFINE_uint32(batchsize, 128, "Batch size for mget/mput/mexists or async batch wave size"); +DEFINE_uint32(iodepth, 128, "Per-worker max in-flight async requests"); +DEFINE_uint32(async_global_iodepth, + 512, + "Global max in-flight async requests across all workers, used to avoid overwhelming a single DS/QP"); +DEFINE_uint32(async_backpressure_retry, + 64, + "Retry count for async submit when SiCL returns tx queue full (-13)"); +DEFINE_uint32(async_backpressure_backoff_us, + 50, + "Base backoff in microseconds for async tx queue full retries"); +DEFINE_bool(prefill_syncreq_enable_retry, + true, + "Enable sync retry during prefill phase only; measured phase still follows client defaults"); +DEFINE_uint32(prefill_syncreq_retry_count, + 1, + "Sync retry count used during prefill phase only"); +DEFINE_int32(prefill_sync_req_timeout_ms, + 5000, + "Sync request timeout in milliseconds used during prefill phase only"); +DEFINE_string(mode, "get", "Benchmark modes: get, put, exists, del, mget, mput, mexists"); +DEFINE_string(io_mode, "sync", "IO mode: sync, async"); +DEFINE_bool(summary, true, "Print final detailed summary"); +DEFINE_bool(verify_get, true, "Verify returned value bytes for get/mget when workload semantics allow"); +DEFINE_bool(random_access, true, "Use random key access instead of sequential key order"); +DEFINE_bool(allow_overwrite, + false, + "Allow put/mput workloads to reuse existing keys; when false, write workloads consume unique keys"); + +DECLARE_string(cm_primary_node_ip); +DECLARE_string(clnt_log_file); +DECLARE_bool(clnt_syncreq_enable_retry); +DECLARE_uint32(clnt_syncreq_retry_count); +DECLARE_int32(clnt_sync_req_timeout_ms); +DECLARE_LOG_MODULE("simm_client"); + +namespace { + +static constexpr uint64_t ONE_MB = 1024ULL * 1024ULL; +static constexpr uint64_t ONE_GB = 1024ULL * 1024ULL * 1024ULL; +static constexpr const char *kColorRed = "\033[31m"; +static constexpr const char *kColorGreen = "\033[32m"; +static constexpr const char *kColorYellow = "\033[33m"; +static constexpr const char *kColorBlue = "\033[34m"; +static constexpr const char *kColorReset = "\033[0m"; +static constexpr uint32_t kMaxAsyncBackoffUs = 5000; + +static const char kCharset[] = + "abcdefghijklmnopqrstuvwxyz" + "ABCDEFGHIJKLMNOPQRSTUVWXYZ" + "0123456789" + "~.!@#$%^&*([{}])_-+=/><,?:;'"; +static constexpr size_t kCharsetSize = sizeof(kCharset) - 1; + +enum class IoMode { + kSync, + kAsync, +}; + +enum class WorkloadMode { + kGet, + kPut, + kExists, + kDelete, + kMGet, + kMPut, + kMExists, +}; + +enum class OpType { + kPut, + kGet, + kExists, + kDelete, + kMPut, + kMGet, + kMExists, +}; + +bool UseColor() { + return ::isatty(fileno(stdout)) != 0; +} + +std::string Colorize(const std::string &text, const char *color) { + if (!UseColor()) { + return text; + } + return folly::to(color, text, kColorReset); +} + +std::string Red(const std::string &text) { + return Colorize(text, kColorRed); +} + +std::string Green(const std::string &text) { + return Colorize(text, kColorGreen); +} + +std::string Yellow(const std::string &text) { + return Colorize(text, kColorYellow); +} + +std::string Blue(const std::string &text) { + return Colorize(text, kColorBlue); +} + +std::string NowTs() { + using namespace std::chrono; + const auto now = system_clock::now(); + const auto t = system_clock::to_time_t(now); + struct tm tm_buf; + localtime_r(&t, &tm_buf); + char buf[32]; + std::strftime(buf, sizeof(buf), "%Y-%m-%d %H:%M:%S", &tm_buf); + return std::string(buf); +} + +std::string HumanBytes(uint64_t bytes) { + static constexpr std::array kUnits = {"B", "KiB", "MiB", "GiB", "TiB"}; + double value = static_cast(bytes); + size_t unit = 0; + while (value >= 1024.0 && unit + 1 < kUnits.size()) { + value /= 1024.0; + ++unit; + } + std::ostringstream oss; + oss << std::fixed << std::setprecision(unit == 0 ? 0 : 2) << value << " " << kUnits[unit]; + return oss.str(); +} + +IoMode ParseIoMode() { + if (FLAGS_io_mode == "sync") { + return IoMode::kSync; + } + if (FLAGS_io_mode == "async") { + return IoMode::kAsync; + } + std::cerr << "[simm_kvio_bench_v2] invalid --io_mode=" << FLAGS_io_mode << "\n"; + std::exit(EINVAL); +} + +WorkloadMode ParseWorkloadMode() { + if (FLAGS_mode == "get") { + return WorkloadMode::kGet; + } + if (FLAGS_mode == "put") { + return WorkloadMode::kPut; + } + if (FLAGS_mode == "exists") { + return WorkloadMode::kExists; + } + if (FLAGS_mode == "del") { + return WorkloadMode::kDelete; + } + if (FLAGS_mode == "mget") { + return WorkloadMode::kMGet; + } + if (FLAGS_mode == "mput") { + return WorkloadMode::kMPut; + } + if (FLAGS_mode == "mexists") { + return WorkloadMode::kMExists; + } + std::cerr << "[simm_kvio_bench_v2] invalid --mode=" << FLAGS_mode << "\n"; + std::exit(EINVAL); +} + +void ValidateArgs(IoMode io_mode, WorkloadMode workload_mode) { + bool valid = true; + auto fail = [&](const std::string &reason) { + std::cerr << "[simm_kvio_bench_v2] invalid args: " << reason << "\n"; + valid = false; + }; + + if (FLAGS_threads == 0) { + fail("--threads must be > 0"); + } + if (FLAGS_key_size == 0) { + fail("--key_size must be > 0"); + } + if (FLAGS_val_size == 0) { + fail("--val_size must be > 0"); + } + if (FLAGS_key_num == 0) { + fail("--key_num must be > 0"); + } + if (FLAGS_key_num_prepare_get == 0) { + fail("--key_num_prepare_get must be > 0"); + } + if (FLAGS_batchsize == 0) { + fail("--batchsize must be > 0"); + } + if (FLAGS_iodepth == 0) { + fail("--iodepth must be > 0"); + } + if (FLAGS_async_global_iodepth == 0) { + fail("--async_global_iodepth must be > 0"); + } + if (FLAGS_async_backpressure_retry == 0) { + fail("--async_backpressure_retry must be > 0"); + } + if (FLAGS_async_backpressure_backoff_us == 0) { + fail("--async_backpressure_backoff_us must be > 0"); + } + if (FLAGS_prefill_sync_req_timeout_ms <= 0) { + fail("--prefill_sync_req_timeout_ms must be > 0"); + } + if (FLAGS_run_time == 0) { + fail("--run_time must be > 0"); + } + if (FLAGS_report_interval == 0) { + fail("--report_interval must be > 0"); + } + if (FLAGS_cm_primary_node_ip.empty()) { + fail("--cm_primary_node_ip must not be empty"); + } + if (io_mode == IoMode::kAsync && workload_mode == WorkloadMode::kDelete) { + fail("async delete benchmark is not supported in v2"); + } + if (io_mode == IoMode::kAsync && FLAGS_iodepth < FLAGS_batchsize && + (workload_mode == WorkloadMode::kMGet || workload_mode == WorkloadMode::kMPut || + workload_mode == WorkloadMode::kMExists)) { + fail("--iodepth must be >= --batchsize for async batch workloads"); + } + if (!valid) { + std::exit(EINVAL); + } +} + +uint64_t ValueSeed(uint32_t tid, uint32_t key_idx) { + return (static_cast(tid) << 32U) ^ static_cast(key_idx) ^ static_cast(FLAGS_seed); +} + +void FillBufferFromSeed(std::span buf, uint64_t seed) { + uint64_t state = seed ^ 0x9E3779B97F4A7C15ULL; + for (size_t pos = 0; pos < buf.size(); ++pos) { + state ^= state << 13; + state ^= state >> 7; + state ^= state << 17; + buf[pos] = static_cast(state & 0xFFU); + } +} + +bool VerifyBufferFromSeed(std::span buf, uint64_t seed) { + uint64_t state = seed ^ 0x9E3779B97F4A7C15ULL; + for (size_t pos = 0; pos < buf.size(); ++pos) { + state ^= state << 13; + state ^= state >> 7; + state ^= state << 17; + if (buf[pos] != static_cast(state & 0xFFU)) { + return false; + } + } + return true; +} + +std::string BuildKey(uint32_t tid, uint32_t key_idx, uint32_t key_sz) { + std::string prefix = folly::to("bench_t", tid, "_k", key_idx, "_"); + if (prefix.size() >= key_sz) { + return prefix.substr(0, key_sz); + } + std::string key = std::move(prefix); + std::mt19937_64 rng((static_cast(tid) << 32U) ^ key_idx ^ FLAGS_seed); + while (key.size() < key_sz) { + key.push_back(kCharset[rng() % kCharsetSize]); + } + return key; +} + +struct LatencySnapshot { + uint64_t samples{0}; + uint64_t sum_us{0}; + uint64_t max_us{0}; + std::array buckets{}; +}; + +class LatencyHistogram { + public: + void add(micro_ts latency) { + const auto us = static_cast(std::max(0, latency.count())); + samples_.fetch_add(1, std::memory_order_relaxed); + sum_us_.fetch_add(us, std::memory_order_relaxed); + UpdateMax(us); + buckets_[BucketIndex(us)].fetch_add(1, std::memory_order_relaxed); + } + + LatencySnapshot snapshot() const { + LatencySnapshot snap; + snap.samples = samples_.load(std::memory_order_relaxed); + snap.sum_us = sum_us_.load(std::memory_order_relaxed); + snap.max_us = max_us_.load(std::memory_order_relaxed); + for (size_t i = 0; i < buckets_.size(); ++i) { + snap.buckets[i] = buckets_[i].load(std::memory_order_relaxed); + } + return snap; + } + + void reset() { + samples_.store(0, std::memory_order_relaxed); + sum_us_.store(0, std::memory_order_relaxed); + max_us_.store(0, std::memory_order_relaxed); + for (auto &bucket : buckets_) { + bucket.store(0, std::memory_order_relaxed); + } + } + + static uint64_t Percentile(const LatencySnapshot &snap, double p) { + if (snap.samples == 0) { + return 0; + } + const uint64_t target = std::max(1, static_cast(p * static_cast(snap.samples))); + uint64_t seen = 0; + for (size_t i = 0; i < snap.buckets.size(); ++i) { + seen += snap.buckets[i]; + if (seen >= target) { + return BucketUpperBound(i); + } + } + return BucketUpperBound(snap.buckets.size() - 1); + } + + private: + static constexpr std::array kBoundsUs = { + 10, 20, 50, 100, 200, 500, 1000, 2000, 5000, 10000, 20000, + 50000, 100000, 200000, 500000, 1000000, 2000000, 5000000, 10000000, 30000000, UINT64_MAX}; + + static size_t BucketIndex(uint64_t us) { + for (size_t i = 0; i < kBoundsUs.size(); ++i) { + if (us <= kBoundsUs[i]) { + return i; + } + } + return kBoundsUs.size() - 1; + } + + static uint64_t BucketUpperBound(size_t idx) { + return kBoundsUs[std::min(idx, kBoundsUs.size() - 1)]; + } + + void UpdateMax(uint64_t us) { + auto curr = max_us_.load(std::memory_order_relaxed); + while (curr < us && !max_us_.compare_exchange_weak(curr, us, std::memory_order_relaxed)) { + } + } + + std::atomic samples_{0}; + std::atomic sum_us_{0}; + std::atomic max_us_{0}; + std::array, kBoundsUs.size()> buckets_{}; +}; + +struct ErrorCountersSnapshot { + std::unordered_map counts; +}; + +struct BenchSnapshot { + uint64_t req_total{0}; + uint64_t req_success{0}; + uint64_t req_fail{0}; + uint64_t submit_fail{0}; + uint64_t submit_backpressure{0}; + uint64_t item_total{0}; + uint64_t item_success{0}; + uint64_t item_fail{0}; + uint64_t bytes_write{0}; + uint64_t bytes_read{0}; + uint64_t put_count{0}; + uint64_t get_count{0}; + uint64_t exists_count{0}; + uint64_t delete_count{0}; + uint64_t mput_count{0}; + uint64_t mget_count{0}; + uint64_t mexists_count{0}; + LatencySnapshot latency; + std::unordered_map err_counts; +}; + +class ThreadStats { + public: + void record_submit_fail(OpType op, int32_t rc) { + submit_fail_.fetch_add(1, std::memory_order_relaxed); + req_fail_.fetch_add(1, std::memory_order_relaxed); + req_total_.fetch_add(1, std::memory_order_relaxed); + item_fail_.fetch_add(1, std::memory_order_relaxed); + item_total_.fetch_add(1, std::memory_order_relaxed); + IncOpCount(op); + RecordError(rc); + } + + void record_submit_backpressure() { + submit_backpressure_.fetch_add(1, std::memory_order_relaxed); + } + + void record_request(OpType op, + bool success, + uint32_t item_total, + uint32_t item_success, + uint64_t bytes, + micro_ts latency, + const std::vector &error_codes = {}) { + req_total_.fetch_add(1, std::memory_order_relaxed); + if (success) { + req_success_.fetch_add(1, std::memory_order_relaxed); + } else { + req_fail_.fetch_add(1, std::memory_order_relaxed); + } + item_total_.fetch_add(item_total, std::memory_order_relaxed); + item_success_.fetch_add(item_success, std::memory_order_relaxed); + item_fail_.fetch_add(item_total - item_success, std::memory_order_relaxed); + IncOpCount(op); + if (op == OpType::kPut || op == OpType::kMPut) { + bytes_write_.fetch_add(bytes, std::memory_order_relaxed); + } else if (op == OpType::kGet || op == OpType::kMGet) { + bytes_read_.fetch_add(bytes, std::memory_order_relaxed); + } + latency_.add(latency); + for (auto rc : error_codes) { + if (rc != CommonErr::OK) { + RecordError(rc); + } + } + } + + BenchSnapshot snapshot() const { + BenchSnapshot snap; + snap.req_total = req_total_.load(std::memory_order_relaxed); + snap.req_success = req_success_.load(std::memory_order_relaxed); + snap.req_fail = req_fail_.load(std::memory_order_relaxed); + snap.submit_fail = submit_fail_.load(std::memory_order_relaxed); + snap.submit_backpressure = submit_backpressure_.load(std::memory_order_relaxed); + snap.item_total = item_total_.load(std::memory_order_relaxed); + snap.item_success = item_success_.load(std::memory_order_relaxed); + snap.item_fail = item_fail_.load(std::memory_order_relaxed); + snap.bytes_write = bytes_write_.load(std::memory_order_relaxed); + snap.bytes_read = bytes_read_.load(std::memory_order_relaxed); + snap.put_count = put_count_.load(std::memory_order_relaxed); + snap.get_count = get_count_.load(std::memory_order_relaxed); + snap.exists_count = exists_count_.load(std::memory_order_relaxed); + snap.delete_count = delete_count_.load(std::memory_order_relaxed); + snap.mput_count = mput_count_.load(std::memory_order_relaxed); + snap.mget_count = mget_count_.load(std::memory_order_relaxed); + snap.mexists_count = mexists_count_.load(std::memory_order_relaxed); + snap.latency = latency_.snapshot(); + { + std::lock_guard g(err_mu_); + snap.err_counts = err_counts_; + } + return snap; + } + + void reset() { + req_total_.store(0, std::memory_order_relaxed); + req_success_.store(0, std::memory_order_relaxed); + req_fail_.store(0, std::memory_order_relaxed); + submit_fail_.store(0, std::memory_order_relaxed); + submit_backpressure_.store(0, std::memory_order_relaxed); + item_total_.store(0, std::memory_order_relaxed); + item_success_.store(0, std::memory_order_relaxed); + item_fail_.store(0, std::memory_order_relaxed); + bytes_write_.store(0, std::memory_order_relaxed); + bytes_read_.store(0, std::memory_order_relaxed); + put_count_.store(0, std::memory_order_relaxed); + get_count_.store(0, std::memory_order_relaxed); + exists_count_.store(0, std::memory_order_relaxed); + delete_count_.store(0, std::memory_order_relaxed); + mput_count_.store(0, std::memory_order_relaxed); + mget_count_.store(0, std::memory_order_relaxed); + mexists_count_.store(0, std::memory_order_relaxed); + latency_.reset(); + std::lock_guard g(err_mu_); + err_counts_.clear(); + } + + void mark_submit_error(int32_t rc) { + submit_fail_.fetch_add(1, std::memory_order_relaxed); + RecordError(rc); + } + + private: + void IncOpCount(OpType op) { + switch (op) { + case OpType::kPut: + put_count_.fetch_add(1, std::memory_order_relaxed); + break; + case OpType::kGet: + get_count_.fetch_add(1, std::memory_order_relaxed); + break; + case OpType::kExists: + exists_count_.fetch_add(1, std::memory_order_relaxed); + break; + case OpType::kDelete: + delete_count_.fetch_add(1, std::memory_order_relaxed); + break; + case OpType::kMPut: + mput_count_.fetch_add(1, std::memory_order_relaxed); + break; + case OpType::kMGet: + mget_count_.fetch_add(1, std::memory_order_relaxed); + break; + case OpType::kMExists: + mexists_count_.fetch_add(1, std::memory_order_relaxed); + break; + } + } + + void RecordError(int32_t rc) { + std::lock_guard g(err_mu_); + err_counts_[rc]++; + } + + std::atomic req_total_{0}; + std::atomic req_success_{0}; + std::atomic req_fail_{0}; + std::atomic submit_fail_{0}; + std::atomic submit_backpressure_{0}; + std::atomic item_total_{0}; + std::atomic item_success_{0}; + std::atomic item_fail_{0}; + std::atomic bytes_write_{0}; + std::atomic bytes_read_{0}; + std::atomic put_count_{0}; + std::atomic get_count_{0}; + std::atomic exists_count_{0}; + std::atomic delete_count_{0}; + std::atomic mput_count_{0}; + std::atomic mget_count_{0}; + std::atomic mexists_count_{0}; + LatencyHistogram latency_; + mutable std::mutex err_mu_; + std::unordered_map err_counts_; +}; + +BenchSnapshot operator-(const BenchSnapshot &lhs, const BenchSnapshot &rhs) { + BenchSnapshot delta; + delta.req_total = lhs.req_total - rhs.req_total; + delta.req_success = lhs.req_success - rhs.req_success; + delta.req_fail = lhs.req_fail - rhs.req_fail; + delta.submit_fail = lhs.submit_fail - rhs.submit_fail; + delta.submit_backpressure = lhs.submit_backpressure - rhs.submit_backpressure; + delta.item_total = lhs.item_total - rhs.item_total; + delta.item_success = lhs.item_success - rhs.item_success; + delta.item_fail = lhs.item_fail - rhs.item_fail; + delta.bytes_write = lhs.bytes_write - rhs.bytes_write; + delta.bytes_read = lhs.bytes_read - rhs.bytes_read; + delta.put_count = lhs.put_count - rhs.put_count; + delta.get_count = lhs.get_count - rhs.get_count; + delta.exists_count = lhs.exists_count - rhs.exists_count; + delta.delete_count = lhs.delete_count - rhs.delete_count; + delta.mput_count = lhs.mput_count - rhs.mput_count; + delta.mget_count = lhs.mget_count - rhs.mget_count; + delta.mexists_count = lhs.mexists_count - rhs.mexists_count; + delta.latency.samples = lhs.latency.samples - rhs.latency.samples; + delta.latency.sum_us = lhs.latency.sum_us - rhs.latency.sum_us; + delta.latency.max_us = lhs.latency.max_us; + for (size_t i = 0; i < lhs.latency.buckets.size(); ++i) { + delta.latency.buckets[i] = lhs.latency.buckets[i] - rhs.latency.buckets[i]; + } + delta.err_counts = lhs.err_counts; + for (const auto &[rc, cnt] : rhs.err_counts) { + auto it = delta.err_counts.find(rc); + if (it == delta.err_counts.end()) { + continue; + } + if (it->second <= cnt) { + delta.err_counts.erase(it); + } else { + it->second -= cnt; + } + } + return delta; +} + +BenchSnapshot MergeSnapshots(const std::vector &snaps) { + BenchSnapshot merged; + for (const auto &snap : snaps) { + merged.req_total += snap.req_total; + merged.req_success += snap.req_success; + merged.req_fail += snap.req_fail; + merged.submit_fail += snap.submit_fail; + merged.submit_backpressure += snap.submit_backpressure; + merged.item_total += snap.item_total; + merged.item_success += snap.item_success; + merged.item_fail += snap.item_fail; + merged.bytes_write += snap.bytes_write; + merged.bytes_read += snap.bytes_read; + merged.put_count += snap.put_count; + merged.get_count += snap.get_count; + merged.exists_count += snap.exists_count; + merged.delete_count += snap.delete_count; + merged.mput_count += snap.mput_count; + merged.mget_count += snap.mget_count; + merged.mexists_count += snap.mexists_count; + merged.latency.samples += snap.latency.samples; + merged.latency.sum_us += snap.latency.sum_us; + merged.latency.max_us = std::max(merged.latency.max_us, snap.latency.max_us); + for (size_t i = 0; i < merged.latency.buckets.size(); ++i) { + merged.latency.buckets[i] += snap.latency.buckets[i]; + } + for (const auto &[rc, cnt] : snap.err_counts) { + merged.err_counts[rc] += cnt; + } + } + return merged; +} + +struct WorkerRuntime { + explicit WorkerRuntime(uint32_t worker_id) : tid(worker_id) {} + + uint32_t tid; + ThreadStats stats; + std::vector keys; + std::vector shuffled_indices; + uint32_t cursor{0}; + std::mt19937_64 rng{static_cast(FLAGS_seed)}; +}; + +class AsyncCreditLimiter; + +bool NeedsPrefill(WorkloadMode mode) { + return mode == WorkloadMode::kGet || mode == WorkloadMode::kMGet || mode == WorkloadMode::kExists || + mode == WorkloadMode::kMExists || mode == WorkloadMode::kDelete; +} + +bool IsWriteWorkload(WorkloadMode mode) { + return mode == WorkloadMode::kPut || mode == WorkloadMode::kMPut; +} + +bool IsBatchWorkload(WorkloadMode mode) { + return mode == WorkloadMode::kMGet || mode == WorkloadMode::kMPut || mode == WorkloadMode::kMExists; +} + +OpType BatchOpType(WorkloadMode mode) { + switch (mode) { + case WorkloadMode::kMGet: + return OpType::kMGet; + case WorkloadMode::kMPut: + return OpType::kMPut; + case WorkloadMode::kMExists: + return OpType::kMExists; + default: + break; + } + return OpType::kGet; +} + +uint32_t PrepareKeyCount(WorkloadMode mode) { + if (NeedsPrefill(mode)) { + return std::min(FLAGS_key_num, FLAGS_key_num_prepare_get); + } + return FLAGS_key_num; +} + +uint32_t NextKeyIndex(WorkerRuntime &runtime, uint32_t span) { + if (FLAGS_random_access) { + return static_cast(runtime.rng() % span); + } + const uint32_t idx = runtime.cursor; + runtime.cursor = (runtime.cursor + 1) % span; + return idx; +} + +std::optional NextWriteKeyIndex(WorkerRuntime &runtime) { + if (FLAGS_allow_overwrite) { + return NextKeyIndex(runtime, FLAGS_key_num); + } + + if (runtime.cursor >= FLAGS_key_num) { + return std::nullopt; + } + + const uint32_t idx = runtime.cursor++; + if (FLAGS_random_access) { + return runtime.shuffled_indices[idx]; + } + return idx; +} + +void PrepareWorkerKeys(std::vector> &workers) { + for (auto &worker_ptr : workers) { + auto &worker = *worker_ptr; + worker.keys.reserve(FLAGS_key_num); + worker.shuffled_indices.reserve(FLAGS_key_num); + worker.rng.seed((static_cast(worker.tid) << 32U) ^ FLAGS_seed); + for (uint32_t i = 0; i < FLAGS_key_num; ++i) { + worker.keys.emplace_back(BuildKey(worker.tid, i, FLAGS_key_size)); + worker.shuffled_indices.emplace_back(i); + } + if (FLAGS_random_access) { + std::shuffle(worker.shuffled_indices.begin(), worker.shuffled_indices.end(), worker.rng); + } + } +} + +void PrepareDataBuffer(simm::clnt::Data &data, uint64_t seed) { + FillBufferFromSeed(data.AsRef().subspan(0, FLAGS_val_size), seed); +} + +bool VerifyDataBuffer(const simm::clnt::Data &data, uint64_t seed, int32_t ret_len) { + if (ret_len < 0) { + return false; + } + if (static_cast(ret_len) != FLAGS_val_size) { + return false; + } + return VerifyBufferFromSeed(data.AsRef().subspan(0, FLAGS_val_size), seed); +} + +struct PrefillFailure { + uint32_t worker_tid{0}; + uint32_t key_idx{0}; + int32_t rc{0}; +}; + +class ScopedPrefillClientFlags { + public: + ScopedPrefillClientFlags() + : old_enable_retry_(FLAGS_clnt_syncreq_enable_retry), + old_retry_count_(FLAGS_clnt_syncreq_retry_count), + old_timeout_ms_(FLAGS_clnt_sync_req_timeout_ms) { + FLAGS_clnt_syncreq_enable_retry = FLAGS_prefill_syncreq_enable_retry; + FLAGS_clnt_syncreq_retry_count = FLAGS_prefill_syncreq_retry_count; + FLAGS_clnt_sync_req_timeout_ms = FLAGS_prefill_sync_req_timeout_ms; + } + + ~ScopedPrefillClientFlags() { + FLAGS_clnt_syncreq_enable_retry = old_enable_retry_; + FLAGS_clnt_syncreq_retry_count = old_retry_count_; + FLAGS_clnt_sync_req_timeout_ms = old_timeout_ms_; + } + + private: + const bool old_enable_retry_; + const uint32_t old_retry_count_; + const int32_t old_timeout_ms_; +}; + +void PrefillWorkerDataset(simm::clnt::KVStore &kvstore, + WorkerRuntime &worker, + WorkloadMode mode, + std::atomic &failed, + std::mutex &failure_mu, + std::optional &first_failure) { + if (!NeedsPrefill(mode)) { + return; + } + + const uint32_t prefill_count = PrepareKeyCount(mode); + simm::clnt::Data put_data = kvstore.Allocate(FLAGS_val_size); + simm::clnt::DataView put_view(put_data); + for (uint32_t idx = 0; idx < prefill_count && !failed.load(std::memory_order_relaxed); ++idx) { + PrepareDataBuffer(put_data, ValueSeed(worker.tid, idx)); + const auto rc = kvstore.Put(worker.keys[idx], put_view); + if (rc != CommonErr::OK) { + failed.store(true, std::memory_order_relaxed); + std::lock_guard g(failure_mu); + if (!first_failure.has_value()) { + first_failure = PrefillFailure{worker.tid, idx, rc}; + } + return; + } + } +} + +bool PrefillDataset(simm::clnt::KVStore &kvstore, + std::vector> &workers, + WorkloadMode mode) { + if (!NeedsPrefill(mode)) { + return true; + } + ScopedPrefillClientFlags scoped_prefill_flags; + std::vector threads; + threads.reserve(workers.size()); + std::atomic failed{false}; + std::mutex failure_mu; + std::optional first_failure; + for (auto &worker_ptr : workers) { + auto *worker = worker_ptr.get(); + threads.emplace_back([&kvstore, worker, mode, &failed, &failure_mu, &first_failure] { + PrefillWorkerDataset(kvstore, *worker, mode, failed, failure_mu, first_failure); + }); + } + for (auto &t : threads) { + t.join(); + } + if (first_failure.has_value()) { + std::cerr << "[simm_kvio_bench_v2] prefill failed for worker " << first_failure->worker_tid + << " key_idx=" << first_failure->key_idx << " rc=" << first_failure->rc << "\n"; + return false; + } + return true; +} + +void RecordSyncSingle(ThreadStats &stats, + OpType op, + int32_t rc, + micro_ts latency, + uint64_t bytes, + bool data_ok = true) { + const bool success = (op == OpType::kGet) ? (rc >= 0 && data_ok) : (rc == CommonErr::OK); + stats.record_request(op, success, 1, success ? 1 : 0, success ? bytes : 0, latency, {rc}); +} + +void RecordSyncBatch(ThreadStats &stats, + OpType op, + const std::vector &ret_codes, + micro_ts latency, + uint64_t bytes_per_item) { + uint32_t succ = 0; + std::vector errs; + errs.reserve(ret_codes.size()); + for (auto rc : ret_codes) { + const bool ok = (op == OpType::kMGet) ? (rc >= 0) : (rc == CommonErr::OK); + succ += ok ? 1U : 0U; + if (!ok) { + errs.push_back(rc); + } + } + stats.record_request(op, + succ == ret_codes.size(), + static_cast(ret_codes.size()), + succ, + succ * bytes_per_item, + latency, + errs); +} + +void RunSyncWorker(simm::clnt::KVStore &kvstore, + WorkerRuntime &runtime, + WorkloadMode workload_mode, + steady_clock_t::time_point deadline) { + const uint32_t working_key_count = PrepareKeyCount(workload_mode); + simm::clnt::Data single_data = kvstore.Allocate(FLAGS_val_size); + simm::clnt::DataView single_view(single_data); + + std::vector batch_data_store; + std::vector batch_views; + std::vector batch_keys; + if (IsBatchWorkload(workload_mode)) { + batch_data_store.reserve(FLAGS_batchsize); + batch_views.reserve(FLAGS_batchsize); + batch_keys.reserve(FLAGS_batchsize); + for (uint32_t i = 0; i < FLAGS_batchsize; ++i) { + batch_data_store.emplace_back(kvstore.Allocate(FLAGS_val_size)); + } + } + + while (steady_clock_t::now() < deadline) { + if (workload_mode == WorkloadMode::kPut) { + const auto key_idx = NextWriteKeyIndex(runtime); + if (!key_idx.has_value()) { + break; + } + PrepareDataBuffer(single_data, ValueSeed(runtime.tid, *key_idx)); + const auto start = steady_clock_t::now(); + const auto rc = kvstore.Put(runtime.keys[*key_idx], single_view); + const auto latency = std::chrono::duration_cast(steady_clock_t::now() - start); + RecordSyncSingle(runtime.stats, OpType::kPut, rc, latency, FLAGS_val_size); + continue; + } + + if (workload_mode == WorkloadMode::kGet) { + const auto key_idx = NextKeyIndex(runtime, working_key_count); + const auto start = steady_clock_t::now(); + const auto rc = kvstore.Get(runtime.keys[key_idx], single_view); + const auto latency = std::chrono::duration_cast(steady_clock_t::now() - start); + const bool verified = !FLAGS_verify_get || VerifyDataBuffer(single_data, ValueSeed(runtime.tid, key_idx), rc); + RecordSyncSingle(runtime.stats, OpType::kGet, verified ? rc : CommonErr::InternalError, latency, FLAGS_val_size, verified); + continue; + } + + if (workload_mode == WorkloadMode::kExists) { + const auto key_idx = NextKeyIndex(runtime, working_key_count); + const auto start = steady_clock_t::now(); + const auto rc = kvstore.Exists(runtime.keys[key_idx]); + const auto latency = std::chrono::duration_cast(steady_clock_t::now() - start); + RecordSyncSingle(runtime.stats, OpType::kExists, rc, latency, 0); + continue; + } + + if (workload_mode == WorkloadMode::kDelete) { + const auto key_idx = NextKeyIndex(runtime, working_key_count); + const auto start = steady_clock_t::now(); + const auto rc = kvstore.Delete(runtime.keys[key_idx]); + const auto latency = std::chrono::duration_cast(steady_clock_t::now() - start); + RecordSyncSingle(runtime.stats, OpType::kDelete, rc, latency, 0); + if (rc == CommonErr::OK) { + PrepareDataBuffer(single_data, ValueSeed(runtime.tid, key_idx)); + const auto refill_rc = kvstore.Put(runtime.keys[key_idx], single_view); + if (refill_rc != CommonErr::OK) { + runtime.stats.record_submit_fail(OpType::kPut, refill_rc); + } + } + continue; + } + + batch_keys.clear(); + batch_views.clear(); + std::vector exists_keys; + exists_keys.reserve(FLAGS_batchsize); + std::vector batch_key_indices; + batch_key_indices.reserve(FLAGS_batchsize); + for (uint32_t i = 0; i < FLAGS_batchsize; ++i) { + std::optional key_idx; + if (workload_mode == WorkloadMode::kMPut) { + key_idx = NextWriteKeyIndex(runtime); + } else { + key_idx = NextKeyIndex(runtime, working_key_count); + } + if (!key_idx.has_value()) { + batch_key_indices.clear(); + break; + } + batch_key_indices.emplace_back(*key_idx); + if (workload_mode == WorkloadMode::kMPut) { + PrepareDataBuffer(batch_data_store[i], ValueSeed(runtime.tid, *key_idx)); + batch_keys.emplace_back(runtime.keys[*key_idx]); + batch_views.emplace_back(batch_data_store[i]); + } else if (workload_mode == WorkloadMode::kMGet) { + batch_keys.emplace_back(runtime.keys[*key_idx]); + batch_views.emplace_back(batch_data_store[i]); + } else { + exists_keys.emplace_back(runtime.keys[*key_idx]); + } + } + if (batch_key_indices.empty()) { + break; + } + + const auto start = steady_clock_t::now(); + if (workload_mode == WorkloadMode::kMPut) { + const auto rets = kvstore.MPut(batch_keys, batch_views); + const auto latency = std::chrono::duration_cast(steady_clock_t::now() - start); + std::vector rc32(rets.begin(), rets.end()); + RecordSyncBatch(runtime.stats, OpType::kMPut, rc32, latency, FLAGS_val_size); + continue; + } + + if (workload_mode == WorkloadMode::kMGet) { + auto rets = kvstore.MGet(batch_keys, batch_views); + const auto latency = std::chrono::duration_cast(steady_clock_t::now() - start); + if (FLAGS_verify_get) { + for (uint32_t i = 0; i < rets.size(); ++i) { + if (rets[i] > 0) { + const bool verified = + VerifyBufferFromSeed(batch_data_store[i].AsRef().subspan(0, FLAGS_val_size), + ValueSeed(runtime.tid, batch_key_indices[i])); + if (!verified) { + rets[i] = CommonErr::InternalError; + } + } + } + } + RecordSyncBatch(runtime.stats, OpType::kMGet, rets, latency, FLAGS_val_size); + continue; + } + + const auto rets = kvstore.MExists(exists_keys); + const auto latency = std::chrono::duration_cast(steady_clock_t::now() - start); + std::vector rc32(rets.begin(), rets.end()); + RecordSyncBatch(runtime.stats, OpType::kMExists, rc32, latency, 0); + } +} + +struct AsyncWorkerContext { + simm::clnt::KVStore *kvstore; + WorkerRuntime *runtime; + WorkloadMode workload_mode; + uint32_t working_key_count; + steady_clock_t::time_point deadline; + AsyncCreditLimiter *limiter; + std::mutex mu; + std::condition_variable cv; + size_t inflight_requests{0}; + std::atomic stop{false}; +}; + +class AsyncCreditLimiter { + public: + explicit AsyncCreditLimiter(size_t credits) : capacity_(credits), available_(credits) {} + + bool acquire(size_t credits, std::atomic &stop, steady_clock_t::time_point deadline) { + std::unique_lock lk(mu_); + cv_.wait_until(lk, deadline, [&] { + return stop.load(std::memory_order_relaxed) || available_ >= credits; + }); + if (stop.load(std::memory_order_relaxed) || available_ < credits) { + return false; + } + available_ -= credits; + return true; + } + + void release(size_t credits) { + { + std::lock_guard lk(mu_); + available_ = std::min(capacity_, available_ + credits); + } + cv_.notify_all(); + } + + private: + const size_t capacity_; + size_t available_; + std::mutex mu_; + std::condition_variable cv_; +}; + +struct AsyncBatchTracker { + AsyncBatchTracker(AsyncWorkerContext *ctx_in, OpType op_in, uint32_t total_in, steady_clock_t::time_point start_in) + : ctx(ctx_in), op(op_in), total(total_in), start(start_in) {} + + AsyncWorkerContext *ctx; + OpType op; + uint32_t total; + steady_clock_t::time_point start; + std::atomic remaining; + std::atomic succ_items; + std::atomic request_ok; + std::mutex err_mu; + std::vector error_codes; +}; + +struct AsyncSlot { + explicit AsyncSlot(simm::clnt::KVStore &kvstore) : data(kvstore.Allocate(FLAGS_val_size)), view(data) {} + + simm::clnt::Data data; + simm::clnt::DataView view; + uint32_t key_idx{0}; + steady_clock_t::time_point start{}; +}; + +void FinishAsyncBatch(const std::shared_ptr &batch, uint64_t bytes_per_item) { + const auto latency = std::chrono::duration_cast(steady_clock_t::now() - batch->start); + std::vector errs; + { + std::lock_guard g(batch->err_mu); + errs = batch->error_codes; + } + batch->ctx->runtime->stats.record_request(batch->op, + batch->request_ok.load(std::memory_order_relaxed), + batch->total, + batch->succ_items.load(std::memory_order_relaxed), + batch->succ_items.load(std::memory_order_relaxed) * bytes_per_item, + latency, + errs); + { + std::lock_guard g(batch->ctx->mu); + batch->ctx->inflight_requests -= batch->total; + } + batch->ctx->limiter->release(batch->total); + batch->ctx->cv.notify_one(); +} + +bool IsAsyncBackpressure(int32_t rc) { + return rc == sicl::transport::SICL_ERR_CH_TX_FULL; +} + +template +int16_t SubmitWithBackpressureRetry(AsyncWorkerContext &ctx, OpType op, SubmitFn &&submit_fn) { + int16_t rc = CommonErr::OK; + for (uint32_t attempt = 0; attempt < FLAGS_async_backpressure_retry; ++attempt) { + rc = submit_fn(); + if (!IsAsyncBackpressure(rc)) { + return rc; + } + ctx.runtime->stats.record_submit_backpressure(); + if (steady_clock_t::now() >= ctx.deadline) { + return rc; + } + const uint32_t backoff_us = + std::min(kMaxAsyncBackoffUs, FLAGS_async_backpressure_backoff_us * (attempt + 1)); + std::this_thread::sleep_for(std::chrono::microseconds(backoff_us)); + } + ctx.runtime->stats.record_submit_fail(op, rc); + return rc; +} + +void SubmitAsyncSingle(AsyncWorkerContext &ctx, OpType op, uint32_t key_idx) { + auto slot = std::make_shared(*ctx.kvstore); + slot->key_idx = key_idx; + slot->start = steady_clock_t::now(); + auto key = ctx.runtime->keys[key_idx]; + + auto finish_single = [slot, &ctx, op](int32_t rc, bool verified, uint64_t bytes) { + const auto latency = std::chrono::duration_cast(steady_clock_t::now() - slot->start); + const int32_t final_rc = verified ? rc : CommonErr::InternalError; + const bool success = (op == OpType::kGet) ? (final_rc >= 0) : (final_rc == CommonErr::OK); + ctx.runtime->stats.record_request(op, + success, + 1, + success ? 1 : 0, + success ? bytes : 0, + latency, + {final_rc}); + { + std::lock_guard g(ctx.mu); + ctx.inflight_requests--; + } + ctx.limiter->release(1); + ctx.cv.notify_one(); + }; + + int16_t submit_rc = CommonErr::OK; + if (op == OpType::kPut) { + PrepareDataBuffer(slot->data, ValueSeed(ctx.runtime->tid, key_idx)); + submit_rc = SubmitWithBackpressureRetry( + ctx, op, [&]() { + return ctx.kvstore->AsyncPut( + key, slot->view, [finish_single](int result) { finish_single(result, true, FLAGS_val_size); }); + }); + } else if (op == OpType::kGet) { + submit_rc = SubmitWithBackpressureRetry( + ctx, op, [&]() { + return ctx.kvstore->AsyncGet( + key, slot->view, [finish_single, slot, &ctx](int result) { + const bool verified = + !FLAGS_verify_get || VerifyDataBuffer(slot->data, ValueSeed(ctx.runtime->tid, slot->key_idx), result); + finish_single(result, verified, FLAGS_val_size); + }); + }); + } else if (op == OpType::kExists) { + submit_rc = SubmitWithBackpressureRetry( + ctx, op, [&]() { + return ctx.kvstore->AsyncExists(key, [finish_single](int result) { finish_single(result, true, 0); }); + }); + } + + if (submit_rc != CommonErr::OK) { + { + std::lock_guard g(ctx.mu); + ctx.inflight_requests--; + } + ctx.limiter->release(1); + if (!IsAsyncBackpressure(submit_rc)) { + ctx.runtime->stats.record_submit_fail(op, submit_rc); + } + ctx.cv.notify_one(); + } +} + +void SubmitAsyncBatchWave(AsyncWorkerContext &ctx, WorkloadMode workload_mode) { + const OpType op = BatchOpType(workload_mode); + auto batch = std::make_shared(&ctx, op, FLAGS_batchsize, steady_clock_t::now()); + batch->remaining.store(FLAGS_batchsize, std::memory_order_relaxed); + batch->succ_items.store(0, std::memory_order_relaxed); + batch->request_ok.store(true, std::memory_order_relaxed); + + for (uint32_t i = 0; i < FLAGS_batchsize; ++i) { + auto slot = std::make_shared(*ctx.kvstore); + std::optional key_idx; + if (workload_mode == WorkloadMode::kMPut) { + key_idx = NextWriteKeyIndex(*ctx.runtime); + } else { + key_idx = NextKeyIndex(*ctx.runtime, ctx.working_key_count); + } + if (!key_idx.has_value()) { + const uint32_t submitted = i; + const uint32_t unsent = FLAGS_batchsize - submitted; + batch->total = submitted; + batch->remaining.store(submitted, std::memory_order_relaxed); + if (unsent > 0) { + { + std::lock_guard g(ctx.mu); + ctx.inflight_requests -= unsent; + } + ctx.limiter->release(unsent); + ctx.cv.notify_one(); + } + if (submitted == 0) { + FinishAsyncBatch(batch, workload_mode == WorkloadMode::kMExists ? 0 : FLAGS_val_size); + } + return; + } + slot->key_idx = *key_idx; + slot->start = batch->start; + auto key = ctx.runtime->keys[*key_idx]; + + auto on_done = [slot, batch, &ctx](int32_t rc, bool verified, uint64_t bytes) { + const int32_t final_rc = verified ? rc : CommonErr::InternalError; + const bool ok = (batch->op == OpType::kMGet) ? (final_rc >= 0) : (final_rc == CommonErr::OK); + if (ok) { + batch->succ_items.fetch_add(1, std::memory_order_relaxed); + } else { + batch->request_ok.store(false, std::memory_order_relaxed); + std::lock_guard g(batch->err_mu); + batch->error_codes.push_back(final_rc); + } + if (batch->remaining.fetch_sub(1, std::memory_order_acq_rel) == 1) { + FinishAsyncBatch(batch, bytes); + } + }; + + int16_t submit_rc = CommonErr::OK; + if (workload_mode == WorkloadMode::kMPut) { + PrepareDataBuffer(slot->data, ValueSeed(ctx.runtime->tid, *key_idx)); + submit_rc = SubmitWithBackpressureRetry( + ctx, op, [&]() { + return ctx.kvstore->AsyncPut( + key, slot->view, [on_done](int result) { on_done(result, true, FLAGS_val_size); }); + }); + } else if (workload_mode == WorkloadMode::kMGet) { + submit_rc = SubmitWithBackpressureRetry( + ctx, op, [&]() { + return ctx.kvstore->AsyncGet( + key, slot->view, [slot, on_done, &ctx](int result) { + const bool verified = + !FLAGS_verify_get || VerifyDataBuffer(slot->data, ValueSeed(ctx.runtime->tid, slot->key_idx), result); + on_done(result, verified, FLAGS_val_size); + }); + }); + } else { + submit_rc = SubmitWithBackpressureRetry( + ctx, op, [&]() { return ctx.kvstore->AsyncExists(key, [on_done](int result) { on_done(result, true, 0); }); }); + } + + if (submit_rc != CommonErr::OK) { + batch->request_ok.store(false, std::memory_order_relaxed); + if (!IsAsyncBackpressure(submit_rc)) { + { + std::lock_guard g(batch->err_mu); + batch->error_codes.push_back(submit_rc); + } + ctx.runtime->stats.mark_submit_error(submit_rc); + } + if (batch->remaining.fetch_sub(1, std::memory_order_acq_rel) == 1) { + FinishAsyncBatch(batch, workload_mode == WorkloadMode::kMExists ? 0 : FLAGS_val_size); + } + } + } +} + +void RunAsyncWorker(AsyncWorkerContext &ctx, steady_clock_t::time_point deadline) { + while (steady_clock_t::now() < deadline) { + const size_t wave = IsBatchWorkload(ctx.workload_mode) ? FLAGS_batchsize : 1; + { + std::unique_lock lk(ctx.mu); + ctx.cv.wait(lk, [&] { + return ctx.stop.load(std::memory_order_relaxed) || (ctx.inflight_requests + wave <= FLAGS_iodepth); + }); + if (ctx.stop.load(std::memory_order_relaxed)) { + break; + } + } + // Acquire global credits outside ctx.mu to avoid deadlock: completion + // callbacks need ctx.mu to decrement inflight_requests, so we must not + // hold ctx.mu while blocking here. + if (!ctx.limiter->acquire(wave, ctx.stop, deadline)) { + break; + } + { + std::lock_guard g(ctx.mu); + ctx.inflight_requests += wave; + } + + if (IsBatchWorkload(ctx.workload_mode)) { + SubmitAsyncBatchWave(ctx, ctx.workload_mode); + continue; + } + + if (ctx.workload_mode == WorkloadMode::kPut) { + const auto key_idx = NextWriteKeyIndex(*ctx.runtime); + if (!key_idx.has_value()) { + // Key space exhausted — undo the credit and inflight acquired above + // before breaking, otherwise the drain phase hangs forever. + { + std::lock_guard g(ctx.mu); + ctx.inflight_requests -= wave; + } + ctx.limiter->release(wave); + ctx.cv.notify_one(); + break; + } + SubmitAsyncSingle(ctx, OpType::kPut, *key_idx); + continue; + } + if (ctx.workload_mode == WorkloadMode::kGet) { + SubmitAsyncSingle(ctx, OpType::kGet, NextKeyIndex(*ctx.runtime, ctx.working_key_count)); + continue; + } + SubmitAsyncSingle(ctx, OpType::kExists, NextKeyIndex(*ctx.runtime, ctx.working_key_count)); + } + + { + std::unique_lock lk(ctx.mu); + ctx.stop.store(true, std::memory_order_relaxed); + ctx.cv.wait(lk, [&] { return ctx.inflight_requests == 0; }); + } +} + +std::vector SnapshotAll(const std::vector> &workers) { + std::vector snaps; + snaps.reserve(workers.size()); + for (const auto &worker_ptr : workers) { + snaps.emplace_back(worker_ptr->stats.snapshot()); + } + return snaps; +} + +void PrintTopErrors(const std::unordered_map &err_counts) { + if (err_counts.empty()) { + return; + } + std::vector> entries(err_counts.begin(), err_counts.end()); + std::sort(entries.begin(), entries.end(), [](const auto &lhs, const auto &rhs) { return lhs.second > rhs.second; }); + std::cout << "TopErrors : "; + for (size_t i = 0; i < std::min(entries.size(), 5); ++i) { + if (i != 0) { + std::cout << ", "; + } + std::cout << entries[i].first << "=" << entries[i].second; + } + std::cout << "\n"; +} + +void PrintReport(const BenchSnapshot &snap, uint32_t threads, uint64_t elapsed_secs, bool delta_report) { + const double req_qps = elapsed_secs == 0 ? 0.0 : static_cast(snap.req_total) / static_cast(elapsed_secs); + const double item_qps = elapsed_secs == 0 ? 0.0 : static_cast(snap.item_total) / static_cast(elapsed_secs); + const double throughput_mib = + elapsed_secs == 0 ? 0.0 + : static_cast(snap.bytes_write + snap.bytes_read) / static_cast(ONE_MB) / + static_cast(elapsed_secs); + const double avg_us = snap.latency.samples == 0 + ? 0.0 + : static_cast(snap.latency.sum_us) / static_cast(snap.latency.samples); + + auto maybe_red = [&](const std::string &label, uint64_t value) { + return folly::to(Red(label), ": ", Red(folly::to(value))); + }; + auto maybe_green = [&](const std::string &label, uint64_t value) { + return folly::to(Green(label), ": ", Green(folly::to(value))); + }; + + std::cout << "++++++++++++++++++++++++++++ " + << (delta_report ? Blue("KVIO BENCH V2 DELTA") : Blue("KVIO BENCH V2 REPORT")) + << " +++++++++++++++++++++++++\n" + << "Timestamp : " << NowTs() << "\n" + << "Mode : " << FLAGS_mode << "\n" + << "IOMode : " << FLAGS_io_mode << "\n" + << "Threads : " << threads << "\n" + << "ElapsedSecs : " << elapsed_secs << "\n" + << "KeySize : " << FLAGS_key_size << " B\n" + << "ValSize : " << FLAGS_val_size << " B\n" + << "BatchSize : " << FLAGS_batchsize << "\n" + << "IODepth : " << FLAGS_iodepth << "\n" + << "ReqTotal : " << snap.req_total << "\n" + << maybe_green("ReqSuccess ", snap.req_success) << "\n" + << maybe_red("ReqFail ", snap.req_fail) << "\n" + << Yellow("SubmitBackpressure") << ": " << Yellow(folly::to(snap.submit_backpressure)) << "\n" + << maybe_red("SubmitFail ", snap.submit_fail) << "\n" + << "ItemTotal : " << snap.item_total << "\n" + << maybe_green("ItemSuccess ", snap.item_success) << "\n" + << maybe_red("ItemFail ", snap.item_fail) << "\n" + << "PutCnt : " << snap.put_count << "\n" + << "GetCnt : " << snap.get_count << "\n" + << "ExistsCnt : " << snap.exists_count << "\n" + << "DeleteCnt : " << snap.delete_count << "\n" + << "MPutCnt : " << snap.mput_count << "\n" + << "MGetCnt : " << snap.mget_count << "\n" + << "MExistsCnt : " << snap.mexists_count << "\n" + << "WriteBytes : " << HumanBytes(snap.bytes_write) << "\n" + << "ReadBytes : " << HumanBytes(snap.bytes_read) << "\n" + << "ReqQPS : " << std::fixed << std::setprecision(2) << req_qps << "\n" + << "ItemQPS : " << item_qps << "\n" + << "ThroughputMiB/s : " << throughput_mib << "\n" + << "AvgLatencyUs : " << avg_us << "\n" + << "P50LatencyUs : " << LatencyHistogram::Percentile(snap.latency, 0.50) << "\n" + << "P95LatencyUs : " << LatencyHistogram::Percentile(snap.latency, 0.95) << "\n" + << "P99LatencyUs : " << LatencyHistogram::Percentile(snap.latency, 0.99) << "\n" + << "P999LatencyUs : " << LatencyHistogram::Percentile(snap.latency, 0.999) << "\n" + << "MaxLatencyUs : " << snap.latency.max_us << "\n"; + PrintTopErrors(snap.err_counts); + if (ParseIoMode() == IoMode::kAsync && IsBatchWorkload(ParseWorkloadMode())) { + std::cout << Yellow("Note : async batch mode is implemented as async wave submission of single-key APIs\n"); + } + std::cout << std::endl; +} + +void RunMeasurePhase(simm::clnt::KVStore &kvstore, + std::vector> &workers, + IoMode io_mode, + WorkloadMode workload_mode) { + std::vector threads; + threads.reserve(workers.size()); + const auto measure_start = steady_clock_t::now(); + const auto deadline = measure_start + std::chrono::seconds(FLAGS_run_time); + auto limiter = std::make_shared(FLAGS_async_global_iodepth); + + if (io_mode == IoMode::kSync) { + for (auto &worker_ptr : workers) { + threads.emplace_back([&kvstore, &worker_ptr, workload_mode, deadline] { + RunSyncWorker(kvstore, *worker_ptr, workload_mode, deadline); + }); + } + } else { + for (auto &worker_ptr : workers) { + threads.emplace_back([&kvstore, &worker_ptr, workload_mode, deadline, limiter] { + AsyncWorkerContext ctx{ + &kvstore, worker_ptr.get(), workload_mode, PrepareKeyCount(workload_mode), deadline, limiter.get(), {}, {}, 0, false}; + RunAsyncWorker(ctx, deadline); + }); + } + } + + if (!FLAGS_summary) { + auto previous = MergeSnapshots(SnapshotAll(workers)); + auto last_report = steady_clock_t::now(); + while (steady_clock_t::now() < deadline) { + std::this_thread::sleep_for(std::chrono::seconds(FLAGS_report_interval)); + const auto current = MergeSnapshots(SnapshotAll(workers)); + const auto delta = current - previous; + const auto now = steady_clock_t::now(); + const auto elapsed = std::max( + 1, std::chrono::duration_cast(now - last_report).count()); + PrintReport(delta, FLAGS_threads, elapsed, true); + previous = current; + last_report = now; + } + } + + for (auto &t : threads) { + t.join(); + } + + if (FLAGS_summary) { + const auto actual_elapsed = std::max( + 1, std::chrono::duration_cast(steady_clock_t::now() - measure_start).count()); + PrintReport(MergeSnapshots(SnapshotAll(workers)), FLAGS_threads, actual_elapsed, false); + } +} + +void RunWarmupPhase(simm::clnt::KVStore &kvstore, + std::vector> &workers, + IoMode io_mode, + WorkloadMode workload_mode) { + if (FLAGS_warmup_time == 0) { + return; + } + std::cout << Yellow("[simm_kvio_bench_v2] starting warmup phase") << std::endl; + std::vector threads; + threads.reserve(workers.size()); + const auto deadline = steady_clock_t::now() + std::chrono::seconds(FLAGS_warmup_time); + auto limiter = std::make_shared(FLAGS_async_global_iodepth); + if (io_mode == IoMode::kSync) { + for (auto &worker_ptr : workers) { + threads.emplace_back([&kvstore, &worker_ptr, workload_mode, deadline] { + RunSyncWorker(kvstore, *worker_ptr, workload_mode, deadline); + }); + } + } else { + for (auto &worker_ptr : workers) { + threads.emplace_back([&kvstore, &worker_ptr, workload_mode, deadline, limiter] { + AsyncWorkerContext ctx{ + &kvstore, worker_ptr.get(), workload_mode, PrepareKeyCount(workload_mode), deadline, limiter.get(), {}, {}, 0, false}; + RunAsyncWorker(ctx, deadline); + }); + } + } + for (auto &t : threads) { + t.join(); + } + for (auto &worker_ptr : workers) { + worker_ptr->stats.reset(); + if (FLAGS_allow_overwrite || !IsWriteWorkload(workload_mode)) { + worker_ptr->cursor = 0; + } + } +} + +} // namespace + +int main(int argc, char **argv) { + gflags::SetUsageMessage( + "simm_kvio_bench_v2 --cm_primary_node_ip= --mode=get --io_mode=sync " + "--threads=32 --run_time=300 --warmup_time=30 --key_size=20 --val_size=4096"); + gflags::ParseCommandLineFlags(&argc, &argv, true); + folly::Init init(&argc, &argv); + +#ifdef NDEBUG + simm::logging::LogConfig clnt_log_config = simm::logging::LogConfig{FLAGS_clnt_log_file, "INFO"}; +#else + simm::logging::LogConfig clnt_log_config = simm::logging::LogConfig{FLAGS_clnt_log_file, "DEBUG"}; +#endif + simm::logging::LoggerManager::Instance().UpdateConfig("simm_client", clnt_log_config); + + const auto io_mode = ParseIoMode(); + const auto workload_mode = ParseWorkloadMode(); + ValidateArgs(io_mode, workload_mode); + + std::cout << Blue("== SIMM KVIO BENCH V2 ==") << "\n" + << "mode=" << FLAGS_mode << ", io_mode=" << FLAGS_io_mode << ", threads=" << FLAGS_threads + << ", key_size=" << FLAGS_key_size << ", val_size=" << FLAGS_val_size << ", batchsize=" << FLAGS_batchsize + << ", iodepth=" << FLAGS_iodepth << ", async_global_iodepth=" << FLAGS_async_global_iodepth + << ", allow_overwrite=" << (FLAGS_allow_overwrite ? "true" : "false") + << ", warmup=" << FLAGS_warmup_time << "s" + << ", run=" << FLAGS_run_time << "s\n"; + + auto kvstore = std::make_unique(); + if (kvstore == nullptr) { + std::cerr << "[simm_kvio_bench_v2] failed to create KVStore\n"; + return ENOMEM; + } + + std::vector> workers; + workers.reserve(FLAGS_threads); + for (uint32_t i = 0; i < FLAGS_threads; ++i) { + workers.emplace_back(std::make_unique(i)); + } + + PrepareWorkerKeys(workers); + if (!PrefillDataset(*kvstore, workers, workload_mode)) { + return EIO; + } + RunWarmupPhase(*kvstore, workers, io_mode, workload_mode); + RunMeasurePhase(*kvstore, workers, io_mode, workload_mode); + return 0; +} diff --git a/tools/simm_stable_test.cc b/tools/simm_stable_test.cc index 6c43d39..b53e2e5 100644 --- a/tools/simm_stable_test.cc +++ b/tools/simm_stable_test.cc @@ -501,8 +501,15 @@ void print_stats(const ThreadStats &s, uint32_t threads, uint64_t elapsed_secs, 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; }; + // Color value red only when non-zero (actual failure); green for success counts + auto fail_val = [&](auto value) -> std::string { + auto s = folly::to(value); + return value != 0 ? std::string(kRed) + s + kReset : s; + }; + auto succ_val = [&](auto value) -> std::string { + auto s = folly::to(value); + return value != 0 ? std::string(kGreen) + s + kReset : s; + }; 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 @@ -514,31 +521,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" - << red_label("Failures ") << ": " << s.total_failures() << "\n" - << red_label("SubmitFails ") << ": " << s.submit_fails_ << "\n" + << "Failures : " << fail_val(s.total_failures()) << "\n" + << "SubmitFails : " << fail_val(s.submit_fails_) << "\n" << "PutCnt : " << s.put_ << "\n" - << red_label("PutFails ") << ": " << s.put_fails_ << "\n" - << green_label("PutSuccs ") << ": " << s.put_succs_ << "\n" + << "PutFails : " << fail_val(s.put_fails_) << "\n" + << "PutSuccs : " << succ_val(s.put_succs_) << "\n" << "OverwriteCnt : " << s.overwrite_put_ << "\n" - << red_label("OverwriteFails ") << ": " << s.overwrite_put_fails_ << "\n" - << green_label("OverwriteSuccs ") << ": " << s.overwrite_put_succs_ << "\n" + << "OverwriteFails : " << fail_val(s.overwrite_put_fails_) << "\n" + << "OverwriteSuccs : " << succ_val(s.overwrite_put_succs_) << "\n" << "GetCnt : " << s.get_ << "\n" - << red_label("GetFails ") << ": " << s.get_fails_ << "\n" - << green_label("GetSuccs ") << ": " << s.get_succs_ << "\n" + << "GetFails : " << fail_val(s.get_fails_) << "\n" + << "GetSuccs : " << succ_val(s.get_succs_) << "\n" << "ExistsCnt : " << s.exists_ << "\n" - << red_label("ExistsFails ") << ": " << s.exists_fails_ << "\n" - << green_label("ExistsSuccs ") << ": " << s.exists_succs_ << "\n" + << "ExistsFails : " << fail_val(s.exists_fails_) << "\n" + << "ExistsSuccs : " << succ_val(s.exists_succs_) << "\n" << "DeleteCnt : " << s.del_ << "\n" - << red_label("DeleteFails ") << ": " << s.del_fails_ << "\n" - << green_label("DeleteSuccs ") << ": " << s.del_succs_ << "\n" + << "DeleteFails : " << fail_val(s.del_fails_) << "\n" + << "DeleteSuccs : " << succ_val(s.del_succs_) << "\n" << "MPutCnt : " << s.mput_ << "\n" - << red_label("MPutFails ") << ": " << s.mput_fails_ << "\n" - << green_label("MPutSuccs ") << ": " << s.mput_succs_ << "\n" + << "MPutFails : " << fail_val(s.mput_fails_) << "\n" + << "MPutSuccs : " << succ_val(s.mput_succs_) << "\n" << "MGetCnt : " << s.mget_ << "\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" + << "MGetFails : " << fail_val(s.mget_fails_) << "\n" + << "MGetSuccs : " << succ_val(s.mget_succs_) << "\n" + << "DataMatch : " << succ_val(s.data_match_) << "\n" + << "DataMismatch : " << fail_val(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" diff --git a/tools/simm_stable_test_upgrade.md b/tools/simm_stable_test_upgrade.md index 35a8211..5577af9 100644 --- a/tools/simm_stable_test_upgrade.md +++ b/tools/simm_stable_test_upgrade.md @@ -1,104 +1,104 @@ -# simm_stable_test 升级说明 +# simm_stable_test Upgrade Notes -## 背景 +## Background -原始版本的 `tools/simm_stable_test.cc` 更偏功能冒烟工具,存在几个明显短板: +The original `tools/simm_stable_test.cc` was more of a functional smoke-test tool, with several notable shortcomings: -- `async` 路径基本未完成,`iodepth` 实际不生效 -- 工作负载模型过于单一,难以支撑长期稳定性验证 -- 缺少周期性统计与延迟指标,不利于长时间运行时观测 -- key 基本一次性生成,无法形成更真实的覆盖写、热点和存活数据压力 -- 停止逻辑较粗,长跑时不够稳妥 +- The `async` path was largely incomplete; `iodepth` had no real effect +- The workload model was too simplistic to support long-running stability testing +- No periodic statistics or latency metrics, making it hard to observe during extended runs +- Keys were essentially generated once, unable to produce realistic overwrite, hotspot, or live-data pressure +- Shutdown logic was too coarse for long-running scenarios -这次改动的目标,是把它提升成更适合持续压测、稳定性回归和故障观测的工具。 +The goal of this change is to elevate the tool into one better suited for sustained stress testing, stability regression, and fault observation. -## 主要改动 +## Key Changes -### 1. 补全 async 模式 +### 1. Complete the async mode -- 实现真正的异步提交与回调处理 -- `iodepth` 现在实际控制每个 worker 的并发异步请求数 -- 支持异步 `put/get/exists/delete` -- `async + batch` 组合当前显式禁止,避免半成品路径误用 +- Implement true asynchronous submission and callback handling +- `iodepth` now actually controls per-worker concurrent async request count +- Support async `put/get/exists/delete` +- `async + batch` combination is explicitly disallowed to prevent use of incomplete paths -### 2. 引入长期稳定 workload 模型 +### 2. Introduce long-running stable workload model -- 每个 worker 持有固定的逻辑 keyspace -- key 在 keyspace 内反复复用,而不是持续只写新 key -- 支持更贴近真实服务的混合流量: +- Each worker owns a fixed logical keyspace +- Keys are reused within the keyspace instead of only writing new keys +- Support mixed traffic closer to real-world services: - `putratio` - `getratio` - `existsratio` - `delratio` - `oputratio` -这样可以覆盖: +This enables coverage of: -- 覆盖写 -- 老数据读取 -- 数据删除后再次访问 -- 热点 key 反复更新 +- Overwrite scenarios +- Stale data reads +- Access after deletion +- Repeated hotspot key updates -### 3. 改进数据校验方式 +### 3. Improve data verification -- 不再为每个 key 保存整份 value 副本 -- 改为记录 `(exists, size, seed)` 元信息 -- value 内容通过 `seed` 按确定性规则生成和校验 +- No longer store a full value copy per key +- Instead track `(exists, size, seed)` metadata +- Value content is deterministically generated and verified via `seed` -收益: +Benefits: -- 降低工具自身内存开销 -- 适合更长时间运行 -- 仍可对 `get` 返回内容做一致性校验 +- Reduced memory footprint of the tool itself +- Suitable for longer runs +- Still able to verify consistency of `get` return data -### 4. 增加延迟与吞吐统计 +### 4. Add latency and throughput statistics -新增统计项: +New metrics: -- 总操作数 -- 成功/失败计数 -- submit 失败计数 -- 数据匹配/不匹配计数 -- put/get 字节量 -- 延迟统计: +- Total operation count +- Success/failure counts +- Submit failure count +- Data match/mismatch counts +- put/get byte volumes +- Latency statistics: - avg - p50 - p95 - p99 - max -### 5. 增加周期性报告 +### 5. Add periodic reporting -新增 `report_interval_inSecs`: +New `report_interval_inSecs` flag: -- 周期性打印增量统计 -- 测试结束时打印最终汇总 +- Periodically print incremental statistics +- Print final summary at test completion -这让工具更适合: +This makes the tool better suited for: -- 长跑观测 -- 问题定位 -- 性能回归对比 +- Long-running observation +- Problem diagnosis +- Performance regression comparison -### 6. 改进退出与收尾 +### 6. Improve shutdown and cleanup -- 注册 `SIGINT` / `SIGTERM` -- 支持收到停止信号后优雅退出 -- async worker 在停止时会等待 inflight 请求完成再退出 +- Register `SIGINT` / `SIGTERM` handlers +- Support graceful shutdown upon receiving stop signals +- Async workers wait for inflight requests to complete before exiting -这能减少: +This reduces: -- 工具自身异常中止 -- 长跑结束时统计不完整 -- 回调尚未收敛时直接退出的问题 +- Abnormal tool termination +- Incomplete statistics at end of long runs +- Issues from exiting before callbacks have settled -### 7. 改善可用性 +### 7. Improve usability -- 补充更完整的 flags -- `--help` 现在有正式 usage 信息 -- 启动时打印关键测试参数 +- Add more complete flags +- `--help` now shows proper usage information +- Print key test parameters at startup -## 新增/调整的重要参数 +## New/Adjusted Parameters - `putratio` - `getratio` @@ -109,33 +109,32 @@ - `report_interval_inSecs` - `strict_verify_exists` -## 当前限制 +## Current Limitations -- `async + batch_mode` 暂不支持 -- 工具仍以 client 视角做稳定性测试,不替代 server 端 profiling/资源分析 -- 延迟统计当前采用分桶估算 percentile,不是全量样本精确分位数 +- `async + batch_mode` is not yet supported +- The tool tests stability from the client perspective only; it does not replace server-side profiling/resource analysis +- Latency statistics use bucketed percentile estimation, not exact quantiles from full samples -## 构建验证 +## Build Verification -已完成以下验证: +The following verifications have been completed: -- `build/release` 下 `simm_stable_test` 构建通过 -- `build/debug` 下 `simm_stable_test` 构建通过 -- `--help` 启动 smoke test 通过 +- `simm_stable_test` builds successfully under `build/release` +- `simm_stable_test` builds successfully under `build/debug` +- `--help` startup smoke test passes -测试前统一使用: +Before testing, set: ```bash export SICL_LOG_LEVEL=WARN ``` -## 建议的后续增强 +## Suggested Future Enhancements -如果后面继续迭代,建议优先考虑: - -- 增加错误码分桶统计 -- 增加更明确的热点分布模型 -- 增加阶段性 summary 文件输出 -- 支持 fault injection 联动场景 -- 增加更细粒度的 async 超时/回调异常观测 +If further iteration is planned, consider prioritizing: +- Add error code bucketed statistics +- Add a more explicit hotspot distribution model +- Add periodic summary file output +- Support fault injection integration scenarios +- Add finer-grained async timeout/callback anomaly observation