Skip to content
7 changes: 7 additions & 0 deletions api/envoy/api/v2/core/config_source.proto
Original file line number Diff line number Diff line change
Expand Up @@ -104,4 +104,11 @@ message ConfigSource {
// source in the bootstrap configuration is used.
AggregatedConfigSource ads = 3;
}

// Optional initialization timeout.
// When this timeout is specified, the xDS API will be considered initialized after the specified
// time, even if the first config is not delivered yet. The timer is activated when the xDS API
// starts initializing, and is disarmed when it is initialized (successfully or not). 0 means no
// timeout - Envoy will wait indefinitely for the first xDS config. Default 0.
google.protobuf.Duration initial_fetch_timeout = 4;
}
21 changes: 21 additions & 0 deletions source/common/router/rds_impl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,24 @@ RdsRouteConfigSubscription::RdsRouteConfigSubscription(
factory_context.clusterManager(), factory_context.random(), *scope_,
"envoy.api.v2.RouteDiscoveryService.FetchRoutes",
"envoy.api.v2.RouteDiscoveryService.StreamRoutes", factory_context.api());

initialization_timeout_ = std::chrono::milliseconds(
PROTOBUF_GET_MS_OR_DEFAULT(rds.config_source(), initial_fetch_timeout, 0));
if (initialization_timeout_.count() > 0) {
initialization_timeout_timer_ = factory_context.dispatcher().createTimer([this]() -> void {
ENVOY_LOG(warn, "rds: initialization timed out for route_config_name={}", route_config_name_);
onConfigUpdateFailed(nullptr);
});
}
}

void RdsRouteConfigSubscription::initialize(std::function<void()> callback) {
initialize_callback_ = callback;
if (initialization_timeout_.count() > 0) {
ASSERT(initialization_timeout_timer_);
initialization_timeout_timer_->enableTimer(initialization_timeout_);
}
subscription_->start({route_config_name_}, *this);
Comment thread
htuch marked this conversation as resolved.
Outdated
}

RdsRouteConfigSubscription::~RdsRouteConfigSubscription() {
Expand Down Expand Up @@ -138,6 +156,9 @@ void RdsRouteConfigSubscription::runInitializeCallbackIfAny() {
initialize_callback_();
initialize_callback_ = nullptr;
}
if (initialization_timeout_timer_) {
initialization_timeout_timer_->disableTimer();
}
}

RdsRouteConfigProviderImpl::RdsRouteConfigProviderImpl(
Expand Down
7 changes: 3 additions & 4 deletions source/common/router/rds_impl.h
Original file line number Diff line number Diff line change
Expand Up @@ -101,10 +101,7 @@ class RdsRouteConfigSubscription
~RdsRouteConfigSubscription();

// Init::Target
void initialize(std::function<void()> callback) override {
initialize_callback_ = callback;
subscription_->start({route_config_name_}, *this);
}
void initialize(std::function<void()> callback) override;

// Config::SubscriptionCallbacks
void onConfigUpdate(const ResourceVector& resources, const std::string& version_info) override;
Expand All @@ -130,6 +127,8 @@ class RdsRouteConfigSubscription

std::unique_ptr<Envoy::Config::Subscription<envoy::api::v2::RouteConfiguration>> subscription_;
std::function<void()> initialize_callback_;
Event::TimerPtr initialization_timeout_timer_;
std::chrono::milliseconds initialization_timeout_;
const std::string route_config_name_;
Stats::ScopePtr scope_;
RdsStats stats_;
Expand Down
21 changes: 21 additions & 0 deletions source/common/upstream/cds_api_impl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,23 @@ CdsApiImpl::CdsApiImpl(const envoy::api::v2::core::ConfigSource& cds_config, Clu
cds_config, local_info, dispatcher, cm, random, *scope_,
"envoy.api.v2.ClusterDiscoveryService.FetchClusters",
"envoy.api.v2.ClusterDiscoveryService.StreamClusters", api);

initialization_timeout_ =
std::chrono::milliseconds(PROTOBUF_GET_MS_OR_DEFAULT(cds_config, initial_fetch_timeout, 0));
if (initialization_timeout_.count() > 0) {
initialization_timeout_timer_ = dispatcher.createTimer([this]() -> void {
ENVOY_LOG(warn, "cds: initialization timed out");
onConfigUpdateFailed(nullptr);
});
}
}

