[SKIP CI] add forget semaphore, add them back not yet implemented

This commit is contained in:
Quentin Dufour
2025-05-01 10:05:04 +02:00
parent f34558af07
commit fa457328c8
8 changed files with 96 additions and 170 deletions
+4 -9
View File
@@ -14,7 +14,6 @@ use futures::stream::FuturesUnordered;
use futures::StreamExt;
use serde::{Deserialize, Serialize};
use tokio::select;
use tokio::time::sleep;
use garage_db as db;
@@ -241,7 +240,7 @@ impl K2VRpcHandler {
let timeout_duration = Duration::from_millis(timeout_msec);
let resps = select! {
r = rpc => r?,
_ = sleep(timeout_duration) => return Ok(None),
_ = tokio::time::sleep(timeout_duration) => return Ok(None),
};
let mut resp: Option<K2VItem> = None;
@@ -378,9 +377,8 @@ impl K2VRpcHandler {
};
// Propagate to rest of network
if let Some((updated, backpressure)) = new {
if let Some(updated) = new {
self.item_table.insert(&updated).await?;
sleep(backpressure).await;
}
Ok(K2VRpc::Ok)
@@ -388,16 +386,14 @@ impl K2VRpcHandler {
async fn handle_insert_many(&self, items: &[InsertedItem]) -> Result<K2VRpc, Error> {
let mut updated_vec = vec![];
let mut backpressure = Duration::ZERO;
{
let local_timestamp_tree = self.local_timestamp_tree.lock().unwrap();
for item in items {
let new = self.local_insert(&local_timestamp_tree, item)?;
if let Some((updated, add_bp)) = new {
if let Some(updated) = new {
updated_vec.push(updated);
backpressure += add_bp;
}
}
}
@@ -406,7 +402,6 @@ impl K2VRpcHandler {
if !updated_vec.is_empty() {
self.item_table.insert_many(&updated_vec).await?;
}
sleep(backpressure).await;
Ok(K2VRpc::Ok)
}
@@ -415,7 +410,7 @@ impl K2VRpcHandler {
&self,
local_timestamp_tree: &MutexGuard<'_, db::Tree>,
item: &InsertedItem,
) -> Result<Option<(K2VItem, Duration)>, Error> {
) -> Result<Option<K2VItem>, Error> {
let now = now_msec();
self.item_table
+24 -38
View File
@@ -1,10 +1,9 @@
use core::borrow::Borrow;
use std::convert::TryInto;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use std::sync::Arc;
use serde_bytes::ByteBuf;
use tokio::sync::Notify;
use tokio::sync::{Notify, Semaphore};
use garage_db as db;
@@ -33,7 +32,7 @@ pub struct TableData<F: TableSchema, R: TableReplication> {
pub(crate) merkle_tree: db::Tree,
pub(crate) merkle_todo: db::Tree,
pub(crate) merkle_todo_notify: Notify,
pub(crate) merkle_todo_sleep: Arc<Mutex<Duration>>,
pub(crate) merkle_todo_bounded_queue: Option<Arc<Semaphore>>,
pub(crate) insert_queue: db::Tree,
pub(crate) insert_queue_notify: Arc<Notify>,
@@ -64,11 +63,12 @@ impl<F: TableSchema, R: TableReplication> TableData<F, R> {
.open_tree(format!("{}:merkle_todo", F::TABLE_NAME))
.expect("Unable to open DB Merkle TODO tree");
let initial = match config {
MerkleBackpressureEnum::None => Duration::ZERO,
MerkleBackpressureEnum::Aimd(aimd) => Duration::from_micros(aimd.initial_us),
let merkle_todo_bounded_queue = match config {
MerkleBackpressureEnum::None => None,
MerkleBackpressureEnum::FixedQueue(p) => {
Some(Arc::new(Semaphore::new(p.max_queue_size)))
}
};
let merkle_todo_sleep = Arc::new(Mutex::new(initial));
let insert_queue = db
.open_tree(format!("{}:insert_queue", F::TABLE_NAME))
@@ -83,7 +83,7 @@ impl<F: TableSchema, R: TableReplication> TableData<F, R> {
store.clone(),
merkle_tree.clone(),
merkle_todo.clone(),
merkle_todo_sleep.clone(),
merkle_todo_bounded_queue.clone(),
gc_todo.clone(),
);
@@ -95,7 +95,7 @@ impl<F: TableSchema, R: TableReplication> TableData<F, R> {
merkle_tree,
merkle_todo,
merkle_todo_notify: Notify::new(),
merkle_todo_sleep,
merkle_todo_bounded_queue,
insert_queue,
insert_queue_notify: Arc::new(Notify::new()),
gc_todo,
@@ -191,18 +191,17 @@ impl<F: TableSchema, R: TableReplication> TableData<F, R> {
// time to enable backpressure (ie. slow down clients).
// - When an entry is updated to be a tombstone, add it to the gc_todo tree
pub(crate) fn update_many<T: Borrow<ByteBuf>>(&self, entries: &[T]) -> Result<Duration, Error> {
let mut backpressure = Duration::ZERO;
pub(crate) fn update_many<T: Borrow<ByteBuf>>(&self, entries: &[T]) -> Result<(), Error> {
for update_bytes in entries.iter() {
backpressure += self.update_entry(update_bytes.borrow().as_slice())?;
self.update_entry(update_bytes.borrow().as_slice())?;
}
Ok(backpressure)
Ok(())
}
pub(crate) fn update_entry(&self, update_bytes: &[u8]) -> Result<Duration, Error> {
pub(crate) fn update_entry(&self, update_bytes: &[u8]) -> Result<(), Error> {
let update = self.decode_entry(update_bytes)?;
let ret = self.update_entry_with(
self.update_entry_with(
update.partition_key(),
update.sort_key(),
|_tx, ent| match ent {
@@ -213,12 +212,7 @@ impl<F: TableSchema, R: TableReplication> TableData<F, R> {
None => Ok(update.clone()),
},
)?;
let backpressure = match ret {
Some((_, d)) => d,
_ => Duration::ZERO,
};
Ok(backpressure)
Ok(())
}
pub fn update_entry_with(
@@ -226,7 +220,7 @@ impl<F: TableSchema, R: TableReplication> TableData<F, R> {
partition_key: &F::P,
sort_key: &F::S,
update_fn: impl Fn(&mut db::Transaction, Option<F::E>) -> db::TxOpResult<F::E>,
) -> Result<Option<(F::E, Duration)>, Error> {
) -> Result<Option<F::E>, Error> {
let tree_key = self.tree_key(partition_key, sort_key);
// transaction begins
@@ -297,16 +291,11 @@ impl<F: TableSchema, R: TableReplication> TableData<F, R> {
// Synchronize with the Merkle Worker
self.merkle_todo_notify.notify_one(); // Wake-up it
let backpressure = self.merkle_todo_sleep.clone().lock().unwrap().clone();
Ok(Some((new_entry, backpressure)))
Ok(Some(new_entry))
}
pub(crate) fn delete_if_equal(
self: &Arc<Self>,
k: &[u8],
v: &[u8],
) -> Result<(bool, Duration), Error> {
pub(crate) fn delete_if_equal(self: &Arc<Self>, k: &[u8], v: &[u8]) -> Result<bool, Error> {
let removed = self
.store
.db()
@@ -324,20 +313,19 @@ impl<F: TableSchema, R: TableReplication> TableData<F, R> {
})?;
if !removed {
return Ok((false, Duration::ZERO));
return Ok(false);
}
self.metrics.internal_delete_counter.add(1);
self.merkle_todo_notify.notify_one();
let backpressure = self.merkle_todo_sleep.clone().lock().unwrap().clone();
Ok((removed, backpressure))
Ok(removed)
}
pub(crate) fn delete_if_equal_hash(
self: &Arc<Self>,
k: &[u8],
vhash: Hash,
) -> Result<(bool, Duration), Error> {
) -> Result<bool, Error> {
let removed = self
.store
.db()
@@ -355,15 +343,13 @@ impl<F: TableSchema, R: TableReplication> TableData<F, R> {
})?;
if !removed {
return Ok((false, Duration::ZERO));
return Ok(false);
}
self.metrics.internal_delete_counter.add(1);
self.merkle_todo_notify.notify_one();
let bp = self.merkle_todo_sleep.clone().lock().unwrap().clone();
Ok((true, bp))
Ok(true)
}
// ---- Insert queue functions ----
+10 -10
View File
@@ -10,7 +10,6 @@ use serde_bytes::ByteBuf;
use futures::future::join_all;
use tokio::sync::watch;
use tokio::time::sleep;
use garage_db as db;
@@ -262,15 +261,13 @@ impl<F: TableSchema, R: TableReplication> TableGc<F, R> {
// GC has been successful for all of these entries.
// We now remove them all from our local table and from the GC todo list.
let mut backpressure = Duration::ZERO;
for item in items {
let (_is_removed, add_bp) = self
let _is_removed = self
.data
.delete_if_equal_hash(&item.key[..], item.value_hash)
.err_context("GC: local delete tombstones")?;
item.remove_if_equal(&self.data.gc_todo)
.err_context("GC: remove from todo list after successful GC")?;
backpressure += add_bp;
}
Ok(())
@@ -279,17 +276,21 @@ impl<F: TableSchema, R: TableReplication> TableGc<F, R> {
impl<F: TableSchema, R: TableReplication> EndpointHandler<GcRpc> for TableGc<F, R> {
async fn handle(self: &Arc<Self>, message: &GcRpc, _from: NodeID) -> Result<GcRpc, Error> {
let maybe_bounded = self.data.merkle_todo_bounded_queue.clone();
match message {
GcRpc::Update(items) => {
let backpressure = self.data.update_many(items)?;
sleep(backpressure).await;
if let Some(b) = maybe_bounded {
b.acquire_many(items.len() as u32).await.unwrap().forget();
}
self.data.update_many(items)?;
Ok(GcRpc::Ok)
}
GcRpc::DeleteIfEqualHash(items) => {
let mut backpressure = Duration::ZERO;
if let Some(b) = maybe_bounded {
b.acquire_many(items.len() as u32).await.unwrap().forget();
}
for (key, vhash) in items.iter() {
let (_is_removed, add_bp) = self.data.delete_if_equal_hash(&key[..], *vhash)?;
backpressure += add_bp;
let _is_removed = self.data.delete_if_equal_hash(&key[..], *vhash)?;
}
Ok(GcRpc::Ok)
}
@@ -336,7 +337,6 @@ impl<F: TableSchema, R: TableReplication> Worker for GcWorker<F, R> {
}
async fn wait_for_work(&mut self) -> WorkerState {
tokio::time::sleep(self.wait_delay).await;
WorkerState::Busy
}
}
+9 -68
View File
@@ -9,7 +9,7 @@ use tokio::sync::watch;
use garage_db as db;
use garage_util::background::*;
use garage_util::config::{MerkleBackpressureAimd, MerkleBackpressureEnum};
use garage_util::config::{MerkleBackpressureEnum, MerkleFixedQueue};
use garage_util::data::*;
use garage_util::encode::{nonversioned_decode, nonversioned_encode};
use garage_util::error::Error;
@@ -73,9 +73,12 @@ impl<F: TableSchema, R: TableReplication> MerkleUpdater<F, R> {
// @FIXME: move in worker
match &data.config {
MerkleBackpressureEnum::None => info!("Merkle Backpressure is not activated"),
MerkleBackpressureEnum::Aimd(v) => info!("Merkle backpressure is activated (initial={}us, max={}us, underload={}us, overload={}x)", v.initial_us, v.max_us, v.underload_us, v.overload_mult),
}
MerkleBackpressureEnum::None => info!("Merkle Backpressure is not activated"),
MerkleBackpressureEnum::FixedQueue(v) => info!(
"Merkle backpressure with a fixed queue size (qlen={}) is activated.",
v.max_queue_size
),
}
Arc::new(Self {
data,
@@ -84,7 +87,7 @@ impl<F: TableSchema, R: TableReplication> MerkleUpdater<F, R> {
}
pub(crate) fn spawn_workers(self: &Arc<Self>, background: &BackgroundRunner) {
background.spawn_worker(MerkleWorker(self.clone(), MerkleWorkerStats::new()));
background.spawn_worker(MerkleWorker(self.clone()));
}
fn updater_loop_iter(&self) -> Result<WorkerState, Error> {
@@ -303,65 +306,7 @@ impl<F: TableSchema, R: TableReplication> MerkleUpdater<F, R> {
}
}
struct MerkleWorker<F: TableSchema, R: TableReplication>(
Arc<MerkleUpdater<F, R>>,
MerkleWorkerStats,
);
struct MerkleWorkerStats {
last_update: SystemTime,
previous_merkle_len: usize,
}
impl MerkleWorkerStats {
fn new() -> MerkleWorkerStats {
Self {
last_update: SystemTime::now(),
previous_merkle_len: 0,
}
}
fn adapt_aimd_backpressure<F: TableSchema, R: TableReplication>(
&mut self,
updater: &MerkleUpdater<F, R>,
config: &MerkleBackpressureAimd,
) -> Result<(), Error> {
// Must have some elapsed time between runs
if self.last_update.elapsed().unwrap() < Duration::from_micros(config.sample_us) {
return Ok(()); // skip update
}
// Capture evolution of the merkle todo length
let prev_merkle_len = self.previous_merkle_len;
let current_merkle_len = updater.data.merkle_todo.len()?;
debug!(
"prev merkle len: {}, new merkle len: {}",
prev_merkle_len, current_merkle_len
);
// Algorithm inspired by Additive Increase Multiplicative Decrease (AIMD)
{
let a = updater.data.merkle_todo_sleep.clone();
let mut v = a.lock().unwrap();
let mut b;
if current_merkle_len <= prev_merkle_len {
// If we decrease the queue size, we can decrease the sleep time
b = v.saturating_sub(Duration::from_micros(config.underload_us));
} else {
// If we are late, we increase the queue size
b = v
.mul_f64(config.overload_mult)
.saturating_add(Duration::from_micros(1));
}
b = b.min(Duration::from_micros(config.max_us));
b = b.max(Duration::from_micros(config.initial_us));
debug!("sleep. before {} -> after {}", v.as_micros(), b.as_micros());
*v = b;
}
self.last_update = SystemTime::now();
self.previous_merkle_len = current_merkle_len;
Ok(())
}
}
struct MerkleWorker<F: TableSchema, R: TableReplication>(Arc<MerkleUpdater<F, R>>);
#[async_trait]
impl<F: TableSchema, R: TableReplication> Worker for MerkleWorker<F, R> {
@@ -378,10 +323,6 @@ impl<F: TableSchema, R: TableReplication> Worker for MerkleWorker<F, R> {
async fn work(&mut self, _must_exit: &mut watch::Receiver<bool>) -> Result<WorkerState, Error> {
let updater = self.0.clone();
match &updater.data.config {
MerkleBackpressureEnum::None => (),
MerkleBackpressureEnum::Aimd(a) => self.1.adapt_aimd_backpressure(&updater, a)?,
};
tokio::task::spawn_blocking(move || {
for _i in 0..100 {
let s = updater.updater_loop_iter();
+14 -10
View File
@@ -1,5 +1,7 @@
use opentelemetry::{global, metrics::*, KeyValue};
use std::sync::{Arc, Mutex};
use std::convert::TryInto;
use std::sync::Arc;
use tokio::sync::Semaphore;
use garage_db as db;
@@ -8,7 +10,7 @@ pub struct TableMetrics {
pub(crate) _table_size: ValueObserver<u64>,
pub(crate) _merkle_tree_size: ValueObserver<u64>,
pub(crate) _merkle_todo_len: ValueObserver<u64>,
pub(crate) _merkle_todo_sleep_ms: ValueObserver<f64>,
pub(crate) _merkle_todo_bounded_queue_free: ValueObserver<u64>,
pub(crate) _gc_todo_len: ValueObserver<u64>,
pub(crate) get_request_counter: BoundCounter<u64>,
@@ -28,7 +30,7 @@ impl TableMetrics {
store: db::Tree,
merkle_tree: db::Tree,
merkle_todo: db::Tree,
merkle_todo_sleep: Arc<Mutex<std::time::Duration>>,
merkle_todo_bounded_queue: Option<Arc<Semaphore>>,
gc_todo: db::Tree,
) -> Self {
let meter = global::meter(table_name);
@@ -75,14 +77,16 @@ impl TableMetrics {
)
.with_description("Merkle tree updater TODO queue length")
.init(),
_merkle_todo_sleep_ms: meter
.f64_value_observer(
"table.merkle_updater_todo_queue_backpressure_ms",
_merkle_todo_bounded_queue_free: meter
.u64_value_observer(
"table.merkle_todo_bounded_queue_free",
move |observer| {
let bp_ref = merkle_todo_sleep.clone();
let bp_val = bp_ref.lock().unwrap();
let bp_millis: f64 = bp_val.as_micros() as f64 / 1000.0f64;
observer.observe(bp_millis, &[KeyValue::new("table_name", table_name)])
let maybe_bounded = merkle_todo_bounded_queue.clone();
let free: u64 = match &maybe_bounded {
Some(v) => v.available_permits().try_into().unwrap(),
None => 0,
};
observer.observe(free, &[KeyValue::new("table_name", table_name)])
}
)
.with_description("Merkle tree updater TODO sleep backpressure to apply in ms")
+7 -4
View File
@@ -245,12 +245,16 @@ impl<F: TableSchema, R: TableReplication> TableSyncer<F, R> {
// All remote nodes have written those items, now we can delete them locally
let mut not_removed = 0;
let maybe_bounded = self.data.merkle_todo_bounded_queue.clone();
if let Some(b) = maybe_bounded {
b.acquire_many(items.len() as u32).await.unwrap().forget();
}
for (k, v) in items.iter() {
let (removed, backpressure) = self.data.delete_if_equal(&k[..], &v[..])?;
let removed = self.data.delete_if_equal(&k[..], &v[..])?;
if !removed {
not_removed += 1;
}
sleep(backpressure).await;
}
if not_removed > 0 {
@@ -471,8 +475,7 @@ impl<F: TableSchema, R: TableReplication> EndpointHandler<SyncRpc> for TableSync
],
);
let backpressure = self.data.update_many(items)?;
sleep(backpressure).await;
self.data.update_many(items)?;
Ok(SyncRpc::Ok)
}
m => Err(Error::unexpected_rpc_message(m)),
+5 -3
View File
@@ -5,7 +5,6 @@ use std::sync::Arc;
use futures::stream::*;
use serde::{Deserialize, Serialize};
use serde_bytes::ByteBuf;
use tokio::time::sleep;
use opentelemetry::{
trace::{FutureExt, TraceContextExt, Tracer},
@@ -536,8 +535,11 @@ impl<F: TableSchema, R: TableReplication> EndpointHandler<TableRpc<F>> for Table
Ok(TableRpc::Update(values))
}
TableRpc::Update(pairs) => {
let backpressure = self.data.update_many(pairs)?;
sleep(backpressure).await;
let maybe_bounded = self.data.merkle_todo_bounded_queue.clone();
if let Some(b) = maybe_bounded {
b.acquire_many(pairs.len() as u32).await.unwrap().forget();
}
self.data.update_many(pairs)?;
Ok(TableRpc::Ok)
}
m => Err(Error::unexpected_rpc_message(m)),
+23 -28
View File
@@ -262,6 +262,7 @@ pub struct KubernetesDiscoveryConfig {
#[derive(Deserialize, Debug, Clone, Default)]
pub struct ExperimentalConfig {
pub merkle_backpressure: MerkleBackpressureEnum,
pub rpc_in_flight_limiters: RpcInFlightLimiterEnum,
}
#[derive(Deserialize, Debug, Clone, Default)]
@@ -269,21 +270,27 @@ pub struct ExperimentalConfig {
pub enum MerkleBackpressureEnum {
#[default]
None,
Aimd(MerkleBackpressureAimd),
FixedQueue(MerkleFixedQueue),
}
#[derive(Deserialize, Debug, Clone, Default)]
pub struct MerkleBackpressureAimd {
#[serde(default = "default_initial_us")]
pub initial_us: u64,
#[serde(default = "default_max_us")]
pub max_us: u64,
#[serde(default = "default_underload_us")]
pub underload_us: u64,
#[serde(default = "default_overload_mult")]
pub overload_mult: f64,
#[serde(default = "default_sample_us")]
pub sample_us: u64,
#[serde(rename_all = "lowercase", tag = "kind")]
pub enum RpcInFlightLimiterEnum {
#[default]
None,
FixedSize(InFlightFixedSize),
}
#[derive(Deserialize, Debug, Clone, Default)]
pub struct InFlightFixedSize {
#[serde(default = "default_max_table_write")]
pub max_table_write: usize,
}
#[derive(Deserialize, Debug, Clone, Default)]
pub struct MerkleFixedQueue {
#[serde(default = "default_max_queue_size")]
pub max_queue_size: usize,
}
/// Read and parse configuration
@@ -312,24 +319,12 @@ fn default_compression() -> Option<i32> {
Some(1)
}
fn default_initial_us() -> u64 {
10
fn default_max_table_write() -> usize {
64
}
fn default_max_us() -> u64 {
30 * 1000 * 1000
}
fn default_underload_us() -> u64 {
10
}
fn default_overload_mult() -> f64 {
1.1
}
fn default_sample_us() -> u64 {
100 * 1000
fn default_max_queue_size() -> usize {
256
}
fn deserialize_compression<'de, D>(deserializer: D) -> Result<Option<i32>, D::Error>