use std::collections::VecDeque;
use std::error::Error;
use std::sync::Arc;
use std::time::{Duration, SystemTime};
use crate::events::event::{try_consume_one, AutoUnregister};
use crate::waiters::waiter::{Signaler, Waiter};
use crate::Event;
type WaitFirst = Waiter<i32>;
impl Signaler for WaitFirst {
fn fire(&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 fire event: {}", error).into());
}
};
*data = pos as i32;
self.condvar.notify_all();
Ok(())
}
fn clear(&self, _: usize) -> Result<(), Box<dyn Error>> {
Ok(())
}
}
pub fn wait_first(events: Vec<&dyn Event>, dur: Duration) -> Result<i32, Box<dyn Error>> {
let signaler = Arc::new(WaitFirst::new(-1));
let mut registered_waits: VecDeque<AutoUnregister> = VecDeque::new();
for (pos, event) in events.iter().enumerate() {
if event.handle().is_set() && try_consume_one(*event)? {
return Ok(pos as i32);
}
registered_waits.push_back(AutoUnregister::register_signaler(
Arc::<WaitFirst>::clone(&signaler),
pos,
*event,
)?);
}
let start = SystemTime::now();
let mut dur_remaining = dur;
let result = loop {
let data = signaler.lock()?;
let lock_result = signaler
.condvar
.wait_timeout_while(data, dur_remaining, |fired| *fired == -1);
let (mut fired, timeout) = match lock_result {
Ok(result) => result,
Err(error) => break Err(format!("Wait failed: {}", error).into()),
};
if timeout.timed_out() {
break Ok(-1);
}
let event = events[*fired as usize];
if event.handle().can_consume() && try_consume_one(event)? {
break Ok(*fired);
}
for (pos, event) in events.iter().enumerate() {
if *fired != pos as i32 && event.handle().can_consume() && try_consume_one(*event)? {
return Ok(pos as i32);
}
}
*fired = -1;
let result = start.elapsed();
match result {
Ok(remaining) => dur_remaining = remaining,
Err(_) => break Ok(-1),
}
};
drop(registered_waits);
result
}