Name server resolver (#20)
* Add StaticNameServerResolver and DynamicNameServerResolver
diff --git a/api/rocketmq/DefaultMQProducer.h b/api/rocketmq/DefaultMQProducer.h
index d3fca1b..697b0b4 100644
--- a/api/rocketmq/DefaultMQProducer.h
+++ b/api/rocketmq/DefaultMQProducer.h
@@ -79,6 +79,8 @@
void setNamesrvAddr(const std::string& name_server_address_list);
+ void setNameServerListDiscoveryEndpoint(const std::string& discovery_endpoint);
+
void setGroupName(const std::string& group_name);
void setInstanceName(const std::string& instance_name);
diff --git a/api/rocketmq/DefaultMQPullConsumer.h b/api/rocketmq/DefaultMQPullConsumer.h
index e0fd331..9b7d3ef 100644
--- a/api/rocketmq/DefaultMQPullConsumer.h
+++ b/api/rocketmq/DefaultMQPullConsumer.h
@@ -36,6 +36,8 @@
void setNamesrvAddr(const std::string& name_srv);
+ void setNameServerListDiscoveryEndpoint(const std::string& discovery_endpoint);
+
void setCredentialsProvider(std::shared_ptr<CredentialsProvider> credentials_provider);
private:
diff --git a/api/rocketmq/DefaultMQPushConsumer.h b/api/rocketmq/DefaultMQPushConsumer.h
index 20af321..176d909 100644
--- a/api/rocketmq/DefaultMQPushConsumer.h
+++ b/api/rocketmq/DefaultMQPushConsumer.h
@@ -36,6 +36,8 @@
void setNamesrvAddr(const std::string& name_srv);
+ void setNameServerListDiscoveryEndpoint(const std::string& discovery_endpoint);
+
void setGroupName(const std::string& group_name);
void setConsumeThreadCount(int thread_count);
@@ -53,8 +55,8 @@
bool isTracingEnabled();
/**
- * SDK of this version always uses asynchronous IO operation. As such, this function is no-op
- * to keep backward compatibility.
+ * SDK of this version always uses asynchronous IO operation. As such, this
+ * function is no-op to keep backward compatibility.
*/
void setAsyncPull(bool);
@@ -65,21 +67,23 @@
void setConsumeMessageBatchMaxSize(int batch_size);
/**
- * Lifecycle of executor is managed by external application. Passed-in executor should remain valid after consumer
- * start and before stopping.
+ * Lifecycle of executor is managed by external application. Passed-in
+ * executor should remain valid after consumer start and before stopping.
* @param executor Executor pool used to invoke consume callback.
*/
void setCustomExecutor(const Executor& executor);
/**
- * This function sets maximum number of message that may be consumed per second.
+ * This function sets maximum number of message that may be consumed per
+ * second.
* @param topic Topic to control
* @param threshold Threshold before throttling is enforced.
*/
void setThrottle(const std::string& topic, uint32_t threshold);
/**
- * Set abstract-resource-namespace, in which canonical name of topic, group remains unique.
+ * Set abstract-resource-namespace, in which canonical name of topic, group
+ * remains unique.
* @param resource_namespace Abstract resource namespace.
*/
void setResourceNamespace(const char* resource_namespace);
diff --git a/src/main/cpp/client/ClientManagerImpl.cpp b/src/main/cpp/client/ClientManagerImpl.cpp
index 69f2c6d..d542352 100644
--- a/src/main/cpp/client/ClientManagerImpl.cpp
+++ b/src/main/cpp/client/ClientManagerImpl.cpp
@@ -948,8 +948,6 @@
Scheduler& ClientManagerImpl::getScheduler() { return scheduler_; }
-TopAddressing& ClientManagerImpl::topAddressing() { return top_addressing_; }
-
void ClientManagerImpl::ack(const std::string& target, const Metadata& metadata, const AckMessageRequest& request,
std::chrono::milliseconds timeout, const std::function<void(bool)>& cb) {
std::string target_host(target.data(), target.length());
diff --git a/src/main/cpp/client/include/ClientManager.h b/src/main/cpp/client/include/ClientManager.h
index 154ad94..8ebb2eb 100644
--- a/src/main/cpp/client/include/ClientManager.h
+++ b/src/main/cpp/client/include/ClientManager.h
@@ -24,8 +24,6 @@
virtual Scheduler& getScheduler() = 0;
- virtual TopAddressing& topAddressing() = 0;
-
virtual std::shared_ptr<grpc::Channel> createChannel(const std::string& target_host) = 0;
virtual void resolveRoute(const std::string& target_host, const Metadata& metadata, const QueryRouteRequest& request,
diff --git a/src/main/cpp/client/include/ClientManagerImpl.h b/src/main/cpp/client/include/ClientManagerImpl.h
index f50ecc6..050ad91 100644
--- a/src/main/cpp/client/include/ClientManagerImpl.h
+++ b/src/main/cpp/client/include/ClientManagerImpl.h
@@ -130,8 +130,6 @@
Scheduler& getScheduler() override;
- TopAddressing& topAddressing() override;
-
/**
* Ack message asynchronously.
* @param target_host Target broker host address.
@@ -248,8 +246,6 @@
grpc::ChannelArguments channel_arguments_;
bool trace_{false};
-
- TopAddressing top_addressing_;
};
ROCKETMQ_NAMESPACE_END
\ No newline at end of file
diff --git a/src/main/cpp/client/mocks/include/ClientManagerMock.h b/src/main/cpp/client/mocks/include/ClientManagerMock.h
index c3c0d4c..7abce6e 100644
--- a/src/main/cpp/client/mocks/include/ClientManagerMock.h
+++ b/src/main/cpp/client/mocks/include/ClientManagerMock.h
@@ -14,8 +14,6 @@
MOCK_METHOD(Scheduler&, getScheduler, (), (override));
- MOCK_METHOD(TopAddressing&, topAddressing, (), (override));
-
MOCK_METHOD((std::shared_ptr<grpc::Channel>), createChannel, (const std::string&), (override));
MOCK_METHOD(void, resolveRoute,
diff --git a/src/main/cpp/rocketmq/ClientImpl.cpp b/src/main/cpp/rocketmq/ClientImpl.cpp
index 4bbc90f..e7806ed 100644
--- a/src/main/cpp/rocketmq/ClientImpl.cpp
+++ b/src/main/cpp/rocketmq/ClientImpl.cpp
@@ -1,6 +1,7 @@
#include <algorithm>
#include <chrono>
#include <cstdint>
+#include <cstdlib>
#include <iterator>
#include <memory>
#include <string>
@@ -30,39 +31,20 @@
return;
}
+ if (!name_server_resolver_) {
+ SPDLOG_ERROR("No name server resolver is configured.");
+ abort();
+ }
+ name_server_resolver_->start();
+
client_manager_ = ClientManagerFactory::getInstance().getClientManager(*this);
client_manager_->start();
exporter_ = std::make_shared<OtlpExporter>(client_manager_, this);
exporter_->start();
- bool update_name_server_list = false;
- {
- absl::MutexLock lock(&name_server_list_mtx_);
- if (name_server_list_.empty()) {
- update_name_server_list = true;
- }
- }
-
std::weak_ptr<ClientImpl> ptr(self());
- if (update_name_server_list) {
- // Acquire name server list immediately
- renewNameServerList();
-
- // Schedule to renew name server list periodically
- SPDLOG_INFO("Name server list was empty. Schedule a task to fetch and renew periodically");
- auto name_server_update_functor = [ptr]() {
- std::shared_ptr<ClientImpl> base = ptr.lock();
- if (base) {
- base->renewNameServerList();
- }
- };
- name_server_update_handle_ =
- client_manager_->getScheduler().schedule(name_server_update_functor, UPDATE_NAME_SERVER_LIST_TASK_NAME,
- std::chrono::milliseconds(0), std::chrono::seconds(30));
- }
-
auto route_update_functor = [ptr]() {
std::shared_ptr<ClientImpl> base = ptr.lock();
if (base) {
@@ -77,9 +59,8 @@
void ClientImpl::shutdown() {
state_.store(State::STOPPING, std::memory_order_relaxed);
- if (name_server_update_handle_) {
- client_manager_->getScheduler().cancel(name_server_update_handle_);
- }
+
+ name_server_resolver_->shutdown();
if (route_update_handle_) {
client_manager_->getScheduler().cancel(route_update_handle_);
@@ -91,7 +72,6 @@
}
const char* ClientImpl::UPDATE_ROUTE_TASK_NAME = "route_updater";
-const char* ClientImpl::UPDATE_NAME_SERVER_LIST_TASK_NAME = "name_server_list_updater";
void ClientImpl::endpointsInUse(absl::flat_hash_set<std::string>& endpoints) {
absl::MutexLock lk(&topic_route_table_mtx_);
@@ -150,81 +130,11 @@
}
}
-void ClientImpl::debugNameServerChanges(const std::vector<std::string>& list) {
- std::string previous;
- bool changed = false;
- {
- absl::MutexLock lock(&name_server_list_mtx_);
- if (name_server_list_ != list) {
- changed = true;
- if (name_server_list_.empty()) {
- previous.append("[]");
- } else {
- previous = absl::StrJoin(name_server_list_.begin(), name_server_list_.end(), ";");
- }
- }
- }
- std::string current = absl::StrJoin(list.begin(), list.end(), ";");
- if (changed) {
- SPDLOG_INFO("Name server list changed. {} --> {}", previous, current);
- } else {
- SPDLOG_DEBUG("Name server list remains the same: {}", current);
- }
-}
-
-void ClientImpl::renewNameServerList() {
- if (State::STARTED != state_.load(std::memory_order_relaxed) &&
- State::STARTING != state_.load(std::memory_order_relaxed)) {
- SPDLOG_WARN("Unexpected client instance state: {}", state_.load(std::memory_order_relaxed));
- return;
- }
-
- std::vector<std::string> list;
- SPDLOG_DEBUG("Begin to renew name server list");
- auto callback = [this](int code, const std::vector<std::string>& name_server_list) {
- if (static_cast<int>(HttpStatus::OK) != code) {
- SPDLOG_WARN("Failed to fetch name server list");
- return;
- }
-
- if (name_server_list.empty()) {
- SPDLOG_WARN("Yuck, got an empty name server list");
- return;
- }
-
- debugNameServerChanges(name_server_list);
- {
- absl::MutexLock lock(&name_server_list_mtx_);
- if (name_server_list_ != name_server_list) {
- name_server_list_.clear();
- name_server_list_.insert(name_server_list_.begin(), name_server_list.begin(), name_server_list.end());
- }
- }
- };
- client_manager_->topAddressing().fetchNameServerAddresses(callback);
-}
-
-bool ClientImpl::selectNameServer(std::string& selected, bool change) {
- static uint32_t index = 0;
- if (change) {
- index++;
- }
- {
- absl::MutexLock lock(&name_server_list_mtx_);
- if (name_server_list_.empty()) {
- return false;
- }
- uint32_t idx = index % name_server_list_.size();
- selected = name_server_list_[idx];
- }
- return true;
-}
-
void ClientImpl::setAccessPoint(rmq::Endpoints* endpoints) {
std::vector<std::pair<std::string, std::uint16_t>> pairs;
{
- absl::MutexLock lk(&name_server_list_mtx_);
- for (const auto& name_server_item : name_server_list_) {
+ std::vector<std::string> name_server_list = name_server_resolver_->resolve();
+ for (const auto& name_server_item : name_server_list) {
std::string::size_type pos = name_server_item.rfind(':');
if (std::string::npos == pos) {
continue;
@@ -254,8 +164,8 @@
}
void ClientImpl::fetchRouteFor(const std::string& topic, const std::function<void(const TopicRouteDataPtr&)>& cb) {
- std::string name_server;
- if (!selectNameServer(name_server)) {
+ std::string name_server = name_server_resolver_->current();
+ if (name_server.empty()) {
SPDLOG_WARN("No name server available");
return;
}
@@ -263,8 +173,8 @@
auto callback = [this, topic, name_server, cb](bool ok, const TopicRouteDataPtr& route) {
if (!ok || !route) {
SPDLOG_WARN("Failed to resolve route for topic={} from {}", topic, name_server);
- std::string name_server_changed;
- if (selectNameServer(name_server_changed, true)) {
+ std::string name_server_changed = name_server_resolver_->next();
+ if (!name_server_changed.empty()) {
SPDLOG_INFO("Change current name server from {} to {}", name_server, name_server_changed);
}
cb(nullptr);
diff --git a/src/main/cpp/rocketmq/DefaultMQProducer.cpp b/src/main/cpp/rocketmq/DefaultMQProducer.cpp
index cdae0c0..539cc6b 100644
--- a/src/main/cpp/rocketmq/DefaultMQProducer.cpp
+++ b/src/main/cpp/rocketmq/DefaultMQProducer.cpp
@@ -1,9 +1,15 @@
#include "rocketmq/DefaultMQProducer.h"
-#include "MixAll.h"
-#include "ProducerImpl.h"
+
+#include <chrono>
+#include <memory>
#include "absl/strings/str_split.h"
+#include "DynamicNameServerResolver.h"
+#include "MixAll.h"
+#include "ProducerImpl.h"
+#include "StaticNameServerResolver.h"
+
ROCKETMQ_NAMESPACE_BEGIN
DefaultMQProducer::DefaultMQProducer(const std::string& group_name)
@@ -26,8 +32,13 @@
}
void DefaultMQProducer::setNamesrvAddr(const std::string& name_server_address_list) {
- std::vector<std::string> name_server_list = absl::StrSplit(name_server_address_list, ';');
- impl_->setNameServerList(name_server_list);
+ auto name_server_resolver = std::make_shared<StaticNameServerResolver>(name_server_address_list);
+ impl_->withNameServerResolver(name_server_resolver);
+}
+
+void DefaultMQProducer::setNameServerListDiscoveryEndpoint(const std::string& discovery_endpoint) {
+ auto name_server_resolver = std::make_shared<DynamicNameServerResolver>(discovery_endpoint, std::chrono::seconds(10));
+ impl_->withNameServerResolver(name_server_resolver);
}
void DefaultMQProducer::setGroupName(const std::string& group_name) { impl_->setGroupName(group_name); }
diff --git a/src/main/cpp/rocketmq/DefaultMQPullConsumer.cpp b/src/main/cpp/rocketmq/DefaultMQPullConsumer.cpp
index 0410662..fd6f244 100644
--- a/src/main/cpp/rocketmq/DefaultMQPullConsumer.cpp
+++ b/src/main/cpp/rocketmq/DefaultMQPullConsumer.cpp
@@ -1,8 +1,13 @@
#include "rocketmq/DefaultMQPullConsumer.h"
-#include "AwaitPullCallback.h"
-#include "PullConsumerImpl.h"
+
#include "absl/strings/str_split.h"
+#include "AwaitPullCallback.h"
+#include "DynamicNameServerResolver.h"
+#include "PullConsumerImpl.h"
+#include "StaticNameServerResolver.h"
+#include <memory>
+
ROCKETMQ_NAMESPACE_BEGIN
DefaultMQPullConsumer::DefaultMQPullConsumer(const std::string& group_name)
@@ -28,15 +33,22 @@
impl_->pull(query, callback);
}
-void DefaultMQPullConsumer::setResourceNamespace(const std::string& resource_namespace) { impl_->resourceNamespace(resource_namespace); }
+void DefaultMQPullConsumer::setResourceNamespace(const std::string& resource_namespace) {
+ impl_->resourceNamespace(resource_namespace);
+}
void DefaultMQPullConsumer::setCredentialsProvider(std::shared_ptr<CredentialsProvider> credentials_provider) {
impl_->setCredentialsProvider(std::move(credentials_provider));
}
void DefaultMQPullConsumer::setNamesrvAddr(const std::string& name_srv) {
- std::vector<std::string> name_server_list = absl::StrSplit(name_srv, absl::ByChar(';'));
- impl_->setNameServerList(name_server_list);
+ auto name_server_resolver = std::make_shared<StaticNameServerResolver>(name_srv);
+ impl_->withNameServerResolver(name_server_resolver);
+}
+
+void DefaultMQPullConsumer::setNameServerListDiscoveryEndpoint(const std::string& discovery_endpoint) {
+ auto name_server_resolver = std::make_shared<DynamicNameServerResolver>(discovery_endpoint, std::chrono::seconds(10));
+ impl_->withNameServerResolver(name_server_resolver);
}
ROCKETMQ_NAMESPACE_END
\ No newline at end of file
diff --git a/src/main/cpp/rocketmq/DefaultMQPushConsumer.cpp b/src/main/cpp/rocketmq/DefaultMQPushConsumer.cpp
index c6f7daa..ab072fd 100644
--- a/src/main/cpp/rocketmq/DefaultMQPushConsumer.cpp
+++ b/src/main/cpp/rocketmq/DefaultMQPushConsumer.cpp
@@ -1,8 +1,14 @@
-#include "rocketmq/DefaultMQPushConsumer.h"
-#include "PushConsumerImpl.h"
-#include "absl/strings/str_split.h"
+#include <chrono>
+#include <memory>
#include <set>
+#include "absl/strings/str_split.h"
+
+#include "DynamicNameServerResolver.h"
+#include "PushConsumerImpl.h"
+#include "StaticNameServerResolver.h"
+#include "rocketmq/DefaultMQPushConsumer.h"
+
ROCKETMQ_NAMESPACE_BEGIN
static std::set<std::string> consumerTable{};
@@ -37,8 +43,17 @@
}
void DefaultMQPushConsumer::setNamesrvAddr(const std::string& name_srv) {
- std::vector<std::string> name_server_list = absl::StrSplit(name_srv, ';');
- impl_->setNameServerList(name_server_list);
+ auto name_server_resolver = std::make_shared<StaticNameServerResolver>(name_srv);
+ impl_->withNameServerResolver(name_server_resolver);
+}
+
+void DefaultMQPushConsumer::setNameServerListDiscoveryEndpoint(const std::string& discovery_endpoint) {
+ if (discovery_endpoint.empty()) {
+ return;
+ }
+
+ auto name_server_resolver = std::make_shared<DynamicNameServerResolver>(discovery_endpoint, std::chrono::seconds(10));
+ impl_->withNameServerResolver(name_server_resolver);
}
void DefaultMQPushConsumer::setGroupName(const std::string& group_name) { impl_->setGroupName(group_name); }
@@ -67,7 +82,9 @@
impl_->setThrottle(topic, threshold);
}
-void DefaultMQPushConsumer::setResourceNamespace(const char* resource_namespace) { impl_->resourceNamespace(resource_namespace); }
+void DefaultMQPushConsumer::setResourceNamespace(const char* resource_namespace) {
+ impl_->resourceNamespace(resource_namespace);
+}
void DefaultMQPushConsumer::setCredentialsProvider(CredentialsProviderPtr credentials_provider) {
impl_->setCredentialsProvider(std::move(credentials_provider));
diff --git a/src/main/cpp/rocketmq/DynamicNameServerResolver.cpp b/src/main/cpp/rocketmq/DynamicNameServerResolver.cpp
new file mode 100644
index 0000000..df1fb2b
--- /dev/null
+++ b/src/main/cpp/rocketmq/DynamicNameServerResolver.cpp
@@ -0,0 +1,126 @@
+#include "DynamicNameServerResolver.h"
+
+#include <atomic>
+#include <chrono>
+#include <cstdint>
+#include <functional>
+#include <memory>
+
+#include "absl/strings/str_join.h"
+#include "spdlog/spdlog.h"
+
+#include "SchedulerImpl.h"
+
+ROCKETMQ_NAMESPACE_BEGIN
+
+DynamicNameServerResolver::DynamicNameServerResolver(absl::string_view endpoint,
+ std::chrono::milliseconds refresh_interval)
+ : endpoint_(endpoint.data(), endpoint.length()), refresh_interval_(refresh_interval),
+ scheduler_(absl::make_unique<SchedulerImpl>()) {
+ absl::string_view remains;
+ if (absl::StartsWith(endpoint_, "https://")) {
+ ssl_ = true;
+ remains = absl::StripPrefix(endpoint_, "https://");
+ } else {
+ remains = absl::StripPrefix(endpoint_, "http://");
+ }
+
+ std::int32_t port = 80;
+ if (ssl_) {
+ port = 443;
+ }
+
+ absl::string_view host;
+ if (absl::StrContains(remains, ':')) {
+ std::vector<absl::string_view> segments = absl::StrSplit(remains, ':');
+ host = segments[0];
+ remains = absl::StripPrefix(remains, host);
+ remains = absl::StripPrefix(remains, ":");
+
+ segments = absl::StrSplit(remains, '/');
+ if (!absl::SimpleAtoi(segments[0], &port)) {
+ SPDLOG_WARN("Failed to parse port of name-server-list discovery service endpoint");
+ abort();
+ }
+ remains = absl::StripPrefix(remains, segments[0]);
+ } else {
+ std::vector<absl::string_view> segments = absl::StrSplit(remains, '/');
+ host = segments[0];
+ remains = absl::StripPrefix(remains, host);
+ }
+
+ top_addressing_ = absl::make_unique<TopAddressing>(std::string(host.data(), host.length()), port,
+ std::string(remains.data(), remains.length()));
+}
+
+std::vector<std::string> DynamicNameServerResolver::resolve() {
+ bool fetch_immediately = false;
+ {
+ absl::MutexLock lk(&name_server_list_mtx_);
+ if (name_server_list_.empty()) {
+ fetch_immediately = true;
+ }
+ }
+
+ if (fetch_immediately) {
+ fetch();
+ }
+
+ {
+ absl::MutexLock lk(&name_server_list_mtx_);
+ return name_server_list_;
+ }
+}
+
+void DynamicNameServerResolver::fetch() {
+ std::weak_ptr<DynamicNameServerResolver> ptr(shared_from_this());
+ auto callback = [ptr](bool success, const std::vector<std::string>& name_server_list) {
+ if (success && !name_server_list.empty()) {
+ std::shared_ptr<DynamicNameServerResolver> resolver = ptr.lock();
+ if (resolver) {
+ resolver->onNameServerListFetched(name_server_list);
+ }
+ }
+ };
+ top_addressing_->fetchNameServerAddresses(callback);
+}
+
+void DynamicNameServerResolver::onNameServerListFetched(const std::vector<std::string>& name_server_list) {
+ if (!name_server_list.empty()) {
+ absl::MutexLock lk(&name_server_list_mtx_);
+ if (name_server_list_ != name_server_list) {
+ SPDLOG_INFO("Name server list changed. {} --> {}", absl::StrJoin(name_server_list_, ";"),
+ absl::StrJoin(name_server_list, ";"));
+ name_server_list_ = name_server_list;
+ }
+ }
+}
+
+void DynamicNameServerResolver::injectHttpClient(std::unique_ptr<HttpClient> http_client) {
+ top_addressing_->injectHttpClient(std::move(http_client));
+}
+
+void DynamicNameServerResolver::start() {
+ scheduler_->start();
+ scheduler_->schedule(std::bind(&DynamicNameServerResolver::fetch, this), "DynamicNameServerResolver",
+ std::chrono::milliseconds(0), refresh_interval_);
+}
+
+void DynamicNameServerResolver::shutdown() { scheduler_->shutdown(); }
+
+std::string DynamicNameServerResolver::current() {
+ absl::MutexLock lk(&name_server_list_mtx_);
+ if (name_server_list_.empty()) {
+ return std::string();
+ }
+
+ std::uint32_t index = index_.load(std::memory_order_relaxed) % name_server_list_.size();
+ return name_server_list_[index];
+}
+
+std::string DynamicNameServerResolver::next() {
+ index_.fetch_add(1, std::memory_order_relaxed);
+ return current();
+}
+
+ROCKETMQ_NAMESPACE_END
\ No newline at end of file
diff --git a/src/main/cpp/rocketmq/StaticNameServerResolver.cpp b/src/main/cpp/rocketmq/StaticNameServerResolver.cpp
new file mode 100644
index 0000000..6e156be
--- /dev/null
+++ b/src/main/cpp/rocketmq/StaticNameServerResolver.cpp
@@ -0,0 +1,24 @@
+#include "StaticNameServerResolver.h"
+
+#include "absl/strings/str_split.h"
+#include <atomic>
+#include <cstdint>
+
+ROCKETMQ_NAMESPACE_BEGIN
+
+StaticNameServerResolver::StaticNameServerResolver(absl::string_view name_server_list)
+ : name_server_list_(absl::StrSplit(name_server_list, ';')) {}
+
+std::string StaticNameServerResolver::current() {
+ std::uint32_t index = index_.load(std::memory_order_relaxed) % name_server_list_.size();
+ return name_server_list_[index];
+}
+
+std::string StaticNameServerResolver::next() {
+ index_.fetch_add(1, std::memory_order_relaxed);
+ return current();
+}
+
+std::vector<std::string> StaticNameServerResolver::resolve() { return name_server_list_; }
+
+ROCKETMQ_NAMESPACE_END
\ No newline at end of file
diff --git a/src/main/cpp/rocketmq/include/ClientImpl.h b/src/main/cpp/rocketmq/include/ClientImpl.h
index 4a2dc6e..a3046cf 100644
--- a/src/main/cpp/rocketmq/include/ClientImpl.h
+++ b/src/main/cpp/rocketmq/include/ClientImpl.h
@@ -1,13 +1,19 @@
+#pragma once
+
+#include <chrono>
+#include <cstdint>
+
+#include "apache/rocketmq/v1/definition.pb.h"
+
#include "Client.h"
#include "ClientConfigImpl.h"
#include "ClientManager.h"
#include "ClientResourceBundle.h"
#include "InvocationContext.h"
+#include "NameServerResolver.h"
#include "OtlpExporter.h"
-#include "apache/rocketmq/v1/definition.pb.h"
#include "rocketmq/MQMessageExt.h"
#include "rocketmq/State.h"
-#include <chrono>
ROCKETMQ_NAMESPACE_BEGIN
@@ -31,11 +37,6 @@
*/
void endpointsInUse(absl::flat_hash_set<std::string>& endpoints) override LOCKS_EXCLUDED(topic_route_table_mtx_);
- void setNameServerList(std::vector<std::string> name_server_list) {
- absl::MutexLock lk(&name_server_list_mtx_);
- name_server_list_ = std::move(name_server_list);
- }
-
void heartbeat() override;
bool active() override {
@@ -48,6 +49,10 @@
void schedule(const std::string& task_name, const std::function<void(void)>& task,
std::chrono::milliseconds delay) override;
+ void withNameServerResolver(std::shared_ptr<NameServerResolver> name_server_resolver) {
+ name_server_resolver_ = std::move(name_server_resolver);
+ }
+
protected:
ClientManagerPtr client_manager_;
std::shared_ptr<OtlpExporter> exporter_;
@@ -60,14 +65,10 @@
inflight_route_requests_ GUARDED_BY(inflight_route_requests_mtx_);
absl::Mutex inflight_route_requests_mtx_ ACQUIRED_BEFORE(topic_route_table_mtx_); // Protects inflight_route_requests_
static const char* UPDATE_ROUTE_TASK_NAME;
- std::uintptr_t route_update_handle_{0};
+ std::uint32_t route_update_handle_{0};
- // Name server list management
- std::vector<std::string> name_server_list_ GUARDED_BY(name_server_list_mtx_);
- absl::Mutex name_server_list_mtx_; // protects name_server_list_
-
- static const char* UPDATE_NAME_SERVER_LIST_TASK_NAME;
- std::uintptr_t name_server_update_handle_{0};
+ // Set Name Server Resolver
+ std::shared_ptr<NameServerResolver> name_server_resolver_;
absl::flat_hash_map<std::string, absl::Time> multiplexing_requests_;
absl::Mutex multiplexing_requests_mtx_;
@@ -75,12 +76,6 @@
absl::flat_hash_set<std::string> isolated_endpoints_ GUARDED_BY(isolated_endpoints_mtx_);
absl::Mutex isolated_endpoints_mtx_;
- void debugNameServerChanges(const std::vector<std::string>& list) LOCKS_EXCLUDED(name_server_list_mtx_);
-
- void renewNameServerList() LOCKS_EXCLUDED(name_server_list_mtx_);
-
- bool selectNameServer(std::string& selected, bool change = false) LOCKS_EXCLUDED(name_server_list_mtx_);
-
void updateRouteInfo() LOCKS_EXCLUDED(topic_route_table_mtx_);
/**
@@ -104,7 +99,7 @@
return resource_bundle;
}
- void setAccessPoint(rmq::Endpoints* endpoints) LOCKS_EXCLUDED(name_server_list_mtx_);
+ void setAccessPoint(rmq::Endpoints* endpoints);
virtual void notifyClientTermination();
diff --git a/src/main/cpp/rocketmq/include/DynamicNameServerResolver.h b/src/main/cpp/rocketmq/include/DynamicNameServerResolver.h
new file mode 100644
index 0000000..7626bd1
--- /dev/null
+++ b/src/main/cpp/rocketmq/include/DynamicNameServerResolver.h
@@ -0,0 +1,59 @@
+#pragma once
+
+#include <atomic>
+#include <chrono>
+#include <cstdint>
+#include <cstdlib>
+#include <memory>
+
+#include "absl/base/thread_annotations.h"
+#include "absl/memory/memory.h"
+#include "absl/strings/numbers.h"
+#include "absl/strings/str_split.h"
+#include "absl/strings/string_view.h"
+#include "absl/synchronization/mutex.h"
+
+#include "NameServerResolver.h"
+#include "Scheduler.h"
+#include "TopAddressing.h"
+
+ROCKETMQ_NAMESPACE_BEGIN
+
+class DynamicNameServerResolver : public NameServerResolver,
+ public std::enable_shared_from_this<DynamicNameServerResolver> {
+public:
+ DynamicNameServerResolver(absl::string_view endpoint, std::chrono::milliseconds refresh_interval);
+
+ void start() override;
+
+ void shutdown() override;
+
+ std::string current() override LOCKS_EXCLUDED(name_server_list_mtx_);
+
+ std::string next() override LOCKS_EXCLUDED(name_server_list_mtx_);
+
+ std::vector<std::string> resolve() override LOCKS_EXCLUDED(name_server_list_mtx_);
+
+ void injectHttpClient(std::unique_ptr<HttpClient> http_client);
+
+private:
+ std::string endpoint_;
+
+ std::chrono::milliseconds refresh_interval_;
+
+ void fetch();
+
+ void onNameServerListFetched(const std::vector<std::string>& name_server_list) LOCKS_EXCLUDED(name_server_list_mtx_);
+
+ std::vector<std::string> name_server_list_ GUARDED_BY(name_server_list_mtx_);
+ absl::Mutex name_server_list_mtx_;
+
+ std::atomic<std::uint32_t> index_{0};
+
+ bool ssl_{false};
+ std::unique_ptr<TopAddressing> top_addressing_;
+
+ std::unique_ptr<Scheduler> scheduler_;
+};
+
+ROCKETMQ_NAMESPACE_END
\ No newline at end of file
diff --git a/src/main/cpp/rocketmq/include/NameServerResolver.h b/src/main/cpp/rocketmq/include/NameServerResolver.h
new file mode 100644
index 0000000..829d2d2
--- /dev/null
+++ b/src/main/cpp/rocketmq/include/NameServerResolver.h
@@ -0,0 +1,25 @@
+#pragma once
+
+#include <string>
+#include <vector>
+
+#include "rocketmq/RocketMQ.h"
+
+ROCKETMQ_NAMESPACE_BEGIN
+
+class NameServerResolver {
+public:
+ virtual ~NameServerResolver() = default;
+
+ virtual void start() = 0;
+
+ virtual void shutdown() = 0;
+
+ virtual std::string next() = 0;
+
+ virtual std::string current() = 0;
+
+ virtual std::vector<std::string> resolve() = 0;
+};
+
+ROCKETMQ_NAMESPACE_END
\ No newline at end of file
diff --git a/src/main/cpp/rocketmq/include/StaticNameServerResolver.h b/src/main/cpp/rocketmq/include/StaticNameServerResolver.h
new file mode 100644
index 0000000..1d62106
--- /dev/null
+++ b/src/main/cpp/rocketmq/include/StaticNameServerResolver.h
@@ -0,0 +1,33 @@
+#pragma once
+
+#include <cstdint>
+#include <vector>
+#include <atomic>
+
+#include "absl/strings/string_view.h"
+
+#include "NameServerResolver.h"
+#include "rocketmq/RocketMQ.h"
+
+ROCKETMQ_NAMESPACE_BEGIN
+
+class StaticNameServerResolver : public NameServerResolver {
+public:
+ StaticNameServerResolver(absl::string_view name_server_list);
+
+ void start() override {}
+
+ void shutdown() override {}
+
+ std::string current() override;
+
+ std::string next() override;
+
+ std::vector<std::string> resolve() override;
+
+private:
+ std::vector<std::string> name_server_list_;
+ std::atomic<std::uint32_t> index_{0};
+};
+
+ROCKETMQ_NAMESPACE_END
\ No newline at end of file
diff --git a/src/main/cpp/rocketmq/mocks/include/NameServerResolverMock.h b/src/main/cpp/rocketmq/mocks/include/NameServerResolverMock.h
new file mode 100644
index 0000000..89d7c7c
--- /dev/null
+++ b/src/main/cpp/rocketmq/mocks/include/NameServerResolverMock.h
@@ -0,0 +1,21 @@
+#pragma once
+
+#include "NameServerResolver.h"
+#include "gmock/gmock.h"
+
+ROCKETMQ_NAMESPACE_BEGIN
+
+class NameServerResolverMock : public NameServerResolver {
+public:
+ MOCK_METHOD(void, start, (), (override));
+
+ MOCK_METHOD(void, shutdown, (), (override));
+
+ MOCK_METHOD(std::string, next, (), (override));
+
+ MOCK_METHOD(std::string, current, (), (override));
+
+ MOCK_METHOD((std::vector<std::string>), resolve, (), (override));
+};
+
+ROCKETMQ_NAMESPACE_END
\ No newline at end of file
diff --git a/src/test/cpp/ut/rocketmq/BUILD.bazel b/src/test/cpp/ut/rocketmq/BUILD.bazel
index c70c146..47fa3d8 100644
--- a/src/test/cpp/ut/rocketmq/BUILD.bazel
+++ b/src/test/cpp/ut/rocketmq/BUILD.bazel
@@ -196,6 +196,31 @@
"//src/main/cpp/base/mocks:base_mocks",
"//src/main/cpp/client/mocks:client_mocks",
"//src/main/cpp/rocketmq:rocketmq_library",
+ "//src/main/cpp/rocketmq/mocks:rocketmq_mocks",
+ "@com_google_googletest//:gtest_main",
+ ],
+)
+
+cc_test(
+ name = "static_name_server_resolver_test",
+ srcs = [
+ "StaticNameServerResolverTest.cpp",
+ ],
+ deps = [
+ "//src/main/cpp/rocketmq:rocketmq_library",
+ "@com_google_googletest//:gtest_main",
+ ],
+)
+
+
+cc_test(
+ name = "dynamic_name_server_resolver_test",
+ srcs = [
+ "DynamicNameServerResolverTest.cpp",
+ ],
+ deps = [
+ "//src/main/cpp/rocketmq:rocketmq_library",
+ "//src/main/cpp/base/mocks:base_mocks",
"@com_google_googletest//:gtest_main",
],
)
\ No newline at end of file
diff --git a/src/test/cpp/ut/rocketmq/ClientImplTest.cpp b/src/test/cpp/ut/rocketmq/ClientImplTest.cpp
index ea071f1..424b4d6 100644
--- a/src/test/cpp/ut/rocketmq/ClientImplTest.cpp
+++ b/src/test/cpp/ut/rocketmq/ClientImplTest.cpp
@@ -1,13 +1,18 @@
+#include <chrono>
+#include <memory>
+#include <string>
+
#include "ClientImpl.h"
#include "ClientManagerFactory.h"
#include "ClientManagerMock.h"
+#include "DynamicNameServerResolver.h"
#include "HttpClientMock.h"
+#include "NameServerResolverMock.h"
#include "SchedulerImpl.h"
#include "TopAddressing.h"
#include "rocketmq/RocketMQ.h"
+
#include "gtest/gtest.h"
-#include <memory>
-#include <string>
ROCKETMQ_NAMESPACE_BEGIN
@@ -24,11 +29,13 @@
public:
void SetUp() override {
grpc_init();
+ name_server_resolver_ = std::make_shared<DynamicNameServerResolver>(endpoint_, std::chrono::seconds(1));
scheduler_.start();
client_manager_ = std::make_shared<testing::NiceMock<ClientManagerMock>>();
ClientManagerFactory::getInstance().addClientManager(resource_namespace_, client_manager_);
ON_CALL(*client_manager_, getScheduler).WillByDefault(testing::ReturnRef(scheduler_));
client_ = std::make_shared<TestClientImpl>(group_);
+ client_->withNameServerResolver(name_server_resolver_);
}
void TearDown() override {
@@ -37,17 +44,17 @@
}
protected:
+ std::string endpoint_{"http://jmenv.tbsite.net:8080/rocketmq/nsaddr"};
std::string resource_namespace_{"mq://test"};
std::string group_{"Group-0"};
std::shared_ptr<testing::NiceMock<ClientManagerMock>> client_manager_;
SchedulerImpl scheduler_;
std::shared_ptr<TestClientImpl> client_;
+ std::shared_ptr<DynamicNameServerResolver> name_server_resolver_;
};
TEST_F(ClientImplTest, testBasic) {
- TopAddressing top_addressing_;
-
auto http_client = absl::make_unique<HttpClientMock>();
std::string once{"10.0.0.1:9876"};
@@ -74,9 +81,8 @@
};
EXPECT_CALL(*http_client, get).WillOnce(testing::Invoke(once_cb)).WillRepeatedly(testing::Invoke(then_cb));
- top_addressing_.injectHttpClient(std::move(http_client));
+ name_server_resolver_->injectHttpClient(std::move(http_client));
- ON_CALL(*client_manager_, topAddressing).WillByDefault(testing::ReturnRef(top_addressing_));
client_->resourceNamespace(resource_namespace_);
client_->start();
{
diff --git a/src/test/cpp/ut/rocketmq/DefaultMQProducerTest.cpp b/src/test/cpp/ut/rocketmq/DefaultMQProducerTest.cpp
index c458519..e827a1e 100644
--- a/src/test/cpp/ut/rocketmq/DefaultMQProducerTest.cpp
+++ b/src/test/cpp/ut/rocketmq/DefaultMQProducerTest.cpp
@@ -132,7 +132,7 @@
TEST_F(DefaultMQProducerUnitTest, testAsyncSendMessage) {
auto producer = std::make_shared<ProducerImpl>(group_name_);
producer->resourceNamespace(resource_namespace_);
- producer->setNameServerList(name_server_list_);
+ producer->withNameServerResolver(name_server_resolver_);
producer->setCredentialsProvider(credentials_provider_);
producer->start();
MQMessage message;
@@ -155,7 +155,7 @@
TEST_F(DefaultMQProducerUnitTest, testSendMessage) {
auto producer = std::make_shared<ProducerImpl>(group_name_);
producer->resourceNamespace(resource_namespace_);
- producer->setNameServerList(name_server_list_);
+ producer->withNameServerResolver(name_server_resolver_);
producer->setCredentialsProvider(credentials_provider_);
producer->start();
MQMessage message;
@@ -167,7 +167,7 @@
TEST_F(DefaultMQProducerUnitTest, testEndpointIsolation) {
auto producer = std::make_shared<ProducerImpl>(group_name_);
producer->resourceNamespace(resource_namespace_);
- producer->setNameServerList(name_server_list_);
+ producer->withNameServerResolver(name_server_resolver_);
producer->setCredentialsProvider(credentials_provider_);
producer->start();
diff --git a/src/test/cpp/ut/rocketmq/DynamicNameServerResolverTest.cpp b/src/test/cpp/ut/rocketmq/DynamicNameServerResolverTest.cpp
new file mode 100644
index 0000000..2435b81
--- /dev/null
+++ b/src/test/cpp/ut/rocketmq/DynamicNameServerResolverTest.cpp
@@ -0,0 +1,61 @@
+#include "DynamicNameServerResolver.h"
+
+#include <chrono>
+#include <map>
+#include <memory>
+
+#include "absl/memory/memory.h"
+#include "absl/strings/str_join.h"
+#include "gmock/gmock.h"
+#include "gtest/gtest.h"
+#include <string>
+
+#include "HttpClientMock.h"
+
+ROCKETMQ_NAMESPACE_BEGIN
+
+class DynamicNameServerResolverTest : public testing::Test {
+public:
+ DynamicNameServerResolverTest()
+ : resolver_(std::make_shared<DynamicNameServerResolver>(endpoint_, std::chrono::seconds(1))) {}
+
+ void SetUp() override {
+ auto http_client = absl::make_unique<testing::NiceMock<HttpClientMock>>();
+
+ auto callback =
+ [this](HttpProtocol, const std::string&, std::uint16_t, const std::string&,
+ const std::function<void(int, const std::multimap<std::string, std::string>&, const std::string&)>& cb) {
+ int code = 200;
+ std::multimap<std::string, std::string> headers;
+ cb(code, headers, name_server_list_);
+ };
+
+ ON_CALL(*http_client, get).WillByDefault(testing::Invoke(callback));
+
+ resolver_->injectHttpClient(std::move(http_client));
+
+ resolver_->start();
+ }
+
+ void TearDown() override { resolver_->shutdown(); }
+
+protected:
+ std::string endpoint_{"http://jmenv.tbsite.net:8080/rocketmq/nsaddr"};
+ std::string name_server_list_{"10.0.0.0:9876;10.0.0.1:9876"};
+ std::shared_ptr<DynamicNameServerResolver> resolver_;
+};
+
+TEST_F(DynamicNameServerResolverTest, testResolve) {
+ auto name_server_list = resolver_->resolve();
+ ASSERT_FALSE(name_server_list.empty());
+ std::string resolved = absl::StrJoin(name_server_list, ";");
+ ASSERT_EQ(name_server_list_, resolved);
+
+ std::string first{"10.0.0.0:9876"};
+ EXPECT_EQ(first, resolver_->current());
+
+ std::string second{"10.0.0.1:9876"};
+ EXPECT_EQ(second, resolver_->next());
+}
+
+ROCKETMQ_NAMESPACE_END
\ No newline at end of file
diff --git a/src/test/cpp/ut/rocketmq/ProducerImplTest.cpp b/src/test/cpp/ut/rocketmq/ProducerImplTest.cpp
index 23f7ae7..a9e253b 100644
--- a/src/test/cpp/ut/rocketmq/ProducerImplTest.cpp
+++ b/src/test/cpp/ut/rocketmq/ProducerImplTest.cpp
@@ -1,14 +1,16 @@
-#include "ProducerImpl.h"
+#include <memory>
+
#include "ClientManagerFactory.h"
#include "ClientManagerMock.h"
+#include "ProducerImpl.h"
#include "SchedulerImpl.h"
+#include "StaticNameServerResolver.h"
#include "TopicRouteData.h"
#include "rocketmq/AsyncCallback.h"
#include "rocketmq/MQMessage.h"
#include "rocketmq/MQSelector.h"
#include "rocketmq/RocketMQ.h"
#include "rocketmq/SendResult.h"
-#include <memory>
ROCKETMQ_NAMESPACE_BEGIN
@@ -19,11 +21,12 @@
void SetUp() override {
grpc_init();
+ name_server_resolver_ = std::make_shared<StaticNameServerResolver>(name_server_list_);
client_manager_ = std::make_shared<testing::NiceMock<ClientManagerMock>>();
ClientManagerFactory::getInstance().addClientManager(resource_namespace_, client_manager_);
producer_ = std::make_shared<ProducerImpl>(group_);
producer_->resourceNamespace(resource_namespace_);
- producer_->setNameServerList(name_server_list_);
+ producer_->withNameServerResolver(name_server_resolver_);
producer_->setCredentialsProvider(credentials_provider_);
{
@@ -44,7 +47,8 @@
protected:
std::shared_ptr<testing::NiceMock<ClientManagerMock>> client_manager_;
std::shared_ptr<ProducerImpl> producer_;
- std::vector<std::string> name_server_list_{"10.0.0.1:9876"};
+ std::string name_server_list_{"10.0.0.1:9876"};
+ std::shared_ptr<NameServerResolver> name_server_resolver_;
std::string resource_namespace_{"mq://test"};
std::string group_{"CID_test"};
std::string topic_{"Topic0"};
diff --git a/src/test/cpp/ut/rocketmq/PullConsumerImplTest.cpp b/src/test/cpp/ut/rocketmq/PullConsumerImplTest.cpp
index a6db0ca..90f32d8 100644
--- a/src/test/cpp/ut/rocketmq/PullConsumerImplTest.cpp
+++ b/src/test/cpp/ut/rocketmq/PullConsumerImplTest.cpp
@@ -1,30 +1,36 @@
-#include "PullConsumerImpl.h"
-#include "ClientManagerFactory.h"
-#include "ClientManagerMock.h"
-#include "InvocationContext.h"
-#include "Scheduler.h"
-#include "rocketmq/AsyncCallback.h"
-#include "rocketmq/ConsumeType.h"
-#include "rocketmq/RocketMQ.h"
-#include "gtest/gtest.h"
-#include <apache/rocketmq/v1/definition.pb.h>
#include <chrono>
#include <memory>
#include <string>
+#include "ClientManagerFactory.h"
+#include "ClientManagerMock.h"
+#include "InvocationContext.h"
+#include "PullConsumerImpl.h"
+#include "Scheduler.h"
+#include "StaticNameServerResolver.h"
+#include "apache/rocketmq/v1/definition.pb.h"
+#include "rocketmq/AsyncCallback.h"
+#include "rocketmq/ConsumeType.h"
+#include "rocketmq/RocketMQ.h"
+
+#include "gtest/gtest.h"
+
ROCKETMQ_NAMESPACE_BEGIN
class PullConsumerImplTest : public testing::Test {
public:
void SetUp() override {
grpc_init();
+
+ name_server_resolver_ = std::make_shared<StaticNameServerResolver>(name_server_list_);
+
scheduler_.start();
client_manager_ = std::make_shared<testing::NiceMock<ClientManagerMock>>();
ON_CALL(*client_manager_, getScheduler).WillByDefault(testing::ReturnRef(scheduler_));
ClientManagerFactory::getInstance().addClientManager(resource_namespace_, client_manager_);
pull_consumer_ = std::make_shared<PullConsumerImpl>(group_);
- pull_consumer_->setNameServerList(name_server_list_);
+ pull_consumer_->withNameServerResolver(name_server_resolver_);
pull_consumer_->resourceNamespace(resource_namespace_);
{
@@ -47,7 +53,8 @@
protected:
std::string resource_namespace_{"mq://test"};
- std::vector<std::string> name_server_list_{"10.0.0.1:9876"};
+ std::string name_server_list_{"10.0.0.1:9876"};
+ std::shared_ptr<NameServerResolver> name_server_resolver_;
std::string group_{"Group-0"};
std::string topic_{"Test"};
std::string tag_{"TagB"};
diff --git a/src/test/cpp/ut/rocketmq/PushConsumerImplTest.cpp b/src/test/cpp/ut/rocketmq/PushConsumerImplTest.cpp
index d98ef51..0789af4 100644
--- a/src/test/cpp/ut/rocketmq/PushConsumerImplTest.cpp
+++ b/src/test/cpp/ut/rocketmq/PushConsumerImplTest.cpp
@@ -1,14 +1,17 @@
-#include "PushConsumerImpl.h"
+#include <memory>
+
+#include "gtest/gtest.h"
+
#include "ClientManagerFactory.h"
#include "ClientManagerMock.h"
#include "InvocationContext.h"
#include "MessageAccessor.h"
+#include "PushConsumerImpl.h"
+#include "StaticNameServerResolver.h"
#include "grpc/grpc.h"
#include "rocketmq/MQMessageExt.h"
#include "rocketmq/MessageListener.h"
#include "rocketmq/RocketMQ.h"
-#include "gtest/gtest.h"
-#include <memory>
ROCKETMQ_NAMESPACE_BEGIN
@@ -25,18 +28,22 @@
void SetUp() override {
grpc_init();
+
+ name_server_resolver_ = std::make_shared<StaticNameServerResolver>(name_server_list_);
+
client_manager_ = std::make_shared<testing::NiceMock<ClientManagerMock>>();
ClientManagerFactory::getInstance().addClientManager(resource_namespace_, client_manager_);
push_consumer_ = std::make_shared<PushConsumerImpl>(group_);
push_consumer_->resourceNamespace(resource_namespace_);
- push_consumer_->setNameServerList(name_server_list_);
+ push_consumer_->withNameServerResolver(name_server_resolver_);
push_consumer_->registerMessageListener(message_listener_.get());
}
void TearDown() override { grpc_shutdown(); }
protected:
- std::vector<std::string> name_server_list_{"10.0.0.1:9876"};
+ std::string name_server_list_{"10.0.0.1:9876"};
+ std::shared_ptr<StaticNameServerResolver> name_server_resolver_;
std::string resource_namespace_{"mq://test"};
std::string group_{"CID_test"};
std::string topic_{"Topic0"};
diff --git a/src/test/cpp/ut/rocketmq/StaticNameServerResolverTest.cpp b/src/test/cpp/ut/rocketmq/StaticNameServerResolverTest.cpp
new file mode 100644
index 0000000..dab7a02
--- /dev/null
+++ b/src/test/cpp/ut/rocketmq/StaticNameServerResolverTest.cpp
@@ -0,0 +1,38 @@
+#include "StaticNameServerResolver.h"
+
+#include "absl/strings/str_split.h"
+
+#include "gtest/gtest.h"
+#include <vector>
+
+ROCKETMQ_NAMESPACE_BEGIN
+
+class StaticNameServerResolverTest : public testing::Test {
+public:
+ StaticNameServerResolverTest() : resolver_(name_server_list_) {}
+
+ void SetUp() override { resolver_.start(); }
+
+ void TearDown() override { resolver_.shutdown(); }
+
+protected:
+ std::string name_server_list_{"10.0.0.1:9876;10.0.0.2:9876"};
+ StaticNameServerResolver resolver_;
+};
+
+TEST_F(StaticNameServerResolverTest, testResolve) {
+ std::vector<std::string> segments = absl::StrSplit(name_server_list_, ';');
+ ASSERT_EQ(segments, resolver_.resolve());
+}
+
+TEST_F(StaticNameServerResolverTest, testCurrentNext) {
+ std::string&& name_server_1 = resolver_.current();
+ std::string expected = "10.0.0.1:9876";
+ EXPECT_EQ(expected, name_server_1);
+
+ expected = "10.0.0.2:9876";
+ std::string&& name_server_2 = resolver_.next();
+ EXPECT_EQ(expected, name_server_2);
+}
+
+ROCKETMQ_NAMESPACE_END
\ No newline at end of file
diff --git a/src/test/cpp/ut/rocketmq/include/MQClientTest.h b/src/test/cpp/ut/rocketmq/include/MQClientTest.h
index 405aca3..6c13f20 100644
--- a/src/test/cpp/ut/rocketmq/include/MQClientTest.h
+++ b/src/test/cpp/ut/rocketmq/include/MQClientTest.h
@@ -1,19 +1,22 @@
#pragma once
-#include "ClientManagerFactory.h"
-#include "RpcClientMock.h"
-#include "ClientManagerImpl.h"
+#include <functional>
+#include <memory>
+
#include "grpc/grpc.h"
#include "gmock/gmock.h"
#include "gtest/gtest.h"
-#include <functional>
-#include <memory>
+
+#include "ClientManagerFactory.h"
+#include "ClientManagerImpl.h"
+#include "RpcClientMock.h"
+#include "StaticNameServerResolver.h"
ROCKETMQ_NAMESPACE_BEGIN
class MQClientTest : public testing::Test {
public:
- MQClientTest() = default;
+ MQClientTest() = default;
void SetUp() override {
grpc_init();
@@ -27,6 +30,8 @@
std::bind(&MQClientTest::mockQueryRoute, this, std::placeholders::_1, std::placeholders::_2)));
client_instance_->addRpcClient(name_server_address_, rpc_client_ns_);
ClientManagerFactory::getInstance().addClientManager(resource_namespace_, client_instance_);
+
+ name_server_resolver_ = std::make_shared<StaticNameServerResolver>(name_server_address_);
}
void TearDown() override {
@@ -43,6 +48,7 @@
const int32_t partition_num_{24};
const int32_t avg_partition_per_host_{8};
std::string name_server_address_{"ipv4:127.0.0.1:9876"};
+ std::shared_ptr<NameServerResolver> name_server_resolver_;
std::string group_name_{"CID_Test"};
std::string topic_{"Topic_Test"};
std::string resource_namespace_{"mq://test"};