Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

323 changes: 184 additions & 139 deletions crates/bifrost/src/background_appender.rs

Large diffs are not rendered by default.

37 changes: 17 additions & 20 deletions crates/bifrost/src/bifrost.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ use restate_types::live::LiveLoadExt;
use restate_core::MetadataWriter;
use restate_core::my_node_id;
use restate_core::{Metadata, ShutdownError};
use restate_memory::NonZeroByteCount;
use restate_types::config::Configuration;
use restate_types::logs::metadata::SealMetadata;
use restate_types::logs::metadata::{LogletParams, Logs, SegmentIndex};
Expand Down Expand Up @@ -206,16 +207,17 @@ impl Bifrost {
))
}

/// See [`BackgroundAppender::new`] for the semantics of `memory_limit`.
pub fn create_background_appender<T: StorageEncode>(
&self,
log_id: LogId,
error_recovery_strategy: ErrorRecoveryStrategy,
queue_capacity: usize,
memory_limit: Option<NonZeroByteCount>,
max_batch_size: usize,
) -> Result<BackgroundAppender<T>> {
Ok(BackgroundAppender::new(
self.create_appender(log_id, error_recovery_strategy)?,
queue_capacity,
memory_limit,
max_batch_size,
))
}
Expand Down Expand Up @@ -1359,16 +1361,16 @@ mod tests {
let bifrost = Bifrost::init_in_memory(env.metadata_writer).await;

let background_appender: crate::BackgroundAppender<String> = bifrost
.create_background_appender(LogId::new(0), ErrorRecoveryStrategy::Wait, 10, 10)?;
.create_background_appender(LogId::new(0), ErrorRecoveryStrategy::Wait, None, 10)?;

let mut handle = background_appender.start("test-appender")?;
let sender = handle.sender();

// A string with 100 bytes
let payload = String::from_utf8(vec![b't'; 100]).unwrap();

// try_enqueue should fail with RecordTooLarge
let result = sender.try_enqueue(payload.clone());
// enqueue should fail with RecordTooLarge
let result = sender.enqueue(payload.clone());
assert_that!(
result,
pat!(Err(pat!(EnqueueError::RecordTooLarge {
Expand All @@ -1378,7 +1380,7 @@ mod tests {
);

// enqueue (async) should also fail with RecordTooLarge
let result = sender.enqueue(payload.clone()).await;
let result = sender.enqueue(payload.clone());
assert_that!(
result,
pat!(Err(pat!(EnqueueError::RecordTooLarge {
Expand All @@ -1387,8 +1389,8 @@ mod tests {
})))
);

// try_enqueue_with_notification should also fail
let result = sender.try_enqueue_with_notification(payload.clone());
// enqueue_with_notification should also fail
let result = sender.enqueue_with_notification(payload.clone());
assert!(matches!(
result,
Err(EnqueueError::RecordTooLarge {
Expand All @@ -1415,14 +1417,14 @@ mod tests {
let bifrost = Bifrost::init_in_memory(env.metadata_writer).await;

let background_appender: crate::BackgroundAppender<String> = bifrost
.create_background_appender(LogId::new(0), ErrorRecoveryStrategy::Wait, 10, 10)?;
.create_background_appender(LogId::new(0), ErrorRecoveryStrategy::Wait, None, 10)?;

let mut handle = background_appender.start("test-appender")?;
let sender = handle.sender();

// With a 10KB limit, the ~2KB estimated record should succeed
let payload = "test".to_string();
sender.enqueue(payload).await?;
sender.enqueue(payload)?;

// Drain and wait for commit
handle.drain().await?;
Expand All @@ -1442,29 +1444,24 @@ mod tests {
let bifrost = Bifrost::init_in_memory(env.metadata_writer).await;

let background_appender: crate::BackgroundAppender<String> = bifrost
.create_background_appender(LogId::new(0), ErrorRecoveryStrategy::Wait, 1000, 100)?;
.create_background_appender(LogId::new(0), ErrorRecoveryStrategy::Wait, None, 100)?;

let mut handle = background_appender.start("test-appender")?;
let sender = handle.sender();

// Rapidly enqueue many records using try_enqueue (non-blocking)
// Rapidly enqueue many records using enqueue
let mut enqueued = 0;
for i in 0..100 {
match sender.try_enqueue(format!("rapid-record-{i}")) {
Ok(()) => enqueued += 1,
Err(EnqueueError::Full(_)) => {
// Queue is full, use async enqueue
sender.enqueue(format!("rapid-record-{i}")).await?;
enqueued += 1;
}
match sender.enqueue(format!("rapid-record-{i}")) {
Ok(_) => enqueued += 1,
Err(e) => return Err(e.into()),
}
}

assert_that!(enqueued, eq(100));

// Wait for all to be committed
let token = sender.notify_committed().await?;
let token = sender.notify_committed()?;
token.await?;

handle.drain().await?;
Expand Down
4 changes: 3 additions & 1 deletion crates/bifrost/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,9 @@ mod types;
mod watchdog;

pub use appender::Appender;
pub use background_appender::{AppenderHandle, BackgroundAppender, CommitToken, LogSender};
pub use background_appender::{
AppenderHandle, BackgroundAppender, CommitToken, EnqueueWithNotificationResult, LogSender,
};
pub use bifrost::{Bifrost, ErrorRecoveryStrategy};
pub use bifrost_admin::{BifrostAdmin, MaybeSealedSegment};
pub use data_record::{DataRecord, DataRecordError};
Expand Down
2 changes: 1 addition & 1 deletion crates/memory/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ tokio = { workspace = true, features = ["sync"] }
tracing = { workspace = true }

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

[lints]
workspace = true
181 changes: 181 additions & 0 deletions crates/memory/src/pool.rs
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,20 @@ impl MemoryPool {
}
}

/// Returns the number of bytes in overdraft (used - capacity). Returns zero if
/// usage is still within capacity.
#[inline]
pub fn overdraft(&self) -> usize {
match &self.inner {
Some(inner) => {
let capacity = inner.capacity.load(Ordering::Relaxed);
let used = inner.used.load(Ordering::Relaxed);
used.saturating_sub(capacity)
}
None => 0,
}
}

/// Tries to reserve `size` bytes without waiting.
///
/// Returns `None` if insufficient capacity.
Expand Down Expand Up @@ -172,6 +186,54 @@ impl MemoryPool {
}
}

/// Waits until there's any available budget. There's no guarantees
/// though that by the time the caller is woken up, that the budget
/// will still be available.
pub async fn wait_until_available(&self) {
let Some(inner) = &self.inner else {
return;
};
loop {
let notified = inner.notify.notified();
if self.available() > 0 {
break;
}
notified.await;
}
}

/// Reserves `size` bytes unconditionally, without checking capacity.
///
/// The pool may go into overdraft: `used` exceeds `capacity`, `available()`
/// reports 0, and ordinary `try_reserve`/`reserve` callers wait until enough
/// leases are returned to repay the debt. Never fails, never waits.
#[inline]
pub fn force_reserve(&self, size: usize) -> MemoryLease {
if size == 0 {
return MemoryLease {
budget: self.clone(),
size,
};
}
match &self.inner {
Some(inner) => {
let prev = inner.used.fetch_add(size, Ordering::Relaxed);
debug_assert!(
prev.checked_add(size).is_some(),
"MemoryPool used counter overflowed"
);
MemoryLease {
budget: self.clone(),
size,
}
}
None => MemoryLease {
budget: self.clone(),
size,
},
}
}

#[inline]
pub fn empty_lease(&self) -> MemoryLease {
MemoryLease {
Expand Down Expand Up @@ -469,6 +531,44 @@ impl PollMemoryPool {
}
}
}

/// Waits until the pool has any available budget, registering `cx` for
/// wakeup while it is exhausted.
///
/// Like [`MemoryPool::wait_until_available`], there is no guarantee that
/// the budget is still available by the time the caller acts on it.
pub fn poll_available(&mut self, cx: &mut std::task::Context<'_>) -> Poll<()> {
// Fast path: unlimited pools always have budget.
if self.pool.is_unlimited() {
return Poll::Ready(());
}

loop {
// Ensure we have a notified future *before* checking availability,
// so we don't miss a concurrent `return_memory()` notification.
let notified = self.notified.get_or_insert_with(|| {
Box::pin(
self.pool
.availability_notified_owned()
.expect("bounded pool must provide notified"),
)
});

if self.pool.available() > 0 {
self.notified = None;
return Poll::Ready(());
}

match notified.as_mut().poll(cx) {
Poll::Pending => return Poll::Pending,
Poll::Ready(()) => {
// We were notified — discard the consumed future and loop
// to re-check availability with a fresh notified.
self.notified = None;
}
}
}
}
}

const _: () = {
Expand All @@ -480,7 +580,9 @@ const _: () = {

#[cfg(test)]
mod tests {
use std::assert_matches;
use std::num::NonZeroUsize;
use std::time::Duration;

use super::*;

Expand Down Expand Up @@ -730,4 +832,83 @@ mod tests {
assert_eq!(count.load(Ordering::Relaxed), 100);
assert_eq!(budget.used(), bytes(0));
}

#[tokio::test(start_paused = true)]
async fn force_reserve() {
let budget = budget(100);
let r1 = budget.try_reserve(90).expect("should succeed");
assert_eq!(r1.size(), bytes(90));
assert_eq!(budget.used(), bytes(90));
assert_eq!(budget.available(), 10);
assert_eq!(budget.overdraft(), 0);

assert_matches!(budget.try_reserve(20), None);
let r2 = budget.force_reserve(20);
assert_eq!(r2.size(), bytes(20));
assert_eq!(budget.used(), bytes(110));
// available() should still report 0 even though we're in overdraft
assert_eq!(budget.available(), 0);
assert_eq!(budget.overdraft(), 10);

// Try to reserve anything while in overdraft will fail
assert_matches!(budget.try_reserve(1), None);

let r3 = budget.force_reserve(50);
assert_eq!(r3.size(), bytes(50));
assert_eq!(budget.used(), bytes(160));
assert_eq!(budget.available(), 0);
assert_eq!(budget.overdraft(), 60);

let mut waiter1 = std::pin::pin!(budget.reserve(10));
let mut waiter2 = std::pin::pin!(budget.wait_until_available());

// Waiters will be blocked
assert!(
tokio::time::timeout(Duration::from_millis(100), waiter1.as_mut())
.await
.is_err()
);
assert!(
tokio::time::timeout(Duration::from_millis(100), waiter2.as_mut())
.await
.is_err()
);

drop(r3);
assert_eq!(budget.used(), bytes(110));
assert_eq!(budget.available(), 0);
assert_eq!(budget.overdraft(), 10);

// Returning capacity while in overdraft won't fullfill waiters
assert!(
tokio::time::timeout(Duration::from_millis(100), waiter1.as_mut())
.await
.is_err()
);
assert!(
tokio::time::timeout(Duration::from_millis(100), waiter2.as_mut())
.await
.is_err()
);

drop(r2);
assert_eq!(budget.used(), bytes(90));
assert_eq!(budget.available(), 10);
assert_eq!(budget.overdraft(), 0);

// Only then will waiters be unblocked
// waiter2 is waiting for available() to be > 0, so it should be
// immediately unblocked.
assert!(
tokio::time::timeout(Duration::from_millis(100), waiter2.as_mut())
.await
.is_ok()
);
// waiter1 is trying to reserve 10 bytes, which it'll be able to acquire
assert!(
tokio::time::timeout(Duration::from_millis(100), waiter1.as_mut())
.await
.is_ok()
);
}
}
Loading
Loading