diff --git a/src/model/k2v/rpc.rs b/src/model/k2v/rpc.rs index 51ff54f9..821f4549 100644 --- a/src/model/k2v/rpc.rs +++ b/src/model/k2v/rpc.rs @@ -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 = 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 { 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, Error> { + ) -> Result, Error> { let now = now_msec(); self.item_table diff --git a/src/table/data.rs b/src/table/data.rs index 46c546f7..2af72253 100644 --- a/src/table/data.rs +++ b/src/table/data.rs @@ -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 { pub(crate) merkle_tree: db::Tree, pub(crate) merkle_todo: db::Tree, pub(crate) merkle_todo_notify: Notify, - pub(crate) merkle_todo_sleep: Arc>, + pub(crate) merkle_todo_bounded_queue: Option>, pub(crate) insert_queue: db::Tree, pub(crate) insert_queue_notify: Arc, @@ -64,11 +63,12 @@ impl TableData { .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 TableData { 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 TableData { 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 TableData { // 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>(&self, entries: &[T]) -> Result { - let mut backpressure = Duration::ZERO; + pub(crate) fn update_many>(&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 { + 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 TableData { 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 TableData { partition_key: &F::P, sort_key: &F::S, update_fn: impl Fn(&mut db::Transaction, Option) -> db::TxOpResult, - ) -> Result, Error> { + ) -> Result, Error> { let tree_key = self.tree_key(partition_key, sort_key); // transaction begins @@ -297,16 +291,11 @@ impl TableData { // 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, - k: &[u8], - v: &[u8], - ) -> Result<(bool, Duration), Error> { + pub(crate) fn delete_if_equal(self: &Arc, k: &[u8], v: &[u8]) -> Result { let removed = self .store .db() @@ -324,20 +313,19 @@ impl TableData { })?; 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, k: &[u8], vhash: Hash, - ) -> Result<(bool, Duration), Error> { + ) -> Result { let removed = self .store .db() @@ -355,15 +343,13 @@ impl TableData { })?; 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 ---- diff --git a/src/table/gc.rs b/src/table/gc.rs index 748b475b..2154d93c 100644 --- a/src/table/gc.rs +++ b/src/table/gc.rs @@ -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 TableGc { // 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 TableGc { impl EndpointHandler for TableGc { async fn handle(self: &Arc, message: &GcRpc, _from: NodeID) -> Result { + 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 Worker for GcWorker { } async fn wait_for_work(&mut self) -> WorkerState { - tokio::time::sleep(self.wait_delay).await; WorkerState::Busy } } diff --git a/src/table/merkle.rs b/src/table/merkle.rs index bfec0696..35233ed4 100644 --- a/src/table/merkle.rs +++ b/src/table/merkle.rs @@ -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 MerkleUpdater { // @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 MerkleUpdater { } pub(crate) fn spawn_workers(self: &Arc, background: &BackgroundRunner) { - background.spawn_worker(MerkleWorker(self.clone(), MerkleWorkerStats::new())); + background.spawn_worker(MerkleWorker(self.clone())); } fn updater_loop_iter(&self) -> Result { @@ -303,65 +306,7 @@ impl MerkleUpdater { } } -struct MerkleWorker( - Arc>, - 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( - &mut self, - updater: &MerkleUpdater, - 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(Arc>); #[async_trait] impl Worker for MerkleWorker { @@ -378,10 +323,6 @@ impl Worker for MerkleWorker { async fn work(&mut self, _must_exit: &mut watch::Receiver) -> Result { 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(); diff --git a/src/table/metrics.rs b/src/table/metrics.rs index aad962fd..c9e6d048 100644 --- a/src/table/metrics.rs +++ b/src/table/metrics.rs @@ -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, pub(crate) _merkle_tree_size: ValueObserver, pub(crate) _merkle_todo_len: ValueObserver, - pub(crate) _merkle_todo_sleep_ms: ValueObserver, + pub(crate) _merkle_todo_bounded_queue_free: ValueObserver, pub(crate) _gc_todo_len: ValueObserver, pub(crate) get_request_counter: BoundCounter, @@ -28,7 +30,7 @@ impl TableMetrics { store: db::Tree, merkle_tree: db::Tree, merkle_todo: db::Tree, - merkle_todo_sleep: Arc>, + merkle_todo_bounded_queue: Option>, 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") diff --git a/src/table/sync.rs b/src/table/sync.rs index b527b869..0db440ba 100644 --- a/src/table/sync.rs +++ b/src/table/sync.rs @@ -245,12 +245,16 @@ impl TableSyncer { // 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 EndpointHandler 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)), diff --git a/src/table/table.rs b/src/table/table.rs index 6d0e9216..3fad75cf 100644 --- a/src/table/table.rs +++ b/src/table/table.rs @@ -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 EndpointHandler> 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)), diff --git a/src/util/config.rs b/src/util/config.rs index fa37111e..b58047b0 100644 --- a/src/util/config.rs +++ b/src/util/config.rs @@ -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 { 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, D::Error>