Skip to main content

ferrox_singleflight/
lib.rs

1//! # Ferrox Singleflight (`ferrox-singleflight`)
2//!
3//! `ferrox-singleflight` provides cache stampede (dogpile effect) suppression. When multiple concurrent requests attempt to compute
4//! or fetch the same missing key simultaneously, `Singleflight` ensures only **one** execution occurs while sharing the result across all callers.
5//!
6//! ## Rationale
7//! High-concurrency systems often suffer from cache stampedes when a popular cache key expires: thousands of incoming requests miss the cache
8//! simultaneously and hammer the database. `Singleflight` intercepts duplicate key lookups in-flight using Tokio broadcast channels.
9//!
10//! ## Key Features
11//! - ⚡ **Duplicate Suppression**: Only 1 worker executes the expensive task; all other waiters receive the cloned result.
12//! - 🛡️ **Memory Efficient**: In-flight computations are freed immediately upon completion.
13
14use dashmap::DashMap;
15use ferrox_errors::AppError;
16use std::sync::Arc;
17use tokio::sync::broadcast;
18use tracing::{info, debug};
19
20/// Singleflight prevents the "Cache Stampede" effect (Dogpile effect).
21/// If multiple requests for the same key arrive simultaneously, only the first one
22/// executes the future. The others wait and receive the result of the first one.
23#[derive(Clone)]
24pub struct Singleflight<T> {
25    // Maps a cache key to a broadcast channel that will receive the result
26    in_flight: Arc<DashMap<String, broadcast::Sender<Result<T, String>>>>,
27}
28
29impl<T: Clone + Send + Sync + 'static> Singleflight<T> {
30    pub fn new() -> Self {
31        Self {
32            in_flight: Arc::new(DashMap::new()),
33        }
34    }
35
36    /// Executes the provided async closure `fut` for the given `key`, or waits for an existing
37    /// execution to finish and returns its result.
38    pub async fn execute<F, Fut>(&self, key: &str, fut: F) -> Result<T, AppError>
39    where
40        F: FnOnce() -> Fut,
41        Fut: std::future::Future<Output = Result<T, AppError>> + Send + 'static,
42    {
43        // 1. Check if the key is already in-flight
44        let rx = {
45            if let Some(tx) = self.in_flight.get(key) {
46                // Another thread is already computing this. We just subscribe to the result.
47                Some(tx.subscribe())
48            } else {
49                // We are the FIRST thread. Create a broadcast channel.
50                let (tx, _rx) = broadcast::channel(1);
51                self.in_flight.insert(key.to_string(), tx);
52                None
53            }
54        };
55
56        // 2. If we got a receiver, await the result broadcasted by the first thread
57        if let Some(mut receiver) = rx {
58            debug!("Singleflight: Suspending execution. Waiting for in-flight result for key: {}", key);
59            return match receiver.recv().await {
60                Ok(Ok(val)) => Ok(val),
61                Ok(Err(e)) => Err(AppError::InternalError(e)),
62                Err(_) => Err(AppError::InternalError("Singleflight sender dropped".into())),
63            };
64        }
65
66        // 3. We are the first thread. We MUST execute the future.
67        info!("Singleflight: Primary execution started for key: {}", key);
68        let result = fut().await;
69
70        // 4. Broadcast the result to all waiters and remove the key
71        if let Some((_, tx)) = self.in_flight.remove(key) {
72            let broadcast_payload = match &result {
73                Ok(val) => Ok(val.clone()),
74                Err(e) => Err(format!("{:?}", e)),
75            };
76            // Ignore error if there are no receivers (it means 0 stampede occurred)
77            let _ = tx.send(broadcast_payload);
78        }
79
80        result
81    }
82}