diff --git a/launchdarkly-server-sdk/src/fdv2/data_system.rs b/launchdarkly-server-sdk/src/fdv2/data_system.rs index 8af62e8..aa391e4 100644 --- a/launchdarkly-server-sdk/src/fdv2/data_system.rs +++ b/launchdarkly-server-sdk/src/fdv2/data_system.rs @@ -10,7 +10,7 @@ use crate::data_system::DataSystem; use crate::stores::store::{DataStore, InMemoryDataStore, TransactionalDataStore}; use super::model::Selector; -use super::source::{FDv2SourceResult, Initializer, Synchronizer}; +use super::source::{FDv2SourceEvent, FDv2SourceResult, Initializer, Synchronizer}; /// Produces a fresh initializer each time the orchestrator starts a run. pub(crate) trait InitializerFactory: Send + Sync { @@ -20,6 +20,11 @@ pub(crate) trait InitializerFactory: Send + Sync { /// Produces a fresh synchronizer each time the orchestrator starts a run. pub(crate) trait SynchronizerFactory: Send + Sync { fn create(&self) -> Box; + + /// Whether this factory builds the FDv1 fallback synchronizer. + fn is_fdv1_fallback(&self) -> bool { + false + } } /// FDv2 orchestrator: owns the memory store and keeps it populated by running @@ -98,7 +103,17 @@ struct SourceManager { impl SourceManager { fn new(factories: Vec>) -> Self { - let states = vec![SourceState::Available; factories.len()]; + // FDv1 fallback factories start blocked; they activate only on a directive. + let states = factories + .iter() + .map(|f| { + if f.is_fdv1_fallback() { + SourceState::Blocked + } else { + SourceState::Available + } + }) + .collect(); Self { factories, states, @@ -155,6 +170,36 @@ impl SourceManager { .filter(|s| **s == SourceState::Available) .count() } + + /// Blocks the FDv2 synchronizers and activates the FDv1 fallback, if any. + fn switch_to_fdv1_fallback(&mut self) { + for (i, factory) in self.factories.iter().enumerate() { + self.states[i] = if factory.is_fdv1_fallback() { + SourceState::Available + } else { + SourceState::Blocked + }; + } + self.synchronizer_index = None; + } + + /// Restores the initial state: FDv2 available, FDv1 fallback blocked. + fn switch_back_to_fdv2(&mut self) { + for (i, factory) in self.factories.iter().enumerate() { + self.states[i] = if factory.is_fdv1_fallback() { + SourceState::Blocked + } else { + SourceState::Available + }; + } + self.synchronizer_index = None; + } + + /// Whether the active factory is the FDv1 fallback. + fn is_current_fdv1_fallback(&self) -> bool { + self.current_factory_index + .is_some_and(|i| self.factories[i].is_fdv1_fallback()) + } } /// Sleeps until `at`, or never when `None` (an inactive timer arm). @@ -176,29 +221,72 @@ async fn run( ) { let mut selector: Selector = None; let mut initialized = false; + let mut fdv2_retry_at: Option = None; - // Initializer phase: try each in order until one yields a basis. + // Initializer phase: try each in order until one yields a basis or an + // FDv1 fallback directive. for mut initializer in initializers { let name = initializer.name().to_string(); let mut shutdown = Box::pin(shutdown_receiver.recv()).fuse(); - futures::select! { + let event = futures::select! { _ = shutdown => return, - event = initializer.run().fuse() => { - if let FDv2SourceResult::ChangeSet(change_set) = event.result { - selector = change_set.selector.clone(); - store.write().apply(change_set); - init_complete(true); - initialized = true; - break; - } - debug!("{name} did not provide a basis"); + event = initializer.run().fuse() => event, + }; + + let FDv2SourceEvent { + result, + fdv1_fallback, + } = event; + if let FDv2SourceResult::ChangeSet(change_set) = result { + selector = change_set.selector.clone(); + store.write().apply(change_set); + init_complete(true); + initialized = true; + } else { + debug!("{name} did not provide a basis"); + } + // An FDv1 fallback directive ends the initializer phase and switches + // to the fallback. + if let Some(fallback_directive) = fdv1_fallback { + info!("FDv2 falling back to the FDv1 protocol"); + source_manager.switch_to_fdv1_fallback(); + if fallback_directive.ttl != Duration::ZERO { + fdv2_retry_at = Some(Instant::now() + fallback_directive.ttl); } + break; + } + if initialized { + break; } } // Synchronizer phase: rotate through synchronizers as the timers fire. let mut current = source_manager.next_synchronizer(); - while let Some(mut active) = current { + loop { + let mut active = match current { + Some(active) => active, + None => { + // No synchronizer is available. While an FDv1 fallback directive's + // retry is pending, wait it out and return to FDv2 rather than exiting + // -- a blocked or terminal fallback must not strand the data system. + // With no retry pending, every source hit an unrecoverable error, so stop. + if fdv2_retry_at.is_none() { + break; + } + let mut shutdown = Box::pin(shutdown_receiver.recv()).fuse(); + let mut fdv2_retry = Box::pin(deadline(fdv2_retry_at)).fuse(); + futures::select! { + _ = shutdown => return, + _ = fdv2_retry => { + source_manager.switch_back_to_fdv2(); + fdv2_retry_at = None; + } + } + current = source_manager.next_synchronizer(); + continue; + } + }; + let name = active.name().to_string(); let has_fallback = source_manager.available_count() > 1; let has_recovery = has_fallback && !source_manager.is_prime(); @@ -211,6 +299,7 @@ async fn run( let mut shutdown = Box::pin(shutdown_receiver.recv()).fuse(); let mut fallback = Box::pin(deadline(fallback_at)).fuse(); let mut recovery = Box::pin(deadline(recovery_at)).fuse(); + let mut fdv2_retry = Box::pin(deadline(fdv2_retry_at)).fuse(); let mut next = active.next(selector.clone()).fuse(); futures::select! { _ = shutdown => return, @@ -221,38 +310,62 @@ async fn run( source_manager.reset_source_index(); break; } - event = next => match event.result { - FDv2SourceResult::ChangeSet(change_set) => { - selector = change_set.selector.clone(); - store.write().apply(change_set); - if !initialized { - init_complete(true); - initialized = true; + // Re-engage FDv2 once the fallback directive's TTL expires. + _ = fdv2_retry => { + source_manager.switch_back_to_fdv2(); + fdv2_retry_at = None; + break; + } + event = next => { + let FDv2SourceEvent { result, fdv1_fallback } = event; + let mut terminal = false; + match result { + FDv2SourceResult::ChangeSet(change_set) => { + selector = change_set.selector.clone(); + store.write().apply(change_set); + if !initialized { + init_complete(true); + initialized = true; + } + // Good data clears the fallback countdown. + fallback_at = None; + interrupted_logged = false; } - // Good data clears the fallback countdown. - fallback_at = None; - interrupted_logged = false; - } - // Sustained interruption starts the fallback countdown. - FDv2SourceResult::Interrupted(error) => { - if !interrupted_logged { - info!("{name} interrupted: {}", error.message); - interrupted_logged = true; + // Sustained interruption starts the fallback countdown. + FDv2SourceResult::Interrupted(error) => { + if !interrupted_logged { + info!("{name} interrupted: {}", error.message); + interrupted_logged = true; + } + if has_fallback && fallback_at.is_none() { + fallback_at = Some(Instant::now() + fallback_timeout); + } + } + // Handled internally by the synchronizer. + FDv2SourceResult::Goodbye => {} + FDv2SourceResult::TerminalError(error) => { + warn!("{name} terminal error: {}", error.message); + terminal = true; } - if has_fallback && fallback_at.is_none() { - fallback_at = Some(Instant::now() + fallback_timeout); + FDv2SourceResult::Shutdown => return, + } + // An FDv1 fallback directive takes precedence over a terminal error. + if let Some(fallback_directive) = fdv1_fallback { + if !source_manager.is_current_fdv1_fallback() { + info!("FDv2 falling back to the FDv1 protocol"); + source_manager.switch_to_fdv1_fallback(); + if fallback_directive.ttl != Duration::ZERO { + fdv2_retry_at = Some(Instant::now() + fallback_directive.ttl); + } + break; } } - // Handled internally by the synchronizer. - FDv2SourceResult::Goodbye => {} - FDv2SourceResult::TerminalError(error) => { - warn!("{name} terminal error: {}", error.message); + if terminal { // Dead source: drop it and advance. source_manager.block_current(); break; } - FDv2SourceResult::Shutdown => return, - }, + } } } @@ -276,7 +389,7 @@ mod tests { use launchdarkly_server_sdk_evaluation::Store; use super::super::model::ChangeSetKind; - use super::super::source::{ErrorInfo, ErrorKind, FDv2SourceEvent}; + use super::super::source::{ErrorInfo, ErrorKind, FDv1FallbackDirective, FDv2SourceEvent}; use crate::stores::change_set::{ChangeSet, ItemChange}; use crate::stores::store_types::StorageItem; use crate::test_common::basic_flag; @@ -337,17 +450,46 @@ mod tests { } } + /// An initializer that reports an FDv1 fallback directive with no TTL. + struct FallbackInitializer; + + impl Initializer for FallbackInitializer { + fn run(&mut self) -> BoxFuture<'_, FDv2SourceEvent> { + Box::pin(async { + FDv2SourceEvent { + result: interrupted(), + fdv1_fallback: Some(FDv1FallbackDirective { + ttl: Duration::ZERO, + }), + } + }) + } + + fn name(&self) -> &str { + "fallback-initializer" + } + } + struct MockSynchronizer { results: VecDeque, selectors_seen: Selectors, hang: bool, + fallback_directive: Option, } impl Synchronizer for MockSynchronizer { fn next(&mut self, selector: Selector) -> BoxFuture<'_, FDv2SourceEvent> { self.selectors_seen.lock().unwrap().push(selector); match self.results.pop_front() { - Some(result) => Box::pin(async move { event(result) }), + Some(result) => { + let fdv1_fallback = self.fallback_directive.clone(); + Box::pin(async move { + FDv2SourceEvent { + result, + fdv1_fallback, + } + }) + } // Out of scripted results: idle if hang, otherwise end the run. None if self.hang => Box::pin(std::future::pending()), None => Box::pin(async move { event(FDv2SourceResult::Shutdown) }), @@ -376,6 +518,8 @@ mod tests { results: Mutex>, selectors_seen: Selectors, hang: bool, + is_fdv1_fallback: bool, + fallback_directive: Option, } impl SynchronizerFactory for MockSynchronizerFactory { @@ -385,8 +529,13 @@ mod tests { results: results.into(), selectors_seen: self.selectors_seen.clone(), hang: self.hang, + fallback_directive: self.fallback_directive.clone(), }) } + + fn is_fdv1_fallback(&self) -> bool { + self.is_fdv1_fallback + } } /// A selector recorder that ignores what it captures. @@ -404,6 +553,34 @@ mod tests { results: Mutex::new(results), selectors_seen, hang, + is_fdv1_fallback: false, + fallback_directive: None, + }) + } + + /// Like `sync_factory`, but flagged as the FDv1 fallback. + fn fdv1_factory( + results: Vec, + selectors_seen: Selectors, + hang: bool, + ) -> Arc { + Arc::new(MockSynchronizerFactory { + results: Mutex::new(results), + selectors_seen, + hang, + is_fdv1_fallback: true, + fallback_directive: None, + }) + } + + /// An FDv2 synchronizer that reports an FDv1 fallback directive. + fn fallback_directive_factory(ttl: Duration) -> Arc { + Arc::new(MockSynchronizerFactory { + results: Mutex::new(vec![interrupted()]), + selectors_seen: no_selectors(), + hang: true, + is_fdv1_fallback: false, + fallback_directive: Some(FDv1FallbackDirective { ttl }), }) } @@ -432,10 +609,43 @@ mod tests { results: results.into(), selectors_seen: no_selectors(), hang, + fallback_directive: None, }) } } + /// Reports an FDv1 fallback directive on its first build, then delivers data. + struct FallbackThenDataFactory { + ttl: Duration, + builds: AtomicUsize, + } + + impl SynchronizerFactory for FallbackThenDataFactory { + fn create(&self) -> Box { + if self.builds.fetch_add(1, Ordering::SeqCst) == 0 { + // First: report the fallback directive, then idle. + Box::new(MockSynchronizer { + results: VecDeque::from(vec![interrupted()]), + selectors_seen: no_selectors(), + hang: true, + fallback_directive: Some(FDv1FallbackDirective { ttl: self.ttl }), + }) + } else { + // After the FDv2 retry: deliver a delta. + Box::new(MockSynchronizer { + results: VecDeque::from(vec![changeset( + ChangeSetKind::Partial, + "fdv2-back", + Some("s".into()), + )]), + selectors_seen: no_selectors(), + hang: false, + fallback_directive: None, + }) + } + } + } + fn recording_init_complete() -> (Arc, InitCalls) { let calls: InitCalls = Arc::new(Mutex::new(Vec::new())); let sink = calls.clone(); @@ -773,6 +983,57 @@ mod tests { assert_eq!(sources.available_count(), 1); } + #[test] + fn fdv1_fallback_factory_starts_blocked() { + let mut sources = SourceManager::new(vec![ + sync_factory(vec![], no_selectors(), false), + fdv1_factory(vec![], no_selectors(), false), + ]); + + // Only the FDv2 synchronizer is available; the FDv1 fallback is dormant. + assert_eq!(sources.available_count(), 1); + sources.next_synchronizer(); + assert!(!sources.is_current_fdv1_fallback()); + } + + #[test] + fn switch_to_and_back_from_fdv1_fallback() { + let mut sources = SourceManager::new(vec![ + sync_factory(vec![], no_selectors(), false), + fdv1_factory(vec![], no_selectors(), false), + ]); + + // The directive activates only the FDv1 fallback. + sources.switch_to_fdv1_fallback(); + assert_eq!(sources.available_count(), 1); + sources.next_synchronizer(); + assert!(sources.is_current_fdv1_fallback()); + + // Switching back restores the FDv2 synchronizer and re-blocks FDv1. + sources.switch_back_to_fdv2(); + assert_eq!(sources.available_count(), 1); + sources.next_synchronizer(); + assert!(!sources.is_current_fdv1_fallback()); + } + + #[test] + fn switch_back_unblocks_terminally_blocked_fdv2() { + let mut sources = SourceManager::new(vec![ + sync_factory(vec![], no_selectors(), false), + fdv1_factory(vec![], no_selectors(), false), + ]); + + // A terminal error blocks the only FDv2 synchronizer. + sources.next_synchronizer(); + sources.block_current(); + assert_eq!(sources.available_count(), 0); + + // A fallback round-trip makes the FDv2 synchronizer usable again. + sources.switch_to_fdv1_fallback(); + sources.switch_back_to_fdv2(); + assert_eq!(sources.available_count(), 1); + } + #[tokio::test(start_paused = true)] async fn fallback_fires_after_sustained_interruption() { // No initializers; the prime interrupts then idles. @@ -910,4 +1171,156 @@ mod tests { assert!(store.read().flag("from-fallback").is_some()); assert!(store.read().flag("prime-recovered").is_some()); } + + #[tokio::test] + async fn fallback_directive_switches_to_fdv1() { + let store = Arc::new(RwLock::new(InMemoryDataStore::new())); + let (init_complete, calls) = recording_init_complete(); + let initializers: Vec> = vec![]; + + // The FDv2 source reports a fallback directive with a zero TTL, so FDv2 + // is never retried; the FDv1 fallback holds the basis. + let source_manager = SourceManager::new(vec![ + fallback_directive_factory(Duration::ZERO), + fdv1_factory( + vec![changeset( + ChangeSetKind::Full, + "from-fdv1", + Some("s".into()), + )], + no_selectors(), + false, + ), + ]); + let (_shutdown_tx, shutdown_rx) = broadcast::channel(1); + + run( + initializers, + source_manager, + store.clone(), + init_complete, + shutdown_rx, + FALLBACK_TIMEOUT, + RECOVERY_TIMEOUT, + ) + .await; + + // Switched to the FDv1 fallback and applied its basis. + assert_eq!(*calls.lock().unwrap(), vec![true]); + assert!(store.read().flag("from-fdv1").is_some()); + } + + #[tokio::test(start_paused = true)] + async fn fdv2_retry_reengages_after_ttl() { + let store = Arc::new(RwLock::new(InMemoryDataStore::new())); + let (init_complete, calls) = recording_init_complete(); + let initializers: Vec> = vec![]; + + // The FDv2 source reports a directive with a 60s TTL; the FDv1 fallback + // supplies a basis and then idles until FDv2 is retried. + let source_manager = SourceManager::new(vec![ + Arc::new(FallbackThenDataFactory { + ttl: Duration::from_secs(60), + builds: AtomicUsize::new(0), + }) as Arc, + fdv1_factory( + vec![changeset( + ChangeSetKind::Full, + "from-fdv1", + Some("s".into()), + )], + no_selectors(), + true, + ), + ]); + let (_shutdown_tx, shutdown_rx) = broadcast::channel(1); + + // Paused time auto-advances past the TTL, re-engaging FDv2. + let handle = tokio::spawn(run( + initializers, + source_manager, + store.clone(), + init_complete, + shutdown_rx, + FALLBACK_TIMEOUT, + RECOVERY_TIMEOUT, + )); + handle.await.unwrap(); + + // Fell back to FDv1, then re-engaged FDv2 once the TTL expired. + assert_eq!(*calls.lock().unwrap(), vec![true]); + assert!(store.read().flag("from-fdv1").is_some()); + assert!(store.read().flag("fdv2-back").is_some()); + } + + #[tokio::test(start_paused = true)] + async fn fdv2_retry_survives_a_terminal_fdv1_fallback() { + let store = Arc::new(RwLock::new(InMemoryDataStore::new())); + let (init_complete, calls) = recording_init_complete(); + let initializers: Vec> = vec![]; + + // The FDv2 source reports a 60s directive; the FDv1 fallback then fails + // terminally, blocking every source before the TTL expires. + let source_manager = SourceManager::new(vec![ + Arc::new(FallbackThenDataFactory { + ttl: Duration::from_secs(60), + builds: AtomicUsize::new(0), + }) as Arc, + fdv1_factory(vec![terminal()], no_selectors(), false), + ]); + let (_shutdown_tx, shutdown_rx) = broadcast::channel(1); + + // Paused time advances past the TTL while no source is available. + let handle = tokio::spawn(run( + initializers, + source_manager, + store.clone(), + init_complete, + shutdown_rx, + FALLBACK_TIMEOUT, + RECOVERY_TIMEOUT, + )); + handle.await.unwrap(); + + // The blocked fallback did not strand the run; FDv2 re-engaged after the TTL. + assert_eq!(*calls.lock().unwrap(), vec![true]); + assert!(store.read().flag("fdv2-back").is_some()); + } + + #[tokio::test] + async fn initializer_fallback_directive_switches_to_fdv1() { + let store = Arc::new(RwLock::new(InMemoryDataStore::new())); + let (init_complete, calls) = recording_init_complete(); + let initializers: Vec> = vec![Box::new(FallbackInitializer)]; + + // The FDv2 synchronizer is never reached; the FDv1 fallback holds the basis. + let source_manager = SourceManager::new(vec![ + sync_factory(vec![], no_selectors(), false), + fdv1_factory( + vec![changeset( + ChangeSetKind::Full, + "from-fdv1", + Some("s".into()), + )], + no_selectors(), + false, + ), + ]); + let (_shutdown_tx, shutdown_rx) = broadcast::channel(1); + + run( + initializers, + source_manager, + store.clone(), + init_complete, + shutdown_rx, + FALLBACK_TIMEOUT, + RECOVERY_TIMEOUT, + ) + .await; + + // The initializer's directive switched straight to the FDv1 fallback. + assert_eq!(*calls.lock().unwrap(), vec![true]); + assert!(store.read().flag("from-fdv1").is_some()); + } }