use std::sync::{Mutex, OnceLock};
pub type ReportFn = Box<dyn Fn(u32, &dyn Fn(&str)) + Send + Sync>;
pub struct DbServer {
pub name: &'static str,
pub report: Option<ReportFn>,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum ServerState {
Registering,
Running,
Paused,
Stopped,
}
impl ServerState {
fn name(self) -> &'static str {
match self {
ServerState::Registering => "registering",
ServerState::Running => "running",
ServerState::Paused => "paused",
ServerState::Stopped => "stopped",
}
}
}
struct Registry {
servers: Vec<DbServer>,
state: ServerState,
}
fn registry() -> &'static Mutex<Registry> {
static REGISTRY: OnceLock<Mutex<Registry>> = OnceLock::new();
REGISTRY.get_or_init(|| {
Mutex::new(Registry {
servers: Vec::new(),
state: ServerState::Registering,
})
})
}
fn ignored_by_env(name: &str) -> bool {
let Ok(ignore) = std::env::var("EPICS_IOC_IGNORE_SERVERS") else {
return false;
};
ignore.split(' ').any(|tok| tok == name)
}
pub fn db_register_server(server: DbServer) -> bool {
admit(&mut registry().lock().unwrap(), server)
}
fn admit(reg: &mut Registry, server: DbServer) -> bool {
if reg.state != ServerState::Registering {
return false;
}
if server.name.contains(' ') {
eprintln!("dbRegisterServer: Bad server name '{}'", server.name);
return false;
}
if ignored_by_env(server.name) {
eprintln!(
"dbRegisterServer: Ignoring '{}', per environment",
server.name
);
return true;
}
if reg.servers.iter().any(|s| s.name == server.name) {
eprintln!("dbRegisterServer: Can't redefine '{}'.", server.name);
return false;
}
reg.servers.push(server);
true
}
pub fn db_run_servers() {
registry().lock().unwrap().state = ServerState::Running;
}
pub fn db_pause_servers() {
registry().lock().unwrap().state = ServerState::Paused;
}
pub fn db_stop_servers() {
registry().lock().unwrap().state = ServerState::Stopped;
}
fn serving() -> &'static (std::sync::atomic::AtomicU64, tokio::sync::Notify) {
static SERVING: OnceLock<(std::sync::atomic::AtomicU64, tokio::sync::Notify)> = OnceLock::new();
SERVING.get_or_init(|| {
(
std::sync::atomic::AtomicU64::new(0),
tokio::sync::Notify::new(),
)
})
}
pub fn serving_generation() -> u64 {
serving().0.load(std::sync::atomic::Ordering::Acquire)
}
pub fn announce_serving() {
serving()
.0
.fetch_add(1, std::sync::atomic::Ordering::Release);
serving().1.notify_waiters();
}
pub async fn serving_after(generation: u64) {
loop {
let notified = serving().1.notified();
if serving_generation() != generation {
return;
}
notified.await;
}
}
pub fn dbsr(level: u32, out: &dyn Fn(&str)) {
render(®istry().lock().unwrap(), level, out);
}
fn render(reg: &Registry, level: u32, out: &dyn Fn(&str)) {
if reg.servers.is_empty() {
out("No server layers registered with IOC");
return;
}
out(&format!("Server state: {}", reg.state.name()));
for srv in ®.servers {
out(&format!("Server '{}'", srv.name));
if reg.state == ServerState::Running {
if let Some(report) = &srv.report {
report(level, out);
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn ignore_servers_matches_whole_names_only() {
unsafe { std::env::set_var("EPICS_IOC_IGNORE_SERVERS", "qsrv2") };
assert!(ignored_by_env("qsrv2"));
assert!(!ignored_by_env("qsrv"), "a prefix must not be suppressed");
assert!(!ignored_by_env("rsrv"));
unsafe { std::env::set_var("EPICS_IOC_IGNORE_SERVERS", "rsrv qsrv2") };
assert!(ignored_by_env("rsrv"));
assert!(ignored_by_env("qsrv2"));
unsafe { std::env::remove_var("EPICS_IOC_IGNORE_SERVERS") };
assert!(!ignored_by_env("rsrv"));
}
#[test]
fn dbsr_with_no_layers_is_one_line() {
let out = std::cell::RefCell::new(Vec::<String>::new());
let reg = Registry {
servers: Vec::new(),
state: ServerState::Running,
};
render(®, 0, &|s: &str| out.borrow_mut().push(s.to_string()));
assert_eq!(*out.borrow(), vec!["No server layers registered with IOC"]);
}
#[test]
fn dbsr_names_the_paused_phase_and_suppresses_the_report() {
let out = std::cell::RefCell::new(Vec::<String>::new());
let reg = Registry {
servers: vec![DbServer {
name: "rsrv",
report: Some(Box::new(|_level, o: &dyn Fn(&str)| {
o("Channel Access Server");
})),
}],
state: ServerState::Paused,
};
render(®, 0, &|s: &str| out.borrow_mut().push(s.to_string()));
assert_eq!(
*out.borrow(),
vec!["Server state: paused", "Server 'rsrv'"],
"C prints no report outside the running phase"
);
}
#[test]
fn dbsr_prints_the_state_then_each_layer() {
let out = std::cell::RefCell::new(Vec::<String>::new());
let sink = |s: &str| out.borrow_mut().push(s.to_string());
let mut reg = Registry {
servers: vec![DbServer {
name: "rsrv",
report: Some(Box::new(|level, o: &dyn Fn(&str)| {
o("Channel Access Server");
if level >= 1 {
o(" detail");
}
})),
}],
state: ServerState::Registering,
};
render(®, 0, &sink);
assert_eq!(
*out.borrow(),
vec!["Server state: registering", "Server 'rsrv'"],
"a layer that is not running yet contributes no report"
);
out.borrow_mut().clear();
reg.state = ServerState::Running;
render(®, 0, &sink);
assert_eq!(
*out.borrow(),
vec![
"Server state: running",
"Server 'rsrv'",
"Channel Access Server"
]
);
out.borrow_mut().clear();
render(®, 1, &sink);
assert_eq!(
*out.borrow(),
vec![
"Server state: running",
"Server 'rsrv'",
"Channel Access Server",
" detail"
],
"the interest level reaches the layer's own report"
);
}
#[test]
fn registration_gates() {
let mut reg = Registry {
servers: Vec::new(),
state: ServerState::Registering,
};
assert!(admit(
&mut reg,
DbServer {
name: "rsrv",
report: None
}
));
assert!(
!admit(
&mut reg,
DbServer {
name: "rsrv",
report: None
}
),
"a duplicate name is refused"
);
assert!(
!admit(
&mut reg,
DbServer {
name: "two words",
report: None
}
),
"a name with a space is refused"
);
reg.state = ServerState::Running;
assert!(
!admit(
&mut reg,
DbServer {
name: "qsrv2",
report: None
}
),
"registration closes once the set is running"
);
}
}