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}