use std::cell::RefCell;
use std::future::Future;
use std::pin::Pin;
use std::rc::Rc;
use rustdv_sim::queue::Queue;
use crate::component::{Component, ComponentNode};
use crate::port::{bind_or_panic, sink_of, GetIf, PeekIf, PortName, PortOwner, PutIf, SinkHandle};
struct FifoInner<T: 'static> {
q: Queue<T>,
size: Option<usize>,
put_taps: RefCell<Vec<Rc<dyn SinkHandle<T>>>>,
get_taps: RefCell<Vec<Rc<dyn SinkHandle<T>>>>,
}
impl<T: 'static> FifoInner<T> {
fn new(size: Option<usize>) -> FifoInner<T> {
FifoInner {
q: match size {
Some(n) => Queue::new(Some(n)),
None => Queue::unbounded(),
},
size,
put_taps: RefCell::new(Vec::new()),
get_taps: RefCell::new(Vec::new()),
}
}
fn is_full(&self) -> bool {
match self.size {
None => false,
Some(s) => self.q.len() >= s,
}
}
fn tap(taps: &RefCell<Vec<Rc<dyn SinkHandle<T>>>>, item: &T) {
let subs: Vec<Rc<dyn SinkHandle<T>>> = taps.borrow().clone();
for sub in subs {
sub.deliver(item);
}
}
async fn put_tapped(&self, item: T) {
self.q.wait_for_space().await;
Self::tap(&self.put_taps, &item);
let _ = self.q.try_put(item);
}
fn try_put_tapped(&self, item: T) -> Result<(), T> {
if !self.q.has_space() {
return Err(item);
}
Self::tap(&self.put_taps, &item);
self.q.try_put(item)
}
fn tap_get(&self, item: Option<T>) -> Option<T> {
if let Some(v) = &item {
Self::tap(&self.get_taps, v);
}
item
}
}
impl<T: 'static> PutIf<T> for FifoInner<T> {
fn put(&self, item: T) -> Pin<Box<dyn Future<Output = ()> + '_>> {
Box::pin(self.put_tapped(item))
}
fn try_put(&self, item: T) -> Result<(), T> {
self.try_put_tapped(item)
}
fn can_put(&self) -> bool {
!self.is_full()
}
}
impl<T: 'static> GetIf<T> for FifoInner<T> {
fn get(&self) -> Pin<Box<dyn Future<Output = T> + '_>> {
Box::pin(async move {
let item = self.q.get().await;
Self::tap(&self.get_taps, &item);
item
})
}
fn try_get(&self) -> Option<T> {
let item = self.q.try_get();
self.tap_get(item)
}
fn can_get(&self) -> bool {
!self.q.is_empty()
}
}
impl<T: Clone + 'static> PeekIf<T> for FifoInner<T> {
fn peek(&self) -> Pin<Box<dyn Future<Output = T> + '_>> {
Box::pin(self.q.peek())
}
fn try_peek(&self) -> Option<T> {
self.q.try_peek()
}
fn can_peek(&self) -> bool {
!self.q.is_empty()
}
}
pub struct PutExport<T: 'static> {
iface: Rc<dyn PutIf<T>>,
}
impl<T: 'static> PutExport<T> {
pub fn connect(&self, owner: &dyn PortOwner, name: PortName<dyn PutIf<T>>) {
bind_or_panic(owner, name, self.iface.clone());
}
}
pub struct GetExport<T: 'static> {
iface: Rc<dyn GetIf<T>>,
}
impl<T: 'static> GetExport<T> {
pub fn connect(&self, owner: &dyn PortOwner, name: PortName<dyn GetIf<T>>) {
bind_or_panic(owner, name, self.iface.clone());
}
}
pub struct TapExport<T: 'static> {
taps: Rc<FifoInner<T>>,
on_put: bool,
}
impl<T: 'static> TapExport<T> {
pub fn connect(&self, owner: &dyn PortOwner, name: PortName<dyn SinkHandle<T>>) {
match sink_of(owner, name) {
Ok(sink) => {
let list = if self.on_put { &self.taps.put_taps } else { &self.taps.get_taps };
list.borrow_mut().push(sink);
}
Err(e) => panic!("{e}"),
}
}
}
pub struct PeekExport<T: 'static> {
iface: Rc<dyn PeekIf<T>>,
}
impl<T: 'static> PeekExport<T> {
pub fn connect(&self, owner: &dyn PortOwner, name: PortName<dyn PeekIf<T>>) {
bind_or_panic(owner, name, self.iface.clone());
}
}
pub struct TlmFifo<T: 'static> {
inner: Rc<FifoInner<T>>,
}
impl<T: 'static> Default for TlmFifo<T> {
fn default() -> Self {
TlmFifo::new(1)
}
}
impl<T: 'static> TlmFifo<T> {
pub fn new(size: usize) -> TlmFifo<T> {
TlmFifo { inner: Rc::new(FifoInner::new(Some(size))) }
}
pub fn unbounded() -> TlmFifo<T> {
TlmFifo { inner: Rc::new(FifoInner::new(None)) }
}
pub fn size(&self) -> Option<usize> {
self.inner.size
}
pub fn used(&self) -> usize {
self.inner.q.len()
}
pub fn is_empty(&self) -> bool {
self.inner.q.is_empty()
}
pub fn is_full(&self) -> bool {
self.inner.is_full()
}
pub fn flush(&self) {
while self.inner.q.try_get().is_some() {}
}
pub fn put_export(&self) -> PutExport<T> {
PutExport { iface: self.inner.clone() }
}
pub fn get_export(&self) -> GetExport<T> {
GetExport { iface: self.inner.clone() }
}
pub fn put_ap(&self) -> TapExport<T> {
TapExport { taps: self.inner.clone(), on_put: true }
}
pub fn get_ap(&self) -> TapExport<T> {
TapExport { taps: self.inner.clone(), on_put: false }
}
pub async fn put(&self, item: T) {
self.inner.put_tapped(item).await
}
pub fn try_put(&self, item: T) -> Result<(), T> {
self.inner.try_put_tapped(item)
}
pub async fn get(&self) -> T {
let item = self.inner.q.get().await;
FifoInner::tap(&self.inner.get_taps, &item);
item
}
pub fn try_get(&self) -> Option<T> {
let item = self.inner.q.try_get();
self.inner.tap_get(item)
}
#[cfg(test)]
pub(crate) fn put_iface_for_test(&self) -> Rc<dyn PutIf<T>> {
self.inner.clone()
}
pub fn handle(&self) -> TlmFifo<T> {
TlmFifo { inner: self.inner.clone() }
}
}
impl<T: Clone + 'static> TlmFifo<T> {
pub fn peek_export(&self) -> PeekExport<T> {
PeekExport { iface: self.inner.clone() }
}
pub async fn peek(&self) -> T {
self.inner.q.peek().await
}
pub fn try_peek(&self) -> Option<T> {
self.inner.q.try_peek()
}
}
impl<T: 'static> Component for TlmFifo<T> {}
impl<T: 'static> ComponentNode for TlmFifo<T> {
fn node_name(&self) -> &'static str {
"TlmFifo"
}
fn children_mut(&mut self) -> Vec<(String, &mut (dyn ComponentNode + 'static))> {
Vec::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::port::{
GetPort, PeekPort, PortField, PortName, PortOwner, PutPort, SubscribePort, Subscriber,
};
use crate::shared::RustdvShared;
use rustdv_sim::testing::block_on;
use std::any::Any;
struct Holder {
put: PutPort<u8>,
get: GetPort<u8>,
peek: PeekPort<u8>,
sub: SubscribePort<u8>,
}
impl Holder {
fn new() -> Holder {
Holder {
put: PutPort::default(),
get: GetPort::default(),
peek: PeekPort::default(),
sub: SubscribePort::default(),
}
}
const PUT: PortName<dyn PutIf<u8>> = PortName::new("put");
const GET: PortName<dyn GetIf<u8>> = PortName::new("get");
const PEEK: PortName<dyn PeekIf<u8>> = PortName::new("peek");
const SUB: PortName<dyn SinkHandle<u8>> = PortName::new("sub");
}
impl PortOwner for Holder {
fn owner_port_slot(&self, name: &str) -> Option<Rc<dyn Any>> {
match name {
"put" => Some(self.put.slot_any()),
"get" => Some(self.get.slot_any()),
"peek" => Some(self.peek.slot_any()),
"sub" => Some(self.sub.slot_any()),
_ => None,
}
}
fn owner_label(&self) -> &'static str {
"Holder"
}
}
#[test]
fn a_port_is_unbound_until_connected() {
let h = Holder::new();
assert!(!h.put.bound());
let fifo: TlmFifo<u8> = TlmFifo::new(1);
fifo.put_export().connect(&h, Holder::PUT);
assert!(h.put.bound());
}
#[test]
fn put_and_get_through_a_fifo() {
block_on(async {
let h = Holder::new();
let fifo: TlmFifo<u8> = TlmFifo::new(2);
fifo.put_export().connect(&h, Holder::PUT);
fifo.get_export().connect(&h, Holder::GET);
h.put.put(1).await;
h.put.put(2).await;
assert_eq!(h.get.get().await, 1, "FIFO order through the ports");
assert_eq!(h.get.get().await, 2);
});
}
#[test]
fn peek_leaves_the_item_for_get() {
block_on(async {
let h = Holder::new();
let fifo: TlmFifo<u8> = TlmFifo::new(1);
fifo.put_export().connect(&h, Holder::PUT);
fifo.peek_export().connect(&h, Holder::PEEK);
fifo.get_export().connect(&h, Holder::GET);
h.put.put(9).await;
assert_eq!(h.peek.peek().await, 9);
assert_eq!(h.get.get().await, 9, "peek did not consume it");
});
}
#[test]
fn try_put_hands_a_refused_item_back() {
block_on(async {
let h = Holder::new();
let fifo: TlmFifo<u8> = TlmFifo::new(1);
fifo.put_export().connect(&h, Holder::PUT);
assert!(h.put.try_put(1).is_ok());
assert_eq!(h.put.try_put(2), Err(2), "the item comes home");
});
}
#[test]
fn can_put_and_can_get_track_the_fifo() {
block_on(async {
let h = Holder::new();
let fifo: TlmFifo<u8> = TlmFifo::new(1);
fifo.put_export().connect(&h, Holder::PUT);
fifo.get_export().connect(&h, Holder::GET);
assert!(h.put.can_put());
assert!(!h.get.can_get());
h.put.put(1).await;
assert!(!h.put.can_put());
assert!(h.get.can_get());
});
}
#[test]
fn a_bad_port_name_is_a_named_error() {
let h = Holder::new();
let nope: PortName<dyn PutIf<u8>> = PortName::new("no_such_port");
let fifo: TlmFifo<u8> = TlmFifo::new(1);
match crate::port::bind(&h, nope, fifo.put_iface_for_test()) {
Err(crate::port::ConnectError::NoSuchPort { owner, name }) => {
assert_eq!(owner, "Holder");
assert_eq!(name, "no_such_port");
}
other => panic!("expected NoSuchPort, got {other:?}"),
}
}
#[test]
fn fifo_size_used_and_flush() {
block_on(async {
let fifo: TlmFifo<u8> = TlmFifo::new(3);
assert_eq!(fifo.size(), Some(3));
assert!(fifo.is_empty());
fifo.put(1).await;
fifo.put(2).await;
assert_eq!(fifo.used(), 2);
fifo.flush();
assert!(fifo.is_empty(), "flush empties it");
});
}
#[test]
fn unbounded_is_never_full() {
block_on(async {
let fifo: TlmFifo<u8> = TlmFifo::unbounded();
assert_eq!(fifo.size(), None);
for n in 0..100 {
assert!(fifo.try_put(n).is_ok());
}
assert!(!fifo.is_full());
});
}
#[test]
fn the_put_tap_sees_every_item_and_consumes_none() {
#[derive(Default)]
struct Log {
seen: Vec<u8>,
}
impl Subscriber<u8> for Log {
fn write(&mut self, item: &u8) {
self.seen.push(*item);
}
}
block_on(async {
let h = Holder::new();
let log: RustdvShared<Log> = RustdvShared::default();
h.sub.subscribe(log.clone());
let fifo: TlmFifo<u8> = TlmFifo::unbounded();
fifo.put_ap().connect(&h, Holder::SUB);
for n in 1..=3u8 {
fifo.put(n).await;
}
assert_eq!(log.get().seen, vec![1, 2, 3], "the tap saw all three");
assert_eq!(fifo.used(), 3, "and took none of them");
});
}
#[test]
fn the_get_tap_fires_as_items_leave() {
#[derive(Default)]
struct Log {
seen: Vec<u8>,
}
impl Subscriber<u8> for Log {
fn write(&mut self, item: &u8) {
self.seen.push(*item);
}
}
block_on(async {
let h = Holder::new();
let log: RustdvShared<Log> = RustdvShared::default();
h.sub.subscribe(log.clone());
let fifo: TlmFifo<u8> = TlmFifo::unbounded();
fifo.get_ap().connect(&h, Holder::SUB);
fifo.put(7).await;
assert!(log.get().seen.is_empty(), "nothing has left yet");
let _ = fifo.get().await;
assert_eq!(log.get().seen, vec![7]);
});
}
#[test]
fn a_handle_is_the_same_fifo() {
block_on(async {
let fifo: TlmFifo<u8> = TlmFifo::unbounded();
let other = fifo.handle();
fifo.put(1).await;
assert_eq!(other.used(), 1, "two handles, one FIFO");
});
}
}