Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions libs/client-sdk/src/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,8 @@ target_sources(${LIBNAME} PRIVATE
data_sources/fdv2/streaming_synchronizer.cpp
data_sources/fdv2/fdv2_data_source.cpp
data_sources/fdv2/cache_initializer.cpp
data_sources/fdv2/source_factories.cpp
data_sources/fdv2/mode_sources.cpp
data_sources/fdv2/fdv1_adapter_synchronizer.cpp
data_sources/data_source_event_handler.cpp
data_sources/polling_data_source.cpp
Expand Down Expand Up @@ -52,6 +54,8 @@ target_sources(${LIBNAME} PRIVATE
data_sources/fdv2/streaming_synchronizer.hpp
data_sources/fdv2/fdv2_data_source.hpp
data_sources/fdv2/cache_initializer.hpp
data_sources/fdv2/source_factories.hpp
data_sources/fdv2/mode_sources.hpp
data_sources/fdv2/fdv1_adapter_synchronizer.hpp
flag_manager/flag_store.hpp
flag_manager/flag_updater.hpp
Expand Down
163 changes: 163 additions & 0 deletions libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,163 @@
#include "mode_sources.hpp"

#include "cache_initializer.hpp"
#include "source_factories.hpp"

#include <launchdarkly/config/shared/defaults.hpp>
#include <launchdarkly/serialization/json_context.hpp>

#include <boost/json.hpp>

#include <utility>
#include <variant>

namespace launchdarkly::client_side::data_sources {

namespace {

// Lets std::visit dispatch to a different lambda per variant alternative.
template <class... Ts>
struct overloaded : Ts... {
using Ts::operator()...;
};
template <class... Ts>
overloaded(Ts...) -> overloaded<Ts...>;

FDv2RequestConfig MakeRequestConfig(std::string base_url,
ModeSourceParams const& params,
std::string serialized_context,
bool use_post) {
return FDv2RequestConfig{std::move(base_url), params.http_properties,
std::move(serialized_context),
use_post ? FDv2ContextTransport::kPostBody
: FDv2ContextTransport::kGetPath,
params.with_reasons};
}

// Maps a wall-clock instant onto the monotonic clock the poll interval is timed
// against. steady_clock is used while running (for monotonicity) but cannot
// persist across a restart, so the last poll is stored as system_clock and
// mapped back here.
std::chrono::steady_clock::time_point ToSteadyClock(
std::chrono::system_clock::time_point instant) {
auto const now = std::chrono::system_clock::now();
return std::chrono::steady_clock::now() - (now - instant);
}

// The FDv1 fallback reuses the client's own FDv1 polling source, which is
// already pointed at the client SDK's FDv1 endpoints.
std::unique_ptr<IFDv2SynchronizerFactory> MakeFDv1Fallback(
FDv2Config::FDv1FallbackConfig const& fallback,
ModeSourceParams const& params) {
auto const defaults =
config::shared::Defaults<config::shared::ClientSDK>::PollingConfig();

config::shared::built::DataSourceConfig<config::shared::ClientSDK> const
fdv1_config{
config::shared::built::PollingConfig<config::shared::ClientSDK>{
fallback.poll_interval, defaults.polling_get_path,
defaults.polling_report_path, defaults.min_polling_interval},
params.with_reasons,
// FDv2 supersedes the REPORT transport, so the option is
// ignored when both are configured.
/* use_report= */ false};

auto endpoints =
fallback.base_url_override
? config::shared::built::
ServiceEndpoints{*fallback.base_url_override,
params.endpoints.StreamingBaseUrl(),
params.endpoints.EventsBaseUrl()}
: params.endpoints;

return std::make_unique<FDv1PollingAdapterFactory>(
params.executor, params.logger, std::move(endpoints), fdv1_config,
params.http_properties, params.context);
}

} // namespace

ModeSources BuildModeSources(FDv2Config const& config,
ConnectionMode mode,
ModeSourceParams const& params) {
ModeSources sources;

auto const definition = config.modes.find(mode);
if (definition == config.modes.end()) {
LD_LOG(params.logger, LogLevel::kError)
<< "fdv2: connection mode "
<< config::shared::GetConnectionModeName(mode)
<< " is not configured";
return sources;
}

auto const serialized_context =
boost::json::serialize(boost::json::value_from(params.context));

for (auto const& entry : definition->second.initializers) {
std::visit(
overloaded{
[&](FDv2Config::CacheConfig const&) {
sources.initializers.push_back(
std::make_unique<FDv2CacheInitializerFactory>(
params.cache, params.context, params.logger));
},
[&](FDv2Config::PollingConfig const& polling) {
sources.initializers.push_back(
std::make_unique<FDv2PollingInitializerFactory>(
params.executor, params.logger,
MakeRequestConfig(
polling.base_url_override.value_or(
params.polling_base_url),
params, serialized_context, config.use_post)));
},
},
entry);
}

// The interval a poll is rate limited against is measured from the last
// time this context was polled, which survives restarts so that repeated
// launches cannot produce a burst of requests.
std::optional<std::chrono::steady_clock::time_point> last_poll;
if (auto const freshness = params.cache->ReadFreshness(params.context)) {
last_poll = ToSteadyClock(*freshness);
}

for (auto const& entry : definition->second.synchronizers) {
std::visit(
overloaded{
[&](FDv2Config::PollingConfig const& polling) {
sources.synchronizers.push_back(
std::make_unique<FDv2PollingSynchronizerFactory>(
params.executor, params.logger,
MakeRequestConfig(
polling.base_url_override.value_or(
params.polling_base_url),
params, serialized_context, config.use_post),
polling.poll_interval, last_poll));
},
[&](FDv2Config::StreamingConfig const& streaming) {
sources.synchronizers.push_back(
std::make_unique<FDv2StreamingSynchronizerFactory>(
params.executor, params.logger,
MakeRequestConfig(
streaming.base_url_override.value_or(
params.streaming_base_url),
params, serialized_context, config.use_post),
MakeRequestConfig(params.polling_base_url, params,
serialized_context,
config.use_post),
streaming.initial_reconnect_delay));
},
},
entry);
}

if (auto const& fallback = definition->second.fdv1_fallback) {
sources.synchronizers.push_back(MakeFDv1Fallback(*fallback, params));
}

return sources;
}

} // namespace launchdarkly::client_side::data_sources
64 changes: 64 additions & 0 deletions libs/client-sdk/src/data_sources/fdv2/mode_sources.hpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
#pragma once

#include "ifdv2_initializer_factory.hpp"
#include "ifdv2_synchronizer_factory.hpp"

#include "../../flag_manager/flag_persistence.hpp"

#include <launchdarkly/config/shared/built/fdv2_config.hpp>
#include <launchdarkly/config/shared/built/http_properties.hpp>
#include <launchdarkly/config/shared/built/service_endpoints.hpp>
#include <launchdarkly/context.hpp>
#include <launchdarkly/logging/logger.hpp>

#include <boost/asio/any_io_executor.hpp>

#include <memory>
#include <string>
#include <vector>

namespace launchdarkly::client_side::data_sources {

using ConnectionMode = config::shared::ConnectionMode;
using FDv2Config = config::shared::built::FDv2Config<config::shared::ClientSDK>;

/** The factories one connection mode calls for. */
struct ModeSources {
std::vector<std::unique_ptr<IFDv2InitializerFactory>> initializers;
std::vector<std::unique_ptr<IFDv2SynchronizerFactory>> synchronizers;
};

/**
* Everything the sources of a mode need that the mode itself does not say:
* where to send requests, which context to evaluate, and where the cache is.
*/
struct ModeSourceParams {
boost::asio::any_io_executor executor;
Logger logger;
std::string polling_base_url;
std::string streaming_base_url;
config::shared::built::HttpProperties http_properties;
/**
* Where the FDv1 fallback source sends its requests. FDv2 sources use
* the resolved base URLs above instead.
*/
config::shared::built::ServiceEndpoints endpoints;
Context context;
/** Whether the application asked for evaluation reasons. */
bool with_reasons;
/** Non-owning. Must outlive sources built from these params. */
flag_manager::FlagPersistence* cache;
};

/**
* Assembles the factories the given mode calls for. Returns empty lists if
* the configuration does not define the mode.
*
* A mode that configures an FDv1 fallback gets its synchronizer appended to
* the synchronizer list.
*/
ModeSources BuildModeSources(FDv2Config const& config,
ConnectionMode mode,
ModeSourceParams const& params);

} // namespace launchdarkly::client_side::data_sources
86 changes: 86 additions & 0 deletions libs/client-sdk/src/data_sources/fdv2/source_factories.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
#include "source_factories.hpp"

#include "../polling_data_source.hpp"
#include "fdv1_adapter_synchronizer.hpp"
#include "polling_initializer.hpp"
#include "polling_synchronizer.hpp"
#include "streaming_synchronizer.hpp"

#include <utility>

namespace launchdarkly::client_side::data_sources {

FDv2PollingInitializerFactory::FDv2PollingInitializerFactory(
boost::asio::any_io_executor executor,
Logger logger,
FDv2RequestConfig request_config)
: executor_(std::move(executor)),
logger_(std::move(logger)),
request_config_(std::move(request_config)) {}

std::unique_ptr<IFDv2Initializer> FDv2PollingInitializerFactory::Build() {
return std::make_unique<FDv2PollingInitializer>(executor_, logger_,
request_config_);
}

FDv2PollingSynchronizerFactory::FDv2PollingSynchronizerFactory(
boost::asio::any_io_executor executor,
Logger logger,
FDv2RequestConfig request_config,
std::chrono::seconds poll_interval,
std::optional<std::chrono::steady_clock::time_point> last_poll)
: executor_(std::move(executor)),
logger_(std::move(logger)),
request_config_(std::move(request_config)),
poll_interval_(poll_interval),
last_poll_(last_poll) {}

std::unique_ptr<IFDv2Synchronizer> FDv2PollingSynchronizerFactory::Build() {
return std::make_unique<FDv2PollingSynchronizer>(
executor_, logger_, request_config_, poll_interval_, last_poll_);
}

FDv2StreamingSynchronizerFactory::FDv2StreamingSynchronizerFactory(
boost::asio::any_io_executor executor,
Logger logger,
FDv2RequestConfig stream_config,
FDv2RequestConfig poll_config,
std::chrono::milliseconds initial_reconnect_delay)
: executor_(std::move(executor)),
logger_(std::move(logger)),
stream_config_(std::move(stream_config)),
poll_config_(std::move(poll_config)),
initial_reconnect_delay_(initial_reconnect_delay) {}

std::unique_ptr<IFDv2Synchronizer> FDv2StreamingSynchronizerFactory::Build() {
return std::make_unique<FDv2StreamingSynchronizer>(
executor_, logger_, stream_config_, poll_config_,
initial_reconnect_delay_);
}

FDv1PollingAdapterFactory::FDv1PollingAdapterFactory(
boost::asio::any_io_executor executor,
Logger logger,
config::shared::built::ServiceEndpoints endpoints,
config::shared::built::DataSourceConfig<config::shared::ClientSDK>
data_source_config,
config::shared::built::HttpProperties http_properties,
Context context)
: executor_(std::move(executor)),
logger_(std::move(logger)),
endpoints_(std::move(endpoints)),
data_source_config_(std::move(data_source_config)),
http_properties_(std::move(http_properties)),
context_(std::move(context)) {}

std::unique_ptr<IFDv2Synchronizer> FDv1PollingAdapterFactory::Build() {
return std::make_unique<FDv1AdapterSynchronizer>(
[this](IDataSourceUpdateSink* sink,
DataSourceStatusManager* status_manager) {
return std::make_shared<PollingDataSource>(
endpoints_, data_source_config_, http_properties_, executor_,
context_, *sink, *status_manager, logger_);
});
}

} // namespace launchdarkly::client_side::data_sources
Loading
Loading