pub(crate) mod breakpoint;
pub mod db_access;
mod field_io;
pub mod filters;
mod link_put_queue;
mod link_set;
mod links;
mod processing;
mod record_lock;
pub(crate) mod scan_index;
mod snapshot;
pub use field_io::ProcessMode;
pub use link_set::{
DynLinkSet, LinkBacking, LinkDbfType, LinkDiagnostics, LinkMetadata, LinkPutOp, LinkSet,
LinkSetRegistry, PutAdmission, RemoteAlarm,
};
pub use processing::{AsyncDbHandle, AsyncToken};
pub use record_lock::{LockSetInfo, LockSetReport, ManyRecordWriteGuard, RecordWriteGuard};
use crate::error::{CaError, CaResult};
use arc_swap::{ArcSwap, ArcSwapOption};
use snapshot::SnapshotCell;
use std::collections::HashMap;
use std::sync::Arc;
use crate::server::pv::ProcessVariable;
use crate::server::record::{Record, RecordInstance};
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")
}
pub(crate) use tsel_stamp::TselStamp;
mod tsel_stamp {
use std::time::SystemTime;
use crate::server::record::CommonFields;
#[derive(Debug, Clone, Copy, PartialEq)]
pub(crate) enum TselStamp {
None,
Time(SystemTime, u64),
Tse(i16),
}
impl TselStamp {
pub(crate) fn stamp(self, name: &str, common: &mut CommonFields, is_soft: bool) {
match self {
TselStamp::Time(time, utag) => {
common.time = time;
common.utag = utag;
return;
}
TselStamp::None => {}
TselStamp::Tse(tse) => common.tse = tse,
}
apply_timestamp(name, common, is_soft);
}
}
fn apply_timestamp(name: &str, common: &mut CommonFields, _is_soft: bool) {
match crate::server::recgbl::get_time_stamp(common.tse, common.time) {
Some(t) => common.time = t,
None => crate::runtime::log::errlog_printf(&format!(
"recGblGetTimeStampSimm: epicsTimeGetEvent failed, {name}.TSE = {}\n",
common.tse
)),
}
}
}
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,
}
#[derive(Clone, PartialEq, Eq, PartialOrd, Ord, Debug)]
struct ScanKey {
phas: i16,
record_type: u32,
load_order: u64,
name: String,
}
impl ScanKey {
fn new(phas: i16, record_type: &str, load_order: u64, name: &str) -> Self {
use crate::server::record::dbd_generated::RECORD_TYPE_ORDER;
Self {
phas,
record_type: RECORD_TYPE_ORDER
.iter()
.position(|t| *t == record_type)
.unwrap_or(RECORD_TYPE_ORDER.len()) as u32,
load_order,
name: name.to_string(),
}
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct DbNode {
pub name: String,
pub alias_of: Option<String>,
}
struct PvDatabaseInner {
simple_pvs:
crate::runtime::sync::PriorityInheritanceMutex<HashMap<String, Arc<ProcessVariable>>>,
pv_destroyed_tx: tokio::sync::broadcast::Sender<()>,
records: parking_lot::RwLock<HashMap<String, Arc<parking_lot::RwLock<RecordInstance>>>>,
scan_index: 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>>>,
deferred_record_inits: std::sync::Mutex<Vec<String>>,
empty_link_assignments: std::sync::Mutex<HashMap<String, Vec<String>>>,
after_ioc_running: std::sync::Mutex<Vec<String>>,
external_resolver: ArcSwapOption<ExternalPvResolver>,
search_resolver: ArcSwapOption<SearchResolver>,
existence_gate: ArcSwapOption<ExistenceGate>,
breakpoints: ArcSwapOption<breakpoint::BreakpointTable>,
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,
record_attributes: crate::runtime::sync::PriorityInheritanceMutex<
std::collections::BTreeMap<String, std::collections::BTreeMap<String, String>>,
>,
subroutine_registry: ArcSwap<HashMap<String, Arc<crate::server::record::SubroutineFn>>>,
device_support_resolver: ArcSwapOption<crate::server::ioc_app::DeviceSupportResolver>,
breaktable_registry: SnapshotCell<crate::server::cvt_bpt::BreakTableRegistry>,
}
thread_local! {
static REGISTRATION_GATE_HELD: std::cell::Cell<Option<&'static str>> =
const { std::cell::Cell::new(None) };
}
#[must_use = "L46 is released as soon as the guard is dropped"]
pub(crate) struct RegistrationGate<'a> {
_guard: crate::runtime::sync::PriorityInheritanceMutexGuard<'a, ()>,
}
impl Drop for RegistrationGate<'_> {
fn drop(&mut self) {
REGISTRATION_GATE_HELD.with(|h| h.set(None));
}
}
#[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>),
Initialising(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(),
}
}
const MAX_ATTRIBUTE_LEN: usize = 39;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum RecordAttributeError {
BadField,
RecordTypeNotFound,
}
impl RecordAttributeError {
#[must_use]
pub const fn status(self) -> u32 {
match self {
Self::BadField => (511 << 16) | 15,
Self::RecordTypeNotFound => (512 << 16) | 1,
}
}
#[must_use]
pub const fn message(self) -> &'static str {
match self {
Self::BadField => "Illegal field value",
Self::RecordTypeNotFound => "Record Type does not exist",
}
}
}
fn seeded_record_attributes()
-> std::collections::BTreeMap<String, std::collections::BTreeMap<String, String>> {
crate::server::record::dbd_generated::RECORD_TYPES
.iter()
.map(|t| {
let mut attrs = std::collections::BTreeMap::new();
attrs.insert("RTYP".to_string(), (*t).to_string());
attrs.insert("VERS".to_string(), "none specified".to_string());
((*t).to_string(), attrs)
})
.collect()
}
impl PvDatabase {
pub(crate) fn lock_registration(&self, site: &'static str) -> RegistrationGate<'_> {
if let Some(holder) = REGISTRATION_GATE_HELD.with(|h| h.get()) {
panic!(
"L46 registration_mutex is not reentrant: `{site}` took it while \
this thread still holds it from `{holder}`. `update_scan_index` \
takes L46 itself and is the single owner of a scan-index \
transition, so a caller must DROP its registration gate before \
reaching it — see `RegistrationGate`."
);
}
let guard = self.inner.registration_mutex.lock();
REGISTRATION_GATE_HELD.with(|h| h.set(Some(site)));
RegistrationGate { _guard: guard }
}
pub fn new() -> Self {
Self {
inner: Arc::new(PvDatabaseInner {
simple_pvs: crate::runtime::sync::PriorityInheritanceMutex::new(HashMap::new()),
pv_destroyed_tx: tokio::sync::broadcast::channel(16).0,
external_resolver: ArcSwapOption::empty(),
breakpoints: 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: 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()),
deferred_record_inits: std::sync::Mutex::new(Vec::new()),
empty_link_assignments: 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(),
record_attributes: crate::runtime::sync::PriorityInheritanceMutex::new(
seeded_record_attributes(),
),
subroutine_registry: ArcSwap::from_pointee(HashMap::new()),
device_support_resolver: ArcSwapOption::empty(),
breaktable_registry: SnapshotCell::new(
crate::server::cvt_bpt::BreakTableRegistry::new(),
),
}),
}
}
pub fn put_record_type_attribute(
&self,
record_type: &str,
name: Option<&str>,
value: Option<&str>,
) -> Result<(), RecordAttributeError> {
let Some(name) = name else {
return Err(RecordAttributeError::BadField);
};
let value = value.unwrap_or("");
if !crate::server::record::dbd_generated::RECORD_TYPES.contains(&record_type) {
return Err(RecordAttributeError::RecordTypeNotFound);
}
let truncated: String = value.chars().take(MAX_ATTRIBUTE_LEN).collect();
self.inner
.record_attributes
.lock()
.entry(record_type.to_string())
.or_default()
.insert(name.to_string(), truncated);
Ok(())
}
pub fn record_type_attribute(&self, record_type: &str, name: &str) -> Option<String> {
self.inner
.record_attributes
.lock()
.get(record_type)?
.get(name)
.cloned()
}
pub fn record_type_attributes(&self, record_type: &str) -> Vec<(String, String)> {
self.inner
.record_attributes
.lock()
.get(record_type)
.map(|m| m.iter().map(|(k, v)| (k.clone(), v.clone())).collect())
.unwrap_or_default()
}
pub async fn add_breaktables(
&self,
tables: Vec<crate::server::cvt_bpt::BrkTable>,
) -> Arc<crate::server::cvt_bpt::BreakTableRegistry> {
let _gate = self.lock_registration("add_breaktables");
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 fn install_device_support_resolver(
&self,
resolver: crate::server::ioc_app::DeviceSupportResolver,
) {
self.inner
.device_support_resolver
.store(Some(Arc::new(resolver)));
}
pub(crate) async fn records_with_device_support(&self) -> usize {
let names = self.all_record_names().await;
names
.iter()
.filter_map(|name| self.get_record(name))
.filter(|rec| rec.read().device.is_some())
.count()
}
pub(crate) fn find_subroutine_named(
&self,
name: &str,
) -> Option<Arc<crate::server::record::SubroutineFn>> {
self.inner.subroutine_registry.load().get(name).cloned()
}
pub(crate) fn subroutine_entries(&self) -> Vec<(String, usize)> {
self.inner
.subroutine_registry
.load()
.iter()
.map(|(name, func)| (name.clone(), Arc::as_ptr(func) as *const () as usize))
.collect()
}
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(crate) fn breakpoints(&self) -> Option<Arc<breakpoint::BreakpointTable>> {
self.inner.breakpoints.load_full()
}
pub(crate) fn breakpoints_or_install(&self) -> Arc<breakpoint::BreakpointTable> {
if let Some(existing) = self.inner.breakpoints.load_full() {
return existing;
}
let fresh = Arc::new(breakpoint::BreakpointTable::new());
let prev = self.inner.breakpoints.compare_and_swap(
&None::<Arc<breakpoint::BreakpointTable>>,
Some(fresh.clone()),
);
match &*prev {
Some(winner) => winner.clone(),
None => fresh,
}
}
pub(crate) fn retire_breakpoints_if_idle(&self) {
if self
.inner
.breakpoints
.load()
.as_ref()
.is_some_and(|t| t.is_empty())
{
self.inner.breakpoints.store(None);
}
}
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.init_ready(name) {
connected += 1;
}
}
if connected == total {
return (connected, total);
}
if std::time::Instant::now() >= deadline {
return (connected, total);
}
crate::runtime::task::sleep_background(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.init_ready(&name) {
names.push(name);
}
}
names
}
fn link_field_texts(
inst: &RecordInstance,
) -> Vec<(String, String, crate::server::record::LinkFieldType)> {
use crate::server::record::LinkFieldType;
use crate::types::DbfLinkClass;
let record_type = inst.record.record_type();
let mut out: Vec<(String, String, LinkFieldType)> = Vec::new();
for desc in crate::server::record::declared_fields(record_type) {
let ftype = match crate::types::dbf_link_class(record_type, desc.name) {
Some(DbfLinkClass::InLink) => LinkFieldType::In,
Some(DbfLinkClass::OutLink) => LinkFieldType::Out,
Some(DbfLinkClass::FwdLink) => LinkFieldType::Fwd,
None => continue,
};
if out.iter().any(|(had, _, _)| had == desc.name) {
continue;
}
let raw = match inst.common_link_text(desc.name) {
Some(raw) => raw.to_string(),
None => match inst.record.get_field(desc.name) {
Some(EpicsValue::String(s)) => s.as_str_lossy().into_owned(),
_ => continue,
},
};
out.push((desc.name.to_string(), raw, ftype));
}
out
}
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 mut out: Vec<(String, String, crate::server::record::ParsedLink)> =
Self::link_field_texts(&rec.read())
.into_iter()
.filter(|(_, raw, _)| !raw.is_empty())
.map(|(field, raw, ftype)| {
let parsed = crate::server::record::parse_link_field(&raw, ftype);
(field, raw, parsed)
})
.filter(|(_, _, parsed)| !matches!(parsed, crate::server::record::ParsedLink::None))
.collect();
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)
}
pub fn external_link_puts_coalesced_for(&self, link_name: &str) -> u64 {
let (target, body) = Self::split_external_link_name(link_name);
self.inner
.link_puts
.coalesced_for(&link_put_queue::LinkKey {
target,
name: body.to_string(),
})
}
pub(crate) fn ca_link_init(&self) -> bool {
if self.inner.link_puts.network().is_none() {
return false;
}
self.inner
.link_puts
.ensure_owner(std::sync::Arc::downgrade(&self.inner));
true
}
fn stage_external_link_open(&self, target: link_put_queue::LinkTarget, name: &str) -> bool {
if self.inner.link_puts.network().is_none() {
return false;
}
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.lock_registration("add_pv");
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.lock_registration("add_pv_with_hooks_full");
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.lock_registration("remove_simple_pv");
let removed = self.inner.simple_pvs.lock().remove(name);
if let Some(pv) = &removed {
pv.destroy();
self.signal_destroyed();
}
removed
}
pub async fn add_simple_pv(&self, pv: ProcessVariable) -> CaResult<()> {
let _gate = self.lock_registration("add_simple_pv");
self.check_name_free(&pv.name)?;
let name = pv.name.clone();
self.inner.simple_pvs.lock().insert(name, Arc::new(pv));
Ok(())
}
pub fn pv_destroyed_events(&self) -> tokio::sync::broadcast::Receiver<()> {
self.inner.pv_destroyed_tx.subscribe()
}
fn signal_destroyed(&self) {
let _ = self.inner.pv_destroyed_tx.send(());
}
#[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::Initialising(_) | DbInitPhase::Running => Err(IocAlreadyInitialized),
}
}
pub fn ioc_is_running(&self) -> bool {
matches!(*self.inner.init_phase.lock().unwrap(), DbInitPhase::Running)
}
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) | DbInitPhase::Initialising(queued) => queued.push(init),
DbInitPhase::Unloaded | DbInitPhase::Running => {
drop(phase);
crate::runtime::task::spawn_background(
crate::runtime::task::CallbackPriority::Medium,
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 mut owed = {
let mut phase = self.inner.init_phase.lock().unwrap();
match std::mem::replace(&mut *phase, DbInitPhase::Initialising(Vec::new())) {
DbInitPhase::Loading(queued) => queued,
DbInitPhase::Unloaded => Vec::new(),
already @ (DbInitPhase::Initialising(_) | DbInitPhase::Running) => {
*phase = already;
return;
}
}
};
self.drain_deferred_record_inits();
self.ca_link_init();
self.db_init_record_links().await;
{
let mut phase = self.inner.init_phase.lock().unwrap();
if let DbInitPhase::Initialising(queued) =
std::mem::replace(&mut *phase, DbInitPhase::Running)
{
owed.extend(queued);
}
}
for init in owed {
init.await;
}
self.build_lock_sets();
}
async fn db_init_record_links(&self) {
use crate::runtime::log::ERL_ERROR;
use crate::server::record::{
DbLinkType, ParsedLink, declared_link_type, link_type_refusal,
};
let empty_assigned =
std::mem::take(&mut *self.inner.empty_link_assignments.lock().unwrap());
let states = crate::server::database::filters::sync::db_state_registry();
for name in self.all_record_names().await {
let Some(rec) = self.get_record(&name) else {
continue;
};
let (record_type, dtyp, fields) = {
let inst = rec.read();
(
inst.record.record_type().to_string(),
inst.common.dtyp.as_str().to_string(),
Self::link_field_texts(&inst),
)
};
let assigned_empty = empty_assigned.get(&name);
for (field, text, ftype) in fields {
let assigned = !text.is_empty()
|| assigned_empty.is_some_and(|fields| fields.iter().any(|f| f == &field));
let mut refused = false;
if assigned {
if let Some(line) =
link_type_refusal(&name, &record_type, Some(&dtyp), &field, &text)
{
eprintln!("{ERL_ERROR}: {line}");
refused = true;
}
}
if !assigned || refused {
let empty = declared_link_type(&record_type, Some(&dtyp), &field)
.map_or("", DbLinkType::empty_link_text);
if empty != text {
Self::set_link_text(&mut rec.write(), &field, empty);
}
continue;
}
if let ParsedLink::State(state) =
crate::server::record::parse_link_field(&text, ftype)
{
states.get_or_create(&state.name);
}
}
let mut inst = rec.write();
let inst = &mut *inst;
inst.record.init_links(&inst.common);
}
}
fn set_link_text(inst: &mut RecordInstance, field: &str, text: &str) {
let value = EpicsValue::String(text.into());
let put = if inst.common_link_text(field).is_some() {
inst.put_common_field_db_load(field, value).map(|_| ())
} else {
inst.record.put_field(field, value)
};
debug_assert!(put.is_ok(), "{field} is a link field with no writer");
let _ = inst.record.special(field, true);
}
pub fn channel_snapshot_for_field(
&self,
record: &Arc<parking_lot::RwLock<RecordInstance>>,
field: &str,
string_view: bool,
) -> Option<crate::server::snapshot::Snapshot> {
self.channel_snapshot_for_field_guarded(
record,
field,
string_view,
&mut std::collections::HashSet::new(),
)
}
pub(crate) fn channel_snapshot_for_field_guarded(
&self,
record: &Arc<parking_lot::RwLock<RecordInstance>>,
field: &str,
string_view: bool,
visited: &mut std::collections::HashSet<String>,
) -> Option<crate::server::snapshot::Snapshot> {
let link = {
let inst = record.read();
match inst.link_backed_metadata_field_of(field) {
None => {
return inst.channel_snapshot_for_field(
field,
string_view,
LinkBacking::none(),
);
}
Some(link_field) => match inst.record.get_field(&link_field) {
Some(EpicsValue::String(text)) => {
let text = text.as_str_lossy().into_owned();
(!text.is_empty()).then_some((link_field, text))
}
_ => None,
},
}
};
let mut resolved = HashMap::new();
if let Some((link_field, text)) = link {
let parsed = crate::server::record::parse_link_v2(&text);
if let Some(meta) = self.link_metadata(&parsed, visited) {
resolved.insert(link_field, meta);
}
}
record.read().channel_snapshot_for_field(
field,
string_view,
LinkBacking::resolved(&resolved),
)
}
pub fn resolve_link_backed_metadata(
&self,
record: &Arc<parking_lot::RwLock<RecordInstance>>,
) -> HashMap<String, LinkMetadata> {
let links: Vec<(String, String)> = {
let inst = record.read();
if inst.link_backed_metadata_links().is_empty() {
return HashMap::new();
}
inst.link_backed_metadata_links()
.iter()
.filter_map(|lf| match inst.record.get_field(lf) {
Some(EpicsValue::String(text)) => {
let text = text.as_str_lossy().into_owned();
(!text.is_empty()).then_some((lf.clone(), text))
}
_ => None,
})
.collect()
};
let mut resolved = HashMap::new();
for (link_field, text) in links {
let parsed = crate::server::record::parse_link_v2(&text);
let mut visited = std::collections::HashSet::new();
if let Some(meta) = self.link_metadata(&parsed, &mut visited) {
resolved.insert(link_field, meta);
}
}
resolved
}
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.lock_registration("add_loaded_record");
let _relink = self.lock_set_membership_change(name);
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);
}
}
let mut refused: Option<String> = None;
let mut empty_links: Vec<String> = Vec::new();
for (field, value) in load.common_fields {
if let crate::types::EpicsValue::String(text) = &value {
if text.as_str_lossy().is_empty()
&& instance.common_link_text(&field.to_uppercase()).is_some()
{
empty_links.push(field.to_uppercase());
}
}
if let crate::types::EpicsValue::String(text) = &value {
if let Some(refusal) = crate::server::db_loader::menu_value_refusal(
instance.record.record_type(),
name,
&field,
&text.as_str_lossy(),
) {
if let Some(notice) = refusal.notice {
eprintln!("{notice}");
}
eprintln!("{}", refusal.line);
if let Some(suggestion) = refusal.suggestion {
eprintln!("{suggestion}");
}
refused.get_or_insert(format!("{name}.{field}"));
continue;
}
}
if let Err(e) = instance.put_common_field_db_load(&field, value) {
eprintln!("put_common_field({field}) failed for {name}: {e}");
}
}
if let Some(what) = refused {
return Err(CaError::BadChoice(what));
}
for (key, value) in &load.info_tags {
instance.set_info(key, value);
}
if !empty_links.is_empty() {
self.inner
.empty_link_assignments
.lock()
.unwrap()
.insert(name.to_string(), empty_links);
}
let defer = self.is_load_deferring();
if !defer {
crate::server::ioc_app::attach_device_support(
&mut instance,
name,
self.inner.device_support_resolver.load().as_deref(),
);
instance.arm_init_subroutines(self.inner.subroutine_registry.load_full());
instance.run_init_passes(name);
super::database::processing::seed_constant_links(&mut instance);
}
let scan = instance.common.scan;
let phas = instance.common.phas;
let record_type = instance.record.record_type();
let rec_arc = Arc::new(parking_lot::RwLock::new(instance));
self.inner
.records
.write()
.insert(name.to_string(), rec_arc.clone());
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 defer {
self.inner
.deferred_record_inits
.lock()
.unwrap()
.push(name.to_string());
return Ok(());
}
self.add_to_scan_list(scan, phas, record_type, seq, name);
drop(gate);
self.rec_gbl_init_simm(&rec_arc);
self.arm_watchdog(name);
Ok(())
}
fn is_load_deferring(&self) -> bool {
matches!(
*self.inner.init_phase.lock().unwrap(),
DbInitPhase::Loading(_)
)
}
pub(crate) fn drain_deferred_record_inits(&self) {
let owed = std::mem::take(&mut *self.inner.deferred_record_inits.lock().unwrap());
for name in owed {
self.init_deferred_record(&name);
}
}
pub(crate) fn record_init_deferred(&self, name: &str) -> bool {
self.inner
.deferred_record_inits
.lock()
.unwrap()
.iter()
.any(|n| n == name)
}
fn init_deferred_record(&self, name: &str) {
let Some(rec_arc) = self.get_record(name) else {
return;
};
let (scan, phas, record_type) = {
let mut guard = rec_arc.write();
let instance = &mut *guard;
crate::server::ioc_app::attach_device_support(
instance,
name,
self.inner.device_support_resolver.load().as_deref(),
);
instance.arm_init_subroutines(self.inner.subroutine_registry.load_full());
instance.run_init_passes(name);
super::database::processing::seed_constant_links(instance);
(
instance.common.scan,
instance.common.phas,
instance.record.record_type(),
)
};
let seq = self
.inner
.load_order
.load()
.get(name)
.copied()
.unwrap_or(u64::MAX);
self.add_to_scan_list(scan, phas, record_type, seq, name);
self.rec_gbl_init_simm(&rec_arc);
self.arm_watchdog(name);
{
let mut guard = rec_arc.write();
let inst = &mut *guard;
inst.record.init_links(&inst.common);
}
self.rec_gbl_init_constant_links(&rec_arc);
}
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,
token: String::new(),
message: format!("name '{name}' is already registered as a {kind}"),
});
}
Ok(())
}
pub async fn remove_record(&self, name: &str) -> bool {
let _gate = self.lock_registration("remove_record");
let _relink = self.lock_set_membership_change(name);
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
};
self.delete_from_scan_list(scan, 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();
let orphaned: Vec<String> = aliases
.iter()
.filter(|(_, target)| *target == name)
.map(|(alias, _)| alias.clone())
.collect();
aliases.retain(|_alias, target| target != name);
drop(aliases);
if !orphaned.is_empty() {
self.inner.load_order.update(|m| {
for alias in &orphaned {
m.remove(alias);
}
});
}
rec_arc.write().destroy();
self.signal_destroyed();
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.lock_registration("add_alias");
let _relink = self.lock_set_membership_change(target);
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());
let seq = self
.inner
.load_order_counter
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
self.inner.load_order.update(|m| {
m.insert(alias.to_string(), seq);
});
Ok(())
}
pub fn breaktable_registry(&self) -> Arc<crate::server::cvt_bpt::BreakTableRegistry> {
self.inner.breaktable_registry.load_full()
}
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;
if self
.inner
.simple_pvs
.lock()
.contains_key(record_path.as_str())
{
return true;
}
let (base, explicit_field) = match record_path.rsplit_once('.') {
Some((base, field)) => (base, Some(field)),
None => (record_path.as_str(), None),
};
let rec = self.inner.records.read().get(base).cloned().or_else(|| {
let target = self.inner.aliases.read().get(base).cloned();
target.and_then(|t| self.inner.records.read().get(&t).cloned())
});
let Some(rec) = rec else {
return false;
};
let Some(field) = explicit_field else {
return true;
};
let instance = rec.read();
match field.strip_suffix('$') {
Some(core) => instance
.resolve_string_view_field(&core.to_ascii_uppercase())
.is_some(),
None => {
let upper = field.to_ascii_uppercase();
instance.resolve_field(&upper).is_some()
|| instance.field_desc(&upper).is_some()
|| instance.resolves_noaccess_name(&upper)
|| self
.record_type_attribute(instance.record.record_type(), &upper)
.is_some()
}
}
}
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 async fn all_db_nodes(&self) -> Vec<DbNode> {
let mut nodes: Vec<DbNode> = {
let records = self.inner.records.read();
records
.keys()
.map(|name| DbNode {
name: name.clone(),
alias_of: None,
})
.collect()
};
nodes.extend({
let aliases = self.inner.aliases.read();
aliases
.iter()
.map(|(alias, target)| DbNode {
name: alias.clone(),
alias_of: Some(target.clone()),
})
.collect::<Vec<_>>()
});
let load_order = self.inner.load_order.load();
nodes.sort_by(|a, b| {
let seq = |n: &DbNode| load_order.get(&n.name).copied().unwrap_or(u64::MAX);
seq(a).cmp(&seq(b)).then_with(|| a.name.cmp(&b.name))
});
nodes
}
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 channel_view_door_tests {
use super::PvDatabase;
use crate::server::records::ai::AiRecord;
#[tokio::test]
async fn the_snapshot_door_refuses_an_ineligible_dollar_view() {
let db = PvDatabase::new();
db.add_record("VD:ai", Box::new(AiRecord::new(1.5)))
.await
.unwrap();
let rec = db.get_record("VD:ai").expect("record");
assert!(
db.channel_snapshot_for_field(&rec, "VAL", false).is_some(),
"the unviewed VAL is an ordinary double snapshot"
);
assert!(
db.channel_snapshot_for_field(&rec, "VAL", true).is_none(),
"`VAL$` on a DBF_DOUBLE is S_dbLib_fieldNotFound, not a double"
);
for eligible in ["DESC", "NAME", "EGU", "FLNK"] {
let snap = db
.channel_snapshot_for_field(&rec, eligible, true)
.unwrap_or_else(|| panic!("`{eligible}$` must be eligible"));
assert!(
matches!(snap.value, crate::types::EpicsValue::String(_)),
"`{eligible}$` serves the string, got {:?}",
snap.value
);
}
}
}
#[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;
TselStamp::None.stamp("REC", &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;
TselStamp::None.stamp("REC", &mut common, false);
assert_eq!(
common.time, device_time,
"TSE=-2 (epicsTimeEventDeviceTime) must preserve device-provided time"
);
}
#[test]
fn apply_timestamp_below_best_time_keeps_the_stale_stamp_and_errlogs() {
use crate::server::record::CommonFields;
use std::io::Write;
use std::sync::{Arc, Mutex};
use std::time::{Duration, SystemTime};
use tracing_subscriber::fmt::MakeWriter;
#[derive(Clone, Default)]
struct CaptureBuf(Arc<Mutex<Vec<u8>>>);
impl Write for CaptureBuf {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.0.lock().unwrap().extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
impl<'a> MakeWriter<'a> for CaptureBuf {
type Writer = CaptureBuf;
fn make_writer(&'a self) -> Self::Writer {
self.clone()
}
}
let buf = CaptureBuf::default();
let subscriber = tracing_subscriber::fmt()
.with_writer(buf.clone())
.with_max_level(tracing::Level::INFO)
.with_target(false)
.without_time()
.finish();
let _guard = tracing::subscriber::set_default(subscriber);
let stale = SystemTime::UNIX_EPOCH + Duration::from_secs(1_000_000);
let mut common = CommonFields::default();
common.tse = -3;
common.time = stale;
TselStamp::None.stamp("X", &mut common, false);
assert_eq!(
common.time, stale,
"TSE below epicsTimeEventBestTime must leave TIME alone, not stamp now"
);
let logged = String::from_utf8_lossy(&buf.0.lock().unwrap()).into_owned();
assert!(
logged.contains("recGblGetTimeStampSimm: epicsTimeGetEvent failed, X.TSE = -3"),
"C errlogs the failed event lookup; captured: {logged:?}"
);
}
#[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());
}
struct ConnectedMetaPendingLset {
names: Vec<String>,
}
#[async_trait::async_trait]
impl link_set::LinkSet for ConnectedMetaPendingLset {
fn is_connected(&self, _: &str) -> bool {
true
}
fn init_ready(&self, _: &str) -> bool {
false
}
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_holds_until_init_ready() {
let db = PvDatabase::new();
db.add_pv("meta:pending", EpicsValue::Long(0))
.await
.unwrap();
let lset = Arc::new(ConnectedMetaPendingLset {
names: vec!["meta:pending".to_string()],
});
db.register_link_set("ca", lset).await;
let (c, t) = db
.wait_for_external_links(std::time::Duration::from_millis(250))
.await;
assert_eq!((c, t), (0, 1));
assert_eq!(
db.unconnected_external_links().await,
vec!["meta:pending".to_string()]
);
}
#[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 search_gate_refuses_a_field_the_record_does_not_have() {
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();
assert!(db.has_name("TARGET.VAL").await);
assert!(db.has_name("TARGET.SEVR").await);
assert!(!db.has_name("TARGET.NOSUCH").await);
assert!(db.has_name("TARGET.MLOK").await);
db.add_record(
"WF",
Box::new(crate::server::records::waveform::WaveformRecord::new(
8,
crate::types::DbFieldType::Double,
)),
)
.await
.unwrap();
assert!(db.has_name("WF.BPTR").await);
assert!(!db.has_name("WF.NOSUCH").await);
assert!(db.has_name("ALIAS_NAME.EGU").await);
assert!(!db.has_name("ALIAS_NAME.NOSUCH").await);
assert!(db.has_name("TARGET.EGU$").await);
assert!(!db.has_name("TARGET.VAL$").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.is_empty(),
"a finished frame leaves no marker behind: {visited:?}",
);
let mut seeded = std::collections::HashSet::new();
seeded.insert("TARGET".to_string());
db.process_record_with_links("ALIAS", &mut seeded, 0)
.await
.unwrap();
assert!(
!seeded.contains("ALIAS"),
"the alias form must never enter the set: {seeded:?}",
);
assert_eq!(
seeded.len(),
1,
"the alias resolved to TARGET and was declined, adding nothing: {seeded:?}",
);
}
#[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 all_db_nodes_interleaves_aliases_at_their_load_position() {
use crate::server::records::ai::AiRecord;
let db = PvDatabase::new();
db.add_record("FIRST", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
db.add_alias("FIRST:ALT", "FIRST").await.unwrap();
db.add_record("SECOND", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
assert_eq!(
db.all_db_nodes().await,
vec![
DbNode {
name: "FIRST".into(),
alias_of: None
},
DbNode {
name: "FIRST:ALT".into(),
alias_of: Some("FIRST".into())
},
DbNode {
name: "SECOND".into(),
alias_of: None
},
]
);
}
#[epics_macros_rs::epics_test]
async fn removing_a_record_drops_its_alias_nodes_from_the_list() {
use crate::server::records::ai::AiRecord;
let db = PvDatabase::new();
db.add_record("GONE", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
db.add_alias("GONE:ALT", "GONE").await.unwrap();
db.add_record("STAYS", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
assert!(db.remove_record("GONE").await);
assert_eq!(
db.all_db_nodes().await,
vec![DbNode {
name: "STAYS".into(),
alias_of: None
}]
);
db.add_record("GONE", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
db.add_alias("GONE:ALT", "GONE").await.unwrap();
assert_eq!(
db.all_db_nodes()
.await
.into_iter()
.map(|node| node.name)
.collect::<Vec<_>>(),
vec!["STAYS", "GONE", "GONE:ALT"]
);
}
#[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 reactor =
crate::runtime::task::Reactor::current().expect("the test driver enters an executor");
let h1 = reactor.spawn(async move { db1.add_pv("RACE", EpicsValue::Double(1.0)).await });
let h2 = reactor
.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());
}
#[test]
fn the_declared_class_sweep_covers_every_name_the_hand_lists_spelled() {
use crate::types::{DbfLinkClass, dbf_link_class};
let cp_inputs: &[(&str, &str)] = &[
("DOL", "ao"),
("DOL0", "seq"),
("DOLF", "seq"),
("DOL1", "sseq"),
("DOLA", "sseq"),
("NVL", "sel"),
("SELL", "sseq"),
("SVL", "histogram"),
];
for (field, record_type) in cp_inputs {
assert_eq!(
dbf_link_class(record_type, field),
Some(DbfLinkClass::InLink),
"{record_type}.{field} was a CP_INPUT_LINK_FIELDS name and must \
still resolve as an input link"
);
}
assert_eq!(dbf_link_class("histogram", "SGNL"), None);
for (field, _) in crate::server::record::record_instance::COMMON_LINK_FIELDS {
assert!(
dbf_link_class("ai", field).is_some() || dbf_link_class("ao", field).is_some(),
"{field} was a COMMON_LINK_FIELDS name and must still resolve"
);
}
let unspelled: &[(&str, &str, DbfLinkClass)] = &[
("fanout", "LNK1", DbfLinkClass::FwdLink),
("dfanout", "OUTA", DbfLinkClass::OutLink),
("ai", "SIML", DbfLinkClass::InLink),
("ai", "SIOL", DbfLinkClass::InLink),
("aSub", "SUBL", DbfLinkClass::InLink),
("aSub", "OUTA", DbfLinkClass::OutLink),
("seq", "LNK1", DbfLinkClass::OutLink),
];
for (record_type, field, class) in unspelled {
assert_eq!(
dbf_link_class(record_type, field),
Some(*class),
"{record_type}.{field}"
);
}
}
#[epics_macros_rs::epics_test]
async fn link_field_texts_is_the_declared_link_set() {
use crate::server::record::LinkFieldType;
use crate::server::records::fanout::FanoutRecord;
let db = PvDatabase::new();
db.add_record("FAN", Box::new(FanoutRecord::default()))
.await
.unwrap();
{
let rec = db.get_record("FAN").unwrap();
let mut inst = rec.write();
inst.record
.put_field("LNK1", EpicsValue::String("TARGET".into()))
.unwrap();
}
let inst = db.get_record("FAN").unwrap();
let inst = inst.read();
let fields = PvDatabase::link_field_texts(&inst);
let lnk1 = fields
.iter()
.find(|(f, _, _)| f == "LNK1")
.unwrap_or_else(|| panic!("fanout.LNK1 must be enumerated, got {fields:?}"));
assert_eq!(lnk1.1, "TARGET");
assert!(
matches!(lnk1.2, LinkFieldType::Fwd),
"fanout.LNK1 is DBF_FWDLINK, got {:?}",
lnk1.2
);
assert!(
!fields.iter().any(|(f, _, _)| f == "VAL"),
"a non-link field must not be enumerated, got {fields:?}"
);
}
#[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.register_link_set(
"pva",
std::sync::Arc::new(DelayedConnectLset {
names: Vec::new(),
connect_at: crate::runtime::task::Instant::now(),
}),
)
.await;
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
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn ioc_init_starts_the_ca_link_owner() {
fn table() -> String {
let out = std::cell::RefCell::new(String::new());
crate::runtime::taskwd::taskwd_show(1, &|line| {
out.borrow_mut().push_str(line);
out.borrow_mut().push('\n');
});
out.into_inner()
}
assert!(
!table().contains("dbCaLink"),
"the link owner was registered before any IOC init"
);
let db = PvDatabase::new();
db.ioc_init().await;
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
while !table().contains("dbCaLink") {
assert!(
std::time::Instant::now() < deadline,
"`dbCaLink` never reached the watchdog table:\n{}",
table()
);
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
}
#[cfg(tokio_backend)]
#[test]
fn ca_link_init_starts_nothing_without_a_reactor() {
let db = PvDatabase::new();
assert!(
!db.ca_link_init(),
"a database with no captured reactor must refuse to start the owner"
);
let out = std::cell::RefCell::new(String::new());
crate::runtime::taskwd::taskwd_show(1, &|line| {
out.borrow_mut().push_str(line);
out.borrow_mut().push('\n');
});
assert!(
!out.into_inner().contains("dbCaLink"),
"the refused owner still reached the watchdog table"
);
}
}