use std::collections::HashSet;
use crate::error::{CaError, CaResult};
use crate::server::snapshot::Snapshot;
use crate::types::EpicsValue;
use super::PvDatabase;
fn check_put_disabled(
instance: &crate::server::record::RecordInstance,
field_upper: &str,
) -> CaResult<()> {
if instance.common.disp != 0 && field_upper != "DISP" {
return Err(CaError::PutDisabled(field_upper.to_string()));
}
Ok(())
}
fn check_no_mod(instance: &crate::server::record::RecordInstance, field: &str) -> CaResult<()> {
if instance.is_no_mod(field) {
return Err(CaError::ReadOnlyField(field.to_string()));
}
Ok(())
}
pub(crate) fn put_drives_processing_of(
instance: &crate::server::record::RecordInstance,
field: &str,
) -> bool {
field == "PROC"
|| (instance.common.scan == crate::server::record::ScanType::Passive
&& (field == "UDF" || instance.record.processes_after_put(field)))
}
fn emit_cycle_posts(instance: &mut crate::server::record::RecordInstance) {
use crate::server::record::{CyclePostMask, EventMask};
for (sf, cycle_mask) in instance.record.take_cycle_posted_fields() {
let mask = match cycle_mask {
CyclePostMask::Value => EventMask::VALUE,
CyclePostMask::ValueLog | CyclePostMask::MonitorValueLog => {
EventMask::VALUE | EventMask::LOG
}
};
instance.notify_field(sf, mask);
}
}
fn special_after_put(
instance: &mut crate::server::record::RecordInstance,
field: &str,
out: &mut Vec<crate::server::record::ProcessAction>,
) -> CaResult<crate::server::record::CommonFieldPutResult> {
let status = instance.record.special(field, true);
out.extend(instance.record.take_special_actions());
if status.is_err() {
emit_cycle_posts(instance);
}
status?;
if let Some(value_field) =
crate::server::record::reseed_constant_input_link(&mut *instance.record, field)
{
let mask = instance.record.special_reseed_post_mask();
instance.notify_field(value_field, mask);
}
Ok(if field == "SIMM" {
instance.rec_gbl_check_simm()
} else {
crate::server::record::CommonFieldPutResult::NoChange
})
}
fn special_before_put(instance: &mut crate::server::record::RecordInstance, field: &str) {
if field == "SIMM" {
instance.rec_gbl_save_simm();
}
}
fn commit_special_reset_alarm(
instance: &mut crate::server::record::RecordInstance,
) -> crate::server::recgbl::EventMask {
use crate::server::recgbl::EventMask;
let alarm_result = crate::server::recgbl::rec_gbl_reset_alarms(&mut instance.common);
for (af, mask) in
crate::server::database::processing::alarm_field_posts(&instance.common, &alarm_result)
{
instance.notify_field(af, mask);
}
if alarm_result.alarm_changed || alarm_result.amsg_changed {
EventMask::ALARM
} else {
EventMask::NONE
}
}
fn coerce_write_value(
record: &dyn crate::server::record::Record,
field: &str,
target: crate::types::DbFieldType,
value: EpicsValue,
) -> CaResult<EpicsValue> {
crate::server::record::coerce_put_value(record, field, target, value)
}
enum PutRequest {
Write(EpicsValue),
EmptyIntoScalar,
}
fn dbput_request(
record: &dyn crate::server::record::Record,
field: &str,
value: EpicsValue,
) -> CaResult<PutRequest> {
let dest_is_array = record.get_field(field).is_some_and(|v| v.is_array());
if value.is_empty_array() && !dest_is_array {
return Ok(PutRequest::EmptyIntoScalar);
}
let target = record
.get_field(field)
.map(|v| v.db_field_type())
.or_else(|| crate::server::record::record_instance::declared_field_type_of(record, field));
let is_char_string_view = matches!(value, EpicsValue::CharArray(_))
&& target == Some(crate::types::DbFieldType::String);
let value = if !dest_is_array && value.is_array() && !is_char_string_view {
value.first_element().unwrap_or(value)
} else {
value
};
match target {
Some(target)
if value.db_field_type() != target || target == crate::types::DbFieldType::String =>
{
Ok(PutRequest::Write(coerce_write_value(
record, field, target, value,
)?))
}
_ => Ok(PutRequest::Write(value)),
}
}
fn set_empty_request_alarm(instance: &mut crate::server::record::RecordInstance) {
crate::server::recgbl::rec_gbl_set_sevr(
&mut instance.common,
crate::server::recgbl::alarm_status::LINK_ALARM,
crate::server::record::AlarmSeverity::Invalid,
);
}
enum NotifyRequest {
None,
New,
Deferred(crate::runtime::sync::oneshot::Sender<()>),
}
impl NotifyRequest {
fn wants_notify(&self) -> bool {
!matches!(self, NotifyRequest::None)
}
#[allow(clippy::type_complexity)]
fn into_completion(
self,
) -> Option<(
crate::runtime::sync::oneshot::Sender<()>,
Option<crate::runtime::sync::oneshot::Receiver<()>>,
)> {
match self {
NotifyRequest::None => None,
NotifyRequest::New => {
let (tx, rx) = crate::runtime::sync::oneshot::channel();
Some((tx, Some(rx)))
}
NotifyRequest::Deferred(tx) => Some((tx, None)),
}
}
}
fn array_nord_before_put(
instance: &crate::server::record::RecordInstance,
field: &str,
) -> Option<EpicsValue> {
if field == "VAL" {
instance.record.get_field("NORD")
} else {
None
}
}
fn post_array_info(
instance: &mut crate::server::record::RecordInstance,
old_nord: &Option<EpicsValue>,
origin: u64,
) {
let Some(old) = old_nord else { return };
let moved = instance
.record
.get_field("NORD")
.is_some_and(|new| new != *old);
if moved {
instance.notify_field_with_origin(
"NORD",
crate::server::recgbl::EventMask::VALUE | crate::server::recgbl::EventMask::LOG,
origin,
);
}
}
impl PvDatabase {
pub fn get_pv_blocking(&self, name: &str) -> CaResult<EpicsValue> {
crate::runtime::task::block_on_sync(self.get_pv(name)).unwrap_or_else(|_| {
Err(CaError::InvalidValue(
"get_pv_blocking cannot block a current-thread runtime; await get_pv() instead"
.into(),
))
})
}
pub async fn get_pv(&self, name: &str) -> CaResult<EpicsValue> {
let (base, field) = super::parse_pv_name(name);
let field = field.to_ascii_uppercase();
if let Some(pv) = self.inner.simple_pvs.read().await.get(name) {
return Ok(pv.get().await);
}
if let Some(rec) = self.get_record(base).await {
let instance = rec.read().await;
return instance
.resolve_field(&field)
.ok_or_else(|| CaError::ChannelNotFound(name.to_string()));
}
Err(CaError::ChannelNotFound(name.to_string()))
}
pub async fn put_pv(&self, name: &str, value: EpicsValue) -> CaResult<()> {
self.put_pv_inner(name, value, true).await
}
pub async fn put_pv_already_locked(&self, name: &str, value: EpicsValue) -> CaResult<()> {
self.put_pv_inner(name, value, false).await
}
pub async fn check_external_put_preconditions(
&self,
record_name: &str,
field: &str,
) -> CaResult<()> {
let field_upper = field.to_ascii_uppercase();
let Some(rec) = self.get_record(record_name).await else {
return Ok(());
};
let instance = rec.read().await;
check_no_mod(&instance, &field_upper)?;
check_put_disabled(&instance, &field_upper)?;
Ok(())
}
pub async fn put_drives_processing(&self, record_name: &str, field: &str) -> bool {
let field_upper = field.to_ascii_uppercase();
let Some(rec) = self.get_record(record_name).await else {
return false;
};
let instance = rec.read().await;
put_drives_processing_of(&instance, &field_upper)
}
async fn put_pv_inner(
&self,
name: &str,
value: EpicsValue,
acquire_gate: bool,
) -> CaResult<()> {
let (base, field) = super::parse_pv_name(name);
let field = field.to_ascii_uppercase();
if let Some(pv) = self.inner.simple_pvs.read().await.get(name) {
pv.set(value).await;
return Ok(());
}
if let Some(rec) = self.get_record(base).await {
let canonical_base: String = self
.resolve_alias(base)
.await
.unwrap_or_else(|| base.to_string());
let _record_gate = if acquire_gate {
Some(self.lock_record(&canonical_base).await)
} else {
None
};
let mut instance = rec.write().await;
check_no_mod(&instance, &field)?;
let request = dbput_request(&*instance.record, &field, value)?;
instance.record.special(&field, false)?;
let prev_value = instance.record.get_field(&field);
let old_nord = array_nord_before_put(&instance, &field);
let mut special_actions = Vec::new();
use crate::server::record::CommonFieldPutResult;
let common_result = match request {
PutRequest::EmptyIntoScalar => {
set_empty_request_alarm(&mut instance);
CommonFieldPutResult::NoChange
}
PutRequest::Write(value) => {
special_before_put(&mut instance, &field);
match instance.record.put_field(&field, value.clone()) {
Ok(()) => {
instance.record.on_put(&field);
special_after_put(&mut instance, &field, &mut special_actions)?
}
Err(CaError::FieldNotFound(_)) => {
instance.put_common_field(&field, value)?
}
Err(e) => return Err(e),
}
}
};
if instance.record.special_commits_alarms(&field) {
let _ = crate::server::recgbl::rec_gbl_reset_alarms(&mut instance.common);
}
if instance.record.special_checks_alarms(&field) {
let inst = &mut *instance;
inst.record.check_alarms(&mut inst.common);
let _ = crate::server::recgbl::rec_gbl_reset_alarms(&mut inst.common);
}
instance.notify_field_written_if_changed(&field, prev_value.as_ref());
post_array_info(&mut instance, &old_nord, 0);
drop(instance);
match common_result {
CommonFieldPutResult::ScanChanged {
old_scan,
new_scan,
phas,
} => {
self.update_scan_index(&canonical_base, old_scan, new_scan, phas, phas)
.await;
}
CommonFieldPutResult::PhasChanged {
scan: s,
old_phas,
new_phas,
} => {
self.update_scan_index(&canonical_base, s, s, old_phas, new_phas)
.await;
}
CommonFieldPutResult::NoChange => {}
}
self.run_special_actions(&canonical_base, &rec, special_actions)
.await;
if field == "ASG" {
crate::server::access_security::notify_asg_field_changed();
}
return Ok(());
}
Err(CaError::ChannelNotFound(name.to_string()))
}
pub async fn put_pv_and_post(&self, name: &str, value: EpicsValue) -> CaResult<()> {
self.put_pv_and_post_with_origin(name, value, 0).await
}
pub async fn post_alarm(&self, name: &str, severity: u16, status: u16) -> CaResult<()> {
if let Some(pv) = self.inner.simple_pvs.read().await.get(name).cloned() {
pv.post_alarm(severity, status).await;
return Ok(());
}
Err(crate::error::CaError::ChannelNotFound(name.to_string()))
}
pub async fn put_pv_and_post_snapshot(&self, name: &str, snapshot: Snapshot) -> CaResult<()> {
if let Some(pv) = self.inner.simple_pvs.read().await.get(name).cloned() {
pv.set_snapshot(snapshot).await;
return Ok(());
}
Err(CaError::ChannelNotFound(name.to_string()))
}
pub async fn set_pv_metadata(&self, name: &str, snapshot: &Snapshot) -> CaResult<()> {
if let Some(pv) = self.inner.simple_pvs.read().await.get(name).cloned() {
pv.set_metadata(metadata_from_snapshot(snapshot));
return Ok(());
}
Err(CaError::ChannelNotFound(name.to_string()))
}
pub async fn post_pv_property(&self, name: &str, snapshot: Snapshot) -> CaResult<()> {
if let Some(pv) = self.inner.simple_pvs.read().await.get(name).cloned() {
pv.set_metadata(metadata_from_snapshot(&snapshot));
pv.post_property(snapshot).await;
return Ok(());
}
Err(CaError::ChannelNotFound(name.to_string()))
}
pub async fn put_pv_and_post_with_origin(
&self,
name: &str,
value: EpicsValue,
origin: u64,
) -> CaResult<()> {
let (base, field) = super::parse_pv_name(name);
let field = field.to_ascii_uppercase();
if let Some(pv) = self.inner.simple_pvs.read().await.get(name).cloned() {
let _ = origin; pv.set(value).await;
return Ok(());
}
if let Some(rec) = self.get_record(base).await {
let canonical_base: String = self
.resolve_alias(base)
.await
.unwrap_or_else(|| base.to_string());
let _record_gate = self.lock_record(&canonical_base).await;
let mut instance = rec.write().await;
check_no_mod(&instance, &field)?;
let request = dbput_request(&*instance.record, &field, value)?;
instance.record.special(&field, false)?;
let old_value = instance.record.get_field(&field);
let old_stat = instance.common.stat;
let old_sevr = instance.common.sevr;
let old_nord = array_nord_before_put(&instance, &field);
let mut special_actions = Vec::new();
use crate::server::record::CommonFieldPutResult;
let common_result = match request {
PutRequest::EmptyIntoScalar => {
set_empty_request_alarm(&mut instance);
CommonFieldPutResult::NoChange
}
PutRequest::Write(value) => {
special_before_put(&mut instance, &field);
match instance.record.put_field(&field, value.clone()) {
Ok(()) => {
instance.record.on_put(&field);
let result =
special_after_put(&mut instance, &field, &mut special_actions)?;
if instance.record.is_udf_defining_put(&field) {
instance.common.udf = 0;
}
result
}
Err(CaError::FieldNotFound(_)) => {
instance.put_common_field(&field, value)?
}
Err(e) => return Err(e),
}
}
};
instance.notify_field_written_if_changed(&field, old_value.as_ref());
let new_value = instance.record.get_field(&field);
let value_changed = old_value != new_value;
let alarm_changed =
old_stat != instance.common.stat || old_sevr != instance.common.sevr;
let nord_changed = old_nord.is_some() && instance.record.get_field("NORD") != old_nord;
if value_changed || alarm_changed || nord_changed {
instance.common.time = crate::runtime::general_time::get_current();
instance.cleanup_subscribers();
if value_changed || alarm_changed {
instance.notify_field_with_origin(
&field,
crate::server::recgbl::EventMask::VALUE
| crate::server::recgbl::EventMask::LOG
| crate::server::recgbl::EventMask::ALARM,
origin,
);
}
post_array_info(&mut instance, &old_nord, origin);
}
drop(instance);
match common_result {
CommonFieldPutResult::ScanChanged {
old_scan,
new_scan,
phas,
} => {
self.update_scan_index(&canonical_base, old_scan, new_scan, phas, phas)
.await;
}
CommonFieldPutResult::PhasChanged {
scan: s,
old_phas,
new_phas,
} => {
self.update_scan_index(&canonical_base, s, s, old_phas, new_phas)
.await;
}
CommonFieldPutResult::NoChange => {}
}
self.run_special_actions(&canonical_base, &rec, special_actions)
.await;
if field == "ASG" {
crate::server::access_security::notify_asg_field_changed();
}
return Ok(());
}
Err(CaError::ChannelNotFound(name.to_string()))
}
async fn run_special_actions(
&self,
record_name: &str,
rec: &std::sync::Arc<crate::runtime::sync::RwLock<crate::server::record::RecordInstance>>,
actions: Vec<crate::server::record::ProcessAction>,
) {
if actions.is_empty() {
return;
}
let mut visited = HashSet::new();
Box::pin(self.execute_process_actions(record_name, rec, actions, &mut visited, 0)).await;
}
pub async fn put_record_field_from_ca(
&self,
record_name: &str,
field: &str,
value: EpicsValue,
) -> CaResult<Option<crate::runtime::sync::oneshot::Receiver<()>>> {
self.put_record_field_from_ca_inner(record_name, field, value, true, NotifyRequest::New)
.await
}
pub async fn put_record_field_from_ca_already_locked(
&self,
record_name: &str,
field: &str,
value: EpicsValue,
) -> CaResult<Option<crate::runtime::sync::oneshot::Receiver<()>>> {
self.put_record_field_from_ca_inner(record_name, field, value, false, NotifyRequest::New)
.await
}
pub async fn put_record_field_from_ca_no_notify(
&self,
record_name: &str,
field: &str,
value: EpicsValue,
) -> CaResult<()> {
self.put_record_field_from_ca_inner(record_name, field, value, true, NotifyRequest::None)
.await
.map(|_| ())
}
pub async fn put_alarm_ack_from_ca(
&self,
record_name: &str,
field: &str,
ack: crate::server::record::AlarmAck,
value: u16,
) -> CaResult<()> {
let field_upper = field.to_ascii_uppercase();
let rec = self
.get_record(record_name)
.await
.ok_or_else(|| CaError::ChannelNotFound(record_name.to_string()))?;
let canonical: String = self
.resolve_alias(record_name)
.await
.unwrap_or_else(|| record_name.to_string());
let _record_gate = self.lock_record(&canonical).await;
let mut instance = rec.write().await;
check_put_disabled(&instance, &field_upper)?;
match ack {
crate::server::record::AlarmAck::Transient => instance.put_ackt(value),
crate::server::record::AlarmAck::Severity => instance.put_acks(value),
}
Ok(())
}
pub async fn put_record_field_from_ca_no_notify_already_locked(
&self,
record_name: &str,
field: &str,
value: EpicsValue,
) -> CaResult<()> {
self.put_record_field_from_ca_inner(record_name, field, value, false, NotifyRequest::None)
.await
.map(|_| ())
}
pub async fn process_record_with_notify(
&self,
record_name: &str,
) -> CaResult<Option<crate::runtime::sync::oneshot::Receiver<()>>> {
let (completion_tx, completion_rx) = crate::runtime::sync::oneshot::channel();
let notify = crate::server::record::NotifyWaitSet::new(completion_tx);
{
let rec_arc = {
let recs = self.inner.records.read().await;
recs.get(record_name).cloned()
};
let Some(rec_arc) = rec_arc else {
return Err(CaError::ChannelNotFound(record_name.to_string()));
};
let mut guard = rec_arc.write().await;
if guard.notify.is_some() {
return Err(CaError::PutCallbackInProgress(record_name.to_string()));
}
guard.notify = Some(notify.clone());
}
let mut visited = HashSet::new();
self.process_record_with_links(record_name, &mut visited, 0)
.await?;
if notify.completed() {
Ok(None)
} else {
Ok(Some(completion_rx))
}
}
pub async fn put_driven_process(&self, record_name: &str) -> CaResult<()> {
self.put_driven_process_inner(record_name, true).await
}
pub async fn put_driven_process_already_locked(&self, record_name: &str) -> CaResult<()> {
self.put_driven_process_inner(record_name, false).await
}
async fn put_driven_process_inner(
&self,
record_name: &str,
acquire_gate: bool,
) -> CaResult<()> {
{
let Some(rec) = self.get_record(record_name).await else {
return Ok(());
};
let mut instance = rec.write().await;
if instance.is_processing() {
instance.common.rpro = 1;
return Ok(());
}
instance.common.putf = true;
}
let mut visited = HashSet::new();
if acquire_gate {
self.process_record_with_links(record_name, &mut visited, 0)
.await
} else {
self.process_record_with_links_already_locked(record_name, &mut visited, 0)
.await
}
}
pub(crate) async fn restart_deferred_notify_put(
&self,
record_name: &str,
put: crate::server::record::DeferredNotifyPut,
) {
let crate::server::record::DeferredNotifyPut {
field,
value,
completion,
} = put;
let _ = self
.put_record_field_from_ca_inner(
record_name,
&field,
value,
true,
NotifyRequest::Deferred(completion),
)
.await;
}
async fn put_record_field_from_ca_inner(
&self,
record_name: &str,
field: &str,
value: EpicsValue,
acquire_gate: bool,
notify_request: NotifyRequest,
) -> CaResult<Option<crate::runtime::sync::oneshot::Receiver<()>>> {
let field = field.to_ascii_uppercase();
let want_notify = notify_request.wants_notify();
let rec = self
.get_record(record_name)
.await
.ok_or_else(|| CaError::ChannelNotFound(record_name.to_string()))?;
let canonical_owned;
let record_name: &str = if let Some(target) = self.resolve_alias(record_name).await {
canonical_owned = target;
&canonical_owned
} else {
record_name
};
let _record_gate = if acquire_gate {
Some(self.lock_record(record_name).await)
} else {
None
};
{
let instance = rec.read().await;
check_put_disabled(&instance, &field)?;
check_no_mod(&instance, &field)?;
}
let snam_registry_reject = {
let name_to_resolve: Option<String> = {
let guard = rec.read().await;
if guard.record.is_subroutine_name_field(&field) {
match &value {
EpicsValue::String(s) => {
let name = s.as_str_lossy();
(!name.is_empty()).then(|| name.into_owned())
}
_ => None,
}
} else {
None
}
};
match name_to_resolve {
Some(name) => self.find_subroutine_named(&name).await.is_none(),
None => false,
}
};
let is_dbf_link_field = {
let guard = rec.read().await;
crate::types::dbf_link_class(guard.record.record_type(), &field).is_some()
};
if want_notify && !is_dbf_link_field {
let mut guard = rec.write().await;
if guard.is_processing() {
let Some((completion, completion_rx)) = notify_request.into_completion() else {
return Ok(None);
};
guard
.park_notify_put(crate::server::record::DeferredNotifyPut {
field,
value,
completion,
})
.map_err(|_| CaError::PutCallbackInProgress(record_name.to_string()))?;
return Ok(completion_rx);
}
}
if field == "PROC" {
let proc_store: CaResult<()> = {
let rec_arc = {
let recs = self.inner.records.read().await;
recs.get(record_name).cloned()
};
if let Some(rec_arc) = rec_arc {
let mut guard = rec_arc.write().await;
match guard.put_common_field("PROC", value) {
Ok(_) => {
guard.notify_field(
"PROC",
crate::server::recgbl::EventMask::VALUE
| crate::server::recgbl::EventMask::LOG,
);
Ok(())
}
Err(e) => Err(e),
}
} else {
Ok(())
}
};
if let Err(e) = proc_store {
if want_notify {
let _ = self.put_driven_process_already_locked(record_name).await;
}
return Err(e);
}
let parked = if let Some((completion_tx, completion_rx)) =
notify_request.into_completion()
{
let notify = crate::server::record::NotifyWaitSet::new(completion_tx);
{
let rec_arc = {
let recs = self.inner.records.read().await;
recs.get(record_name).cloned()
};
if let Some(rec_arc) = rec_arc {
let mut guard = rec_arc.write().await;
if guard.notify.is_some() {
return Err(CaError::PutCallbackInProgress(record_name.to_string()));
}
guard.notify = Some(notify.clone());
}
}
Some((notify, completion_rx))
} else {
None
};
let _ = self.put_driven_process_already_locked(record_name).await;
return match parked {
Some((notify, completion_rx)) => {
if notify.completed() {
Ok(None)
} else {
Ok(completion_rx)
}
}
None => Ok(None),
};
}
let mut special_actions = Vec::new();
let mut instance = rec.write().await;
let block_result: CaResult<crate::server::record::CommonFieldPutResult> = (|| {
let request = dbput_request(&*instance.record, &field, value)?;
instance.record.special(&field, false)?;
special_before_put(&mut instance, &field);
let prev_value = instance.record.get_field(&field);
let old_nord = array_nord_before_put(&instance, &field);
use crate::server::record::CommonFieldPutResult;
let common_result = match request {
PutRequest::EmptyIntoScalar => {
set_empty_request_alarm(&mut instance);
CommonFieldPutResult::NoChange
}
PutRequest::Write(value) => {
match instance.record.put_field(&field, value.clone()) {
Ok(()) => {
instance.record.on_put(&field);
let result =
special_after_put(&mut instance, &field, &mut special_actions)?;
if snam_registry_reject {
return Err(CaError::BadField("SNAM: Subroutine not found".into()));
}
if instance.record.is_udf_defining_put(&field) {
instance.common.udf = 0;
}
result
}
Err(CaError::FieldNotFound(_)) => {
instance.put_common_field(&field, value)?
}
Err(e) => return Err(e),
}
}
};
if instance.record.special_checks_alarms(&field) {
let inst = &mut *instance;
inst.record.check_alarms(&mut inst.common);
let _ = crate::server::recgbl::rec_gbl_reset_alarms(&mut inst.common);
}
instance.notify_field_written_if_changed(&field, prev_value.as_ref());
instance.cleanup_subscribers();
let suppress_value_field_post = field == instance.record.primary_field()
&& instance
.record
.process_passive_fields()
.iter()
.any(|f| f.eq_ignore_ascii_case(&field));
if !suppress_value_field_post {
instance.notify_field(
&field,
crate::server::recgbl::EventMask::VALUE | crate::server::recgbl::EventMask::LOG,
);
}
post_array_info(&mut instance, &old_nord, 0);
let side_effect_alarm_mask = if instance.record.special_commits_alarms(&field) {
commit_special_reset_alarm(&mut instance)
} else {
crate::server::recgbl::EventMask::NONE
};
let side_effect_value_only = instance.record.value_only_change_fields();
for sf in instance.record.monitor_side_effect_fields(&field) {
use crate::server::recgbl::EventMask;
let mask = if side_effect_value_only
.iter()
.any(|f| f.eq_ignore_ascii_case(sf))
{
EventMask::VALUE
} else {
EventMask::VALUE | EventMask::LOG
};
instance.notify_field(sf, mask | side_effect_alarm_mask);
}
emit_cycle_posts(&mut instance);
Ok(common_result)
})();
let common_result = match block_result {
Ok(cr) => {
drop(instance);
cr
}
Err(e) => {
if instance.record.special_commits_alarms(&field) {
let _ = instance.record.special(&field, true);
let alarm_mask = commit_special_reset_alarm(&mut instance);
let value_only = instance.record.value_only_change_fields();
for sf in instance.record.monitor_side_effect_fields(&field) {
use crate::server::recgbl::EventMask;
let mask = if value_only.iter().any(|f| f.eq_ignore_ascii_case(sf)) {
EventMask::VALUE
} else {
EventMask::VALUE | EventMask::LOG
};
instance.notify_field(sf, mask | alarm_mask);
}
}
if instance.record.special_checks_alarms(&field) {
let inst = &mut *instance;
inst.record.check_alarms(&mut inst.common);
let _ = crate::server::recgbl::rec_gbl_reset_alarms(&mut inst.common);
}
if instance.record.is_udf_defining_put(&field)
&& field != instance.record.primary_field()
{
instance.common.udf = 0;
}
if want_notify && put_drives_processing_of(&instance, &field) {
drop(instance);
let _ = self.put_driven_process_already_locked(record_name).await;
}
return Err(e);
}
};
if field == "ASG" {
crate::server::access_security::notify_asg_field_changed();
}
self.run_special_actions(record_name, &rec, std::mem::take(&mut special_actions))
.await;
match common_result {
crate::server::record::CommonFieldPutResult::ScanChanged {
old_scan,
new_scan,
phas,
} => {
self.update_scan_index(record_name, old_scan, new_scan, phas, phas)
.await;
}
crate::server::record::CommonFieldPutResult::PhasChanged {
scan: s,
old_phas,
new_phas,
} => {
self.update_scan_index(record_name, s, s, old_phas, new_phas)
.await;
}
crate::server::record::CommonFieldPutResult::NoChange => {}
}
let should_process = {
let instance = rec.read().await;
put_drives_processing_of(&instance, &field)
};
if !should_process {
return Ok(None);
}
let parked = if let Some((completion_tx, completion_rx)) = notify_request.into_completion()
{
let notify = crate::server::record::NotifyWaitSet::new(completion_tx);
{
let rec_arc = {
let recs = self.inner.records.read().await;
recs.get(record_name).cloned()
};
if let Some(rec_arc) = rec_arc {
let mut guard = rec_arc.write().await;
if guard.notify.is_some() {
return Err(CaError::PutCallbackInProgress(record_name.to_string()));
}
guard.notify = Some(notify.clone());
}
}
Some((notify, completion_rx))
} else {
None
};
if field == "VAL" {
let rec_arc = {
let recs = self.inner.records.read().await;
recs.get(record_name).cloned()
};
if let Some(rec_arc) = rec_arc {
let mut guard = rec_arc.write().await;
if guard.record.soft_channel_skips_convert() {
guard.record.set_device_did_compute(true);
}
}
}
let _ = self.put_driven_process_already_locked(record_name).await;
let originating_pending = want_notify && {
let rec = self.inner.records.read().await;
if let Some(rec_arc) = rec.get(record_name) {
rec_arc.read().await.notify.is_some()
} else {
false
}
};
if !originating_pending {
let rec_arc = {
let recs = self.inner.records.read().await;
recs.get(record_name).cloned()
};
if let Some(rec_arc) = rec_arc {
let mut guard = rec_arc.write().await;
if !guard.is_processing() {
guard.common.putf = false;
}
}
}
match parked {
Some((notify, completion_rx)) => {
if notify.completed() {
Ok(None)
} else {
Ok(completion_rx)
}
}
None => Ok(None),
}
}
pub async fn put_pv_no_process(&self, name: &str, value: EpicsValue) -> CaResult<()> {
let (base, field) = super::parse_pv_name(name);
let field = field.to_ascii_uppercase();
if let Some(pv) = self.inner.simple_pvs.read().await.get(name) {
pv.set(value).await;
return Ok(());
}
if let Some(rec) = self.get_record(base).await {
let canonical_base: String = self
.resolve_alias(base)
.await
.unwrap_or_else(|| base.to_string());
let _record_gate = self.lock_record(&canonical_base).await;
let mut instance = rec.write().await;
let prev_value = instance.record.get_field(&field);
match instance.record.put_field(&field, value.clone()) {
Ok(()) => {}
Err(CaError::FieldNotFound(_)) => {
instance.put_common_field(&field, value)?;
}
Err(e) => return Err(e),
}
instance.notify_field_written_if_changed(&field, prev_value.as_ref());
if field == "ASG" {
crate::server::access_security::notify_asg_field_changed();
}
return Ok(());
}
Err(CaError::ChannelNotFound(name.to_string()))
}
}
fn metadata_from_snapshot(snapshot: &Snapshot) -> crate::server::pv::PvMetadata {
crate::server::pv::PvMetadata {
display: snapshot.display.clone(),
control: snapshot.control.clone(),
enums: snapshot.enums.clone(),
}
}
#[cfg(test)]
mod tests {
use super::super::PvDatabase;
use crate::types::EpicsValue;
#[tokio::test]
async fn put_pv_and_post_handles_simple_pv() {
let db = PvDatabase::new();
db.add_pv("gw:test", EpicsValue::Double(0.0)).await.unwrap();
db.put_pv_and_post("gw:test", EpicsValue::Double(42.0))
.await
.expect("simple PV put_pv_and_post must succeed");
let pv = db.find_pv("gw:test").await.expect("PV exists");
assert!(matches!(pv.get().await, EpicsValue::Double(v) if v == 42.0));
}
#[tokio::test]
async fn field_io_entry_points_accept_aliases() {
use crate::server::records::ai::AiRecord;
let db = PvDatabase::new();
db.add_record("CANON", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
db.add_alias("ALT", "CANON").await.unwrap();
db.put_pv("CANON.VAL", EpicsValue::Double(1.5))
.await
.unwrap();
let v = db.get_pv("ALT.VAL").await.unwrap();
assert!(matches!(v, EpicsValue::Double(x) if x == 1.5));
db.put_pv("ALT.VAL", EpicsValue::Double(7.0)).await.unwrap();
let v = db.get_pv("CANON.VAL").await.unwrap();
assert!(matches!(v, EpicsValue::Double(x) if x == 7.0));
db.put_pv_and_post("ALT.VAL", EpicsValue::Double(11.0))
.await
.unwrap();
let v = db.get_pv("CANON.VAL").await.unwrap();
assert!(matches!(v, EpicsValue::Double(x) if x == 11.0));
db.put_pv_no_process("ALT.VAL", EpicsValue::Double(13.0))
.await
.unwrap();
let v = db.get_pv("ALT.VAL").await.unwrap();
assert!(matches!(v, EpicsValue::Double(x) if x == 13.0));
}
#[tokio::test]
async fn write_path_menu_label_resolves_against_field_menu() {
use crate::server::records::sel::SelRecord;
let db = PvDatabase::new();
db.add_record("SEL", Box::new(SelRecord::default()))
.await
.unwrap();
db.put_pv("SEL.SELM", EpicsValue::String("Specified".into()))
.await
.unwrap();
assert_eq!(db.get_pv("SEL.SELM").await.unwrap(), EpicsValue::Enum(0));
db.put_pv_and_post("SEL.SELM", EpicsValue::String("High Signal".into()))
.await
.unwrap();
assert_eq!(db.get_pv("SEL.SELM").await.unwrap(), EpicsValue::Enum(1));
db.put_pv("SEL.SELM", EpicsValue::String("2".into()))
.await
.unwrap();
assert_eq!(db.get_pv("SEL.SELM").await.unwrap(), EpicsValue::Enum(2));
}
#[tokio::test]
async fn set_pv_metadata_installs_without_posting() {
use crate::error::CaError;
use crate::server::snapshot::{DisplayInfo, Snapshot};
use crate::types::DbFieldType;
use std::time::SystemTime;
let db = PvDatabase::new();
db.add_pv("gw:meta", EpicsValue::Double(0.0)).await.unwrap();
const DBE_PROPERTY: u16 = 8;
let pv = db.find_pv("gw:meta").await.expect("PV exists");
let mut prop_rx = pv
.add_subscriber(1, DbFieldType::Double, DBE_PROPERTY)
.await
.expect("subscriber added");
let mut ctrl = Snapshot::new(EpicsValue::Double(0.0), 0, 0, SystemTime::UNIX_EPOCH);
ctrl.display = Some(DisplayInfo {
units: "mm".into(),
precision: 3,
upper_disp_limit: 10.0,
lower_disp_limit: -10.0,
..Default::default()
});
db.set_pv_metadata("gw:meta", &ctrl)
.await
.expect("simple PV set_pv_metadata must succeed");
let installed = pv.metadata();
assert_eq!(
installed.display.expect("display metadata installed").units,
"mm"
);
assert!(
prop_rx.try_recv().is_err(),
"set_pv_metadata must not post a DBE_PROPERTY event"
);
assert!(matches!(
db.set_pv_metadata("no:such:pv", &ctrl).await,
Err(CaError::ChannelNotFound(_))
));
}
#[tokio::test]
async fn post_pv_property_refreshes_and_posts_property_event() {
use crate::error::CaError;
use crate::server::snapshot::{DisplayInfo, Snapshot};
use crate::types::{DbFieldType, WallTime};
const DBE_PROPERTY: u16 = 8;
const DBE_VALUE: u16 = 1;
const MAJOR: u16 = 2;
const HIGH: u16 = 3;
let db = PvDatabase::new();
db.add_pv("gw:prop", EpicsValue::Double(0.0)).await.unwrap();
let pv = db.find_pv("gw:prop").await.expect("PV exists");
let mut prop_rx = pv
.add_subscriber(1, DbFieldType::Double, DBE_PROPERTY)
.await
.expect("property subscriber added");
let mut val_rx = pv
.add_subscriber(2, DbFieldType::Double, DBE_VALUE)
.await
.expect("value subscriber added");
let upstream_ts = WallTime::from_unix(2_000_000, 0);
let mut ctrl = Snapshot::new(EpicsValue::Double(5.0), HIGH, MAJOR, upstream_ts);
ctrl.display = Some(DisplayInfo {
units: "V".into(),
precision: 1,
..Default::default()
});
db.post_pv_property("gw:prop", ctrl)
.await
.expect("simple PV post_pv_property must succeed");
assert_eq!(
pv.metadata().display.expect("metadata refreshed").units,
"V"
);
let ev = prop_rx
.try_recv()
.expect("DBE_PROPERTY subscriber receives the property event");
assert_eq!(
ev.snapshot.display.expect("event carries metadata").units,
"V"
);
assert_eq!(
ev.snapshot.alarm.severity, MAJOR,
"upstream severity preserved"
);
assert_eq!(ev.snapshot.alarm.status, HIGH, "upstream status preserved");
assert_eq!(
ev.snapshot.timestamp, upstream_ts,
"control-DBR timestamp preserved, not a fresh wall clock"
);
assert!(
val_rx.try_recv().is_err(),
"DBE_VALUE-only subscriber must not receive a property post"
);
let again = Snapshot::new(EpicsValue::Double(0.0), 0, 0, WallTime::UNIX_EPOCH);
assert!(matches!(
db.post_pv_property("no:such:pv", again).await,
Err(CaError::ChannelNotFound(_))
));
}
#[tokio::test]
async fn put_record_field_from_ca_accepts_alias() {
use crate::server::records::ai::AiRecord;
let db = PvDatabase::new();
db.add_record("CANON", Box::new(AiRecord::new(0.0)))
.await
.unwrap();
db.add_alias("ALT", "CANON").await.unwrap();
let _ = db
.put_record_field_from_ca("ALT", "VAL", EpicsValue::Double(2.5))
.await
.expect("put via alias must succeed");
let v = db.get_pv("CANON.VAL").await.unwrap();
assert!(matches!(v, EpicsValue::Double(x) if x == 2.5));
}
#[tokio::test]
async fn alarm_ack_put_posts_record_wide_dbe_alarm() {
use crate::server::recgbl::EventMask;
use crate::server::records::ai::AiRecord;
use crate::types::DbFieldType;
let db = PvDatabase::new();
db.add_record("A:REC", Box::new(AiRecord::new(1.0)))
.await
.unwrap();
let rec = db.get_record("A:REC").await.expect("record exists");
let (mut alarm_rx, mut value_rx) = {
let mut inst = rec.write().await;
let a = inst
.add_subscriber("VAL", 1, DbFieldType::Double, EventMask::ALARM.bits())
.expect("alarm subscriber");
let v = inst
.add_subscriber("VAL", 2, DbFieldType::Double, EventMask::VALUE.bits())
.expect("value subscriber");
(a, v)
};
db.put_alarm_ack_from_ca(
"A:REC",
"VAL",
crate::server::record::AlarmAck::Transient,
0,
)
.await
.expect("ackt put");
assert!(
alarm_rx.try_recv().is_ok(),
"DBE_ALARM subscriber must receive the record-wide alarm post"
);
assert!(
value_rx.try_recv().is_err(),
"DBE_VALUE-only subscriber must not receive the alarm post"
);
db.put_alarm_ack_from_ca(
"A:REC",
"VAL",
crate::server::record::AlarmAck::Transient,
0,
)
.await
.expect("ackt re-put");
assert!(
alarm_rx.try_recv().is_err(),
"unchanged ACKT must post nothing"
);
}
#[tokio::test]
async fn post_property_fields_writes_and_posts_dbe_property_only() {
use crate::server::recgbl::EventMask;
use crate::server::records::mbbi::MbbiRecord;
use crate::types::DbFieldType;
let db = PvDatabase::new();
db.add_record("M:ENUM", Box::new(MbbiRecord::new(0)))
.await
.unwrap();
let rec = db.get_record("M:ENUM").await.expect("record exists");
let (mut prop_rx, mut val_rx) = {
let mut inst = rec.write().await;
let p = inst
.add_subscriber("ZRST", 1, DbFieldType::String, EventMask::PROPERTY.bits())
.expect("property subscriber");
let v = inst
.add_subscriber("ZRST", 2, DbFieldType::String, EventMask::VALUE.bits())
.expect("value subscriber");
(p, v)
};
let posted = db
.post_property_fields(
"M:ENUM",
vec![("ZRST".to_string(), EpicsValue::String("LABEL".into()))],
)
.await
.expect("post_property_fields succeeds");
assert_eq!(posted, vec!["ZRST".to_string()]);
assert_eq!(
db.get_pv("M:ENUM.ZRST").await.unwrap(),
EpicsValue::String("LABEL".into())
);
assert!(
prop_rx.try_recv().is_ok(),
"DBE_PROPERTY subscriber must receive the property post"
);
assert!(
val_rx.try_recv().is_err(),
"DBE_VALUE-only subscriber must not receive a property post"
);
}
#[tokio::test]
async fn ca_put_to_non_pp_val_posts_monitor() {
use crate::server::database::db_access::DbSubscription;
use crate::server::records::calc::CalcRecord;
let db = PvDatabase::new();
db.add_record("CALC1", Box::new(CalcRecord::new("0")))
.await
.unwrap();
let mut sub = DbSubscription::subscribe(&db, "CALC1.VAL")
.await
.expect("subscribe to CALC1.VAL");
db.put_record_field_from_ca("CALC1", "VAL", EpicsValue::Double(5.0))
.await
.expect("CA put to CALC1.VAL must succeed");
let got = tokio::time::timeout(std::time::Duration::from_secs(1), sub.recv_f64())
.await
.expect("a DBE_VALUE monitor must fire for a direct VAL put to a non-pp record");
assert_eq!(got, Some(5.0));
}
#[tokio::test]
async fn r8_22_db_burst_keeps_earlier_distinct_updates() {
use crate::server::database::db_access::DbSubscription;
use crate::server::event_queue::{event_que_size, events_per_que};
use crate::server::records::calc::CalcRecord;
let db = PvDatabase::new();
db.add_record("CALC1", Box::new(CalcRecord::new("0")))
.await
.unwrap();
let mut sub = DbSubscription::subscribe(&db, "CALC1.VAL")
.await
.expect("subscribe to CALC1.VAL");
let appended = event_que_size() - events_per_que();
let burst = appended + 40;
for i in 1..=burst {
db.put_record_field_from_ca("CALC1", "VAL", EpicsValue::Double(i as f64))
.await
.expect("CA put to CALC1.VAL must succeed");
}
let mut seq = Vec::new();
while let Ok(Some(v)) =
tokio::time::timeout(std::time::Duration::from_millis(200), sub.recv_f64()).await
{
seq.push(v);
}
let want: Vec<f64> = (1..appended)
.map(|i| i as f64)
.chain(std::iter::once(burst as f64))
.collect();
assert_eq!(
seq, want,
"record burst delivery must be {{earlier distinct backlog…, coalesced tail}}"
);
}
}