win-events 0.0.2

A poor attempt to implement WaitForMultipleObjects and the ManualReset, AutoReset and Pulse types
Documentation
use std::collections::VecDeque;
use std::error::Error;
use std::sync::Arc;
use std::time::{Duration, SystemTime};

use crate::events::event::{try_consume_all, AutoUnregister};
use crate::waiters::waiter::{Signaler, Waiter};
use crate::Event;

struct Positions {
    vec: Vec<bool>,
    remaining: usize,
}

type WaitAll = Waiter<Positions>;

impl Signaler for WaitAll {
    fn fire(&self, pos: usize) -> Result<(), Box<dyn Error>> {
        let completed = {
            let lock = self.lock();
            let mut data = match lock {
                Ok(data) => data,
                Err(error) => {
                    return Err(format!("Failed to set event: {}", error).into());
                }
            };

            if data.vec[pos] == true {
                return Ok(());
            }

            data.vec[pos] = true;
            data.remaining -= 1;
            data.remaining == 0
        };

        // the last one should notify the waiting thread
        if completed {
            self.condvar.notify_all();
        }

        Ok(())
    }

    fn clear(&self, pos: usize) -> Result<(), Box<dyn Error>> {
        let lock = self.lock();
        let mut data = match lock {
            Ok(data) => data,
            Err(error) => return Err(format!("Failed to set event: {}", error).into()),
        };

        if data.vec[pos] == false {
            return Ok(());
        }

        data.vec[pos] = false;
        data.remaining += 1;
        Ok(())
    }
}

/// Waits for the all events to fire and returns true on success
///
/// # Examples
///
/// This example spawns a thread that will set all of the the events created after a short delay.
///
/// The main thread will wait on all the events to be set and return a remaining count of zero
///
///```
/// use std::thread::{sleep, spawn};
/// use std::time::Duration;
/// use win_events::{ManualResetEvent, wait_all};
///
/// let evt0 = ManualResetEvent::new(false);
/// let evt1 = ManualResetEvent::new(false);
/// let evt2 = ManualResetEvent::new(false);
/// let evt_inner0 = evt0.clone();
/// let evt_inner1 = evt1.clone();
/// let evt_inner2 = evt2.clone();
///
/// let worker = spawn(move || {
///    sleep(Duration::from_millis(10));
///    evt_inner0.set();
///    evt_inner1.set();
///    evt_inner2.set()
/// });
///
/// let wait_result = wait_all(
///    vec![&evt0, &evt1, &evt2],
///    Duration::from_millis(100),
/// ).unwrap();
///
/// assert_eq!(true, wait_result);
/// worker.join().unwrap()
///```
pub fn wait_all(events: Vec<&dyn Event>, dur: Duration) -> Result<bool, Box<dyn Error>> {
    // sort the events by the arc internal heap address, this will make the consume all
    // ordered when trying to lock the wake locks for each event. This should prevent
    // dead locks with multiple threads waiting on mutexes to be freed.
    let events = {
        let mut sorted_events = events;
        sorted_events.sort_by_key(|e| Arc::as_ptr(e.handle()));
        sorted_events
    };

    // create a waiter
    let signaler = Arc::new(WaitAll::new(Positions {
        remaining: events.len(),
        vec: vec![false; events.len()],
    }));

    // Register the waiter with the events
    let mut registered_locks = VecDeque::new();
    for (pos, event) in events.iter().enumerate() {
        registered_locks.push_back(AutoUnregister::register_signaler(
            Arc::<WaitAll>::clone(&signaler),
            pos,
            *event,
        )?);
    }

    let start = SystemTime::now();
    let mut dur_remaining = dur;
    let result = loop {
        let lock = signaler.lock()?;
        let lock_result = signaler
            .condvar
            .wait_timeout_while(lock, dur_remaining, |pos| pos.remaining > 0);
        let (mut positions, timeout) = match lock_result {
            Ok(result) => result,
            Err(error) => break Err(format!("Wait failed: {}", error).into()),
        };

        // check if there are any remaining events, if there are there was a time out
        if timeout.timed_out() {
            break Ok(false);
        }

        // Try to consume all the events, if not recalculate the remaining values
        if try_consume_all(&events)? {
            break Ok(true);
        }

        // reset the wait information
        positions.remaining = 0;
        for (pos, event) in events.iter().enumerate() {
            let set = event.handle().is_set();
            if !set {
                positions.remaining += 1
            }
            positions.vec[pos] = set
        }

        // Calculate the remaining time
        match start.elapsed() {
            Ok(elapsed) => {
                if elapsed > dur_remaining {
                    break Ok(false);
                }
                dur_remaining = dur - elapsed;
            }
            // The system time can change and cause this to return early... I'm sure it's fine
            Err(_) => break Ok(false),
        }
    };

    drop(registered_locks);
    result
}