Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ Mea (Make Easy Async) is a runtime-agnostic library providing essential synchron
* [**Condvar**](https://docs.rs/mea/*/mea/condvar/struct.Condvar.html): A condition variable that allows tasks to wait for a notification.
* [**Latch**](https://docs.rs/mea/*/mea/latch/struct.Latch.html): A synchronization primitive that allows one or more tasks to wait until a set of operations completes.
* [**Mutex**](https://docs.rs/mea/*/mea/mutex/struct.Mutex.html): A mutual exclusion primitive for protecting shared data.
* [**OnceCell**](https://docs.rs/mea/*/mea/once/struct.OnceCell.html): A cell that can be written to only once, providing safe, lazy initialization.
* [**RwLock**](https://docs.rs/mea/*/mea/rwlock/struct.RwLock.html): A reader-writer lock that allows multiple readers or a single writer at a time.
* [**Semaphore**](https://docs.rs/mea/*/mea/semaphore/struct.Semaphore.html): A synchronization primitive that controls access to a shared resource.
* [**ShutdownSend & ShutdownRecv**](https://docs.rs/mea/*/mea/shutdown/): A composite synchronization primitive for managing shutdown signals.
Expand Down Expand Up @@ -69,6 +70,7 @@ This crate collects runtime-agnostic synchronization primitives from spare parts
* **Condvar** is inspired by `std::sync::Condvar` and `async_std::sync::Condvar`, with a different implementation based on the internal `Semaphore` primitive. Different from the async_std implementation, this condvar is fair.
* **Latch** is inspired by [`latches`](https://github.com/mirromutth/latches), with a different implementation based on the internal `CountdownState` primitive. No `wait` or `watch` method is provided, since it can be easily implemented by [composing delay futures](https://docs.rs/fastimer/*/fastimer/fn.timeout.html). No sync variant is provided, since it can be easily implemented with block_on of any runtime.
* **Mutex** is derived from `tokio::sync::Mutex`. No blocking method is provided, since it can be easily implemented with block_on of any runtime.
* **OnceCell** is derived from `tokio::sync::OnceCell`, but using our own semaphore implementation.
* **RwLock** is derived from `tokio::sync::RwLock`, but the `max_readers` can be any `NonZeroUsize` (effectively any positive `usize`) instead of `[0, u32::MAX >> 3]`. No blocking method is provided, since it can be easily implemented with block_on of any runtime.
* **Semaphore** is derived from `tokio::sync::Semaphore`, without `close` method since it is quite tricky to use. And thus, this semaphore doesn't have the limitation of max permits. Besides, new methods like `forget_exact` are added to fit the specific use case.
* **WaitGroup** is inspired by [`waitgroup-rs`](https://github.com/laizy/waitgroup-rs), with a different implementation based on the internal `CountdownState` primitive. It fixes the unsound issue as described [here](https://github.com/rust-lang/futures-rs/issues/2880#issuecomment-2333842804).
Expand Down
4 changes: 3 additions & 1 deletion mea/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.

#![cfg_attr(docsrs, feature(doc_auto_cfg))]
#![cfg_attr(docsrs, feature(doc_cfg))]
#![deny(missing_docs)]

//! # Mea - Make Easy Async
Expand All @@ -28,6 +28,7 @@
//! * [`Condvar`]: A condition variable that allows tasks to wait for a notification
//! * [`Latch`]: A single-use barrier that allows one or more tasks to wait until a signal is given
//! * [`Mutex`]: A mutual exclusion primitive for protecting shared data
//! * [`OnceCell`]: A cell that can be initialized only once and provides safe concurrent access
//! * [`RwLock`]: A reader-writer lock that allows multiple readers or a single writer at a time
//! * [`Semaphore`]: A synchronization primitive that controls access to a shared resource
//! * [`ShutdownSend`] & [`ShutdownRecv`]: A composite synchronization primitive for managing
Expand Down Expand Up @@ -56,6 +57,7 @@
//! [`Condvar`]: condvar::Condvar
//! [`Latch`]: latch::Latch
//! [`Mutex`]: mutex::Mutex
//! [`OnceCell`]: once::OnceCell
//! [`RwLock`]: rwlock::RwLock
//! [`Semaphore`]: semaphore::Semaphore
//! [`ShutdownSend`]: shutdown::ShutdownSend
Expand Down
44 changes: 18 additions & 26 deletions mea/src/once/once_cell.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ use std::mem::MaybeUninit;
use std::sync::atomic::AtomicBool;
use std::sync::atomic::Ordering;

use crate::internal;
use crate::semaphore::Semaphore;

/// A thread-safe cell which can nominally be written to only once.
///
Expand Down Expand Up @@ -46,7 +46,7 @@ use crate::internal;
pub struct OnceCell<T> {
value_set: AtomicBool,
value: UnsafeCell<MaybeUninit<T>>,
semaphore: internal::Semaphore,
semaphore: Semaphore,

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Previously I was considering using the inernal Semaphore's notify_all method, but it seems a simply P/V pattern is enough here. Switch to the higher level encapsulation.

}

// SAFETY: OnceCell<T> can be shared between threads as long as T is Sync + Send.
Expand All @@ -67,7 +67,7 @@ impl<T> OnceCell<T> {
Self {
value_set: AtomicBool::new(false),
value: UnsafeCell::new(MaybeUninit::uninit()),
semaphore: internal::Semaphore::new(1),
semaphore: Semaphore::new(1),
}
}

Expand All @@ -78,7 +78,8 @@ impl<T> OnceCell<T> {

/// Returns the reference to the internal value or `None` if it is not set yet.
pub fn get(&self) -> Option<&T> {
if self.value_set.load(Ordering::Acquire) {
if self.initialized() {
// SAFETY: value was initialized
Some(unsafe { self.get_unchecked() })
} else {
None
Expand All @@ -101,27 +102,28 @@ impl<T> OnceCell<T> {
F: FnOnce() -> Fut,
Fut: Future<Output = T>,
{
if self.initialized() {
// SAFETY: We just checked that the value is initialized.
return unsafe { self.get_unchecked() };
if let Some(v) = self.get() {
return v;
}
self.semaphore.acquire(1).await;
let _guard = Guard {
semaphore: &self.semaphore,
};
if self.initialized() {
// Another task initialized the value while we were waiting for the semaphore.
// SAFETY: We just checked that the value is initialized.
return unsafe { self.get_unchecked() };

let _permit = self.semaphore.acquire(1).await;

if let Some(v) = self.get() {
// double-checked: another task initialized the value
// while we were waiting for the permit
return v;
}

let value = init().await;

let value_ptr = self.value.get();
unsafe { value_ptr.write(MaybeUninit::new(value)) };

// Use `store` with `Release` ordering to ensure that when loading it with `Acquire`
// ordering, the initialized value is visible.
self.value_set.store(true, Ordering::Release);
// SAFETY: value initialized one line above

// SAFETY: value initialized above
unsafe { self.get_unchecked() }
}

Expand All @@ -140,13 +142,3 @@ impl<T> Drop for OnceCell<T> {
}
}
}

struct Guard<'a> {
semaphore: &'a internal::Semaphore,
}

impl<'a> Drop for Guard<'a> {
fn drop(&mut self) {
self.semaphore.release(1);
}
}