use std::sync::Arc;
use std::time::Duration;
use futures::FutureExt;
use parking_lot::RwLock;
use tokio::sync::broadcast;
use tokio::time::{sleep_until, Instant};
use crate::data_system::DataSystem;
use crate::stores::store::{DataStore, InMemoryDataStore, TransactionalDataStore};
use super::model::{ChangeSetKind, Selector};
use super::source::{FDv2SourceEvent, FDv2SourceResult, Initializer, Synchronizer};
pub trait InitializerFactory: Send + Sync {
fn create(&self) -> Box<dyn Initializer>;
}
pub trait SynchronizerFactory: Send + Sync {
fn create(&self) -> Box<dyn Synchronizer>;
fn is_fdv1_fallback(&self) -> bool {
false
}
}
pub(crate) struct FDv2DataSystem {
initializer_factories: Vec<Arc<dyn InitializerFactory>>,
synchronizer_factories: Vec<Arc<dyn SynchronizerFactory>>,
fallback_timeout: Duration,
recovery_timeout: Duration,
store: Arc<RwLock<InMemoryDataStore>>,
}
impl FDv2DataSystem {
pub(crate) fn new(
initializer_factories: Vec<Arc<dyn InitializerFactory>>,
synchronizer_factories: Vec<Arc<dyn SynchronizerFactory>>,
fallback_timeout: Duration,
recovery_timeout: Duration,
) -> Self {
Self {
initializer_factories,
synchronizer_factories,
fallback_timeout,
recovery_timeout,
store: Arc::new(RwLock::new(InMemoryDataStore::new())),
}
}
}
impl DataSystem for FDv2DataSystem {
fn start(
&self,
init_complete: Arc<dyn Fn(bool) + Send + Sync>,
shutdown_receiver: broadcast::Receiver<()>,
) {
let initializer_factories = self.initializer_factories.clone();
let source_manager = SourceManager::new(self.synchronizer_factories.clone());
let store = self.store.clone();
tokio::spawn(run(
initializer_factories,
source_manager,
store,
init_complete,
shutdown_receiver,
self.fallback_timeout,
self.recovery_timeout,
));
}
fn store(&self) -> Arc<RwLock<dyn DataStore>> {
self.store.clone()
}
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum SourceState {
Available,
Blocked,
}
struct SourceManager {
factories: Vec<Arc<dyn SynchronizerFactory>>,
states: Vec<SourceState>,
synchronizer_index: Option<usize>,
current_factory_index: Option<usize>,
}
impl SourceManager {
fn new(factories: Vec<Arc<dyn SynchronizerFactory>>) -> Self {
let states = factories
.iter()
.map(|f| {
if f.is_fdv1_fallback() {
SourceState::Blocked
} else {
SourceState::Available
}
})
.collect();
Self {
factories,
states,
synchronizer_index: None,
current_factory_index: None,
}
}
fn next_synchronizer(&mut self) -> Option<Box<dyn Synchronizer>> {
let n = self.factories.len();
if n == 0 {
self.current_factory_index = None;
return None;
}
let mut i = self.synchronizer_index.map_or(0, |c| (c + 1) % n);
for _ in 0..n {
if self.states[i] == SourceState::Available {
self.synchronizer_index = Some(i);
self.current_factory_index = Some(i);
return Some(self.factories[i].create());
}
i = (i + 1) % n;
}
self.current_factory_index = None;
None
}
fn block_current(&mut self) {
if let Some(i) = self.current_factory_index {
self.states[i] = SourceState::Blocked;
}
}
fn reset_source_index(&mut self) {
self.synchronizer_index = None;
}
fn is_prime(&self) -> bool {
let first = self
.states
.iter()
.position(|s| *s == SourceState::Available);
first == self.current_factory_index && first.is_some()
}
fn available_count(&self) -> usize {
self.states
.iter()
.filter(|s| **s == SourceState::Available)
.count()
}
fn switch_to_fdv1_fallback(&mut self) {
for (i, factory) in self.factories.iter().enumerate() {
self.states[i] = if factory.is_fdv1_fallback() {
SourceState::Available
} else {
SourceState::Blocked
};
}
self.synchronizer_index = None;
}
fn switch_back_to_fdv2(&mut self) {
for (i, factory) in self.factories.iter().enumerate() {
self.states[i] = if factory.is_fdv1_fallback() {
SourceState::Blocked
} else {
SourceState::Available
};
}
self.synchronizer_index = None;
}
fn is_current_fdv1_fallback(&self) -> bool {
self.current_factory_index
.is_some_and(|i| self.factories[i].is_fdv1_fallback())
}
}
async fn deadline(at: Option<Instant>) {
match at {
Some(t) => sleep_until(t).await,
None => std::future::pending::<()>().await,
}
}
async fn run(
initializer_factories: Vec<Arc<dyn InitializerFactory>>,
mut source_manager: SourceManager,
store: Arc<RwLock<InMemoryDataStore>>,
init_complete: Arc<dyn Fn(bool) + Send + Sync>,
mut shutdown_receiver: broadcast::Receiver<()>,
fallback_timeout: Duration,
recovery_timeout: Duration,
) {
let mut selector: Selector = None;
let mut initialized = false;
let mut got_full = false;
let mut fdv2_retry_at: Option<Instant> = None;
for factory in initializer_factories {
let mut initializer = factory.create();
let name = initializer.name().to_string();
let mut shutdown = Box::pin(shutdown_receiver.recv()).fuse();
let event = futures::select! {
_ = shutdown => return,
event = initializer.run().fuse() => event,
};
let FDv2SourceEvent {
result,
fdv1_fallback,
} = event;
let mut has_basis = false;
match result {
FDv2SourceResult::ChangeSet(change_set) => {
let is_full = matches!(change_set.kind, ChangeSetKind::Full);
has_basis = is_full && change_set.selector.is_some();
if !matches!(change_set.kind, ChangeSetKind::None) {
selector = change_set.selector.clone();
}
store.write().apply(change_set);
if is_full {
got_full = true;
}
}
_ => debug!("{name} did not provide a basis"),
}
if let Some(fallback_directive) = fdv1_fallback {
info!("FDv2 falling back to the FDv1 protocol");
source_manager.switch_to_fdv1_fallback();
fdv2_retry_at = Some(Instant::now() + fallback_directive.ttl);
break;
}
if has_basis {
break;
}
}
if got_full && !initialized {
init_complete(true);
initialized = true;
}
let mut current = source_manager.next_synchronizer();
loop {
let mut active = match current {
Some(active) => active,
None => {
if fdv2_retry_at.is_none() {
break;
}
let mut shutdown = Box::pin(shutdown_receiver.recv()).fuse();
let mut fdv2_retry = Box::pin(deadline(fdv2_retry_at)).fuse();
futures::select! {
_ = shutdown => return,
_ = fdv2_retry => {
source_manager.switch_back_to_fdv2();
fdv2_retry_at = None;
}
}
current = source_manager.next_synchronizer();
continue;
}
};
let name = active.name().to_string();
let has_fallback = source_manager.available_count() > 1;
let has_recovery = has_fallback && !source_manager.is_prime();
let mut fallback_at: Option<Instant> = None;
let recovery_at = has_recovery.then(|| Instant::now() + recovery_timeout);
let mut interrupted_logged = false;
loop {
let mut shutdown = Box::pin(shutdown_receiver.recv()).fuse();
let mut fallback = Box::pin(deadline(fallback_at)).fuse();
let mut recovery = Box::pin(deadline(recovery_at)).fuse();
let mut fdv2_retry = Box::pin(deadline(fdv2_retry_at)).fuse();
let mut next = active.next(selector.clone()).fuse();
futures::select! {
_ = shutdown => return,
_ = fallback => break,
_ = recovery => {
source_manager.reset_source_index();
break;
}
_ = fdv2_retry => {
source_manager.switch_back_to_fdv2();
fdv2_retry_at = None;
break;
}
event = next => {
let FDv2SourceEvent { result, fdv1_fallback } = event;
let mut terminal = false;
match result {
FDv2SourceResult::ChangeSet(change_set) => {
let is_full = matches!(change_set.kind, ChangeSetKind::Full);
if !matches!(change_set.kind, ChangeSetKind::None) {
selector = change_set.selector.clone();
store.write().apply(change_set);
if is_full && !initialized {
init_complete(true);
initialized = true;
}
}
fallback_at = None;
interrupted_logged = false;
}
FDv2SourceResult::Interrupted(error) => {
if !interrupted_logged {
info!("{name} interrupted: {}", error.message);
interrupted_logged = true;
}
if has_fallback && fallback_at.is_none() {
fallback_at = Some(Instant::now() + fallback_timeout);
}
}
FDv2SourceResult::Goodbye => {}
FDv2SourceResult::TerminalError(error) => {
warn!("{name} terminal error: {}", error.message);
terminal = true;
}
}
if let Some(fallback_directive) = fdv1_fallback {
if !source_manager.is_current_fdv1_fallback() {
info!("FDv2 falling back to the FDv1 protocol");
source_manager.switch_to_fdv1_fallback();
fdv2_retry_at = Some(Instant::now() + fallback_directive.ttl);
break;
}
}
if terminal {
source_manager.block_current();
break;
}
}
}
}
current = source_manager.next_synchronizer();
}
if !initialized {
init_complete(false);
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::VecDeque;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Mutex;
use futures::future::BoxFuture;
use launchdarkly_server_sdk_evaluation::Store;
use super::super::model::ChangeSetKind;
use super::super::source::{ErrorInfo, ErrorKind, FDv1FallbackDirective, FDv2SourceEvent};
use crate::stores::change_set::{ChangeSet, ItemChange};
use crate::stores::store_types::StorageItem;
use crate::test_common::basic_flag;
const FALLBACK_TIMEOUT: Duration = Duration::from_secs(120);
const RECOVERY_TIMEOUT: Duration = Duration::from_secs(300);
type Selectors = Arc<Mutex<Vec<Selector>>>;
type InitCalls = Arc<Mutex<Vec<bool>>>;
fn changeset(kind: ChangeSetKind, key: &str, selector: Selector) -> FDv2SourceResult {
FDv2SourceResult::ChangeSet(ChangeSet {
kind,
changes: vec![ItemChange::Flag {
key: key.to_string(),
item: StorageItem::Item(basic_flag(key)),
}],
selector,
})
}
fn interrupted() -> FDv2SourceResult {
FDv2SourceResult::Interrupted(ErrorInfo {
kind: ErrorKind::Unknown,
message: "test".into(),
})
}
fn terminal() -> FDv2SourceResult {
FDv2SourceResult::TerminalError(ErrorInfo {
kind: ErrorKind::Unknown,
message: "test".into(),
})
}
fn event(result: FDv2SourceResult) -> FDv2SourceEvent {
FDv2SourceEvent {
result,
fdv1_fallback: None,
}
}
struct MockInitializer {
results: VecDeque<FDv2SourceResult>,
}
impl Initializer for MockInitializer {
fn run(&mut self) -> BoxFuture<'_, FDv2SourceEvent> {
let result = self.results.pop_front().unwrap_or_else(interrupted);
Box::pin(async move { event(result) })
}
fn name(&self) -> &str {
"mock-initializer"
}
}
struct FallbackInitializer {
ttl: Duration,
}
impl Initializer for FallbackInitializer {
fn run(&mut self) -> BoxFuture<'_, FDv2SourceEvent> {
let ttl = self.ttl;
Box::pin(async move {
FDv2SourceEvent {
result: interrupted(),
fdv1_fallback: Some(FDv1FallbackDirective { ttl }),
}
})
}
fn name(&self) -> &str {
"fallback-initializer"
}
}
struct FallbackInitializerFactory {
ttl: Duration,
}
impl InitializerFactory for FallbackInitializerFactory {
fn create(&self) -> Box<dyn Initializer> {
Box::new(FallbackInitializer { ttl: self.ttl })
}
}
struct MockSynchronizer {
results: VecDeque<FDv2SourceResult>,
selectors_seen: Selectors,
shutdown: Option<broadcast::Sender<()>>,
fallback_directive: Option<FDv1FallbackDirective>,
}
impl Synchronizer for MockSynchronizer {
fn next(&mut self, selector: Selector) -> BoxFuture<'_, FDv2SourceEvent> {
self.selectors_seen.lock().unwrap().push(selector);
match self.results.pop_front() {
Some(result) => {
let fdv1_fallback = self.fallback_directive.clone();
Box::pin(async move {
FDv2SourceEvent {
result,
fdv1_fallback,
}
})
}
None => {
if let Some(shutdown) = &self.shutdown {
let _ = shutdown.send(());
}
Box::pin(std::future::pending())
}
}
}
fn name(&self) -> &str {
"mock-synchronizer"
}
}
struct MockInitializerFactory {
results: Mutex<Vec<FDv2SourceResult>>,
}
impl InitializerFactory for MockInitializerFactory {
fn create(&self) -> Box<dyn Initializer> {
let results = std::mem::take(&mut *self.results.lock().unwrap());
Box::new(MockInitializer {
results: results.into(),
})
}
}
fn init_factory(results: Vec<FDv2SourceResult>) -> Arc<dyn InitializerFactory> {
Arc::new(MockInitializerFactory {
results: Mutex::new(results),
})
}
struct MockSynchronizerFactory {
results: Mutex<Vec<FDv2SourceResult>>,
selectors_seen: Selectors,
shutdown: Option<broadcast::Sender<()>>,
is_fdv1_fallback: bool,
fallback_directive: Option<FDv1FallbackDirective>,
}
impl SynchronizerFactory for MockSynchronizerFactory {
fn create(&self) -> Box<dyn Synchronizer> {
let results = std::mem::take(&mut *self.results.lock().unwrap());
Box::new(MockSynchronizer {
results: results.into(),
selectors_seen: self.selectors_seen.clone(),
shutdown: self.shutdown.clone(),
fallback_directive: self.fallback_directive.clone(),
})
}
fn is_fdv1_fallback(&self) -> bool {
self.is_fdv1_fallback
}
}
fn no_selectors() -> Selectors {
Arc::new(Mutex::new(Vec::new()))
}
fn sync_factory(
results: Vec<FDv2SourceResult>,
selectors_seen: Selectors,
shutdown: Option<broadcast::Sender<()>>,
is_fdv1_fallback: bool,
) -> Arc<dyn SynchronizerFactory> {
Arc::new(MockSynchronizerFactory {
results: Mutex::new(results),
selectors_seen,
shutdown,
is_fdv1_fallback,
fallback_directive: None,
})
}
fn fallback_directive_factory(ttl: Duration) -> Arc<dyn SynchronizerFactory> {
Arc::new(MockSynchronizerFactory {
results: Mutex::new(vec![interrupted()]),
selectors_seen: no_selectors(),
shutdown: None,
is_fdv1_fallback: false,
fallback_directive: Some(FDv1FallbackDirective { ttl }),
})
}
struct DownThenDataFactory {
builds: AtomicUsize,
shutdown: broadcast::Sender<()>,
}
impl SynchronizerFactory for DownThenDataFactory {
fn create(&self) -> Box<dyn Synchronizer> {
let (results, shutdown) = if self.builds.fetch_add(1, Ordering::SeqCst) == 0 {
(vec![interrupted()], None)
} else {
(
vec![changeset(
ChangeSetKind::Partial,
"prime-recovered",
Some("s".into()),
)],
Some(self.shutdown.clone()),
)
};
Box::new(MockSynchronizer {
results: results.into(),
selectors_seen: no_selectors(),
shutdown,
fallback_directive: None,
})
}
}
struct FallbackThenDataFactory {
ttl: Duration,
reengage_kind: ChangeSetKind,
builds: AtomicUsize,
shutdown: broadcast::Sender<()>,
}
impl SynchronizerFactory for FallbackThenDataFactory {
fn create(&self) -> Box<dyn Synchronizer> {
if self.builds.fetch_add(1, Ordering::SeqCst) == 0 {
Box::new(MockSynchronizer {
results: VecDeque::from(vec![interrupted()]),
selectors_seen: no_selectors(),
shutdown: None,
fallback_directive: Some(FDv1FallbackDirective { ttl: self.ttl }),
})
} else {
Box::new(MockSynchronizer {
results: VecDeque::from(vec![changeset(
self.reengage_kind,
"fdv2-back",
Some("s".into()),
)]),
selectors_seen: no_selectors(),
shutdown: Some(self.shutdown.clone()),
fallback_directive: None,
})
}
}
}
fn recording_init_complete() -> (Arc<dyn Fn(bool) + Send + Sync>, InitCalls) {
let calls: InitCalls = Arc::new(Mutex::new(Vec::new()));
let sink = calls.clone();
let cb: Arc<dyn Fn(bool) + Send + Sync> =
Arc::new(move |success| sink.lock().unwrap().push(success));
(cb, calls)
}
#[tokio::test]
async fn start_applies_basis_and_exposes_it_via_store_handle() {
let system = FDv2DataSystem::new(
vec![init_factory(vec![changeset(
ChangeSetKind::Full,
"f1",
Some("s1".into()),
)])],
vec![sync_factory(
vec![],
no_selectors(),
None,
false,
)],
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
);
let calls: InitCalls = Arc::new(Mutex::new(Vec::new()));
let notify = Arc::new(tokio::sync::Notify::new());
let sink = calls.clone();
let waker = notify.clone();
let init_complete: Arc<dyn Fn(bool) + Send + Sync> = Arc::new(move |success| {
sink.lock().unwrap().push(success);
waker.notify_one();
});
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
system.start(init_complete, shutdown_rx);
notify.notified().await;
assert_eq!(*calls.lock().unwrap(), vec![true]);
assert!(system.store().read().flag("f1").is_some());
drop(shutdown_tx);
}
#[tokio::test]
async fn initializer_basis_signals_once_and_propagates_selector() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let selectors_seen: Selectors = Arc::new(Mutex::new(Vec::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> =
vec![init_factory(vec![changeset(
ChangeSetKind::Full,
"init-flag",
Some("sel-1".into()),
)])];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![sync_factory(
vec![changeset(
ChangeSetKind::Partial,
"sync-flag",
Some("sel-2".into()),
)],
selectors_seen.clone(),
Some(shutdown_tx),
false,
)]);
run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
assert_eq!(*calls.lock().unwrap(), vec![true]);
assert!(store.read().flag("init-flag").is_some());
assert!(store.read().flag("sync-flag").is_some());
assert_eq!(selectors_seen.lock().unwrap()[0], Some("sel-1".into()));
}
#[tokio::test]
async fn selectorless_full_continues_to_next_initializer() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let selectors_seen: Selectors = Arc::new(Mutex::new(Vec::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![
init_factory(vec![changeset(ChangeSetKind::Full, "no-basis-flag", None)]),
init_factory(vec![changeset(
ChangeSetKind::Partial,
"merged-flag",
Some("sel-2".into()),
)]),
];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![sync_factory(
vec![],
selectors_seen.clone(),
Some(shutdown_tx),
false,
)]);
run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
assert!(store.read().flag("no-basis-flag").is_some());
assert!(store.read().flag("merged-flag").is_some());
assert_eq!(selectors_seen.lock().unwrap()[0], Some("sel-2".into()));
assert_eq!(*calls.lock().unwrap(), vec![true]);
}
#[tokio::test]
async fn initializer_selectorless_full_defers_signal() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let flag_at_signal = Arc::new(Mutex::new(None));
let probe = store.clone();
let sink = flag_at_signal.clone();
let init_complete: Arc<dyn Fn(bool) + Send + Sync> = Arc::new(move |_| {
*sink.lock().unwrap() = Some(probe.read().flag("from-second").is_some());
});
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![
init_factory(vec![changeset(ChangeSetKind::Full, "from-first", None)]),
init_factory(vec![changeset(ChangeSetKind::Full, "from-second", None)]),
];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![sync_factory(
vec![],
no_selectors(),
Some(shutdown_tx),
false,
)]);
run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
assert_eq!(*flag_at_signal.lock().unwrap(), Some(true));
}
#[tokio::test]
async fn basis_stops_later_initializers() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let selectors_seen: Selectors = Arc::new(Mutex::new(Vec::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![
init_factory(vec![changeset(
ChangeSetKind::Full,
"from-first",
Some("s1".into()),
)]),
init_factory(vec![changeset(
ChangeSetKind::Full,
"from-second",
Some("s2".into()),
)]),
];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![sync_factory(
vec![],
selectors_seen.clone(),
Some(shutdown_tx),
false,
)]);
run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
assert!(store.read().flag("from-first").is_some());
assert!(store.read().flag("from-second").is_none());
assert_eq!(*calls.lock().unwrap(), vec![true]);
assert_eq!(selectors_seen.lock().unwrap()[0], Some("s1".into()));
}
#[tokio::test]
async fn none_changeset_does_not_clobber_the_selector() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let selectors_seen: Selectors = Arc::new(Mutex::new(Vec::new()));
let (init_complete, _calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![sync_factory(
vec![
changeset(ChangeSetKind::Full, "flag", Some("s1".into())),
changeset(ChangeSetKind::None, "flag", None),
],
selectors_seen.clone(),
Some(shutdown_tx),
false,
)]);
run(
initializer_factories,
source_manager,
store,
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
let seen = selectors_seen.lock().unwrap();
assert_eq!(*seen, vec![None, Some("s1".into()), Some("s1".into())]);
}
#[tokio::test]
async fn selectorless_change_clears_the_selector() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let selectors_seen: Selectors = Arc::new(Mutex::new(Vec::new()));
let (init_complete, _calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![sync_factory(
vec![
changeset(ChangeSetKind::Full, "flag", Some("s1".into())),
changeset(ChangeSetKind::Partial, "flag", None),
],
selectors_seen.clone(),
Some(shutdown_tx),
false,
)]);
run(
initializer_factories,
source_manager,
store,
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
let seen = selectors_seen.lock().unwrap();
assert_eq!(*seen, vec![None, Some("s1".into()), None]);
}
#[tokio::test]
async fn initializer_delta_does_not_initialize() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> =
vec![init_factory(vec![changeset(
ChangeSetKind::Partial,
"delta-flag",
Some("s1".into()),
)])];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![sync_factory(
vec![terminal()],
no_selectors(),
Some(shutdown_tx),
false,
)]);
run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
assert!(store.read().flag("delta-flag").is_some());
assert_eq!(*calls.lock().unwrap(), vec![false]);
}
#[tokio::test]
async fn synchronizer_delta_does_not_initialize() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![sync_factory(
vec![
changeset(ChangeSetKind::Partial, "delta-flag", Some("s1".into())),
terminal(),
],
no_selectors(),
Some(shutdown_tx),
false,
)]);
run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
assert!(store.read().flag("delta-flag").is_some());
assert_eq!(*calls.lock().unwrap(), vec![false]);
}
#[tokio::test]
async fn synchronizer_selectorless_full_initializes() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![sync_factory(
vec![changeset(ChangeSetKind::Full, "full-flag", None)],
no_selectors(),
Some(shutdown_tx),
false,
)]);
run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
assert_eq!(*calls.lock().unwrap(), vec![true]);
assert!(store.read().flag("full-flag").is_some());
}
#[tokio::test]
async fn failed_initializers_let_synchronizer_provide_the_basis() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![
init_factory(vec![interrupted()]),
init_factory(vec![terminal()]),
];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![sync_factory(
vec![changeset(
ChangeSetKind::Full,
"sync-flag",
Some("s".into()),
)],
no_selectors(),
Some(shutdown_tx),
false,
)]);
run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
assert_eq!(*calls.lock().unwrap(), vec![true]);
assert!(store.read().flag("sync-flag").is_some());
}
#[tokio::test]
async fn exhausting_all_sources_signals_failure_once() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> =
vec![init_factory(vec![terminal()])];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![sync_factory(
vec![interrupted(), terminal()],
no_selectors(),
Some(shutdown_tx),
false,
)]);
run(
initializer_factories,
source_manager,
store,
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
assert_eq!(*calls.lock().unwrap(), vec![false]);
}
#[tokio::test]
async fn synchronizer_terminal_error_advances_to_next() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![
sync_factory(
vec![terminal()],
no_selectors(),
Some(shutdown_tx.clone()),
false,
),
sync_factory(
vec![changeset(
ChangeSetKind::Full,
"from-second",
Some("s".into()),
)],
no_selectors(),
Some(shutdown_tx),
false,
),
]);
run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
assert_eq!(*calls.lock().unwrap(), vec![true]);
assert!(store.read().flag("from-second").is_some());
}
#[tokio::test]
async fn synchronizer_interrupted_retries_same_source() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![sync_factory(
vec![
interrupted(),
changeset(ChangeSetKind::Full, "after-retry", Some("s".into())),
],
no_selectors(),
Some(shutdown_tx),
false,
)]);
run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
assert_eq!(*calls.lock().unwrap(), vec![true]);
assert!(store.read().flag("after-retry").is_some());
}
#[tokio::test]
async fn shutdown_ends_run_without_signaling() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
let source_manager = SourceManager::new(vec![sync_factory(
vec![],
no_selectors(),
None,
false,
)]);
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let handle = tokio::spawn(run(
initializer_factories,
source_manager,
store,
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
));
shutdown_tx.send(()).unwrap();
handle.await.unwrap();
assert!(calls.lock().unwrap().is_empty());
}
#[test]
fn rotates_cyclically_skips_blocked_and_exhausts() {
let mut sources = SourceManager::new(vec![
sync_factory(
vec![],
no_selectors(),
None,
false,
),
sync_factory(
vec![],
no_selectors(),
None,
false,
),
sync_factory(
vec![],
no_selectors(),
None,
false,
),
]);
sources.next_synchronizer();
assert_eq!(sources.current_factory_index, Some(0));
sources.next_synchronizer();
assert_eq!(sources.current_factory_index, Some(1));
sources.block_current();
sources.next_synchronizer();
assert_eq!(sources.current_factory_index, Some(2));
sources.next_synchronizer();
assert_eq!(sources.current_factory_index, Some(0));
sources.block_current();
sources.next_synchronizer();
sources.block_current();
assert!(sources.next_synchronizer().is_none());
}
#[test]
fn reset_source_index_returns_to_prime() {
let mut sources = SourceManager::new(vec![
sync_factory(
vec![],
no_selectors(),
None,
false,
),
sync_factory(
vec![],
no_selectors(),
None,
false,
),
]);
sources.next_synchronizer();
sources.next_synchronizer();
assert_eq!(sources.current_factory_index, Some(1));
sources.reset_source_index();
sources.next_synchronizer();
assert_eq!(sources.current_factory_index, Some(0));
}
#[test]
fn is_prime_and_available_count_track_state() {
let mut sources = SourceManager::new(vec![
sync_factory(
vec![],
no_selectors(),
None,
false,
),
sync_factory(
vec![],
no_selectors(),
None,
false,
),
]);
assert_eq!(sources.available_count(), 2);
sources.next_synchronizer();
assert!(sources.is_prime());
sources.next_synchronizer();
assert!(!sources.is_prime());
sources.block_current();
assert_eq!(sources.available_count(), 1);
}
#[test]
fn fdv1_fallback_factory_starts_blocked() {
let mut sources = SourceManager::new(vec![
sync_factory(
vec![],
no_selectors(),
None,
false,
),
sync_factory(
vec![],
no_selectors(),
None,
true,
),
]);
assert_eq!(sources.available_count(), 1);
sources.next_synchronizer();
assert!(!sources.is_current_fdv1_fallback());
}
#[test]
fn switch_to_and_back_from_fdv1_fallback() {
let mut sources = SourceManager::new(vec![
sync_factory(
vec![],
no_selectors(),
None,
false,
),
sync_factory(
vec![],
no_selectors(),
None,
true,
),
]);
sources.switch_to_fdv1_fallback();
assert_eq!(sources.available_count(), 1);
sources.next_synchronizer();
assert!(sources.is_current_fdv1_fallback());
sources.switch_back_to_fdv2();
assert_eq!(sources.available_count(), 1);
sources.next_synchronizer();
assert!(!sources.is_current_fdv1_fallback());
}
#[test]
fn switch_back_unblocks_terminally_blocked_fdv2() {
let mut sources = SourceManager::new(vec![
sync_factory(
vec![],
no_selectors(),
None,
false,
),
sync_factory(
vec![],
no_selectors(),
None,
true,
),
]);
sources.next_synchronizer();
sources.block_current();
assert_eq!(sources.available_count(), 0);
sources.switch_to_fdv1_fallback();
sources.switch_back_to_fdv2();
assert_eq!(sources.available_count(), 1);
}
#[tokio::test(start_paused = true)]
async fn fallback_fires_after_sustained_interruption() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![
sync_factory(
vec![interrupted()],
no_selectors(),
None,
false,
),
sync_factory(
vec![changeset(
ChangeSetKind::Full,
"from-fallback",
Some("s".into()),
)],
no_selectors(),
Some(shutdown_tx),
false,
),
]);
let handle = tokio::spawn(run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
));
handle.await.unwrap();
assert_eq!(*calls.lock().unwrap(), vec![true]);
assert!(store.read().flag("from-fallback").is_some());
}
#[tokio::test(start_paused = true)]
async fn changeset_cancels_the_fallback_timer() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let calls: InitCalls = Arc::new(Mutex::new(Vec::new()));
let notify = Arc::new(tokio::sync::Notify::new());
let sink = calls.clone();
let waker = notify.clone();
let init_complete: Arc<dyn Fn(bool) + Send + Sync> = Arc::new(move |success| {
sink.lock().unwrap().push(success);
waker.notify_one();
});
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![
sync_factory(
vec![
interrupted(),
changeset(ChangeSetKind::Full, "from-prime", Some("s".into())),
],
no_selectors(),
None,
false,
),
sync_factory(
vec![changeset(
ChangeSetKind::Full,
"from-fallback",
Some("s".into()),
)],
no_selectors(),
Some(shutdown_tx.clone()),
false,
),
]);
let handle = tokio::spawn(run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
));
notify.notified().await;
tokio::time::advance(FALLBACK_TIMEOUT * 2).await;
shutdown_tx.send(()).unwrap();
handle.await.unwrap();
assert!(store.read().flag("from-prime").is_some());
assert!(store.read().flag("from-fallback").is_none());
}
#[tokio::test(start_paused = true)]
async fn recovery_fires_and_returns_to_the_prime() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![
Arc::new(DownThenDataFactory {
builds: AtomicUsize::new(0),
shutdown: shutdown_tx,
}) as Arc<dyn SynchronizerFactory>,
sync_factory(
vec![changeset(
ChangeSetKind::Full,
"from-fallback",
Some("s".into()),
)],
no_selectors(),
None,
false,
),
]);
let handle = tokio::spawn(run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
));
handle.await.unwrap();
assert_eq!(*calls.lock().unwrap(), vec![true]);
assert!(store.read().flag("from-fallback").is_some());
assert!(store.read().flag("prime-recovered").is_some());
}
#[tokio::test]
async fn fallback_directive_switches_to_fdv1() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![
fallback_directive_factory(Duration::from_secs(60)),
sync_factory(
vec![changeset(
ChangeSetKind::Full,
"from-fdv1",
Some("s".into()),
)],
no_selectors(),
Some(shutdown_tx),
true,
),
]);
run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
assert_eq!(*calls.lock().unwrap(), vec![true]);
assert!(store.read().flag("from-fdv1").is_some());
}
#[tokio::test(start_paused = true)]
async fn fdv2_retry_reengages_after_ttl() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![
Arc::new(FallbackThenDataFactory {
ttl: Duration::from_secs(60),
reengage_kind: ChangeSetKind::Partial,
builds: AtomicUsize::new(0),
shutdown: shutdown_tx,
}) as Arc<dyn SynchronizerFactory>,
sync_factory(
vec![changeset(
ChangeSetKind::Full,
"from-fdv1",
Some("s".into()),
)],
no_selectors(),
None,
true,
),
]);
let handle = tokio::spawn(run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
));
handle.await.unwrap();
assert_eq!(*calls.lock().unwrap(), vec![true]);
assert!(store.read().flag("from-fdv1").is_some());
assert!(store.read().flag("fdv2-back").is_some());
}
#[tokio::test(start_paused = true)]
async fn fdv2_retry_survives_a_terminal_fdv1_fallback() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> = vec![];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![
Arc::new(FallbackThenDataFactory {
ttl: Duration::from_secs(60),
reengage_kind: ChangeSetKind::Full,
builds: AtomicUsize::new(0),
shutdown: shutdown_tx.clone(),
}) as Arc<dyn SynchronizerFactory>,
sync_factory(
vec![terminal()],
no_selectors(),
Some(shutdown_tx),
true,
),
]);
let handle = tokio::spawn(run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
));
handle.await.unwrap();
assert_eq!(*calls.lock().unwrap(), vec![true]);
assert!(store.read().flag("fdv2-back").is_some());
}
#[tokio::test]
async fn initializer_fallback_directive_switches_to_fdv1() {
let store = Arc::new(RwLock::new(InMemoryDataStore::new()));
let (init_complete, calls) = recording_init_complete();
let initializer_factories: Vec<Arc<dyn InitializerFactory>> =
vec![Arc::new(FallbackInitializerFactory {
ttl: Duration::from_secs(60),
})];
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
let source_manager = SourceManager::new(vec![
sync_factory(
vec![],
no_selectors(),
Some(shutdown_tx.clone()),
false,
),
sync_factory(
vec![changeset(
ChangeSetKind::Full,
"from-fdv1",
Some("s".into()),
)],
no_selectors(),
Some(shutdown_tx),
true,
),
]);
run(
initializer_factories,
source_manager,
store.clone(),
init_complete,
shutdown_rx,
FALLBACK_TIMEOUT,
RECOVERY_TIMEOUT,
)
.await;
assert_eq!(*calls.lock().unwrap(), vec![true]);
assert!(store.read().flag("from-fdv1").is_some());
}
}