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/// # compio::runtime::Runtime::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 = compio::runtime::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 = compio::runtime::Runtime::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 = compio::runtime::Runtime::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 = compio::runtime::Runtime::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 = compio::runtime::Runtime::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 = compio::runtime::Runtime::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}