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 {}