Skip to main content

urushi_tui_app/delivery/
admission.rs

1//! Public admission policy and sender values.
2
3use std::error::Error;
4use std::fmt;
5use std::future::Future;
6use std::pin::Pin;
7use std::sync::Arc;
8
9use super::super::effect::Mapper;
10
11/// The default number of unaccepted items a bounded source may hold.
12const DEFAULT_CAPACITY: usize = 64;
13
14/// What a source does with an item the runtime has not accepted yet.
15///
16/// A source declares its policy through the `_with` form of its subscription
17/// constructor and otherwise gets [`Admission::bounded`] at the default
18/// capacity. The runtime's own sources carry the policies the runtime gives
19/// them and take no admission from the application.
20///
21/// ```
22/// use urushi_tui_app::Admission;
23///
24/// assert_eq!(Admission::default(), Admission::bounded(64));
25/// assert_ne!(Admission::latest(), Admission::bounded(1));
26/// ```
27#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
28pub struct Admission {
29    policy: Policy,
30}
31
32#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
33pub(crate) enum Policy {
34    /// At most `capacity` unaccepted items, oldest accepted first; the source
35    /// waits while it is full.
36    Bounded { capacity: usize },
37    /// One unaccepted slot, which a newer item replaces; the source never
38    /// waits.
39    Latest,
40}
41
42impl Admission {
43    /// A FIFO of at most `capacity` unaccepted items.
44    ///
45    /// The source waits while the queue is full, which is what makes a fast
46    /// source slow down rather than the runtime grow without bound. A capacity
47    /// of zero would admit nothing, so it is read as one.
48    pub const fn bounded(capacity: usize) -> Self {
49        let capacity = if capacity == 0 { 1 } else { capacity };
50        Self {
51            policy: Policy::Bounded { capacity },
52        }
53    }
54
55    /// One unaccepted slot, which a newer item replaces.
56    ///
57    /// The source never waits, and an item the runtime did not accept before
58    /// the next arrived is dropped. This suits an observation whose older value
59    /// carries nothing once a newer one exists.
60    pub const fn latest() -> Self {
61        Self {
62            policy: Policy::Latest,
63        }
64    }
65
66    pub(crate) const fn policy(self) -> Policy {
67        self.policy
68    }
69}
70
71impl Default for Admission {
72    fn default() -> Self {
73        Self::bounded(DEFAULT_CAPACITY)
74    }
75}
76
77/// Where an application-defined source puts its messages.
78///
79/// The source is handed one by [`Subscription::run`](crate::Subscription::run)
80/// or [`run_blocking`](crate::Subscription::run_blocking). Whether a send waits
81/// is the source's [`Admission`]: under [`Admission::bounded`] a send waits
82/// while the queue is full, and under [`Admission::latest`] it never waits.
83pub struct Sender<Message> {
84    sink: Arc<dyn Sink<Message>>,
85}
86
87/// The capability hidden behind a public [`Sender`].
88///
89/// The source-local inbox implements it; tests for mapping may substitute a
90/// collector without depending on admission internals.
91pub(crate) trait Sink<Message>: Send + Sync + 'static {
92    fn send<'a>(&'a self, message: Message) -> BoxSend<'a>;
93
94    fn blocking_send(&self, message: Message) -> Result<(), SendError>;
95}
96
97pub(crate) type BoxSend<'a> = Pin<Box<dyn Future<Output = Result<(), SendError>> + Send + 'a>>;
98
99impl<Message: Send + 'static> Sender<Message> {
100    pub(crate) fn new(sink: Arc<dyn Sink<Message>>) -> Self {
101        Self { sink }
102    }
103
104    /// Sends a message, waiting where the source's admission says to.
105    pub async fn send(&self, message: Message) -> Result<(), SendError> {
106        self.sink.send(message).await
107    }
108
109    /// Sends a message from a thread that may block, for a source that is not
110    /// asynchronous.
111    pub fn blocking_send(&self, message: Message) -> Result<(), SendError> {
112        self.sink.blocking_send(message)
113    }
114
115    /// A sender that takes the source's own message and passes it through `f`
116    /// on the way in. This is how `map` reaches a source that is given a sender
117    /// rather than one that returns values.
118    pub(crate) fn contramap<From>(self, f: Mapper<From, Message>) -> Sender<From>
119    where
120        From: Send + 'static,
121    {
122        Sender::new(Arc::new(Mapped {
123            inner: self,
124            map: f,
125        }))
126    }
127}
128
129impl<Message> Clone for Sender<Message> {
130    fn clone(&self) -> Self {
131        Self {
132            sink: Arc::clone(&self.sink),
133        }
134    }
135}
136
137impl<Message> fmt::Debug for Sender<Message> {
138    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
139        formatter.write_str("Sender")
140    }
141}
142
143struct Mapped<From, To> {
144    inner: Sender<To>,
145    map: Mapper<From, To>,
146}
147
148impl<From, To> Sink<From> for Mapped<From, To>
149where
150    From: Send + 'static,
151    To: Send + 'static,
152{
153    fn send<'a>(&'a self, message: From) -> BoxSend<'a> {
154        let mapped = (self.map)(message);
155        Box::pin(self.inner.send(mapped))
156    }
157
158    fn blocking_send(&self, message: From) -> Result<(), SendError> {
159        self.inner.blocking_send((self.map)(message))
160    }
161}
162
163/// The runtime is no longer taking this source's messages.
164///
165/// A source that meets this has been stopped — the application stopped
166/// declaring it, or the runtime is shutting down — and should return.
167#[derive(Clone, Copy, Debug, PartialEq, Eq)]
168pub struct SendError;
169
170impl fmt::Display for SendError {
171    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
172        formatter.write_str("the runtime stopped taking this source's messages")
173    }
174}
175
176impl Error for SendError {}