use std::collections::HashMap;
use std::collections::VecDeque;
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, Ordering};
use crate::runtime::sync::{Arc, Notify};
use crate::server::pv::MonitorEvent;
pub const EVENT_ENTRIES: usize = 4;
pub fn events_per_que() -> usize {
crate::runtime::env::get("EPICS_CAS_MAX_EVENTS_PER_CHAN")
.and_then(|v| v.parse::<usize>().ok())
.filter(|n| *n >= EVENT_ENTRIES)
.unwrap_or(36)
}
pub fn event_que_size() -> usize {
EVENT_ENTRIES * events_per_que()
}
#[derive(Default, Debug)]
struct FlowCtrl {
on: AtomicBool,
}
impl FlowCtrl {
fn is_on(&self) -> bool {
self.on.load(Ordering::Acquire)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PostOutcome {
Appended { first_event: bool },
Replaced,
Closed,
}
struct SubQ {
events: VecDeque<MonitorEvent>,
nreplace: u64,
producer_gone: bool,
reader_gone: bool,
}
impl SubQ {
fn new() -> Self {
Self {
events: VecDeque::new(),
nreplace: 0,
producer_gone: false,
reader_gone: false,
}
}
}
struct QueInner {
total_pending: usize,
n_duplicates: usize,
draining: bool,
poll_wakers: Vec<std::task::Waker>,
subs: HashMap<u32, SubQ>,
quota: usize,
size: usize,
replace_threshold: usize,
}
impl QueInner {
fn ring_space(&self) -> usize {
self.size - self.total_pending
}
fn remove_front(&mut self, sid: u32) -> Option<MonitorEvent> {
let sub = self.subs.get_mut(&sid)?;
let event = sub.events.pop_front()?;
if !sub.events.is_empty() {
debug_assert!(self.n_duplicates >= 1, "nDuplicates underflow");
self.n_duplicates -= 1;
}
self.total_pending -= 1;
self.draining = self.total_pending > 0;
Some(event)
}
fn detach(&mut self, sid: u32) {
while self.remove_front(sid).is_some() {}
if self.subs.remove(&sid).is_some() {
self.quota -= EVENT_ENTRIES;
}
}
fn try_attach(&mut self, sid: u32) -> bool {
if self.quota >= self.size - EVENT_ENTRIES {
return false;
}
self.quota += EVENT_ENTRIES;
self.subs.insert(sid, SubQ::new());
true
}
fn may_drain(&self, flow_ctrl_on: bool) -> bool {
self.draining || !flow_ctrl_on || self.n_duplicates > 0
}
}
pub struct EvQue {
flow: Arc<FlowCtrl>,
inner: Mutex<QueInner>,
wake: Notify,
}
impl EvQue {
fn new(flow: Arc<FlowCtrl>) -> Self {
Self {
flow,
inner: Mutex::new(QueInner {
total_pending: 0,
n_duplicates: 0,
draining: false,
poll_wakers: Vec::new(),
subs: HashMap::new(),
quota: 0,
size: event_que_size(),
replace_threshold: events_per_que(),
}),
wake: Notify::new(),
}
}
fn lock(&self) -> std::sync::MutexGuard<'_, QueInner> {
self.inner
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
fn wake_readers(&self) {
let wakers = std::mem::take(&mut self.lock().poll_wakers);
self.wake.notify_waiters();
for waker in wakers {
waker.wake();
}
}
fn post(&self, sid: u32, event: MonitorEvent) -> PostOutcome {
let outcome = {
let mut q = self.lock();
let flow_on = self.flow.is_on();
let ring_space = q.ring_space();
let size = q.size;
let threshold = q.replace_threshold;
let Some(sub) = q.subs.get_mut(&sid) else {
return PostOutcome::Closed;
};
if sub.reader_gone {
return PostOutcome::Closed;
}
let npend = sub.events.len();
if npend > 0 && (flow_on || ring_space <= threshold) {
let last = sub
.events
.back_mut()
.expect("npend > 0 ⇒ this monitor has a last log");
let displaced = last.mask;
*last = event;
last.mask |= displaced;
sub.nreplace += 1;
PostOutcome::Replaced
} else {
sub.events.push_back(event);
if npend > 0 {
q.n_duplicates += 1;
}
q.total_pending += 1;
debug_assert!(
q.total_pending < size,
"the quota reservation must keep the ring from filling"
);
PostOutcome::Appended {
first_event: ring_space == size,
}
}
};
if outcome != PostOutcome::Closed {
self.wake_readers();
}
outcome
}
async fn next(&self, sid: u32) -> Option<MonitorEvent> {
loop {
let wake = self.wake.notified();
tokio::pin!(wake);
wake.as_mut().enable();
{
let mut q = self.lock();
let flow_on = self.flow.is_on();
let sub = q.subs.get(&sid)?;
let has_entry = !sub.events.is_empty();
let producer_gone = sub.producer_gone;
if has_entry {
if q.may_drain(flow_on) {
q.draining = true;
return q.remove_front(sid);
}
} else if producer_gone {
return None;
}
}
wake.await;
}
}
fn poll_next(
&self,
sid: u32,
cx: &mut std::task::Context<'_>,
) -> std::task::Poll<Option<MonitorEvent>> {
use std::task::Poll;
let mut q = self.lock();
let flow_on = self.flow.is_on();
let Some(sub) = q.subs.get(&sid) else {
return Poll::Ready(None);
};
let has_entry = !sub.events.is_empty();
let producer_gone = sub.producer_gone;
if has_entry {
if q.may_drain(flow_on) {
q.draining = true;
return Poll::Ready(q.remove_front(sid));
}
} else if producer_gone {
return Poll::Ready(None);
}
if !q.poll_wakers.iter().any(|w| w.will_wake(cx.waker())) {
q.poll_wakers.push(cx.waker().clone());
}
Poll::Pending
}
fn try_next(&self, sid: u32) -> Result<MonitorEvent, TryRecvError> {
let mut q = self.lock();
let flow_on = self.flow.is_on();
let Some(sub) = q.subs.get(&sid) else {
return Err(TryRecvError::Disconnected);
};
if sub.events.is_empty() {
return Err(if sub.producer_gone {
TryRecvError::Disconnected
} else {
TryRecvError::Empty
});
}
if !q.may_drain(flow_on) {
return Err(TryRecvError::Empty);
}
q.draining = true;
q.remove_front(sid).ok_or(TryRecvError::Empty)
}
fn reader_gone(&self, sid: u32) -> bool {
self.lock().subs.get(&sid).is_none_or(|s| s.reader_gone)
}
fn close_producer(&self, sid: u32) {
{
let mut q = self.lock();
let Some(sub) = q.subs.get_mut(&sid) else {
return;
};
sub.producer_gone = true;
if sub.reader_gone {
q.detach(sid);
}
}
self.wake_readers();
}
fn close_reader(&self, sid: u32) {
{
let mut q = self.lock();
let Some(sub) = q.subs.get_mut(&sid) else {
return;
};
sub.reader_gone = true;
q.detach(sid);
}
self.wake_readers();
}
pub fn n_duplicates(&self) -> usize {
self.lock().n_duplicates
}
pub fn nreplace(&self, sid: u32) -> u64 {
self.lock().subs.get(&sid).map_or(0, |s| s.nreplace)
}
pub fn npend(&self, sid: u32) -> usize {
self.lock().subs.get(&sid).map_or(0, |s| s.events.len())
}
pub fn quota(&self) -> usize {
self.lock().quota
}
}
pub struct EventUser {
flow: Arc<FlowCtrl>,
ques: Mutex<Vec<Arc<EvQue>>>,
}
impl Default for EventUser {
fn default() -> Self {
Self::new()
}
}
impl EventUser {
pub fn new() -> Self {
let flow = Arc::new(FlowCtrl::default());
Self {
ques: Mutex::new(vec![Arc::new(EvQue::new(flow.clone()))]),
flow,
}
}
fn attach_que(&self, sid: u32) -> Arc<EvQue> {
let mut ques = self
.ques
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
for que in ques.iter() {
if que.lock().try_attach(sid) {
return que.clone();
}
}
let que = Arc::new(EvQue::new(self.flow.clone()));
let attached = que.lock().try_attach(sid);
debug_assert!(attached, "a fresh queue must admit its first monitor");
ques.push(que.clone());
que
}
pub fn flow_ctrl_on(&self) {
self.flow.on.store(true, Ordering::Release);
}
pub fn flow_ctrl_off(&self) {
self.flow.on.store(false, Ordering::Release);
let ques = self
.ques
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
for que in ques.iter() {
que.wake_readers();
}
}
pub fn is_flow_ctrl_on(&self) -> bool {
self.flow.is_on()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TryRecvError {
Empty,
Disconnected,
}
impl std::fmt::Display for TryRecvError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Empty => write!(f, "no monitor event available"),
Self::Disconnected => write!(f, "monitor producer gone"),
}
}
}
impl std::error::Error for TryRecvError {}
pub struct EventReader {
que: Arc<EvQue>,
sid: u32,
}
impl EventReader {
pub async fn recv(&mut self) -> Option<MonitorEvent> {
self.que.next(self.sid).await
}
pub fn try_recv(&mut self) -> Result<MonitorEvent, TryRecvError> {
self.que.try_next(self.sid)
}
pub fn poll_recv(
&mut self,
cx: &mut std::task::Context<'_>,
) -> std::task::Poll<Option<MonitorEvent>> {
self.que.poll_next(self.sid, cx)
}
pub fn queue(&self) -> Arc<EvQue> {
self.que.clone()
}
pub fn npend(&self) -> usize {
self.que.npend(self.sid)
}
}
impl Drop for EventReader {
fn drop(&mut self) {
self.que.close_reader(self.sid);
}
}
pub struct EventSink {
que: Arc<EvQue>,
sid: u32,
}
impl EventSink {
pub fn post(&self, event: MonitorEvent) -> PostOutcome {
self.que.post(self.sid, event)
}
pub fn is_closed(&self) -> bool {
self.que.reader_gone(self.sid)
}
}
impl Drop for EventSink {
fn drop(&mut self) {
self.que.close_producer(self.sid);
}
}
pub fn attach(user: &EventUser, sid: u32) -> (EventSink, EventReader) {
let que = user.attach_que(sid);
(
EventSink {
que: que.clone(),
sid,
},
EventReader { que, sid },
)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::server::recgbl::EventMask;
use crate::server::snapshot::Snapshot;
use crate::types::EpicsValue;
fn ev(v: i32) -> MonitorEvent {
MonitorEvent {
snapshot: Snapshot::new(EpicsValue::Long(v), 0, 0, std::time::SystemTime::UNIX_EPOCH),
origin: 0,
mask: EventMask::VALUE,
}
}
fn val(e: &MonitorEvent) -> i32 {
match e.snapshot.value {
EpicsValue::Long(v) => v,
ref other => panic!("expected Long, got {other:?}"),
}
}
#[epics_macros_rs::epics_test]
async fn npend_zero_appends_even_under_flow_control() {
let user = EventUser::new();
user.flow_ctrl_on();
let (sink, reader) = attach(&user, 1);
assert_eq!(
sink.post(ev(1)),
PostOutcome::Appended { first_event: true }
);
assert_eq!(reader.npend(), 1);
assert_eq!(
reader.queue().n_duplicates(),
0,
"a monitor's first entry is not a duplicate"
);
}
#[epics_macros_rs::epics_test]
async fn npend_positive_under_flow_control_replaces_in_place() {
let user = EventUser::new();
user.flow_ctrl_on();
let (sink, mut reader) = attach(&user, 1);
sink.post(ev(1));
assert_eq!(sink.post(ev(2)), PostOutcome::Replaced);
assert_eq!(sink.post(ev(3)), PostOutcome::Replaced);
assert_eq!(reader.npend(), 1, "the queue never grew past one entry");
assert_eq!(reader.queue().nreplace(1), 2);
assert_eq!(reader.queue().n_duplicates(), 0);
user.flow_ctrl_off();
assert_eq!(
val(&reader.recv().await.unwrap()),
3,
"the latest value survives in the held entry"
);
assert!(matches!(reader.try_recv(), Err(TryRecvError::Empty)));
}
#[epics_macros_rs::epics_test]
async fn ring_space_above_threshold_appends_distinct_entries() {
let user = EventUser::new();
let (sink, mut reader) = attach(&user, 1);
for v in 1..=3 {
assert!(matches!(sink.post(ev(v)), PostOutcome::Appended { .. }));
}
assert_eq!(reader.npend(), 3);
assert_eq!(reader.queue().n_duplicates(), 2, "npend 3 ⇒ 2 duplicates");
let got: Vec<i32> = (0..3).map(|_| val(&reader.try_recv().unwrap())).collect();
assert_eq!(got, vec![1, 2, 3]);
assert_eq!(reader.queue().n_duplicates(), 0, "symmetric on drain");
}
#[epics_macros_rs::epics_test]
async fn ring_space_at_threshold_replaces_only_the_last_entry() {
let user = EventUser::new();
let (sink, mut reader) = attach(&user, 1);
let appended = event_que_size() - events_per_que();
for v in 0..appended as i32 {
assert!(
matches!(sink.post(ev(v)), PostOutcome::Appended { .. }),
"post {v} must append while ring space is above the threshold"
);
}
assert_eq!(reader.npend(), appended);
for v in 100..110 {
assert_eq!(sink.post(ev(v)), PostOutcome::Replaced);
}
assert_eq!(reader.npend(), appended, "the backlog did not grow");
assert_eq!(reader.queue().nreplace(1), 10);
let got: Vec<i32> = (0..appended)
.map(|_| val(&reader.try_recv().unwrap()))
.collect();
let mut want: Vec<i32> = (0..appended as i32 - 1).collect();
want.push(109);
assert_eq!(
got, want,
"earlier distinct entries survive; only the tail coalesced"
);
assert!(matches!(reader.try_recv(), Err(TryRecvError::Empty)));
}
#[epics_macros_rs::epics_test]
async fn detach_with_queued_entries_keeps_the_duplicate_count_symmetric() {
let user = EventUser::new();
let (sink, reader) = attach(&user, 1);
let que = reader.queue();
for v in 1..=3 {
sink.post(ev(v)); }
assert_eq!(que.n_duplicates(), 2);
assert_eq!(que.npend(1), 3);
drop(reader);
assert_eq!(que.n_duplicates(), 0, "teardown removed its duplicates");
assert_eq!(que.npend(1), 0, "its entries left the ring");
assert_eq!(sink.post(ev(4)), PostOutcome::Closed);
assert!(sink.is_closed());
}
#[epics_macros_rs::epics_test]
async fn producer_drop_drains_then_ends_the_stream() {
let user = EventUser::new();
let (sink, mut reader) = attach(&user, 1);
sink.post(ev(1));
sink.post(ev(2));
drop(sink);
assert_eq!(val(&reader.recv().await.unwrap()), 1);
assert_eq!(val(&reader.recv().await.unwrap()), 2);
assert!(reader.recv().await.is_none(), "drained ⇒ end of stream");
assert!(matches!(reader.try_recv(), Err(TryRecvError::Disconnected)));
}
#[epics_macros_rs::epics_test]
async fn flow_ctrl_off_wakes_a_suspended_reader() {
let user = Arc::new(EventUser::new());
user.flow_ctrl_on();
let (sink, mut reader) = attach(&user, 1);
sink.post(ev(1));
sink.post(ev(2)); let u2 = user.clone();
let waker = crate::runtime::task::spawn(async move {
crate::runtime::task::yield_now().await;
u2.flow_ctrl_off();
});
let got = crate::runtime::task::timeout(std::time::Duration::from_secs(2), reader.recv())
.await
.expect("EVENTS_ON must wake the suspended reader")
.expect("the held entry is delivered");
waker.await.unwrap();
assert_eq!(val(&got), 2, "the held entry carries the latest value");
}
#[epics_macros_rs::epics_test]
async fn duplicate_on_a_sibling_subscription_releases_the_events_off_drain() {
let user = EventUser::new();
let (sink_a, mut reader_a) = attach(&user, 1);
let (sink_b, _reader_b) = attach(&user, 2);
assert!(
Arc::ptr_eq(&reader_a.queue(), &_reader_b.queue()),
"the quota admits both monitors to one queue — C's sharing granularity"
);
sink_a.post(ev(1)); sink_b.post(ev(10));
sink_b.post(ev(11)); assert_eq!(reader_a.queue().n_duplicates(), 1);
user.flow_ctrl_on();
let got = crate::runtime::task::timeout(std::time::Duration::from_secs(2), reader_a.recv())
.await
.expect("a duplicate anywhere on the queue must release the drain")
.expect("A's entry is delivered");
assert_eq!(val(&got), 1);
}
#[epics_macros_rs::epics_test]
async fn flow_control_without_duplicates_suspends_every_reader_on_the_queue() {
let user = EventUser::new();
user.flow_ctrl_on();
let (sink_a, mut reader_a) = attach(&user, 1);
let (sink_b, mut reader_b) = attach(&user, 2);
sink_a.post(ev(1));
sink_b.post(ev(2));
assert_eq!(reader_a.queue().n_duplicates(), 0);
assert!(matches!(reader_a.try_recv(), Err(TryRecvError::Empty)));
assert!(matches!(reader_b.try_recv(), Err(TryRecvError::Empty)));
user.flow_ctrl_off();
assert_eq!(val(&reader_a.try_recv().unwrap()), 1);
assert_eq!(val(&reader_b.try_recv().unwrap()), 2);
}
#[epics_macros_rs::epics_test]
async fn quota_exhaustion_chains_a_second_queue() {
let user = EventUser::new();
let cap = event_que_size() / EVENT_ENTRIES - 1;
let mut held = Vec::new();
for sid in 0..cap as u32 {
held.push(attach(&user, sid));
}
let first = held[0].1.queue();
for (_, reader) in &held {
assert!(Arc::ptr_eq(&reader.queue(), &first), "all within quota");
}
assert_eq!(first.quota(), event_que_size() - EVENT_ENTRIES);
let (_sink, overflow) = attach(&user, cap as u32);
assert!(
!Arc::ptr_eq(&overflow.queue(), &first),
"the {cap}th monitor exhausts the quota; the next one chains a queue"
);
assert_eq!(overflow.queue().quota(), EVENT_ENTRIES);
}
#[epics_macros_rs::epics_test]
async fn detaching_a_monitor_releases_its_quota() {
let user = EventUser::new();
let cap = event_que_size() / EVENT_ENTRIES - 1;
let mut held: Vec<_> = (0..cap as u32).map(|sid| attach(&user, sid)).collect();
let first = held[0].1.queue();
assert_eq!(first.quota(), event_que_size() - EVENT_ENTRIES);
held.pop(); assert_eq!(first.quota(), event_que_size() - 2 * EVENT_ENTRIES);
let (_sink, reader) = attach(&user, 900);
assert!(
Arc::ptr_eq(&reader.queue(), &first),
"the released quota must be reusable"
);
}
#[epics_macros_rs::epics_test]
async fn detaching_one_monitor_leaves_a_siblings_duplicates_counted() {
let user = EventUser::new();
let (sink_a, reader_a) = attach(&user, 1);
let (sink_b, reader_b) = attach(&user, 2);
let que = reader_a.queue();
for v in 1..=3 {
sink_a.post(ev(v)); }
for v in 10..=11 {
sink_b.post(ev(v)); }
assert_eq!(que.n_duplicates(), 3);
drop(reader_a);
drop(sink_a);
assert_eq!(que.n_duplicates(), 1, "only A's duplicates left the ring");
assert_eq!(que.npend(2), 2, "B's entries are untouched");
drop(reader_b);
drop(sink_b);
assert_eq!(que.n_duplicates(), 0);
assert_eq!(que.quota(), 0, "both monitors released their reservation");
}
struct CountWaker(std::sync::atomic::AtomicUsize);
impl std::task::Wake for CountWaker {
fn wake(self: Arc<Self>) {
self.0.fetch_add(1, Ordering::SeqCst);
}
}
fn count_waker() -> (Arc<CountWaker>, std::task::Waker) {
let counter = Arc::new(CountWaker(std::sync::atomic::AtomicUsize::new(0)));
let waker = std::task::Waker::from(counter.clone());
(counter, waker)
}
#[test]
fn poll_recv_parks_then_delivers_on_post() {
let user = EventUser::new();
let (sink, mut reader) = attach(&user, 7);
let (counter, waker) = count_waker();
let mut cx = std::task::Context::from_waker(&waker);
assert!(reader.poll_recv(&mut cx).is_pending(), "empty queue parks");
assert_eq!(counter.0.load(Ordering::SeqCst), 0);
sink.post(ev(41));
assert_eq!(
counter.0.load(Ordering::SeqCst),
1,
"the post must flush the registered waker"
);
match reader.poll_recv(&mut cx) {
std::task::Poll::Ready(Some(event)) => assert_eq!(val(&event), 41),
other => panic!("expected the posted event, got {other:?}"),
}
assert!(
reader.poll_recv(&mut cx).is_pending(),
"drained queue parks"
);
assert_eq!(counter.0.load(Ordering::SeqCst), 1);
}
#[test]
fn poll_recv_drains_backlog_then_reports_disconnect() {
let user = EventUser::new();
let (sink, mut reader) = attach(&user, 7);
let (counter, waker) = count_waker();
let mut cx = std::task::Context::from_waker(&waker);
sink.post(ev(1));
sink.post(ev(2));
match reader.poll_recv(&mut cx) {
std::task::Poll::Ready(Some(event)) => assert_eq!(val(&event), 1),
other => panic!("expected first entry, got {other:?}"),
}
match reader.poll_recv(&mut cx) {
std::task::Poll::Ready(Some(event)) => assert_eq!(val(&event), 2),
other => panic!("expected second entry, got {other:?}"),
}
assert!(reader.poll_recv(&mut cx).is_pending());
let woken_before = counter.0.load(Ordering::SeqCst);
drop(sink);
assert!(
counter.0.load(Ordering::SeqCst) > woken_before,
"producer teardown must wake the parked poller"
);
assert!(
matches!(reader.poll_recv(&mut cx), std::task::Poll::Ready(None)),
"producer gone + queue drained ⇒ end of stream"
);
}
#[test]
fn poll_recv_respects_events_off_and_wakes_on_events_on() {
let user = EventUser::new();
let (sink, mut reader) = attach(&user, 7);
let (counter, waker) = count_waker();
let mut cx = std::task::Context::from_waker(&waker);
user.flow_ctrl_on();
sink.post(ev(5)); assert!(
reader.poll_recv(&mut cx).is_pending(),
"EVENTS_OFF with no duplicates suspends the poll-based reader too"
);
user.flow_ctrl_off();
assert_eq!(
counter.0.load(Ordering::SeqCst),
1,
"the post preceded registration (woke nobody); EVENTS_ON must wake \
the parked poller"
);
match reader.poll_recv(&mut cx) {
std::task::Poll::Ready(Some(event)) => assert_eq!(val(&event), 5),
other => panic!("expected the withheld entry after EVENTS_ON, got {other:?}"),
}
}
}