Skip to main content

ferrox_singleflight/
lib.rs

1use dashmap::DashMap;
2use ferrox_errors::AppError;
3use std::sync::Arc;
4use tokio::sync::broadcast;
5use tracing::{info, debug};
6
7/// Singleflight prevents the "Cache Stampede" effect (Dogpile effect).
8/// If multiple requests for the same key arrive simultaneously, only the first one
9/// executes the future. The others wait and receive the result of the first one.
10#[derive(Clone)]
11pub struct Singleflight<T> {
12    // Maps a cache key to a broadcast channel that will receive the result
13    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    /// Executes the provided async closure `fut` for the given `key`, or waits for an existing
24    /// execution to finish and returns its result.
25    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        // 1. Check if the key is already in-flight
31        let rx = {
32            if let Some(tx) = self.in_flight.get(key) {
33                // Another thread is already computing this. We just subscribe to the result.
34                Some(tx.subscribe())
35            } else {
36                // We are the FIRST thread. Create a broadcast channel.
37                let (tx, _rx) = broadcast::channel(1);
38                self.in_flight.insert(key.to_string(), tx);
39                None
40            }
41        };
42
43        // 2. If we got a receiver, await the result broadcasted by the first thread
44        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        // 3. We are the first thread. We MUST execute the future.
54        info!("Singleflight: Primary execution started for key: {}", key);
55        let result = fut().await;
56
57        // 4. Broadcast the result to all waiters and remove the key
58        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            // Ignore error if there are no receivers (it means 0 stampede occurred)
64            let _ = tx.send(broadcast_payload);
65        }
66
67        result
68    }
69}