Skip to main content

monocoque_core/
backpressure.rs

1//! Backpressure: `BytePermits`
2//!
3//! Byte-based flow control for write pumps.
4//!
5//! Design principle:
6//! - Backpressure scales with **bytes**, not message count
7//! - One giant message should not starve other connections
8//! - Pluggable: `NoOp` (default) → Semaphore → dynamic policy
9//!
10//! Usage:
11//! ```rust,ignore
12//! let permits = SemaphorePermits::new(10 * 1024 * 1024); // 10MB limit
13//! let permit = permits.acquire(n_bytes).await;
14//! writer.write(buf).await;
15//! drop(permit); // releases automatically
16//! ```
17
18use async_trait::async_trait;
19use parking_lot::{Condvar, Mutex};
20use std::sync::Arc;
21
22/// Backpressure permit trait.
23///
24/// Implementations control write pump flow based on byte counts.
25#[async_trait]
26pub trait BytePermits: Send + Sync {
27    /// Acquire permission to write `n_bytes`.
28    ///
29    /// This may block if the system is under memory pressure.
30    async fn acquire(&self, n_bytes: usize) -> Permit;
31}
32
33/// Internal state for the byte semaphore.
34struct SemInner {
35    available: usize,
36    /// Total capacity; used to clamp oversized acquires so we never deadlock.
37    max_bytes: usize,
38}
39
40/// RAII permit guard.
41///
42/// Releases the permit when dropped.
43pub struct Permit {
44    inner: Option<PermitInner>,
45}
46
47enum PermitInner {
48    /// Byte-counting semaphore backed by `parking_lot` primitives (usable in `Drop`).
49    ByteSem(Arc<(Mutex<SemInner>, Condvar)>, usize),
50    NoOp,
51}
52
53impl Drop for Permit {
54    fn drop(&mut self) {
55        match self.inner.take() {
56            Some(PermitInner::ByteSem(inner, n_bytes)) => {
57                let (mutex, condvar) = &*inner;
58                let mut guard = mutex.lock();
59                guard.available += n_bytes;
60                drop(guard);
61                condvar.notify_all();
62            }
63            Some(PermitInner::NoOp) | None => {}
64        }
65    }
66}
67
68impl Permit {
69    pub(crate) const fn noop() -> Self {
70        Self {
71            inner: Some(PermitInner::NoOp),
72        }
73    }
74
75    fn byte_sem(inner: Arc<(Mutex<SemInner>, Condvar)>, n_bytes: usize) -> Self {
76        Self {
77            inner: Some(PermitInner::ByteSem(inner, n_bytes)),
78        }
79    }
80}
81
82/// No-op implementation (Phase 0).
83///
84/// Always grants permits immediately.
85/// Use this until memory pressure becomes an issue.
86#[derive(Debug, Clone, Copy, Default)]
87pub struct NoOpPermits;
88
89#[async_trait]
90impl BytePermits for NoOpPermits {
91    async fn acquire(&self, _n_bytes: usize) -> Permit {
92        Permit::noop()
93    }
94}
95
96/// Semaphore-based backpressure implementation.
97///
98/// Enforces a maximum number of bytes that can be buffered at once.
99/// When the limit is reached, `acquire()` will block until space is available.
100/// Acquires all N bytes in a single atomic operation (O(1), not O(N)).
101///
102/// # Example
103///
104/// ```
105/// use monocoque_core::backpressure::{BytePermits, SemaphorePermits};
106///
107/// # monocoque_core::rt::LocalRuntime::new().unwrap().block_on(async {
108/// // Allow up to 10MB of buffered data
109/// let permits = SemaphorePermits::new(10 * 1024 * 1024);
110///
111/// // Acquire permit for 1KB write
112/// let permit = permits.acquire(1024).await;
113/// // ... perform write ...
114/// drop(permit); // releases 1024 bytes back to the pool
115/// # });
116/// ```
117#[derive(Clone)]
118pub struct SemaphorePermits {
119    inner: Arc<(Mutex<SemInner>, Condvar)>,
120}
121
122impl SemaphorePermits {
123    /// Create a new semaphore-based backpressure controller.
124    ///
125    /// # Arguments
126    ///
127    /// * `max_bytes` - Maximum number of bytes that can be buffered
128    #[must_use]
129    pub fn new(max_bytes: usize) -> Self {
130        Self {
131            inner: Arc::new((
132                Mutex::new(SemInner {
133                    available: max_bytes,
134                    max_bytes,
135                }),
136                Condvar::new(),
137            )),
138        }
139    }
140}
141
142#[async_trait]
143impl BytePermits for SemaphorePermits {
144    async fn acquire(&self, n_bytes: usize) -> Permit {
145        if n_bytes == 0 {
146            return Permit::noop();
147        }
148
149        // Blocking wait is performed on a dedicated thread so we don't block
150        // the async executor. parking_lot::Condvar::wait is synchronous and
151        // safe to use here because SemInner uses parking_lot::Mutex.
152        let inner = self.inner.clone();
153        let actual = crate::rt::spawn_blocking(move || {
154            let (mutex, condvar) = &*inner;
155            let mut guard = mutex.lock();
156            // Clamp to max_bytes so a single oversized message never deadlocks:
157            // the message will consume the entire capacity instead of waiting
158            // forever for capacity that can never exist.
159            let claim = n_bytes.min(guard.max_bytes);
160            // Wait until enough capacity is available.
161            while guard.available < claim {
162                condvar.wait(&mut guard);
163            }
164            guard.available -= claim;
165            // Return the actual bytes claimed so the Permit releases the right amount.
166            claim
167        })
168        .await;
169
170        Permit::byte_sem(self.inner.clone(), actual)
171    }
172}
173
174#[cfg(test)]
175mod tests {
176    use super::*;
177
178    #[test]
179    fn noop_permits_always_succeed() {
180        let permits = NoOpPermits;
181        let rt = crate::rt::LocalRuntime::new().unwrap();
182        rt.block_on(async {
183            let _p1 = permits.acquire(1024).await;
184            let _p2 = permits.acquire(1_000_000).await;
185            // Should not block
186        });
187    }
188
189    #[test]
190    fn semaphore_permits_enforce_limit() {
191        let permits = SemaphorePermits::new(1024);
192        let rt = crate::rt::LocalRuntime::new().unwrap();
193
194        rt.block_on(async {
195            // First 1024 bytes should succeed
196            let p1 = permits.acquire(1024).await;
197
198            // Try to acquire more - this would block, so we test the behavior
199            // by checking we can acquire after dropping
200            drop(p1);
201
202            let _p2 = permits.acquire(512).await;
203            let _p3 = permits.acquire(512).await;
204            // Should succeed with 1024 total
205        });
206    }
207
208    #[test]
209    fn semaphore_permits_release_on_drop() {
210        let permits = SemaphorePermits::new(1000);
211        let rt = crate::rt::LocalRuntime::new().unwrap();
212
213        rt.block_on(async {
214            {
215                let _p1 = permits.acquire(500).await;
216                let _p2 = permits.acquire(500).await;
217                // Full capacity used
218            } // Permits dropped here
219
220            // Should be able to acquire again after drop
221            let _p3 = permits.acquire(1000).await;
222        });
223    }
224
225    #[test]
226    fn semaphore_permits_oversized_acquire_does_not_deadlock() {
227        // A single acquire larger than max_bytes must complete (clamped to max_bytes)
228        // rather than deadlocking forever waiting for capacity that can never exist.
229        let permits = SemaphorePermits::new(1024);
230        let rt = crate::rt::LocalRuntime::new().unwrap();
231
232        rt.block_on(async {
233            let permit = permits.acquire(2048).await; // 2× max - must not deadlock
234            drop(permit);
235            // After release, we can acquire up to max_bytes again.
236            let _p = permits.acquire(1024).await;
237        });
238    }
239
240    #[test]
241    fn semaphore_permits_single_atomic_acquire() {
242        // Verify that acquiring N bytes is done atomically (not O(N) individual acquires)
243        let permits = SemaphorePermits::new(1024 * 1024); // 1MB
244        let rt = crate::rt::LocalRuntime::new().unwrap();
245
246        rt.block_on(async {
247            // Acquire a large block in one shot - this should not loop N times
248            let permit = permits.acquire(512 * 1024).await; // 512KB
249            drop(permit);
250        });
251    }
252}