use std::collections::{HashMap, VecDeque};
use std::sync::Arc;
use std::sync::Weak;
use super::link_set::{DynLinkSet, LinkPutOp};
use crate::types::EpicsValue;
#[derive(Clone, PartialEq, Eq, Hash, Debug)]
pub(crate) enum LinkTarget {
Scheme(String),
Any,
}
#[derive(Clone, PartialEq, Eq, Hash, Debug)]
pub(crate) struct LinkKey {
pub(crate) target: LinkTarget,
pub(crate) name: String,
}
struct PutCompletion(Option<tokio::sync::oneshot::Sender<Result<(), String>>>);
impl PutCompletion {
fn resolve(mut self, result: Result<(), String>) {
if let Some(tx) = self.0.take() {
let _ = tx.send(result);
}
}
}
impl Drop for PutCompletion {
fn drop(&mut self) {
if let Some(tx) = self.0.take() {
let _ = tx.send(Err(
"external link put dropped before completion (link queue torn down)".to_string(),
));
}
}
}
struct StagedPut {
value: EpicsValue,
op: LinkPutOp,
completion: Option<PutCompletion>,
}
enum LinkState {
Queued(StagedPut),
InFlight,
InFlightRestaged(StagedPut),
}
#[derive(PartialEq, Eq, Debug)]
enum OpenState {
Queued,
InFlight,
Done,
}
#[derive(Default)]
struct QueueInner {
links: HashMap<LinkKey, LinkState>,
ready: VecDeque<LinkKey>,
opens: HashMap<LinkKey, OpenState>,
ready_opens: VecDeque<LinkKey>,
coalesced: u64,
completed: u64,
opened: u64,
}
impl QueueInner {
fn is_idle(&self) -> bool {
self.links.is_empty()
&& !self
.opens
.values()
.any(|s| matches!(s, OpenState::Queued | OpenState::InFlight))
}
}
pub(crate) struct LinkPutQueue {
inner: parking_lot::Mutex<QueueInner>,
work: Arc<tokio::sync::Notify>,
idle: tokio::sync::Notify,
owner_started: std::sync::atomic::AtomicBool,
}
impl Default for LinkPutQueue {
fn default() -> Self {
Self {
inner: parking_lot::Mutex::new(QueueInner::default()),
work: Arc::new(tokio::sync::Notify::new()),
idle: tokio::sync::Notify::new(),
owner_started: std::sync::atomic::AtomicBool::new(false),
}
}
}
impl Drop for LinkPutQueue {
fn drop(&mut self) {
self.inner.get_mut().links.clear();
self.work.notify_waiters();
}
}
impl LinkPutQueue {
pub(crate) fn stage_put(
&self,
key: LinkKey,
value: EpicsValue,
op: LinkPutOp,
) -> Option<tokio::sync::oneshot::Receiver<Result<(), String>>> {
let (completion, rx) = match op {
LinkPutOp::Plain => (None, None),
LinkPutOp::Async => {
let (tx, rx) = tokio::sync::oneshot::channel();
(Some(PutCompletion(Some(tx))), Some(rx))
}
};
let staged = StagedPut {
value,
op,
completion,
};
let wake = {
let mut inner = self.inner.lock();
match inner.links.remove(&key) {
None => {
inner.links.insert(key.clone(), LinkState::Queued(staged));
inner.ready.push_back(key);
true
}
Some(LinkState::Queued(prev)) => {
Self::supersede(&mut inner, prev);
inner.links.insert(key, LinkState::Queued(staged));
false
}
Some(LinkState::InFlight) => {
inner.links.insert(key, LinkState::InFlightRestaged(staged));
false
}
Some(LinkState::InFlightRestaged(prev)) => {
Self::supersede(&mut inner, prev);
inner.links.insert(key, LinkState::InFlightRestaged(staged));
false
}
}
};
if wake {
self.work.notify_one();
}
rx
}
fn supersede(inner: &mut QueueInner, prev: StagedPut) {
inner.coalesced += 1;
if let Some(c) = prev.completion {
c.resolve(Err(
"external link put superseded by a newer put on the same link".to_string(),
));
}
}
fn take_ready(&self) -> Option<(LinkKey, StagedPut)> {
let mut inner = self.inner.lock();
loop {
let key = inner.ready.pop_front()?;
match inner.links.remove(&key) {
Some(LinkState::Queued(staged)) => {
inner.links.insert(key.clone(), LinkState::InFlight);
return Some((key, staged));
}
Some(other) => {
inner.links.insert(key, other);
}
None => {}
}
}
}
fn finish(&self, key: LinkKey) {
let (wake_work, wake_idle) = {
let mut inner = self.inner.lock();
inner.completed += 1;
let wake_work = match inner.links.remove(&key) {
Some(LinkState::InFlight) | None => false,
Some(LinkState::InFlightRestaged(staged)) => {
inner.links.insert(key.clone(), LinkState::Queued(staged));
inner.ready.push_back(key);
true
}
Some(other) => {
inner.links.insert(key, other);
false
}
};
(wake_work, inner.is_idle())
};
if wake_work {
self.work.notify_one();
}
if wake_idle {
self.idle.notify_waiters();
}
}
pub(crate) fn stage_open(&self, key: LinkKey) -> bool {
let staged = {
let mut inner = self.inner.lock();
if inner.opens.contains_key(&key) {
false
} else {
inner.opens.insert(key.clone(), OpenState::Queued);
inner.ready_opens.push_back(key);
true
}
};
if staged {
self.work.notify_one();
}
staged
}
fn take_ready_open(&self) -> Option<LinkKey> {
let mut inner = self.inner.lock();
loop {
let key = inner.ready_opens.pop_front()?;
match inner.opens.get_mut(&key) {
Some(state @ OpenState::Queued) => {
*state = OpenState::InFlight;
return Some(key);
}
_ => continue,
}
}
}
fn finish_open(&self, key: LinkKey) {
let wake_idle = {
let mut inner = self.inner.lock();
inner.opened += 1;
inner.opens.insert(key, OpenState::Done);
inner.is_idle()
};
if wake_idle {
self.idle.notify_waiters();
}
}
pub(crate) fn is_idle(&self) -> bool {
self.inner.lock().is_idle()
}
pub(crate) fn opened_count(&self) -> u64 {
self.inner.lock().opened
}
pub(crate) fn coalesced_count(&self) -> u64 {
self.inner.lock().coalesced
}
pub(crate) fn completed_count(&self) -> u64 {
self.inner.lock().completed
}
pub(crate) async fn sync(&self) {
loop {
let notified = self.idle.notified();
if self.is_idle() {
return;
}
notified.await;
}
}
pub(crate) fn ensure_owner(self: &Arc<Self>, db: Weak<super::PvDatabaseInner>) {
if self
.owner_started
.swap(true, std::sync::atomic::Ordering::AcqRel)
{
return;
}
let queue = Arc::downgrade(self);
let work = self.work.clone();
crate::runtime::task::spawn(owner_loop(queue, work, db));
}
}
async fn owner_loop(
queue: Weak<LinkPutQueue>,
work: Arc<tokio::sync::Notify>,
db: Weak<super::PvDatabaseInner>,
) {
loop {
let notified = work.notified();
tokio::pin!(notified);
notified.as_mut().enable();
{
let Some(q) = queue.upgrade() else { return };
while let Some(key) = q.take_ready_open() {
let Some(inner) = db.upgrade() else { return };
let lsets = resolve_lsets(&inner, &key.target);
drop(inner);
let q2 = q.clone();
crate::runtime::task::spawn(async move {
for lset in &lsets {
lset.connect_link(&key.name).await;
}
q2.finish_open(key);
});
}
while let Some((key, staged)) = q.take_ready() {
let Some(inner) = db.upgrade() else { return };
let lsets = resolve_lsets(&inner, &key.target);
drop(inner);
let q2 = q.clone();
crate::runtime::task::spawn(async move {
run_put(&lsets, &key, staged).await;
q2.finish(key);
});
}
}
notified.await;
}
}
pub(super) fn resolve_lsets(
inner: &super::PvDatabaseInner,
target: &LinkTarget,
) -> Vec<DynLinkSet> {
let registry = inner.link_sets.load();
match target {
LinkTarget::Scheme(s) => registry.get(s).into_iter().collect(),
LinkTarget::Any => registry
.schemes()
.iter()
.filter_map(|s| registry.get(s))
.collect(),
}
}
async fn run_put(lsets: &[DynLinkSet], key: &LinkKey, staged: StagedPut) {
let StagedPut {
value,
op,
completion,
} = staged;
let mut result = Err(format!(
"no link set accepted the write to external link '{}'",
key.name
));
for lset in lsets {
match lset.put_value(&key.name, value.clone(), op).await {
Ok(()) => {
lset.flush_puts().await;
result = Ok(());
break;
}
Err(e) => result = Err(e),
}
}
if let Err(e) = &result {
eprintln!("dbCa: external link put to '{}' failed: {e}", key.name);
}
if let Some(c) = completion {
c.resolve(result);
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::future::Future;
use std::task::{Context, Poll, Waker};
fn key(name: &str) -> LinkKey {
LinkKey {
target: LinkTarget::Scheme("ca".to_string()),
name: name.to_string(),
}
}
fn resolved(rx: tokio::sync::oneshot::Receiver<Result<(), String>>) -> Result<(), String> {
let mut rx = std::pin::pin!(rx);
match rx.as_mut().poll(&mut Context::from_waker(Waker::noop())) {
Poll::Ready(out) => out.expect("completion resolved, not dropped"),
Poll::Pending => panic!("completion was left unresolved"),
}
}
#[test]
fn second_put_on_same_link_coalesces_latest_wins() {
let q = LinkPutQueue::default();
assert!(
q.stage_put(key("A"), EpicsValue::Long(1), LinkPutOp::Plain)
.is_none()
);
assert!(
q.stage_put(key("A"), EpicsValue::Long(2), LinkPutOp::Plain)
.is_none()
);
assert_eq!(q.coalesced_count(), 1, "one put must be coalesced away");
let (k, staged) = q.take_ready().expect("one ready link");
assert_eq!(k, key("A"));
assert_eq!(staged.value, EpicsValue::Long(2), "latest value wins");
assert!(
q.take_ready().is_none(),
"the link is enqueued exactly once"
);
}
#[test]
fn distinct_links_keep_their_own_pending_value() {
let q = LinkPutQueue::default();
q.stage_put(key("A"), EpicsValue::Long(1), LinkPutOp::Plain);
q.stage_put(key("B"), EpicsValue::Long(2), LinkPutOp::Plain);
assert_eq!(q.coalesced_count(), 0);
assert_eq!(q.take_ready().expect("A first").0, key("A"));
assert_eq!(q.take_ready().expect("B second").0, key("B"));
assert!(q.take_ready().is_none());
}
#[test]
fn put_staged_during_flight_is_requeued_on_finish() {
let q = LinkPutQueue::default();
q.stage_put(key("A"), EpicsValue::Long(1), LinkPutOp::Plain);
let (k, first) = q.take_ready().expect("ready");
assert_eq!(first.value, EpicsValue::Long(1));
assert!(q.take_ready().is_none(), "in flight, nothing ready");
q.stage_put(key("A"), EpicsValue::Long(9), LinkPutOp::Plain);
assert!(
q.take_ready().is_none(),
"still in flight — at most one write per link keeps write order"
);
q.finish(k);
let (_, second) = q.take_ready().expect("re-queued after finish");
assert_eq!(second.value, EpicsValue::Long(9));
}
#[test]
fn superseded_completion_is_resolved_not_dropped() {
let q = LinkPutQueue::default();
let rx = q
.stage_put(key("A"), EpicsValue::Long(1), LinkPutOp::Async)
.expect("Async put yields a completion");
q.stage_put(key("A"), EpicsValue::Long(2), LinkPutOp::Async);
assert!(
resolved(rx).unwrap_err().contains("superseded"),
"a superseded completion must say so"
);
}
#[test]
fn teardown_resolves_pending_completions() {
let q = LinkPutQueue::default();
let rx = q
.stage_put(key("A"), EpicsValue::Long(1), LinkPutOp::Async)
.expect("Async put yields a completion");
drop(q);
assert!(
resolved(rx).is_err(),
"a torn-down pending put must not report success"
);
}
#[test]
fn plain_put_has_no_completion() {
let q = LinkPutQueue::default();
assert!(
q.stage_put(key("A"), EpicsValue::Long(1), LinkPutOp::Plain)
.is_none()
);
}
#[test]
fn idle_tracks_staged_and_in_flight() {
let q = LinkPutQueue::default();
assert!(q.is_idle());
q.stage_put(key("A"), EpicsValue::Long(1), LinkPutOp::Plain);
assert!(!q.is_idle(), "staged but not run");
let (k, _) = q.take_ready().expect("ready");
assert!(!q.is_idle(), "in flight");
q.finish(k);
assert!(q.is_idle());
assert_eq!(q.completed_count(), 1);
}
}