ferrox_singleflight/
lib.rs1use dashmap::DashMap;
2use ferrox_errors::AppError;
3use std::sync::Arc;
4use tokio::sync::broadcast;
5use tracing::{info, debug};
6
7#[derive(Clone)]
11pub struct Singleflight<T> {
12 in_flight: Arc<DashMap<String, broadcast::Sender<Result<T, String>>>>,
14}
15
16impl<T: Clone + Send + Sync + 'static> Singleflight<T> {
17 pub fn new() -> Self {
18 Self {
19 in_flight: Arc::new(DashMap::new()),
20 }
21 }
22
23 pub async fn execute<F, Fut>(&self, key: &str, fut: F) -> Result<T, AppError>
26 where
27 F: FnOnce() -> Fut,
28 Fut: std::future::Future<Output = Result<T, AppError>> + Send + 'static,
29 {
30 let rx = {
32 if let Some(tx) = self.in_flight.get(key) {
33 Some(tx.subscribe())
35 } else {
36 let (tx, _rx) = broadcast::channel(1);
38 self.in_flight.insert(key.to_string(), tx);
39 None
40 }
41 };
42
43 if let Some(mut receiver) = rx {
45 debug!("Singleflight: Suspending execution. Waiting for in-flight result for key: {}", key);
46 return match receiver.recv().await {
47 Ok(Ok(val)) => Ok(val),
48 Ok(Err(e)) => Err(AppError::InternalError(e)),
49 Err(_) => Err(AppError::InternalError("Singleflight sender dropped".into())),
50 };
51 }
52
53 info!("Singleflight: Primary execution started for key: {}", key);
55 let result = fut().await;
56
57 if let Some((_, tx)) = self.in_flight.remove(key) {
59 let broadcast_payload = match &result {
60 Ok(val) => Ok(val.clone()),
61 Err(e) => Err(format!("{:?}", e)),
62 };
63 let _ = tx.send(broadcast_payload);
65 }
66
67 result
68 }
69}