use std::collections::HashSet;
use std::sync::Arc;
use crate::runtime::sync::RwLock;
use crate::server::record::{AlarmSeverity, NotifyWaitSet, OutTarget, RecordInstance, ScanType};
use crate::types::{DbFieldType, EpicsValue, PvString};
use super::link_set::LinkDbfType;
use super::processing::join_put_notify;
use super::{LinkPutOp, PvDatabase, SelmKind, SelmResult, dbr_ushort_cast, select_link_indices_ex};
fn empty_read_fetch(
link: &crate::server::record::ParsedLink,
) -> crate::server::recgbl::simm::LinkFetch {
use crate::server::recgbl::simm::LinkFetch;
if crate::server::recgbl::simm::is_constant(link) {
LinkFetch::NoData
} else {
LinkFetch::Failed
}
}
fn char_bytes_as_string(value: &EpicsValue, max_elements: usize) -> Option<PvString> {
let bytes: &[u8] = match value {
EpicsValue::CharArray(b) | EpicsValue::UCharArray(b) => b,
EpicsValue::Char(b) | EpicsValue::UChar(b) => std::slice::from_ref(b),
_ => return None,
};
let n = bytes.len().min(max_elements);
let text = &bytes[..n];
let end = text.iter().position(|&b| b == 0).unwrap_or(n);
Some(PvString::from_bytes(text[..end].to_vec()))
}
fn link_dbf_to_field_type(t: LinkDbfType) -> DbFieldType {
match t {
LinkDbfType::Char => DbFieldType::Char,
LinkDbfType::UChar => DbFieldType::UChar,
LinkDbfType::Short => DbFieldType::Short,
LinkDbfType::UShort => DbFieldType::UShort,
LinkDbfType::Long => DbFieldType::Long,
LinkDbfType::ULong => DbFieldType::ULong,
LinkDbfType::Int64 => DbFieldType::Int64,
LinkDbfType::UInt64 => DbFieldType::UInt64,
LinkDbfType::Float => DbFieldType::Float,
LinkDbfType::Double => DbFieldType::Double,
LinkDbfType::String => DbFieldType::String,
LinkDbfType::Enum => DbFieldType::Enum,
}
}
fn field_type_to_link_dbf(t: DbFieldType) -> LinkDbfType {
match t {
DbFieldType::Char => LinkDbfType::Char,
DbFieldType::UChar => LinkDbfType::UChar,
DbFieldType::Short => LinkDbfType::Short,
DbFieldType::UShort => LinkDbfType::UShort,
DbFieldType::Long => LinkDbfType::Long,
DbFieldType::ULong => LinkDbfType::ULong,
DbFieldType::Int64 => LinkDbfType::Int64,
DbFieldType::UInt64 => LinkDbfType::UInt64,
DbFieldType::Float => LinkDbfType::Float,
DbFieldType::Double => LinkDbfType::Double,
DbFieldType::String => LinkDbfType::String,
DbFieldType::Enum => LinkDbfType::Enum,
}
}
pub(crate) const DFANOUT_LINK_FIELDS: [&str; 16] = [
"OUTA", "OUTB", "OUTC", "OUTD", "OUTE", "OUTF", "OUTG", "OUTH", "OUTI", "OUTJ", "OUTK", "OUTL",
"OUTM", "OUTN", "OUTO", "OUTP",
];
pub(crate) const LNK_LINK_FIELDS: [&str; 16] = [
"LNK0", "LNK1", "LNK2", "LNK3", "LNK4", "LNK5", "LNK6", "LNK7", "LNK8", "LNK9", "LNKA", "LNKB",
"LNKC", "LNKD", "LNKE", "LNKF",
];
pub(crate) const CP_INPUT_LINK_FIELDS: &[&str] = &[
"DOL", "DOL0", "DOL1", "DOL2", "DOL3", "DOL4", "DOL5", "DOL6", "DOL7", "DOL8", "DOL9", "DOLA",
"DOLB", "DOLC", "DOLD", "DOLE", "DOLF", "NVL", "SELL", "SVL",
];
#[derive(Clone, Debug)]
pub(crate) struct LinkAlarm {
pub stat: u16,
pub sevr: AlarmSeverity,
pub amsg: String,
}
impl LinkAlarm {
pub(crate) fn pending(common: &crate::server::record::CommonFields) -> Self {
LinkAlarm {
stat: common.nsta,
sevr: common.nsev,
amsg: common.namsg.clone(),
}
}
pub(crate) fn committed(common: &crate::server::record::CommonFields) -> Self {
LinkAlarm {
stat: common.stat,
sevr: common.sevr,
amsg: common.amsg.clone(),
}
}
}
pub(crate) fn inherit_sevr_msg(
dest: &mut crate::server::record::CommonFields,
ms: crate::server::record::MonitorSwitch,
src: &LinkAlarm,
) {
use crate::server::recgbl::{alarm_status, rec_gbl_set_sevr, rec_gbl_set_sevr_msg};
use crate::server::record::{AlarmSeverity, MonitorSwitch};
match ms {
MonitorSwitch::Maximize => {
rec_gbl_set_sevr(dest, alarm_status::LINK_ALARM, src.sevr);
}
MonitorSwitch::MaximizeIfInvalid => {
if src.sevr == AlarmSeverity::Invalid {
rec_gbl_set_sevr(dest, alarm_status::LINK_ALARM, src.sevr);
}
}
MonitorSwitch::MaximizeStatus => {
rec_gbl_set_sevr_msg(dest, src.stat, src.sevr, src.amsg.clone());
}
MonitorSwitch::NoMaximize => {} }
}
#[derive(Clone, Copy)]
pub(crate) struct OutLinkSrc<'a> {
pub putf: bool,
pub notify: Option<&'a Arc<NotifyWaitSet>>,
pub alarm: &'a LinkAlarm,
pub field: &'a str,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub(crate) enum ProcessTargetGate {
ScanPassive,
ProcField,
}
impl ProcessTargetGate {
fn admits(self, target_scan: ScanType) -> bool {
match self {
Self::ScanPassive => target_scan == ScanType::Passive,
Self::ProcField => true,
}
}
}
#[derive(Clone, Debug)]
pub(crate) struct SeqGroup {
pub dol: String,
pub lnk: String,
pub dly: f64,
pub dov: f64,
}
pub(crate) enum MultiOut {
Fanout(Vec<String>),
Dfanout(Vec<String>),
Seq(Vec<SeqGroup>),
}
impl MultiOut {
fn len(&self) -> usize {
match self {
MultiOut::Fanout(v) => v.len(),
MultiOut::Dfanout(v) => v.len(),
MultiOut::Seq(v) => v.len(),
}
}
}
pub(crate) fn multi_output_dispatch_owned(record_type: &str) -> bool {
matches!(record_type, "fanout" | "dfanout" | "seq")
}
#[derive(Clone, Copy, PartialEq, Eq)]
pub(crate) enum MultiOutPhaseKind {
Output,
ForwardLink,
}
#[derive(Clone, Copy)]
pub(crate) enum MultiOutPhase {
Output { skip_out: bool },
ForwardLink,
}
pub(crate) fn multi_out_phase_of(record_type: &str) -> MultiOutPhaseKind {
match record_type {
"dfanout" | "seq" => MultiOutPhaseKind::Output,
_ => MultiOutPhaseKind::ForwardLink,
}
}
impl PvDatabase {
async fn read_db_link_value(&self, db: &crate::server::record::DbLink) -> Option<EpicsValue> {
self.read_target_value(&db.record, &db.field).await
}
async fn read_target_value(&self, record: &str, field: &str) -> Option<EpicsValue> {
let pv_name = if field == "VAL" {
record.to_string()
} else {
format!("{record}.{field}")
};
if self.has_name_no_resolve(record).await {
self.get_pv(&pv_name).await.ok()
} else {
self.resolve_external_pv(&pv_name).await
}
}
pub(crate) async fn read_link_value(
&self,
link: &crate::server::record::ParsedLink,
visited: &mut HashSet<String>,
depth: usize,
) -> Option<EpicsValue> {
match link {
crate::server::record::ParsedLink::None => None,
crate::server::record::ParsedLink::Ca(ca) => self.resolve_external_pv(&ca.pv).await,
crate::server::record::ParsedLink::Pva(name) => self.resolve_external_pv(name).await,
crate::server::record::ParsedLink::PvaJson(j) => {
self.resolve_external_pv(&j.link_identity_key()).await
}
crate::server::record::ParsedLink::Constant(_) => None,
crate::server::record::ParsedLink::Db(db) => {
self.process_passive_db_source(db, visited, depth).await;
self.read_db_link_value(db).await
}
crate::server::record::ParsedLink::Hw(_) => None,
crate::server::record::ParsedLink::Calc(calc) => self.evaluate_calc_link(calc).await,
}
}
pub(crate) async fn read_link_value_as(
&self,
link: &crate::server::record::ParsedLink,
read_as: crate::server::record::LinkReadAs,
visited: &mut HashSet<String>,
depth: usize,
) -> crate::server::recgbl::simm::LinkFetch {
use crate::server::recgbl::simm::LinkFetch;
use crate::server::record::LinkReadAs;
let Some(value) = self.read_link_value(link, visited, depth).await else {
return empty_read_fetch(link);
};
let converted = match read_as {
LinkReadAs::Native => Some(value),
LinkReadAs::Double => {
let scalar = if value.is_array() {
value.first_element()
} else {
Some(value)
};
scalar.and_then(|s| s.to_f64()).map(EpicsValue::Double)
}
LinkReadAs::String => self
.dbr_string_of(link, &value)
.await
.map(EpicsValue::String),
LinkReadAs::CharArrayAsString { max_elements } => {
char_bytes_as_string(&value, max_elements).map(EpicsValue::String)
}
};
match converted {
Some(v) => LinkFetch::Value(v),
None => LinkFetch::Failed,
}
}
async fn dbr_string_of(
&self,
link: &crate::server::record::ParsedLink,
value: &EpicsValue,
) -> Option<PvString> {
if let crate::server::record::ParsedLink::Db(db) = link
&& self.has_name_no_resolve(&db.record).await
&& let Some(rec) = self.get_record(&db.record).await
{
let guard = rec.read().await;
return guard.field_as_dbr_string(&db.field);
}
crate::server::record::value_as_dbr_string(value)
}
pub(crate) async fn read_link_value_no_process(
&self,
link: &crate::server::record::ParsedLink,
) -> Option<EpicsValue> {
match link {
crate::server::record::ParsedLink::None => None,
crate::server::record::ParsedLink::Ca(ca) => self.resolve_external_pv(&ca.pv).await,
crate::server::record::ParsedLink::Pva(name) => self.resolve_external_pv(name).await,
crate::server::record::ParsedLink::PvaJson(j) => {
self.resolve_external_pv(&j.link_identity_key()).await
}
crate::server::record::ParsedLink::Constant(_) => link.constant_value(),
crate::server::record::ParsedLink::Db(db) => self.read_db_link_value(db).await,
crate::server::record::ParsedLink::Hw(_) => None,
crate::server::record::ParsedLink::Calc(calc) => self.evaluate_calc_link(calc).await,
}
}
pub async fn evaluate_calc_link(
&self,
calc: &crate::server::record::CalcLink,
) -> Option<EpicsValue> {
use crate::calc::engine::{CALC_NARGS, NumericInputs};
if calc.args.len() > CALC_NARGS {
return None;
}
let mut vars = [0.0f64; CALC_NARGS];
for (i, arg) in calc.args.iter().enumerate() {
let (record, field) = match arg.rsplit_once('.') {
Some((r, f)) => (r, f),
None => (arg.as_str(), "VAL"),
};
let v = self.read_target_value(record, field).await?;
vars[i] = v.to_f64()?;
}
let compiled = crate::calc::compile(&calc.expr).ok()?;
let mut inputs = NumericInputs::with_vars(vars);
let result = crate::calc::eval(&compiled, &mut inputs).ok()?;
Some(EpicsValue::Double(result))
}
pub async fn evaluate_calc_link_with_time(
&self,
calc: &crate::server::record::CalcLink,
) -> Option<(EpicsValue, Option<std::time::SystemTime>)> {
let value = self.evaluate_calc_link(calc).await?;
let time = match calc.time_source {
Some(letter) => {
let idx = (letter as u8).saturating_sub(b'A') as usize;
let src = calc.args.get(idx)?;
let record_name = src.rsplit_once('.').map(|(r, _)| r).unwrap_or(src);
if self.has_name_no_resolve(record_name).await {
let rec = self.get_record(record_name).await?;
let inst = rec.read().await;
Some(inst.common.time)
} else {
let (secs, ns, _utag) = self
.external_link_time(&format!("ca://{record_name}"))
.await?;
let secs = secs.max(0) as u64;
let ns = (ns.max(0) as u32).min(999_999_999);
Some(std::time::UNIX_EPOCH + std::time::Duration::new(secs, ns))
}
}
None => None,
};
Some((value, time))
}
pub(crate) async fn read_link_with_alarm(
&self,
link: &crate::server::record::ParsedLink,
) -> (crate::server::recgbl::simm::LinkFetch, Option<LinkAlarm>) {
use crate::server::recgbl::simm::LinkFetch;
let (value, alarm) = self.read_link_value_and_alarm(link).await;
let fetch = match value {
Some(v) => LinkFetch::Value(v),
None => empty_read_fetch(link),
};
(fetch, alarm)
}
pub(crate) async fn input_link_inheritance(
&self,
reader_name: &str,
link: &crate::server::record::ParsedLink,
alarm: Option<LinkAlarm>,
) -> Option<(crate::server::record::MonitorSwitch, LinkAlarm)> {
let alarm = alarm?;
match link {
crate::server::record::ParsedLink::Db(db) => {
let target = self
.resolve_alias(&db.record)
.await
.unwrap_or_else(|| db.record.clone());
let reader = self
.resolve_alias(reader_name)
.await
.unwrap_or_else(|| reader_name.to_string());
if target == reader {
return None;
}
Some((db.monitor_switch, alarm))
}
crate::server::record::ParsedLink::Ca(ca) => Some((ca.monitor_switch, alarm)),
crate::server::record::ParsedLink::Pva(_)
| crate::server::record::ParsedLink::PvaJson(_) => {
Some((crate::server::record::MonitorSwitch::MaximizeStatus, alarm))
}
_ => None,
}
}
async fn read_link_value_and_alarm(
&self,
link: &crate::server::record::ParsedLink,
) -> (Option<EpicsValue>, Option<LinkAlarm>) {
match link {
crate::server::record::ParsedLink::Db(db) => {
let pv_name = if db.field == "VAL" {
db.record.clone()
} else {
format!("{}.{}", db.record, db.field)
};
if !self.has_name_no_resolve(&db.record).await {
return (
self.resolve_external_pv(&pv_name).await,
self.external_link_alarm(&pv_name).await,
);
}
let value = self.get_pv(&pv_name).await.ok();
let alarm = if let Some(rec) = self.get_record(&db.record).await {
let inst = rec.read().await;
Some(LinkAlarm::committed(&inst.common))
} else {
None
};
(value, alarm)
}
crate::server::record::ParsedLink::Constant(_) => (None, None),
crate::server::record::ParsedLink::Pva(_)
| crate::server::record::ParsedLink::PvaJson(_)
| crate::server::record::ParsedLink::Ca(_) => {
let name = link
.external_pv_name()
.expect("Ca/Pva/PvaJson link carries a PV name");
let value = self.resolve_external_pv(&name).await;
let alarm = self.external_link_alarm(&name).await;
(value, alarm)
}
crate::server::record::ParsedLink::Calc(calc) => {
(self.evaluate_calc_link(calc).await, None)
}
crate::server::record::ParsedLink::Hw(_) | crate::server::record::ParsedLink::None => {
(None, None)
}
}
}
async fn registered_link_sets(&self) -> Vec<crate::server::database::DynLinkSet> {
let registry = self.inner.link_sets.read().await;
registry
.schemes()
.iter()
.filter_map(|s| registry.get(s))
.collect()
}
pub(crate) async fn external_link_time(&self, name: &str) -> Option<(i64, i32, u64)> {
let (scheme, body) = if let Some(rest) = name.strip_prefix("pva://") {
("pva", rest)
} else if let Some(rest) = name.strip_prefix("ca://") {
("ca", rest)
} else {
for lset in self.registered_link_sets().await {
if let Some(ts) = lset.time_stamp(name).await {
return Some(ts);
}
}
return None;
};
let lset = self.inner.link_sets.read().await.get(scheme)?;
lset.time_stamp(body).await
}
async fn external_link_alarm(&self, name: &str) -> Option<LinkAlarm> {
let (scheme, body) = if let Some(rest) = name.strip_prefix("pva://") {
("pva", rest)
} else if let Some(rest) = name.strip_prefix("ca://") {
("ca", rest)
} else {
for lset in self.registered_link_sets().await {
if let Some(sev) = lset.alarm_severity(name).await {
return Some(LinkAlarm {
stat: lset
.alarm_status(name)
.await
.map(|s| s as u16)
.unwrap_or(crate::server::recgbl::alarm_status::LINK_ALARM),
sevr: crate::server::record::AlarmSeverity::from_u16(sev as u16),
amsg: lset.alarm_message(name).await.unwrap_or_default(),
});
}
}
return None;
};
let lset = self.inner.link_sets.read().await.get(scheme)?;
let sev = lset.alarm_severity(body).await?;
Some(LinkAlarm {
stat: lset
.alarm_status(body)
.await
.map(|s| s as u16)
.unwrap_or(crate::server::recgbl::alarm_status::LINK_ALARM),
sevr: crate::server::record::AlarmSeverity::from_u16(sev as u16),
amsg: lset.alarm_message(body).await.unwrap_or_default(),
})
}
pub async fn external_link_alarm_snapshot(
&self,
name: &str,
) -> Option<crate::server::database::RemoteAlarm> {
let (scheme, body) = if let Some(rest) = name.strip_prefix("pva://") {
("pva", rest)
} else if let Some(rest) = name.strip_prefix("ca://") {
("ca", rest)
} else {
for lset in self.registered_link_sets().await {
if let Some(snap) = lset.remote_alarm(name).await {
return Some(snap);
}
}
return None;
};
let lset = self.inner.link_sets.read().await.get(scheme)?;
lset.remote_alarm(body).await
}
pub async fn external_link_metadata(
&self,
name: &str,
) -> Option<crate::server::database::LinkMetadata> {
let (scheme, body) = if let Some(rest) = name.strip_prefix("pva://") {
("pva", rest)
} else if let Some(rest) = name.strip_prefix("ca://") {
("ca", rest)
} else {
for lset in self.registered_link_sets().await {
if let Some(meta) = lset.link_metadata(name).await {
return Some(meta);
}
}
return None;
};
let lset = self.inner.link_sets.read().await.get(scheme)?;
lset.link_metadata(body).await
}
pub async fn link_metadata(
&self,
link: &crate::server::record::ParsedLink,
visited: &mut HashSet<String>,
) -> Option<crate::server::database::LinkMetadata> {
use crate::server::record::ParsedLink;
match link {
ParsedLink::Db(db) => self.db_link_metadata(db, visited).await,
ParsedLink::Ca(_) | ParsedLink::Pva(_) | ParsedLink::PvaJson(_) => {
let name = link
.external_pv_name()
.expect("Ca/Pva/PvaJson link carries a PV name");
self.external_link_metadata(&name).await
}
_ => None,
}
}
async fn db_link_metadata(
&self,
db: &crate::server::record::DbLink,
visited: &mut HashSet<String>,
) -> Option<crate::server::database::LinkMetadata> {
let key = format!("{}.{}", db.record, db.field);
if !visited.insert(key.clone()) {
return None;
}
let meta = self.db_target_metadata(db).await;
visited.remove(&key);
meta
}
async fn db_target_metadata(
&self,
db: &crate::server::record::DbLink,
) -> Option<crate::server::database::LinkMetadata> {
if !self.has_name_no_resolve(&db.record).await {
let pv = if db.field == "VAL" {
db.record.clone()
} else {
format!("{}.{}", db.record, db.field)
};
return self.external_link_metadata(&pv).await;
}
let record = self.get_record(&db.record).await?;
let snapshot = record.read().await.snapshot_for_field(&db.field)?;
Some(crate::server::database::LinkMetadata {
dbf_type: Some(field_type_to_link_dbf(snapshot.value.db_field_type())),
element_count: Some(snapshot.value.count() as i64),
graphic_limits: Some(snapshot.graphic_limits().unwrap_or((0.0, 0.0))),
control_limits: Some(snapshot.control_limits().unwrap_or((0.0, 0.0))),
alarm_limits: Some(snapshot.alarm_limits().unwrap_or((
f64::NAN,
f64::NAN,
f64::NAN,
f64::NAN,
))),
precision: Some(snapshot.precision().unwrap_or(0)),
units: Some(
snapshot
.units()
.map(|u| u.as_str_lossy().into_owned())
.unwrap_or_default(),
),
description: None,
})
}
pub(crate) async fn process_passive_db_source(
&self,
db: &crate::server::record::DbLink,
visited: &mut HashSet<String>,
depth: usize,
) {
if db.policy != crate::server::record::LinkProcessPolicy::ProcessPassive {
return;
}
if let Some(src) = self.get_record(&db.record).await {
let is_passive =
src.read().await.common.scan == crate::server::record::ScanType::Passive;
if is_passive {
let _ = self
.process_record_with_links_recursive(&db.record, visited, depth + 1)
.await;
}
}
}
pub async fn read_link_value_soft(
&self,
link: &crate::server::record::ParsedLink,
is_soft: bool,
visited: &mut HashSet<String>,
depth: usize,
) -> Option<EpicsValue> {
match link {
crate::server::record::ParsedLink::Constant(_) => None,
crate::server::record::ParsedLink::Db(db) if is_soft => {
self.process_passive_db_source(db, visited, depth).await;
self.read_db_link_value(db).await
}
crate::server::record::ParsedLink::Ca(_)
| crate::server::record::ParsedLink::Pva(_)
| crate::server::record::ParsedLink::PvaJson(_)
if is_soft =>
{
let name = link
.external_pv_name()
.expect("Ca/Pva/PvaJson link carries a PV name");
self.resolve_external_pv(&name).await
}
crate::server::record::ParsedLink::Calc(calc) => self.evaluate_calc_link(calc).await,
_ => None,
}
}
pub(crate) async fn process_target(
&self,
target_name: &str,
gate: ProcessTargetGate,
src_putf: bool,
src_notify: Option<&Arc<NotifyWaitSet>>,
visited: &mut HashSet<String>,
depth: usize,
) {
let Some(target_rec) = self.get_record(target_name).await else {
return;
};
let process = {
let mut tg = target_rec.write().await;
if !gate.admits(tg.common.scan) {
return;
}
let pact = tg.is_processing();
if !pact {
tg.common.putf = src_putf;
join_put_notify(&mut tg, src_notify);
} else if src_putf && !visited.contains(target_name) {
tg.common.rpro = 1;
tg.common.putf = false;
}
!pact
};
if process {
let _ = self
.process_record_with_links_recursive(target_name, visited, depth + 1)
.await;
}
}
pub(crate) async fn write_db_link_value(
&self,
link: &crate::server::record::DbLink,
value: EpicsValue,
src: OutLinkSrc<'_>,
visited: &mut HashSet<String>,
depth: usize,
) -> bool {
let target_name = if link.field == "VAL" {
link.record.clone()
} else {
format!("{}.{}", link.record, link.field)
};
if !self.has_name_no_resolve(&link.record).await {
let op = Self::external_put_op(src.notify);
if let Err(e) = self.write_external_pv(&target_name, value, op).await {
eprintln!("OUT-link write to external PV '{target_name}' failed: {e}");
return true;
}
return false;
}
let put_result = self.put_pv_already_locked(&target_name, value).await;
if link.monitor_switch != crate::server::record::MonitorSwitch::NoMaximize {
if let Some(target_rec) = self.get_record(&link.record).await {
let mut tg = target_rec.write().await;
inherit_sevr_msg(&mut tg.common, link.monitor_switch, src.alarm);
}
}
if put_result.is_err() {
return true;
}
let gate = if link.field == "PROC" {
Some(ProcessTargetGate::ProcField)
} else if link.policy == crate::server::record::LinkProcessPolicy::ProcessPassive {
Some(ProcessTargetGate::ScanPassive)
} else {
None
};
if let Some(gate) = gate {
self.process_target(&link.record, gate, src.putf, src.notify, visited, depth)
.await;
}
false
}
pub(crate) async fn write_external_pv(
&self,
name: &str,
value: EpicsValue,
op: LinkPutOp,
) -> Result<(), String> {
let (scheme, body) = if let Some(rest) = name.strip_prefix("pva://") {
("pva", rest)
} else if let Some(rest) = name.strip_prefix("ca://") {
("ca", rest)
} else {
let lsets = self.registered_link_sets().await;
if lsets.is_empty() {
return Err(format!("no link set registered for external link '{name}'"));
}
let mut last_err = String::new();
for lset in lsets {
match lset.put_value(name, value.clone(), op).await {
Ok(()) => {
lset.flush_puts().await;
return Ok(());
}
Err(e) => last_err = e,
}
}
return Err(last_err);
};
let lset = self
.inner
.link_sets
.read()
.await
.get(scheme)
.ok_or_else(|| format!("no '{scheme}' link set registered for '{name}'"))?;
let result = lset.put_value(body, value, op).await;
if result.is_ok() {
lset.flush_puts().await;
}
result
}
pub(crate) async fn scan_forward_external_pv(&self, name: &str) -> Result<(), String> {
let (scheme, body) = if let Some(rest) = name.strip_prefix("pva://") {
("pva", rest)
} else if let Some(rest) = name.strip_prefix("ca://") {
("ca", rest)
} else {
let lsets = self.registered_link_sets().await;
if lsets.is_empty() {
return Err(format!("no link set registered for forward link '{name}'"));
}
let mut last_err = String::new();
for lset in lsets {
match lset.scan_forward(name).await {
Ok(()) => return Ok(()),
Err(e) => last_err = e,
}
}
return Err(last_err);
};
let lset = self
.inner
.link_sets
.read()
.await
.get(scheme)
.ok_or_else(|| format!("no '{scheme}' link set registered for '{name}'"))?;
lset.scan_forward(body).await
}
fn external_put_op(src_notify: Option<&Arc<NotifyWaitSet>>) -> LinkPutOp {
if src_notify.is_some() {
LinkPutOp::Async
} else {
LinkPutOp::Plain
}
}
pub(crate) async fn resolve_out_target(
&self,
link: &crate::server::record::ParsedLink,
) -> OutTarget {
let external = |name: String| async move {
match self.external_link_metadata(&name).await {
Some(m) => {
let field_type = m.dbf_type.map(link_dbf_to_field_type);
OutTarget {
field_type,
element_count: m.element_count.unwrap_or(1).max(1),
is_ca_link: true,
puts_as_string: matches!(
field_type,
Some(DbFieldType::String) | Some(DbFieldType::Enum)
),
}
}
None => OutTarget {
is_ca_link: true,
..OutTarget::UNRESOLVED
},
}
};
match link {
crate::server::record::ParsedLink::Db(db) => {
if !self.has_name_no_resolve(&db.record).await {
let target_name = if db.field == "VAL" {
db.record.clone()
} else {
format!("{}.{}", db.record, db.field)
};
return external(target_name).await;
}
let Some(target) = self.get_record(&db.record).await else {
return OutTarget::UNRESOLVED;
};
let guard = target.read().await;
let field_type = crate::server::record::record_instance::declared_field_type_of(
guard.record.as_ref(),
&db.field,
)
.or_else(|| guard.record.get_field(&db.field).map(|v| v.db_field_type()));
let element_count = match guard.record.get_field(&db.field) {
Some(v) if v.is_array() => guard
.record
.get_field("NELM")
.and_then(|n| n.as_int_i64())
.filter(|n| *n > 0)
.unwrap_or(v.count() as i64),
_ => 1,
};
OutTarget {
field_type,
element_count: element_count.max(1),
is_ca_link: false,
puts_as_string: guard.field_puts_as_string(&db.field),
}
}
crate::server::record::ParsedLink::Ca(_)
| crate::server::record::ParsedLink::Pva(_)
| crate::server::record::ParsedLink::PvaJson(_) => {
let name = link
.external_pv_name()
.expect("Ca/Pva/PvaJson link carries a PV name");
external(name.to_string()).await
}
_ => OutTarget::UNRESOLVED,
}
}
pub(crate) async fn multi_out_buffer_choice(
&self,
rec: &Arc<RwLock<RecordInstance>>,
link_field: &str,
link: &crate::server::record::ParsedLink,
staged: EpicsValue,
) -> EpicsValue {
let target = self.resolve_out_target(link).await;
let guard = rec.read().await;
guard
.record
.multi_output_buffer(link_field, staged, &target)
}
pub(crate) async fn dispatch_multi_output_values(
&self,
rec: &Arc<RwLock<RecordInstance>>,
src: OutLinkSrc<'_>,
skip_out: bool,
visited: &mut HashSet<String>,
depth: usize,
) {
let pairs = {
let instance = rec.read().await;
let links = if skip_out || multi_output_dispatch_owned(instance.record.record_type()) {
&[][..]
} else {
instance.record.multi_output_links()
};
let mut pairs = Vec::new();
for &(link_field, val_field) in links {
let link_str = match instance.record.get_field(link_field) {
Some(EpicsValue::String(s)) => s,
_ => continue,
};
if link_str.is_empty() {
continue;
}
if let Some(val) = instance.record.get_field(val_field) {
pairs.push((link_field, link_str, val));
}
}
pairs
};
for (link_field, link_str, val) in pairs {
let parsed =
crate::server::record::parse_output_link_v2(link_str.as_str_lossy().as_ref());
let val = self
.multi_out_buffer_choice(rec, link_field, &parsed, val)
.await;
self.write_out_link_value(
rec,
&parsed,
val,
OutLinkSrc {
field: link_field,
..src
},
visited,
depth,
)
.await;
}
}
pub(crate) async fn write_out_link_value(
&self,
src_rec: &Arc<RwLock<RecordInstance>>,
link: &crate::server::record::ParsedLink,
value: EpicsValue,
src: OutLinkSrc<'_>,
visited: &mut HashSet<String>,
depth: usize,
) -> bool {
let failed = match link {
crate::server::record::ParsedLink::Db(db) => {
self.write_db_link_value(db, value, src, visited, depth)
.await
}
crate::server::record::ParsedLink::Ca(_)
| crate::server::record::ParsedLink::Pva(_)
| crate::server::record::ParsedLink::PvaJson(_) => {
let name = link
.external_pv_name()
.expect("Ca/Pva/PvaJson link carries a PV name");
let op = Self::external_put_op(src.notify);
match self.write_external_pv(&name, value, op).await {
Ok(()) => false,
Err(e) => {
eprintln!("OUT-link write to external PV '{name}' failed: {e}");
true
}
}
}
_ => false,
};
if failed {
let mut inst = src_rec.write().await;
crate::server::recgbl::rec_gbl_set_link_alarm(&mut inst.common, src.field);
}
failed
}
fn field_str(instance: &RecordInstance, field: &str) -> String {
match instance.record.get_field(field) {
Some(EpicsValue::String(s)) => s.as_str_lossy().into_owned(),
_ => String::new(),
}
}
fn field_i16(instance: &RecordInstance, field: &str) -> i16 {
instance
.record
.get_field(field)
.and_then(|v| v.to_f64())
.unwrap_or(0.0) as i16
}
fn field_u16(instance: &RecordInstance, field: &str) -> u16 {
match instance.record.get_field(field) {
Some(EpicsValue::UShort(v)) => v,
Some(other) => other.as_int_i64().unwrap_or(0) as u16,
None => 0,
}
}
async fn apply_selm_alarm(
rec: &Arc<RwLock<RecordInstance>>,
alarm: Option<(u16, AlarmSeverity)>,
) {
let Some((stat, sevr)) = alarm else {
return;
};
let posted = {
let mut inst = rec.write().await;
if (sevr as u16) > (inst.common.sevr as u16) {
inst.common.sevr = sevr;
inst.common.stat = stat;
true
} else {
false
}
};
if posted {
let mut inst = rec.write().await;
inst.notify_field("SEVR", crate::server::recgbl::EventMask::ALARM);
inst.notify_field("STAT", crate::server::recgbl::EventMask::VALUE);
}
}
pub(crate) async fn dispatch_multi_output(
&self,
rec: &Arc<RwLock<RecordInstance>>,
phase: MultiOutPhase,
visited: &mut HashSet<String>,
depth: usize,
) -> Option<(u16, AlarmSeverity)> {
let record_type = rec.read().await.record.record_type().to_string();
let is_value_phase = matches!(phase, MultiOutPhase::Output { .. });
if matches!(multi_out_phase_of(&record_type), MultiOutPhaseKind::Output) != is_value_phase {
return None;
}
let (src_putf, src_notify, src_alarm) = {
let guard = rec.read().await;
(
guard.common.putf,
guard.notify.clone(),
LinkAlarm::pending(&guard.common),
)
};
let out_src = OutLinkSrc {
putf: src_putf,
notify: src_notify.as_ref(),
alarm: &src_alarm,
field: "",
};
{
let sell = {
let instance = rec.read().await;
match instance.record.record_type() {
"fanout" | "dfanout" => Some(Self::field_str(&instance, "SELL")),
"seq" if Self::field_i16(&instance, "SELM") != 0 => {
Some(Self::field_str(&instance, "SELL"))
}
_ => None,
}
};
if let Some(sell) = sell {
if !sell.is_empty() {
let parsed = crate::server::record::parse_link_v2(&sell);
if let Some(val) = self.fetch_link(rec, &parsed).await.value() {
let seln = dbr_ushort_cast(&val);
let mut instance = rec.write().await;
let _ = instance.record.put_field("SELN", EpicsValue::UShort(seln));
}
}
}
}
let dispatch_info: Option<(SelmResult, MultiOut, Option<EpicsValue>)> = {
let instance = rec.read().await;
match instance.record.record_type() {
"fanout" => {
let selm = Self::field_i16(&instance, "SELM");
let seln = Self::field_u16(&instance, "SELN");
let offs = Self::field_i16(&instance, "OFFS");
let shft = Self::field_i16(&instance, "SHFT");
let links: Vec<String> = LNK_LINK_FIELDS
.iter()
.map(|f| Self::field_str(&instance, f))
.collect();
let sel = select_link_indices_ex(
SelmKind::FanoutSeq,
selm,
seln,
offs,
shft,
links.len(),
);
Some((sel, MultiOut::Fanout(links), None))
}
"dfanout" => {
let selm = Self::field_i16(&instance, "SELM");
let seln = Self::field_u16(&instance, "SELN");
let val = match phase {
MultiOutPhase::Output { skip_out: true } => None,
MultiOutPhase::Output { skip_out: false } => instance.record.val(),
MultiOutPhase::ForwardLink => return None,
};
let links: Vec<String> = DFANOUT_LINK_FIELDS
.iter()
.map(|f| Self::field_str(&instance, f))
.collect();
let sel =
select_link_indices_ex(SelmKind::Dfanout, selm, seln, 0, 0, links.len());
Some((sel, MultiOut::Dfanout(links), val))
}
"seq" => {
let selm = Self::field_i16(&instance, "SELM");
let seln = Self::field_u16(&instance, "SELN");
let offs = Self::field_i16(&instance, "OFFS");
let shft = Self::field_i16(&instance, "SHFT");
let dol_names = [
"DOL0", "DOL1", "DOL2", "DOL3", "DOL4", "DOL5", "DOL6", "DOL7", "DOL8",
"DOL9", "DOLA", "DOLB", "DOLC", "DOLD", "DOLE", "DOLF",
];
let lnk_names = LNK_LINK_FIELDS;
let dly_names = [
"DLY0", "DLY1", "DLY2", "DLY3", "DLY4", "DLY5", "DLY6", "DLY7", "DLY8",
"DLY9", "DLYA", "DLYB", "DLYC", "DLYD", "DLYE", "DLYF",
];
let do_names = [
"DO0", "DO1", "DO2", "DO3", "DO4", "DO5", "DO6", "DO7", "DO8", "DO9",
"DOA", "DOB", "DOC", "DOD", "DOE", "DOF",
];
let groups: Vec<SeqGroup> = (0..16)
.map(|i| SeqGroup {
dol: Self::field_str(&instance, dol_names[i]),
lnk: Self::field_str(&instance, lnk_names[i]),
dly: instance
.record
.get_field(dly_names[i])
.and_then(|v| v.to_f64())
.unwrap_or(0.0),
dov: instance
.record
.get_field(do_names[i])
.and_then(|v| v.to_f64())
.unwrap_or(0.0),
})
.collect();
let sel = select_link_indices_ex(
SelmKind::FanoutSeq,
selm,
seln,
offs,
shft,
groups.len(),
);
Some((sel, MultiOut::Seq(groups), None))
}
_ => None,
}
};
let (sel, payload, val) = match dispatch_info {
Some(info) => info,
None => return None,
};
debug_assert!(sel.indices.iter().all(|&i| i < payload.len()));
debug_assert!(multi_output_dispatch_owned(
rec.read().await.record.record_type()
));
let pending_selm_alarm = sel.alarm;
if !is_value_phase {
Self::apply_selm_alarm(rec, sel.alarm).await;
}
let indices = sel.indices;
let mut link_failed = false;
match payload {
MultiOut::Fanout(links) => {
for idx in indices {
let link_str = &links[idx];
if link_str.is_empty() {
continue;
}
let parsed = crate::server::record::parse_forward_link_v2(link_str);
if let crate::server::record::ParsedLink::Db(ref db) = parsed {
self.process_target(
&db.record,
ProcessTargetGate::ScanPassive,
src_putf,
src_notify.as_ref(),
visited,
depth,
)
.await;
}
}
}
MultiOut::Dfanout(links) => {
if let Some(ref val) = val {
for idx in indices {
let link_str = &links[idx];
if link_str.is_empty() {
continue;
}
let parsed = crate::server::record::parse_output_link_v2(link_str);
if self
.write_out_link_value(
rec,
&parsed,
val.clone(),
OutLinkSrc {
field: DFANOUT_LINK_FIELDS[idx],
..out_src
},
visited,
depth,
)
.await
{
link_failed = true;
}
}
}
}
MultiOut::Seq(groups) => {
const DO_NAMES: [&str; 16] = [
"DO0", "DO1", "DO2", "DO3", "DO4", "DO5", "DO6", "DO7", "DO8", "DO9", "DOA",
"DOB", "DOC", "DOD", "DOE", "DOF",
];
let is_real = |s: &str| {
!s.is_empty()
&& !matches!(
crate::server::record::parse_link_v2(s),
crate::server::record::ParsedLink::Constant(_)
)
};
let rec_name = rec.read().await.name.clone();
for idx in indices {
let grp = &groups[idx];
let lnk_real = is_real(&grp.lnk);
let dol_real = is_real(&grp.dol);
if !lnk_real && !dol_real {
continue;
}
if grp.dly > 0.0 {
tokio::time::sleep(std::time::Duration::from_secs_f64(grp.dly)).await;
}
let new_dov = if dol_real {
let dol_parsed = crate::server::record::parse_link_v2(&grp.dol);
self.read_link_value(&dol_parsed, visited, depth)
.await
.and_then(|v| v.to_f64())
.unwrap_or(grp.dov)
} else {
grp.dov
};
if lnk_real {
let lnk_parsed = crate::server::record::parse_output_link_v2(&grp.lnk);
self.write_out_link_value(
rec,
&lnk_parsed,
EpicsValue::Double(new_dov),
OutLinkSrc {
field: LNK_LINK_FIELDS[idx],
..out_src
},
visited,
depth,
)
.await;
}
if new_dov != grp.dov {
let _ = self
.post_fields(
&rec_name,
vec![(DO_NAMES[idx].to_string(), EpicsValue::Double(new_dov))],
)
.await;
}
}
}
}
if is_value_phase {
let link_alarm = if link_failed {
Some((
crate::server::recgbl::alarm_status::LINK_ALARM,
AlarmSeverity::Major,
))
} else {
None
};
return match (pending_selm_alarm, link_alarm) {
(Some(a), Some(b)) => Some(if (a.1 as u16) >= (b.1 as u16) { a } else { b }),
(a, b) => a.or(b),
};
}
None
}
pub(crate) async fn dispatch_event_record(&self, rec: &Arc<RwLock<RecordInstance>>) {
let event_name = {
let instance = rec.read().await;
if instance.record.record_type() != "event" {
return;
}
match instance.record.get_field("VAL") {
Some(EpicsValue::String(s)) => s.as_str_lossy().into_owned(),
_ => return,
}
};
if event_name.trim().is_empty() {
return;
}
let db = self.clone();
crate::runtime::task::spawn(async move {
db.post_event_named(&event_name).await;
});
}
pub async fn register_cp_link(
&self,
source_record: &str,
target_record: &str,
passive_only: bool,
) {
let source = self
.resolve_alias(source_record)
.await
.unwrap_or_else(|| source_record.to_string());
let target = self
.resolve_alias(target_record)
.await
.unwrap_or_else(|| target_record.to_string());
let mut cp = self.inner.cp_links.write().await;
let targets = cp.entry(source).or_default();
if let Some(existing) = targets.iter_mut().find(|t| t.record == target) {
existing.passive_only = existing.passive_only && passive_only;
} else {
targets.push(super::CpTarget {
record: target,
passive_only,
});
}
}
pub async fn get_cp_targets(&self, source_record: &str) -> Vec<super::CpTarget> {
self.inner
.cp_links
.read()
.await
.get(source_record)
.cloned()
.unwrap_or_default()
}
pub async fn register_external_cp_link(
&self,
external_pv: &str,
target_record: &str,
passive_only: bool,
) {
let key = external_pv
.strip_prefix("ca://")
.or_else(|| external_pv.strip_prefix("pva://"))
.unwrap_or(external_pv);
let target = self
.resolve_alias(target_record)
.await
.unwrap_or_else(|| target_record.to_string());
let mut cp = self.inner.external_cp_links.write().await;
let targets = cp.entry(key.to_string()).or_default();
if let Some(existing) = targets.iter_mut().find(|t| t.record == target) {
existing.passive_only = existing.passive_only && passive_only;
} else {
targets.push(super::CpTarget {
record: target,
passive_only,
});
}
}
pub async fn get_external_cp_targets(&self, external_pv: &str) -> Vec<super::CpTarget> {
self.inner
.external_cp_links
.read()
.await
.get(external_pv)
.cloned()
.unwrap_or_default()
}
pub async fn external_cp_pv_names(&self) -> Vec<String> {
self.inner
.external_cp_links
.read()
.await
.keys()
.cloned()
.collect()
}
async fn classify_cp_link(
&self,
field: &str,
parsed: crate::server::record::ParsedLink,
target_name: &str,
db_links: &mut Vec<(String, String, bool)>,
ext_links: &mut Vec<(String, String, bool)>,
convert_to_ca: &mut Vec<(String, crate::server::record::CaLink)>,
) {
match parsed {
crate::server::record::ParsedLink::Db(db) => {
let Some(passive_only) = db.policy.cp_passive_only() else {
return;
};
if self.has_name_no_resolve(&db.record).await {
db_links.push((db.record, target_name.to_string(), passive_only));
} else {
let pv = if db.field == "VAL" {
db.record.clone()
} else {
format!("{}.{}", db.record, db.field)
};
ext_links.push((pv.clone(), target_name.to_string(), passive_only));
convert_to_ca.push((
field.to_string(),
crate::server::record::CaLink {
pv,
monitor_switch: db.monitor_switch,
policy: db.policy,
},
));
}
}
crate::server::record::ParsedLink::Ca(ca) => {
if let Some(passive_only) = ca.policy.cp_passive_only() {
ext_links.push((ca.pv, target_name.to_string(), passive_only));
}
}
_ => {}
}
}
pub async fn setup_cp_links(&self) {
let names = self.all_record_names().await;
let mut db_links: Vec<(String, String, bool)> = Vec::new();
let mut ext_links: Vec<(String, String, bool)> = Vec::new();
for target_name in &names {
let mut convert_to_ca: Vec<(String, crate::server::record::CaLink)> = Vec::new();
for (field, _raw, parsed) in self.record_link_fields(target_name).await {
self.classify_cp_link(
&field,
parsed,
target_name,
&mut db_links,
&mut ext_links,
&mut convert_to_ca,
)
.await;
}
if !convert_to_ca.is_empty() {
if let Some(rec_arc) = self.get_record(target_name).await {
let mut inst = rec_arc.write().await;
for (field, calink) in convert_to_ca {
let link = crate::server::record::ParsedLink::Ca(calink);
match field.as_str() {
"INP" => inst.parsed_inp = link,
"OUT" => inst.parsed_out = link,
"TSEL" => inst.parsed_tsel = link,
"SDIS" => inst.parsed_sdis = link,
_ => {}
}
}
}
}
}
let db_count = db_links.len();
for (source, target, passive_only) in db_links {
self.register_cp_link(&source, &target, passive_only).await;
}
let ext_count = ext_links.len();
for (external_pv, target, passive_only) in ext_links {
self.register_external_cp_link(&external_pv, &target, passive_only)
.await;
}
if db_count > 0 {
eprintln!("iocInit: {db_count} CP link subscriptions");
}
if ext_count > 0 {
let ext_pvs = self.external_cp_pv_names().await;
for pv in &ext_pvs {
let _ = self.resolve_external_pv(pv).await;
}
eprintln!(
"iocInit: {ext_count} external CP link subscriptions ({} PVs warmed)",
ext_pvs.len()
);
}
}
}
#[cfg(test)]
mod out_link_put_fail_tests {
use super::{LinkAlarm, OutLinkSrc};
use crate::server::database::PvDatabase;
use crate::server::record::{AlarmSeverity, DbLink, LinkProcessPolicy, MonitorSwitch};
use crate::server::records::calc::CalcRecord;
use crate::types::EpicsValue;
use std::collections::HashSet;
#[tokio::test]
async fn pp_out_link_failed_write_does_not_process_target() {
let db = PvDatabase::new();
db.add_record("TGT", Box::new(CalcRecord::new("7")))
.await
.unwrap();
assert!(
matches!(db.get_pv("TGT.VAL").await.unwrap(), EpicsValue::Double(v) if v == 0.0),
"calc VAL must start at its Default 0.0 before any process"
);
let link = DbLink {
record: "TGT".to_string(),
field: "PINI".to_string(),
policy: LinkProcessPolicy::ProcessPassive, monitor_switch: MonitorSwitch::NoMaximize,
};
let alarm = LinkAlarm {
stat: 0,
sevr: AlarmSeverity::NoAlarm,
amsg: String::new(),
};
let src = OutLinkSrc {
putf: false,
notify: None,
alarm: &alarm,
field: "OUT",
};
let mut visited = HashSet::new();
db.write_db_link_value(
&link,
EpicsValue::String("NOT_A_MENU_CHOICE".into()),
src,
&mut visited,
0,
)
.await;
assert!(
matches!(db.get_pv("TGT.VAL").await.unwrap(), EpicsValue::Double(v) if v == 0.0),
"a failed OUT-link write must NOT process the PP target \
(VAL must stay 0.0, not become 7.0)"
);
}
#[tokio::test]
async fn pp_out_link_empty_array_alarms_target_and_still_processes_it() {
use crate::server::recgbl::alarm_status;
let db = PvDatabase::new();
db.add_record("ETGT", Box::new(CalcRecord::new("7")))
.await
.unwrap();
let link = DbLink {
record: "ETGT".to_string(),
field: "VAL".to_string(),
policy: LinkProcessPolicy::ProcessPassive, monitor_switch: MonitorSwitch::NoMaximize,
};
let alarm = LinkAlarm {
stat: 0,
sevr: AlarmSeverity::NoAlarm,
amsg: String::new(),
};
let src = OutLinkSrc {
putf: false,
notify: None,
alarm: &alarm,
field: "OUT",
};
let mut visited = HashSet::new();
db.write_db_link_value(&link, EpicsValue::DoubleArray(vec![]), src, &mut visited, 0)
.await;
assert!(
matches!(db.get_pv("ETGT.VAL").await.unwrap(), EpicsValue::Double(v) if v == 7.0),
"an empty-array put is accepted by C, so the PP target must process"
);
let inst = db.get_record("ETGT").await.unwrap();
let inst = inst.read().await;
assert_eq!(inst.common.stat, alarm_status::LINK_ALARM);
assert_eq!(inst.common.sevr, AlarmSeverity::Invalid);
}
}
#[cfg(test)]
mod nonlocal_db_link_write_tests {
use super::OutLinkSrc;
use crate::server::database::{LinkPutOp, LinkSet, PvDatabase};
use crate::server::record::{AlarmSeverity, DbLink, LinkProcessPolicy, MonitorSwitch};
use crate::server::records::calc::CalcRecord;
use crate::types::EpicsValue;
use std::collections::HashSet;
use std::sync::{Arc, Mutex};
struct RecordingLset {
puts: Arc<Mutex<Vec<(String, EpicsValue)>>>,
}
#[async_trait::async_trait]
impl LinkSet for RecordingLset {
async fn is_connected(&self, _: &str) -> bool {
true
}
async fn get_value(&self, _: &str) -> Option<EpicsValue> {
None
}
async fn put_value(
&self,
name: &str,
value: EpicsValue,
_op: LinkPutOp,
) -> Result<(), String> {
self.puts.lock().unwrap().push((name.to_string(), value));
Ok(())
}
}
fn out_src(alarm: &super::LinkAlarm) -> OutLinkSrc<'_> {
OutLinkSrc {
putf: false,
notify: None,
alarm,
field: "OUT",
}
}
fn no_alarm() -> super::LinkAlarm {
super::LinkAlarm {
stat: 0,
sevr: AlarmSeverity::NoAlarm,
amsg: String::new(),
}
}
#[tokio::test]
async fn nonlocal_db_out_link_writes_through_external_put() {
let db = PvDatabase::new();
let puts = Arc::new(Mutex::new(Vec::new()));
db.register_link_set("ca", Arc::new(RecordingLset { puts: puts.clone() }))
.await;
let link = DbLink {
record: "OTHER:PV".to_string(),
field: "VAL".to_string(),
policy: LinkProcessPolicy::NoProcess,
monitor_switch: MonitorSwitch::NoMaximize,
};
let alarm = no_alarm();
let mut visited = HashSet::new();
db.write_db_link_value(
&link,
EpicsValue::Double(42.0),
out_src(&alarm),
&mut visited,
0,
)
.await;
let captured = puts.lock().unwrap();
assert_eq!(
captured.len(),
1,
"non-local OUT-link write must reach the external put path exactly once"
);
assert_eq!(captured[0].0, "OTHER:PV");
assert!(matches!(captured[0].1, EpicsValue::Double(v) if v == 42.0));
}
#[tokio::test]
async fn local_db_out_link_writes_local_not_external() {
let db = PvDatabase::new();
let puts = Arc::new(Mutex::new(Vec::new()));
db.register_link_set("ca", Arc::new(RecordingLset { puts: puts.clone() }))
.await;
db.add_record("TGT", Box::new(CalcRecord::new("0")))
.await
.unwrap();
let link = DbLink {
record: "TGT".to_string(),
field: "VAL".to_string(),
policy: LinkProcessPolicy::NoProcess,
monitor_switch: MonitorSwitch::NoMaximize,
};
let alarm = no_alarm();
let mut visited = HashSet::new();
db.write_db_link_value(
&link,
EpicsValue::Double(7.0),
out_src(&alarm),
&mut visited,
0,
)
.await;
assert!(
puts.lock().unwrap().is_empty(),
"a local OUT-link write must not reach the external put path"
);
assert!(
matches!(db.get_pv("TGT.VAL").await.unwrap(), EpicsValue::Double(v) if v == 7.0),
"the local target must hold the written value"
);
}
struct ForwardingLset {
forwards: Arc<Mutex<Vec<String>>>,
connected: bool,
}
#[async_trait::async_trait]
impl LinkSet for ForwardingLset {
async fn is_connected(&self, _: &str) -> bool {
self.connected
}
async fn get_value(&self, _: &str) -> Option<EpicsValue> {
None
}
async fn scan_forward(&self, name: &str) -> Result<(), String> {
self.forwards.lock().unwrap().push(name.to_string());
if self.connected {
Ok(())
} else {
Err("Disconn".into())
}
}
}
#[tokio::test]
async fn external_forward_link_dispatches_through_scan_forward() {
let db = PvDatabase::new();
let forwards = Arc::new(Mutex::new(Vec::new()));
db.register_link_set(
"pva",
Arc::new(ForwardingLset {
forwards: forwards.clone(),
connected: true,
}),
)
.await;
db.scan_forward_external_pv("pva://OTHER:PROC")
.await
.expect("a connected forward must succeed");
let captured = forwards.lock().unwrap();
assert_eq!(captured.len(), 1);
assert_eq!(captured[0], "OTHER:PROC");
}
#[tokio::test]
async fn record_processing_fires_external_forward_link() {
let db = PvDatabase::new();
let forwards = Arc::new(Mutex::new(Vec::new()));
db.register_link_set(
"pva",
Arc::new(ForwardingLset {
forwards: forwards.clone(),
connected: true,
}),
)
.await;
db.add_record("SRC", Box::new(CalcRecord::new("0")))
.await
.unwrap();
if let Some(rec) = db.get_record("SRC").await {
let mut inst = rec.write().await;
inst.put_common_field("FLNK", EpicsValue::String("pva://TARGET".into()))
.unwrap();
}
let mut visited = HashSet::new();
db.process_record_with_links("SRC", &mut visited, 0)
.await
.unwrap();
let captured = forwards.lock().unwrap();
assert_eq!(
captured.len(),
1,
"an external pva:// FLNK must fire scan_forward exactly once"
);
assert_eq!(captured[0], "TARGET");
}
#[tokio::test]
async fn disconnected_external_forward_link_raises_pending_link_invalid() {
let db = PvDatabase::new();
let forwards = Arc::new(Mutex::new(Vec::new()));
db.register_link_set(
"pva",
Arc::new(ForwardingLset {
forwards: forwards.clone(),
connected: false,
}),
)
.await;
db.add_record("SRC", Box::new(CalcRecord::new("0")))
.await
.unwrap();
if let Some(rec) = db.get_record("SRC").await {
let mut inst = rec.write().await;
inst.put_common_field("FLNK", EpicsValue::String("pva://TARGET".into()))
.unwrap();
}
let mut visited = HashSet::new();
db.process_record_with_links("SRC", &mut visited, 0)
.await
.unwrap();
assert_eq!(
forwards.lock().unwrap().len(),
1,
"scan_forward is still attempted on a disconnected link"
);
let rec = db.get_record("SRC").await.unwrap();
let inst = rec.read().await;
assert_eq!(inst.common.nsev, AlarmSeverity::Invalid);
assert_eq!(
inst.common.nsta,
crate::server::recgbl::alarm_status::LINK_ALARM
);
assert_eq!(inst.common.namsg, "Disconn");
}
}
#[cfg(test)]
mod cp_link_locality_tests {
use crate::server::database::PvDatabase;
use crate::server::record::ParsedLink;
use crate::server::records::ai::AiRecord;
#[tokio::test]
async fn cp_link_to_nonlocal_target_forced_external() {
let db = PvDatabase::new();
db.add_record("HOLDER", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
{
let rec = db.get_record("HOLDER").await.unwrap();
rec.write().await.common.inp = "OTHER:PV CP".to_string();
}
db.setup_cp_links().await;
let inp = db
.get_record("HOLDER")
.await
.unwrap()
.read()
.await
.parsed_inp
.clone();
match inp {
ParsedLink::Ca(ca) => assert_eq!(
ca.pv, "OTHER:PV",
"non-local CP link must carry the verbatim PV name"
),
other => panic!("non-local CP link must be forced to Ca, got {other:?}"),
}
assert!(
db.external_cp_pv_names()
.await
.contains(&"OTHER:PV".to_string()),
"non-local CP link must be registered as an external CP trigger"
);
}
#[tokio::test]
async fn cp_link_to_local_target_stays_db() {
let db = PvDatabase::new();
db.add_record("SRC", Box::new(AiRecord::new(1.0)))
.await
.unwrap();
db.add_record("HOLDER", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
{
let rec = db.get_record("HOLDER").await.unwrap();
let mut inst = rec.write().await;
inst.common.inp = "SRC CP".to_string();
inst.parsed_inp = crate::server::record::parse_link_v2("SRC CP");
}
db.setup_cp_links().await;
let inp = db
.get_record("HOLDER")
.await
.unwrap()
.read()
.await
.parsed_inp
.clone();
assert!(
matches!(inp, ParsedLink::Db(_)),
"local CP link must stay a Db link, got {inp:?}"
);
assert!(
!db.external_cp_pv_names().await.contains(&"SRC".to_string()),
"a local CP target must not be registered as an external CP link"
);
}
}
#[cfg(test)]
mod nonlocal_db_link_read_tests {
use crate::server::database::PvDatabase;
use crate::server::record::{ParsedLink, parse_link_v2};
use crate::types::EpicsValue;
use std::collections::HashSet;
use std::sync::Arc;
#[tokio::test]
async fn plain_nonlocal_db_link_reads_via_external_resolver() {
let db = PvDatabase::new();
db.set_external_resolver(Arc::new(|name: &str| {
let hit = name == "OTHER:PV";
Box::pin(async move {
if hit {
Some(EpicsValue::Double(42.0))
} else {
None
}
})
}))
.await;
let link = parse_link_v2("OTHER:PV");
assert!(
matches!(link, ParsedLink::Db(_)),
"a bare non-scheme name parses to a Db link, got {link:?}"
);
let mut visited = HashSet::new();
let v = db.read_link_value(&link, &mut visited, 0).await;
assert_eq!(
v,
Some(EpicsValue::Double(42.0)),
"a plain non-local Db link must resolve through the external resolver (C dbCaAddLink fallback)"
);
}
#[tokio::test]
async fn nonlocal_db_link_value_and_alarm_via_external() {
let db = PvDatabase::new();
db.set_external_resolver(Arc::new(|name: &str| {
let hit = name == "OTHER:PV";
Box::pin(async move {
if hit {
Some(EpicsValue::Double(5.0))
} else {
None
}
})
}))
.await;
let link = parse_link_v2("OTHER:PV");
let (fetch, alarm) = db.read_link_with_alarm(&link).await;
assert_eq!(
fetch.value(),
Some(EpicsValue::Double(5.0)),
"non-local Db link value must come from the external resolver"
);
assert!(
alarm.is_none(),
"non-local Db link must not fabricate a local-record alarm, got {alarm:?}"
);
}
#[tokio::test]
async fn local_db_link_reads_local_not_external() {
let db = PvDatabase::new();
db.add_pv("SRC", EpicsValue::Double(7.0)).await.unwrap();
db.set_external_resolver(Arc::new(|_name: &str| {
Box::pin(async { Some(EpicsValue::Double(-1.0)) })
}))
.await;
let link = parse_link_v2("SRC");
let mut visited = HashSet::new();
let v = db.read_link_value(&link, &mut visited, 0).await;
assert_eq!(
v,
Some(EpicsValue::Double(7.0)),
"a local Db link must read the local DB, not the external resolver"
);
}
}
#[cfg(test)]
mod link_metadata_tests {
use crate::server::database::PvDatabase;
use crate::server::record::parse_link_v2;
use crate::server::records::ai::AiRecord;
use crate::server::records::stringout::StringoutRecord;
use crate::server::records::waveform::WaveformRecord;
use crate::types::DbFieldType;
use std::collections::HashSet;
#[tokio::test]
async fn constant_link_reports_no_metadata() {
let db = PvDatabase::new();
let mut visited = HashSet::new();
let link = parse_link_v2("5");
assert!(
matches!(link, crate::server::record::ParsedLink::Constant(_)),
"fixture must actually be a constant link, got {link:?}"
);
assert_eq!(
db.link_metadata(&link, &mut visited).await,
None,
"a constant link has no metadata lset slots: C returns S_db_noLSET"
);
}
#[tokio::test]
async fn db_link_to_supported_field_propagates_target_limits() {
let db = PvDatabase::new();
let mut src = AiRecord::new(1.0);
src.hopr = 10.0;
src.lopr = -10.0;
src.egu = "mm".into();
src.prec = 3;
db.add_record("SRC", Box::new(src)).await.unwrap();
let mut visited = HashSet::new();
let meta = db
.link_metadata(&parse_link_v2("SRC"), &mut visited)
.await
.expect("a local db link to an existing field reports metadata");
assert_eq!(meta.graphic_limits, Some((-10.0, 10.0)));
assert_eq!(meta.units.as_deref(), Some("mm"));
assert_eq!(meta.precision, Some(3));
}
#[tokio::test]
async fn db_link_to_unsupported_field_yields_c_defaults_not_none() {
let db = PvDatabase::new();
db.add_record("STR", Box::new(StringoutRecord::new("hello")))
.await
.unwrap();
let mut visited = HashSet::new();
let meta = db
.link_metadata(&parse_link_v2("STR"), &mut visited)
.await
.expect("dbGet returns 0 even when every option is turned off");
assert_eq!(
meta.graphic_limits,
Some((0.0, 0.0)),
"no get_graphic_double: C memsets the buffer to zero and returns 0"
);
assert_eq!(
meta.control_limits,
Some((0.0, 0.0)),
"no get_control_double: C memsets the buffer to zero and returns 0"
);
let (lolo, low, high, hihi) = meta.alarm_limits.expect("alarm limits are written");
assert!(
lolo.is_nan() && low.is_nan() && high.is_nan() && hihi.is_nan(),
"no get_alarm_double: C's pre-filled epicsNAN survives, NOT zero \
(dbAccess.c:290) — got ({lolo}, {low}, {high}, {hihi})"
);
assert_eq!(meta.precision, Some(0), "no get_precision: C memsets to 0");
assert_eq!(
meta.units.as_deref(),
Some(""),
"no get_units: C memsets the units buffer"
);
}
#[tokio::test]
async fn db_link_alarm_default_is_per_slot_not_per_record() {
let db = PvDatabase::new();
db.add_record("WF", Box::new(WaveformRecord::new(8, DbFieldType::Double)))
.await
.unwrap();
let mut visited = HashSet::new();
let meta = db
.link_metadata(&parse_link_v2("WF"), &mut visited)
.await
.expect("waveform reports metadata");
let (lolo, low, high, hihi) = meta.alarm_limits.expect("alarm limits are written");
assert!(
lolo.is_nan() && low.is_nan() && high.is_nan() && hihi.is_nan(),
"waveform NULLs get_alarm_double: NaN, not zero"
);
assert!(
meta.graphic_limits.is_some(),
"waveform keeps get_graphic_double, so that slot is still served"
);
}
#[tokio::test]
async fn revisited_db_link_target_is_refused() {
let db = PvDatabase::new();
db.add_record("SRC", Box::new(AiRecord::new(1.0)))
.await
.unwrap();
let link = parse_link_v2("SRC");
let mut visited = HashSet::new();
assert!(
db.link_metadata(&link, &mut visited).await.is_some(),
"first visit resolves"
);
assert!(
visited.is_empty(),
"the guard must be CLEARED after the fetch (C clears the flag at \
dbDbLink.c:257), or a diamond onto one target would fail"
);
visited.insert("SRC.VAL".to_string());
assert_eq!(
db.link_metadata(&link, &mut visited).await,
None,
"a link already on the chain must write nothing"
);
}
#[tokio::test]
async fn db_link_to_missing_field_reports_no_metadata() {
let db = PvDatabase::new();
db.add_record("SRC", Box::new(AiRecord::new(1.0)))
.await
.unwrap();
let mut visited = HashSet::new();
assert_eq!(
db.link_metadata(&parse_link_v2("SRC.NOSUCHFIELD"), &mut visited)
.await,
None,
);
}
}