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
15pub trait InitializerFactory: Send + Sync {
17 fn create(&self) -> Box<dyn Initializer>;
19}
20
21pub trait SynchronizerFactory: Send + Sync {
23 fn create(&self) -> Box<dyn Synchronizer>;
25
26 fn is_fdv1_fallback(&self) -> bool {
28 false
29 }
30}
31
32pub(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#[derive(Clone, Copy, PartialEq, Eq)]
87enum SourceState {
88 Available,
89 Blocked,
90}
91
92struct SourceManager {
94 factories: Vec<Arc<dyn SynchronizerFactory>>,
95 states: Vec<SourceState>,
96 synchronizer_index: Option<usize>,
98 current_factory_index: Option<usize>,
100}
101
102impl SourceManager {
103 fn new(factories: Vec<Arc<dyn SynchronizerFactory>>) -> Self {
104 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 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 fn block_current(&mut self) {
146 if let Some(i) = self.current_factory_index {
147 self.states[i] = SourceState::Blocked;
148 }
149 }
150
151 fn reset_source_index(&mut self) {
153 self.synchronizer_index = None;
154 }
155
156 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 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 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 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
203async 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 let mut got_full = false;
224 let mut fdv2_retry_at: Option<Instant> = None;
225
226 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 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 let mut current = source_manager.next_synchronizer();
274 loop {
275 let mut active = match current {
276 Some(active) => active,
277 None => {
278 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 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 _ = fallback => break,
317 _ = recovery => {
319 source_manager.reset_source_index();
320 break;
321 }
322 _ = 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 fallback_at = None;
344 interrupted_logged = false;
345 }
346 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 FDv2SourceResult::Goodbye => {}
358 FDv2SourceResult::TerminalError(error) => {
359 warn!("{name} terminal error: {}", error.message);
360 terminal = true;
361 }
362 }
363 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 source_manager.block_current();
375 break;
376 }
377 }
378 }
379 }
380
381 current = source_manager.next_synchronizer();
382 }
383
384 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 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 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 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 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 fn no_selectors() -> Selectors {
574 Arc::new(Mutex::new(Vec::new()))
575 }
576
577 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 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 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 (vec![interrupted()], None)
615 } else {
616 (
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 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 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 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 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 None,
690 false,
691 )],
692 FALLBACK_TIMEOUT,
693 RECOVERY_TIMEOUT,
694 );
695
696 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 system.start(init_complete, shutdown_rx);
709 notify.notified().await;
710
711 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 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 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 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 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 assert_eq!(*calls.lock().unwrap(), vec![true]);
758
759 assert!(store.read().flag("init-flag").is_some());
761 assert!(store.read().flag("sync-flag").is_some());
762
763 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 init_factory(vec![changeset(ChangeSetKind::Full, "no-basis-flag", None)]),
776 init_factory(vec![changeset(
778 ChangeSetKind::Partial,
779 "merged-flag",
780 Some("sel-2".into()),
781 )]),
782 ];
783
784 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 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 assert!(store.read().flag("no-basis-flag").is_some());
806 assert!(store.read().flag("merged-flag").is_some());
807
808 assert_eq!(selectors_seen.lock().unwrap()[0], Some("sel-2".into()));
811
812 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1087 let (init_complete, calls) = recording_init_complete();
1088
1089 let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![
1091 init_factory(vec![interrupted()]),
1092 init_factory(vec![terminal()]),
1093 ];
1094
1095 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 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 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 let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
1128 let (init_complete, calls) = recording_init_complete();
1129
1130 let initializer_factories: Vec<Arc<dyn InitializerFactory>> =
1132 vec![init_factory(vec![terminal()])];
1133
1134 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 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 assert_eq!(*calls.lock().unwrap(), vec![false]);
1156 }
1157
1158 #[tokio::test]
1159 async fn synchronizer_terminal_error_advances_to_next() {
1160 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 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 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 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 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 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 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 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 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 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 let source_manager = SourceManager::new(vec![sync_factory(
1246 vec![],
1247 no_selectors(),
1248 None,
1249 false,
1250 )]);
1251 let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
1252
1253 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 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 None,
1277 false,
1278 ),
1279 sync_factory(
1280 vec![],
1281 no_selectors(),
1282 None,
1283 false,
1284 ),
1285 sync_factory(
1286 vec![],
1287 no_selectors(),
1288 None,
1289 false,
1290 ),
1291 ]);
1292
1293 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 sources.block_current();
1301 sources.next_synchronizer();
1302 assert_eq!(sources.current_factory_index, Some(2));
1303
1304 sources.next_synchronizer();
1306 assert_eq!(sources.current_factory_index, Some(0));
1307
1308 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 None,
1322 false,
1323 ),
1324 sync_factory(
1325 vec![],
1326 no_selectors(),
1327 None,
1328 false,
1329 ),
1330 ]);
1331
1332 sources.next_synchronizer();
1333 sources.next_synchronizer();
1334 assert_eq!(sources.current_factory_index, Some(1));
1335
1336 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 None,
1349 false,
1350 ),
1351 sync_factory(
1352 vec![],
1353 no_selectors(),
1354 None,
1355 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 None,
1376 false,
1377 ),
1378 sync_factory(
1379 vec![],
1380 no_selectors(),
1381 None,
1382 true,
1383 ),
1384 ]);
1385
1386 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 None,
1399 false,
1400 ),
1401 sync_factory(
1402 vec![],
1403 no_selectors(),
1404 None,
1405 true,
1406 ),
1407 ]);
1408
1409 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 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 None,
1429 false,
1430 ),
1431 sync_factory(
1432 vec![],
1433 no_selectors(),
1434 None,
1435 true,
1436 ),
1437 ]);
1438
1439 sources.next_synchronizer();
1441 sources.block_current();
1442 assert_eq!(sources.available_count(), 0);
1443
1444 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 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 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 None,
1464 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 false,
1475 ),
1476 ]);
1477
1478 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 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 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 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 None,
1519 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 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 notify.notified().await;
1545 tokio::time::advance(FALLBACK_TIMEOUT * 2).await;
1546 shutdown_tx.send(()).unwrap();
1547 handle.await.unwrap();
1548
1549 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 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 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 None,
1576 false,
1577 ),
1578 ]);
1579
1580 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 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 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 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 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 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 None,
1661 true,
1662 ),
1663 ]);
1664
1665 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 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 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 true,
1704 ),
1705 ]);
1706
1707 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 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 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 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 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 assert_eq!(*calls.lock().unwrap(), vec![true]);
1768 assert!(store.read().flag("from-fdv1").is_some());
1769 }
1770}