use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use asyn_rs::error::AsynResult;
use asyn_rs::param::ParamType;
use asyn_rs::port::{PortDriver, PortDriverBase, PortFlags};
use asyn_rs::port_handle::PortHandle;
use asyn_rs::request::RequestOp;
use asyn_rs::runtime::config::RuntimeConfig;
use asyn_rs::runtime::port::create_port_runtime;
use asyn_rs::user::AsynUser;
const WATCHDOG: Duration = Duration::from_secs(10);
fn run_with_watchdog<F>(what: &str, f: F) -> Result<(), String>
where
F: FnOnce() + Send + 'static,
{
let (tx, rx) = std::sync::mpsc::channel::<std::thread::Result<()>>();
std::thread::Builder::new()
.name(format!("watchdog-{what}"))
.spawn(move || {
let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(f));
let _ = tx.send(outcome);
})
.expect("spawn");
match rx.recv_timeout(WATCHDOG) {
Ok(Ok(())) => Ok(()),
Ok(Err(_)) => Err(format!("{what} panicked")),
Err(_) => Err(format!("{what} hung (no result within {WATCHDOG:?})")),
}
}
struct PlainDriver {
base: PortDriverBase,
}
impl PlainDriver {
fn new(name: &str) -> Self {
let mut base = PortDriverBase::new(name, 1, PortFlags::default());
base.create_param("VAL", ParamType::Int32).unwrap();
Self { base }
}
}
impl PortDriver for PlainDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
}
#[test]
fn submit_blocking_from_foreign_current_thread_runtime_succeeds() {
let (handle, _join) = create_port_runtime(PlainDriver::new("ctx_a"), RuntimeConfig::default());
let port: PortHandle = handle.port_handle().clone();
run_with_watchdog("ctx-a", move || {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
rt.block_on(async {
port.write_int32_blocking(0, 0, 42)
.expect("write from a foreign current-thread runtime must succeed");
assert_eq!(
port.read_int32_blocking(0, 0)
.expect("read from a foreign current-thread runtime must succeed"),
42
);
});
})
.unwrap();
}
struct SelfCallingDriver {
base: PortDriverBase,
self_handle: Arc<parking_lot::Mutex<Option<PortHandle>>>,
got_error: Arc<AtomicBool>,
}
impl PortDriver for SelfCallingDriver {
fn base(&self) -> &PortDriverBase {
&self.base
}
fn base_mut(&mut self) -> &mut PortDriverBase {
&mut self.base
}
fn io_read_int32(&mut self, _user: &AsynUser) -> AsynResult<i32> {
let me = self.self_handle.lock().clone().expect("handle installed");
match me.submit_blocking(RequestOp::Int32Read, AsynUser::new(0)) {
Ok(_) => panic!("re-entrant submit_blocking must not report success"),
Err(_) => {
self.got_error.store(true, Ordering::SeqCst);
Ok(7)
}
}
}
}
#[test]
fn submit_blocking_from_own_actor_thread_returns_error() {
let self_handle = Arc::new(parking_lot::Mutex::new(None));
let got_error = Arc::new(AtomicBool::new(false));
let mut base = PortDriverBase::new("ctx_b", 1, PortFlags::default());
base.create_param("VAL", ParamType::Int32).unwrap();
let driver = SelfCallingDriver {
base,
self_handle: self_handle.clone(),
got_error: got_error.clone(),
};
let (handle, _join) = create_port_runtime(driver, RuntimeConfig::default());
let port: PortHandle = handle.port_handle().clone();
*self_handle.lock() = Some(port.clone());
let got_error_probe = got_error.clone();
run_with_watchdog("ctx-b", move || {
let v = port
.read_int32_blocking(0, 0)
.expect("the outer request itself must still complete");
assert_eq!(v, 7, "driver returned its sentinel after the inner error");
assert!(
got_error_probe.load(Ordering::SeqCst),
"re-entrant submit_blocking must return an error to the driver"
);
})
.unwrap();
let port2: PortHandle = handle.port_handle().clone();
run_with_watchdog("ctx-b-alive", move || {
port2
.write_int32_blocking(0, 0, 5)
.expect("port must still serve requests after refusing a re-entrant call");
})
.unwrap();
}