use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use serde_json::Value;
use tokio::sync::oneshot;
use trust_tasks_rs::TrustTask;
#[derive(Clone, Default)]
pub struct PendingReplies {
inner: Arc<Mutex<HashMap<String, oneshot::Sender<TrustTask<Value>>>>>,
}
impl PendingReplies {
pub fn new() -> Self {
Self::default()
}
pub fn register(&self, request_id: &str) -> oneshot::Receiver<TrustTask<Value>> {
let (tx, rx) = oneshot::channel();
self.lock().insert(request_id.to_string(), tx);
rx
}
pub fn abandon(&self, request_id: &str) {
self.lock().remove(request_id);
}
pub fn complete(&self, document: TrustTask<Value>) -> bool {
let Some(thread_id) = document.thread_id.clone() else {
return false;
};
let waiter = self.lock().remove(&thread_id);
match waiter {
Some(tx) => tx.send(document).is_ok(),
None => false,
}
}
fn lock(
&self,
) -> std::sync::MutexGuard<'_, HashMap<String, oneshot::Sender<TrustTask<Value>>>> {
self.inner.lock().unwrap_or_else(|p| p.into_inner())
}
}