pub mod db_access;
mod field_io;
pub mod filters;
mod link_put_queue;
mod link_set;
mod links;
mod processing;
mod record_lock;
mod scan_index;
mod snapshot;
pub use link_set::{
DynLinkSet, LinkDbfType, LinkMetadata, LinkPutOp, LinkSet, LinkSetRegistry, PutAdmission,
RemoteAlarm,
};
pub use processing::{AsyncDbHandle, AsyncToken};
pub use record_lock::{ManyRecordWriteGuard, RecordWriteGuard};
use crate::error::{CaError, CaResult};
use arc_swap::{ArcSwap, ArcSwapOption};
use snapshot::SnapshotCell;
use std::collections::{BTreeSet, HashMap};
use std::sync::Arc;
use crate::server::pv::ProcessVariable;
use crate::server::record::{Record, RecordInstance, ScanList};
use crate::types::EpicsValue;
#[derive(Default, Debug, Clone)]
pub struct RecordLoad {
pub common_fields: Vec<(String, EpicsValue)>,
pub info_tags: Vec<(String, String)>,
}
impl RecordLoad {
pub fn from_common_fields(common_fields: Vec<(String, EpicsValue)>) -> Self {
Self {
common_fields,
info_tags: Vec::new(),
}
}
}
pub fn parse_pv_name(name: &str) -> (&str, &str) {
match name.rsplit_once('.') {
Some((base, field)) => (base, field),
None => (name, "VAL"),
}
}
pub fn is_value_field(field: &str) -> bool {
field.eq_ignore_ascii_case("VAL")
}
fn apply_timestamp(common: &mut super::record::CommonFields, _is_soft: bool) {
common.time = crate::server::recgbl::get_time_stamp(common.tse, common.time);
}
pub enum PvEntry {
Simple(Arc<ProcessVariable>),
Record(Arc<parking_lot::RwLock<RecordInstance>>),
}
pub type ExternalPvResolver = Arc<dyn Fn(&str) -> Option<EpicsValue> + Send + Sync>;
pub type SearchResolver = Arc<
dyn Fn(
String,
Option<std::net::SocketAddr>,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = bool> + Send>>
+ Send
+ Sync,
>;
pub type ExistenceGate = Arc<
dyn Fn(
String,
Option<std::net::SocketAddr>,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = bool> + Send>>
+ Send
+ Sync,
>;
#[derive(Clone, Debug)]
pub struct CpTarget {
pub record: String,
pub passive_only: bool,
}
struct ScanIndex {
buckets: [crate::runtime::sync::PriorityInheritanceMutex<BTreeSet<(i16, u64, String)>>;
ScanList::COUNT],
}
impl ScanIndex {
fn new() -> Self {
Self {
buckets: std::array::from_fn(|_| {
crate::runtime::sync::PriorityInheritanceMutex::new(BTreeSet::new())
}),
}
}
fn bucket(
&self,
list: ScanList,
) -> &crate::runtime::sync::PriorityInheritanceMutex<BTreeSet<(i16, u64, String)>> {
&self.buckets[list.slot()]
}
}
struct PvDatabaseInner {
simple_pvs:
crate::runtime::sync::PriorityInheritanceMutex<HashMap<String, Arc<ProcessVariable>>>,
records: parking_lot::RwLock<HashMap<String, Arc<parking_lot::RwLock<RecordInstance>>>>,
scan_index: ScanIndex,
load_order: SnapshotCell<HashMap<String, u64>>,
load_order_counter: std::sync::atomic::AtomicU64,
cp_links: SnapshotCell<HashMap<String, Vec<CpTarget>>>,
external_cp_links: SnapshotCell<HashMap<String, Vec<CpTarget>>>,
aliases: parking_lot::RwLock<HashMap<String, String>>,
registration_mutex: crate::runtime::sync::PriorityInheritanceMutex<()>,
init_phase: std::sync::Mutex<DbInitPhase>,
record_init_waiting: std::sync::Mutex<HashMap<String, Vec<RecordInit>>>,
after_ioc_running: std::sync::Mutex<Vec<String>>,
external_resolver: ArcSwapOption<ExternalPvResolver>,
search_resolver: ArcSwapOption<SearchResolver>,
existence_gate: ArcSwapOption<ExistenceGate>,
link_sets: SnapshotCell<link_set::LinkSetRegistry>,
link_puts: Arc<link_put_queue::LinkPutQueue>,
scan_started: std::sync::atomic::AtomicBool,
pini_done: std::sync::atomic::AtomicBool,
pini_notify: tokio::sync::Notify,
record_locks: record_lock::RecordLockRegistry,
subroutine_registry: ArcSwap<HashMap<String, Arc<crate::server::record::SubroutineFn>>>,
breaktable_registry: SnapshotCell<crate::server::cvt_bpt::BreakTableRegistry>,
}
#[derive(Clone)]
pub struct PvDatabase {
inner: Arc<PvDatabaseInner>,
}
type RecordInit = std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send + 'static>>;
enum DbInitPhase {
Unloaded,
Loading(Vec<RecordInit>),
Running,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct IocAlreadyInitialized;
impl std::fmt::Display for IocAlreadyInitialized {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("IOC already initialized - No new records can be added")
}
}
impl std::error::Error for IocAlreadyInitialized {}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub(crate) enum SelmKind {
FanoutSeq,
Dfanout,
}
#[derive(Clone, Debug, Default)]
pub(crate) struct SelmResult {
pub indices: Vec<usize>,
pub alarm: Option<(u16, crate::server::record::AlarmSeverity)>,
}
pub(crate) fn dbr_ushort_cast(value: &EpicsValue) -> u16 {
match value.convert_to(crate::types::DbFieldType::UShort) {
EpicsValue::UShort(v) => v,
EpicsValue::UShortArray(v) => v.first().copied().unwrap_or(0),
_ => 0,
}
}
pub(crate) fn select_link_indices_ex(
kind: SelmKind,
selm: i16,
seln: u16,
offs: i16,
shft: i16,
count: usize,
) -> SelmResult {
use crate::server::recgbl::alarm_status::SOFT_ALARM;
use crate::server::record::AlarmSeverity;
let invalid = || SelmResult {
indices: Vec::new(),
alarm: Some((SOFT_ALARM, AlarmSeverity::Invalid)),
};
let ok = |indices: Vec<usize>| SelmResult {
indices,
alarm: None,
};
match selm {
0 => ok((0..count).collect()),
1 => match kind {
SelmKind::FanoutSeq => {
let i = seln as i32 + offs as i32;
if i < 0 || i >= count as i32 {
invalid()
} else {
ok(vec![i as usize])
}
}
SelmKind::Dfanout => {
let seln_i = seln as i32;
if seln_i > count as i32 {
invalid()
} else if seln_i == 0 {
ok(Vec::new())
} else {
ok(vec![(seln_i - 1) as usize])
}
}
},
2 => {
let mask: u32 = match kind {
SelmKind::FanoutSeq => {
if !(-15..=15).contains(&shft) {
return invalid();
}
let raw = seln as u32;
if shft >= 0 {
raw >> shft
} else {
raw << (-shft)
}
}
SelmKind::Dfanout => seln as u32,
};
ok((0..count).filter(|i| mask & (1 << i) != 0).collect())
}
_ => invalid(),
}
}
impl PvDatabase {
pub fn new() -> Self {
Self {
inner: Arc::new(PvDatabaseInner {
simple_pvs: crate::runtime::sync::PriorityInheritanceMutex::new(HashMap::new()),
external_resolver: ArcSwapOption::empty(),
search_resolver: ArcSwapOption::empty(),
existence_gate: ArcSwapOption::empty(),
link_sets: SnapshotCell::new(link_set::LinkSetRegistry::new()),
link_puts: Arc::new(link_put_queue::LinkPutQueue::default()),
records: parking_lot::RwLock::new(HashMap::new()),
scan_index: ScanIndex::new(),
load_order: SnapshotCell::new(HashMap::new()),
load_order_counter: std::sync::atomic::AtomicU64::new(0),
cp_links: SnapshotCell::new(HashMap::new()),
external_cp_links: SnapshotCell::new(HashMap::new()),
aliases: parking_lot::RwLock::new(HashMap::new()),
registration_mutex: crate::runtime::sync::PriorityInheritanceMutex::new(()),
init_phase: std::sync::Mutex::new(DbInitPhase::Unloaded),
record_init_waiting: std::sync::Mutex::new(HashMap::new()),
after_ioc_running: std::sync::Mutex::new(Vec::new()),
scan_started: std::sync::atomic::AtomicBool::new(false),
pini_done: std::sync::atomic::AtomicBool::new(false),
pini_notify: tokio::sync::Notify::new(),
record_locks: record_lock::RecordLockRegistry::default(),
subroutine_registry: ArcSwap::from_pointee(HashMap::new()),
breaktable_registry: SnapshotCell::new(
crate::server::cvt_bpt::BreakTableRegistry::new(),
),
}),
}
}
pub async fn add_breaktables(
&self,
tables: Vec<crate::server::cvt_bpt::BrkTable>,
) -> Arc<crate::server::cvt_bpt::BreakTableRegistry> {
let _gate = self.inner.registration_mutex.lock();
if tables.is_empty() {
return self.inner.breaktable_registry.load_full();
}
let snapshot = self.inner.breaktable_registry.update(|next| {
for table in tables {
next.insert(table);
}
});
let instances: Vec<_> = self.inner.records.read().values().cloned().collect();
for inst in instances {
inst.write()
.record
.install_breaktable_registry(snapshot.clone());
}
snapshot
}
pub async fn install_subroutine_registry(
&self,
registry: HashMap<String, Arc<crate::server::record::SubroutineFn>>,
) {
self.inner.subroutine_registry.store(Arc::new(registry));
}
pub(crate) fn find_subroutine_named(
&self,
name: &str,
) -> Option<Arc<crate::server::record::SubroutineFn>> {
self.inner.subroutine_registry.load().get(name).cloned()
}
pub fn try_claim_scan_start(&self) -> bool {
self.inner
.scan_started
.compare_exchange(
false,
true,
std::sync::atomic::Ordering::AcqRel,
std::sync::atomic::Ordering::Acquire,
)
.is_ok()
}
pub fn mark_pini_done(&self) {
self.inner
.pini_done
.store(true, std::sync::atomic::Ordering::Release);
self.inner.pini_notify.notify_waiters();
}
pub fn pini_done(&self) -> bool {
self.inner
.pini_done
.load(std::sync::atomic::Ordering::Acquire)
}
pub async fn wait_for_pini(&self) {
if self
.inner
.pini_done
.load(std::sync::atomic::Ordering::Acquire)
{
return;
}
let notified = self.inner.pini_notify.notified();
if self
.inner
.pini_done
.load(std::sync::atomic::Ordering::Acquire)
{
return;
}
notified.await;
}
pub async fn set_search_resolver(&self, resolver: SearchResolver) {
self.inner.search_resolver.store(Some(Arc::new(resolver)));
}
pub async fn clear_search_resolver(&self) {
self.inner.search_resolver.store(None);
}
pub async fn set_existence_gate(&self, gate: ExistenceGate) {
self.inner.existence_gate.store(Some(Arc::new(gate)));
}
pub async fn clear_existence_gate(&self) {
self.inner.existence_gate.store(None);
}
async fn simple_pv_gate_denies(&self, name: &str, peer: Option<std::net::SocketAddr>) -> bool {
let Some(gate) = self.inner.existence_gate.load_full() else {
return false;
};
let gate = (*gate).clone();
let record_path = filters::split_channel_name(name).record_path;
let known = self
.inner
.simple_pvs
.lock()
.contains_key(record_path.as_str());
if !known {
return false;
}
!gate(record_path, peer).await
}
pub async fn set_external_resolver(&self, resolver: ExternalPvResolver) {
self.inner.external_resolver.store(Some(Arc::new(resolver)));
}
pub async fn register_link_set(&self, scheme: &str, lset: link_set::DynLinkSet) {
self.inner.link_sets.update(|r| r.register(scheme, lset));
}
pub async fn link_set(&self, scheme: &str) -> Option<link_set::DynLinkSet> {
self.inner.link_sets.load().get(scheme)
}
pub async fn registered_link_schemes(&self) -> Vec<String> {
let mut s = self.inner.link_sets.load().schemes();
s.sort();
s
}
pub async fn wait_for_external_links(&self, timeout: std::time::Duration) -> (usize, usize) {
let targets = self.external_link_targets().await;
let total = targets.len();
if total == 0 {
return (0, 0);
}
let deadline = std::time::Instant::now() + timeout;
loop {
let mut connected = 0usize;
for (lset, name) in &targets {
if lset.is_connected(name) {
connected += 1;
}
}
if connected == total {
return (connected, total);
}
if std::time::Instant::now() >= deadline {
return (connected, total);
}
crate::runtime::task::sleep(std::time::Duration::from_millis(100)).await;
}
}
async fn external_link_targets(&self) -> Vec<(link_set::DynLinkSet, String)> {
let Some(ca_lset) = self.inner.link_sets.load().get("ca") else {
return Vec::new();
};
let mut targets: Vec<(link_set::DynLinkSet, String)> = Vec::new();
for n in ca_lset.link_names() {
if self.has_name_no_resolve(&n) {
targets.push((ca_lset.clone(), n));
}
}
targets
}
pub async fn unconnected_external_links(&self) -> Vec<String> {
let mut names = Vec::new();
for (lset, name) in self.external_link_targets().await {
if !lset.is_connected(&name) {
names.push(name);
}
}
names
}
pub fn record_link_fields(
&self,
record_name: &str,
) -> Vec<(String, String, crate::server::record::ParsedLink)> {
let rec = match self.get_record(record_name) {
Some(r) => r,
None => return Vec::new(),
};
let inst = rec.read();
let mut out = Vec::new();
let push = |field: &str,
raw: &str,
ftype: crate::server::record::LinkFieldType,
out: &mut Vec<_>| {
if raw.is_empty() {
return;
}
let parsed = crate::server::record::parse_link_field(raw, ftype);
if !matches!(parsed, crate::server::record::ParsedLink::None) {
out.push((field.to_string(), raw.to_string(), parsed));
}
};
use crate::server::record::LinkFieldType;
for (field, ftype) in crate::server::record::record_instance::COMMON_LINK_FIELDS {
let Some(raw) = inst.common_link_text(field) else {
continue;
};
push(field, raw, ftype, &mut out);
}
let mut field_names: Vec<&str> = inst
.record
.multi_input_links()
.iter()
.map(|(lf, _vf)| *lf)
.collect();
field_names.extend_from_slice(crate::server::database::links::CP_INPUT_LINK_FIELDS);
for field in field_names {
if let Some(EpicsValue::String(s)) = inst.record.get_field(field) {
push(field, &s.as_str_lossy(), LinkFieldType::In, &mut out);
}
}
drop(inst);
for entry in &mut out {
entry.2 = self.db_init_link_locality(std::mem::replace(
&mut entry.2,
crate::server::record::ParsedLink::None,
));
}
out
}
pub(crate) fn resolve_external_pv(&self, name: &str) -> Option<EpicsValue> {
let (target, body) = Self::split_external_link_name(name);
for lset in link_put_queue::resolve_lsets(&self.inner, &target) {
if let Some(v) = lset.get_cached_value(body) {
return Some(v);
}
}
self.stage_external_link_open_by_name(name);
let resolver = self
.inner
.external_resolver
.load_full()
.map(|r| (*r).clone());
match resolver {
Some(r) => r(name),
None => None,
}
}
fn split_external_link_name(name: &str) -> (link_put_queue::LinkTarget, &str) {
if let Some(rest) = name.strip_prefix("pva://") {
(link_put_queue::LinkTarget::Scheme("pva".to_string()), rest)
} else if let Some(rest) = name.strip_prefix("ca://") {
(link_put_queue::LinkTarget::Scheme("ca".to_string()), rest)
} else {
(link_put_queue::LinkTarget::Any, name)
}
}
pub(crate) fn stage_external_link_open_by_name(&self, name: &str) -> bool {
let (target, body) = Self::split_external_link_name(name);
if link_put_queue::resolve_lsets(&self.inner, &target).is_empty() {
return false;
}
self.stage_external_link_open(target, body)
}
fn stage_external_link_open(&self, target: link_put_queue::LinkTarget, name: &str) -> bool {
self.inner
.link_puts
.ensure_owner(std::sync::Arc::downgrade(&self.inner));
self.inner.link_puts.stage_open(link_put_queue::LinkKey {
target,
name: name.to_string(),
})
}
pub fn external_link_opens_completed(&self) -> u64 {
self.inner.link_puts.opened_count()
}
pub async fn add_pv(&self, name: &str, initial: EpicsValue) -> CaResult<()> {
let _gate = self.inner.registration_mutex.lock();
self.check_name_free(name)?;
let pv = Arc::new(ProcessVariable::new(name.to_string(), initial));
self.inner.simple_pvs.lock().insert(name.to_string(), pv);
Ok(())
}
pub async fn add_pv_with_hook(
&self,
name: &str,
initial: EpicsValue,
hook: crate::server::pv::WriteHook,
) -> CaResult<()> {
self.add_pv_with_hooks(name, initial, hook, None).await
}
pub async fn add_pv_with_hooks(
&self,
name: &str,
initial: EpicsValue,
write_hook: crate::server::pv::WriteHook,
access_hook: Option<crate::server::pv::AccessHook>,
) -> CaResult<()> {
self.add_pv_with_hooks_full(name, initial, write_hook, access_hook, None)
.await
}
pub async fn add_pv_with_hooks_full(
&self,
name: &str,
initial: EpicsValue,
write_hook: crate::server::pv::WriteHook,
access_hook: Option<crate::server::pv::AccessHook>,
read_hook: Option<crate::server::pv::ReadHook>,
) -> CaResult<()> {
let _gate = self.inner.registration_mutex.lock();
self.check_name_free(name)?;
let pv = Arc::new(ProcessVariable::new(name.to_string(), initial));
pv.set_write_hook(write_hook);
if let Some(access) = access_hook {
pv.set_access_hook(access);
}
if let Some(read) = read_hook {
pv.set_read_hook(read);
}
self.inner.simple_pvs.lock().insert(name.to_string(), pv);
Ok(())
}
pub async fn remove_simple_pv(&self, name: &str) -> Option<Arc<ProcessVariable>> {
let _gate = self.inner.registration_mutex.lock();
self.inner.simple_pvs.lock().remove(name)
}
#[must_use = "C refuses a load after iocInit (dbReadCOM, dbLexRoutines.c:236); \
the refusal must be reported and no record created"]
pub fn begin_load(&self) -> Result<(), IocAlreadyInitialized> {
let mut phase = self.inner.init_phase.lock().unwrap();
match *phase {
DbInitPhase::Unloaded => {
*phase = DbInitPhase::Loading(Vec::new());
Ok(())
}
DbInitPhase::Loading(_) => Ok(()),
DbInitPhase::Running => Err(IocAlreadyInitialized),
}
}
pub(crate) fn schedule_record_init(
&self,
record: &str,
init: impl std::future::Future<Output = ()> + Send + 'static,
) {
if !self.inner.records.read().contains_key(record) {
self.inner
.record_init_waiting
.lock()
.unwrap()
.entry(record.to_string())
.or_default()
.push(Box::pin(init));
return;
}
self.dispatch_record_init(Box::pin(init));
}
fn dispatch_record_init(&self, init: RecordInit) {
let mut phase = self.inner.init_phase.lock().unwrap();
match &mut *phase {
DbInitPhase::Loading(queued) => queued.push(init),
DbInitPhase::Unloaded | DbInitPhase::Running => {
drop(phase);
crate::runtime::task::spawn(init);
}
}
}
fn release_record_inits(&self, record: &str) {
let parked = self
.inner
.record_init_waiting
.lock()
.unwrap()
.remove(record);
for init in parked.into_iter().flatten() {
self.dispatch_record_init(init);
}
}
pub async fn ioc_init(&self) {
let owed = {
let mut phase = self.inner.init_phase.lock().unwrap();
match std::mem::replace(&mut *phase, DbInitPhase::Running) {
DbInitPhase::Loading(queued) => queued,
DbInitPhase::Unloaded => return,
DbInitPhase::Running => return,
}
};
for init in owed {
init.await;
}
}
pub async fn add_record(&self, name: &str, record: Box<dyn Record>) -> CaResult<()> {
self.add_loaded_record(name, record, RecordLoad::default())
.await
}
pub async fn add_loaded_record(
&self,
name: &str,
record: Box<dyn Record>,
load: RecordLoad,
) -> CaResult<()> {
let _gate = self.inner.registration_mutex.lock();
self.check_name_free(name)?;
let mut instance = RecordInstance::new_boxed(name.to_string(), record);
instance
.record
.set_async_context(name.to_string(), self.async_handle());
{
let snapshot = self.inner.breaktable_registry.load_full();
if !snapshot.is_empty() {
instance.record.install_breaktable_registry(snapshot);
}
}
for (field, value) in load.common_fields {
if let Err(e) = instance.put_common_field_db_load(&field, value) {
eprintln!("put_common_field({field}) failed for {name}: {e}");
}
}
for (key, value) in &load.info_tags {
instance.set_info(key, value);
}
instance.run_init_passes(name);
super::database::processing::seed_constant_links(&mut instance);
let scan = instance.common.scan;
let phas = instance.common.phas;
self.inner.records.write().insert(
name.to_string(),
Arc::new(parking_lot::RwLock::new(instance)),
);
self.release_record_inits(name);
let seq = self
.inner
.load_order_counter
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
self.inner.load_order.update(|m| {
m.insert(name.to_string(), seq);
});
if let Some(list) = scan.scan_list() {
self.inner
.scan_index
.bucket(list)
.lock()
.insert((phas, seq, name.to_string()));
}
Ok(())
}
fn check_name_free(&self, name: &str) -> CaResult<()> {
let kind = if self.inner.simple_pvs.lock().contains_key(name) {
Some("simple PV")
} else if self.inner.records.read().contains_key(name) {
Some("record")
} else if self.inner.aliases.read().contains_key(name) {
Some("alias")
} else {
None
};
if let Some(kind) = kind {
return Err(CaError::DbParseError {
line: 0,
column: 0,
message: format!("name '{name}' is already registered as a {kind}"),
});
}
Ok(())
}
pub async fn remove_record(&self, name: &str) -> bool {
let _gate = self.inner.registration_mutex.lock();
let removed = self.inner.records.write().remove(name);
let Some(rec_arc) = removed else {
return false;
};
let scan = {
let inst = rec_arc.read();
inst.common.scan
};
if let Some(list) = scan.scan_list() {
self.inner
.scan_index
.bucket(list)
.lock()
.retain(|(_, _, n)| n != name);
}
self.inner.load_order.update(|m| {
m.remove(name);
});
self.inner.cp_links.update(|cp| {
cp.remove(name);
for targets in cp.values_mut() {
targets.retain(|t| t.record != name);
}
});
let mut aliases = self.inner.aliases.write();
aliases.retain(|_alias, target| target != name);
true
}
async fn find_entry_no_resolve(&self, name: &str) -> Option<PvEntry> {
let record_path = filters::split_channel_name(name).record_path;
let (base, _field) = parse_pv_name(&record_path);
let simple = self
.inner
.simple_pvs
.lock()
.get(record_path.as_str())
.cloned();
if let Some(pv) = simple {
return Some(PvEntry::Simple(pv));
}
if let Some(rec) = self.inner.records.read().get(base) {
return Some(PvEntry::Record(rec.clone()));
}
if let Some(target) = self.inner.aliases.read().get(base).cloned() {
if let Some(rec) = self.inner.records.read().get(&target) {
return Some(PvEntry::Record(rec.clone()));
}
}
None
}
pub async fn add_alias(&self, alias: &str, target: &str) -> CaResult<()> {
let _gate = self.inner.registration_mutex.lock();
if !self.inner.records.read().contains_key(target) {
return Err(CaError::ChannelNotFound(format!(
"alias target '{target}' is not a registered record"
)));
}
self.check_name_free(alias)?;
self.inner
.aliases
.write()
.insert(alias.to_string(), target.to_string());
Ok(())
}
pub fn resolve_alias(&self, name: &str) -> Option<String> {
self.inner.aliases.read().get(name).cloned()
}
pub fn queue_after_ioc_running(&self, line: impl Into<String>) {
self.inner
.after_ioc_running
.lock()
.unwrap()
.push(line.into());
}
pub fn take_after_ioc_running(&self) -> Vec<String> {
std::mem::take(&mut *self.inner.after_ioc_running.lock().unwrap())
}
fn has_name_no_resolve(&self, name: &str) -> bool {
let record_path = filters::split_channel_name(name).record_path;
let (base, _) = parse_pv_name(&record_path);
if self
.inner
.simple_pvs
.lock()
.contains_key(record_path.as_str())
{
return true;
}
if self.inner.records.read().contains_key(base) {
return true;
}
if let Some(target) = self.inner.aliases.read().get(base) {
return self.inner.records.read().contains_key(target);
}
false
}
pub async fn find_entry(&self, name: &str) -> Option<PvEntry> {
self.find_entry_from(name, None).await
}
pub async fn find_entry_from(
&self,
name: &str,
peer: Option<std::net::SocketAddr>,
) -> Option<PvEntry> {
if let Some(entry) = self.find_entry_no_resolve(name).await {
if matches!(entry, PvEntry::Simple(_)) && self.simple_pv_gate_denies(name, peer).await {
return None;
}
return Some(entry);
}
let resolver = self.inner.search_resolver.load_full().map(|r| (*r).clone());
if let Some(r) = resolver {
if r(name.to_string(), peer).await {
return self.find_entry_no_resolve(name).await;
}
}
None
}
pub async fn has_name(&self, name: &str) -> bool {
self.has_name_from(name, None).await
}
pub async fn has_name_from(&self, name: &str, peer: Option<std::net::SocketAddr>) -> bool {
if self.has_name_no_resolve(name) {
if self.simple_pv_gate_denies(name, peer).await {
return false;
}
return true;
}
let resolver = self.inner.search_resolver.load_full().map(|r| (*r).clone());
if let Some(r) = resolver {
if r(name.to_string(), peer).await {
return self.has_name_no_resolve(name);
}
}
false
}
pub async fn find_pv(&self, name: &str) -> Option<Arc<ProcessVariable>> {
self.inner.simple_pvs.lock().get(name).cloned()
}
pub fn get_record(&self, name: &str) -> Option<Arc<parking_lot::RwLock<RecordInstance>>> {
if let Some(rec) = self.inner.records.read().get(name).cloned() {
return Some(rec);
}
let target = self.inner.aliases.read().get(name).cloned()?;
self.inner.records.read().get(&target).cloned()
}
pub fn get_record_no_resolve(
&self,
name: &str,
) -> Option<Arc<parking_lot::RwLock<RecordInstance>>> {
self.inner.records.read().get(name).cloned()
}
pub async fn all_record_names(&self) -> Vec<String> {
let mut names: Vec<String> = {
let records = self.inner.records.read();
records.keys().cloned().collect()
};
let load_order = self.inner.load_order.load();
names.sort_by(|a, b| {
let seq = |n: &String| load_order.get(n).copied().unwrap_or(u64::MAX);
seq(a).cmp(&seq(b)).then_with(|| a.cmp(b))
});
names
}
pub fn all_alias_names(&self) -> Vec<String> {
self.inner.aliases.read().keys().cloned().collect()
}
pub fn aliases_for_record(&self, canonical: &str) -> Vec<String> {
let aliases = self.inner.aliases.read();
let mut hits: Vec<String> = aliases
.iter()
.filter_map(|(alias, target)| {
if target == canonical {
Some(alias.clone())
} else {
None
}
})
.collect();
hits.sort();
hits
}
pub async fn all_simple_pv_names(&self) -> Vec<String> {
self.inner.simple_pvs.lock().keys().cloned().collect()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn apply_timestamp_tse_minus_one_always_overwrites_with_best_time() {
use crate::server::record::CommonFields;
use std::time::{Duration, SystemTime};
let stale = SystemTime::UNIX_EPOCH + Duration::from_secs(1_000_000);
let mut common = CommonFields::default();
common.tse = -1;
common.time = stale;
apply_timestamp(&mut common, false);
assert_ne!(
common.time, stale,
"TSE=-1 must always overwrite via generalTime BestTime, \
matching C epicsTimeGetEvent(-1) called unconditionally"
);
}
#[test]
fn apply_timestamp_tse_minus_two_preserves_device_provided_time() {
use crate::server::record::CommonFields;
use std::time::{Duration, SystemTime};
let device_time = SystemTime::UNIX_EPOCH + Duration::from_secs(2_000_000);
let mut common = CommonFields::default();
common.tse = -2;
common.time = device_time;
apply_timestamp(&mut common, false);
assert_eq!(
common.time, device_time,
"TSE=-2 (epicsTimeEventDeviceTime) must preserve device-provided time"
);
}
#[test]
fn select_link_indices_fanout_all_specified_mask() {
use crate::server::record::AlarmSeverity;
let r = select_link_indices_ex(SelmKind::FanoutSeq, 0, 0, 0, 0, 16);
assert_eq!(r.indices, (0..16).collect::<Vec<_>>());
assert!(r.alarm.is_none());
let r = select_link_indices_ex(SelmKind::FanoutSeq, 1, 0, 0, 0, 16);
assert_eq!(r.indices, vec![0]);
let r = select_link_indices_ex(SelmKind::FanoutSeq, 1, 2, 3, 0, 16);
assert_eq!(r.indices, vec![5]);
let r = select_link_indices_ex(SelmKind::FanoutSeq, 1, 20, 0, 0, 16);
assert!(r.indices.is_empty());
assert_eq!(r.alarm, Some((15, AlarmSeverity::Invalid)));
let r = select_link_indices_ex(SelmKind::FanoutSeq, 1, 0, -1, 0, 16);
assert_eq!(r.alarm, Some((15, AlarmSeverity::Invalid)));
let r = select_link_indices_ex(SelmKind::FanoutSeq, 2, 5, 0, 0, 16);
assert_eq!(r.indices, vec![0, 2]);
let r = select_link_indices_ex(SelmKind::FanoutSeq, 2, 5, 0, 1, 16);
assert_eq!(r.indices, vec![1]);
let r = select_link_indices_ex(SelmKind::FanoutSeq, 2, 5, 0, -1, 16);
assert_eq!(r.indices, vec![1, 3]);
let r = select_link_indices_ex(SelmKind::FanoutSeq, 2, 5, 0, 16, 16);
assert_eq!(r.alarm, Some((15, AlarmSeverity::Invalid)));
let r = select_link_indices_ex(SelmKind::FanoutSeq, 9, 0, 0, 0, 16);
assert_eq!(r.alarm, Some((15, AlarmSeverity::Invalid)));
}
#[test]
fn select_link_indices_dfanout_specified_is_one_based() {
use crate::server::record::AlarmSeverity;
let r = select_link_indices_ex(SelmKind::Dfanout, 1, 1, 0, 0, 16);
assert_eq!(r.indices, vec![0]);
let r = select_link_indices_ex(SelmKind::Dfanout, 1, 2, 0, 0, 16);
assert_eq!(r.indices, vec![1]);
let r = select_link_indices_ex(SelmKind::Dfanout, 1, 0, 0, 0, 16);
assert!(r.indices.is_empty());
assert!(r.alarm.is_none());
let r = select_link_indices_ex(SelmKind::Dfanout, 1, 17, 0, 0, 16);
assert_eq!(r.alarm, Some((15, AlarmSeverity::Invalid)));
let r = select_link_indices_ex(SelmKind::Dfanout, 2, 5, 0, 7, 16);
assert_eq!(r.indices, vec![0, 2]);
}
#[test]
fn seln_cast_follows_the_source_type() {
assert_eq!(dbr_ushort_cast(&EpicsValue::Long(-1)), 65535);
assert_eq!(dbr_ushort_cast(&EpicsValue::Long(65536)), 0);
assert_eq!(dbr_ushort_cast(&EpicsValue::Short(-1)), 65535);
assert_eq!(dbr_ushort_cast(&EpicsValue::Int64(-1)), 65535);
assert_eq!(dbr_ushort_cast(&EpicsValue::Long(3)), 3);
assert_eq!(
dbr_ushort_cast(&EpicsValue::Double(-1.0)),
crate::types::c_cast::f64_to_u16(-1.0)
);
assert_eq!(
dbr_ushort_cast(&EpicsValue::Double(65536.0)),
crate::types::c_cast::f64_to_u16(65536.0)
);
assert_eq!(dbr_ushort_cast(&EpicsValue::Double(3.7)), 3);
}
#[test]
fn seln_at_the_unsigned_maximum_selects_by_selm() {
use crate::server::record::AlarmSeverity;
let seln_max = 65535u16;
let r = select_link_indices_ex(SelmKind::FanoutSeq, 1, seln_max, 0, 0, 16);
assert!(r.indices.is_empty());
assert_eq!(r.alarm, Some((15, AlarmSeverity::Invalid)));
let r = select_link_indices_ex(SelmKind::FanoutSeq, 2, seln_max, 0, 0, 16);
assert_eq!(r.indices, (0..16).collect::<Vec<_>>());
let r = select_link_indices_ex(SelmKind::Dfanout, 1, seln_max, 0, 0, 16);
assert!(r.indices.is_empty());
assert_eq!(r.alarm, Some((15, AlarmSeverity::Invalid)));
}
struct DelayedConnectLset {
names: Vec<String>,
connect_at: crate::runtime::task::Instant,
}
#[async_trait::async_trait]
impl link_set::LinkSet for DelayedConnectLset {
fn is_connected(&self, _: &str) -> bool {
crate::runtime::task::Instant::now() >= self.connect_at
}
fn get_cached_value(&self, _: &str) -> Option<EpicsValue> {
None
}
async fn get_value(&self, name: &str) -> Option<EpicsValue> {
self.get_cached_value(name)
}
fn link_names(&self) -> Vec<String> {
self.names.clone()
}
}
#[epics_macros_rs::epics_test]
async fn wait_for_external_links_returns_zero_zero_when_no_lsets() {
let db = PvDatabase::new();
let (c, t) = db
.wait_for_external_links(std::time::Duration::from_millis(50))
.await;
assert_eq!((c, t), (0, 0));
}
#[epics_macros_rs::epics_test]
async fn wait_for_external_links_connected_quickly() {
let db = PvDatabase::new();
db.add_pv("pv:A", EpicsValue::Long(0)).await.unwrap();
db.add_pv("pv:B", EpicsValue::Long(0)).await.unwrap();
let lset = Arc::new(DelayedConnectLset {
names: vec!["pv:A".to_string(), "pv:B".to_string()],
connect_at: crate::runtime::task::Instant::now(),
});
db.register_link_set("ca", lset).await;
let (c, t) = db
.wait_for_external_links(std::time::Duration::from_secs(1))
.await;
assert_eq!((c, t), (2, 2));
}
#[epics_macros_rs::epics_test]
async fn wait_for_external_links_returns_partial_on_timeout() {
let db = PvDatabase::new();
db.add_pv("slow:pv", EpicsValue::Long(0)).await.unwrap();
let lset = Arc::new(DelayedConnectLset {
names: vec!["slow:pv".to_string()],
connect_at: crate::runtime::task::Instant::now() + std::time::Duration::from_secs(60),
});
db.register_link_set("ca", lset).await;
let started = crate::runtime::task::Instant::now();
let (c, t) = db
.wait_for_external_links(std::time::Duration::from_millis(250))
.await;
let elapsed = started.elapsed();
assert_eq!((c, t), (0, 1));
assert!(
elapsed >= std::time::Duration::from_millis(200),
"wait must consume at least the configured budget, got {:?}",
elapsed
);
assert!(
elapsed < std::time::Duration::from_secs(2),
"wait must not exceed the budget by much, got {:?}",
elapsed
);
}
#[epics_macros_rs::epics_test]
async fn wait_for_external_links_skips_nonlocal_targets() {
let db = PvDatabase::new();
let lset = Arc::new(DelayedConnectLset {
names: vec!["test".to_string()],
connect_at: crate::runtime::task::Instant::now() + std::time::Duration::from_secs(60),
});
db.register_link_set("ca", lset).await;
let started = crate::runtime::task::Instant::now();
let (c, t) = db
.wait_for_external_links(std::time::Duration::from_secs(10))
.await;
assert_eq!((c, t), (0, 0));
assert!(
started.elapsed() < std::time::Duration::from_secs(1),
"non-local link must not be waited on, got {:?}",
started.elapsed()
);
assert!(db.unconnected_external_links().await.is_empty());
}
#[epics_macros_rs::epics_test]
async fn alias_resolves_through_find_entry() {
let db = PvDatabase::new();
db.add_record(
"TARGET",
Box::new(crate::server::records::ai::AiRecord::new(42.0)),
)
.await
.unwrap();
db.add_alias("ALIAS_NAME", "TARGET").await.unwrap();
let via_alias = db.find_entry("ALIAS_NAME").await;
let via_target = db.find_entry("TARGET").await;
assert!(via_alias.is_some());
assert!(via_target.is_some());
assert!(db.has_name("ALIAS_NAME").await);
assert!(db.has_name("TARGET").await);
assert!(!db.has_name("NOT:THERE").await);
}
#[epics_macros_rs::epics_test]
async fn alias_target_must_exist() {
let db = PvDatabase::new();
let err = db.add_alias("DANGLING", "MISSING_TARGET").await;
assert!(err.is_err(), "alias to missing target must be rejected");
}
#[epics_macros_rs::epics_test]
async fn alias_collision_with_existing_record_rejected() {
let db = PvDatabase::new();
db.add_record(
"EXISTING",
Box::new(crate::server::records::ai::AiRecord::new(0.0)),
)
.await
.unwrap();
db.add_record(
"OTHER",
Box::new(crate::server::records::ai::AiRecord::new(0.0)),
)
.await
.unwrap();
let err = db.add_alias("EXISTING", "OTHER").await;
assert!(
err.is_err(),
"alias name colliding with record must be rejected"
);
}
#[epics_macros_rs::epics_test]
async fn get_record_resolves_alias() {
let db = PvDatabase::new();
db.add_record(
"TARGET",
Box::new(crate::server::records::ai::AiRecord::new(0.0)),
)
.await
.unwrap();
db.add_alias("ALIAS", "TARGET").await.unwrap();
let via_canonical = db.get_record("TARGET");
let via_alias = db.get_record("ALIAS");
assert!(via_canonical.is_some());
assert!(via_alias.is_some(), "get_record must resolve alias");
assert!(Arc::ptr_eq(&via_canonical.unwrap(), &via_alias.unwrap()));
}
#[epics_macros_rs::epics_test]
async fn add_record_installs_breaktable_registry_from_snapshot() {
let db = PvDatabase::new();
let ramp = crate::server::cvt_bpt::BrkTable::build(
"ramp",
&[(0.0, 0.0), (100.0, 10.0), (300.0, 30.0)],
)
.unwrap();
db.add_breaktables(vec![ramp]).await;
let mut rec = crate::server::records::ai::AiRecord::new(0.0);
rec.put_field("LINR", EpicsValue::Short(15)).unwrap(); db.add_record("AI:BPT", Box::new(rec)).await.unwrap();
let arc = db.get_record("AI:BPT").unwrap();
let mut inst = arc.write();
inst.record.put_field("RVAL", EpicsValue::Long(50)).unwrap();
inst.record.process().unwrap();
assert_eq!(inst.record.get_field("VAL"), Some(EpicsValue::Double(5.0)));
}
#[epics_macros_rs::epics_test]
async fn add_breaktables_reinstalls_registry_into_existing_records() {
let db = PvDatabase::new();
let mut rec = crate::server::records::ai::AiRecord::new(0.0);
rec.put_field("LINR", EpicsValue::Short(15)).unwrap(); db.add_record("AI:BPT", Box::new(rec)).await.unwrap();
let ramp = crate::server::cvt_bpt::BrkTable::build(
"ramp",
&[(0.0, 0.0), (100.0, 10.0), (300.0, 30.0)],
)
.unwrap();
db.add_breaktables(vec![ramp]).await;
let arc = db.get_record("AI:BPT").unwrap();
let mut inst = arc.write();
inst.record.put_field("RVAL", EpicsValue::Long(50)).unwrap();
inst.record.process().unwrap();
assert_eq!(inst.record.get_field("VAL"), Some(EpicsValue::Double(5.0)));
}
#[epics_macros_rs::epics_test]
async fn get_record_no_resolve_skips_alias_table() {
let db = PvDatabase::new();
db.add_record(
"TARGET",
Box::new(crate::server::records::ai::AiRecord::new(0.0)),
)
.await
.unwrap();
db.add_alias("ALIAS", "TARGET").await.unwrap();
assert!(db.get_record_no_resolve("TARGET").is_some());
assert!(
db.get_record_no_resolve("ALIAS").is_none(),
"get_record_no_resolve must not follow alias table"
);
}
#[epics_macros_rs::epics_test]
async fn register_cp_link_normalises_alias_to_canonical() {
let db = PvDatabase::new();
db.add_record(
"SRC_REAL",
Box::new(crate::server::records::ai::AiRecord::new(0.0)),
)
.await
.unwrap();
db.add_record(
"DST_REAL",
Box::new(crate::server::records::ai::AiRecord::new(0.0)),
)
.await
.unwrap();
db.add_alias("SRC_ALIAS", "SRC_REAL").await.unwrap();
db.add_alias("DST_ALIAS", "DST_REAL").await.unwrap();
db.register_cp_link("SRC_ALIAS", "DST_ALIAS", false).await;
let targets = db.get_cp_targets("SRC_REAL");
assert_eq!(targets.len(), 1);
assert_eq!(targets[0].record, "DST_REAL");
assert!(!targets[0].passive_only);
let alias_lookup = db.get_cp_targets("SRC_ALIAS");
assert!(alias_lookup.is_empty());
}
#[epics_macros_rs::epics_test]
async fn aliases_for_record_returns_sorted_targets_only() {
let db = PvDatabase::new();
db.add_record(
"TARGET",
Box::new(crate::server::records::ai::AiRecord::new(0.0)),
)
.await
.unwrap();
db.add_record(
"OTHER",
Box::new(crate::server::records::ai::AiRecord::new(0.0)),
)
.await
.unwrap();
db.add_alias("ZZ", "TARGET").await.unwrap();
db.add_alias("AA", "TARGET").await.unwrap();
db.add_alias("MM", "OTHER").await.unwrap();
assert_eq!(
db.aliases_for_record("TARGET"),
vec!["AA".to_string(), "ZZ".to_string()]
);
assert_eq!(db.aliases_for_record("OTHER"), vec!["MM".to_string()]);
assert!(db.aliases_for_record("MISSING").is_empty());
}
#[epics_macros_rs::epics_test]
async fn all_alias_names_returns_registered_aliases() {
let db = PvDatabase::new();
db.add_record(
"TARGET",
Box::new(crate::server::records::ai::AiRecord::new(0.0)),
)
.await
.unwrap();
db.add_alias("ALIAS_A", "TARGET").await.unwrap();
db.add_alias("ALIAS_B", "TARGET").await.unwrap();
let mut aliases = db.all_alias_names();
aliases.sort();
assert_eq!(aliases, vec!["ALIAS_A".to_string(), "ALIAS_B".to_string()]);
assert!(!aliases.contains(&"TARGET".to_string()));
}
#[epics_macros_rs::epics_test]
async fn complete_async_record_accepts_alias() {
let db = PvDatabase::new();
db.add_record(
"TARGET",
Box::new(crate::server::records::ai::AiRecord::new(0.0)),
)
.await
.unwrap();
db.add_alias("ALIAS", "TARGET").await.unwrap();
db.complete_async_record("ALIAS").await.unwrap();
db.complete_async_record("TARGET").await.unwrap();
}
#[epics_macros_rs::epics_test]
async fn process_record_accepts_alias() {
let db = PvDatabase::new();
db.add_record(
"TARGET",
Box::new(crate::server::records::ai::AiRecord::new(0.0)),
)
.await
.unwrap();
db.add_alias("ALIAS", "TARGET").await.unwrap();
db.process_record("TARGET").await.unwrap();
db.process_record("ALIAS").await.unwrap();
assert!(db.process_record("MISSING").await.is_err());
}
#[epics_macros_rs::epics_test]
async fn process_record_with_links_accepts_alias_and_avoids_cycle() {
let db = PvDatabase::new();
db.add_record(
"TARGET",
Box::new(crate::server::records::ai::AiRecord::new(0.0)),
)
.await
.unwrap();
db.add_alias("ALIAS", "TARGET").await.unwrap();
let mut visited = std::collections::HashSet::new();
db.process_record_with_links("ALIAS", &mut visited, 0)
.await
.unwrap();
assert!(
visited.contains("TARGET"),
"visited must record the canonical name: {visited:?}",
);
assert!(
!visited.contains("ALIAS"),
"visited must NOT record the alias form: {visited:?}",
);
}
#[epics_macros_rs::epics_test]
async fn alias_duplicate_rejected() {
let db = PvDatabase::new();
db.add_record(
"TARGET",
Box::new(crate::server::records::ai::AiRecord::new(0.0)),
)
.await
.unwrap();
db.add_alias("ALIAS", "TARGET").await.unwrap();
let err = db.add_alias("ALIAS", "TARGET").await;
assert!(err.is_err(), "duplicate alias name must be rejected");
}
#[epics_macros_rs::epics_test]
async fn add_pv_and_add_record_reject_duplicates_across_namespaces() {
use crate::server::records::ai::AiRecord;
let db = PvDatabase::new();
db.add_pv("A", EpicsValue::Double(1.0)).await.unwrap();
assert!(db.add_pv("A", EpicsValue::Double(2.0)).await.is_err());
let noop_hook: crate::server::pv::WriteHook =
std::sync::Arc::new(|_v, _ctx| Box::pin(async { Ok(()) }));
assert!(
db.add_pv_with_hook("A", EpicsValue::Double(2.0), noop_hook)
.await
.is_err()
);
assert!(
db.add_record("A", Box::new(AiRecord::new(0.0)))
.await
.is_err()
);
assert!(db.add_alias("A", "A").await.is_err());
db.add_record("R", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
assert!(
db.add_record("R", Box::new(AiRecord::new(1.0)))
.await
.is_err()
);
assert!(db.add_pv("R", EpicsValue::Double(0.0)).await.is_err());
assert!(db.add_alias("R", "R").await.is_err());
db.add_alias("AL", "R").await.unwrap();
assert!(db.add_pv("AL", EpicsValue::Double(0.0)).await.is_err());
assert!(
db.add_record("AL", Box::new(AiRecord::new(0.0)))
.await
.is_err()
);
}
#[epics_macros_rs::epics_test]
async fn remove_record_purges_dangling_aliases() {
use crate::server::records::ai::AiRecord;
let db = PvDatabase::new();
db.add_record("R", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
db.add_alias("ALT1", "R").await.unwrap();
db.add_alias("ALT2", "R").await.unwrap();
db.add_record("OTHER", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
db.add_alias("KEEPER", "OTHER").await.unwrap();
assert!(db.remove_record("R").await);
db.add_pv("ALT1", EpicsValue::Double(0.0)).await.unwrap();
db.add_pv("ALT2", EpicsValue::Double(0.0)).await.unwrap();
assert_eq!(db.resolve_alias("KEEPER"), Some("OTHER".to_string()));
}
#[epics_macros_rs::epics_test]
async fn add_alias_rejects_simple_pv_collision() {
use crate::server::records::ai::AiRecord;
let db = PvDatabase::new();
db.add_pv("PVX", EpicsValue::Double(0.0)).await.unwrap();
db.add_record("TARGET", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
assert!(db.add_alias("PVX", "TARGET").await.is_err());
}
#[epics_macros_rs::epics_test]
async fn concurrent_add_pv_and_add_record_do_not_deadlock() {
use crate::server::records::ai::AiRecord;
let db = std::sync::Arc::new(PvDatabase::new());
let db1 = db.clone();
let db2 = db.clone();
let h1 = crate::runtime::task::spawn(async move {
db1.add_pv("RACE", EpicsValue::Double(1.0)).await
});
let h2 = crate::runtime::task::spawn(async move {
db2.add_record("RACE", Box::new(AiRecord::new(0.0))).await
});
let r1 = crate::runtime::task::timeout(std::time::Duration::from_secs(2), h1)
.await
.expect("add_pv must not block on add_record");
let r2 = crate::runtime::task::timeout(std::time::Duration::from_secs(2), h2)
.await
.expect("add_record must not block on add_pv");
let r1 = r1.unwrap();
let r2 = r2.unwrap();
assert!(
(r1.is_ok() && r2.is_err()) || (r1.is_err() && r2.is_ok()),
"exactly one of the racing inserts must succeed: r1={r1:?} r2={r2:?}",
);
}
#[epics_macros_rs::epics_test]
async fn existence_gate_blocks_cached_simple_pv_per_request() {
use std::net::SocketAddr;
let db = PvDatabase::new();
db.add_pv("SHADOW:x", EpicsValue::Double(1.0))
.await
.unwrap();
db.add_record(
"REC",
Box::new(crate::server::records::ai::AiRecord::new(0.0)),
)
.await
.unwrap();
let denied: SocketAddr = "127.0.0.1:5064".parse().unwrap();
let allowed: SocketAddr = "192.0.2.5:5064".parse().unwrap();
assert!(db.has_name_from("SHADOW:x", Some(denied)).await);
assert!(db.find_entry_from("SHADOW:x", Some(denied)).await.is_some());
let gate: ExistenceGate = Arc::new(move |name, peer| {
Box::pin(async move { !(name == "SHADOW:x" && peer == Some(denied)) })
});
db.set_existence_gate(gate).await;
assert!(!db.has_name_from("SHADOW:x", Some(denied)).await);
assert!(db.find_entry_from("SHADOW:x", Some(denied)).await.is_none());
assert!(db.has_name_from("SHADOW:x", Some(allowed)).await);
assert!(
db.find_entry_from("SHADOW:x", Some(allowed))
.await
.is_some()
);
assert!(db.has_name_from("REC", Some(denied)).await);
assert!(db.find_entry_from("REC", Some(denied)).await.is_some());
}
#[epics_macros_rs::epics_test]
async fn record_link_fields_surfaces_device_support_inp() {
use crate::server::record::ParsedLink;
use crate::server::records::ai::AiRecord;
let db = PvDatabase::new();
db.add_record("AI", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
{
let rec = db.get_record("AI").unwrap();
rec.write().common.inp = "pva://mini:current?proc=CP".to_string();
}
let links = db.record_link_fields("AI");
let inp = links
.iter()
.find(|(f, _, _)| f == "INP")
.unwrap_or_else(|| panic!("INP link must be surfaced, got {links:?}"));
assert_eq!(inp.1, "pva://mini:current?proc=CP");
assert!(
matches!(inp.2, ParsedLink::Pva(_)),
"a pva:// INP must parse to ParsedLink::Pva, got {:?}",
inp.2
);
}
}