use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use epics_base_rs::server::database::{LinkSet, PvDatabase};
use epics_base_rs::server::record::ParsedLink;
use epics_base_rs::server::records::ai::AiRecord;
use epics_base_rs::server::records::ao::AoRecord;
use epics_base_rs::server::records::calc::CalcRecord;
use epics_base_rs::types::EpicsValue;
struct CountingLset {
cache: parking_lot::Mutex<HashMap<String, EpicsValue>>,
opened: parking_lot::Mutex<Vec<String>>,
network_reads: Arc<AtomicUsize>,
}
impl CountingLset {
fn new(reads: &Arc<AtomicUsize>) -> Arc<Self> {
Arc::new(Self {
cache: parking_lot::Mutex::new(HashMap::new()),
opened: parking_lot::Mutex::new(Vec::new()),
network_reads: Arc::clone(reads),
})
}
fn opened_names(&self) -> Vec<String> {
let mut v = self.opened.lock().clone();
v.sort();
v
}
}
#[async_trait::async_trait]
impl LinkSet for CountingLset {
fn is_connected(&self, name: &str) -> bool {
self.cache.lock().contains_key(name)
}
async fn get_value(&self, name: &str) -> Option<EpicsValue> {
self.network_reads.fetch_add(1, Ordering::SeqCst);
self.connect_link(name).await;
self.cache.lock().get(name).cloned()
}
fn get_cached_value(&self, name: &str) -> Option<EpicsValue> {
self.cache.lock().get(name).cloned()
}
async fn connect_link(&self, name: &str) {
self.opened.lock().push(name.to_string());
self.cache
.lock()
.insert(name.to_string(), EpicsValue::Double(12.5));
}
}
async fn add_ai(db: &PvDatabase, name: &str, inp: &str) {
db.add_record(name, Box::new(AiRecord::new(0.0)))
.await
.unwrap();
let rec = db.get_record(name).expect("just added");
let mut inst = rec.write();
inst.put_common_field("INP", EpicsValue::String(inp.into()))
.unwrap();
inst.common.udf = 0;
}
async fn add_ao(db: &PvDatabase, name: &str, out: &str) {
db.add_record(name, Box::new(AoRecord::new(0.0)))
.await
.unwrap();
let rec = db.get_record(name).expect("just added");
let mut inst = rec.write();
inst.put_common_field("OUT", EpicsValue::String(out.into()))
.unwrap();
inst.common.udf = 0;
}
async fn process(db: &PvDatabase, name: &str) {
let mut visited = HashSet::new();
db.process_record_with_links(name, &mut visited, 0)
.await
.unwrap();
}
async fn val_of(db: &PvDatabase, name: &str) -> Option<f64> {
let rec = db.get_record(name).expect("record exists");
let inst = rec.read();
inst.record.val().and_then(|v| v.to_f64())
}
#[epics_macros_rs::epics_test]
async fn non_cp_external_inp_link_is_opened_at_init_not_on_first_scan() {
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
let db = PvDatabase::new();
db.register_link_set("ca", lset.clone()).await;
add_ai(&db, "AI_PLAIN", "ca://REMOTE:IN").await;
let staged = db.setup_external_link_opens().await;
assert_eq!(staged, 1, "the plain external INP link must be staged");
db.sync_external_link_puts().await;
assert_eq!(
lset.opened_names(),
vec!["REMOTE:IN".to_string()],
"the work owner opened the link at init, before any record scan"
);
process(&db, "AI_PLAIN").await;
assert_eq!(
val_of(&db, "AI_PLAIN").await,
Some(12.5),
"the first scan of a link opened at init reads a warm cache"
);
assert_eq!(
reads.load(Ordering::SeqCst),
0,
"the record thread must never reach the network-capable get_value"
);
}
#[epics_macros_rs::epics_test]
async fn out_only_external_link_is_opened_at_init() {
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
let db = PvDatabase::new();
db.register_link_set("ca", lset.clone()).await;
add_ao(&db, "AO_PLAIN", "ca://REMOTE:OUT").await;
let staged = db.setup_external_link_opens().await;
assert_eq!(staged, 1, "an OUT-only external link must be staged");
db.sync_external_link_puts().await;
assert_eq!(
lset.opened_names(),
vec!["REMOTE:OUT".to_string()],
"C's open-at-init is direction-agnostic"
);
}
#[epics_macros_rs::epics_test]
async fn record_specific_input_link_fields_are_opened_at_init() {
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
let db = PvDatabase::new();
db.register_link_set("ca", lset.clone()).await;
db.add_record("CALC1", Box::new(CalcRecord::new("A+B")))
.await
.unwrap();
{
let rec = db.get_record("CALC1").expect("just added");
let mut inst = rec.write();
inst.record
.put_field("INPA", EpicsValue::String("ca://REMOTE:A".into()))
.unwrap();
inst.record
.put_field("INPB", EpicsValue::String("ca://REMOTE:B".into()))
.unwrap();
}
let staged = db.setup_external_link_opens().await;
assert_eq!(staged, 2, "both INPA and INPB must be staged");
db.sync_external_link_puts().await;
assert_eq!(
lset.opened_names(),
vec!["REMOTE:A".to_string(), "REMOTE:B".to_string()],
);
}
#[epics_macros_rs::epics_test]
async fn local_and_constant_links_stage_nothing() {
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
let db = PvDatabase::new();
db.register_link_set("ca", lset.clone()).await;
db.add_record("LOCAL_TGT", Box::new(AiRecord::new(1.0)))
.await
.unwrap();
add_ai(&db, "AI_LOCAL", "LOCAL_TGT").await; add_ai(&db, "AI_CONST", "3.5").await;
let staged = db.setup_external_link_opens().await;
assert_eq!(staged, 0, "no external link exists in this database");
db.sync_external_link_puts().await;
assert!(
lset.opened_names().is_empty(),
"a local or constant link must not open an external channel, got {:?}",
lset.opened_names()
);
assert_eq!(db.external_link_opens_completed(), 0);
}
#[epics_macros_rs::epics_test]
async fn plain_non_local_db_link_converts_and_opens_at_init() {
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
let db = PvDatabase::new();
db.register_link_set("ca", lset.clone()).await;
add_ai(&db, "AI_NONLOCAL", "OTHER:PV").await;
assert_eq!(
db.initialize_link_locality().await,
1,
"the non-local plain Db link must be converted"
);
let inp = db
.get_record("AI_NONLOCAL")
.unwrap()
.read()
.parsed_inp
.clone();
match inp {
ParsedLink::Ca(ca) => assert_eq!(ca.pv, "OTHER:PV"),
other => panic!("expected the link to become Ca, got {other:?}"),
}
let staged = db.setup_external_link_opens().await;
assert_eq!(staged, 1, "the converted link must be opened at init");
db.sync_external_link_puts().await;
assert_eq!(lset.opened_names(), vec!["OTHER:PV".to_string()]);
process(&db, "AI_NONLOCAL").await;
assert_eq!(val_of(&db, "AI_NONLOCAL").await, Some(12.5));
assert_eq!(
reads.load(Ordering::SeqCst),
0,
"the record thread must never reach the network-capable get_value"
);
}
#[epics_macros_rs::epics_test]
async fn plain_non_local_db_link_keeps_its_field_in_the_channel_name() {
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
let db = PvDatabase::new();
db.register_link_set("ca", lset.clone()).await;
add_ai(&db, "AI_NONLOCAL_FLD", "OTHER:PV.SEVR").await;
assert_eq!(db.initialize_link_locality().await, 1);
db.setup_external_link_opens().await;
db.sync_external_link_puts().await;
assert_eq!(lset.opened_names(), vec!["OTHER:PV.SEVR".to_string()]);
}
#[epics_macros_rs::epics_test]
async fn local_db_link_is_not_converted() {
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
let db = PvDatabase::new();
db.register_link_set("ca", lset.clone()).await;
db.add_record("LOCAL_SRC", Box::new(AiRecord::new(7.0)))
.await
.unwrap();
add_ai(&db, "AI_LOCAL_DB", "LOCAL_SRC").await;
assert_eq!(
db.initialize_link_locality().await,
0,
"a local Db link must not be converted"
);
let inp = db
.get_record("AI_LOCAL_DB")
.unwrap()
.read()
.parsed_inp
.clone();
assert!(
matches!(inp, ParsedLink::Db(_)),
"a local Db link must stay a Db link, got {inp:?}"
);
assert_eq!(db.setup_external_link_opens().await, 0);
db.sync_external_link_puts().await;
assert!(lset.opened_names().is_empty());
}
#[epics_macros_rs::epics_test]
async fn non_local_cp_link_keeps_its_conversion_and_its_external_trigger() {
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
let db = PvDatabase::new();
db.register_link_set("ca", lset.clone()).await;
add_ai(&db, "AI_CP_NONLOCAL", "OTHER:CP CP").await;
db.initialize_link_locality().await;
db.setup_cp_links().await;
let inp = db
.get_record("AI_CP_NONLOCAL")
.unwrap()
.read()
.parsed_inp
.clone();
match inp {
ParsedLink::Ca(ca) => assert_eq!(ca.pv, "OTHER:CP"),
other => panic!("a non-local CP link must be Ca, got {other:?}"),
}
assert!(
db.external_cp_pv_names()
.await
.contains(&"OTHER:CP".to_string()),
"the external CP trigger must still be registered"
);
}
#[epics_macros_rs::epics_test]
async fn a_record_added_after_init_does_not_un_convert_the_link() {
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
let db = PvDatabase::new();
db.register_link_set("ca", lset.clone()).await;
add_ai(&db, "AI_EARLY", "LATE_SRC").await;
assert_eq!(db.initialize_link_locality().await, 1);
db.add_record("LATE_SRC", Box::new(AiRecord::new(3.0)))
.await
.unwrap();
let inp = db.get_record("AI_EARLY").unwrap().read().parsed_inp.clone();
match inp {
ParsedLink::Ca(ca) => assert_eq!(
ca.pv, "LATE_SRC",
"the link stays external; C never un-converts"
),
other => panic!("a converted link must stay Ca, got {other:?}"),
}
}
#[epics_macros_rs::epics_test]
async fn constant_link_is_not_converted() {
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
let db = PvDatabase::new();
db.register_link_set("ca", lset.clone()).await;
add_ai(&db, "AI_CONST2", "3.5").await;
assert_eq!(db.initialize_link_locality().await, 0);
assert_eq!(db.setup_external_link_opens().await, 0);
db.sync_external_link_puts().await;
assert!(lset.opened_names().is_empty());
}
#[epics_macros_rs::epics_test]
async fn external_forward_link_is_opened_at_init() {
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
let db = PvDatabase::new();
db.register_link_set("ca", lset.clone()).await;
db.add_record("AI_FWD", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
{
let rec = db.get_record("AI_FWD").expect("just added");
let mut inst = rec.write();
inst.put_common_field("FLNK", EpicsValue::String("ca://REMOTE:FWD.PROC".into()))
.unwrap();
}
let staged = db.setup_external_link_opens().await;
assert_eq!(staged, 1, "an external FLNK must be staged at init");
db.sync_external_link_puts().await;
assert_eq!(
lset.opened_names(),
vec!["REMOTE:FWD.PROC".to_string()],
"C's open-at-init covers DBF_FWDLINK"
);
}
#[epics_macros_rs::epics_test]
async fn local_forward_link_stages_nothing() {
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
let db = PvDatabase::new();
db.register_link_set("ca", lset.clone()).await;
db.add_record("FWD_TGT", Box::new(AiRecord::new(1.0)))
.await
.unwrap();
db.add_record("AI_FWD_LOCAL", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
{
let rec = db.get_record("AI_FWD_LOCAL").expect("just added");
let mut inst = rec.write();
inst.put_common_field("FLNK", EpicsValue::String("FWD_TGT".into()))
.unwrap();
}
let staged = db.setup_external_link_opens().await;
assert_eq!(staged, 0, "a local FLNK resolves as a DB link");
db.sync_external_link_puts().await;
assert!(
lset.opened_names().is_empty(),
"a local FLNK must not open an external channel, got {:?}",
lset.opened_names()
);
}
#[epics_macros_rs::epics_test]
async fn forward_link_never_registers_a_cp_holder() {
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
let db = PvDatabase::new();
db.register_link_set("ca", lset.clone()).await;
db.add_record("CP_SRC", Box::new(AiRecord::new(1.0)))
.await
.unwrap();
db.add_record("AI_FWD_CP", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
{
let rec = db.get_record("AI_FWD_CP").expect("just added");
let mut inst = rec.write();
inst.put_common_field("FLNK", EpicsValue::String("CP_SRC CP".into()))
.unwrap();
}
db.setup_cp_links().await;
assert!(
db.get_cp_targets("CP_SRC").is_empty(),
"a forward link's CP modifier is masked off, so no holder edge exists"
);
}
#[epics_macros_rs::epics_test]
async fn external_cp_link_already_warmed_is_not_staged_again() {
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
let db = PvDatabase::new();
db.register_link_set("ca", lset.clone()).await;
add_ai(&db, "AI_CP", "ca://REMOTE:CP CP").await;
db.setup_cp_links().await; let staged = db.setup_external_link_opens().await;
assert_eq!(
staged, 0,
"the CP pass already staged this link's open; the init pass must not \
stage a second one"
);
db.sync_external_link_puts().await;
assert_eq!(
lset.opened_names(),
vec!["REMOTE:CP".to_string()],
"exactly one connect, from exactly one staged open"
);
assert_eq!(
db.external_link_opens_completed(),
1,
"the work owner completed one open, not two"
);
}
#[epics_macros_rs::epics_test]
async fn no_link_set_registered_stages_nothing_and_leaves_the_link_openable() {
let db = PvDatabase::new();
add_ai(&db, "AI_NOLSET", "ca://REMOTE:IN").await;
assert_eq!(
db.setup_external_link_opens().await,
0,
"no lset addresses this link yet, so nothing may be staged"
);
let reads = Arc::new(AtomicUsize::new(0));
let lset = CountingLset::new(&reads);
db.register_link_set("ca", lset.clone()).await;
assert_eq!(db.setup_external_link_opens().await, 1);
db.sync_external_link_puts().await;
assert_eq!(lset.opened_names(), vec!["REMOTE:IN".to_string()]);
}