saluki_core/observability/metrics/
reflector.rs

1//! Mechanisms for processing a data source and sharing the processed results.
2use std::sync::Arc;
3
4use async_trait::async_trait;
5use futures::{Stream, StreamExt};
6use saluki_common::sync::shutdown::ShutdownHandle;
7use tokio::{
8    select,
9    sync::{Mutex, Notify},
10};
11
12use crate::runtime::{InitializationError, Supervisable, SupervisorFuture};
13
14/// Processes input data and modifies shared state based on the result.
15pub trait Processor: Send + Sync {
16    /// The type of input to the processor.
17    type Input;
18
19    /// The state that the processor acts on.
20    type State: Send + Sync;
21
22    /// Builds the initial state for the processor.
23    fn build_initial_state(&self) -> Self::State;
24
25    /// Processes the input, potentially updating the reflector state.
26    fn process(&self, input: Self::Input, state: &Self::State);
27}
28
29struct StoreInner<P: Processor> {
30    processor: P,
31    state: P::State,
32    notify_update: Notify,
33}
34
35/// Shared state based on a processor.
36///
37/// `Store` acts as the glue between a processor and the data that it processes. It acts as the entrypoint for taking in
38/// a group of items from a data source, running them through the configured processor, and then notifying callers
39/// that an update has taken place.
40struct Store<P: Processor> {
41    inner: Arc<StoreInner<P>>,
42}
43
44impl<P: Processor> Store<P> {
45    fn from_processor(processor: P) -> Self {
46        let state = processor.build_initial_state();
47        Self {
48            inner: Arc::new(StoreInner {
49                processor,
50                state,
51                notify_update: Notify::const_new(),
52            }),
53        }
54    }
55
56    pub fn process<I>(&self, inputs: I)
57    where
58        I: IntoIterator<Item = P::Input>,
59    {
60        for input in inputs {
61            self.inner.processor.process(input, &self.inner.state);
62        }
63        self.inner.notify_update.notify_waiters();
64    }
65
66    pub async fn wait_for_update(&self) {
67        self.inner.notify_update.notified().await;
68    }
69
70    pub fn state(&self) -> &P::State {
71        &self.inner.state
72    }
73}
74
75impl<P: Processor> Clone for Store<P> {
76    fn clone(&self) -> Self {
77        Self {
78            inner: Arc::clone(&self.inner),
79        }
80    }
81}
82
83/// `Reflector` composes a source of data with a processor that's used to transform the data, and then stores the
84/// results and allows for shared access by multiple callers.
85///
86/// Reflectors are a term often found in the context of custom Kubernetes controllers, where they're used to reduce the
87/// load on the Kubernetes API server by caching the state of resources in memory. `Reflector` provides comparable
88/// functionality, allowing for a single data source to be consumed, and then shared amongst multiple callers. However,
89///
90/// `Reflector` utilizes the concept of a _processor_, which dictates both the type of data that can be consumed and
91/// data that gets stored. This means that `Reflector` is more than just a cache of the data source, but also
92/// potentially a mapped version of it, allowing for transforming the data in whatever way is necessary.
93pub struct Reflector<P: Processor> {
94    store: Store<P>,
95}
96
97impl<P: Processor> Reflector<P> {
98    /// Creates a new reflector with the given data source and processor.
99    ///
100    /// A reflector composes a source of data with a processor that's used to transform the data, and then stores
101    /// the processed results. It can be listened to for updates, and cheaply shared. This allows multiple interested
102    /// components to subscribe to the same data source without having to duplicate the processing or storage of the
103    /// data.
104    ///
105    /// Returns the reflector alongside a [`ReflectorWorker`] that consumes the data source and feeds the processed
106    /// items into the shared state. The reflector is usable immediately, but reports only the initial state until
107    /// the worker is added to a [`Supervisor`][crate::runtime::Supervisor] and starts running. Register the worker
108    /// as transient, for the reasons given on [`ReflectorWorker`].
109    ///
110    /// `Reflector` is cheaply cloneable and can either be cloned for each caller or shared between them (for example, via
111    /// `Arc<T>`).
112    pub fn new<S, I>(source: S, processor: P) -> (Self, ReflectorWorker<P, S>)
113    where
114        S: Stream<Item = I> + Unpin + Send + 'static,
115        I: IntoIterator<Item = P::Input> + Send,
116        P: 'static,
117    {
118        let store = Store::from_processor(processor);
119        let worker = ReflectorWorker {
120            store: store.clone(),
121            source: Arc::new(Mutex::new(source)),
122        };
123
124        (Self { store }, worker)
125    }
126}
127
128impl<P: Processor> Clone for Reflector<P> {
129    fn clone(&self) -> Self {
130        Self {
131            store: self.store.clone(),
132        }
133    }
134}
135
136/// A worker that drives a [`Reflector`]'s data source.
137///
138/// Consumes items from the source and feeds them through the processor into the reflector's shared state. Until this
139/// worker runs, the reflector it was created with reports only the initial state its processor built.
140///
141/// The source is retained across restarts rather than being rebuilt, so a worker that fails and is restarted resumes
142/// from wherever the source left off. That matters for a source backed by a subscription: rebuilding it would drop
143/// whatever accumulated while the worker was down.
144///
145/// Register this worker as [`transient`][crate::runtime::ChildBuilder::transient] rather than with the permanent
146/// default of [`Supervisor::add_worker`][crate::runtime::Supervisor::add_worker]. An exhausted source is terminal
147/// here: the worker returns normally once the source ends, and the retained source means a restart would only feed
148/// it the same dead source, exit immediately again, and burn through the supervisor's restart budget. Transient
149/// leaves it stopped instead, which is the intended outcome. A source that cannot end -- a subscription held open
150/// for the life of the process, say -- never reaches this case, but nothing about the type guarantees that.
151pub struct ReflectorWorker<P: Processor, S> {
152    store: Store<P>,
153    source: Arc<Mutex<S>>,
154}
155
156#[async_trait]
157impl<P, S, I> Supervisable for ReflectorWorker<P, S>
158where
159    P: Processor + 'static,
160    S: Stream<Item = I> + Unpin + Send + 'static,
161    I: IntoIterator<Item = P::Input> + Send,
162{
163    fn name(&self) -> &str {
164        "reflector"
165    }
166
167    async fn initialize(&self, process_shutdown: ShutdownHandle) -> Result<SupervisorFuture, InitializationError> {
168        let store = self.store.clone();
169        let source = Arc::clone(&self.source);
170
171        Ok(Box::pin(async move {
172            let mut source = source.lock_owned().await;
173
174            select! {
175                _ = process_shutdown => {},
176                _ = drive_source(&mut *source, &store) => {},
177            }
178
179            Ok(())
180        }))
181    }
182}
183
184/// Feeds every item the source yields through the store, returning once the source is exhausted.
185async fn drive_source<P, S, I>(source: &mut S, store: &Store<P>)
186where
187    P: Processor,
188    S: Stream<Item = I> + Unpin,
189    I: IntoIterator<Item = P::Input>,
190{
191    while let Some(inputs) = source.next().await {
192        store.process(inputs);
193    }
194}
195
196impl<P: Processor> Reflector<P> {
197    /// Waits for the next update to the reflector.
198    ///
199    /// When this method completes, callers must query the reflector to acquire the latest state.
200    pub async fn wait_for_update(&self) {
201        self.store.wait_for_update().await;
202    }
203
204    /// Returns a reference a to the reflector's state.
205    pub fn state(&self) -> &P::State {
206        self.store.state()
207    }
208}
209
210#[cfg(test)]
211mod tests {
212    use std::{
213        pin::Pin,
214        sync::Mutex as StdMutex,
215        task::{Context, Poll},
216        time::Duration,
217    };
218
219    use tokio::{sync::mpsc, time::timeout};
220
221    use super::*;
222
223    /// Accumulates every input it is handed, in order.
224    struct TestProcessor;
225
226    impl Processor for TestProcessor {
227        type Input = u32;
228        type State = StdMutex<Vec<u32>>;
229
230        fn build_initial_state(&self) -> Self::State {
231            StdMutex::new(Vec::new())
232        }
233
234        fn process(&self, input: Self::Input, state: &Self::State) {
235            state.lock().unwrap().push(input);
236        }
237    }
238
239    /// A source fed by a channel, so tests control exactly when items become available.
240    struct TestSource(mpsc::UnboundedReceiver<Vec<u32>>);
241
242    impl Stream for TestSource {
243        type Item = Vec<u32>;
244
245        fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
246            self.0.poll_recv(cx)
247        }
248    }
249
250    fn build() -> (
251        mpsc::UnboundedSender<Vec<u32>>,
252        Reflector<TestProcessor>,
253        ReflectorWorker<TestProcessor, TestSource>,
254    ) {
255        let (tx, rx) = mpsc::unbounded_channel();
256        let (reflector, worker) = Reflector::new(TestSource(rx), TestProcessor);
257        (tx, reflector, worker)
258    }
259
260    fn observed(reflector: &Reflector<TestProcessor>) -> Vec<u32> {
261        reflector.state().lock().unwrap().clone()
262    }
263
264    /// Waits until the reflector has observed at least `expected` items, so that tests synchronize on the worker
265    /// having made progress rather than on a fixed delay.
266    async fn wait_for_len(reflector: &Reflector<TestProcessor>, expected: usize) {
267        timeout(Duration::from_secs(5), async {
268            while observed(reflector).len() < expected {
269                tokio::task::yield_now().await;
270            }
271        })
272        .await
273        .expect("timed out waiting for the reflector to observe the expected items");
274    }
275
276    #[tokio::test]
277    async fn reports_initial_state_before_worker_runs() {
278        // Creating a reflector must not depend on its worker: callers acquire the handle during bootstrap, well
279        // before the supervisor that drives the worker is running.
280        let (tx, reflector, _worker) = build();
281
282        tx.send(vec![1, 2, 3]).unwrap();
283
284        assert_eq!(observed(&reflector), Vec::<u32>::new());
285    }
286
287    #[tokio::test]
288    async fn worker_feeds_source_items_into_shared_state() {
289        let (tx, reflector, worker) = build();
290
291        tx.send(vec![1, 2]).unwrap();
292        tx.send(vec![3]).unwrap();
293        drop(tx);
294
295        let fut = worker.initialize(ShutdownHandle::noop()).await.unwrap();
296        fut.await.unwrap();
297
298        assert_eq!(observed(&reflector), vec![1, 2, 3]);
299    }
300
301    #[tokio::test]
302    async fn worker_completes_when_source_ends() {
303        let (tx, _reflector, worker) = build();
304        drop(tx);
305
306        let fut = worker.initialize(ShutdownHandle::noop()).await.unwrap();
307
308        // An exhausted source is a terminal condition, so the worker returns rather than waiting for shutdown.
309        timeout(Duration::from_secs(5), fut)
310            .await
311            .expect("worker should complete once the source is exhausted")
312            .unwrap();
313    }
314
315    #[tokio::test]
316    async fn worker_stops_on_shutdown_signal() {
317        // The source stays open here, so only the shutdown signal can end the worker.
318        let (_tx, _reflector, worker) = build();
319        let (coordinator, shutdown) = ShutdownHandle::paired();
320
321        let fut = worker.initialize(shutdown).await.unwrap();
322        let handle = tokio::spawn(fut);
323
324        coordinator.shutdown();
325
326        timeout(Duration::from_secs(5), handle)
327            .await
328            .expect("worker should stop once shutdown is signalled")
329            .unwrap()
330            .unwrap();
331    }
332
333    #[tokio::test]
334    async fn restarted_worker_resumes_from_the_same_source() {
335        // A restart re-enters `initialize`, which must not rebuild the source: doing so would drop everything the
336        // source accumulated while the worker was down.
337        let (tx, reflector, worker) = build();
338
339        tx.send(vec![1]).unwrap();
340
341        let fut = worker.initialize(ShutdownHandle::noop()).await.unwrap();
342        let handle = tokio::spawn(fut);
343        wait_for_len(&reflector, 1).await;
344
345        // Abort rather than shut down cleanly, standing in for the worker failing mid-flight.
346        handle.abort();
347        let _ = handle.await;
348
349        // Sent while nothing is draining the source, so it can only be observed if the restart reuses it.
350        tx.send(vec![2]).unwrap();
351        drop(tx);
352
353        let fut = worker.initialize(ShutdownHandle::noop()).await.unwrap();
354        fut.await.unwrap();
355
356        assert_eq!(observed(&reflector), vec![1, 2]);
357    }
358
359    #[tokio::test]
360    async fn state_is_shared_across_clones() {
361        let (tx, reflector, worker) = build();
362        let cloned = reflector.clone();
363
364        tx.send(vec![7]).unwrap();
365        drop(tx);
366
367        let fut = worker.initialize(ShutdownHandle::noop()).await.unwrap();
368        fut.await.unwrap();
369
370        assert_eq!(observed(&cloned), vec![7]);
371    }
372}