sync: reduce memory size of watch::Receiver (#2191)

This reduces the `mem::size_of::<watch::Receiver>()` from 4 words to 2.

- The `id` is now the pointer of the `Arc<WatchInner>`.
- The `ver` is moved into the `WatchInner`.
This commit is contained in:
Sean McArthur
2020-01-29 12:22:21 -08:00
committed by GitHub
parent 9d6b99494b
commit 116a18b849
+59 -60
View File
@@ -54,10 +54,10 @@
use crate::future::poll_fn; use crate::future::poll_fn;
use crate::sync::task::AtomicWaker; use crate::sync::task::AtomicWaker;
use fnv::FnvHashMap; use fnv::FnvHashSet;
use std::ops; use std::ops;
use std::sync::atomic::AtomicUsize; use std::sync::atomic::AtomicUsize;
use std::sync::atomic::Ordering::SeqCst; use std::sync::atomic::Ordering::{Relaxed, SeqCst};
use std::sync::{Arc, Mutex, RwLock, RwLockReadGuard, Weak}; use std::sync::{Arc, Mutex, RwLock, RwLockReadGuard, Weak};
use std::task::Poll::{Pending, Ready}; use std::task::Poll::{Pending, Ready};
use std::task::{Context, Poll}; use std::task::{Context, Poll};
@@ -71,13 +71,7 @@ pub struct Receiver<T> {
shared: Arc<Shared<T>>, shared: Arc<Shared<T>>,
/// Pointer to the watcher's internal state /// Pointer to the watcher's internal state
inner: Arc<WatchInner>, inner: Watcher,
/// Watcher ID.
id: u64,
/// Last observed version
ver: usize,
} }
/// Sends values to the associated [`Receiver`](struct.Receiver.html). /// Sends values to the associated [`Receiver`](struct.Receiver.html).
@@ -138,14 +132,16 @@ struct Shared<T> {
cancel: AtomicWaker, cancel: AtomicWaker,
} }
#[derive(Debug)] type Watchers = FnvHashSet<Watcher>;
struct Watchers {
next_id: u64, /// The watcher's ID is based on the Arc's pointer.
watchers: FnvHashMap<u64, Arc<WatchInner>>, #[derive(Clone, Debug)]
} struct Watcher(Arc<WatchInner>);
#[derive(Debug)] #[derive(Debug)]
struct WatchInner { struct WatchInner {
/// Last observed version
version: AtomicUsize,
waker: AtomicWaker, waker: AtomicWaker,
} }
@@ -179,21 +175,20 @@ const CLOSED: usize = 1;
/// [`Sender`]: struct.Sender.html /// [`Sender`]: struct.Sender.html
/// [`Receiver`]: struct.Receiver.html /// [`Receiver`]: struct.Receiver.html
pub fn channel<T: Clone>(init: T) -> (Sender<T>, Receiver<T>) { pub fn channel<T: Clone>(init: T) -> (Sender<T>, Receiver<T>) {
const INIT_ID: u64 = 0; const VERSION_0: usize = 0;
const VERSION_1: usize = 2;
let inner = Arc::new(WatchInner::new()); // We don't start knowing VERSION_1
let inner = Watcher::new_version(VERSION_0);
// Insert the watcher // Insert the watcher
let mut watchers = FnvHashMap::with_capacity_and_hasher(0, Default::default()); let mut watchers = FnvHashSet::with_capacity_and_hasher(0, Default::default());
watchers.insert(INIT_ID, inner.clone()); watchers.insert(inner.clone());
let shared = Arc::new(Shared { let shared = Arc::new(Shared {
value: RwLock::new(init), value: RwLock::new(init),
version: AtomicUsize::new(2), version: AtomicUsize::new(VERSION_1),
watchers: Mutex::new(Watchers { watchers: Mutex::new(watchers),
next_id: INIT_ID + 1,
watchers,
}),
cancel: AtomicWaker::new(), cancel: AtomicWaker::new(),
}); });
@@ -201,12 +196,7 @@ pub fn channel<T: Clone>(init: T) -> (Sender<T>, Receiver<T>) {
shared: Arc::downgrade(&shared), shared: Arc::downgrade(&shared),
}; };
let rx = Receiver { let rx = Receiver { shared, inner };
shared,
inner,
id: INIT_ID,
ver: 0,
};
(tx, rx) (tx, rx)
} }
@@ -240,9 +230,8 @@ impl<T> Receiver<T> {
let state = self.shared.version.load(SeqCst); let state = self.shared.version.load(SeqCst);
let version = state & !CLOSED; let version = state & !CLOSED;
if version != self.ver { if self.inner.version.swap(version, Relaxed) != version {
let inner = self.shared.value.read().unwrap(); let inner = self.shared.value.read().unwrap();
self.ver = version;
return Ready(Some(Ref { inner })); return Ready(Some(Ref { inner }));
} }
@@ -312,42 +301,19 @@ impl<T: Clone> crate::stream::Stream for Receiver<T> {
impl<T> Clone for Receiver<T> { impl<T> Clone for Receiver<T> {
fn clone(&self) -> Self { fn clone(&self) -> Self {
let inner = Arc::new(WatchInner::new()); let ver = self.inner.version.load(Relaxed);
let inner = Watcher::new_version(ver);
let shared = self.shared.clone(); let shared = self.shared.clone();
let id = { shared.watchers.lock().unwrap().insert(inner.clone());
let mut watchers = shared.watchers.lock().unwrap();
let id = watchers.next_id;
watchers.next_id += 1; Receiver { shared, inner }
watchers.watchers.insert(id, inner.clone());
id
};
let ver = self.ver;
Receiver {
shared,
inner,
id,
ver,
}
} }
} }
impl<T> Drop for Receiver<T> { impl<T> Drop for Receiver<T> {
fn drop(&mut self) { fn drop(&mut self) {
let mut watchers = self.shared.watchers.lock().unwrap(); self.shared.watchers.lock().unwrap().remove(&self.inner);
watchers.watchers.remove(&self.id);
}
}
impl WatchInner {
fn new() -> Self {
WatchInner {
waker: AtomicWaker::new(),
}
} }
} }
@@ -399,7 +365,7 @@ impl<T> Sender<T> {
fn notify_all<T>(shared: &Shared<T>) { fn notify_all<T>(shared: &Shared<T>) {
let watchers = shared.watchers.lock().unwrap(); let watchers = shared.watchers.lock().unwrap();
for watcher in watchers.watchers.values() { for watcher in watchers.iter() {
// Notify the task // Notify the task
watcher.waker.wake(); watcher.waker.wake();
} }
@@ -431,3 +397,36 @@ impl<T> Drop for Shared<T> {
self.cancel.wake(); self.cancel.wake();
} }
} }
// ===== impl Watcher =====
impl Watcher {
fn new_version(version: usize) -> Self {
Watcher(Arc::new(WatchInner {
version: AtomicUsize::new(version),
waker: AtomicWaker::new(),
}))
}
}
impl std::cmp::PartialEq for Watcher {
fn eq(&self, other: &Watcher) -> bool {
Arc::ptr_eq(&self.0, &other.0)
}
}
impl std::cmp::Eq for Watcher {}
impl std::hash::Hash for Watcher {
fn hash<H: std::hash::Hasher>(&self, state: &mut H) {
(&*self.0 as *const WatchInner).hash(state)
}
}
impl std::ops::Deref for Watcher {
type Target = WatchInner;
fn deref(&self) -> &Self::Target {
&self.0
}
}