use std::alloc::Layout;
use std::ops::{Deref, DerefMut};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
use std::time::{Duration, Instant};
use bytes::Bytes;
use offset_allocator::{Allocation, Allocator};
use parking_lot::{Mutex, RwLock};
use super::backend::{RdmaBackend, RdmaError};
pub(crate) const GRANULE: usize = 4096;
#[derive(Debug, Clone)]
pub struct RdmaPoolConfig {
pub initial_arena_bytes: u64,
pub max_arena_bytes: u64,
pub dedicated_arena_min: u64,
pub registered_bytes_budget: u64,
pub arena_reclaim_after: Option<Duration>,
pub retain_arena_bytes: u64,
}
impl Default for RdmaPoolConfig {
fn default() -> Self {
Self {
initial_arena_bytes: 64 << 20,
max_arena_bytes: 1 << 30,
dedicated_arena_min: 64 << 20,
registered_bytes_budget: 1 << 30,
arena_reclaim_after: None,
retain_arena_bytes: 64 << 20,
}
}
}
fn coarse_now_ms() -> u64 {
static BASE: std::sync::OnceLock<Instant> = std::sync::OnceLock::new();
BASE.get_or_init(Instant::now).elapsed().as_millis() as u64
}
const NOT_EMPTY: u64 = u64::MAX;
#[derive(Debug, Clone)]
pub(crate) struct RemoteRef {
pub addr: u64,
pub len: u64,
pub packed_key: Bytes,
pub generation: u64,
}
struct PageMemory {
ptr: *mut u8,
len: usize,
freeable: AtomicBool,
}
unsafe impl Send for PageMemory {}
unsafe impl Sync for PageMemory {}
impl PageMemory {
fn new(len: usize) -> Result<Self, RdmaError> {
let len = len
.checked_next_multiple_of(GRANULE)
.filter(|n| *n > 0)
.ok_or(RdmaError::OutOfRange)?;
let layout = Layout::from_size_align(len, GRANULE).map_err(|_| RdmaError::OutOfRange)?;
let ptr = unsafe { std::alloc::alloc_zeroed(layout) };
if ptr.is_null() {
return Err(RdmaError::Backend(format!(
"allocating a {len} B rdma arena failed"
)));
}
Ok(Self {
ptr,
len,
freeable: AtomicBool::new(false),
})
}
fn addr(&self) -> usize {
self.ptr as usize
}
fn mark_unmapped(&self) {
self.freeable.store(true, Ordering::Release);
}
fn mark_pinned(&self) {
self.freeable.store(false, Ordering::Release);
}
}
impl Drop for PageMemory {
fn drop(&mut self) {
if !self.freeable.load(Ordering::Acquire) {
tracing::error!(
bytes = self.len,
"rdma arena dropped without a confirmed unmap; leaking its pages rather than \
freeing memory the backend may still have pinned. The registry was torn down \
without `RdmaRegistry::shutdown`."
);
return;
}
let layout = Layout::from_size_align(self.len, GRANULE).expect("layout was valid to alloc");
unsafe { std::alloc::dealloc(self.ptr, layout) };
}
}
pub(crate) struct Arena {
memory: PageMemory,
backend_region_id: u64,
packed_key: Bytes,
generation: u64,
charged: u64,
granules: u32,
free: Mutex<Allocator<u32>>,
live: AtomicUsize,
empty_since_ms: AtomicU64,
in_flight: AtomicUsize,
dedicated: bool,
}
impl Arena {
fn base(&self) -> *mut u8 {
self.memory.ptr
}
pub(crate) fn len(&self) -> usize {
self.memory.len
}
pub(crate) fn live(&self) -> usize {
self.live.load(Ordering::Acquire)
}
pub(crate) fn in_flight(&self) -> usize {
self.in_flight.load(Ordering::Acquire)
}
fn charged(&self) -> u64 {
self.charged
}
fn empty_for_ms(&self, now_ms: u64) -> Option<u64> {
match self.empty_since_ms.load(Ordering::Acquire) {
NOT_EMPTY => None,
since => Some(now_ms.saturating_sub(since)),
}
}
fn try_alloc(self: &Arc<Self>, len: usize) -> Option<PinnedBuf> {
let granules = u32::try_from(len.div_ceil(GRANULE)).ok()?;
if granules == 0 || granules > self.granules {
return None;
}
let allocation = self.free.lock().allocate(granules)?;
self.live.fetch_add(1, Ordering::AcqRel);
self.empty_since_ms.store(NOT_EMPTY, Ordering::Release);
Some(PinnedBuf {
inner: Arc::new(Suballoc {
offset: allocation.offset as usize * GRANULE,
len,
allocation,
arena: Arc::clone(self),
}),
})
}
fn release(&self, allocation: Allocation<u32>) {
self.free.lock().free(allocation);
if self.live.fetch_sub(1, Ordering::AcqRel) == 1 {
self.empty_since_ms
.store(coarse_now_ms(), Ordering::Release);
}
}
}
pub(crate) struct Suballoc {
arena: Arc<Arena>,
allocation: Allocation<u32>,
offset: usize,
len: usize,
}
impl Drop for Suballoc {
fn drop(&mut self) {
self.arena.release(self.allocation);
}
}
pub(crate) struct TransferHold(Arc<Suballoc>);
impl Drop for TransferHold {
fn drop(&mut self) {
self.0.arena.in_flight.fetch_sub(1, Ordering::AcqRel);
}
}
pub struct PinnedBuf {
inner: Arc<Suballoc>,
}
impl PinnedBuf {
pub(crate) fn remote(&self) -> RemoteRef {
RemoteRef {
addr: self.addr(),
len: self.inner.len as u64,
packed_key: self.inner.arena.packed_key.clone(),
generation: self.inner.arena.generation,
}
}
pub(crate) fn addr(&self) -> u64 {
self.inner.arena.base() as u64 + self.inner.offset as u64
}
pub(crate) fn arena_offset(&self) -> u64 {
self.inner.offset as u64
}
pub(crate) fn backend_region_id(&self) -> u64 {
self.inner.arena.backend_region_id
}
pub(crate) fn hold(&self) -> TransferHold {
self.inner.arena.in_flight.fetch_add(1, Ordering::AcqRel);
TransferHold(Arc::clone(&self.inner))
}
pub fn len(&self) -> usize {
self.inner.len
}
pub fn is_empty(&self) -> bool {
self.inner.len == 0
}
}
impl std::fmt::Debug for PinnedBuf {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PinnedBuf")
.field("addr", &self.addr())
.field("len", &self.inner.len)
.field("generation", &self.inner.arena.generation)
.finish()
}
}
impl Deref for PinnedBuf {
type Target = [u8];
fn deref(&self) -> &[u8] {
unsafe {
std::slice::from_raw_parts(
self.inner.arena.base().add(self.inner.offset),
self.inner.len,
)
}
}
}
impl DerefMut for PinnedBuf {
fn deref_mut(&mut self) -> &mut [u8] {
unsafe {
std::slice::from_raw_parts_mut(
self.inner.arena.base().add(self.inner.offset),
self.inner.len,
)
}
}
}
pub(crate) struct Budget {
registered: AtomicU64,
limit: u64,
metrics: Option<Arc<crate::observability::VeloMetrics>>,
}
impl Budget {
pub(crate) fn new(limit: u64, metrics: Option<Arc<crate::observability::VeloMetrics>>) -> Self {
Self {
registered: AtomicU64::new(0),
limit,
metrics,
}
}
pub(crate) fn try_reserve(self: &Arc<Self>, bytes: u64) -> Result<Reservation, RdmaError> {
let mut current = self.registered.load(Ordering::Acquire);
loop {
let exceeded = || RdmaError::BudgetExceeded {
requested: bytes,
registered: current,
budget: self.limit,
};
let next = current.checked_add(bytes).ok_or_else(exceeded)?;
if next > self.limit {
return Err(exceeded());
}
match self.registered.compare_exchange_weak(
current,
next,
Ordering::AcqRel,
Ordering::Acquire,
) {
Ok(_) => {
self.publish(next);
return Ok(Reservation {
budget: Arc::clone(self),
bytes,
committed: false,
});
}
Err(observed) => current = observed,
}
}
}
pub(crate) fn release(&self, bytes: u64) {
let next = self
.registered
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |current| {
Some(current.saturating_sub(bytes))
})
.unwrap_or(0);
self.publish(next.saturating_sub(bytes));
}
pub(crate) fn charge(&self, bytes: u64) -> bool {
let next = self
.registered
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |current| {
Some(current.saturating_add(bytes))
})
.unwrap_or(0)
.saturating_add(bytes);
self.publish(next);
next <= self.limit
}
pub(crate) fn registered(&self) -> u64 {
self.registered.load(Ordering::Acquire)
}
fn publish(&self, bytes: u64) {
if let Some(m) = &self.metrics {
m.set_rdma_registered_bytes(bytes);
}
}
}
#[must_use = "an uncommitted Reservation releases its claim when dropped"]
pub(crate) struct Reservation {
budget: Arc<Budget>,
bytes: u64,
committed: bool,
}
impl Reservation {
pub(crate) fn commit(mut self) {
self.committed = true;
}
pub(crate) fn bytes(&self) -> u64 {
self.bytes
}
pub(crate) fn raise_to(&mut self, bytes: u64) {
debug_assert!(
bytes >= self.bytes,
"raise_to must never shrink a claim: {} to {bytes}",
self.bytes
);
let Some(extra) = bytes.checked_sub(self.bytes).filter(|e| *e != 0) else {
return;
};
if !self.budget.charge(extra) {
tracing::warn!(
extra,
registered = self.budget.registered(),
budget = self.budget.limit,
"rdma: the backend pinned more than the registered-bytes budget allows; \
accepting the overshoot rather than unmapping a live registration"
);
}
self.bytes = bytes;
}
}
pub(crate) fn page_enclosing_len(ptr: usize, len: usize) -> Option<u64> {
let start = ptr & !(GRANULE - 1);
let end = ptr.checked_add(len)?.checked_next_multiple_of(GRANULE)?;
u64::try_from(end.checked_sub(start)?).ok()
}
impl Drop for Reservation {
fn drop(&mut self) {
if !self.committed {
self.budget.release(self.bytes);
}
}
}
pub(crate) fn pool_arena_target(initial: u64, max: u64, pooled: usize) -> u64 {
u32::try_from(pooled)
.ok()
.and_then(|shift| 1u64.checked_shl(shift))
.and_then(|factor| initial.checked_mul(factor))
.unwrap_or(max)
.min(max)
}
pub(crate) struct ArenaSet {
backend: Arc<dyn RdmaBackend>,
cfg: RdmaPoolConfig,
budget: Arc<Budget>,
generations: Arc<AtomicU64>,
metrics: Option<Arc<crate::observability::VeloMetrics>>,
arenas: RwLock<Vec<Arc<Arena>>>,
grow: tokio::sync::Mutex<()>,
unconfirmed: Mutex<Vec<Arc<Arena>>>,
}
impl ArenaSet {
pub(crate) fn new(
backend: Arc<dyn RdmaBackend>,
cfg: RdmaPoolConfig,
budget: Arc<Budget>,
generations: Arc<AtomicU64>,
metrics: Option<Arc<crate::observability::VeloMetrics>>,
) -> Self {
Self {
backend,
cfg,
budget,
generations,
metrics,
arenas: RwLock::new(Vec::new()),
grow: tokio::sync::Mutex::new(()),
unconfirmed: Mutex::new(Vec::new()),
}
}
pub(crate) async fn alloc(&self, len: usize) -> Result<PinnedBuf, RdmaError> {
if len == 0 {
return Err(RdmaError::OutOfRange);
}
let dedicated = len as u64 >= self.cfg.dedicated_arena_min;
if !dedicated && let Some(buf) = self.try_existing(len) {
return Ok(buf);
}
let _grow = self.grow.lock().await;
if !dedicated && let Some(buf) = self.try_existing(len) {
return Ok(buf);
}
let arena_bytes = if dedicated {
(len as u64)
.checked_next_multiple_of(GRANULE as u64)
.ok_or(RdmaError::OutOfRange)?
} else {
self.next_pool_arena_bytes(len)?
};
let arena = self.map_arena(arena_bytes, dedicated).await?;
let buf = arena
.try_alloc(len)
.ok_or_else(|| RdmaError::Backend("a fresh arena refused its own request".into()))?;
self.arenas.write().push(arena);
Ok(buf)
}
pub(crate) fn try_alloc_existing(&self, len: usize) -> Option<PinnedBuf> {
if len == 0 || len as u64 >= self.cfg.dedicated_arena_min {
return None;
}
self.try_existing(len)
}
fn try_existing(&self, len: usize) -> Option<PinnedBuf> {
let arenas = self.arenas.read();
arenas
.iter()
.filter(|a| !a.dedicated)
.find_map(|arena| arena.try_alloc(len))
}
pub(crate) async fn reclaim_idle(&self) -> usize {
let now_ms = coarse_now_ms();
let victims: Vec<Arc<Arena>> = {
let mut arenas = self.arenas.write();
let mut pooled_bytes: u64 = arenas
.iter()
.filter(|a| !a.dedicated)
.map(|a| a.charged())
.sum();
let mut victims = Vec::new();
let mut keep = Vec::with_capacity(arenas.len());
for arena in std::mem::take(&mut *arenas).into_iter().rev() {
if self.is_reclaimable(&arena, now_ms, pooled_bytes) {
if !arena.dedicated {
pooled_bytes -= arena.charged();
}
victims.push(arena);
} else {
keep.push(arena);
}
}
keep.reverse();
*arenas = keep;
victims
};
if victims.is_empty() {
return 0;
}
let mut reclaimed = 0usize;
let mut unconfirmed = Vec::new();
for arena in victims {
match self.backend.unmap(arena.backend_region_id).await {
Ok(()) => {
arena.memory.mark_unmapped();
self.budget.release(arena.charged());
reclaimed += 1;
tracing::debug!(
bytes = arena.len(),
dedicated = arena.dedicated,
generation = arena.generation,
"rdma: reclaimed an empty arena"
);
}
Err(e) => {
tracing::warn!(
%e,
bytes = arena.len(),
"rdma: reclaim unmap unconfirmed; the arena's pages and its budget stay \
held until velo shutdown completes"
);
unconfirmed.push(arena);
}
}
}
if !unconfirmed.is_empty() {
self.unconfirmed.lock().extend(unconfirmed);
}
reclaimed
}
fn is_reclaimable(&self, arena: &Arena, now_ms: u64, pooled_bytes: u64) -> bool {
if arena.live() != 0 || arena.in_flight() != 0 {
return false;
}
if arena.dedicated {
return true;
}
let Some(after) = self.cfg.arena_reclaim_after else {
return false;
};
let Some(empty_for) = arena.empty_for_ms(now_ms) else {
return false;
};
u128::from(empty_for) >= after.as_millis()
&& pooled_bytes.saturating_sub(arena.charged()) >= self.cfg.retain_arena_bytes
}
fn next_pool_arena_bytes(&self, len: usize) -> Result<u64, RdmaError> {
let pooled = self.arenas.read().iter().filter(|a| !a.dedicated).count();
let grown = pool_arena_target(
self.cfg.initial_arena_bytes,
self.cfg.max_arena_bytes,
pooled,
);
let needed = (len as u64)
.checked_next_multiple_of(GRANULE as u64)
.ok_or(RdmaError::OutOfRange)?;
Ok(grown.max(needed).max(GRANULE as u64))
}
async fn map_arena(&self, bytes: u64, dedicated: bool) -> Result<Arc<Arena>, RdmaError> {
let len = usize::try_from(bytes)
.ok()
.and_then(|len| len.checked_next_multiple_of(GRANULE))
.ok_or(RdmaError::OutOfRange)?;
let granules = u32::try_from(len / GRANULE).map_err(|_| RdmaError::OutOfRange)?;
let mut reservation = self.budget.try_reserve(len as u64)?;
let memory = PageMemory::new(len)?;
debug_assert_eq!(
memory.len as u64,
reservation.bytes(),
"the budget claim and the mapped length must be the same number"
);
memory.mark_unmapped();
let region = self.backend.map(memory.addr(), memory.len).await?;
memory.mark_pinned();
reservation.raise_to(region.effective_len.max(memory.len as u64));
debug_assert_eq!(
memory.len / GRANULE,
granules as usize,
"the granule count was validated against a different length than was mapped"
);
let arena = Arc::new(Arena {
backend_region_id: region.backend_region_id,
packed_key: region.packed_key,
generation: self.generations.fetch_add(1, Ordering::Relaxed),
charged: reservation.bytes(),
granules,
free: Mutex::new(Allocator::with_max_allocs(granules, granules)),
live: AtomicUsize::new(0),
empty_since_ms: AtomicU64::new(NOT_EMPTY),
in_flight: AtomicUsize::new(0),
dedicated,
memory,
});
if let Some(m) = &self.metrics {
m.record_rdma_registration(crate::observability::RdmaRegistrationKind::Arena);
}
tracing::debug!(
bytes = arena.len(),
dedicated,
generation = arena.generation,
"rdma: mapped a new arena"
);
reservation.commit();
Ok(arena)
}
pub(crate) async fn unmap_all(&self, deadline: Instant) -> usize {
let arenas: Vec<Arc<Arena>> = std::mem::take(&mut *self.arenas.write());
let mut unconfirmed = Vec::new();
for arena in arenas {
let budget_spent = Instant::now() >= deadline;
if budget_spent {
if arena.in_flight() != 0 {
tracing::warn!(
in_flight = arena.in_flight(),
bytes = arena.len(),
"rdma: the shutdown budget was already spent before this arena, so its \
transfers were not waited for at all; force-unmapping. Transport \
teardown is the backstop and a straggling transfer fails at its own end."
);
}
} else {
while arena.in_flight() != 0 && Instant::now() < deadline {
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
}
if arena.in_flight() != 0 {
tracing::warn!(
in_flight = arena.in_flight(),
bytes = arena.len(),
"rdma: transfers into this arena outlasted the shutdown budget; \
unmapping anyway, and they will fail at their own end"
);
}
}
let live = arena.live();
if live != 0 {
tracing::debug!(
live,
bytes = arena.len(),
"rdma: unmapping an arena a caller still holds buffers from; their memory \
stays valid, it is simply no longer registered"
);
}
let remaining = deadline.saturating_duration_since(Instant::now());
let unmap = self.backend.unmap(arena.backend_region_id);
let outcome = match tokio::time::timeout(remaining, unmap).await {
Ok(result) => result,
Err(_) => Err(RdmaError::Timeout),
};
match outcome {
Ok(()) => {
arena.memory.mark_unmapped();
self.budget.release(arena.charged);
}
Err(e) => {
let bytes = arena.len();
tracing::warn!(
%e,
bytes,
"rdma: arena unmap unconfirmed; its pages stay leaked until velo \
shutdown completes"
);
unconfirmed.push(arena);
}
}
}
let count = unconfirmed.len();
self.unconfirmed.lock().extend(unconfirmed);
count
}
pub(crate) fn release_unconfirmed(&self) {
for arena in self.unconfirmed.lock().drain(..) {
arena.memory.mark_unmapped();
self.budget.release(arena.charged);
}
}
#[cfg(test)]
pub(crate) fn arena_count(&self) -> usize {
self.arenas.read().len()
}
#[cfg(test)]
pub(crate) fn live_allocations(&self) -> usize {
self.arenas.read().iter().map(|a| a.live()).sum()
}
#[cfg(feature = "test-helpers")]
pub(crate) fn in_flight_transfers(&self) -> usize {
self.arenas.read().iter().map(|a| a.in_flight()).sum()
}
#[cfg(test)]
pub(crate) fn registered_bytes(&self) -> u64 {
self.budget.registered()
}
}