Files
MetaCrate/crates/libremetaverse/src/threading_runtime.rs
Chili Palmer c9a1170a27
Some checks failed
API and SemVer surface / api-surface (push) Failing after 1m34s
Native code generation / deterministic (push) Has been cancelled
Concurrency and resource soak audit / soak (push) Has been cancelled
Documentation / documentation (push) Has been cancelled
performance evidence / audit (push) Has been cancelled
First release candidate / non-fuzz-release-gate (push) Has been cancelled
Release platform and feature matrix / audit (push) Has been cancelled
Release platform and feature matrix / matrix (false, linux-stable-minimal, x86_64-unknown-linux-gnu, stable) (push) Has been cancelled
Release platform and feature matrix / matrix (false, macos-stable-portable, x86_64-apple-darwin, stable) (push) Has been cancelled
Release platform and feature matrix / matrix (false, windows-stable-portable, x86_64-pc-windows-gnu, stable) (push) Has been cancelled
Release platform and feature matrix / matrix (true, linux-msrv-portable, x86_64-unknown-linux-gnu, 1.96.0) (push) Has been cancelled
Release platform and feature matrix / matrix (true, linux-stable-default, x86_64-unknown-linux-gnu, stable) (push) Has been cancelled
Release platform and feature matrix / matrix (true, linux-stable-features, x86_64-unknown-linux-gnu, stable) (push) Has been cancelled
Release platform and feature matrix / matrix (true, linux-stable-release-surface, x86_64-unknown-linux-gnu, stable) (push) Has been cancelled
Dependency and supply-chain audit / audit (push) Has been cancelled
Native release artifact audit / audit (push) Has been cancelled
Imaging and meshing gate / native (push) Failing after 53s
JPEG 2000 feature / linux (push) Successful in 2m49s
Native Rust workspace compile / compile (push) Failing after 59s
Skia feature / linux (push) Has been cancelled
Complete first release candidate audit (#107)
2026-08-12 14:44:28 +00:00

1009 lines
30 KiB
Rust

//! Cross-platform implementations of the pinned managed synchronization APIs.
use crate::Error;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Condvar, Mutex, MutexGuard};
use std::time::{Duration, Instant};
/// Compact non-reentrant spin lock matching the mapped C# value type.
#[derive(Debug, Default)]
pub struct SpinLockSlim {
locked: AtomicBool,
}
impl SpinLockSlim {
pub fn enter(&mut self) -> Result<(), Error> {
let mut spins = 0_u32;
while self
.locked
.compare_exchange_weak(false, true, Ordering::Acquire, Ordering::Relaxed)
.is_err()
{
if spins < 64 {
std::hint::spin_loop();
} else {
std::thread::yield_now();
}
spins = spins.saturating_add(1);
}
Ok(())
}
pub fn exit(&mut self) -> Result<(), Error> {
self.locked.store(false, Ordering::Release);
Ok(())
}
}
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
#[derive(Debug)]
struct EventState {
signalled: bool,
disposed: bool,
}
#[derive(Debug)]
struct ManagedEvent {
state: Mutex<EventState>,
changed: Condvar,
auto_reset: bool,
}
impl ManagedEvent {
fn new(initial_state: bool, auto_reset: bool) -> Self {
Self {
state: Mutex::new(EventState {
signalled: initial_state,
disposed: false,
}),
changed: Condvar::new(),
auto_reset,
}
}
fn dispose(&self) {
let mut state = lock(&self.state);
state.disposed = true;
state.signalled = true;
self.changed.notify_all();
}
fn reset(&self) {
let mut state = lock(&self.state);
if !state.disposed {
state.signalled = false;
}
}
fn set(&self) {
let mut state = lock(&self.state);
state.signalled = true;
if self.auto_reset {
self.changed.notify_one();
} else {
self.changed.notify_all();
}
}
fn wait(&self, timeout: Option<Duration>) -> bool {
let mut state = lock(&self.state);
if let Some(timeout) = timeout {
let deadline = Instant::now().checked_add(timeout);
while !state.signalled {
let Some(remaining) =
deadline.and_then(|value| value.checked_duration_since(Instant::now()))
else {
return false;
};
if remaining.is_zero() {
return false;
}
let (next, result) = self
.changed
.wait_timeout(state, remaining)
.unwrap_or_else(std::sync::PoisonError::into_inner);
state = next;
if result.timed_out() && !state.signalled {
return false;
}
}
} else {
while !state.signalled {
state = self
.changed
.wait(state)
.unwrap_or_else(std::sync::PoisonError::into_inner);
}
}
if self.auto_reset && !state.disposed {
state.signalled = false;
}
true
}
fn is_set(&self) -> bool {
lock(&self.state).signalled
}
fn was_disposed(&self) -> bool {
lock(&self.state).disposed
}
}
pub trait IEventWait: Send + Sync {
fn reset(&self) -> Result<(), Error>;
fn set(&self) -> Result<(), Error>;
fn wait_one_with_method(&self) -> Result<(), Error>;
fn wait_one_with_int32(&self, milliseconds_timeout: i32) -> Result<bool, Error>;
fn wait_one_with_time_span(&self, timeout: Duration) -> Result<bool, Error>;
}
pub trait IAdvancedDisposable: Send + Sync {
fn was_disposed(&self) -> bool;
}
macro_rules! managed_event {
($name:ident, $auto_reset:literal) => {
#[derive(Debug)]
pub struct $name(ManagedEvent);
impl $name {
pub fn new_with_constructor() -> Result<Self, Error> {
Self::new_with_boolean(false)
}
pub fn new_with_boolean(initial_state: bool) -> Result<Self, Error> {
Ok(Self(ManagedEvent::new(initial_state, $auto_reset)))
}
pub fn dispose(&self) -> Result<(), Error> {
self.0.dispose();
Ok(())
}
pub fn reset(&self) -> Result<(), Error> {
self.0.reset();
Ok(())
}
pub fn set(&self) -> Result<(), Error> {
self.0.set();
Ok(())
}
pub fn wait_one_with_method(&self) -> Result<(), Error> {
self.0.wait(None);
Ok(())
}
pub fn wait_one_with_int32(&self, milliseconds_timeout: i32) -> Result<bool, Error> {
match milliseconds_timeout {
-1 => Ok(self.0.wait(None)),
value if value >= 0 => {
Ok(self.0.wait(Some(Duration::from_millis(value as u64))))
}
_ => Err(Error::Argument),
}
}
pub fn wait_one_with_time_span(&self, timeout: Duration) -> Result<bool, Error> {
Ok(self.0.wait(Some(timeout)))
}
#[must_use]
pub fn is_set(&self) -> bool {
self.0.is_set()
}
#[must_use]
pub fn was_disposed(&self) -> bool {
self.0.was_disposed()
}
}
impl IEventWait for $name {
fn reset(&self) -> Result<(), Error> {
self.reset()
}
fn set(&self) -> Result<(), Error> {
self.set()
}
fn wait_one_with_method(&self) -> Result<(), Error> {
self.wait_one_with_method()
}
fn wait_one_with_int32(&self, timeout: i32) -> Result<bool, Error> {
self.wait_one_with_int32(timeout)
}
fn wait_one_with_time_span(&self, timeout: Duration) -> Result<bool, Error> {
self.wait_one_with_time_span(timeout)
}
}
impl IAdvancedDisposable for $name {
fn was_disposed(&self) -> bool {
self.was_disposed()
}
}
};
}
managed_event!(ManagedAutoResetEvent, true);
managed_event!(ManagedManualResetEvent, false);
#[derive(Debug)]
struct SemaphoreState {
available: i32,
disposed: bool,
}
#[derive(Debug)]
pub struct ManagedSemaphore {
state: Mutex<SemaphoreState>,
changed: Condvar,
}
impl ManagedSemaphore {
pub fn new(available_count: i32) -> Result<Self, Error> {
if available_count < 1 {
return Err(Error::Argument);
}
Ok(Self {
state: Mutex::new(SemaphoreState {
available: available_count,
disposed: false,
}),
changed: Condvar::new(),
})
}
pub fn dispose(&self) -> Result<(), Error> {
let mut state = lock(&self.state);
state.disposed = true;
self.changed.notify_all();
Ok(())
}
pub fn enter_with_method(&self) -> Result<(), Error> {
self.enter_with_int32(1)
}
pub fn enter_with_int32(&self, count: i32) -> Result<(), Error> {
if count < 1 {
return Err(Error::Argument);
}
let mut state = lock(&self.state);
while !state.disposed && state.available < count {
state = self
.changed
.wait(state)
.unwrap_or_else(std::sync::PoisonError::into_inner);
}
if !state.disposed {
state.available -= count;
}
Ok(())
}
pub fn exit_with_method(&self) -> Result<(), Error> {
self.exit_with_int32(1)
}
pub fn exit_with_int32(&self, count: i32) -> Result<(), Error> {
if count < 1 {
return Err(Error::Argument);
}
let mut state = lock(&self.state);
if state.disposed {
return Ok(());
}
state.available = state
.available
.checked_add(count)
.ok_or(Error::InvalidOperation)?;
if count == 1 {
self.changed.notify_one();
} else {
self.changed.notify_all();
}
Ok(())
}
#[must_use]
pub fn was_disposed(&self) -> bool {
lock(&self.state).disposed
}
}
impl IAdvancedDisposable for ManagedSemaphore {
fn was_disposed(&self) -> bool {
self.was_disposed()
}
}
#[derive(Debug, Default)]
struct ReaderWriterState {
readers: usize,
writer: bool,
upgradeable: bool,
waiting_writers: usize,
}
#[derive(Debug, Default)]
struct ReaderWriterCore {
state: Mutex<ReaderWriterState>,
changed: Condvar,
}
impl ReaderWriterCore {
fn enter_read(&self) {
let mut state = lock(&self.state);
while state.writer || state.waiting_writers != 0 {
state = self
.changed
.wait(state)
.unwrap_or_else(std::sync::PoisonError::into_inner);
}
state.readers += 1;
}
fn exit_read(&self) -> Result<(), Error> {
let mut state = lock(&self.state);
if state.readers == 0 {
return Err(Error::InvalidOperation);
}
state.readers -= 1;
if state.readers == 0 {
self.changed.notify_all();
}
Ok(())
}
fn enter_upgradeable(&self) {
let mut state = lock(&self.state);
while state.writer || state.upgradeable || state.waiting_writers != 0 {
state = self
.changed
.wait(state)
.unwrap_or_else(std::sync::PoisonError::into_inner);
}
state.upgradeable = true;
}
fn exit_upgradeable(&self, upgraded: bool) -> Result<(), Error> {
let mut state = lock(&self.state);
if !state.upgradeable || upgraded != state.writer {
return Err(Error::InvalidOperation);
}
state.writer = false;
state.upgradeable = false;
self.changed.notify_all();
Ok(())
}
fn upgrade(&self) -> Result<(), Error> {
let mut state = lock(&self.state);
if !state.upgradeable || state.writer {
return Err(Error::InvalidOperation);
}
while state.readers != 0 {
state = self
.changed
.wait(state)
.unwrap_or_else(std::sync::PoisonError::into_inner);
}
state.writer = true;
Ok(())
}
fn enter_write(&self) {
let mut state = lock(&self.state);
state.waiting_writers += 1;
while state.writer || state.upgradeable || state.readers != 0 {
state = self
.changed
.wait(state)
.unwrap_or_else(std::sync::PoisonError::into_inner);
}
state.waiting_writers -= 1;
state.writer = true;
}
fn exit_write(&self) -> Result<(), Error> {
let mut state = lock(&self.state);
if !state.writer {
return Err(Error::InvalidOperation);
}
state.writer = false;
self.changed.notify_all();
Ok(())
}
fn exit_upgraded(&self) -> Result<(), Error> {
let mut state = lock(&self.state);
if !state.writer || !state.upgradeable {
return Err(Error::InvalidOperation);
}
state.writer = false;
state.upgradeable = false;
self.changed.notify_all();
Ok(())
}
}
pub trait IReaderWriterLockSlim: Send + Sync {
fn enter_read_lock(&self) -> Result<(), Error>;
fn exit_read_lock(&self) -> Result<(), Error>;
fn enter_upgradeable_lock(&self) -> Result<(), Error>;
fn exit_upgradeable_lock(&self, upgraded: bool) -> Result<(), Error>;
fn unchecked_exit_upgradeable_lock(&self) -> Result<(), Error>;
fn unchecked_upgrade_to_write_lock(&self) -> Result<(), Error>;
fn upgrade_to_write_lock(&self, upgraded: &mut bool) -> Result<(), Error>;
fn unchecked_exit_upgraded_lock(&self) -> Result<(), Error>;
fn enter_write_lock(&self) -> Result<(), Error>;
fn exit_write_lock(&self) -> Result<(), Error>;
}
macro_rules! reader_writer_methods {
($this:ident, $field:tt) => {
pub fn enter_read_lock(&$this) -> Result<(), Error> {
$this.$field.enter_read();
Ok(())
}
pub fn exit_read_lock(&$this) -> Result<(), Error> {
$this.$field.exit_read()
}
pub fn enter_upgradeable_lock(&$this) -> Result<(), Error> {
$this.$field.enter_upgradeable();
Ok(())
}
pub fn exit_upgradeable_lock(&$this, upgraded: bool) -> Result<(), Error> {
$this.$field.exit_upgradeable(upgraded)
}
pub fn unchecked_exit_upgradeable_lock(&$this) -> Result<(), Error> {
$this.$field.exit_upgradeable(false)
}
pub fn unchecked_upgrade_to_write_lock(&$this) -> Result<(), Error> {
$this.$field.upgrade()
}
pub fn upgrade_to_write_lock(&$this, upgraded: &mut bool) -> Result<(), Error> {
if !*upgraded {
$this.$field.upgrade()?;
*upgraded = true;
}
Ok(())
}
pub fn unchecked_exit_upgraded_lock(&$this) -> Result<(), Error> {
$this.$field.exit_upgraded()
}
pub fn enter_write_lock(&$this) -> Result<(), Error> {
$this.$field.enter_write();
Ok(())
}
pub fn exit_write_lock(&$this) -> Result<(), Error> {
$this.$field.exit_write()
}
};
}
#[derive(Debug, Clone, Default)]
pub struct OptimisticReaderWriterLock(Arc<ReaderWriterCore>);
impl OptimisticReaderWriterLock {
pub fn new() -> Result<Self, Error> {
Ok(Self::default())
}
reader_writer_methods!(self, 0);
pub fn read_lock(&self) -> Result<OptimisticReadLock, Error> {
self.enter_read_lock()?;
Ok(OptimisticReadLock {
lock: Some(self.clone()),
})
}
pub fn upgradeable_lock(&self) -> Result<OptimisticUpgradeableLock, Error> {
self.enter_upgradeable_lock()?;
Ok(OptimisticUpgradeableLock {
lock: Some(self.clone()),
upgraded: false,
})
}
pub fn write_lock(&self) -> Result<OptimisticWriteLock, Error> {
self.enter_write_lock()?;
Ok(OptimisticWriteLock {
lock: Some(self.clone()),
})
}
}
impl IReaderWriterLockSlim for OptimisticReaderWriterLock {
fn enter_read_lock(&self) -> Result<(), Error> {
self.enter_read_lock()
}
fn exit_read_lock(&self) -> Result<(), Error> {
self.exit_read_lock()
}
fn enter_upgradeable_lock(&self) -> Result<(), Error> {
self.enter_upgradeable_lock()
}
fn exit_upgradeable_lock(&self, upgraded: bool) -> Result<(), Error> {
self.exit_upgradeable_lock(upgraded)
}
fn unchecked_exit_upgradeable_lock(&self) -> Result<(), Error> {
self.unchecked_exit_upgradeable_lock()
}
fn unchecked_upgrade_to_write_lock(&self) -> Result<(), Error> {
self.unchecked_upgrade_to_write_lock()
}
fn upgrade_to_write_lock(&self, upgraded: &mut bool) -> Result<(), Error> {
self.upgrade_to_write_lock(upgraded)
}
fn unchecked_exit_upgraded_lock(&self) -> Result<(), Error> {
self.unchecked_exit_upgraded_lock()
}
fn enter_write_lock(&self) -> Result<(), Error> {
self.enter_write_lock()
}
fn exit_write_lock(&self) -> Result<(), Error> {
self.exit_write_lock()
}
}
#[derive(Debug, Default)]
pub struct SpinReaderWriterLockSlim(ReaderWriterCore);
impl SpinReaderWriterLockSlim {
reader_writer_methods!(self, 0);
}
impl IReaderWriterLockSlim for SpinReaderWriterLockSlim {
fn enter_read_lock(&self) -> Result<(), Error> {
self.enter_read_lock()
}
fn exit_read_lock(&self) -> Result<(), Error> {
self.exit_read_lock()
}
fn enter_upgradeable_lock(&self) -> Result<(), Error> {
self.enter_upgradeable_lock()
}
fn exit_upgradeable_lock(&self, upgraded: bool) -> Result<(), Error> {
self.exit_upgradeable_lock(upgraded)
}
fn unchecked_exit_upgradeable_lock(&self) -> Result<(), Error> {
self.unchecked_exit_upgradeable_lock()
}
fn unchecked_upgrade_to_write_lock(&self) -> Result<(), Error> {
self.unchecked_upgrade_to_write_lock()
}
fn upgrade_to_write_lock(&self, upgraded: &mut bool) -> Result<(), Error> {
self.upgrade_to_write_lock(upgraded)
}
fn unchecked_exit_upgraded_lock(&self) -> Result<(), Error> {
self.unchecked_exit_upgraded_lock()
}
fn enter_write_lock(&self) -> Result<(), Error> {
self.enter_write_lock()
}
fn exit_write_lock(&self) -> Result<(), Error> {
self.exit_write_lock()
}
}
#[derive(Debug)]
pub struct OptimisticReadLock {
lock: Option<OptimisticReaderWriterLock>,
}
impl OptimisticReadLock {
pub fn dispose(&mut self) -> Result<(), Error> {
self.lock
.take()
.map_or(Ok(()), |lock| lock.exit_read_lock())
}
}
impl Drop for OptimisticReadLock {
fn drop(&mut self) {
let _ = self.dispose();
}
}
#[derive(Debug)]
pub struct OptimisticWriteLock {
lock: Option<OptimisticReaderWriterLock>,
}
impl OptimisticWriteLock {
pub fn dispose(&mut self) -> Result<(), Error> {
self.lock
.take()
.map_or(Ok(()), |lock| lock.exit_write_lock())
}
}
impl Drop for OptimisticWriteLock {
fn drop(&mut self) {
let _ = self.dispose();
}
}
#[derive(Debug)]
pub struct OptimisticUpgradeableLock {
lock: Option<OptimisticReaderWriterLock>,
upgraded: bool,
}
impl OptimisticUpgradeableLock {
pub fn disposable_upgrade(&mut self) -> Result<OptimisticWriteLock, Error> {
let lock = self.lock.as_ref().ok_or(Error::InvalidOperation)?;
lock.unchecked_upgrade_to_write_lock()?;
Ok(OptimisticWriteLock {
lock: Some(lock.clone()),
})
}
pub fn upgrade(&mut self) -> Result<bool, Error> {
let lock = self.lock.as_ref().ok_or(Error::InvalidOperation)?;
if self.upgraded {
return Ok(false);
}
lock.unchecked_upgrade_to_write_lock()?;
self.upgraded = true;
Ok(true)
}
pub fn dispose(&mut self) -> Result<(), Error> {
self.lock
.take()
.map_or(Ok(()), |lock| lock.exit_upgradeable_lock(self.upgraded))
}
}
impl Drop for OptimisticUpgradeableLock {
fn drop(&mut self) {
let _ = self.dispose();
}
}
/// Owned compatibility wrapper for the C# `SpinReaderWriterLock` API.
///
/// The synchronization state is shared by every returned guard. Rust uses a
/// blocking condition variable after contention instead of burning a CPU in a
/// spin loop, while preserving the reader/writer/upgradeable exclusion rules.
#[derive(Clone, Debug, Default)]
pub struct SpinReaderWriterLock(OptimisticReaderWriterLock);
impl SpinReaderWriterLock {
pub fn new() -> Result<Self, Error> {
Ok(Self::default())
}
pub fn read_lock(&self) -> Result<SpinReadLock, Error> {
Ok(SpinReadLock(self.0.read_lock()?))
}
pub fn upgradeable_lock(&self) -> Result<SpinUpgradeableLock, Error> {
Ok(SpinUpgradeableLock {
inner: Mutex::new(Some(self.0.upgradeable_lock()?)),
})
}
pub fn write_lock(&self) -> Result<SpinWriteLock, Error> {
Ok(SpinWriteLock(Mutex::new(Some(self.0.write_lock()?))))
}
}
#[derive(Debug)]
pub struct SpinReadLock(OptimisticReadLock);
impl SpinReadLock {
pub fn dispose(&mut self) -> Result<(), Error> {
self.0.dispose()
}
}
#[derive(Debug)]
pub struct SpinWriteLock(Mutex<Option<OptimisticWriteLock>>);
impl SpinWriteLock {
pub fn dispose(&mut self) -> Result<(), Error> {
self.0
.get_mut()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.take()
.map_or(Ok(()), |mut guard| guard.dispose())
}
}
impl Drop for SpinWriteLock {
fn drop(&mut self) {
let _ = self.dispose();
}
}
#[derive(Debug)]
pub struct SpinUpgradeableLock {
inner: Mutex<Option<OptimisticUpgradeableLock>>,
}
impl SpinUpgradeableLock {
pub fn disposable_upgrade(&mut self) -> Result<SpinWriteLock, Error> {
self.disposable_upgrade_shared()
}
fn disposable_upgrade_shared(&self) -> Result<SpinWriteLock, Error> {
let mut inner = lock(&self.inner);
let guard = inner.as_mut().ok_or(Error::InvalidOperation)?;
Ok(SpinWriteLock(Mutex::new(Some(guard.disposable_upgrade()?))))
}
pub fn upgrade(&mut self) -> Result<bool, Error> {
self.upgrade_shared()
}
fn upgrade_shared(&self) -> Result<bool, Error> {
lock(&self.inner)
.as_mut()
.ok_or(Error::InvalidOperation)?
.upgrade()
}
pub fn dispose(&mut self) -> Result<(), Error> {
self.inner
.get_mut()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.take()
.map_or(Ok(()), |mut guard| guard.dispose())
}
}
impl Drop for SpinUpgradeableLock {
fn drop(&mut self) {
let _ = self.dispose();
}
}
fn close_error(error: Error) -> libremetaverse_types::compat::ExternalError {
libremetaverse_types::compat::ExternalError(error.to_string())
}
impl libremetaverse_types::compat::Close for SpinReadLock {
fn close(&mut self) -> Result<(), libremetaverse_types::compat::ExternalError> {
self.dispose().map_err(close_error)
}
}
impl libremetaverse_types::compat::Close for SpinWriteLock {
fn close(&mut self) -> Result<(), libremetaverse_types::compat::ExternalError> {
self.dispose().map_err(close_error)
}
}
impl crate::threading::disposers::IUpgradeableLock for SpinUpgradeableLock {
fn disposable_upgrade(&self) -> Result<Box<dyn libremetaverse_types::compat::Close>, Error> {
Ok(Box::new(self.disposable_upgrade_shared()?))
}
fn upgrade(&self) -> Result<bool, Error> {
self.upgrade_shared()
}
}
impl crate::threading::IReaderWriterLock for SpinReaderWriterLock {
fn read_lock(&self) -> Result<Box<dyn libremetaverse_types::compat::Close>, Error> {
Ok(Box::new(self.read_lock()?))
}
fn upgradeable_lock(
&self,
) -> Result<Box<dyn crate::threading::disposers::IUpgradeableLock>, Error> {
Ok(Box::new(self.upgradeable_lock()?))
}
fn write_lock(&self) -> Result<Box<dyn libremetaverse_types::compat::Close>, Error> {
Ok(Box::new(self.write_lock()?))
}
}
/// Converts native wait handles into executor-neutral futures.
pub struct WaitHandleAsyncFactory;
impl WaitHandleAsyncFactory {
async fn observe(
handle: libremetaverse_types::compat::WaitHandle,
timeout: Option<Duration>,
token: Option<libremetaverse_types::compat::CancellationToken>,
) -> Result<bool, Error> {
if handle.is_signalled() {
return Ok(true);
}
if timeout == Some(Duration::ZERO) {
return Ok(false);
}
if let Some(token) = token.as_ref() {
token.throw_if_cancellation_requested()?;
}
let completion = libremetaverse_types::compat::TaskCompletionSource::new();
let worker_completion = completion.clone();
let worker_token = token.clone();
std::thread::Builder::new()
.name("wait-handle-observer".to_owned())
.spawn(move || {
let deadline = timeout.and_then(|timeout| Instant::now().checked_add(timeout));
loop {
if worker_token.as_ref().is_some_and(
libremetaverse_types::compat::CancellationToken::is_cancellation_requested,
) {
let _ = worker_completion.try_set_cancelled();
break;
}
let slice = match deadline {
Some(deadline) => {
let now = Instant::now();
if now >= deadline {
let _ = worker_completion.try_set_result(false);
break;
}
deadline
.saturating_duration_since(now)
.min(Duration::from_millis(25))
}
None if worker_token.is_some() => Duration::from_millis(25),
None => {
let _ = worker_completion.try_set_result(handle.wait_timeout(None));
break;
}
};
if handle.wait_timeout(Some(slice)) {
let _ = worker_completion.try_set_result(true);
break;
}
}
})
.map_err(|_| Error::InvalidOperation)?;
let cancellation = token.map(|token| {
let completion = completion.clone();
token.register_callback(Arc::new(move || {
let _ = completion.try_set_cancelled();
}))
});
let result = completion.future().await;
drop(cancellation);
result
}
pub async fn from_wait_handle_with_wait_handle(
handle: libremetaverse_types::compat::WaitHandle,
) -> Result<(), Error> {
Self::observe(handle, None, None).await.map(|_| ())
}
pub async fn from_wait_handle_with_wait_handle_cancellation_token(
handle: libremetaverse_types::compat::WaitHandle,
token: libremetaverse_types::compat::CancellationToken,
) -> Result<(), Error> {
Self::observe(handle, None, Some(token)).await.map(|_| ())
}
pub async fn from_wait_handle_with_wait_handle_time_span(
handle: libremetaverse_types::compat::WaitHandle,
timeout: Duration,
) -> Result<bool, Error> {
Self::observe(handle, Some(timeout), None).await
}
pub async fn from_wait_handle_with_wait_handle_time_span_cancellation_token(
handle: libremetaverse_types::compat::WaitHandle,
timeout: Duration,
token: libremetaverse_types::compat::CancellationToken,
) -> Result<bool, Error> {
Self::observe(handle, Some(timeout), Some(token)).await
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicBool, Ordering};
#[test]
fn auto_reset_releases_exactly_one_waiter_per_signal() {
let event = Arc::new(ManagedAutoResetEvent::new_with_constructor().unwrap());
let first = Arc::clone(&event);
let second = Arc::clone(&event);
let first_done = Arc::new(AtomicBool::new(false));
let second_done = Arc::new(AtomicBool::new(false));
let first_flag = Arc::clone(&first_done);
let second_flag = Arc::clone(&second_done);
let a = std::thread::spawn(move || {
first.wait_one_with_method().unwrap();
first_flag.store(true, Ordering::SeqCst);
});
let b = std::thread::spawn(move || {
second.wait_one_with_method().unwrap();
second_flag.store(true, Ordering::SeqCst);
});
std::thread::sleep(Duration::from_millis(20));
event.set().unwrap();
std::thread::sleep(Duration::from_millis(20));
assert_ne!(
first_done.load(Ordering::SeqCst),
second_done.load(Ordering::SeqCst)
);
event.set().unwrap();
a.join().unwrap();
b.join().unwrap();
}
#[test]
fn upgradeable_guard_excludes_writer_and_releases_on_drop() {
let lock = OptimisticReaderWriterLock::new().unwrap();
let mut guard = lock.upgradeable_lock().unwrap();
assert!(guard.upgrade().unwrap());
assert!(!guard.upgrade().unwrap());
guard.dispose().unwrap();
let mut writer = lock.write_lock().unwrap();
writer.dispose().unwrap();
}
#[test]
fn semaphore_validates_counts_and_unblocks_after_exit() {
assert!(matches!(ManagedSemaphore::new(0), Err(Error::Argument)));
let semaphore = Arc::new(ManagedSemaphore::new(1).unwrap());
semaphore.enter_with_method().unwrap();
let acquired = Arc::new(AtomicBool::new(false));
let thread_semaphore = Arc::clone(&semaphore);
let thread_acquired = Arc::clone(&acquired);
let waiter = std::thread::spawn(move || {
thread_semaphore.enter_with_method().unwrap();
thread_acquired.store(true, Ordering::SeqCst);
});
std::thread::sleep(Duration::from_millis(20));
assert!(!acquired.load(Ordering::SeqCst));
semaphore.exit_with_method().unwrap();
waiter.join().unwrap();
}
#[test]
fn wait_handle_future_observes_signal_and_timeout() {
let handle = libremetaverse_types::compat::WaitHandle::default();
let signal = handle.clone();
std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(10));
signal.set();
});
let task = libremetaverse_types::compat::Task::from_future(async move {
WaitHandleAsyncFactory::from_wait_handle_with_wait_handle_time_span(
handle,
Duration::from_secs(1),
)
.await
});
assert_eq!(
task.wait_timeout(Duration::from_secs(2)).unwrap(),
Some(true)
);
let handle = libremetaverse_types::compat::WaitHandle::default();
let task = libremetaverse_types::compat::Task::from_future(async move {
WaitHandleAsyncFactory::from_wait_handle_with_wait_handle_time_span(
handle,
Duration::from_millis(1),
)
.await
});
assert_eq!(
task.wait_timeout(Duration::from_secs(1)).unwrap(),
Some(false)
);
}
}