sync.rsannotatedsync.rssource65 lines · 2.1 KB · raw
1use std::sync::atomic::{AtomicBool, Ordering};
2use std::sync::{Arc, Mutex};
3
4struct SharedVarState<T> {
5    locked: Mutex<Option<T>>,
6    updated: AtomicBool,
7}
8
9pub struct SharedVarReceiver<T> {
10    latest: Option<T>,
11    state: Arc<SharedVarState<T>>,
12}
13
14#[derive(Clone)]
15pub struct SharedVarSender<T> {
16    state: Arc<SharedVarState<T>>,
17}
18
19impl<T> SharedVarReceiver<T> {

Returns the most recently received value, or None if no value has been received yet.

Returns a mutable reference for convenience, but any mutations will only affect the current value and will get discarded when a new value is received.

24    #[must_use]
25    #[allow(clippy::missing_panics_doc)] // Mutex poisoning is impossible here
26    pub fn get(&mut self) -> Option<&mut T> {
27        if self.state.updated.load(Ordering::Relaxed)
28            && self.state.updated.compare_exchange(true, false, Ordering::AcqRel, Ordering::Relaxed)
29                == Ok(true)
30            && let Some(value) = self.state.locked.lock().unwrap().take()
31        {
32            self.latest = Some(value);
33        }
34
35        self.latest.as_mut()
36    }
37}
39impl<T> SharedVarSender<T> {

Replace the current value with a new one.

41    #[allow(clippy::missing_panics_doc)] // Mutex poisoning is impossible here
42    pub fn update(&self, value: T) {
43        *self.state.locked.lock().unwrap() = Some(value);
44        self.state.updated.store(true, Ordering::Release);
45    }
46}

Creates a shared var.

This is similar to an SPSC channel, but the receiver retains the most recent value received and will return it on repeated [SharedVarReceiver::get] calls. It only retains the most recent value, so the receiver will miss values if the sender sends multiple values in between [SharedVarReceiver::get] calls.

54pub fn new_shared_var<T>() -> (SharedVarSender<T>, SharedVarReceiver<T>) {
55    let sender = SharedVarSender {
56        state: Arc::new(SharedVarState {
57            locked: Mutex::new(None),
58            updated: AtomicBool::new(false),
59        }),
60    };
61
62    let receiver = SharedVarReceiver { latest: None, state: Arc::clone(&sender.state) };
63
64    (sender, receiver)
65}