Skip to main content

tycho_simulation/snapshot_feed/
publisher.rs

1//! What a feed publishes through, and the only way it reaches a consumer.
2
3use std::{future::Future, time::Duration};
4
5use futures::{Stream, StreamExt};
6use tokio::{
7    select,
8    sync::watch,
9    time::{sleep_until, Instant},
10};
11use tracing::{debug, info, warn};
12
13/// A feed's end of the channel a consumer reads it through, and the age of what it holds.
14///
15/// A feed receives one of these from whichever consumer type runs it
16/// ([`SnapshotFeedStream`](crate::snapshot_feed::SnapshotFeedStream),
17/// [`SnapshotFeedStreams`](crate::snapshot_feed::SnapshotFeedStreams) or
18/// [`SnapshotFeedWatch`](crate::snapshot_feed::SnapshotFeedWatch)) and cannot make one: the
19/// channel, its `None` seed and the task the feed runs in all belong to whoever reads the feed,
20/// so no consumer can hold a feed's future and give it their own pace.
21///
22/// [`publishing`](Self::publishing) is the whole of what a feed does with it — hand it the
23/// snapshots and the age at which one goes stale, and everything else is taken care of.
24pub struct Publisher<T> {
25    tx: watch::Sender<Option<T>>,
26    /// When the published snapshot was published; `None` while nothing servable is published.
27    published_at: Option<Instant>,
28}
29
30impl<T> Publisher<T> {
31    /// A publisher and the receiver that reads it, seeded with `None` — nothing servable yet.
32    /// Crate-private on purpose: it is what keeps a feed's future out of consumer hands.
33    pub(crate) fn channel() -> (Self, watch::Receiver<Option<T>>) {
34        let (tx, rx) = watch::channel(None);
35        (Publisher { tx, published_at: None }, rx)
36    }
37
38    /// Publishes every snapshot `snapshots` produces, withdrawing one that goes `max_age` without
39    /// being refreshed, and resolves with the error the feed ended on. A `max_age` of `None` keeps
40    /// whatever was published until the feed ends.
41    ///
42    /// Resolves `Ok(())` when the stream runs out, which a feed holding a connection never does,
43    /// and as soon as the last receiver goes away, since everything the feed produces from then
44    /// on would reach nobody.
45    ///
46    /// Everything a feed does between two snapshots — ticking, fetching, connecting, reading,
47    /// backing off — belongs inside `snapshots`: this awaits nothing else, so that is the wait
48    /// staleness and the last receiver leaving can cut into. Work a feed does elsewhere runs to
49    /// its end before either is noticed.
50    pub async fn publishing<E>(
51        mut self,
52        max_age: Option<Duration>,
53        snapshots: impl Stream<Item = Result<T, E>>,
54    ) -> Result<(), E> {
55        tokio::pin!(snapshots);
56        loop {
57            match self
58                .next_snapshot(max_age, snapshots.next())
59                .await
60            {
61                Some(Ok(snapshot)) => self.publish(snapshot),
62                Some(Err(e)) => return Err(e),
63                None => return Ok(()),
64            }
65        }
66    }
67
68    /// The next snapshot the feed produces, or `None` when there will not be a useful one: the
69    /// stream ran out, or the last receiver went away and nothing the feed produces can reach
70    /// anyone.
71    ///
72    /// Withdraws the published snapshot along the way if it reaches `max_age` first. Going stale
73    /// never cuts the wait short — only the stream, or the last receiver leaving, decides when
74    /// this returns — because a snapshot goes stale by the clock rather than by anything the feed
75    /// is waiting for.
76    async fn next_snapshot<E>(
77        &mut self,
78        max_age: Option<Duration>,
79        next: impl Future<Output = Option<Result<T, E>>>,
80    ) -> Option<Result<T, E>> {
81        tokio::pin!(next);
82        // Split so the wait can watch the receivers while the deadline still updates what is
83        // published.
84        let Publisher { tx, published_at } = self;
85
86        select! {
87            biased;
88            _ = tx.closed() => None,
89            max_age = going_stale(*published_at, max_age) => {
90                warn!(
91                    ?max_age,
92                    "no snapshot received within max_snapshot_age, withdrawing the published one"
93                );
94                withdraw(tx, published_at);
95
96                // Nothing is published any more, so nothing else can go stale before the
97                // snapshot this is still waiting for arrives.
98                next.await
99            }
100            out = &mut next => out,
101        }
102    }
103
104    /// Publishes a snapshot, replacing whatever was being served.
105    fn publish(&mut self, snapshot: T) {
106        if self.tx.send(Some(snapshot)).is_err() {
107            // The last receiver went away between the wait that produced this snapshot and here:
108            // it was not published, so there is nothing to announce and nothing whose age to
109            // track. The next wait ends the feed.
110            return;
111        }
112
113        match self.published_at {
114            // A consumer that was being served now has something newer; one that was not can
115            // start.
116            Some(_) => debug!("snapshot refreshed"),
117            None => info!("serving a snapshot"),
118        }
119
120        self.published_at = Some(Instant::now());
121    }
122}
123
124/// Completes once a snapshot published at `published_at` has reached `max_age`, and never if
125/// there is no snapshot, or no age for one to reach. Hands back the age it waited out, which the
126/// caller reports.
127async fn going_stale(published_at: Option<Instant>, max_age: Option<Duration>) -> Duration {
128    let Some((published_at, max_age)) = Option::zip(published_at, max_age) else {
129        return std::future::pending().await;
130    };
131    sleep_until(published_at + max_age).await;
132    max_age
133}
134
135/// Takes the published snapshot away and stops its clock. An already-empty channel stays quiet:
136/// receivers are notified only if there was a snapshot to take.
137fn withdraw<T>(tx: &watch::Sender<Option<T>>, published_at: &mut Option<Instant>) {
138    *published_at = None;
139    tx.send_if_modified(|snapshot| snapshot.take().is_some());
140}
141
142impl<T> Drop for Publisher<T> {
143    /// Whatever ends a feed ends its snapshot: it gave up, ran out, or its task was dropped
144    /// because nobody was reading. Nothing refreshes the published snapshot after that, so it is
145    /// withdrawn rather than left standing for consumers to keep using.
146    fn drop(&mut self) {
147        withdraw(&self.tx, &mut self.published_at);
148    }
149}
150
151#[cfg(test)]
152mod tests {
153    use std::sync::{
154        atomic::{AtomicBool, Ordering},
155        Arc,
156    };
157
158    use futures::{stream, StreamExt};
159    use rstest::rstest;
160    use tokio::time::sleep;
161    use tokio_stream::wrappers::WatchStream;
162
163    use super::{super::expect_to_finish, *};
164
165    #[tokio::test(start_paused = true)]
166    async fn next_snapshot_withdraws_at_max_age_and_keeps_waiting() {
167        let (mut publisher, mut rx) = Publisher::channel();
168        publisher.publish(1);
169        assert!(rx.borrow_and_update().is_some(), "a snapshot must be published to go stale");
170
171        let start = Instant::now();
172        let next = publisher
173            .next_snapshot(Some(Duration::from_millis(100)), async {
174                rx.changed()
175                    .await
176                    .expect("the publisher outlives this wait");
177                assert!(rx.borrow_and_update().is_none(), "the change must be the withdrawal");
178                assert_eq!(
179                    start.elapsed(),
180                    Duration::from_millis(100),
181                    "the snapshot must be withdrawn once it reaches max_age"
182                );
183
184                sleep(Duration::from_millis(200)).await;
185                Some(Ok::<u32, String>(2))
186            })
187            .await;
188
189        assert_eq!(next, Some(Ok(2)), "the snapshot the feed was waiting for still arrives");
190        assert_eq!(
191            start.elapsed(),
192            Duration::from_millis(300),
193            "going stale must not cut the wait short: the call returns when the feed answers"
194        );
195    }
196
197    #[rstest]
198    #[case::no_max_age_never_goes_stale(None, Duration::from_secs(3600))]
199    #[case::snapshot_arrives_before_max_age(
200        Some(Duration::from_millis(100)),
201        Duration::from_millis(50)
202    )]
203    #[tokio::test(start_paused = true)]
204    async fn next_snapshot_keeps_a_snapshot_that_has_not_gone_stale(
205        #[case] max_age: Option<Duration>,
206        #[case] wait: Duration,
207    ) {
208        let (mut publisher, mut rx) = Publisher::channel();
209        publisher.publish(1);
210        rx.borrow_and_update();
211
212        publisher
213            .next_snapshot(max_age, async move {
214                sleep(wait).await;
215                Some(Ok::<u32, String>(2))
216            })
217            .await;
218
219        assert!(rx.borrow().is_some());
220        assert!(!rx.has_changed().unwrap(), "the watch must not be disturbed");
221    }
222
223    #[tokio::test(start_paused = true)]
224    async fn publishing_stops_as_soon_as_the_last_receiver_is_gone() {
225        // A stream that never answers, watched by nobody: the feed must not be left waiting on
226        // the venue for a snapshot that could reach no one.
227        let asked = Arc::new(AtomicBool::new(false));
228        let asked_clone = Arc::clone(&asked);
229        let snapshots = stream::once(async move {
230            asked_clone.store(true, Ordering::SeqCst);
231            std::future::pending::<Result<u32, String>>().await
232        });
233
234        let (publisher, rx) = Publisher::channel();
235        drop(rx);
236
237        expect_to_finish(
238            "the feed kept waiting although nobody was reading",
239            publisher.publishing(None, snapshots),
240        )
241        .await
242        .expect("no reader left is not a failure of the feed");
243
244        assert!(!asked.load(Ordering::SeqCst), "the stream is not even asked for a snapshot");
245    }
246
247    #[tokio::test]
248    async fn dropping_the_publisher_withdraws_what_it_published() {
249        // Read as a consumer does, through the changes: the publisher owns the sender, so once
250        // it is gone the channel is closed and only what it sent on the way out is left.
251        let (mut publisher, rx) = Publisher::channel();
252        let mut changes = WatchStream::from_changes(rx);
253        publisher.publish(1);
254        assert_eq!(changes.next().await, Some(Some(1)));
255
256        drop(publisher);
257
258        assert_eq!(changes.next().await, Some(None), "dropping it withdraws what it published");
259        assert_eq!(changes.next().await, None, "and ends the stream");
260    }
261
262    #[tokio::test]
263    async fn dropping_a_publisher_that_published_nothing_leaves_the_watch_quiet() {
264        let (publisher, rx) = Publisher::<u32>::channel();
265        let mut changes = WatchStream::from_changes(rx);
266        drop(publisher);
267
268        assert_eq!(changes.next().await, None, "there was no snapshot to take away");
269    }
270}