void CdsApiImpl::initialize() {
if (initialization_timeout_.count() > 0) {
ASSERT(initialization_timeout_timer_);
initialization_timeout_timer_->enableTimer(initialization_timeout_);
}
subscription_->start({}, *this);
}

void CdsApiImpl::onConfigUpdate(const ResourceVector& resources, const std::string& version_info) {
Expand Down Expand Up @@ -117,6 +134,10 @@ void CdsApiImpl::runInitializeCallbackIfAny() {
initialize_callback_();
initialize_callback_ = nullptr;
}
if (initialization_timeout_timer_) {
initialization_timeout_timer_->disableTimer();
initialization_timeout_timer_.reset();
}
}

} // namespace Upstream
Expand Down
4 changes: 3 additions & 1 deletion source/common/upstream/cds_api_impl.h
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ class CdsApiImpl : public CdsApi,
Api::Api& api);

// Upstream::CdsApi
void initialize() override { subscription_->start({}, *this); }
void initialize() override;
void setInitializedCb(std::function<void()> callback) override {
initialize_callback_ = callback;
}
Expand All @@ -52,6 +52,8 @@ class CdsApiImpl : public CdsApi,
std::string version_info_;
std::function<void()> initialize_callback_;
Stats::ScopePtr scope_;
Event::TimerPtr initialization_timeout_timer_;
std::chrono::milliseconds initialization_timeout_;
};

} // namespace Upstream
Expand Down
16 changes: 16 additions & 0 deletions source/server/lds_api.cc
Original file line number Diff line number Diff line change
Expand Up @@ -28,10 +28,23 @@ LdsApiImpl::LdsApiImpl(const envoy::api::v2::core::ConfigSource& lds_config,
"envoy.api.v2.ListenerDiscoveryService.StreamListeners", api);
Config::Utility::checkLocalInfo("lds", local_info);
init_manager.registerTarget(*this, "LDS");

initialization_timeout_ =
std::chrono::milliseconds(PROTOBUF_GET_MS_OR_DEFAULT(lds_config, initial_fetch_timeout, 0));
if (initialization_timeout_.count() > 0) {
initialization_timeout_timer_ = dispatcher.createTimer([this]() -> void {
ENVOY_LOG(warn, "lds: initialization timed out");
onConfigUpdateFailed(nullptr);
});
}
}

void LdsApiImpl::initialize(std::function<void()> callback) {
initialize_callback_ = callback;
if (initialization_timeout_.count() > 0) {
ASSERT(initialization_timeout_timer_);
initialization_timeout_timer_->enableTimer(initialization_timeout_);
}
subscription_->start({}, *this);
}

