mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-25 00:00:18 +02:00
sync: add watch::Sender::send_replace (#3962)
This commit is contained in:
+26
-4
@@ -58,6 +58,7 @@ use crate::sync::notify::Notify;
|
||||
use crate::loom::sync::atomic::AtomicUsize;
|
||||
use crate::loom::sync::atomic::Ordering::Relaxed;
|
||||
use crate::loom::sync::{Arc, RwLock, RwLockReadGuard};
|
||||
use std::mem;
|
||||
use std::ops;
|
||||
|
||||
/// Receives values from the associated [`Sender`](struct@Sender).
|
||||
@@ -427,15 +428,34 @@ impl<T> Sender<T> {
|
||||
/// This method fails if the channel has been closed, which happens when
|
||||
/// every receiver has been dropped.
|
||||
pub fn send(&self, value: T) -> Result<(), error::SendError<T>> {
|
||||
self.send_replace(value)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Sends a new value via the channel, notifying all receivers and returning
|
||||
/// the previous value in the channel.
|
||||
///
|
||||
/// This can be useful for reusing the buffers inside a watched value.
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// ```
|
||||
/// use tokio::sync::watch;
|
||||
///
|
||||
/// let (tx, _rx) = watch::channel(1);
|
||||
/// assert_eq!(tx.send_replace(2).unwrap(), 1);
|
||||
/// assert_eq!(tx.send_replace(3).unwrap(), 2);
|
||||
/// ```
|
||||
pub fn send_replace(&self, value: T) -> Result<T, error::SendError<T>> {
|
||||
// This is pretty much only useful as a hint anyway, so synchronization isn't critical.
|
||||
if 0 == self.receiver_count() {
|
||||
return Err(error::SendError(value));
|
||||
}
|
||||
|
||||
{
|
||||
let old = {
|
||||
// Acquire the write lock and update the value.
|
||||
let mut lock = self.shared.value.write().unwrap();
|
||||
*lock = value;
|
||||
let old = mem::replace(&mut *lock, value);
|
||||
|
||||
self.shared.state.increment_version();
|
||||
|
||||
@@ -445,12 +465,14 @@ impl<T> Sender<T> {
|
||||
// that receivers are able to figure out the version number of the
|
||||
// value they are currently looking at.
|
||||
drop(lock);
|
||||
}
|
||||
|
||||
old
|
||||
};
|
||||
|
||||
// Notify all watchers
|
||||
self.shared.notify_rx.notify_waiters();
|
||||
|
||||
Ok(())
|
||||
Ok(old)
|
||||
}
|
||||
|
||||
/// Returns a reference to the most recently sent value
|
||||
|
||||
Reference in New Issue
Block a user