use std::{
collections::{HashMap, VecDeque},
sync::{Arc, Mutex},
};
use scv_core::ToolError;
use tokio::sync::Notify;
use tokio_util::sync::CancellationToken;
use crate::sync::lock;
#[derive(Debug, Default)]
pub(super) struct Lanes {
state: Mutex<State>,
changed: Notify,
}
#[derive(Debug, Default)]
struct State {
next: u64,
lanes: HashMap<String, VecDeque<u64>>,
}
#[derive(Debug)]
pub(super) struct Place {
lanes: Arc<Lanes>,
handle: String,
id: u64,
pub(super) ahead: usize,
}
impl Lanes {
pub(super) fn join(self: &Arc<Self>, handle: &str) -> Place {
let mut state = lock(&self.state);
state.next += 1;
let id = state.next;
let lane = state.lanes.entry(handle.to_owned()).or_default();
let ahead = lane.len();
lane.push_back(id);
Place {
lanes: Arc::clone(self),
handle: handle.to_owned(),
id,
ahead,
}
}
}
impl Place {
pub(super) fn handle(&self) -> &str {
&self.handle
}
fn is_first(&self) -> bool {
lock(&self.lanes.state)
.lanes
.get(&self.handle)
.and_then(VecDeque::front)
== Some(&self.id)
}
pub(super) async fn wait_first(
&self,
cancellation: &CancellationToken,
) -> Result<(), ToolError> {
loop {
let changed = self.lanes.changed.notified();
if self.is_first() {
return Ok(());
}
tokio::select! {
() = changed => {}
() = cancellation.cancelled() => {
return Err(ToolError::cancelled(format!(
"cancelled while waiting for conversation {}",
self.handle
)));
}
}
}
}
}
impl Drop for Place {
fn drop(&mut self) {
let mut state = lock(&self.lanes.state);
if let Some(lane) = state.lanes.get_mut(&self.handle) {
lane.retain(|id| *id != self.id);
if lane.is_empty() {
state.lanes.remove(&self.handle);
}
}
drop(state);
self.lanes.changed.notify_waiters();
}
}