Expand Down Expand Up @@ -95,6 +108,9 @@ void LdsApiImpl::runInitializeCallbackIfAny() {
initialize_callback_();
initialize_callback_ = nullptr;
}
if (initialization_timeout_timer_) {
initialization_timeout_timer_->disableTimer();
}
}

} // namespace Server
Expand Down
2 changes: 2 additions & 0 deletions source/server/lds_api.h
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,8 @@ class LdsApiImpl : public LdsApi,
Stats::ScopePtr scope_;
Upstream::ClusterManager& cm_;
std::function<void()> initialize_callback_;
Event::TimerPtr initialization_timeout_timer_;
std::chrono::milliseconds initialization_timeout_;
};

} // namespace Server
Expand Down
54 changes: 46 additions & 8 deletions test/common/upstream/cds_api_impl_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,20 @@ class CdsApiImplTest : public testing::Test {
CdsApiImplTest() : request_(&cm_.async_client_), api_(Api::createApiForTest(store_)) {}

void setup() {
envoy::api::v2::core::ConfigSource cds_config;
setupConfig(cds_config);
Upstream::ClusterManager::ClusterInfoMap cluster_map;
Upstream::MockClusterMockPrioritySet cluster;
setupClusters(cluster_map, cluster);

cds_ = CdsApiImpl::create(cds_config, cm_, dispatcher_, random_, local_info_, store_, *api_);
resetCdsInitializedCb();

expectRequest();
cds_->initialize();
}

void setupConfig(envoy::api::v2::core::ConfigSource& cds_config) {
const std::string config_json = R"EOF(
{
"cluster": {
Expand All @@ -43,23 +57,19 @@ class CdsApiImplTest : public testing::Test {
)EOF";

Json::ObjectSharedPtr config = Json::Factory::loadFromString(config_json);
envoy::api::v2::core::ConfigSource cds_config;
Config::Utility::translateCdsConfig(*config, cds_config);
cds_config.mutable_api_config_source()->set_api_type(
envoy::api::v2::core::ApiConfigSource::REST);
Upstream::ClusterManager::ClusterInfoMap cluster_map;
Upstream::MockClusterMockPrioritySet cluster;
}

void setupClusters(Upstream::ClusterManager::ClusterInfoMap& cluster_map,
Upstream::MockClusterMockPrioritySet& cluster) {
cluster_map.emplace("foo_cluster", cluster);
EXPECT_CALL(cm_, clusters()).WillOnce(Return(cluster_map));
EXPECT_CALL(cluster, info());
EXPECT_CALL(*cluster.info_, addedViaApi());
EXPECT_CALL(cluster, info());
EXPECT_CALL(*cluster.info_, type());
cds_ = CdsApiImpl::create(cds_config, cm_, dispatcher_, random_, local_info_, store_, *api_);
resetCdsInitializedCb();

expectRequest();
cds_->initialize();
}

void resetCdsInitializedCb() {
Expand Down Expand Up @@ -173,6 +183,7 @@ class CdsApiImplTest : public testing::Test {
Http::MockAsyncClientRequest request_;
CdsApiPtr cds_;
Event::MockTimer* interval_timer_;
Event::MockTimer* initialization_timeout_timer_;
Http::AsyncClient::Callbacks* callbacks_{};
ReadyWatcher initialized_;
Api::ApiPtr api_;
Expand Down Expand Up @@ -543,5 +554,32 @@ version_info: '0'
EXPECT_EQ(0UL, store_.gauge("cluster_manager.cds.version").value());
}

TEST_F(CdsApiImplTest, InitializationTimeout) {
initialization_timeout_timer_ = new Event::MockTimer(&dispatcher_);
interval_timer_ = new Event::MockTimer(&dispatcher_);

envoy::api::v2::core::ConfigSource cds_config;
setupConfig(cds_config);
cds_config.mutable_initial_fetch_timeout()->set_seconds(10);

Upstream::ClusterManager::ClusterInfoMap cluster_map;
Upstream::MockClusterMockPrioritySet cluster;
{
InSequence s;
setupClusters(cluster_map, cluster);
}

auto cds = CdsApiImpl::create(cds_config, cm_, dispatcher_, random_, local_info_, store_, *api_);
cds->setInitializedCb([this]() -> void { initialized_.ready(); });

EXPECT_CALL(*initialization_timeout_timer_, enableTimer(std::chrono::milliseconds(10 * 1000)));
EXPECT_CALL(initialized_, ready()).Times(0);
cds->initialize();

EXPECT_CALL(initialized_, ready());
EXPECT_CALL(*initialization_timeout_timer_, disableTimer());
initialization_timeout_timer_->callback_();
}

} // namespace Upstream
} // namespace Envoy