From 69b404506b7e25c739f66e86530b6026d13137bf Mon Sep 17 00:00:00 2001 From: Bee Klimt Date: Wed, 2 Sep 2026 03:32:23 +0000 Subject: [PATCH 1/6] feat: Discard the FDv2 selector when the evaluation context changes --- libs/client-sdk/src/client_impl.cpp | 4 + .../src/flag_manager/flag_manager.cpp | 4 + .../src/flag_manager/flag_manager.hpp | 7 + .../tests/fdv2_context_switch_test.cpp | 175 ++++++++++++++++++ 4 files changed, 190 insertions(+) create mode 100644 libs/client-sdk/tests/fdv2_context_switch_test.cpp diff --git a/libs/client-sdk/src/client_impl.cpp b/libs/client-sdk/src/client_impl.cpp index 82a09db35..fd5a60e02 100644 --- a/libs/client-sdk/src/client_impl.cpp +++ b/libs/client-sdk/src/client_impl.cpp @@ -160,6 +160,10 @@ static bool IsInitializedSuccessfully(DataSourceStatus::DataSourceState state) { std::future ClientImpl::IdentifyAsync(Context context) { UpdateContextSynchronized(context); + // A selector describes one context's data, so it is never carried over. + // Any flag data already loaded stays available for evaluation until a + // full data set arrives for the new context. + flag_manager_.ClearSelector(); flag_manager_.LoadCache(context); event_processor_->SendAsync(events::IdentifyEventParams{ std::chrono::system_clock::now(), std::move(context)}); diff --git a/libs/client-sdk/src/flag_manager/flag_manager.cpp b/libs/client-sdk/src/flag_manager/flag_manager.cpp index bdeba00b8..68a0334fa 100644 --- a/libs/client-sdk/src/flag_manager/flag_manager.cpp +++ b/libs/client-sdk/src/flag_manager/flag_manager.cpp @@ -37,4 +37,8 @@ void FlagManager::LoadCache(Context const& context) { persistence_updater_.LoadCached(context); } +void FlagManager::ClearSelector() { + flag_store_.ClearSelector(); +} + } // namespace launchdarkly::client_side::flag_manager diff --git a/libs/client-sdk/src/flag_manager/flag_manager.hpp b/libs/client-sdk/src/flag_manager/flag_manager.hpp index ad09bb55e..ec08a057e 100644 --- a/libs/client-sdk/src/flag_manager/flag_manager.hpp +++ b/libs/client-sdk/src/flag_manager/flag_manager.hpp @@ -28,6 +28,13 @@ class FlagManager { void LoadCache(Context const& context); + /** + * Forgets the selector for the data currently held, leaving the data in + * place. Called when the evaluation context changes, since a selector + * describes one context's data and is never reused for another. + */ + void ClearSelector(); + private: FlagStore flag_store_; FlagUpdater flag_updater_; diff --git a/libs/client-sdk/tests/fdv2_context_switch_test.cpp b/libs/client-sdk/tests/fdv2_context_switch_test.cpp new file mode 100644 index 000000000..9999a2635 --- /dev/null +++ b/libs/client-sdk/tests/fdv2_context_switch_test.cpp @@ -0,0 +1,175 @@ +#include + +#include +#include + +#include +#include + +#include +#include +#include +#include + +using launchdarkly::Context; +using launchdarkly::ContextBuilder; +using launchdarkly::EvaluationDetailInternal; +using launchdarkly::EvaluationResult; +using launchdarkly::Value; +using launchdarkly::client_side::FlagChange; +using launchdarkly::client_side::FlagChangeSet; +using launchdarkly::client_side::ItemDescriptor; +using launchdarkly::client_side::flag_manager::FlagManager; +using launchdarkly::client_side::flag_manager::PersistenceEncodeKey; +using launchdarkly::data_model::ChangeSetType; +using launchdarkly::data_model::Selector; + +namespace { + +class TestPersistence : public IPersistence { + public: + using StoreType = + std::map>>; + + explicit TestPersistence(StoreType store) : store_(std::move(store)) {} + + void Set(std::string storageNamespace, + std::string key, + std::string data) noexcept override { + store_[storageNamespace][key] = data; + } + + void Remove(std::string storageNamespace, + std::string key) noexcept override { + store_[storageNamespace].erase(key); + } + + std::optional Read(std::string storageNamespace, + std::string key) noexcept override { + return store_[storageNamespace][key]; + } + + StoreType store_; +}; + +ItemDescriptor Flag(std::uint64_t version, Value value) { + return ItemDescriptor{ + EvaluationResult{version, std::nullopt, false, false, std::nullopt, + EvaluationDetailInternal{std::move(value), + std::nullopt, std::nullopt}}}; +} + +char const* const kEnvironment = + "LaunchDarkly_rUTcjlHPv6Vegd27YmtGYkEGkEUGaEbn5M0JYTFQUpA="; + +} // namespace + +// A selector names a state the service can compute changes against for one +// context. Sending it for another would ask for the wrong delta. +TEST(FDv2ContextSwitchTest, ClearSelectorForgetsTheBasisButKeepsTheData) { + auto logger = launchdarkly::logging::NullLogger(); + FlagManager flag_manager("the-key", logger, 5, nullptr); + + flag_manager.Updater().Apply( + ContextBuilder().Kind("user", "first").Build(), + FlagChangeSet{ChangeSetType::kFull, + {FlagChange{"flagA", Flag(1, Value("a"))}}, + Selector{Selector::State{1, "state-1"}}}, + /* from_cache= */ false); + + ASSERT_TRUE(flag_manager.Store().CurrentSelector().value.has_value()); + + flag_manager.ClearSelector(); + + EXPECT_FALSE(flag_manager.Store().CurrentSelector().value.has_value()); + ASSERT_TRUE(flag_manager.Store().Get("flagA")); + EXPECT_EQ(Value("a"), + flag_manager.Store().Get("flagA")->item->Detail().Value()); +} + +// Nothing is available to evaluate against for the new context yet, so the +// previous context's data has to stay until a full data set arrives. +TEST(FDv2ContextSwitchTest, CacheMissRetainsThePreviousContextsData) { + auto logger = launchdarkly::logging::NullLogger(); + auto persistence = + std::make_shared(TestPersistence::StoreType()); + FlagManager flag_manager("the-key", logger, 5, persistence); + + auto first = ContextBuilder().Kind("user", "first").Build(); + flag_manager.Updater().Apply( + first, + FlagChangeSet{ChangeSetType::kFull, + {FlagChange{"flagA", Flag(1, Value("first-value"))}}, + Selector{Selector::State{1, "state-1"}}}, + /* from_cache= */ false); + + auto second = ContextBuilder().Kind("user", "second").Build(); + flag_manager.ClearSelector(); + flag_manager.LoadCache(second); + + ASSERT_TRUE(flag_manager.Store().Get("flagA")); + EXPECT_EQ(Value("first-value"), + flag_manager.Store().Get("flagA")->item->Detail().Value()); +} + +TEST(FDv2ContextSwitchTest, CacheHitReplacesThePreviousContextsData) { + auto logger = launchdarkly::logging::NullLogger(); + auto second = ContextBuilder().Kind("user", "second").Build(); + auto persistence = + std::make_shared(TestPersistence::StoreType{ + {kEnvironment, + {{PersistenceEncodeKey(second.CanonicalKey()), + R"({"flagB":{"version":1,"value":"second-value"}})"}}}}); + FlagManager flag_manager("the-key", logger, 5, persistence); + + flag_manager.Updater().Apply( + ContextBuilder().Kind("user", "first").Build(), + FlagChangeSet{ChangeSetType::kFull, + {FlagChange{"flagA", Flag(1, Value("first-value"))}}, + Selector{Selector::State{1, "state-1"}}}, + /* from_cache= */ false); + + flag_manager.ClearSelector(); + flag_manager.LoadCache(second); + + EXPECT_FALSE(flag_manager.Store().Get("flagA")); + ASSERT_TRUE(flag_manager.Store().Get("flagB")); + EXPECT_EQ(Value("second-value"), + flag_manager.Store().Get("flagB")->item->Detail().Value()); +} + +// The cache initializer for a new context reports a hit or a miss for that +// context alone, whatever the store currently holds. +TEST(FDv2ContextSwitchTest, CacheInitializerReadsTheNewContext) { + auto logger = launchdarkly::logging::NullLogger(); + auto first = ContextBuilder().Kind("user", "first").Build(); + auto second = ContextBuilder().Kind("user", "second").Build(); + auto persistence = + std::make_shared(TestPersistence::StoreType{ + {kEnvironment, + {{PersistenceEncodeKey(first.CanonicalKey()), + R"({"flagA":{"version":1,"value":"first-value"}})"}}}}); + FlagManager flag_manager("the-key", logger, 5, persistence); + + using launchdarkly::client_side::data_sources::FDv2CacheInitializer; + using launchdarkly::client_side::data_sources::FDv2SourceResult; + + auto second_result = + FDv2CacheInitializer(&flag_manager.Cache(), second, logger) + .Run() + .GetResult(); + auto* second_change_set = + std::get_if(&second_result->value); + ASSERT_NE(nullptr, second_change_set); + EXPECT_EQ(ChangeSetType::kNone, second_change_set->change_set.type); + + auto first_result = + FDv2CacheInitializer(&flag_manager.Cache(), first, logger) + .Run() + .GetResult(); + auto* first_change_set = + std::get_if(&first_result->value); + ASSERT_NE(nullptr, first_change_set); + EXPECT_EQ(ChangeSetType::kFull, first_change_set->change_set.type); +} From 40375cbe9902a7ecc44b3919e401e408182f3d75 Mon Sep 17 00:00:00 2001 From: Bee Klimt Date: Wed, 2 Sep 2026 03:40:10 +0000 Subject: [PATCH 2/6] feat: Honor the service's FDv1 fallback directive in the client --- libs/client-sdk/src/CMakeLists.txt | 2 + .../fdv2/fdv1_adapter_synchronizer.cpp | 206 ++++++++++++++++ .../fdv2/fdv1_adapter_synchronizer.hpp | 128 ++++++++++ .../data_sources/fdv2/fdv2_data_source.cpp | 97 +++++++- .../data_sources/fdv2/fdv2_data_source.hpp | 14 ++ .../tests/fdv1_adapter_synchronizer_test.cpp | 230 ++++++++++++++++++ .../tests/fdv2_data_source_test.cpp | 155 ++++++++++++ 7 files changed, 831 insertions(+), 1 deletion(-) create mode 100644 libs/client-sdk/src/data_sources/fdv2/fdv1_adapter_synchronizer.cpp create mode 100644 libs/client-sdk/src/data_sources/fdv2/fdv1_adapter_synchronizer.hpp create mode 100644 libs/client-sdk/tests/fdv1_adapter_synchronizer_test.cpp diff --git a/libs/client-sdk/src/CMakeLists.txt b/libs/client-sdk/src/CMakeLists.txt index 625775ca5..90e74c4fa 100644 --- a/libs/client-sdk/src/CMakeLists.txt +++ b/libs/client-sdk/src/CMakeLists.txt @@ -23,6 +23,7 @@ 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/fdv1_adapter_synchronizer.cpp data_sources/data_source_event_handler.cpp data_sources/polling_data_source.cpp flag_manager/flag_store.cpp @@ -51,6 +52,7 @@ 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/fdv1_adapter_synchronizer.hpp flag_manager/flag_store.hpp flag_manager/flag_updater.hpp bindings/c/sdk.cpp diff --git a/libs/client-sdk/src/data_sources/fdv2/fdv1_adapter_synchronizer.cpp b/libs/client-sdk/src/data_sources/fdv2/fdv1_adapter_synchronizer.cpp new file mode 100644 index 000000000..9468a807c --- /dev/null +++ b/libs/client-sdk/src/data_sources/fdv2/fdv1_adapter_synchronizer.cpp @@ -0,0 +1,206 @@ +#include "fdv1_adapter_synchronizer.hpp" + +#include + +namespace launchdarkly::client_side::data_sources { + +using DataSourceState = DataSourceStatus::DataSourceState; + +// ----- State ----- + +FDv1AdapterSynchronizer::State::State(async::Future closed) + : closed_(std::move(closed)) {} + +async::Future FDv1AdapterSynchronizer::State::GetNext() { + std::lock_guard lock(mutex_); + if (!result_queue_.empty()) { + auto result = std::move(result_queue_.front()); + result_queue_.pop_front(); + return async::MakeFuture(std::move(result)); + } + return pending_promise_.emplace().GetFuture(); +} + +void FDv1AdapterSynchronizer::State::ResolvePendingAsShutdown() { + std::optional> promise; + { + std::lock_guard lock(mutex_); + if (pending_promise_) { + promise = std::move(pending_promise_); + pending_promise_.reset(); + } + } + if (promise) { + promise->Resolve(FDv2SourceResult{FDv2SourceResult::Shutdown{}}); + } +} + +void FDv1AdapterSynchronizer::State::Notify(FDv2SourceResult result) { + std::optional> promise; + { + std::lock_guard lock(mutex_); + if (closed_.IsFinished()) { + return; + } + if (pending_promise_) { + promise = std::move(pending_promise_); + pending_promise_.reset(); + } else { + result_queue_.push_back(std::move(result)); + return; + } + } + // Resolve outside the lock. Promise::Resolve may invoke inline + // continuations that could call back into Notify or GetNext. + promise->Resolve(std::move(result)); +} + +// ----- ConvertingSink ----- + +FDv1AdapterSynchronizer::ConvertingSink::ConvertingSink( + std::weak_ptr state) + : state_(std::move(state)) {} + +void FDv1AdapterSynchronizer::ConvertingSink::Init( + Context const& /* context */, + std::unordered_map data) { + auto state = state_.lock(); + if (!state) { + return; + } + FlagChangeSetData changes; + changes.reserve(data.size()); + for (auto& [key, item] : data) { + changes.push_back(FlagChange{key, std::move(item)}); + } + state->Notify(FDv2SourceResult{FDv2SourceResult::ChangeSet{ + FlagChangeSet{data_model::ChangeSetType::kFull, std::move(changes), + data_model::Selector{}}}}); +} + +void FDv1AdapterSynchronizer::ConvertingSink::Upsert( + Context const& /* context */, + std::string key, + ItemDescriptor item) { + auto state = state_.lock(); + if (!state) { + return; + } + state->Notify(FDv2SourceResult{FDv2SourceResult::ChangeSet{ + FlagChangeSet{data_model::ChangeSetType::kPartial, + {FlagChange{std::move(key), std::move(item)}}, + data_model::Selector{}}}}); +} + +void FDv1AdapterSynchronizer::ConvertingSink::Apply( + Context const& /* context */, + FlagChangeSet change_set, + bool /* from_cache */) { + auto state = state_.lock(); + if (!state) { + return; + } + state->Notify( + FDv2SourceResult{FDv2SourceResult::ChangeSet{std::move(change_set)}}); +} + +// ----- FDv1AdapterSynchronizer ----- + +namespace { + +// Turns the wrapped source's status into the result the orchestrator acts on. +// A valid status carries no error and needs no result. The changeset that +// accompanied it already reported the recovery. +std::optional ResultForStatus( + DataSourceStatus const& status) { + auto const error = status.LastError(); + if (!error) { + return std::nullopt; + } + switch (status.State()) { + case DataSourceState::kInterrupted: + // An error encountered before the source ever became valid is + // reported as still initializing, but it is the same recoverable + // failure. + case DataSourceState::kInitializing: + return FDv2SourceResult{FDv2SourceResult::Interrupted{*error}}; + case DataSourceState::kShutdown: + return FDv2SourceResult{FDv2SourceResult::TerminalError{*error}}; + case DataSourceState::kValid: + case DataSourceState::kSetOffline: + return std::nullopt; + } + return std::nullopt; +} + +} // namespace + +FDv1AdapterSynchronizer::FDv1AdapterSynchronizer(SourceBuilder source_builder) + : state_(std::make_shared(close_promise_.GetFuture())), + sink_(std::make_shared(state_)), + status_manager_(std::make_shared()), + status_subscription_(status_manager_->OnDataSourceStatusChange( + [state = state_](DataSourceStatus status) { + if (auto result = ResultForStatus(status)) { + state->Notify(std::move(*result)); + } + })), + fdv1_source_(source_builder(sink_.get(), status_manager_.get())) {} + +FDv1AdapterSynchronizer::~FDv1AdapterSynchronizer() { + Close(); +} + +async::Future FDv1AdapterSynchronizer::Next( + data_model::Selector /* selector */) { + auto closed = close_promise_.GetFuture(); + if (closed.IsFinished()) { + return async::MakeFuture( + FDv2SourceResult{FDv2SourceResult::Shutdown{}}); + } + { + std::lock_guard lock(lifecycle_mutex_); + if (!started_) { + started_ = true; + fdv1_source_->Start(); + } + } + auto result_future = state_->GetNext(); + if (result_future.IsFinished()) { + return result_future; + } + return async::WhenAny(closed, result_future) + .Then( + [state = state_, result_future](std::size_t const& idx) mutable + -> async::Future { + if (idx == 0) { + state->ResolvePendingAsShutdown(); + return async::MakeFuture( + FDv2SourceResult{FDv2SourceResult::Shutdown{}}); + } + return result_future; + }, + async::kInlineExecutor); +} + +void FDv1AdapterSynchronizer::Close() { + if (!close_promise_.Resolve(std::monostate{})) { + return; + } + std::lock_guard lock(lifecycle_mutex_); + bool const was_started = started_; + started_ = true; + if (was_started) { + // The sink and status manager are captured so that they outlive any + // callback the source has already queued. + fdv1_source_->ShutdownAsync( + [sink = sink_, status = status_manager_, source = fdv1_source_] {}); + } +} + +std::string const& FDv1AdapterSynchronizer::Identity() const { + static std::string const identity = "FDv1 fallback adapter"; + return identity; +} + +} // namespace launchdarkly::client_side::data_sources diff --git a/libs/client-sdk/src/data_sources/fdv2/fdv1_adapter_synchronizer.hpp b/libs/client-sdk/src/data_sources/fdv2/fdv1_adapter_synchronizer.hpp new file mode 100644 index 000000000..2d5154301 --- /dev/null +++ b/libs/client-sdk/src/data_sources/fdv2/fdv1_adapter_synchronizer.hpp @@ -0,0 +1,128 @@ +#pragma once + +#include "../data_source.hpp" +#include "../data_source_status_manager.hpp" +#include "../data_source_update_sink.hpp" +#include "ifdv2_synchronizer.hpp" + +#include +#include +#include + +#include +#include +#include +#include +#include +#include + +namespace launchdarkly::client_side::data_sources { + +/** + * Presents an FDv1 data source as an FDv2 synchronizer, so that the + * orchestrator can run it while the service has directed the SDK away from + * FDv2. + * + * FDv1 has no selectors, so the changesets this reports carry none. The + * orchestrator therefore never asks the service for a delta against data + * FDv1 supplied. + * + * Thread safety: Next() and Close() may be called from any thread. Only one + * Next() may be outstanding at a time. + */ +class FDv1AdapterSynchronizer final : public IFDv2Synchronizer { + public: + /** + * Builds the wrapped FDv1 source. Called once during construction with + * the sink and status manager the source must report through, both of + * which the adapter keeps alive for the source's lifetime. + */ + using SourceBuilder = + std::function(IDataSourceUpdateSink*, + DataSourceStatusManager*)>; + + explicit FDv1AdapterSynchronizer(SourceBuilder source_builder); + + ~FDv1AdapterSynchronizer() override; + + async::Future Next( + data_model::Selector selector) override; + + void Close() override; + + [[nodiscard]] std::string const& Identity() const override; + + private: + /** + * Holds the result queue and the pending Next() promise. Shared with the + * wrapped source's sink and status subscription. Thread-safe. + */ + class State { + public: + explicit State(async::Future closed); + + async::Future GetNext(); + + /** + * Resolves any pending Next() with Shutdown and clears it, so that a + * caller abandoned by Close() is not left waiting. + */ + void ResolvePendingAsShutdown(); + + void Notify(FDv2SourceResult result); + + private: + // Finished once the owning adapter's Close() has run. Read in Notify + // to drop late results. + async::Future const closed_; + + mutable std::mutex mutex_; + // Both protected by mutex_. + std::optional> pending_promise_; + std::deque result_queue_; + }; + + /** + * Turns the FDv1 source's Init and Upsert calls into FDv2 changesets + * queued on State. Thread-safe (delegates to State). + */ + class ConvertingSink final : public IDataSourceUpdateSink { + public: + explicit ConvertingSink(std::weak_ptr state); + + void Init( + Context const& context, + std::unordered_map data) override; + void Upsert(Context const& context, + std::string key, + ItemDescriptor item) override; + void Apply(Context const& context, + FlagChangeSet change_set, + bool from_cache) override; + + private: + std::weak_ptr state_; + }; + + // Thread-safe primitive. Declared before state_ so state_'s constructor + // can take a future from it. + async::Promise close_promise_; + + // shared_ptr so async callbacks that fire after this is destroyed can + // hold their own reference. + std::shared_ptr const state_; + std::shared_ptr const sink_; + std::shared_ptr const status_manager_; + std::unique_ptr const status_subscription_; + + std::shared_ptr const fdv1_source_; + + // Serializes Start and ShutdownAsync on fdv1_source_ across concurrent + // Next() and Close() calls. + std::mutex lifecycle_mutex_; + // Protected by lifecycle_mutex_. Set when Next() starts the source, or + // when Close() runs first, so that a later Next() cannot start it. + bool started_ = false; +}; + +} // namespace launchdarkly::client_side::data_sources diff --git a/libs/client-sdk/src/data_sources/fdv2/fdv2_data_source.cpp b/libs/client-sdk/src/data_sources/fdv2/fdv2_data_source.cpp index 539ac0743..248d2c860 100644 --- a/libs/client-sdk/src/data_sources/fdv2/fdv2_data_source.cpp +++ b/libs/client-sdk/src/data_sources/fdv2/fdv2_data_source.cpp @@ -1,6 +1,7 @@ #include "fdv2_data_source.hpp" #include +#include #include @@ -11,6 +12,10 @@ namespace launchdarkly::client_side::data_sources { +static char const* const kNoFDv1FallbackConfigured = + "the service directed the SDK to FDv1, but no FDv1 fallback is " + "configured"; + namespace { // Lets std::visit dispatch to a different lambda per variant alternative. @@ -83,6 +88,7 @@ FDv2DataSource::~FDv2DataSource() { void FDv2DataSource::Close() { std::lock_guard lock(mutex_); closed_ = true; + fdv2_retry_cancel_.Cancel(); if (active_initializer_) { active_initializer_->Close(); } @@ -291,14 +297,31 @@ void FDv2DataSource::OnInitializerResult(FDv2SourceResult result) { }, result.value); + bool disconnected = false; { std::lock_guard lock(mutex_); active_initializer_.reset(); if (closed_ || got_shutdown) { return; } + if (result.fdv1_fallback) { + if (EngageFDv1FallbackLocked(*result.fdv1_fallback)) { + // No basis yet, but the FDv1 tier can supply one, so hand off + // to the synchronizer phase rather than continuing the chain. + got_basis = true; + } else { + disconnected = true; + } + } } + if (disconnected) { + status_manager_->SetState( + DataSourceStatus::DataSourceState::kInterrupted, + DataSourceStatus::ErrorInfo::ErrorKind::kUnknown, + kNoFDv1FallbackConfigured); + return; + } if (got_basis) { StartSynchronizers(); } else { @@ -478,6 +501,7 @@ void FDv2DataSource::OnSynchronizerResult(FDv2SourceResult result) { }, result.value); + bool disconnected = false; { std::lock_guard lock(mutex_); if (closed_ || got_shutdown) { @@ -485,13 +509,30 @@ void FDv2DataSource::OnSynchronizerResult(FDv2SourceResult result) { active_conditions_.reset(); return; } - if (advance) { + if (result.fdv1_fallback && + !source_manager_.IsCurrentSynchronizerFDv1Fallback()) { + active_synchronizer_.reset(); + active_conditions_.reset(); + if (EngageFDv1FallbackLocked(*result.fdv1_fallback)) { + advance = true; + } else { + advance = false; + disconnected = true; + } + } else if (advance) { source_manager_.BlockCurrentSynchronizer(); active_synchronizer_.reset(); active_conditions_.reset(); } } + if (disconnected) { + status_manager_->SetState( + DataSourceStatus::DataSourceState::kInterrupted, + DataSourceStatus::ErrorInfo::ErrorKind::kUnknown, + kNoFDv1FallbackConfigured); + return; + } if (advance) { StartSynchronizers(); } else { @@ -499,6 +540,60 @@ void FDv2DataSource::OnSynchronizerResult(FDv2SourceResult result) { } } +bool FDv2DataSource::EngageFDv1FallbackLocked( + FDv1FallbackDirective const& directive) { + source_manager_.SwitchToFDv1Fallback(); + + // Cancel any attempt already scheduled and start fresh. A + // CancellationSource is one-shot, so reusing it would leak the prior + // timer. + fdv2_retry_cancel_.Cancel(); + fdv2_retry_cancel_ = async::CancellationSource{}; + LD_LOG(logger_, LogLevel::kInfo) + << "fdv2: will attempt FDv2 again in " << directive.ttl.count() << "s"; + async::Delay(executor_, directive.ttl, fdv2_retry_cancel_.GetToken()) + .Then( + [weak = weak_from_this()](bool const& fired) -> std::monostate { + if (!fired) { + return {}; + } + if (auto self = weak.lock()) { + self->OnFDv2RetryTimer(); + } + return {}; + }, + [executor = executor_](async::Continuation work) { + boost::asio::post(executor, std::move(work)); + }); + + bool const available = source_manager_.AvailableSynchronizerCount() > 0; + if (available) { + LD_LOG(logger_, LogLevel::kInfo) + << "fdv2: falling back to the FDv1 synchronizer"; + } else { + LD_LOG(logger_, LogLevel::kWarn) + << "fdv2: " << kNoFDv1FallbackConfigured; + } + return available; +} + +void FDv2DataSource::OnFDv2RetryTimer() { + { + std::lock_guard lock(mutex_); + if (closed_) { + return; + } + LD_LOG(logger_, LogLevel::kInfo) << "fdv2: re-attempting FDv2"; + source_manager_.SwitchBackToFDv2(); + if (active_synchronizer_) { + active_synchronizer_->Close(); + active_synchronizer_.reset(); + } + active_conditions_.reset(); + } + StartSynchronizers(); +} + void FDv2DataSource::ApplyResult(FDv2SourceResult::ChangeSet change_set, std::optional environment_id, bool from_cache) { diff --git a/libs/client-sdk/src/data_sources/fdv2/fdv2_data_source.hpp b/libs/client-sdk/src/data_sources/fdv2/fdv2_data_source.hpp index f5baf8280..18dc580a3 100644 --- a/libs/client-sdk/src/data_sources/fdv2/fdv2_data_source.hpp +++ b/libs/client-sdk/src/data_sources/fdv2/fdv2_data_source.hpp @@ -8,6 +8,7 @@ #include "../../flag_manager/flag_store.hpp" +#include #include #include #include @@ -193,6 +194,15 @@ class FDv2DataSource final // reflects why. void ReportExhausted(bool any_synchronizers_configured); + // Moves the synchronizer list onto the FDv1 tier and schedules the + // attempt to return to FDv2. Called with mutex_ held. Returns true if an + // FDv1 synchronizer is available to start. + bool EngageFDv1FallbackLocked(FDv1FallbackDirective const& directive); + + // Switches the synchronizer list back to FDv2 and restarts the + // synchronizer phase. + void OnFDv2RetryTimer(); + // Logger is itself thread-safe and cheap to copy. Logger logger_; @@ -233,6 +243,10 @@ class FDv2DataSource final // Any outstanding work to complete before signaling shutdown complete. async::Future closing_ = async::MakeFuture({}); + // Cancelled in Close() to abort a pending attempt to return to FDv2, and + // replaced when a new attempt is scheduled. The replacement is why this + // needs mutex_ despite the source being internally synchronized. + async::CancellationSource fdv2_retry_cancel_; }; } // namespace launchdarkly::client_side::data_sources diff --git a/libs/client-sdk/tests/fdv1_adapter_synchronizer_test.cpp b/libs/client-sdk/tests/fdv1_adapter_synchronizer_test.cpp new file mode 100644 index 000000000..1070d2ede --- /dev/null +++ b/libs/client-sdk/tests/fdv1_adapter_synchronizer_test.cpp @@ -0,0 +1,230 @@ +#include + +#include + +#include +#include +#include + +#include +#include + +#include +#include +#include + +using launchdarkly::ContextBuilder; +using launchdarkly::EvaluationDetailInternal; +using launchdarkly::EvaluationResult; +using launchdarkly::Value; +using launchdarkly::client_side::IDataSourceUpdateSink; +using launchdarkly::client_side::ItemDescriptor; + +using namespace launchdarkly::client_side::data_sources; + +namespace { + +ItemDescriptor Flag(std::uint64_t version, Value value) { + return ItemDescriptor{ + EvaluationResult{version, std::nullopt, false, false, std::nullopt, + EvaluationDetailInternal{std::move(value), + std::nullopt, std::nullopt}}}; +} + +// Stands in for an FDv1 data source. Records its lifecycle and lets tests +// drive the sink and status manager it was handed. +class FakeFDv1Source : public IDataSource { + public: + FakeFDv1Source(IDataSourceUpdateSink* sink, + DataSourceStatusManager* status_manager) + : sink_(sink), status_manager_(status_manager) {} + + void Start() override { ++start_count_; } + + void ShutdownAsync(std::function completion) override { + ++shutdown_count_; + if (completion) { + completion(); + } + } + + IDataSourceUpdateSink* Sink() { return sink_; } + DataSourceStatusManager* StatusManager() { return status_manager_; } + + int start_count_ = 0; + int shutdown_count_ = 0; + + private: + IDataSourceUpdateSink* const sink_; + DataSourceStatusManager* const status_manager_; +}; + +// Builds an adapter over a FakeFDv1Source and keeps a handle on it. +class AdapterFixture { + public: + AdapterFixture() { + adapter_ = std::make_unique( + [this](IDataSourceUpdateSink* sink, + DataSourceStatusManager* status_manager) { + auto source = + std::make_shared(sink, status_manager); + source_ = source; + return source; + }); + } + + FDv1AdapterSynchronizer& Adapter() { return *adapter_; } + FakeFDv1Source& Source() { return *source_; } + + std::optional Next() { + return adapter_->Next(launchdarkly::data_model::Selector{}) + .WaitForResult(std::chrono::seconds{2}); + } + + private: + std::shared_ptr source_; + std::unique_ptr adapter_; +}; + +} // namespace + +TEST(FDv1AdapterSynchronizerTest, TheFirstNextStartsTheWrappedSource) { + AdapterFixture f; + + EXPECT_EQ(0, f.Source().start_count_); + + f.Source().Sink()->Init(ContextBuilder().Kind("user", "user-key").Build(), + {{"flagA", Flag(1, Value("a"))}}); + f.Next(); + + EXPECT_EQ(1, f.Source().start_count_); +} + +TEST(FDv1AdapterSynchronizerTest, InitBecomesAFullChangeSet) { + AdapterFixture f; + + f.Source().Sink()->Init(ContextBuilder().Kind("user", "user-key").Build(), + {{"flagA", Flag(1, Value("a"))}}); + auto result = f.Next(); + + ASSERT_TRUE(result.has_value()); + auto* change_set = std::get_if(&result->value); + ASSERT_NE(nullptr, change_set); + EXPECT_EQ(launchdarkly::data_model::ChangeSetType::kFull, + change_set->change_set.type); + ASSERT_EQ(1u, change_set->change_set.data.size()); + EXPECT_EQ("flagA", change_set->change_set.data[0].key); +} + +TEST(FDv1AdapterSynchronizerTest, UpsertBecomesAPartialChangeSet) { + AdapterFixture f; + + f.Source().Sink()->Upsert(ContextBuilder().Kind("user", "user-key").Build(), + "flagA", Flag(2, Value("a2"))); + auto result = f.Next(); + + ASSERT_TRUE(result.has_value()); + auto* change_set = std::get_if(&result->value); + ASSERT_NE(nullptr, change_set); + EXPECT_EQ(launchdarkly::data_model::ChangeSetType::kPartial, + change_set->change_set.type); + ASSERT_EQ(1u, change_set->change_set.data.size()); + EXPECT_EQ("flagA", change_set->change_set.data[0].key); +} + +// FDv1 has no selectors, so the orchestrator must never end up asking the +// service for a delta against data FDv1 supplied. +TEST(FDv1AdapterSynchronizerTest, ChangeSetsCarryNoSelector) { + AdapterFixture f; + + f.Source().Sink()->Init(ContextBuilder().Kind("user", "user-key").Build(), + {{"flagA", Flag(1, Value("a"))}}); + auto result = f.Next(); + + auto* change_set = std::get_if(&result->value); + ASSERT_NE(nullptr, change_set); + EXPECT_FALSE(change_set->change_set.selector.value.has_value()); +} + +TEST(FDv1AdapterSynchronizerTest, ARecoverableErrorBecomesInterrupted) { + AdapterFixture f; + + f.Source().StatusManager()->SetState( + DataSourceStatus::DataSourceState::kInterrupted, + DataSourceStatus::ErrorInfo::ErrorKind::kNetworkError, "boom"); + auto result = f.Next(); + + ASSERT_TRUE(result.has_value()); + EXPECT_TRUE( + std::holds_alternative(result->value)); +} + +TEST(FDv1AdapterSynchronizerTest, AnUnrecoverableErrorBecomesTerminal) { + AdapterFixture f; + + f.Source().StatusManager()->SetState( + DataSourceStatus::DataSourceState::kShutdown, + DataSourceStatus::ErrorInfo::ErrorKind::kErrorResponse, "unauthorized"); + auto result = f.Next(); + + ASSERT_TRUE(result.has_value()); + EXPECT_TRUE( + std::holds_alternative(result->value)); +} + +// The changeset that accompanies a recovery already reports it. +TEST(FDv1AdapterSynchronizerTest, BecomingValidReportsNothing) { + AdapterFixture f; + + f.Source().StatusManager()->SetState( + DataSourceStatus::DataSourceState::kValid); + + auto future = f.Adapter().Next(launchdarkly::data_model::Selector{}); + EXPECT_FALSE(future.IsFinished()); +} + +TEST(FDv1AdapterSynchronizerTest, CloseShutsDownTheWrappedSource) { + AdapterFixture f; + + f.Source().Sink()->Init(ContextBuilder().Kind("user", "user-key").Build(), + {}); + f.Next(); + f.Adapter().Close(); + + EXPECT_EQ(1, f.Source().shutdown_count_); +} + +// Nothing was started, so there is nothing to shut down. +TEST(FDv1AdapterSynchronizerTest, CloseBeforeAnyNextDoesNotShutDown) { + AdapterFixture f; + + f.Adapter().Close(); + + EXPECT_EQ(0, f.Source().shutdown_count_); + EXPECT_EQ(0, f.Source().start_count_); +} + +TEST(FDv1AdapterSynchronizerTest, NextAfterCloseIsShutdown) { + AdapterFixture f; + + f.Adapter().Close(); + auto result = f.Next(); + + ASSERT_TRUE(result.has_value()); + EXPECT_TRUE( + std::holds_alternative(result->value)); + EXPECT_EQ(0, f.Source().start_count_); +} + +TEST(FDv1AdapterSynchronizerTest, CloseUnblocksAPendingNext) { + AdapterFixture f; + + auto future = f.Adapter().Next(launchdarkly::data_model::Selector{}); + ASSERT_FALSE(future.IsFinished()); + + f.Adapter().Close(); + + ASSERT_TRUE(future.IsFinished()); + EXPECT_TRUE(std::holds_alternative( + future.GetResult()->value)); +} diff --git a/libs/client-sdk/tests/fdv2_data_source_test.cpp b/libs/client-sdk/tests/fdv2_data_source_test.cpp index bf365618d..0cbdd0b9c 100644 --- a/libs/client-sdk/tests/fdv2_data_source_test.cpp +++ b/libs/client-sdk/tests/fdv2_data_source_test.cpp @@ -169,6 +169,16 @@ class MultiShotSynchronizerFactory : public IFDv2SynchronizerFactory { std::vector> sources_; }; +// Stands in for the FDv1 tier, which the orchestrator keeps in reserve. +class FDv1FallbackFactory : public MultiShotSynchronizerFactory { + public: + explicit FDv1FallbackFactory( + std::vector> sources) + : MultiShotSynchronizerFactory(std::move(sources)) {} + + [[nodiscard]] bool IsFDv1Fallback() const override { return true; } +}; + // Initializer whose Run() stays pending until Deliver() resolves it, so // orchestration can be examined in flight. class StalledInitializer : public IFDv2Initializer { @@ -280,6 +290,10 @@ class Harness { } boost::asio::io_context& Context() { return ioc_; } + + // Runs the orchestration to a standstill without waiting on timers, for + // tests where a scheduled retry is not the point. + void Drain() { ioc_.poll(); } DataSourceStatusManager& StatusManager() { return status_manager_; } flag_manager::FlagStore const& Store() { return flag_manager_.Store(); } @@ -801,3 +815,144 @@ TEST(ClientFDv2DataSourceTest, CacheSourcedDataIsMarkedAsSuch) { EXPECT_TRUE(h.Applies()[0]); EXPECT_FALSE(h.Applies()[1]); } + +// ============================================================================ +// FDv1 fallback +// ============================================================================ + +namespace { + +FDv2SourceResult WithFallbackDirective(FDv2SourceResult result, + std::chrono::seconds ttl) { + result.fdv1_fallback = FDv1FallbackDirective{ttl}; + return result; +} + +} // namespace + +// The payload that arrived alongside the directive is still applied before the +// SDK moves off FDv2. +TEST(ClientFDv2DataSourceTest, FallbackDirectiveOnAnInitializerAppliesItsData) { + Harness h; + + std::vector> initializers; + initializers.push_back(std::make_unique( + std::make_unique(WithFallbackDirective( + MakeChangeSetResult(data_model::ChangeSetType::kFull, + {FlagChange{"flagA", MakeFlag(1, Value("a"))}}, + MakeSelector(1, "state-1")), + std::chrono::seconds{3600})))); + + std::vector> synchronizers; + synchronizers.push_back(std::make_unique( + std::make_unique(std::vector{}))); + std::vector> fdv1_sources; + fdv1_sources.push_back( + std::make_unique(std::vector{})); + synchronizers.push_back( + std::make_unique(std::move(fdv1_sources))); + + auto source = + h.MakeDataSource(std::move(initializers), std::move(synchronizers)); + source->Start(); + h.Drain(); + + EXPECT_TRUE(h.Store().Get("flagA")); +} + +TEST(ClientFDv2DataSourceTest, FallbackDirectiveStartsTheFDv1Tier) { + Harness h; + + std::vector fdv2_results; + fdv2_results.push_back(WithFallbackDirective( + MakeChangeSetResult(data_model::ChangeSetType::kFull, + {FlagChange{"flagA", MakeFlag(1, Value("a"))}}, + MakeSelector(1, "state-1")), + std::chrono::seconds{3600})); + + std::vector fdv1_results; + fdv1_results.push_back( + MakeChangeSetResult(data_model::ChangeSetType::kFull, + {FlagChange{"from-fdv1", MakeFlag(1, Value("b"))}}, + data_model::Selector{})); + + std::vector> synchronizers; + synchronizers.push_back(std::make_unique( + std::make_unique(std::move(fdv2_results)))); + std::vector> fdv1_sources; + fdv1_sources.push_back( + std::make_unique(std::move(fdv1_results))); + auto fdv1 = std::make_unique(std::move(fdv1_sources)); + auto* fdv1_ptr = fdv1.get(); + synchronizers.push_back(std::move(fdv1)); + + auto source = h.MakeDataSource({}, std::move(synchronizers)); + source->Start(); + h.Drain(); + + EXPECT_EQ(1, fdv1_ptr->build_count_); + EXPECT_TRUE(h.Store().Get("from-fdv1")); +} + +// Continuing to attempt FDv2 after the service has said not to would be +// pointless, so the SDK disconnects instead. +TEST(ClientFDv2DataSourceTest, FallbackWithNoFDv1TierDisconnects) { + Harness h; + + std::vector fdv2_results; + fdv2_results.push_back(WithFallbackDirective( + MakeChangeSetResult(data_model::ChangeSetType::kFull, + {FlagChange{"flagA", MakeFlag(1, Value("a"))}}, + MakeSelector(1, "state-1")), + std::chrono::seconds{3600})); + + std::vector> synchronizers; + synchronizers.push_back(std::make_unique( + std::make_unique(std::move(fdv2_results)))); + + auto source = h.MakeDataSource({}, std::move(synchronizers)); + source->Start(); + h.Drain(); + + EXPECT_EQ(DataSourceStatus::DataSourceState::kInterrupted, h.State()); + // Flag data already received stays available for evaluation. + EXPECT_TRUE(h.Store().Get("flagA")); +} + +// Once the TTL elapses the SDK tries FDv2 again, so a fallback is never +// permanent. +TEST(ClientFDv2DataSourceTest, FDv2IsRetriedAfterTheFallbackTtl) { + Harness h; + + std::vector first_fdv2; + first_fdv2.push_back(WithFallbackDirective( + MakeChangeSetResult(data_model::ChangeSetType::kFull, + {FlagChange{"flagA", MakeFlag(1, Value("a"))}}, + MakeSelector(1, "state-1")), + std::chrono::seconds{1})); + + std::vector> fdv2_sources; + fdv2_sources.push_back( + std::make_unique(std::move(first_fdv2))); + fdv2_sources.push_back( + std::make_unique(std::vector{})); + + std::vector> fdv1_sources; + fdv1_sources.push_back(std::make_unique( + std::vector{}, nullptr, nullptr, + /* stall_after_results= */ true)); + + std::vector> synchronizers; + auto fdv2 = + std::make_unique(std::move(fdv2_sources)); + auto* fdv2_ptr = fdv2.get(); + synchronizers.push_back(std::move(fdv2)); + synchronizers.push_back( + std::make_unique(std::move(fdv1_sources))); + + auto source = h.MakeDataSource({}, std::move(synchronizers)); + source->Start(); + h.Context().run(); + + EXPECT_EQ(2, fdv2_ptr->build_count_); +} From 1c261fe1571aa4b82416262d8292a8213337c753 Mon Sep 17 00:00:00 2001 From: Bee Klimt Date: Wed, 2 Sep 2026 03:30:25 +0000 Subject: [PATCH 3/6] feat: Represent the client FDv2 connection modes and assemble their sources --- libs/client-sdk/src/CMakeLists.txt | 4 + .../src/data_sources/fdv2/mode_sources.cpp | 132 +++++++++++++++ .../src/data_sources/fdv2/mode_sources.hpp | 60 +++++++ .../data_sources/fdv2/source_factories.cpp | 59 +++++++ .../data_sources/fdv2/source_factories.hpp | 96 +++++++++++ .../tests/fdv2_mode_sources_test.cpp | 160 ++++++++++++++++++ .../config/shared/built/fdv2_config.hpp | 97 +++++++++++ .../config/shared/connection_mode.hpp | 23 +++ .../launchdarkly/config/shared/defaults.hpp | 47 +++++ libs/common/src/CMakeLists.txt | 1 + libs/common/src/config/connection_mode.cpp | 17 ++ 11 files changed, 696 insertions(+) create mode 100644 libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp create mode 100644 libs/client-sdk/src/data_sources/fdv2/mode_sources.hpp create mode 100644 libs/client-sdk/src/data_sources/fdv2/source_factories.cpp create mode 100644 libs/client-sdk/src/data_sources/fdv2/source_factories.hpp create mode 100644 libs/client-sdk/tests/fdv2_mode_sources_test.cpp create mode 100644 libs/common/include/launchdarkly/config/shared/built/fdv2_config.hpp create mode 100644 libs/common/include/launchdarkly/config/shared/connection_mode.hpp create mode 100644 libs/common/src/config/connection_mode.cpp diff --git a/libs/client-sdk/src/CMakeLists.txt b/libs/client-sdk/src/CMakeLists.txt index 90e74c4fa..8db6518b4 100644 --- a/libs/client-sdk/src/CMakeLists.txt +++ b/libs/client-sdk/src/CMakeLists.txt @@ -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 @@ -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 diff --git a/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp b/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp new file mode 100644 index 000000000..2d3bfa98c --- /dev/null +++ b/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp @@ -0,0 +1,132 @@ +#include "mode_sources.hpp" + +#include "cache_initializer.hpp" +#include "source_factories.hpp" + +#include + +#include + +#include +#include + +namespace launchdarkly::client_side::data_sources { + +namespace { + +// Lets std::visit dispatch to a different lambda per variant alternative. +template +struct overloaded : Ts... { + using Ts::operator()...; +}; +template +overloaded(Ts...) -> overloaded; + +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}; +} + +// The poll interval is measured on the monotonic clock, but the last poll was +// recorded on the wall clock so that it could be persisted. An instant in the +// future means the wall clock moved backwards since it was written, which +// says nothing usable about when the last poll happened. +std::optional ToSteadyClock( + std::optional instant) { + if (!instant) { + return std::nullopt; + } + auto const now = std::chrono::system_clock::now(); + if (*instant > now) { + return std::nullopt; + } + return std::chrono::steady_clock::now() - (now - *instant); +} + +} // 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)); + + auto polling_config = [&](FDv2Config::PollingConfig const& polling) { + return MakeRequestConfig( + polling.base_url_override.value_or(params.polling_base_url), params, + serialized_context, config.use_post); + }; + auto streaming_config = [&](FDv2Config::StreamingConfig const& streaming) { + return MakeRequestConfig( + streaming.base_url_override.value_or(params.streaming_base_url), + params, serialized_context, config.use_post); + }; + + for (auto const& entry : definition->second.initializers) { + std::visit( + overloaded{ + [&](FDv2Config::CacheConfig const&) { + sources.initializers.push_back( + std::make_unique( + params.cache, params.context, params.logger)); + }, + [&](FDv2Config::PollingConfig const& polling) { + sources.initializers.push_back( + std::make_unique( + params.executor, params.logger, + polling_config(polling))); + }, + }, + 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. + auto const last_poll = + ToSteadyClock(params.cache->FreshnessFor(params.context)); + + for (auto const& entry : definition->second.synchronizers) { + std::visit( + overloaded{ + [&](FDv2Config::PollingConfig const& polling) { + sources.synchronizers.push_back( + std::make_unique( + params.executor, params.logger, + polling_config(polling), polling.poll_interval, + last_poll)); + }, + [&](FDv2Config::StreamingConfig const& streaming) { + sources.synchronizers.push_back( + std::make_unique( + params.executor, params.logger, + streaming_config(streaming), + polling_config(FDv2Config::PollingConfig{ + std::chrono::seconds::zero(), std::nullopt}), + streaming.initial_reconnect_delay)); + }, + }, + entry); + } + + return sources; +} + +} // namespace launchdarkly::client_side::data_sources diff --git a/libs/client-sdk/src/data_sources/fdv2/mode_sources.hpp b/libs/client-sdk/src/data_sources/fdv2/mode_sources.hpp new file mode 100644 index 000000000..b22837ec0 --- /dev/null +++ b/libs/client-sdk/src/data_sources/fdv2/mode_sources.hpp @@ -0,0 +1,60 @@ +#pragma once + +#include "ifdv2_initializer_factory.hpp" +#include "ifdv2_synchronizer_factory.hpp" + +#include "../../flag_manager/flag_persistence.hpp" + +#include +#include +#include +#include + +#include + +#include +#include +#include + +namespace launchdarkly::client_side::data_sources { + +using ConnectionMode = config::shared::ConnectionMode; +using FDv2Config = config::shared::built::FDv2Config; + +/** The factories one connection mode calls for. */ +struct ModeSources { + std::vector> initializers; + std::vector> 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; + /** Used by any source that does not configure a URL of its own. */ + std::string polling_base_url; + std::string streaming_base_url; + config::shared::built::HttpProperties http_properties; + Context context; + /** Whether the application asked for evaluation reasons. */ + bool with_reasons; + /** + * The local cache, read by the cache initializer and for the last time + * this context was polled. Non-owning. Must outlive the 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. + */ +ModeSources BuildModeSources(FDv2Config const& config, + ConnectionMode mode, + ModeSourceParams const& params); + +} // namespace launchdarkly::client_side::data_sources diff --git a/libs/client-sdk/src/data_sources/fdv2/source_factories.cpp b/libs/client-sdk/src/data_sources/fdv2/source_factories.cpp new file mode 100644 index 000000000..593350a1b --- /dev/null +++ b/libs/client-sdk/src/data_sources/fdv2/source_factories.cpp @@ -0,0 +1,59 @@ +#include "source_factories.hpp" + +#include "polling_initializer.hpp" +#include "polling_synchronizer.hpp" +#include "streaming_synchronizer.hpp" + +#include + +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 FDv2PollingInitializerFactory::Build() { + return std::make_unique(executor_, logger_, + request_config_); +} + +FDv2PollingSynchronizerFactory::FDv2PollingSynchronizerFactory( + boost::asio::any_io_executor executor, + Logger logger, + FDv2RequestConfig request_config, + std::chrono::seconds poll_interval, + std::optional 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 FDv2PollingSynchronizerFactory::Build() { + return std::make_unique( + 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 FDv2StreamingSynchronizerFactory::Build() { + return std::make_unique( + executor_, logger_, stream_config_, poll_config_, + initial_reconnect_delay_); +} + +} // namespace launchdarkly::client_side::data_sources diff --git a/libs/client-sdk/src/data_sources/fdv2/source_factories.hpp b/libs/client-sdk/src/data_sources/fdv2/source_factories.hpp new file mode 100644 index 000000000..f5cc48afe --- /dev/null +++ b/libs/client-sdk/src/data_sources/fdv2/source_factories.hpp @@ -0,0 +1,96 @@ +#pragma once + +#include "fdv2_request_config.hpp" +#include "ifdv2_initializer_factory.hpp" +#include "ifdv2_synchronizer_factory.hpp" + +#include + +#include + +#include +#include +#include + +namespace launchdarkly::client_side::data_sources { + +/** + * Builds fresh FDv2PollingInitializer instances on demand. + * + * Thread-safe: Build() may be called from any thread, and the + * configuration it hands to each source is fixed at construction. + */ +class FDv2PollingInitializerFactory final : public IFDv2InitializerFactory { + public: + FDv2PollingInitializerFactory(boost::asio::any_io_executor executor, + Logger logger, + FDv2RequestConfig request_config); + + std::unique_ptr Build() override; + + private: + boost::asio::any_io_executor const executor_; + Logger const logger_; + FDv2RequestConfig const request_config_; +}; + +/** + * Builds fresh FDv2PollingSynchronizer instances on demand. + * + * Thread-safe: Build() may be called from any thread, and the + * configuration it hands to each source is fixed at construction. + */ +class FDv2PollingSynchronizerFactory final : public IFDv2SynchronizerFactory { + public: + /** + * @param last_poll When this context was last polled, if that is known, + * so that the first poll after the synchronizer starts still respects the + * interval. + */ + FDv2PollingSynchronizerFactory( + boost::asio::any_io_executor executor, + Logger logger, + FDv2RequestConfig request_config, + std::chrono::seconds poll_interval, + std::optional last_poll); + + std::unique_ptr Build() override; + + private: + boost::asio::any_io_executor const executor_; + Logger const logger_; + FDv2RequestConfig const request_config_; + std::chrono::seconds const poll_interval_; + std::optional const last_poll_; +}; + +/** + * Builds fresh FDv2StreamingSynchronizer instances on demand. + * + * Thread-safe: Build() may be called from any thread, and the + * configuration it hands to each source is fixed at construction. + */ +class FDv2StreamingSynchronizerFactory final : public IFDv2SynchronizerFactory { + public: + /** + * @param poll_config Where to poll in answer to a `ping` event on the + * stream. Must describe the same context as stream_config. + */ + FDv2StreamingSynchronizerFactory( + boost::asio::any_io_executor executor, + Logger logger, + FDv2RequestConfig stream_config, + FDv2RequestConfig poll_config, + std::chrono::milliseconds initial_reconnect_delay); + + std::unique_ptr Build() override; + + private: + boost::asio::any_io_executor const executor_; + Logger const logger_; + FDv2RequestConfig const stream_config_; + FDv2RequestConfig const poll_config_; + std::chrono::milliseconds const initial_reconnect_delay_; +}; + +} // namespace launchdarkly::client_side::data_sources diff --git a/libs/client-sdk/tests/fdv2_mode_sources_test.cpp b/libs/client-sdk/tests/fdv2_mode_sources_test.cpp new file mode 100644 index 000000000..8a0780052 --- /dev/null +++ b/libs/client-sdk/tests/fdv2_mode_sources_test.cpp @@ -0,0 +1,160 @@ +#include + +#include +#include + +#include +#include +#include + +#include + +using launchdarkly::ContextBuilder; +using launchdarkly::client_side::flag_manager::FlagManager; +using launchdarkly::config::shared::GetConnectionModeName; + +using namespace launchdarkly::client_side::data_sources; + +namespace { + +class ModeSourcesFixture : public ::testing::Test { + public: + ModeSourcesFixture() + : logger_(launchdarkly::logging::NullLogger()), + flag_manager_("sdk-key", logger_, 5, nullptr) {} + + static FDv2Config Defaults() { + return launchdarkly::config::shared::Defaults< + launchdarkly::config::shared::ClientSDK>::FDv2Config(); + } + + ModeSourceParams Params() { + return ModeSourceParams{ + ioc_.get_executor(), + logger_, + "https://polling.example.com", + "https://streaming.example.com", + launchdarkly::config::shared::Defaults< + launchdarkly::config::shared::ClientSDK>::HttpProperties(), + ContextBuilder().Kind("user", "user-key").Build(), + /* with_reasons= */ false, + &flag_manager_.Cache()}; + } + + private: + boost::asio::io_context ioc_; + launchdarkly::Logger logger_; + FlagManager flag_manager_; +}; + +} // namespace + +TEST_F(ModeSourcesFixture, StreamingModeInitializesFromCacheThenPolls) { + auto sources = + BuildModeSources(Defaults(), ConnectionMode::kStreaming, Params()); + + ASSERT_EQ(2u, sources.initializers.size()); + EXPECT_TRUE(sources.initializers[0]->IsFromCache()); + EXPECT_FALSE(sources.initializers[1]->IsFromCache()); +} + +// Streaming is the primary tier and polling the fallback, so the SDK keeps +// receiving updates when a stream cannot be maintained. +TEST_F(ModeSourcesFixture, StreamingModeFallsBackToPolling) { + auto sources = + BuildModeSources(Defaults(), ConnectionMode::kStreaming, Params()); + + ASSERT_EQ(2u, sources.synchronizers.size()); + EXPECT_EQ("FDv2 streaming synchronizer", + sources.synchronizers[0]->Build()->Identity()); + EXPECT_EQ("FDv2 polling synchronizer", + sources.synchronizers[1]->Build()->Identity()); +} + +TEST_F(ModeSourcesFixture, PollingModeInitializesFromCacheOnly) { + auto sources = + BuildModeSources(Defaults(), ConnectionMode::kPolling, Params()); + + ASSERT_EQ(1u, sources.initializers.size()); + EXPECT_TRUE(sources.initializers[0]->IsFromCache()); + ASSERT_EQ(1u, sources.synchronizers.size()); + EXPECT_EQ("FDv2 polling synchronizer", + sources.synchronizers[0]->Build()->Identity()); +} + +// Offline still loads persisted flags, but makes no requests. +TEST_F(ModeSourcesFixture, OfflineModeHasOnlyTheCacheInitializer) { + auto sources = + BuildModeSources(Defaults(), ConnectionMode::kOffline, Params()); + + ASSERT_EQ(1u, sources.initializers.size()); + EXPECT_TRUE(sources.initializers[0]->IsFromCache()); + EXPECT_TRUE(sources.synchronizers.empty()); +} + +TEST_F(ModeSourcesFixture, UnconfiguredModeProducesNoSources) { + auto config = launchdarkly::config::shared::Defaults< + launchdarkly::config::shared::ClientSDK>::FDv2Config(); + config.modes.erase(ConnectionMode::kPolling); + + auto sources = BuildModeSources(config, ConnectionMode::kPolling, Params()); + + EXPECT_TRUE(sources.initializers.empty()); + EXPECT_TRUE(sources.synchronizers.empty()); +} + +TEST_F(ModeSourcesFixture, AnOverriddenModeReplacesTheBuiltInPipeline) { + auto config = launchdarkly::config::shared::Defaults< + launchdarkly::config::shared::ClientSDK>::FDv2Config(); + config.modes[ConnectionMode::kStreaming] = FDv2Config::ModeDefinition{ + {FDv2Config::CacheConfig{}}, + {FDv2Config::StreamingConfig{std::chrono::seconds{1}, std::nullopt}}, + std::nullopt}; + + auto sources = + BuildModeSources(config, ConnectionMode::kStreaming, Params()); + + EXPECT_EQ(1u, sources.initializers.size()); + ASSERT_EQ(1u, sources.synchronizers.size()); + EXPECT_EQ("FDv2 streaming synchronizer", + sources.synchronizers[0]->Build()->Identity()); +} + +TEST(ConnectionModeTest, ModesAreNamedAsTheConfigurationSpellsThem) { + EXPECT_STREQ("streaming", + GetConnectionModeName(ConnectionMode::kStreaming)); + EXPECT_STREQ("polling", GetConnectionModeName(ConnectionMode::kPolling)); + EXPECT_STREQ("offline", GetConnectionModeName(ConnectionMode::kOffline)); +} + +// Desktop starts in streaming, and provides no background mode. +TEST(FDv2ConfigTest, DefaultsProvideTheThreeDesktopModes) { + auto const config = launchdarkly::config::shared::Defaults< + launchdarkly::config::shared::ClientSDK>::FDv2Config(); + + EXPECT_EQ(ConnectionMode::kStreaming, config.initial_mode); + EXPECT_EQ(3u, config.modes.size()); + EXPECT_EQ(1u, config.modes.count(ConnectionMode::kStreaming)); + EXPECT_EQ(1u, config.modes.count(ConnectionMode::kPolling)); + EXPECT_EQ(1u, config.modes.count(ConnectionMode::kOffline)); +} + +TEST(FDv2ConfigTest, DefaultTimeoutsAreTwoAndFiveMinutes) { + auto const config = launchdarkly::config::shared::Defaults< + launchdarkly::config::shared::ClientSDK>::FDv2Config(); + + EXPECT_EQ(std::chrono::seconds{120}, config.fallback_timeout); + EXPECT_EQ(std::chrono::seconds{300}, config.recovery_timeout); +} + +TEST(FDv2ConfigTest, ModesThatMakeRequestsConfigureAnFDv1Fallback) { + auto const config = launchdarkly::config::shared::Defaults< + launchdarkly::config::shared::ClientSDK>::FDv2Config(); + + EXPECT_TRUE( + config.modes.at(ConnectionMode::kStreaming).fdv1_fallback.has_value()); + EXPECT_TRUE( + config.modes.at(ConnectionMode::kPolling).fdv1_fallback.has_value()); + EXPECT_FALSE( + config.modes.at(ConnectionMode::kOffline).fdv1_fallback.has_value()); +} diff --git a/libs/common/include/launchdarkly/config/shared/built/fdv2_config.hpp b/libs/common/include/launchdarkly/config/shared/built/fdv2_config.hpp new file mode 100644 index 000000000..eb55c52f2 --- /dev/null +++ b/libs/common/include/launchdarkly/config/shared/built/fdv2_config.hpp @@ -0,0 +1,97 @@ +#pragma once + +#include +#include + +#include +#include +#include +#include +#include +#include + +namespace launchdarkly::config::shared::built { + +template +struct FDv2Config; + +/** + * The FDv2 data system configuration: which connection mode the SDK starts + * in, what each mode's sources are, and the settings shared across them. + */ +template <> +struct FDv2Config { + /** Reads flag data the SDK persisted on a previous run. */ + struct CacheConfig {}; + + struct StreamingConfig { + std::chrono::milliseconds initial_reconnect_delay; + /** Overrides the streaming base URL for this source alone. */ + std::optional base_url_override; + }; + + struct PollingConfig { + std::chrono::seconds poll_interval; + /** Overrides the polling base URL for this source alone. */ + std::optional base_url_override; + }; + + /** + * The FDv1 polling source used while the service has directed the SDK + * away from FDv2. + */ + struct FDv1FallbackConfig { + std::chrono::seconds poll_interval; + /** Overrides the polling base URL for the fallback alone. */ + std::optional base_url_override; + }; + + using InitializerEntry = std::variant; + using SynchronizerEntry = std::variant; + + /** + * What one connection mode does. The synchronizer list is ordered by + * preference. The first entry is the primary, and later entries are the + * tiers the SDK falls back to. + */ + struct ModeDefinition { + std::vector initializers; + std::vector synchronizers; + std::optional fdv1_fallback; + }; + + /** The mode the SDK starts in. */ + ConnectionMode initial_mode; + + /** + * Where a source sends its requests when it does not override the URL + * itself. FDv2's endpoints are not the ones FDv1 uses, so these hold + * FDv2's own defaults. When the application configures its own endpoints, + * these are resolved to those instead. + */ + std::string polling_base_url; + std::string streaming_base_url; + + /** What each mode does. Modes absent from the map are unavailable. */ + std::map modes; + + /** + * Whether to send the evaluation context in a request body rather than + * base64url-encoded into the request path. + */ + bool use_post; + + /** + * How long the active synchronizer may remain interrupted before the SDK + * falls back to the next tier. + */ + std::chrono::milliseconds fallback_timeout; + + /** + * How long a fallback tier must run before the SDK attempts to return to + * the preferred one. + */ + std::chrono::milliseconds recovery_timeout; +}; + +} // namespace launchdarkly::config::shared::built diff --git a/libs/common/include/launchdarkly/config/shared/connection_mode.hpp b/libs/common/include/launchdarkly/config/shared/connection_mode.hpp new file mode 100644 index 000000000..6aa5feda2 --- /dev/null +++ b/libs/common/include/launchdarkly/config/shared/connection_mode.hpp @@ -0,0 +1,23 @@ +#pragma once + +namespace launchdarkly::config::shared { + +/** + * A named data system configuration: which sources the SDK uses to load flag + * data, and which it uses to keep that data current. + */ +enum class ConnectionMode { + /** Stream updates, falling back to polling. */ + kStreaming, + /** Poll for updates on an interval. */ + kPolling, + /** Evaluate against whatever is cached, and make no requests. */ + kOffline, +}; + +/** + * The mode's name as the configuration API spells it. + */ +char const* GetConnectionModeName(ConnectionMode mode); + +} // namespace launchdarkly::config::shared diff --git a/libs/common/include/launchdarkly/config/shared/defaults.hpp b/libs/common/include/launchdarkly/config/shared/defaults.hpp index 0fc0cb9d9..ab3c6eda8 100644 --- a/libs/common/include/launchdarkly/config/shared/defaults.hpp +++ b/libs/common/include/launchdarkly/config/shared/defaults.hpp @@ -2,6 +2,7 @@ #include #include +#include #include #include #include @@ -75,6 +76,52 @@ struct Defaults { "/msdk/evalx/context", std::chrono::minutes(5)}; } + /** + * The three connection modes a desktop SDK provides, starting in + * streaming, with automatic mode switching off. + */ + static auto FDv2Config() -> shared::built::FDv2Config { + using Config = shared::built::FDv2Config; + + // Both timeouts are chosen for consistency with the other + // LaunchDarkly SDKs. + auto const fallback_timeout = std::chrono::seconds(120); + auto const recovery_timeout = std::chrono::seconds(300); + + Config::StreamingConfig const streaming{std::chrono::seconds(1), + std::nullopt}; + Config::PollingConfig const polling{std::chrono::minutes(5), + std::nullopt}; + Config::FDv1FallbackConfig const fdv1_fallback{std::chrono::minutes(5), + std::nullopt}; + + return { + shared::ConnectionMode::kStreaming, + "https://sdk.launchdarkly.com", + "https://clientstream.launchdarkly.com", + { + // Streaming initializes from the cache, then polls for a + // basis so that the stream can deliver only what changed, + // and falls back to polling if the stream cannot be kept up. + {shared::ConnectionMode::kStreaming, + Config::ModeDefinition{{Config::CacheConfig{}, polling}, + {streaming, polling}, + fdv1_fallback}}, + {shared::ConnectionMode::kPolling, + Config::ModeDefinition{ + {Config::CacheConfig{}}, {polling}, fdv1_fallback}}, + // Offline evaluates against the cache and makes no requests, + // so it has nothing to fall back to. + {shared::ConnectionMode::kOffline, + Config::ModeDefinition{ + {Config::CacheConfig{}}, {}, std::nullopt}}, + }, + /* use_post= */ false, + fallback_timeout, + recovery_timeout, + }; + } + static std::size_t MaxCachedContexts() { return 5; } }; diff --git a/libs/common/src/CMakeLists.txt b/libs/common/src/CMakeLists.txt index b0a9446b2..c39644fd0 100644 --- a/libs/common/src/CMakeLists.txt +++ b/libs/common/src/CMakeLists.txt @@ -28,6 +28,7 @@ add_library(${LIBNAME} OBJECT attributes_builder.cpp error.cpp config/service_endpoints.cpp + config/connection_mode.cpp config/events.cpp config/endpoints_builder.cpp config/events_builder.cpp diff --git a/libs/common/src/config/connection_mode.cpp b/libs/common/src/config/connection_mode.cpp new file mode 100644 index 000000000..2bd40630b --- /dev/null +++ b/libs/common/src/config/connection_mode.cpp @@ -0,0 +1,17 @@ +#include + +namespace launchdarkly::config::shared { + +char const* GetConnectionModeName(ConnectionMode mode) { + switch (mode) { + case ConnectionMode::kStreaming: + return "streaming"; + case ConnectionMode::kPolling: + return "polling"; + case ConnectionMode::kOffline: + return "offline"; + } + return "unknown"; +} + +} // namespace launchdarkly::config::shared From 06f4074c196e46197db108b7916bccd49abe81dc Mon Sep 17 00:00:00 2001 From: Bee Klimt Date: Mon, 14 Sep 2026 15:19:45 -0700 Subject: [PATCH 4/6] refactor: Call ReadFreshness from the FDv2 connection-mode sources --- libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp b/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp index 2d3bfa98c..f159096c2 100644 --- a/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp +++ b/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp @@ -101,7 +101,7 @@ ModeSources BuildModeSources(FDv2Config const& config, // time this context was polled, which survives restarts so that repeated // launches cannot produce a burst of requests. auto const last_poll = - ToSteadyClock(params.cache->FreshnessFor(params.context)); + ToSteadyClock(params.cache->ReadFreshness(params.context)); for (auto const& entry : definition->second.synchronizers) { std::visit( From 6f06cfb7410ed69b85fd6e5d94bbf58e82d82e31 Mon Sep 17 00:00:00 2001 From: Bee Klimt Date: Fri, 25 Sep 2026 18:05:11 -0700 Subject: [PATCH 5/6] refactor: Simplify the FDv2 connection-mode source assembly --- .../src/data_sources/fdv2/mode_sources.cpp | 59 +++++++++---------- .../config/shared/built/fdv2_config.hpp | 8 +-- 2 files changed, 29 insertions(+), 38 deletions(-) diff --git a/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp b/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp index f159096c2..1647db20e 100644 --- a/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp +++ b/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp @@ -33,20 +33,14 @@ FDv2RequestConfig MakeRequestConfig(std::string base_url, params.with_reasons}; } -// The poll interval is measured on the monotonic clock, but the last poll was -// recorded on the wall clock so that it could be persisted. An instant in the -// future means the wall clock moved backwards since it was written, which -// says nothing usable about when the last poll happened. -std::optional ToSteadyClock( - std::optional instant) { - if (!instant) { - return std::nullopt; - } +// 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(); - if (*instant > now) { - return std::nullopt; - } - return std::chrono::steady_clock::now() - (now - *instant); + return std::chrono::steady_clock::now() - (now - instant); } } // namespace @@ -68,17 +62,6 @@ ModeSources BuildModeSources(FDv2Config const& config, auto const serialized_context = boost::json::serialize(boost::json::value_from(params.context)); - auto polling_config = [&](FDv2Config::PollingConfig const& polling) { - return MakeRequestConfig( - polling.base_url_override.value_or(params.polling_base_url), params, - serialized_context, config.use_post); - }; - auto streaming_config = [&](FDv2Config::StreamingConfig const& streaming) { - return MakeRequestConfig( - streaming.base_url_override.value_or(params.streaming_base_url), - params, serialized_context, config.use_post); - }; - for (auto const& entry : definition->second.initializers) { std::visit( overloaded{ @@ -91,7 +74,10 @@ ModeSources BuildModeSources(FDv2Config const& config, sources.initializers.push_back( std::make_unique( params.executor, params.logger, - polling_config(polling))); + MakeRequestConfig( + polling.base_url_override.value_or( + params.polling_base_url), + params, serialized_context, config.use_post))); }, }, entry); @@ -100,8 +86,10 @@ ModeSources BuildModeSources(FDv2Config const& config, // 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. - auto const last_poll = - ToSteadyClock(params.cache->ReadFreshness(params.context)); + std::optional last_poll; + if (auto const freshness = params.cache->ReadFreshness(params.context)) { + last_poll = ToSteadyClock(*freshness); + } for (auto const& entry : definition->second.synchronizers) { std::visit( @@ -110,16 +98,23 @@ ModeSources BuildModeSources(FDv2Config const& config, sources.synchronizers.push_back( std::make_unique( params.executor, params.logger, - polling_config(polling), polling.poll_interval, - last_poll)); + 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( params.executor, params.logger, - streaming_config(streaming), - polling_config(FDv2Config::PollingConfig{ - std::chrono::seconds::zero(), std::nullopt}), + 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)); }, }, diff --git a/libs/common/include/launchdarkly/config/shared/built/fdv2_config.hpp b/libs/common/include/launchdarkly/config/shared/built/fdv2_config.hpp index eb55c52f2..b0df26091 100644 --- a/libs/common/include/launchdarkly/config/shared/built/fdv2_config.hpp +++ b/libs/common/include/launchdarkly/config/shared/built/fdv2_config.hpp @@ -63,13 +63,9 @@ struct FDv2Config { /** The mode the SDK starts in. */ ConnectionMode initial_mode; - /** - * Where a source sends its requests when it does not override the URL - * itself. FDv2's endpoints are not the ones FDv1 uses, so these hold - * FDv2's own defaults. When the application configures its own endpoints, - * these are resolved to those instead. - */ + /** Base URL for polling sources that do not override it. */ std::string polling_base_url; + /** Base URL for streaming sources that do not override it. */ std::string streaming_base_url; /** What each mode does. Modes absent from the map are unavailable. */ From 58353ce1f161a852b82bd083d51c8bf3e93b68fb Mon Sep 17 00:00:00 2001 From: Bee Klimt Date: Sat, 26 Sep 2026 00:31:50 -0700 Subject: [PATCH 6/6] feat: Wire the FDv1 fallback tier into the connection-mode assembly --- .../src/data_sources/fdv2/mode_sources.cpp | 36 +++++++++++++++++ .../src/data_sources/fdv2/mode_sources.hpp | 11 +++++ .../data_sources/fdv2/source_factories.cpp | 27 +++++++++++++ .../data_sources/fdv2/source_factories.hpp | 40 +++++++++++++++++++ .../tests/fdv2_mode_sources_test.cpp | 25 +++++++++++- 5 files changed, 137 insertions(+), 2 deletions(-) diff --git a/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp b/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp index 1647db20e..acc2d45bc 100644 --- a/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp +++ b/libs/client-sdk/src/data_sources/fdv2/mode_sources.cpp @@ -3,6 +3,7 @@ #include "cache_initializer.hpp" #include "source_factories.hpp" +#include #include #include @@ -43,6 +44,37 @@ std::chrono::steady_clock::time_point ToSteadyClock( 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 MakeFDv1Fallback( + FDv2Config::FDv1FallbackConfig const& fallback, + ModeSourceParams const& params) { + auto const defaults = + config::shared::Defaults::PollingConfig(); + + config::shared::built::DataSourceConfig const + fdv1_config{ + config::shared::built::PollingConfig{ + 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( + params.executor, params.logger, std::move(endpoints), fdv1_config, + params.http_properties, params.context); +} + } // namespace ModeSources BuildModeSources(FDv2Config const& config, @@ -121,6 +153,10 @@ ModeSources BuildModeSources(FDv2Config const& config, entry); } + if (auto const& fallback = definition->second.fdv1_fallback) { + sources.synchronizers.push_back(MakeFDv1Fallback(*fallback, params)); + } + return sources; } diff --git a/libs/client-sdk/src/data_sources/fdv2/mode_sources.hpp b/libs/client-sdk/src/data_sources/fdv2/mode_sources.hpp index b22837ec0..a66837ccc 100644 --- a/libs/client-sdk/src/data_sources/fdv2/mode_sources.hpp +++ b/libs/client-sdk/src/data_sources/fdv2/mode_sources.hpp @@ -7,6 +7,7 @@ #include #include +#include #include #include @@ -38,6 +39,11 @@ struct ModeSourceParams { 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; @@ -52,6 +58,11 @@ struct ModeSourceParams { /** * 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 list. The orchestrator holds that tier in reserve rather than using it + * in rotation. It is started only while the service has directed the SDK away + * from FDv2. */ ModeSources BuildModeSources(FDv2Config const& config, ConnectionMode mode, diff --git a/libs/client-sdk/src/data_sources/fdv2/source_factories.cpp b/libs/client-sdk/src/data_sources/fdv2/source_factories.cpp index 593350a1b..18690808a 100644 --- a/libs/client-sdk/src/data_sources/fdv2/source_factories.cpp +++ b/libs/client-sdk/src/data_sources/fdv2/source_factories.cpp @@ -1,5 +1,7 @@ #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" @@ -56,4 +58,29 @@ std::unique_ptr FDv2StreamingSynchronizerFactory::Build() { initial_reconnect_delay_); } +FDv1PollingAdapterFactory::FDv1PollingAdapterFactory( + boost::asio::any_io_executor executor, + Logger logger, + config::shared::built::ServiceEndpoints endpoints, + config::shared::built::DataSourceConfig + 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 FDv1PollingAdapterFactory::Build() { + return std::make_unique( + [this](IDataSourceUpdateSink* sink, + DataSourceStatusManager* status_manager) { + return std::make_shared( + endpoints_, data_source_config_, http_properties_, executor_, + context_, *sink, *status_manager, logger_); + }); +} + } // namespace launchdarkly::client_side::data_sources diff --git a/libs/client-sdk/src/data_sources/fdv2/source_factories.hpp b/libs/client-sdk/src/data_sources/fdv2/source_factories.hpp index f5cc48afe..59a6b0690 100644 --- a/libs/client-sdk/src/data_sources/fdv2/source_factories.hpp +++ b/libs/client-sdk/src/data_sources/fdv2/source_factories.hpp @@ -4,6 +4,10 @@ #include "ifdv2_initializer_factory.hpp" #include "ifdv2_synchronizer_factory.hpp" +#include +#include +#include +#include #include #include @@ -93,4 +97,40 @@ class FDv2StreamingSynchronizerFactory final : public IFDv2SynchronizerFactory { std::chrono::milliseconds const initial_reconnect_delay_; }; +/** + * Builds fresh FDv1AdapterSynchronizer instances wrapping a freshly-built + * FDv1 polling source. + * + * The synchronizers it builds report themselves as the FDv1 tier, which the + * orchestrator keeps in reserve until the service directs the SDK away from + * FDv2. + * + * Thread-safe: Build() may be called from any thread, and the + * configuration it hands to each source is fixed at construction. + */ +class FDv1PollingAdapterFactory final : public IFDv2SynchronizerFactory { + public: + FDv1PollingAdapterFactory( + boost::asio::any_io_executor executor, + Logger logger, + config::shared::built::ServiceEndpoints endpoints, + config::shared::built::DataSourceConfig + data_source_config, + config::shared::built::HttpProperties http_properties, + Context context); + + std::unique_ptr Build() override; + + [[nodiscard]] bool IsFDv1Fallback() const override { return true; } + + private: + boost::asio::any_io_executor const executor_; + Logger const logger_; + config::shared::built::ServiceEndpoints const endpoints_; + config::shared::built::DataSourceConfig const + data_source_config_; + config::shared::built::HttpProperties const http_properties_; + Context const context_; +}; + } // namespace launchdarkly::client_side::data_sources diff --git a/libs/client-sdk/tests/fdv2_mode_sources_test.cpp b/libs/client-sdk/tests/fdv2_mode_sources_test.cpp index 8a0780052..250565ff8 100644 --- a/libs/client-sdk/tests/fdv2_mode_sources_test.cpp +++ b/libs/client-sdk/tests/fdv2_mode_sources_test.cpp @@ -36,6 +36,8 @@ class ModeSourcesFixture : public ::testing::Test { "https://streaming.example.com", launchdarkly::config::shared::Defaults< launchdarkly::config::shared::ClientSDK>::HttpProperties(), + launchdarkly::config::shared::Defaults< + launchdarkly::config::shared::ClientSDK>::ServiceEndpoints(), ContextBuilder().Kind("user", "user-key").Build(), /* with_reasons= */ false, &flag_manager_.Cache()}; @@ -64,7 +66,7 @@ TEST_F(ModeSourcesFixture, StreamingModeFallsBackToPolling) { auto sources = BuildModeSources(Defaults(), ConnectionMode::kStreaming, Params()); - ASSERT_EQ(2u, sources.synchronizers.size()); + ASSERT_EQ(3u, sources.synchronizers.size()); EXPECT_EQ("FDv2 streaming synchronizer", sources.synchronizers[0]->Build()->Identity()); EXPECT_EQ("FDv2 polling synchronizer", @@ -77,7 +79,7 @@ TEST_F(ModeSourcesFixture, PollingModeInitializesFromCacheOnly) { ASSERT_EQ(1u, sources.initializers.size()); EXPECT_TRUE(sources.initializers[0]->IsFromCache()); - ASSERT_EQ(1u, sources.synchronizers.size()); + ASSERT_EQ(2u, sources.synchronizers.size()); EXPECT_EQ("FDv2 polling synchronizer", sources.synchronizers[0]->Build()->Identity()); } @@ -158,3 +160,22 @@ TEST(FDv2ConfigTest, ModesThatMakeRequestsConfigureAnFDv1Fallback) { EXPECT_FALSE( config.modes.at(ConnectionMode::kOffline).fdv1_fallback.has_value()); } + +// The FDv1 tier is appended last, and the orchestrator keeps it blocked until +// the service directs the SDK away from FDv2. +TEST_F(ModeSourcesFixture, ModesWithAFallbackAppendTheFDv1Tier) { + auto sources = + BuildModeSources(Defaults(), ConnectionMode::kStreaming, Params()); + + ASSERT_FALSE(sources.synchronizers.empty()); + auto const& last = sources.synchronizers.back(); + EXPECT_TRUE(last->IsFDv1Fallback()); + EXPECT_EQ("FDv1 fallback adapter", last->Build()->Identity()); +} + +TEST_F(ModeSourcesFixture, ModesWithNoFallbackAppendNothing) { + auto sources = + BuildModeSources(Defaults(), ConnectionMode::kOffline, Params()); + + EXPECT_TRUE(sources.synchronizers.empty()); +}