From 15578a89572d2c54078b853af5d131c8b080585a Mon Sep 17 00:00:00 2001 From: Bee Klimt Date: Wed, 2 Sep 2026 03:40:10 +0000 Subject: [PATCH 1/2] 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 4faf302f4..b26e3cb58 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. @@ -82,6 +87,7 @@ FDv2DataSource::~FDv2DataSource() { void FDv2DataSource::Close() { std::lock_guard lock(mutex_); closed_ = true; + fdv2_retry_cancel_.Cancel(); if (active_initializer_) { active_initializer_->Close(); } @@ -293,14 +299,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 { @@ -482,6 +505,7 @@ void FDv2DataSource::OnSynchronizerResult(FDv2SourceResult result) { }, result.value); + bool disconnected = false; { std::lock_guard lock(mutex_); if (closed_ || got_shutdown) { @@ -489,13 +513,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 { @@ -503,6 +544,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 ed6bc75ee..12ce8f7d4 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_; @@ -231,6 +241,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 ece389e3e..a44dc1236 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(); } @@ -888,3 +902,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 f0596a82e34e0a5882577d163407cca26a301c52 Mon Sep 17 00:00:00 2001 From: Bee Klimt Date: Mon, 5 Oct 2026 22:03:33 -0700 Subject: [PATCH 2/2] chore: Tighten the FDv1 fallback comments --- .../fdv2/fdv1_adapter_synchronizer.cpp | 16 +++++------ .../fdv2/fdv1_adapter_synchronizer.hpp | 25 +++++++---------- .../data_sources/fdv2/fdv2_data_source.cpp | 5 ++-- .../tests/fdv1_adapter_synchronizer_test.cpp | 27 ++++++++++++++++--- .../tests/fdv2_data_source_test.cpp | 12 ++++++--- 5 files changed, 52 insertions(+), 33 deletions(-) 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 index 9468a807c..7feaee7c3 100644 --- a/libs/client-sdk/src/data_sources/fdv2/fdv1_adapter_synchronizer.cpp +++ b/libs/client-sdk/src/data_sources/fdv2/fdv1_adapter_synchronizer.cpp @@ -108,9 +108,8 @@ void FDv1AdapterSynchronizer::ConvertingSink::Apply( 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. +// Turns the wrapped source's status into a result, or into nothing when the +// status carries no error. std::optional ResultForStatus( DataSourceStatus const& status) { auto const error = status.LastError(); @@ -119,9 +118,8 @@ std::optional ResultForStatus( } 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. + // An error 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: @@ -172,7 +170,7 @@ async::Future FDv1AdapterSynchronizer::Next( return async::WhenAny(closed, result_future) .Then( [state = state_, result_future](std::size_t const& idx) mutable - -> async::Future { + -> async::Future { if (idx == 0) { state->ResolvePendingAsShutdown(); return async::MakeFuture( @@ -191,8 +189,8 @@ void FDv1AdapterSynchronizer::Close() { 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. + // The completion does no work. It exists to hold these alive until + // shutdown finishes. fdv1_source_->ShutdownAsync( [sink = sink_, status = status_manager_, source = fdv1_source_] {}); } 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 index 2d5154301..92356e075 100644 --- a/libs/client-sdk/src/data_sources/fdv2/fdv1_adapter_synchronizer.hpp +++ b/libs/client-sdk/src/data_sources/fdv2/fdv1_adapter_synchronizer.hpp @@ -19,13 +19,9 @@ 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. + * Presents an FDv1 data source as an FDv2 synchronizer. * - * 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. + * FDv1 has no selectors, so the changesets this reports carry none. * * Thread safety: Next() and Close() may be called from any thread. Only one * Next() may be outstanding at a time. @@ -33,9 +29,8 @@ namespace launchdarkly::client_side::data_sources { 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. + * Builds the wrapped FDv1 source. Called once during construction, with a + * sink and status manager that stay valid for the source's lifetime. */ using SourceBuilder = std::function(IDataSourceUpdateSink*, @@ -63,17 +58,17 @@ class FDv1AdapterSynchronizer final : public IFDv2Synchronizer { async::Future GetNext(); - /** - * Resolves any pending Next() with Shutdown and clears it, so that a - * caller abandoned by Close() is not left waiting. - */ + /** Resolves any pending Next() with Shutdown and clears it. */ void ResolvePendingAsShutdown(); + /** + * Hands the result to a pending Next(), or queues it for the next + * one. Results that arrive after Close() are dropped. + */ void Notify(FDv2SourceResult result); private: - // Finished once the owning adapter's Close() has run. Read in Notify - // to drop late results. + // Finished once the owning adapter's Close() has run. async::Future const closed_; mutable std::mutex mutex_; 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 b26e3cb58..adea662b5 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 @@ -548,9 +548,8 @@ 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. + // Cancel any attempt already scheduled. 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) diff --git a/libs/client-sdk/tests/fdv1_adapter_synchronizer_test.cpp b/libs/client-sdk/tests/fdv1_adapter_synchronizer_test.cpp index 1070d2ede..8b0dee274 100644 --- a/libs/client-sdk/tests/fdv1_adapter_synchronizer_test.cpp +++ b/libs/client-sdk/tests/fdv1_adapter_synchronizer_test.cpp @@ -91,22 +91,27 @@ class AdapterFixture { TEST(FDv1AdapterSynchronizerTest, TheFirstNextStartsTheWrappedSource) { AdapterFixture f; + // Building the adapter leaves the source alone. EXPECT_EQ(0, f.Source().start_count_); + // Ask for the first result. f.Source().Sink()->Init(ContextBuilder().Kind("user", "user-key").Build(), {{"flagA", Flag(1, Value("a"))}}); f.Next(); + // The source is now running. EXPECT_EQ(1, f.Source().start_count_); } TEST(FDv1AdapterSynchronizerTest, InitBecomesAFullChangeSet) { AdapterFixture f; + // Report a data set through the FDv1 sink. f.Source().Sink()->Init(ContextBuilder().Kind("user", "user-key").Build(), {{"flagA", Flag(1, Value("a"))}}); auto result = f.Next(); + // It arrives as a full changeset carrying the flag. ASSERT_TRUE(result.has_value()); auto* change_set = std::get_if(&result->value); ASSERT_NE(nullptr, change_set); @@ -119,10 +124,12 @@ TEST(FDv1AdapterSynchronizerTest, InitBecomesAFullChangeSet) { TEST(FDv1AdapterSynchronizerTest, UpsertBecomesAPartialChangeSet) { AdapterFixture f; + // Report a single flag through the FDv1 sink. f.Source().Sink()->Upsert(ContextBuilder().Kind("user", "user-key").Build(), "flagA", Flag(2, Value("a2"))); auto result = f.Next(); + // It arrives as a partial changeset carrying that flag. ASSERT_TRUE(result.has_value()); auto* change_set = std::get_if(&result->value); ASSERT_NE(nullptr, change_set); @@ -132,15 +139,15 @@ TEST(FDv1AdapterSynchronizerTest, UpsertBecomesAPartialChangeSet) { 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; + // Report a data set through the FDv1 sink. f.Source().Sink()->Init(ContextBuilder().Kind("user", "user-key").Build(), {{"flagA", Flag(1, Value("a"))}}); auto result = f.Next(); + // FDv1 supplies no selector, so the changeset carries none. auto* change_set = std::get_if(&result->value); ASSERT_NE(nullptr, change_set); EXPECT_FALSE(change_set->change_set.selector.value.has_value()); @@ -149,11 +156,13 @@ TEST(FDv1AdapterSynchronizerTest, ChangeSetsCarryNoSelector) { TEST(FDv1AdapterSynchronizerTest, ARecoverableErrorBecomesInterrupted) { AdapterFixture f; + // Report a recoverable failure through the status manager. f.Source().StatusManager()->SetState( DataSourceStatus::DataSourceState::kInterrupted, DataSourceStatus::ErrorInfo::ErrorKind::kNetworkError, "boom"); auto result = f.Next(); + // It arrives as an interruption. ASSERT_TRUE(result.has_value()); EXPECT_TRUE( std::holds_alternative(result->value)); @@ -162,11 +171,13 @@ TEST(FDv1AdapterSynchronizerTest, ARecoverableErrorBecomesInterrupted) { TEST(FDv1AdapterSynchronizerTest, AnUnrecoverableErrorBecomesTerminal) { AdapterFixture f; + // Report an unrecoverable failure through the status manager. f.Source().StatusManager()->SetState( DataSourceStatus::DataSourceState::kShutdown, DataSourceStatus::ErrorInfo::ErrorKind::kErrorResponse, "unauthorized"); auto result = f.Next(); + // It arrives as a terminal error. ASSERT_TRUE(result.has_value()); EXPECT_TRUE( std::holds_alternative(result->value)); @@ -176,9 +187,11 @@ TEST(FDv1AdapterSynchronizerTest, AnUnrecoverableErrorBecomesTerminal) { TEST(FDv1AdapterSynchronizerTest, BecomingValidReportsNothing) { AdapterFixture f; + // Report the source becoming valid. f.Source().StatusManager()->SetState( DataSourceStatus::DataSourceState::kValid); + // A status carrying no error is not a result, so Next() stays pending. auto future = f.Adapter().Next(launchdarkly::data_model::Selector{}); EXPECT_FALSE(future.IsFinished()); } @@ -186,20 +199,23 @@ TEST(FDv1AdapterSynchronizerTest, BecomingValidReportsNothing) { TEST(FDv1AdapterSynchronizerTest, CloseShutsDownTheWrappedSource) { AdapterFixture f; + // Start the source by asking for a result, then close the adapter. f.Source().Sink()->Init(ContextBuilder().Kind("user", "user-key").Build(), {}); f.Next(); f.Adapter().Close(); + // The wrapped source is shut down too. EXPECT_EQ(1, f.Source().shutdown_count_); } -// Nothing was started, so there is nothing to shut down. TEST(FDv1AdapterSynchronizerTest, CloseBeforeAnyNextDoesNotShutDown) { AdapterFixture f; + // Close without ever asking for a result. f.Adapter().Close(); + // The source was never started, so it is left alone. EXPECT_EQ(0, f.Source().shutdown_count_); EXPECT_EQ(0, f.Source().start_count_); } @@ -207,9 +223,11 @@ TEST(FDv1AdapterSynchronizerTest, CloseBeforeAnyNextDoesNotShutDown) { TEST(FDv1AdapterSynchronizerTest, NextAfterCloseIsShutdown) { AdapterFixture f; + // Ask for a result after closing. f.Adapter().Close(); auto result = f.Next(); + // The result is Shutdown, and the source is never started. ASSERT_TRUE(result.has_value()); EXPECT_TRUE( std::holds_alternative(result->value)); @@ -219,11 +237,14 @@ TEST(FDv1AdapterSynchronizerTest, NextAfterCloseIsShutdown) { TEST(FDv1AdapterSynchronizerTest, CloseUnblocksAPendingNext) { AdapterFixture f; + // Leave a Next() outstanding with nothing to deliver. auto future = f.Adapter().Next(launchdarkly::data_model::Selector{}); ASSERT_FALSE(future.IsFinished()); + // Close the adapter. f.Adapter().Close(); + // The outstanding Next() resolves with Shutdown. 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 a44dc1236..3fecf95c7 100644 --- a/libs/client-sdk/tests/fdv2_data_source_test.cpp +++ b/libs/client-sdk/tests/fdv2_data_source_test.cpp @@ -169,7 +169,6 @@ 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( @@ -291,8 +290,7 @@ 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. + // Runs the orchestration to a standstill without waiting on timers. void Drain() { ioc_.poll(); } DataSourceStatusManager& StatusManager() { return status_manager_; } flag_manager::FlagStore const& Store() { return flag_manager_.Store(); } @@ -939,11 +937,13 @@ TEST(ClientFDv2DataSourceTest, FallbackDirectiveOnAnInitializerAppliesItsData) { synchronizers.push_back( std::make_unique(std::move(fdv1_sources))); + // Run the initializer whose result carries the directive. auto source = h.MakeDataSource(std::move(initializers), std::move(synchronizers)); source->Start(); h.Drain(); + // Its payload is applied, even though the SDK moves off FDv2. EXPECT_TRUE(h.Store().Get("flagA")); } @@ -973,10 +973,12 @@ TEST(ClientFDv2DataSourceTest, FallbackDirectiveStartsTheFDv1Tier) { auto* fdv1_ptr = fdv1.get(); synchronizers.push_back(std::move(fdv1)); + // Run the synchronizer whose result carries the directive. auto source = h.MakeDataSource({}, std::move(synchronizers)); source->Start(); h.Drain(); + // The FDv1 tier is started and its data is applied. EXPECT_EQ(1, fdv1_ptr->build_count_); EXPECT_TRUE(h.Store().Get("from-fdv1")); } @@ -997,10 +999,12 @@ TEST(ClientFDv2DataSourceTest, FallbackWithNoFDv1TierDisconnects) { synchronizers.push_back(std::make_unique( std::make_unique(std::move(fdv2_results)))); + // Run with a directive but no FDv1 tier configured. auto source = h.MakeDataSource({}, std::move(synchronizers)); source->Start(); h.Drain(); + // There is nowhere to fall back to, so the SDK reports an interruption. EXPECT_EQ(DataSourceStatus::DataSourceState::kInterrupted, h.State()); // Flag data already received stays available for evaluation. EXPECT_TRUE(h.Store().Get("flagA")); @@ -1037,9 +1041,11 @@ TEST(ClientFDv2DataSourceTest, FDv2IsRetriedAfterTheFallbackTtl) { synchronizers.push_back( std::make_unique(std::move(fdv1_sources))); + // Run past the directive's one-second TTL. auto source = h.MakeDataSource({}, std::move(synchronizers)); source->Start(); h.Context().run(); + // FDv2 is started a second time. EXPECT_EQ(2, fdv2_ptr->build_count_); }