Skip to main content

launchdarkly_server_sdk/fdv2/
data_system.rs

1use std::sync::Arc;
2use std::time::Duration;
3
4use futures::FutureExt;
5use parking_lot::RwLock;
6use tokio::sync::broadcast;
7use tokio::time::{sleep_until, Instant};
8
9use crate::data_system::DataSystem;
10use crate::stores::store::{DataStore, InMemoryDataStore, TransactionalDataStore};
11
12use super::model::{ChangeSetKind, Selector};
13use super::source::{FDv2SourceEvent, FDv2SourceResult, Initializer, Synchronizer};
14
15/// Produces a fresh initializer each time the orchestrator starts a run.
16pub trait InitializerFactory: Send + Sync {
17    /// Builds a new initializer instance.
18    fn create(&self) -> Box<dyn Initializer>;
19}
20
21/// Produces a fresh synchronizer each time the orchestrator starts a run.
22pub trait SynchronizerFactory: Send + Sync {
23    /// Builds a new synchronizer instance.
24    fn create(&self) -> Box<dyn Synchronizer>;
25
26    /// Whether this factory builds the FDv1 fallback synchronizer.
27    fn is_fdv1_fallback(&self) -> bool {
28        false
29    }
30}
31
32/// FDv2 orchestrator: owns the memory store and keeps it populated by running
33/// initializers to obtain a basis, then synchronizers for ongoing changes.
34pub(crate) struct FDv2DataSystem {
35    initializer_factories: Vec<Arc<dyn InitializerFactory>>,
36    synchronizer_factories: Vec<Arc<dyn SynchronizerFactory>>,
37    fallback_timeout: Duration,
38    recovery_timeout: Duration,
39    store: Arc<RwLock<InMemoryDataStore>>,
40}
41
42impl FDv2DataSystem {
43    pub(crate) fn new(
44        initializer_factories: Vec<Arc<dyn InitializerFactory>>,
45        synchronizer_factories: Vec<Arc<dyn SynchronizerFactory>>,
46        fallback_timeout: Duration,
47        recovery_timeout: Duration,
48    ) -> Self {
49        Self {
50            initializer_factories,
51            synchronizer_factories,
52            fallback_timeout,
53            recovery_timeout,
54            store: Arc::new(RwLock::new(InMemoryDataStore::new())),
55        }
56    }
57}
58
59impl DataSystem for FDv2DataSystem {
60    fn start(
61        &self,
62        init_complete: Arc<dyn Fn(bool) + Send + Sync>,
63        shutdown_receiver: broadcast::Receiver<()>,
64    ) {
65        let initializer_factories = self.initializer_factories.clone();
66        let source_manager = SourceManager::new(self.synchronizer_factories.clone());
67        let store = self.store.clone();
68
69        tokio::spawn(run(
70            initializer_factories,
71            source_manager,
72            store,
73            init_complete,
74            shutdown_receiver,
75            self.fallback_timeout,
76            self.recovery_timeout,
77        ));
78    }
79
80    fn store(&self) -> Arc<RwLock<dyn DataStore>> {
81        self.store.clone()
82    }
83}
84
85/// Per-factory availability used by the synchronizer rotation.
86#[derive(Clone, Copy, PartialEq, Eq)]
87enum SourceState {
88    Available,
89    Blocked,
90}
91
92/// Owns the synchronizer factories and tracks which one is currently active.
93struct SourceManager {
94    factories: Vec<Arc<dyn SynchronizerFactory>>,
95    states: Vec<SourceState>,
96    /// Iteration cursor; `None` restarts the search from the prime.
97    synchronizer_index: Option<usize>,
98    /// Index of the most recently returned factory, for blocking and prime checks.
99    current_factory_index: Option<usize>,
100}
101
102impl SourceManager {
103    fn new(factories: Vec<Arc<dyn SynchronizerFactory>>) -> Self {
104        // FDv1 fallback factories start blocked; they activate only on a directive.
105        let states = factories
106            .iter()
107            .map(|f| {
108                if f.is_fdv1_fallback() {
109                    SourceState::Blocked
110                } else {
111                    SourceState::Available
112                }
113            })
114            .collect();
115        Self {
116            factories,
117            states,
118            synchronizer_index: None,
119            current_factory_index: None,
120        }
121    }
122
123    /// Builds the next available synchronizer, advancing cyclically past the
124    /// active one and skipping blocked factories. `None` when all are blocked.
125    fn next_synchronizer(&mut self) -> Option<Box<dyn Synchronizer>> {
126        let n = self.factories.len();
127        if n == 0 {
128            self.current_factory_index = None;
129            return None;
130        }
131        let mut i = self.synchronizer_index.map_or(0, |c| (c + 1) % n);
132        for _ in 0..n {
133            if self.states[i] == SourceState::Available {
134                self.synchronizer_index = Some(i);
135                self.current_factory_index = Some(i);
136                return Some(self.factories[i].create());
137            }
138            i = (i + 1) % n;
139        }
140        self.current_factory_index = None;
141        None
142    }
143
144    /// Marks the active factory blocked, used on a terminal error.
145    fn block_current(&mut self) {
146        if let Some(i) = self.current_factory_index {
147            self.states[i] = SourceState::Blocked;
148        }
149    }
150
151    /// Makes the next `next_synchronizer` restart from the prime, used on recovery.
152    fn reset_source_index(&mut self) {
153        self.synchronizer_index = None;
154    }
155
156    /// Whether the active factory is the first available one.
157    fn is_prime(&self) -> bool {
158        let first = self
159            .states
160            .iter()
161            .position(|s| *s == SourceState::Available);
162        first == self.current_factory_index && first.is_some()
163    }
164
165    fn available_count(&self) -> usize {
166        self.states
167            .iter()
168            .filter(|s| **s == SourceState::Available)
169            .count()
170    }
171
172    /// Blocks the FDv2 synchronizers and activates the FDv1 fallback, if any.
173    fn switch_to_fdv1_fallback(&mut self) {
174        for (i, factory) in self.factories.iter().enumerate() {
175            self.states[i] = if factory.is_fdv1_fallback() {
176                SourceState::Available
177            } else {
178                SourceState::Blocked
179            };
180        }
181        self.synchronizer_index = None;
182    }
183
184    /// Restores the initial state: FDv2 available, FDv1 fallback blocked.
185    fn switch_back_to_fdv2(&mut self) {
186        for (i, factory) in self.factories.iter().enumerate() {
187            self.states[i] = if factory.is_fdv1_fallback() {
188                SourceState::Blocked
189            } else {
190                SourceState::Available
191            };
192        }
193        self.synchronizer_index = None;
194    }
195
196    /// Whether the active factory is the FDv1 fallback.
197    fn is_current_fdv1_fallback(&self) -> bool {
198        self.current_factory_index
199            .is_some_and(|i| self.factories[i].is_fdv1_fallback())
200    }
201}
202
203/// Sleeps until `at`, or never when `None` (an inactive timer arm).
204async fn deadline(at: Option<Instant>) {
205    match at {
206        Some(t) => sleep_until(t).await,
207        None => std::future::pending::<()>().await,
208    }
209}
210
211async fn run(
212    initializer_factories: Vec<Arc<dyn InitializerFactory>>,
213    mut source_manager: SourceManager,
214    store: Arc<RwLock<InMemoryDataStore>>,
215    init_complete: Arc<dyn Fn(bool) + Send + Sync>,
216    mut shutdown_receiver: broadcast::Receiver<()>,
217    fallback_timeout: Duration,
218    recovery_timeout: Duration,
219) {
220    let mut selector: Selector = None;
221    let mut initialized = false;
222    // Whether an initializer produced a full payload.
223    let mut got_full = false;
224    let mut fdv2_retry_at: Option<Instant> = None;
225
226    // Initializer phase: try each until one yields a basis or an FDv1 fallback directive.
227    for factory in initializer_factories {
228        let mut initializer = factory.create();
229        let name = initializer.name().to_string();
230        let mut shutdown = Box::pin(shutdown_receiver.recv()).fuse();
231        let event = futures::select! {
232            _ = shutdown => return,
233            event = initializer.run().fuse() => event,
234        };
235
236        let FDv2SourceEvent {
237            result,
238            fdv1_fallback,
239        } = event;
240        let mut has_basis = false;
241        match result {
242            FDv2SourceResult::ChangeSet(change_set) => {
243                let is_full = matches!(change_set.kind, ChangeSetKind::Full);
244                has_basis = is_full && change_set.selector.is_some();
245                if !matches!(change_set.kind, ChangeSetKind::None) {
246                    selector = change_set.selector.clone();
247                }
248                store.write().apply(change_set);
249                if is_full {
250                    got_full = true;
251                }
252            }
253            _ => debug!("{name} did not provide a basis"),
254        }
255        // An FDv1 fallback directive ends the initializer phase and activates the fallback.
256        if let Some(fallback_directive) = fdv1_fallback {
257            info!("FDv2 falling back to the FDv1 protocol");
258            source_manager.switch_to_fdv1_fallback();
259            fdv2_retry_at = Some(Instant::now() + fallback_directive.ttl);
260            break;
261        }
262        if has_basis {
263            break;
264        }
265    }
266
267    if got_full && !initialized {
268        init_complete(true);
269        initialized = true;
270    }
271
272    // Synchronizer phase: rotate through synchronizers as the timers fire.
273    let mut current = source_manager.next_synchronizer();
274    loop {
275        let mut active = match current {
276            Some(active) => active,
277            None => {
278                // No synchronizer is available. While an FDv1 fallback directive's
279                // retry is pending, wait it out and return to FDv2 rather than exiting.
280                // A blocked or terminal fallback must not strand the data system.
281                // With no retry pending, every source hit an unrecoverable error, so stop.
282                if fdv2_retry_at.is_none() {
283                    break;
284                }
285                let mut shutdown = Box::pin(shutdown_receiver.recv()).fuse();
286                let mut fdv2_retry = Box::pin(deadline(fdv2_retry_at)).fuse();
287                futures::select! {
288                    _ = shutdown => return,
289                    _ = fdv2_retry => {
290                        source_manager.switch_back_to_fdv2();
291                        fdv2_retry_at = None;
292                    }
293                }
294                current = source_manager.next_synchronizer();
295                continue;
296            }
297        };
298
299        let name = active.name().to_string();
300        let has_fallback = source_manager.available_count() > 1;
301        let has_recovery = has_fallback && !source_manager.is_prime();
302        let mut fallback_at: Option<Instant> = None;
303        let recovery_at = has_recovery.then(|| Instant::now() + recovery_timeout);
304        let mut interrupted_logged = false;
305
306        // Drive the active synchronizer, racing its events against the timers.
307        loop {
308            let mut shutdown = Box::pin(shutdown_receiver.recv()).fuse();
309            let mut fallback = Box::pin(deadline(fallback_at)).fuse();
310            let mut recovery = Box::pin(deadline(recovery_at)).fuse();
311            let mut fdv2_retry = Box::pin(deadline(fdv2_retry_at)).fuse();
312            let mut next = active.next(selector.clone()).fuse();
313            futures::select! {
314                _ = shutdown => return,
315                // Fall back to the next synchronizer.
316                _ = fallback => break,
317                // Recover to the prime.
318                _ = recovery => {
319                    source_manager.reset_source_index();
320                    break;
321                }
322                // Re-engage FDv2 once the fallback directive's TTL expires.
323                _ = fdv2_retry => {
324                    source_manager.switch_back_to_fdv2();
325                    fdv2_retry_at = None;
326                    break;
327                }
328                event = next => {
329                    let FDv2SourceEvent { result, fdv1_fallback } = event;
330                    let mut terminal = false;
331                    match result {
332                        FDv2SourceResult::ChangeSet(change_set) => {
333                            let is_full = matches!(change_set.kind, ChangeSetKind::Full);
334                            if !matches!(change_set.kind, ChangeSetKind::None) {
335                                selector = change_set.selector.clone();
336                                store.write().apply(change_set);
337                                if is_full && !initialized {
338                                    init_complete(true);
339                                    initialized = true;
340                                }
341                            }
342                            // A successful response clears the countdown.
343                            fallback_at = None;
344                            interrupted_logged = false;
345                        }
346                        // Sustained interruption starts the fallback countdown.
347                        FDv2SourceResult::Interrupted(error) => {
348                            if !interrupted_logged {
349                                info!("{name} interrupted: {}", error.message);
350                                interrupted_logged = true;
351                            }
352                            if has_fallback && fallback_at.is_none() {
353                                fallback_at = Some(Instant::now() + fallback_timeout);
354                            }
355                        }
356                        // Handled internally by the synchronizer.
357                        FDv2SourceResult::Goodbye => {}
358                        FDv2SourceResult::TerminalError(error) => {
359                            warn!("{name} terminal error: {}", error.message);
360                            terminal = true;
361                        }
362                    }
363                    // An FDv1 fallback directive takes precedence over a terminal error.
364                    if let Some(fallback_directive) = fdv1_fallback {
365                        if !source_manager.is_current_fdv1_fallback() {
366                            info!("FDv2 falling back to the FDv1 protocol");
367                            source_manager.switch_to_fdv1_fallback();
368                            fdv2_retry_at = Some(Instant::now() + fallback_directive.ttl);
369                            break;
370                        }
371                    }
372                    if terminal {
373                        // Dead source: drop it and advance.
374                        source_manager.block_current();
375                        break;
376                    }
377                }
378            }
379        }
380
381        current = source_manager.next_synchronizer();
382    }
383
384    // Every source blocked without ever obtaining a full payload.
385    if !initialized {
386        init_complete(false);
387    }
388}
389
390#[cfg(test)]
391mod tests {
392    use super::*;
393    use std::collections::VecDeque;
394    use std::sync::atomic::{AtomicUsize, Ordering};
395    use std::sync::Mutex;
396
397    use futures::future::BoxFuture;
398    use launchdarkly_server_sdk_evaluation::Store;
399
400    use super::super::model::ChangeSetKind;
401    use super::super::source::{ErrorInfo, ErrorKind, FDv1FallbackDirective, FDv2SourceEvent};
402    use crate::stores::change_set::{ChangeSet, ItemChange};
403    use crate::stores::store_types::StorageItem;
404    use crate::test_common::basic_flag;
405
406    const FALLBACK_TIMEOUT: Duration = Duration::from_secs(120);
407    const RECOVERY_TIMEOUT: Duration = Duration::from_secs(300);
408
409    type Selectors = Arc<Mutex<Vec<Selector>>>;
410    type InitCalls = Arc<Mutex<Vec<bool>>>;
411
412    fn changeset(kind: ChangeSetKind, key: &str, selector: Selector) -> FDv2SourceResult {
413        FDv2SourceResult::ChangeSet(ChangeSet {
414            kind,
415            changes: vec![ItemChange::Flag {
416                key: key.to_string(),
417                item: StorageItem::Item(basic_flag(key)),
418            }],
419            selector,
420        })
421    }
422
423    fn interrupted() -> FDv2SourceResult {
424        FDv2SourceResult::Interrupted(ErrorInfo {
425            kind: ErrorKind::Unknown,
426            message: "test".into(),
427        })
428    }
429
430    fn terminal() -> FDv2SourceResult {
431        FDv2SourceResult::TerminalError(ErrorInfo {
432            kind: ErrorKind::Unknown,
433            message: "test".into(),
434        })
435    }
436
437    fn event(result: FDv2SourceResult) -> FDv2SourceEvent {
438        FDv2SourceEvent {
439            result,
440            fdv1_fallback: None,
441        }
442    }
443
444    struct MockInitializer {
445        results: VecDeque<FDv2SourceResult>,
446    }
447
448    impl Initializer for MockInitializer {
449        fn run(&mut self) -> BoxFuture<'_, FDv2SourceEvent> {
450            let result = self.results.pop_front().unwrap_or_else(interrupted);
451            Box::pin(async move { event(result) })
452        }
453
454        fn name(&self) -> &str {
455            "mock-initializer"
456        }
457    }
458
459    /// An initializer that reports an FDv1 fallback directive with the given TTL.
460    struct FallbackInitializer {
461        ttl: Duration,
462    }
463
464    impl Initializer for FallbackInitializer {
465        fn run(&mut self) -> BoxFuture<'_, FDv2SourceEvent> {
466            let ttl = self.ttl;
467            Box::pin(async move {
468                FDv2SourceEvent {
469                    result: interrupted(),
470                    fdv1_fallback: Some(FDv1FallbackDirective { ttl }),
471                }
472            })
473        }
474
475        fn name(&self) -> &str {
476            "fallback-initializer"
477        }
478    }
479
480    struct FallbackInitializerFactory {
481        ttl: Duration,
482    }
483
484    impl InitializerFactory for FallbackInitializerFactory {
485        fn create(&self) -> Box<dyn Initializer> {
486            Box::new(FallbackInitializer { ttl: self.ttl })
487        }
488    }
489
490    struct MockSynchronizer {
491        results: VecDeque<FDv2SourceResult>,
492        selectors_seen: Selectors,
493        // When set, the run ends through this once the script is exhausted. When None,
494        // the synchronizer idles so the test's timers can fire.
495        shutdown: Option<broadcast::Sender<()>>,
496        fallback_directive: Option<FDv1FallbackDirective>,
497    }
498
499    impl Synchronizer for MockSynchronizer {
500        fn next(&mut self, selector: Selector) -> BoxFuture<'_, FDv2SourceEvent> {
501            self.selectors_seen.lock().unwrap().push(selector);
502            match self.results.pop_front() {
503                Some(result) => {
504                    let fdv1_fallback = self.fallback_directive.clone();
505                    Box::pin(async move {
506                        FDv2SourceEvent {
507                            result,
508                            fdv1_fallback,
509                        }
510                    })
511                }
512                // Out of scripted results: end the run through the shutdown channel,
513                // or idle when there is no sender.
514                None => {
515                    if let Some(shutdown) = &self.shutdown {
516                        let _ = shutdown.send(());
517                    }
518                    Box::pin(std::future::pending())
519                }
520            }
521        }
522
523        fn name(&self) -> &str {
524            "mock-synchronizer"
525        }
526    }
527
528    struct MockInitializerFactory {
529        results: Mutex<Vec<FDv2SourceResult>>,
530    }
531
532    impl InitializerFactory for MockInitializerFactory {
533        fn create(&self) -> Box<dyn Initializer> {
534            let results = std::mem::take(&mut *self.results.lock().unwrap());
535            Box::new(MockInitializer {
536                results: results.into(),
537            })
538        }
539    }
540
541    /// A single-initializer factory scripted with the given results.
542    fn init_factory(results: Vec<FDv2SourceResult>) -> Arc<dyn InitializerFactory> {
543        Arc::new(MockInitializerFactory {
544            results: Mutex::new(results),
545        })
546    }
547
548    struct MockSynchronizerFactory {
549        results: Mutex<Vec<FDv2SourceResult>>,
550        selectors_seen: Selectors,
551        shutdown: Option<broadcast::Sender<()>>,
552        is_fdv1_fallback: bool,
553        fallback_directive: Option<FDv1FallbackDirective>,
554    }
555
556    impl SynchronizerFactory for MockSynchronizerFactory {
557        fn create(&self) -> Box<dyn Synchronizer> {
558            let results = std::mem::take(&mut *self.results.lock().unwrap());
559            Box::new(MockSynchronizer {
560                results: results.into(),
561                selectors_seen: self.selectors_seen.clone(),
562                shutdown: self.shutdown.clone(),
563                fallback_directive: self.fallback_directive.clone(),
564            })
565        }
566
567        fn is_fdv1_fallback(&self) -> bool {
568            self.is_fdv1_fallback
569        }
570    }
571
572    /// A selector recorder that ignores what it captures.
573    fn no_selectors() -> Selectors {
574        Arc::new(Mutex::new(Vec::new()))
575    }
576
577    /// A single-synchronizer factory scripted with the given results.
578    fn sync_factory(
579        results: Vec<FDv2SourceResult>,
580        selectors_seen: Selectors,
581        shutdown: Option<broadcast::Sender<()>>,
582        is_fdv1_fallback: bool,
583    ) -> Arc<dyn SynchronizerFactory> {
584        Arc::new(MockSynchronizerFactory {
585            results: Mutex::new(results),
586            selectors_seen,
587            shutdown,
588            is_fdv1_fallback,
589            fallback_directive: None,
590        })
591    }
592
593    /// An FDv2 synchronizer that reports an FDv1 fallback directive.
594    fn fallback_directive_factory(ttl: Duration) -> Arc<dyn SynchronizerFactory> {
595        Arc::new(MockSynchronizerFactory {
596            results: Mutex::new(vec![interrupted()]),
597            selectors_seen: no_selectors(),
598            shutdown: None,
599            is_fdv1_fallback: false,
600            fallback_directive: Some(FDv1FallbackDirective { ttl }),
601        })
602    }
603
604    /// A factory that is down on its first build and delivers data on rebuild.
605    struct DownThenDataFactory {
606        builds: AtomicUsize,
607        shutdown: broadcast::Sender<()>,
608    }
609
610    impl SynchronizerFactory for DownThenDataFactory {
611        fn create(&self) -> Box<dyn Synchronizer> {
612            let (results, shutdown) = if self.builds.fetch_add(1, Ordering::SeqCst) == 0 {
613                // Down: interrupt, then idle so the fallback timer fires.
614                (vec![interrupted()], None)
615            } else {
616                // Rebuilt: deliver a delta, then end the run.
617                (
618                    vec![changeset(
619                        ChangeSetKind::Partial,
620                        "prime-recovered",
621                        Some("s".into()),
622                    )],
623                    Some(self.shutdown.clone()),
624                )
625            };
626            Box::new(MockSynchronizer {
627                results: results.into(),
628                selectors_seen: no_selectors(),
629                shutdown,
630                fallback_directive: None,
631            })
632        }
633    }
634
635    /// Reports an FDv1 fallback directive on its first build, then delivers data.
636    struct FallbackThenDataFactory {
637        ttl: Duration,
638        reengage_kind: ChangeSetKind,
639        builds: AtomicUsize,
640        shutdown: broadcast::Sender<()>,
641    }
642
643    impl SynchronizerFactory for FallbackThenDataFactory {
644        fn create(&self) -> Box<dyn Synchronizer> {
645            if self.builds.fetch_add(1, Ordering::SeqCst) == 0 {
646                // First: report the fallback directive, then idle.
647                Box::new(MockSynchronizer {
648                    results: VecDeque::from(vec![interrupted()]),
649                    selectors_seen: no_selectors(),
650                    shutdown: None,
651                    fallback_directive: Some(FDv1FallbackDirective { ttl: self.ttl }),
652                })
653            } else {
654                // After the FDv2 retry: deliver the re-engagement payload, then end the run.
655                Box::new(MockSynchronizer {
656                    results: VecDeque::from(vec![changeset(
657                        self.reengage_kind,
658                        "fdv2-back",
659                        Some("s".into()),
660                    )]),
661                    selectors_seen: no_selectors(),
662                    shutdown: Some(self.shutdown.clone()),
663                    fallback_directive: None,
664                })
665            }
666        }
667    }
668
669    fn recording_init_complete() -> (Arc<dyn Fn(bool) + Send + Sync>, InitCalls) {
670        let calls: InitCalls = Arc::new(Mutex::new(Vec::new()));
671        let sink = calls.clone();
672        let cb: Arc<dyn Fn(bool) + Send + Sync> =
673            Arc::new(move |success| sink.lock().unwrap().push(success));
674        (cb, calls)
675    }
676
677    #[tokio::test]
678    async fn start_applies_basis_and_exposes_it_via_store_handle() {
679        // A data system whose sole initializer yields one full basis.
680        let system = FDv2DataSystem::new(
681            vec![init_factory(vec![changeset(
682                ChangeSetKind::Full,
683                "f1",
684                Some("s1".into()),
685            )])],
686            vec![sync_factory(
687                vec![],
688                no_selectors(),
689                /* idle */ None,
690                /* is_fdv1_fallback = */ false,
691            )],
692            FALLBACK_TIMEOUT,
693            RECOVERY_TIMEOUT,
694        );
695
696        // Record init-complete calls and wake the test when one arrives.
697        let calls: InitCalls = Arc::new(Mutex::new(Vec::new()));
698        let notify = Arc::new(tokio::sync::Notify::new());
699        let sink = calls.clone();
700        let waker = notify.clone();
701        let init_complete: Arc<dyn Fn(bool) + Send + Sync> = Arc::new(move |success| {
702            sink.lock().unwrap().push(success);
703            waker.notify_one();
704        });
705        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
706
707        // Start the system and wait for initialization to finish.
708        system.start(init_complete, shutdown_rx);
709        notify.notified().await;
710
711        // The basis was applied and is readable through the store() handle.
712        assert_eq!(*calls.lock().unwrap(), vec![true]);
713        assert!(system.store().read().flag("f1").is_some());
714        drop(shutdown_tx);
715    }
716
717    #[tokio::test]
718    async fn initializer_basis_signals_once_and_propagates_selector() {
719        // Store and recorders shared with the mock sources.
720        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
721        let selectors_seen: Selectors = Arc::new(Mutex::new(Vec::new()));
722        let (init_complete, calls) = recording_init_complete();
723
724        // Initializer delivers a full basis carrying selector "sel-1".
725        let initializer_factories: Vec<Arc<dyn InitializerFactory>> =
726            vec![init_factory(vec![changeset(
727                ChangeSetKind::Full,
728                "init-flag",
729                Some("sel-1".into()),
730            )])];
731
732        // Synchronizer delivers a partial change carrying selector "sel-2".
733        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
734        let source_manager = SourceManager::new(vec![sync_factory(
735            vec![changeset(
736                ChangeSetKind::Partial,
737                "sync-flag",
738                Some("sel-2".into()),
739            )],
740            selectors_seen.clone(),
741            Some(shutdown_tx),
742            /* is_fdv1_fallback = */ false,
743        )]);
744
745        run(
746            initializer_factories,
747            source_manager,
748            store.clone(),
749            init_complete,
750            shutdown_rx,
751            FALLBACK_TIMEOUT,
752            RECOVERY_TIMEOUT,
753        )
754        .await;
755
756        // init_complete fired exactly once despite two successful applies.
757        assert_eq!(*calls.lock().unwrap(), vec![true]);
758
759        // Both the basis flag and the later partial change are in the store.
760        assert!(store.read().flag("init-flag").is_some());
761        assert!(store.read().flag("sync-flag").is_some());
762
763        // The synchronizer's first call received the basis selector.
764        assert_eq!(selectors_seen.lock().unwrap()[0], Some("sel-1".into()));
765    }
766
767    #[tokio::test]
768    async fn selectorless_full_continues_to_next_initializer() {
769        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
770        let selectors_seen: Selectors = Arc::new(Mutex::new(Vec::new()));
771        let (init_complete, calls) = recording_init_complete();
772
773        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![
774            // The first initializer delivers a full payload with no selector.
775            init_factory(vec![changeset(ChangeSetKind::Full, "no-basis-flag", None)]),
776            // The second delivers a delta carrying a selector, so it merges over the first.
777            init_factory(vec![changeset(
778                ChangeSetKind::Partial,
779                "merged-flag",
780                Some("sel-2".into()),
781            )]),
782        ];
783
784        // The synchronizer only records the selector it is started with.
785        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
786        let source_manager = SourceManager::new(vec![sync_factory(
787            vec![],
788            selectors_seen.clone(),
789            Some(shutdown_tx),
790            /* is_fdv1_fallback = */ false,
791        )]);
792
793        run(
794            initializer_factories,
795            source_manager,
796            store.clone(),
797            init_complete,
798            shutdown_rx,
799            FALLBACK_TIMEOUT,
800            RECOVERY_TIMEOUT,
801        )
802        .await;
803
804        // The selector-less payload survives and the delta merged over it.
805        assert!(store.read().flag("no-basis-flag").is_some());
806        assert!(store.read().flag("merged-flag").is_some());
807
808        // The selector-less payload did not stop the initializers, so the synchronizer
809        // starts from the second initializer's selector.
810        assert_eq!(selectors_seen.lock().unwrap()[0], Some("sel-2".into()));
811
812        // init_complete fired exactly once.
813        assert_eq!(*calls.lock().unwrap(), vec![true]);
814    }
815
816    #[tokio::test]
817    async fn initializer_selectorless_full_defers_signal() {
818        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
819
820        // Capture whether the second initializer's flag is present when the signal fires.
821        let flag_at_signal = Arc::new(Mutex::new(None));
822        let probe = store.clone();
823        let sink = flag_at_signal.clone();
824        let init_complete: Arc<dyn Fn(bool) + Send + Sync> = Arc::new(move |_| {
825            *sink.lock().unwrap() = Some(probe.read().flag("from-second").is_some());
826        });
827
828        // Neither initializer produces a basis.
829        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![
830            init_factory(vec![changeset(ChangeSetKind::Full, "from-first", None)]),
831            init_factory(vec![changeset(ChangeSetKind::Full, "from-second", None)]),
832        ];
833
834        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
835        let source_manager = SourceManager::new(vec![sync_factory(
836            vec![],
837            no_selectors(),
838            Some(shutdown_tx),
839            /* is_fdv1_fallback = */ false,
840        )]);
841
842        run(
843            initializer_factories,
844            source_manager,
845            store.clone(),
846            init_complete,
847            shutdown_rx,
848            FALLBACK_TIMEOUT,
849            RECOVERY_TIMEOUT,
850        )
851        .await;
852
853        // The signal waited until the initializers were exhausted, so the second
854        // initializer's payload was already applied when it fired.
855        assert_eq!(*flag_at_signal.lock().unwrap(), Some(true));
856    }
857
858    #[tokio::test]
859    async fn basis_stops_later_initializers() {
860        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
861        let selectors_seen: Selectors = Arc::new(Mutex::new(Vec::new()));
862        let (init_complete, calls) = recording_init_complete();
863
864        // The first initializer yields a basis; the second would apply another flag.
865        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![
866            init_factory(vec![changeset(
867                ChangeSetKind::Full,
868                "from-first",
869                Some("s1".into()),
870            )]),
871            init_factory(vec![changeset(
872                ChangeSetKind::Full,
873                "from-second",
874                Some("s2".into()),
875            )]),
876        ];
877
878        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
879        let source_manager = SourceManager::new(vec![sync_factory(
880            vec![],
881            selectors_seen.clone(),
882            Some(shutdown_tx),
883            /* is_fdv1_fallback = */ false,
884        )]);
885
886        run(
887            initializer_factories,
888            source_manager,
889            store.clone(),
890            init_complete,
891            shutdown_rx,
892            FALLBACK_TIMEOUT,
893            RECOVERY_TIMEOUT,
894        )
895        .await;
896
897        // The basis ended the initializer phase, so the second initializer never ran.
898        assert!(store.read().flag("from-first").is_some());
899        assert!(store.read().flag("from-second").is_none());
900        assert_eq!(*calls.lock().unwrap(), vec![true]);
901
902        // The synchronizer starts from the basis selector.
903        assert_eq!(selectors_seen.lock().unwrap()[0], Some("s1".into()));
904    }
905
906    #[tokio::test]
907    async fn none_changeset_does_not_clobber_the_selector() {
908        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
909        let selectors_seen: Selectors = Arc::new(Mutex::new(Vec::new()));
910        let (init_complete, _calls) = recording_init_complete();
911        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
912
913        // A full basis carrying selector "s1", then a "no changes" (None) changeset.
914        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
915        let source_manager = SourceManager::new(vec![sync_factory(
916            vec![
917                changeset(ChangeSetKind::Full, "flag", Some("s1".into())),
918                changeset(ChangeSetKind::None, "flag", None),
919            ],
920            selectors_seen.clone(),
921            Some(shutdown_tx),
922            /* is_fdv1_fallback = */ false,
923        )]);
924
925        run(
926            initializer_factories,
927            source_manager,
928            store,
929            init_complete,
930            shutdown_rx,
931            FALLBACK_TIMEOUT,
932            RECOVERY_TIMEOUT,
933        )
934        .await;
935
936        // Exactly three requests, and the one issued after the None changeset still
937        // carries "s1". The None must not reset the selector to None.
938        let seen = selectors_seen.lock().unwrap();
939        assert_eq!(*seen, vec![None, Some("s1".into()), Some("s1".into())]);
940    }
941
942    #[tokio::test]
943    async fn selectorless_change_clears_the_selector() {
944        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
945        let selectors_seen: Selectors = Arc::new(Mutex::new(Vec::new()));
946        let (init_complete, _calls) = recording_init_complete();
947        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
948
949        // A full payload carrying selector "s1", then a partial change with no selector.
950        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
951        let source_manager = SourceManager::new(vec![sync_factory(
952            vec![
953                changeset(ChangeSetKind::Full, "flag", Some("s1".into())),
954                changeset(ChangeSetKind::Partial, "flag", None),
955            ],
956            selectors_seen.clone(),
957            Some(shutdown_tx),
958            /* is_fdv1_fallback = */ false,
959        )]);
960
961        run(
962            initializer_factories,
963            source_manager,
964            store,
965            init_complete,
966            shutdown_rx,
967            FALLBACK_TIMEOUT,
968            RECOVERY_TIMEOUT,
969        )
970        .await;
971
972        // The selectorless change advanced state, so the stale selector is dropped.
973        let seen = selectors_seen.lock().unwrap();
974        assert_eq!(*seen, vec![None, Some("s1".into()), None]);
975    }
976
977    #[tokio::test]
978    async fn initializer_delta_does_not_initialize() {
979        // A spec-compliant backend never returns a delta to an initializer.
980        // Here the sole initializer returns one anyway, which is not a full payload.
981        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
982        let (init_complete, calls) = recording_init_complete();
983        let initializer_factories: Vec<Arc<dyn InitializerFactory>> =
984            vec![init_factory(vec![changeset(
985                ChangeSetKind::Partial,
986                "delta-flag",
987                Some("s1".into()),
988            )])];
989
990        // The synchronizer then fails terminally, exhausting every source.
991        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
992        let source_manager = SourceManager::new(vec![sync_factory(
993            vec![terminal()],
994            no_selectors(),
995            Some(shutdown_tx),
996            /* is_fdv1_fallback = */ false,
997        )]);
998
999        run(
1000            initializer_factories,
1001            source_manager,
1002            store.clone(),
1003            init_complete,
1004            shutdown_rx,
1005            FALLBACK_TIMEOUT,
1006            RECOVERY_TIMEOUT,
1007        )
1008        .await;
1009
1010        // The delta applied, but a delta is not a full payload, so failure is signaled.
1011        assert!(store.read().flag("delta-flag").is_some());
1012        assert_eq!(*calls.lock().unwrap(), vec![false]);
1013    }
1014
1015    #[tokio::test]
1016    async fn synchronizer_delta_does_not_initialize() {
1017        // No initializers, so the store starts uninitialized.
1018        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1019        let (init_complete, calls) = recording_init_complete();
1020        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
1021
1022        // A spec-compliant backend never sends a delta before a full payload.
1023        // Here the synchronizer delivers one first anyway, then fails terminally.
1024        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1025        let source_manager = SourceManager::new(vec![sync_factory(
1026            vec![
1027                changeset(ChangeSetKind::Partial, "delta-flag", Some("s1".into())),
1028                terminal(),
1029            ],
1030            no_selectors(),
1031            Some(shutdown_tx),
1032            /* is_fdv1_fallback = */ false,
1033        )]);
1034
1035        run(
1036            initializer_factories,
1037            source_manager,
1038            store.clone(),
1039            init_complete,
1040            shutdown_rx,
1041            FALLBACK_TIMEOUT,
1042            RECOVERY_TIMEOUT,
1043        )
1044        .await;
1045
1046        // The delta applied, but a delta is not a full payload, so failure is signaled.
1047        assert!(store.read().flag("delta-flag").is_some());
1048        assert_eq!(*calls.lock().unwrap(), vec![false]);
1049    }
1050
1051    #[tokio::test]
1052    async fn synchronizer_selectorless_full_initializes() {
1053        // No initializers, so the synchronizer provides the full payload.
1054        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1055        let (init_complete, calls) = recording_init_complete();
1056        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
1057
1058        // A full payload with no selector still initializes, as the FDv1 adapter produces.
1059        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1060        let source_manager = SourceManager::new(vec![sync_factory(
1061            vec![changeset(ChangeSetKind::Full, "full-flag", None)],
1062            no_selectors(),
1063            Some(shutdown_tx),
1064            /* is_fdv1_fallback = */ false,
1065        )]);
1066
1067        run(
1068            initializer_factories,
1069            source_manager,
1070            store.clone(),
1071            init_complete,
1072            shutdown_rx,
1073            FALLBACK_TIMEOUT,
1074            RECOVERY_TIMEOUT,
1075        )
1076        .await;
1077
1078        // A selectorless full still initializes.
1079        assert_eq!(*calls.lock().unwrap(), vec![true]);
1080        assert!(store.read().flag("full-flag").is_some());
1081    }
1082
1083    #[tokio::test]
1084    async fn failed_initializers_let_synchronizer_provide_the_basis() {
1085        // Store and init-complete recorder.
1086        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1087        let (init_complete, calls) = recording_init_complete();
1088
1089        // Both initializers fail without producing a basis.
1090        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![
1091            init_factory(vec![interrupted()]),
1092            init_factory(vec![terminal()]),
1093        ];
1094
1095        // The synchronizer then delivers the basis.
1096        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1097        let source_manager = SourceManager::new(vec![sync_factory(
1098            vec![changeset(
1099                ChangeSetKind::Full,
1100                "sync-flag",
1101                Some("s".into()),
1102            )],
1103            no_selectors(),
1104            Some(shutdown_tx),
1105            /* is_fdv1_fallback = */ false,
1106        )]);
1107
1108        run(
1109            initializer_factories,
1110            source_manager,
1111            store.clone(),
1112            init_complete,
1113            shutdown_rx,
1114            FALLBACK_TIMEOUT,
1115            RECOVERY_TIMEOUT,
1116        )
1117        .await;
1118
1119        // Initialization succeeded via the synchronizer.
1120        assert_eq!(*calls.lock().unwrap(), vec![true]);
1121        assert!(store.read().flag("sync-flag").is_some());
1122    }
1123
1124    #[tokio::test]
1125    async fn exhausting_all_sources_signals_failure_once() {
1126        // Store and init-complete recorder.
1127        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1128        let (init_complete, calls) = recording_init_complete();
1129
1130        // The initializer fails.
1131        let initializer_factories: Vec<Arc<dyn InitializerFactory>> =
1132            vec![init_factory(vec![terminal()])];
1133
1134        // The synchronizer reports an interruption, then fails terminally.
1135        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1136        let source_manager = SourceManager::new(vec![sync_factory(
1137            vec![interrupted(), terminal()],
1138            no_selectors(),
1139            Some(shutdown_tx),
1140            /* is_fdv1_fallback = */ false,
1141        )]);
1142
1143        run(
1144            initializer_factories,
1145            source_manager,
1146            store,
1147            init_complete,
1148            shutdown_rx,
1149            FALLBACK_TIMEOUT,
1150            RECOVERY_TIMEOUT,
1151        )
1152        .await;
1153
1154        // Failure was reported exactly once.
1155        assert_eq!(*calls.lock().unwrap(), vec![false]);
1156    }
1157
1158    #[tokio::test]
1159    async fn synchronizer_terminal_error_advances_to_next() {
1160        // Store and init-complete recorder; no initializers.
1161        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1162        let (init_complete, calls) = recording_init_complete();
1163        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
1164
1165        // First synchronizer fails terminally; the second provides the basis.
1166        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1167        let source_manager = SourceManager::new(vec![
1168            sync_factory(
1169                vec![terminal()],
1170                no_selectors(),
1171                Some(shutdown_tx.clone()),
1172                /* is_fdv1_fallback = */ false,
1173            ),
1174            sync_factory(
1175                vec![changeset(
1176                    ChangeSetKind::Full,
1177                    "from-second",
1178                    Some("s".into()),
1179                )],
1180                no_selectors(),
1181                Some(shutdown_tx),
1182                /* is_fdv1_fallback = */ false,
1183            ),
1184        ]);
1185
1186        run(
1187            initializer_factories,
1188            source_manager,
1189            store.clone(),
1190            init_complete,
1191            shutdown_rx,
1192            FALLBACK_TIMEOUT,
1193            RECOVERY_TIMEOUT,
1194        )
1195        .await;
1196
1197        // The run advanced past the terminal source and applied the second's basis.
1198        assert_eq!(*calls.lock().unwrap(), vec![true]);
1199        assert!(store.read().flag("from-second").is_some());
1200    }
1201
1202    #[tokio::test]
1203    async fn synchronizer_interrupted_retries_same_source() {
1204        // Store and init-complete recorder; no initializers.
1205        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1206        let (init_complete, calls) = recording_init_complete();
1207        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
1208
1209        // A single synchronizer: an interruption, then a basis on the retry.
1210        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1211        let source_manager = SourceManager::new(vec![sync_factory(
1212            vec![
1213                interrupted(),
1214                changeset(ChangeSetKind::Full, "after-retry", Some("s".into())),
1215            ],
1216            no_selectors(),
1217            Some(shutdown_tx),
1218            /* is_fdv1_fallback = */ false,
1219        )]);
1220
1221        run(
1222            initializer_factories,
1223            source_manager,
1224            store.clone(),
1225            init_complete,
1226            shutdown_rx,
1227            FALLBACK_TIMEOUT,
1228            RECOVERY_TIMEOUT,
1229        )
1230        .await;
1231
1232        // The retry on the same synchronizer delivered the basis.
1233        assert_eq!(*calls.lock().unwrap(), vec![true]);
1234        assert!(store.read().flag("after-retry").is_some());
1235    }
1236
1237    #[tokio::test]
1238    async fn shutdown_ends_run_without_signaling() {
1239        // Store and init-complete recorder; no initializers.
1240        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1241        let (init_complete, calls) = recording_init_complete();
1242        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
1243
1244        // A synchronizer whose next() never resolves.
1245        let source_manager = SourceManager::new(vec![sync_factory(
1246            vec![],
1247            no_selectors(),
1248            /* idle */ None,
1249            /* is_fdv1_fallback = */ false,
1250        )]);
1251        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1252
1253        // Drive the run on a task, then signal shutdown.
1254        let handle = tokio::spawn(run(
1255            initializer_factories,
1256            source_manager,
1257            store,
1258            init_complete,
1259            shutdown_rx,
1260            FALLBACK_TIMEOUT,
1261            RECOVERY_TIMEOUT,
1262        ));
1263        shutdown_tx.send(()).unwrap();
1264        handle.await.unwrap();
1265
1266        // The run returned before any initialization signal.
1267        assert!(calls.lock().unwrap().is_empty());
1268    }
1269
1270    #[test]
1271    fn rotates_cyclically_skips_blocked_and_exhausts() {
1272        let mut sources = SourceManager::new(vec![
1273            sync_factory(
1274                vec![],
1275                no_selectors(),
1276                /* idle */ None,
1277                /* is_fdv1_fallback = */ false,
1278            ),
1279            sync_factory(
1280                vec![],
1281                no_selectors(),
1282                /* idle */ None,
1283                /* is_fdv1_fallback = */ false,
1284            ),
1285            sync_factory(
1286                vec![],
1287                no_selectors(),
1288                /* idle */ None,
1289                /* is_fdv1_fallback = */ false,
1290            ),
1291        ]);
1292
1293        // Advances cyclically from the prime.
1294        sources.next_synchronizer();
1295        assert_eq!(sources.current_factory_index, Some(0));
1296        sources.next_synchronizer();
1297        assert_eq!(sources.current_factory_index, Some(1));
1298
1299        // Blocking the active factory drops it from the rotation.
1300        sources.block_current();
1301        sources.next_synchronizer();
1302        assert_eq!(sources.current_factory_index, Some(2));
1303
1304        // The next pass wraps around and skips the blocked factory.
1305        sources.next_synchronizer();
1306        assert_eq!(sources.current_factory_index, Some(0));
1307
1308        // With every factory blocked there is nothing left to return.
1309        sources.block_current();
1310        sources.next_synchronizer();
1311        sources.block_current();
1312        assert!(sources.next_synchronizer().is_none());
1313    }
1314
1315    #[test]
1316    fn reset_source_index_returns_to_prime() {
1317        let mut sources = SourceManager::new(vec![
1318            sync_factory(
1319                vec![],
1320                no_selectors(),
1321                /* idle */ None,
1322                /* is_fdv1_fallback = */ false,
1323            ),
1324            sync_factory(
1325                vec![],
1326                no_selectors(),
1327                /* idle */ None,
1328                /* is_fdv1_fallback = */ false,
1329            ),
1330        ]);
1331
1332        sources.next_synchronizer();
1333        sources.next_synchronizer();
1334        assert_eq!(sources.current_factory_index, Some(1));
1335
1336        // Recovery rewinds so the next search starts from the prime.
1337        sources.reset_source_index();
1338        sources.next_synchronizer();
1339        assert_eq!(sources.current_factory_index, Some(0));
1340    }
1341
1342    #[test]
1343    fn is_prime_and_available_count_track_state() {
1344        let mut sources = SourceManager::new(vec![
1345            sync_factory(
1346                vec![],
1347                no_selectors(),
1348                /* idle */ None,
1349                /* is_fdv1_fallback = */ false,
1350            ),
1351            sync_factory(
1352                vec![],
1353                no_selectors(),
1354                /* idle */ None,
1355                /* is_fdv1_fallback = */ false,
1356            ),
1357        ]);
1358        assert_eq!(sources.available_count(), 2);
1359
1360        sources.next_synchronizer();
1361        assert!(sources.is_prime());
1362        sources.next_synchronizer();
1363        assert!(!sources.is_prime());
1364
1365        sources.block_current();
1366        assert_eq!(sources.available_count(), 1);
1367    }
1368
1369    #[test]
1370    fn fdv1_fallback_factory_starts_blocked() {
1371        let mut sources = SourceManager::new(vec![
1372            sync_factory(
1373                vec![],
1374                no_selectors(),
1375                /* idle */ None,
1376                /* is_fdv1_fallback = */ false,
1377            ),
1378            sync_factory(
1379                vec![],
1380                no_selectors(),
1381                /* idle */ None,
1382                /* is_fdv1_fallback = */ true,
1383            ),
1384        ]);
1385
1386        // Only the FDv2 synchronizer is available; the FDv1 fallback is dormant.
1387        assert_eq!(sources.available_count(), 1);
1388        sources.next_synchronizer();
1389        assert!(!sources.is_current_fdv1_fallback());
1390    }
1391
1392    #[test]
1393    fn switch_to_and_back_from_fdv1_fallback() {
1394        let mut sources = SourceManager::new(vec![
1395            sync_factory(
1396                vec![],
1397                no_selectors(),
1398                /* idle */ None,
1399                /* is_fdv1_fallback = */ false,
1400            ),
1401            sync_factory(
1402                vec![],
1403                no_selectors(),
1404                /* idle */ None,
1405                /* is_fdv1_fallback = */ true,
1406            ),
1407        ]);
1408
1409        // The directive activates only the FDv1 fallback.
1410        sources.switch_to_fdv1_fallback();
1411        assert_eq!(sources.available_count(), 1);
1412        sources.next_synchronizer();
1413        assert!(sources.is_current_fdv1_fallback());
1414
1415        // Switching back restores the FDv2 synchronizer and re-blocks FDv1.
1416        sources.switch_back_to_fdv2();
1417        assert_eq!(sources.available_count(), 1);
1418        sources.next_synchronizer();
1419        assert!(!sources.is_current_fdv1_fallback());
1420    }
1421
1422    #[test]
1423    fn switch_back_unblocks_terminally_blocked_fdv2() {
1424        let mut sources = SourceManager::new(vec![
1425            sync_factory(
1426                vec![],
1427                no_selectors(),
1428                /* idle */ None,
1429                /* is_fdv1_fallback = */ false,
1430            ),
1431            sync_factory(
1432                vec![],
1433                no_selectors(),
1434                /* idle */ None,
1435                /* is_fdv1_fallback = */ true,
1436            ),
1437        ]);
1438
1439        // A terminal error blocks the only FDv2 synchronizer.
1440        sources.next_synchronizer();
1441        sources.block_current();
1442        assert_eq!(sources.available_count(), 0);
1443
1444        // A fallback round-trip makes the FDv2 synchronizer usable again.
1445        sources.switch_to_fdv1_fallback();
1446        sources.switch_back_to_fdv2();
1447        assert_eq!(sources.available_count(), 1);
1448    }
1449
1450    #[tokio::test(start_paused = true)]
1451    async fn fallback_fires_after_sustained_interruption() {
1452        // No initializers; the prime interrupts then idles.
1453        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1454        let (init_complete, calls) = recording_init_complete();
1455        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
1456
1457        // Prime only ever interrupts; the fallback stands by with a basis.
1458        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1459        let source_manager = SourceManager::new(vec![
1460            sync_factory(
1461                vec![interrupted()],
1462                no_selectors(),
1463                /* idle */ None,
1464                /* is_fdv1_fallback = */ false,
1465            ),
1466            sync_factory(
1467                vec![changeset(
1468                    ChangeSetKind::Full,
1469                    "from-fallback",
1470                    Some("s".into()),
1471                )],
1472                no_selectors(),
1473                Some(shutdown_tx),
1474                /* is_fdv1_fallback = */ false,
1475            ),
1476        ]);
1477
1478        // Paused time auto-advances past the fallback timeout while the prime idles.
1479        let handle = tokio::spawn(run(
1480            initializer_factories,
1481            source_manager,
1482            store.clone(),
1483            init_complete,
1484            shutdown_rx,
1485            FALLBACK_TIMEOUT,
1486            RECOVERY_TIMEOUT,
1487        ));
1488        handle.await.unwrap();
1489
1490        // The run fell back to the second synchronizer and applied its basis.
1491        assert_eq!(*calls.lock().unwrap(), vec![true]);
1492        assert!(store.read().flag("from-fallback").is_some());
1493    }
1494
1495    #[tokio::test(start_paused = true)]
1496    async fn changeset_cancels_the_fallback_timer() {
1497        // Notify on init so the test can act once the prime's basis lands.
1498        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1499        let calls: InitCalls = Arc::new(Mutex::new(Vec::new()));
1500        let notify = Arc::new(tokio::sync::Notify::new());
1501        let sink = calls.clone();
1502        let waker = notify.clone();
1503        let init_complete: Arc<dyn Fn(bool) + Send + Sync> = Arc::new(move |success| {
1504            sink.lock().unwrap().push(success);
1505            waker.notify_one();
1506        });
1507        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
1508
1509        // Prime interrupts (arming the timer), delivers a basis (clearing it), then idles.
1510        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1511        let source_manager = SourceManager::new(vec![
1512            sync_factory(
1513                vec![
1514                    interrupted(),
1515                    changeset(ChangeSetKind::Full, "from-prime", Some("s".into())),
1516                ],
1517                no_selectors(),
1518                /* idle */ None,
1519                /* is_fdv1_fallback = */ false,
1520            ),
1521            sync_factory(
1522                vec![changeset(
1523                    ChangeSetKind::Full,
1524                    "from-fallback",
1525                    Some("s".into()),
1526                )],
1527                no_selectors(),
1528                Some(shutdown_tx.clone()),
1529                /* is_fdv1_fallback = */ false,
1530            ),
1531        ]);
1532
1533        let handle = tokio::spawn(run(
1534            initializer_factories,
1535            source_manager,
1536            store.clone(),
1537            init_complete,
1538            shutdown_rx,
1539            FALLBACK_TIMEOUT,
1540            RECOVERY_TIMEOUT,
1541        ));
1542
1543        // Wait for the basis, let well past the fallback timeout elapse, then stop.
1544        notify.notified().await;
1545        tokio::time::advance(FALLBACK_TIMEOUT * 2).await;
1546        shutdown_tx.send(()).unwrap();
1547        handle.await.unwrap();
1548
1549        // The basis kept the run on the prime; it never fell back.
1550        assert!(store.read().flag("from-prime").is_some());
1551        assert!(store.read().flag("from-fallback").is_none());
1552    }
1553
1554    #[tokio::test(start_paused = true)]
1555    async fn recovery_fires_and_returns_to_the_prime() {
1556        // No initializers; the prime is down, so the run falls back then recovers.
1557        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1558        let (init_complete, calls) = recording_init_complete();
1559        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
1560
1561        // Prime recovers on rebuild; the fallback supplies a basis then idles.
1562        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1563        let source_manager = SourceManager::new(vec![
1564            Arc::new(DownThenDataFactory {
1565                builds: AtomicUsize::new(0),
1566                shutdown: shutdown_tx,
1567            }) as Arc<dyn SynchronizerFactory>,
1568            sync_factory(
1569                vec![changeset(
1570                    ChangeSetKind::Full,
1571                    "from-fallback",
1572                    Some("s".into()),
1573                )],
1574                no_selectors(),
1575                /* idle */ None,
1576                /* is_fdv1_fallback = */ false,
1577            ),
1578        ]);
1579
1580        // Paused time auto-advances through the fallback then recovery timeouts.
1581        let handle = tokio::spawn(run(
1582            initializer_factories,
1583            source_manager,
1584            store.clone(),
1585            init_complete,
1586            shutdown_rx,
1587            FALLBACK_TIMEOUT,
1588            RECOVERY_TIMEOUT,
1589        ));
1590        handle.await.unwrap();
1591
1592        // Fell back to the fallback's basis, then recovered to the prime's delta.
1593        assert_eq!(*calls.lock().unwrap(), vec![true]);
1594        assert!(store.read().flag("from-fallback").is_some());
1595        assert!(store.read().flag("prime-recovered").is_some());
1596    }
1597
1598    #[tokio::test]
1599    async fn fallback_directive_switches_to_fdv1() {
1600        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1601        let (init_complete, calls) = recording_init_complete();
1602        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
1603
1604        // The FDv2 source reports a fallback directive; the FDv1 fallback then
1605        // supplies the basis and ends the run.
1606        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1607        let source_manager = SourceManager::new(vec![
1608            fallback_directive_factory(Duration::from_secs(60)),
1609            sync_factory(
1610                vec![changeset(
1611                    ChangeSetKind::Full,
1612                    "from-fdv1",
1613                    Some("s".into()),
1614                )],
1615                no_selectors(),
1616                Some(shutdown_tx),
1617                /* is_fdv1_fallback = */ true,
1618            ),
1619        ]);
1620
1621        run(
1622            initializer_factories,
1623            source_manager,
1624            store.clone(),
1625            init_complete,
1626            shutdown_rx,
1627            FALLBACK_TIMEOUT,
1628            RECOVERY_TIMEOUT,
1629        )
1630        .await;
1631
1632        // Switched to the FDv1 fallback and applied its basis.
1633        assert_eq!(*calls.lock().unwrap(), vec![true]);
1634        assert!(store.read().flag("from-fdv1").is_some());
1635    }
1636
1637    #[tokio::test(start_paused = true)]
1638    async fn fdv2_retry_reengages_after_ttl() {
1639        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1640        let (init_complete, calls) = recording_init_complete();
1641        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
1642
1643        // The FDv2 source reports a directive with a 60s TTL; the FDv1 fallback
1644        // supplies a basis and then idles until FDv2 is retried.
1645        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1646        let source_manager = SourceManager::new(vec![
1647            Arc::new(FallbackThenDataFactory {
1648                ttl: Duration::from_secs(60),
1649                reengage_kind: ChangeSetKind::Partial,
1650                builds: AtomicUsize::new(0),
1651                shutdown: shutdown_tx,
1652            }) as Arc<dyn SynchronizerFactory>,
1653            sync_factory(
1654                vec![changeset(
1655                    ChangeSetKind::Full,
1656                    "from-fdv1",
1657                    Some("s".into()),
1658                )],
1659                no_selectors(),
1660                /* idle */ None,
1661                /* is_fdv1_fallback = */ true,
1662            ),
1663        ]);
1664
1665        // Paused time auto-advances past the TTL, re-engaging FDv2.
1666        let handle = tokio::spawn(run(
1667            initializer_factories,
1668            source_manager,
1669            store.clone(),
1670            init_complete,
1671            shutdown_rx,
1672            FALLBACK_TIMEOUT,
1673            RECOVERY_TIMEOUT,
1674        ));
1675        handle.await.unwrap();
1676
1677        // Fell back to FDv1, then re-engaged FDv2 once the TTL expired.
1678        assert_eq!(*calls.lock().unwrap(), vec![true]);
1679        assert!(store.read().flag("from-fdv1").is_some());
1680        assert!(store.read().flag("fdv2-back").is_some());
1681    }
1682
1683    #[tokio::test(start_paused = true)]
1684    async fn fdv2_retry_survives_a_terminal_fdv1_fallback() {
1685        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1686        let (init_complete, calls) = recording_init_complete();
1687        let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
1688
1689        // The FDv2 source reports a 60s directive; the FDv1 fallback then fails
1690        // terminally, blocking every source before the TTL expires.
1691        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1692        let source_manager = SourceManager::new(vec![
1693            Arc::new(FallbackThenDataFactory {
1694                ttl: Duration::from_secs(60),
1695                reengage_kind: ChangeSetKind::Full,
1696                builds: AtomicUsize::new(0),
1697                shutdown: shutdown_tx.clone(),
1698            }) as Arc<dyn SynchronizerFactory>,
1699            sync_factory(
1700                vec![terminal()],
1701                no_selectors(),
1702                Some(shutdown_tx),
1703                /* is_fdv1_fallback = */ true,
1704            ),
1705        ]);
1706
1707        // Paused time advances past the TTL while no source is available.
1708        let handle = tokio::spawn(run(
1709            initializer_factories,
1710            source_manager,
1711            store.clone(),
1712            init_complete,
1713            shutdown_rx,
1714            FALLBACK_TIMEOUT,
1715            RECOVERY_TIMEOUT,
1716        ));
1717        handle.await.unwrap();
1718
1719        // The blocked fallback did not strand the run; FDv2 re-engaged after the TTL.
1720        assert_eq!(*calls.lock().unwrap(), vec![true]);
1721        assert!(store.read().flag("fdv2-back").is_some());
1722    }
1723
1724    #[tokio::test]
1725    async fn initializer_fallback_directive_switches_to_fdv1() {
1726        let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1727        let (init_complete, calls) = recording_init_complete();
1728        let initializer_factories: Vec<Arc<dyn InitializerFactory>> =
1729            vec![Arc::new(FallbackInitializerFactory {
1730                ttl: Duration::from_secs(60),
1731            })];
1732
1733        // The initializer's directive switches to the FDv1 fallback, which supplies
1734        // the basis and ends the run.
1735        let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1736        let source_manager = SourceManager::new(vec![
1737            sync_factory(
1738                vec![],
1739                no_selectors(),
1740                Some(shutdown_tx.clone()),
1741                /* is_fdv1_fallback = */ false,
1742            ),
1743            sync_factory(
1744                vec![changeset(
1745                    ChangeSetKind::Full,
1746                    "from-fdv1",
1747                    Some("s".into()),
1748                )],
1749                no_selectors(),
1750                Some(shutdown_tx),
1751                /* is_fdv1_fallback = */ true,
1752            ),
1753        ]);
1754
1755        run(
1756            initializer_factories,
1757            source_manager,
1758            store.clone(),
1759            init_complete,
1760            shutdown_rx,
1761            FALLBACK_TIMEOUT,
1762            RECOVERY_TIMEOUT,
1763        )
1764        .await;
1765
1766        // The initializer's directive switched straight to the FDv1 fallback.
1767        assert_eq!(*calls.lock().unwrap(), vec![true]);
1768        assert!(store.read().flag("from-fdv1").is_some());
1769    }
1770}