use std::future::Future;
use crate::runtime::task::BlockingBridge;
use crate::server::database::PvDatabase;
pub fn spawn_program<F, Fut, E>(
bridge: &BlockingBridge,
db: &PvDatabase,
program: &'static str,
run: F,
) where
F: FnOnce(PvDatabase) -> Fut + Send + 'static,
Fut: Future<Output = Result<(), E>> + Send + 'static,
E: std::fmt::Display,
{
let db = db.clone();
bridge.spawn(async move {
db.wait_for_pini().await;
if let Err(e) = run(db).await {
eprintln!("{program} error: {e}");
}
});
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use super::*;
async fn eventually(probe: impl Fn() -> bool) -> bool {
for _ in 0..300 {
if probe() {
return true;
}
crate::runtime::task::sleep(Duration::from_millis(10)).await;
}
false
}
fn started_flag() -> (
Arc<AtomicBool>,
impl FnOnce(PvDatabase) -> std::future::Ready<Result<(), String>> + Send + 'static,
) {
let started = Arc::new(AtomicBool::new(false));
let mark = Arc::clone(&started);
(started, move |_db: PvDatabase| {
mark.store(true, Ordering::Release);
std::future::ready(Ok(()))
})
}
#[epics_macros_rs::epics_test]
async fn a_program_started_before_pini_runs_only_after_it() {
let db = PvDatabase::new();
let bridge = BlockingBridge::capture();
let (started, run) = started_flag();
spawn_program(&bridge, &db, "probe", run);
crate::runtime::task::sleep(Duration::from_millis(100)).await;
assert!(
!started.load(Ordering::Acquire),
"the program ran before the PINI pass was published"
);
db.mark_pini_done();
assert!(
eventually(|| started.load(Ordering::Acquire)).await,
"the program did not start once PINI was published"
);
}
#[epics_macros_rs::epics_test]
async fn a_program_started_after_pini_runs_at_once() {
let db = PvDatabase::new();
db.mark_pini_done();
let bridge = BlockingBridge::capture();
let (started, run) = started_flag();
spawn_program(&bridge, &db, "probe", run);
assert!(eventually(|| started.load(Ordering::Acquire)).await);
}
}