use std::{collections::HashMap, sync::Arc};
use tokio::sync::{oneshot, Mutex, OwnedSemaphorePermit, Semaphore};
use crate::DigMessage;
#[derive(Debug)]
pub(crate) struct Request {
sender: oneshot::Sender<DigMessage>,
_permit: OwnedSemaphorePermit,
}
impl Request {
pub(crate) fn send(self, message: DigMessage) {
self.sender.send(message).ok();
}
}
#[derive(Debug)]
pub(crate) struct RequestMap {
items: Mutex<HashMap<u16, Request>>,
capacity: Arc<Semaphore>,
next_id: Mutex<u16>,
}
impl RequestMap {
pub(crate) fn new() -> Self {
Self {
items: Mutex::new(HashMap::new()),
capacity: Arc::new(Semaphore::new(u16::MAX as usize)),
next_id: Mutex::new(0),
}
}
pub(crate) async fn insert(&self, sender: oneshot::Sender<DigMessage>) -> u16 {
let permit = self
.capacity
.clone()
.acquire_owned()
.await
.expect("request capacity semaphore is never closed");
let mut items = self.items.lock().await;
items.retain(|_, request| !request.sender.is_closed());
let mut next_id = self.next_id.lock().await;
let mut id = *next_id;
loop {
if !items.contains_key(&id) {
break;
}
id = id.wrapping_add(1);
}
*next_id = id.wrapping_add(1);
drop(next_id);
items.insert(
id,
Request {
sender,
_permit: permit,
},
);
id
}
pub(crate) async fn remove(&self, id: u16) -> Option<Request> {
self.items.lock().await.remove(&id)
}
}