Skip to content

Commit d657d84

Browse files
committed
[MemoryPool] Add support for overdrafts in MemoryPool
This PR adds support for overdrafts in `MemoryPool` via the `force_reserve` API. This allows the memory pool to issue leases for more capacity than it holds. It also introduces a new `overdraft` API to query how much in the negative the memory pool is. Also, it introduces a new `wait_until_available` which waits until the pool is out of the overdraft mode (just as a notification without reserving anything).
1 parent 05a3446 commit d657d84

2 files changed

Lines changed: 144 additions & 1 deletion

File tree

crates/memory/Cargo.toml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,7 @@ tokio = { workspace = true, features = ["sync"] }
2525
tracing = { workspace = true }
2626

2727
[dev-dependencies]
28-
tokio = { workspace = true, features = ["rt", "macros", "time"] }
28+
tokio = { workspace = true, features = ["rt", "macros", "time", "test-util"] }
2929

3030
[lints]
3131
workspace = true

crates/memory/src/pool.rs

Lines changed: 143 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -108,6 +108,20 @@ impl MemoryPool {
108108
}
109109
}
110110

111+
/// Returns the number of bytes in overdraft (used - capacity). Returns zero if
112+
/// usage is still within capacity.
113+
#[inline]
114+
pub fn overdraft(&self) -> usize {
115+
match &self.inner {
116+
Some(inner) => {
117+
let capacity = inner.capacity.load(Ordering::Relaxed);
118+
let used = inner.used.load(Ordering::Relaxed);
119+
used.saturating_sub(capacity)
120+
}
121+
None => 0,
122+
}
123+
}
124+
111125
/// Tries to reserve `size` bytes without waiting.
112126
///
113127
/// Returns `None` if insufficient capacity.
@@ -172,6 +186,54 @@ impl MemoryPool {
172186
}
173187
}
174188

189+
/// Waits until there's any available budget. There's no guarantees
190+
/// though that by the time the caller is woken up, that the budget
191+
/// will still be available.
192+
pub async fn wait_until_available(&self) {
193+
let Some(inner) = &self.inner else {
194+
return;
195+
};
196+
loop {
197+
let notified = inner.notify.notified();
198+
if self.available() > 0 {
199+
break;
200+
}
201+
notified.await;
202+
}
203+
}
204+
205+
/// Reserves `size` bytes unconditionally, without checking capacity.
206+
///
207+
/// The pool may go into overdraft: `used` exceeds `capacity`, `available()`
208+
/// reports 0, and ordinary `try_reserve`/`reserve` callers wait until enough
209+
/// leases are returned to repay the debt. Never fails, never waits.
210+
#[inline]
211+
pub fn force_reserve(&self, size: usize) -> MemoryLease {
212+
if size == 0 {
213+
return MemoryLease {
214+
budget: self.clone(),
215+
size,
216+
};
217+
}
218+
match &self.inner {
219+
Some(inner) => {
220+
let prev = inner.used.fetch_add(size, Ordering::Relaxed);
221+
debug_assert!(
222+
prev.checked_add(size).is_some(),
223+
"MemoryPool used counter overflowed"
224+
);
225+
MemoryLease {
226+
budget: self.clone(),
227+
size,
228+
}
229+
}
230+
None => MemoryLease {
231+
budget: self.clone(),
232+
size,
233+
},
234+
}
235+
}
236+
175237
#[inline]
176238
pub fn empty_lease(&self) -> MemoryLease {
177239
MemoryLease {
@@ -480,7 +542,9 @@ const _: () = {
480542

481543
#[cfg(test)]
482544
mod tests {
545+
use std::assert_matches;
483546
use std::num::NonZeroUsize;
547+
use std::time::Duration;
484548

485549
use super::*;
486550

@@ -730,4 +794,83 @@ mod tests {
730794
assert_eq!(count.load(Ordering::Relaxed), 100);
731795
assert_eq!(budget.used(), bytes(0));
732796
}
797+
798+
#[tokio::test(start_paused = true)]
799+
async fn force_reserve() {
800+
let budget = budget(100);
801+
let r1 = budget.try_reserve(90).expect("should succeed");
802+
assert_eq!(r1.size(), bytes(90));
803+
assert_eq!(budget.used(), bytes(90));
804+
assert_eq!(budget.available(), 10);
805+
assert_eq!(budget.overdraft(), 0);
806+
807+
assert_matches!(budget.try_reserve(20), None);
808+
let r2 = budget.force_reserve(20);
809+
assert_eq!(r2.size(), bytes(20));
810+
assert_eq!(budget.used(), bytes(110));
811+
// available() should still report 0 even though we're in overdraft
812+
assert_eq!(budget.available(), 0);
813+
assert_eq!(budget.overdraft(), 10);
814+
815+
// Try to reserve anything while in overdraft will fail
816+
assert_matches!(budget.try_reserve(1), None);
817+
818+
let r3 = budget.force_reserve(50);
819+
assert_eq!(r3.size(), bytes(50));
820+
assert_eq!(budget.used(), bytes(160));
821+
assert_eq!(budget.available(), 0);
822+
assert_eq!(budget.overdraft(), 60);
823+
824+
let mut waiter1 = std::pin::pin!(budget.reserve(10));
825+
let mut waiter2 = std::pin::pin!(budget.wait_until_available());
826+
827+
// Waiters will be blocked
828+
assert!(
829+
tokio::time::timeout(Duration::from_millis(100), waiter1.as_mut())
830+
.await
831+
.is_err()
832+
);
833+
assert!(
834+
tokio::time::timeout(Duration::from_millis(100), waiter2.as_mut())
835+
.await
836+
.is_err()
837+
);
838+
839+
drop(r3);
840+
assert_eq!(budget.used(), bytes(110));
841+
assert_eq!(budget.available(), 0);
842+
assert_eq!(budget.overdraft(), 10);
843+
844+
// Returning capacity while in overdraft won't fullfill waiters
845+
assert!(
846+
tokio::time::timeout(Duration::from_millis(100), waiter1.as_mut())
847+
.await
848+
.is_err()
849+
);
850+
assert!(
851+
tokio::time::timeout(Duration::from_millis(100), waiter2.as_mut())
852+
.await
853+
.is_err()
854+
);
855+
856+
drop(r2);
857+
assert_eq!(budget.used(), bytes(90));
858+
assert_eq!(budget.available(), 10);
859+
assert_eq!(budget.overdraft(), 0);
860+
861+
// Only then will waiters be unblocked
862+
// waiter2 is waiting for available() to be > 0, so it should be
863+
// immediately unblocked.
864+
assert!(
865+
tokio::time::timeout(Duration::from_millis(100), waiter2.as_mut())
866+
.await
867+
.is_ok()
868+
);
869+
// waiter1 is trying to reserve 10 bytes, which it'll be able to acquire
870+
assert!(
871+
tokio::time::timeout(Duration::from_millis(100), waiter1.as_mut())
872+
.await
873+
.is_ok()
874+
);
875+
}
733876
}

0 commit comments

Comments
 (0)