220 lines
9.3 KiB
Rust
220 lines
9.3 KiB
Rust
// file: crates/ksp-worker-api/tests/snapshot_source.rs
|
|
// version: 1
|
|
|
|
//! External std-only latest-value source, object-safety and resynchronization canaries.
|
|
|
|
#[derive(Clone)]
|
|
struct TestSnapshotSource {
|
|
state: std::sync::Arc<TestState>,
|
|
}
|
|
|
|
struct TestState {
|
|
current: std::sync::Mutex<TestInner>,
|
|
}
|
|
|
|
struct TestInner {
|
|
snapshot: ksp_worker_api::WorkerSnapshot,
|
|
waiters: std::vec::Vec<std::task::Waker>,
|
|
}
|
|
|
|
impl TestSnapshotSource {
|
|
fn new(snapshot: ksp_worker_api::WorkerSnapshot) -> Self {
|
|
return Self {
|
|
state: std::sync::Arc::new(TestState { current: std::sync::Mutex::new(TestInner { snapshot, waiters: std::vec::Vec::new() }) }),
|
|
};
|
|
}
|
|
|
|
fn publish(&self, state: ksp_worker_api::WorkerState, health: ksp_worker_api::WorkerHealth, activity: ksp_worker_api::WorkerActivity) -> bool {
|
|
let mut current = match self.state.current.lock() {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(poisoned) => poisoned.into_inner(),
|
|
};
|
|
let sequence = match current.snapshot.sequence().next() {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return false,
|
|
};
|
|
let id = current.snapshot.id().clone();
|
|
let kind = current.snapshot.kind().clone();
|
|
current.snapshot = ksp_worker_api::WorkerSnapshot::new(id, kind, sequence, state, health, activity);
|
|
let waiters = std::mem::take(&mut current.waiters);
|
|
std::mem::drop(current);
|
|
for waiter in waiters {
|
|
waiter.wake();
|
|
}
|
|
return true;
|
|
}
|
|
}
|
|
|
|
struct TestWaitFuture<'a> {
|
|
source: &'a TestSnapshotSource,
|
|
observed: ksp_worker_api::WorkerSnapshotSequence,
|
|
}
|
|
|
|
impl std::future::Future for TestWaitFuture<'_> {
|
|
type Output = ksp_worker_api::WorkerSnapshot;
|
|
|
|
fn poll(self: std::pin::Pin<&mut Self>, context: &mut std::task::Context<'_>) -> std::task::Poll<Self::Output> {
|
|
let mut current = match self.source.state.current.lock() {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(poisoned) => poisoned.into_inner(),
|
|
};
|
|
if current.snapshot.sequence().is_after(self.observed) {
|
|
return std::task::Poll::Ready(current.snapshot.clone());
|
|
}
|
|
if !current.waiters.iter().any(|registered| return registered.will_wake(context.waker())) {
|
|
current.waiters.push(context.waker().clone());
|
|
}
|
|
return std::task::Poll::Pending;
|
|
}
|
|
}
|
|
|
|
impl ksp_worker_api::WorkerSnapshotSource for TestSnapshotSource {
|
|
fn current(&self) -> ksp_worker_api::WorkerSnapshot {
|
|
let current = match self.state.current.lock() {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(poisoned) => poisoned.into_inner(),
|
|
};
|
|
return current.snapshot.clone();
|
|
}
|
|
|
|
fn wait_for_change(&self, observed: ksp_worker_api::WorkerSnapshotSequence) -> ksp_worker_api::WorkerSnapshotFuture<'_> {
|
|
return std::boxed::Box::pin(TestWaitFuture { source: self, observed });
|
|
}
|
|
}
|
|
|
|
struct WakeProbe {
|
|
woken: std::sync::atomic::AtomicBool,
|
|
}
|
|
|
|
impl WakeProbe {
|
|
fn new() -> Self {
|
|
return Self { woken: std::sync::atomic::AtomicBool::new(false) };
|
|
}
|
|
|
|
fn is_woken(&self) -> bool {
|
|
return self.woken.load(std::sync::atomic::Ordering::Acquire);
|
|
}
|
|
}
|
|
|
|
impl std::task::Wake for WakeProbe {
|
|
fn wake(self: std::sync::Arc<Self>) {
|
|
self.woken.store(true, std::sync::atomic::Ordering::Release);
|
|
return;
|
|
}
|
|
}
|
|
|
|
fn poll_snapshot(future: &mut ksp_worker_api::WorkerSnapshotFuture<'_>, wake: &std::sync::Arc<WakeProbe>) -> std::task::Poll<ksp_worker_api::WorkerSnapshot> {
|
|
let waker = std::task::Waker::from(wake.clone());
|
|
let mut context = std::task::Context::from_waker(&waker);
|
|
return std::future::Future::poll(future.as_mut(), &mut context);
|
|
}
|
|
|
|
fn initial_snapshot() -> std::option::Option<ksp_worker_api::WorkerSnapshot> {
|
|
let id = match ksp_worker_api::WorkerId::new("external-source-001") {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
let kind = match ksp_worker_api::WorkerKindCode::new("external_test_worker") {
|
|
std::result::Result::Ok(value) => value,
|
|
std::result::Result::Err(_) => return std::option::Option::None,
|
|
};
|
|
return std::option::Option::Some(ksp_worker_api::WorkerSnapshot::new(
|
|
id,
|
|
kind,
|
|
ksp_worker_api::WorkerSnapshotSequence::initial(),
|
|
ksp_worker_api::WorkerState::Running,
|
|
ksp_worker_api::WorkerHealth::Healthy,
|
|
ksp_worker_api::WorkerActivity::Idle,
|
|
));
|
|
}
|
|
|
|
#[test]
|
|
fn pre_003_snapshot_source_is_object_safe_send_sync_and_externally_implementable() {
|
|
fn require_send_sync<T: std::marker::Send + std::marker::Sync>() {}
|
|
require_send_sync::<TestSnapshotSource>();
|
|
let snapshot = initial_snapshot();
|
|
assert!(snapshot.is_some());
|
|
let source = match snapshot {
|
|
std::option::Option::Some(value) => TestSnapshotSource::new(value),
|
|
std::option::Option::None => return,
|
|
};
|
|
let object: &dyn ksp_worker_api::WorkerSnapshotSource = &source;
|
|
let current = object.current();
|
|
assert_eq!(current.sequence().value(), 0);
|
|
let mut wait = object.wait_for_change(current.sequence());
|
|
fn require_send<T: std::marker::Send>(_: &T) {}
|
|
require_send(&wait);
|
|
let wake = std::sync::Arc::new(WakeProbe::new());
|
|
assert!(matches!(poll_snapshot(&mut wait, &wake), std::task::Poll::Pending));
|
|
assert!(source.publish(ksp_worker_api::WorkerState::Running, ksp_worker_api::WorkerHealth::Healthy, ksp_worker_api::WorkerActivity::Active));
|
|
assert!(wake.is_woken());
|
|
let changed = poll_snapshot(&mut wait, &wake);
|
|
assert!(matches!(changed, std::task::Poll::Ready(_)));
|
|
return;
|
|
}
|
|
|
|
#[test]
|
|
fn pre_003_slow_and_independent_listeners_coalesce_to_latest_value() {
|
|
let snapshot = initial_snapshot();
|
|
assert!(snapshot.is_some());
|
|
let source = match snapshot {
|
|
std::option::Option::Some(value) => TestSnapshotSource::new(value),
|
|
std::option::Option::None => return,
|
|
};
|
|
let observed = ksp_worker_api::WorkerSnapshotSource::current(&source).sequence();
|
|
let mut slow = ksp_worker_api::WorkerSnapshotSource::wait_for_change(&source, observed);
|
|
let mut fast = ksp_worker_api::WorkerSnapshotSource::wait_for_change(&source, observed);
|
|
let slow_wake = std::sync::Arc::new(WakeProbe::new());
|
|
let fast_wake = std::sync::Arc::new(WakeProbe::new());
|
|
assert!(matches!(poll_snapshot(&mut slow, &slow_wake), std::task::Poll::Pending));
|
|
assert!(matches!(poll_snapshot(&mut fast, &fast_wake), std::task::Poll::Pending));
|
|
assert!(source.publish(ksp_worker_api::WorkerState::Running, ksp_worker_api::WorkerHealth::Degraded, ksp_worker_api::WorkerActivity::Active));
|
|
assert!(slow_wake.is_woken());
|
|
assert!(fast_wake.is_woken());
|
|
let fast_value = match poll_snapshot(&mut fast, &fast_wake) {
|
|
std::task::Poll::Ready(value) => value,
|
|
std::task::Poll::Pending => return,
|
|
};
|
|
assert_eq!(fast_value.sequence().value(), 1);
|
|
assert_eq!(fast_value.health(), ksp_worker_api::WorkerHealth::Degraded);
|
|
assert!(source.publish(ksp_worker_api::WorkerState::Running, ksp_worker_api::WorkerHealth::Healthy, ksp_worker_api::WorkerActivity::Idle));
|
|
let slow_value = match poll_snapshot(&mut slow, &slow_wake) {
|
|
std::task::Poll::Ready(value) => value,
|
|
std::task::Poll::Pending => return,
|
|
};
|
|
assert_eq!(slow_value.sequence().value(), 2);
|
|
assert_eq!(slow_value.health(), ksp_worker_api::WorkerHealth::Healthy);
|
|
assert_eq!(slow_value.activity(), ksp_worker_api::WorkerActivity::Idle);
|
|
return;
|
|
}
|
|
|
|
#[test]
|
|
fn pre_003_late_listener_resynchronizes_and_terminal_snapshot_remains_current() {
|
|
let snapshot = initial_snapshot();
|
|
assert!(snapshot.is_some());
|
|
let source = match snapshot {
|
|
std::option::Option::Some(value) => TestSnapshotSource::new(value),
|
|
std::option::Option::None => return,
|
|
};
|
|
assert!(source.publish(ksp_worker_api::WorkerState::Running, ksp_worker_api::WorkerHealth::Degraded, ksp_worker_api::WorkerActivity::Active));
|
|
assert!(source.publish(ksp_worker_api::WorkerState::Stopping, ksp_worker_api::WorkerHealth::Healthy, ksp_worker_api::WorkerActivity::Idle));
|
|
let current = ksp_worker_api::WorkerSnapshotSource::current(&source);
|
|
assert_eq!(current.sequence().value(), 2);
|
|
assert_eq!(current.state(), ksp_worker_api::WorkerState::Stopping);
|
|
let mut wait = ksp_worker_api::WorkerSnapshotSource::wait_for_change(&source, current.sequence());
|
|
let wake = std::sync::Arc::new(WakeProbe::new());
|
|
assert!(matches!(poll_snapshot(&mut wait, &wake), std::task::Poll::Pending));
|
|
assert!(source.publish(ksp_worker_api::WorkerState::Stopped, ksp_worker_api::WorkerHealth::Healthy, ksp_worker_api::WorkerActivity::Idle));
|
|
assert!(wake.is_woken());
|
|
let terminal = match poll_snapshot(&mut wait, &wake) {
|
|
std::task::Poll::Ready(value) => value,
|
|
std::task::Poll::Pending => return,
|
|
};
|
|
assert_eq!(terminal.sequence().value(), 3);
|
|
assert_eq!(terminal.state(), ksp_worker_api::WorkerState::Stopped);
|
|
assert!(terminal.state().is_terminal());
|
|
let retained = ksp_worker_api::WorkerSnapshotSource::current(&source);
|
|
assert_eq!(retained, terminal);
|
|
return;
|
|
}
|