Skip to main content

linera_views/backends/
metering.rs

1// Copyright (c) Zefchain Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Adds metrics to a key-value store.
5
6use std::{
7    collections::{btree_map::Entry, BTreeMap},
8    sync::{Arc, LazyLock, Mutex},
9};
10
11use convert_case::{Case, Casing};
12use linera_base::prometheus_util::{
13    exponential_bucket_latencies, register_histogram_vec, register_int_counter_vec,
14    MeasureLatency as _,
15};
16use prometheus::{exponential_buckets, HistogramVec, IntCounterVec};
17
18#[cfg(with_testing)]
19use crate::store::TestKeyValueDatabase;
20use crate::{
21    batch::Batch,
22    store::{KeyValueDatabase, ReadableKeyValueStore, WithError, WritableKeyValueStore},
23};
24
25#[derive(Clone)]
26/// The implementation of the `KeyValueStoreMetrics` for the `KeyValueStore`.
27pub struct KeyValueStoreMetrics {
28    read_value_bytes_latency: HistogramVec,
29    contains_key_latency: HistogramVec,
30    contains_keys_latency: HistogramVec,
31    read_multi_values_bytes_latency: HistogramVec,
32    find_keys_by_prefix_latency: HistogramVec,
33    find_key_values_by_prefix_latency: HistogramVec,
34    write_batch_latency: HistogramVec,
35    clear_journal_latency: HistogramVec,
36    connect_latency: HistogramVec,
37    open_shared_latency: HistogramVec,
38    open_exclusive_latency: HistogramVec,
39    list_all_latency: HistogramVec,
40    list_root_keys_latency: HistogramVec,
41    delete_all_latency: HistogramVec,
42    exists_latency: HistogramVec,
43    create_latency: HistogramVec,
44    delete_latency: HistogramVec,
45    read_value_none_cases: IntCounterVec,
46    read_value_key_size: HistogramVec,
47    read_value_value_size: HistogramVec,
48    read_multi_values_num_entries: HistogramVec,
49    read_multi_values_key_sizes: HistogramVec,
50    contains_keys_num_entries: HistogramVec,
51    contains_keys_key_sizes: HistogramVec,
52    contains_key_key_size: HistogramVec,
53    find_keys_by_prefix_prefix_size: HistogramVec,
54    find_keys_by_prefix_num_keys: HistogramVec,
55    find_keys_by_prefix_keys_size: HistogramVec,
56    find_key_values_by_prefix_prefix_size: HistogramVec,
57    find_key_values_by_prefix_num_keys: HistogramVec,
58    find_key_values_by_prefix_key_values_size: HistogramVec,
59    write_batch_size: HistogramVec,
60    list_all_sizes: HistogramVec,
61    exists_true_cases: IntCounterVec,
62}
63
64#[derive(Default)]
65struct StoreMetrics {
66    stores: BTreeMap<String, Arc<KeyValueStoreMetrics>>,
67}
68
69/// The global variables of the RocksDB stores
70static STORE_COUNTERS: LazyLock<Mutex<StoreMetrics>> =
71    LazyLock::new(|| Mutex::new(StoreMetrics::default()));
72
73fn get_counter(name: &str) -> Arc<KeyValueStoreMetrics> {
74    let mut store_metrics = STORE_COUNTERS.lock().unwrap();
75    let key = name.to_string();
76    match store_metrics.stores.entry(key) {
77        Entry::Occupied(entry) => {
78            let entry = entry.into_mut();
79            entry.clone()
80        }
81        Entry::Vacant(entry) => {
82            let store_metric = Arc::new(KeyValueStoreMetrics::new(name));
83            entry.insert(store_metric.clone());
84            store_metric
85        }
86    }
87}
88
89impl KeyValueStoreMetrics {
90    /// Creation of a named Metered counter.
91    pub fn new(name: &str) -> Self {
92        // name can be "rocks db". Then var_name = "rocks_db" and title_name = "RocksDb"
93        let var_name = name.replace(' ', "_");
94        let title_name = name.to_case(Case::Snake);
95
96        // Latency buckets in milliseconds: up to 10 seconds
97        let latency_buckets = exponential_bucket_latencies(10000.0);
98        // Size buckets in bytes: 1B to 10MB
99        let size_buckets =
100            Some(exponential_buckets(1.0, 4.0, 12).expect("Size buckets creation should not fail"));
101        // Count buckets: 1 to 100,000
102        let count_buckets = Some(
103            exponential_buckets(1.0, 3.0, 11).expect("Count buckets creation should not fail"),
104        );
105
106        let entry1 = format!("{var_name}_read_value_bytes_latency");
107        let entry2 = format!("{title_name} read value bytes latency");
108        let read_value_bytes_latency =
109            register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
110
111        let entry1 = format!("{var_name}_contains_key_latency");
112        let entry2 = format!("{title_name} contains key latency");
113        let contains_key_latency =
114            register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
115
116        let entry1 = format!("{var_name}_contains_keys_latency");
117        let entry2 = format!("{title_name} contains keys latency");
118        let contains_keys_latency =
119            register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
120
121        let entry1 = format!("{var_name}_read_multi_value_bytes_latency");
122        let entry2 = format!("{title_name} read multi value bytes latency");
123        let read_multi_values_bytes_latency =
124            register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
125
126        let entry1 = format!("{var_name}_find_keys_by_prefix_latency");
127        let entry2 = format!("{title_name} find keys by prefix latency");
128        let find_keys_by_prefix_latency =
129            register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
130
131        let entry1 = format!("{var_name}_find_key_values_by_prefix_latency");
132        let entry2 = format!("{title_name} find key values by prefix latency");
133        let find_key_values_by_prefix_latency =
134            register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
135
136        let entry1 = format!("{var_name}_write_batch_latency");
137        let entry2 = format!("{title_name} write batch latency");
138        let write_batch_latency =
139            register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
140
141        let entry1 = format!("{var_name}_clear_journal_latency");
142        let entry2 = format!("{title_name} clear journal latency");
143        let clear_journal_latency =
144            register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
145
146        let entry1 = format!("{var_name}_connect_latency");
147        let entry2 = format!("{title_name} connect latency");
148        let connect_latency =
149            register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
150
151        let entry1 = format!("{var_name}_open_shared_latency");
152        let entry2 = format!("{title_name} open shared partition");
153        let open_shared_latency =
154            register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
155
156        let entry1 = format!("{var_name}_open_exclusive_latency");
157        let entry2 = format!("{title_name} open exclusive partition");
158        let open_exclusive_latency =
159            register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
160
161        let entry1 = format!("{var_name}_list_all_latency");
162        let entry2 = format!("{title_name} list all latency");
163        let list_all_latency =
164            register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
165
166        let entry1 = format!("{var_name}_list_root_keys_latency");
167        let entry2 = format!("{title_name} list root keys latency");
168        let list_root_keys_latency =
169            register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
170
171        let entry1 = format!("{var_name}_delete_all_latency");
172        let entry2 = format!("{title_name} delete all latency");
173        let delete_all_latency =
174            register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
175
176        let entry1 = format!("{var_name}_exists_latency");
177        let entry2 = format!("{title_name} exists latency");
178        let exists_latency = register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
179
180        let entry1 = format!("{var_name}_create_latency");
181        let entry2 = format!("{title_name} create latency");
182        let create_latency = register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
183
184        let entry1 = format!("{var_name}_delete_latency");
185        let entry2 = format!("{title_name} delete latency");
186        let delete_latency = register_histogram_vec(&entry1, &entry2, &[], latency_buckets.clone());
187
188        let entry1 = format!("{var_name}_read_value_none_cases");
189        let entry2 = format!("{title_name} read value none cases");
190        let read_value_none_cases = register_int_counter_vec(&entry1, &entry2, &[]);
191
192        let entry1 = format!("{var_name}_read_value_key_size");
193        let entry2 = format!("{title_name} read value key size");
194        let read_value_key_size =
195            register_histogram_vec(&entry1, &entry2, &[], size_buckets.clone());
196
197        let entry1 = format!("{var_name}_read_value_value_size");
198        let entry2 = format!("{title_name} read value value size");
199        let read_value_value_size =
200            register_histogram_vec(&entry1, &entry2, &[], size_buckets.clone());
201
202        let entry1 = format!("{var_name}_read_multi_values_num_entries");
203        let entry2 = format!("{title_name} read multi values num entries");
204        let read_multi_values_num_entries =
205            register_histogram_vec(&entry1, &entry2, &[], count_buckets.clone());
206
207        let entry1 = format!("{var_name}_read_multi_values_key_sizes");
208        let entry2 = format!("{title_name} read multi values key sizes");
209        let read_multi_values_key_sizes =
210            register_histogram_vec(&entry1, &entry2, &[], size_buckets.clone());
211
212        let entry1 = format!("{var_name}_contains_keys_num_entries");
213        let entry2 = format!("{title_name} contains keys num entries");
214        let contains_keys_num_entries =
215            register_histogram_vec(&entry1, &entry2, &[], count_buckets.clone());
216
217        let entry1 = format!("{var_name}_contains_keys_key_sizes");
218        let entry2 = format!("{title_name} contains keys key sizes");
219        let contains_keys_key_sizes =
220            register_histogram_vec(&entry1, &entry2, &[], size_buckets.clone());
221
222        let entry1 = format!("{var_name}_contains_key_key_size");
223        let entry2 = format!("{title_name} contains key key size");
224        let contains_key_key_size =
225            register_histogram_vec(&entry1, &entry2, &[], size_buckets.clone());
226
227        let entry1 = format!("{var_name}_find_keys_by_prefix_prefix_size");
228        let entry2 = format!("{title_name} find keys by prefix prefix size");
229        let find_keys_by_prefix_prefix_size =
230            register_histogram_vec(&entry1, &entry2, &[], size_buckets.clone());
231
232        let entry1 = format!("{var_name}_find_keys_by_prefix_num_keys");
233        let entry2 = format!("{title_name} find keys by prefix num keys");
234        let find_keys_by_prefix_num_keys =
235            register_histogram_vec(&entry1, &entry2, &[], count_buckets.clone());
236
237        let entry1 = format!("{var_name}_find_keys_by_prefix_keys_size");
238        let entry2 = format!("{title_name} find keys by prefix keys size");
239        let find_keys_by_prefix_keys_size =
240            register_histogram_vec(&entry1, &entry2, &[], size_buckets.clone());
241
242        let entry1 = format!("{var_name}_find_key_values_by_prefix_prefix_size");
243        let entry2 = format!("{title_name} find key values by prefix prefix size");
244        let find_key_values_by_prefix_prefix_size =
245            register_histogram_vec(&entry1, &entry2, &[], size_buckets.clone());
246
247        let entry1 = format!("{var_name}_find_key_values_by_prefix_num_keys");
248        let entry2 = format!("{title_name} find key values by prefix num keys");
249        let find_key_values_by_prefix_num_keys =
250            register_histogram_vec(&entry1, &entry2, &[], count_buckets.clone());
251
252        let entry1 = format!("{var_name}_find_key_values_by_prefix_key_values_size");
253        let entry2 = format!("{title_name} find key values by prefix key values size");
254        let find_key_values_by_prefix_key_values_size =
255            register_histogram_vec(&entry1, &entry2, &[], size_buckets.clone());
256
257        let entry1 = format!("{var_name}_write_batch_size");
258        let entry2 = format!("{title_name} write batch size");
259        let write_batch_size = register_histogram_vec(&entry1, &entry2, &[], size_buckets);
260
261        let entry1 = format!("{var_name}_list_all_sizes");
262        let entry2 = format!("{title_name} list all sizes");
263        let list_all_sizes = register_histogram_vec(&entry1, &entry2, &[], count_buckets);
264
265        let entry1 = format!("{var_name}_exists_true_cases");
266        let entry2 = format!("{title_name} exists true cases");
267        let exists_true_cases = register_int_counter_vec(&entry1, &entry2, &[]);
268
269        KeyValueStoreMetrics {
270            read_value_bytes_latency,
271            contains_key_latency,
272            contains_keys_latency,
273            read_multi_values_bytes_latency,
274            find_keys_by_prefix_latency,
275            find_key_values_by_prefix_latency,
276            write_batch_latency,
277            clear_journal_latency,
278            connect_latency,
279            open_shared_latency,
280            open_exclusive_latency,
281            list_all_latency,
282            list_root_keys_latency,
283            delete_all_latency,
284            exists_latency,
285            create_latency,
286            delete_latency,
287            read_value_none_cases,
288            read_value_key_size,
289            read_value_value_size,
290            read_multi_values_num_entries,
291            read_multi_values_key_sizes,
292            contains_keys_num_entries,
293            contains_keys_key_sizes,
294            contains_key_key_size,
295            find_keys_by_prefix_prefix_size,
296            find_keys_by_prefix_num_keys,
297            find_keys_by_prefix_keys_size,
298            find_key_values_by_prefix_prefix_size,
299            find_key_values_by_prefix_num_keys,
300            find_key_values_by_prefix_key_values_size,
301            write_batch_size,
302            list_all_sizes,
303            exists_true_cases,
304        }
305    }
306}
307
308/// A metered database that keeps track of every operation.
309#[derive(Clone)]
310pub struct MeteredDatabase<D> {
311    /// The metrics being computed.
312    counter: Arc<KeyValueStoreMetrics>,
313    /// The underlying database.
314    database: D,
315}
316
317/// A metered store that keeps track of every operation.
318#[derive(Clone)]
319pub struct MeteredStore<S> {
320    /// The metrics being computed.
321    counter: Arc<KeyValueStoreMetrics>,
322    /// The underlying store.
323    store: S,
324}
325
326impl<D> WithError for MeteredDatabase<D>
327where
328    D: WithError,
329{
330    type Error = D::Error;
331}
332
333impl<S> WithError for MeteredStore<S>
334where
335    S: WithError,
336{
337    type Error = S::Error;
338}
339
340impl<S> ReadableKeyValueStore for MeteredStore<S>
341where
342    S: ReadableKeyValueStore,
343{
344    const MAX_KEY_SIZE: usize = S::MAX_KEY_SIZE;
345
346    fn root_key(&self) -> Result<Vec<u8>, Self::Error> {
347        self.store.root_key()
348    }
349
350    async fn read_value_bytes(&self, key: &[u8]) -> Result<Option<Vec<u8>>, Self::Error> {
351        let _latency = self.counter.read_value_bytes_latency.measure_latency();
352        self.counter
353            .read_value_key_size
354            .with_label_values(&[])
355            .observe(key.len() as f64);
356        let result = self.store.read_value_bytes(key).await?;
357        match &result {
358            None => self
359                .counter
360                .read_value_none_cases
361                .with_label_values(&[])
362                .inc(),
363            Some(value) => self
364                .counter
365                .read_value_value_size
366                .with_label_values(&[])
367                .observe(value.len() as f64),
368        }
369        Ok(result)
370    }
371
372    async fn contains_key(&self, key: &[u8]) -> Result<bool, Self::Error> {
373        let _latency = self.counter.contains_key_latency.measure_latency();
374        self.counter
375            .contains_key_key_size
376            .with_label_values(&[])
377            .observe(key.len() as f64);
378        self.store.contains_key(key).await
379    }
380
381    async fn contains_keys(&self, keys: &[Vec<u8>]) -> Result<Vec<bool>, Self::Error> {
382        let _latency = self.counter.contains_keys_latency.measure_latency();
383        self.counter
384            .contains_keys_num_entries
385            .with_label_values(&[])
386            .observe(keys.len() as f64);
387        let key_sizes = keys.iter().map(|k| k.len()).sum::<usize>();
388        self.counter
389            .contains_keys_key_sizes
390            .with_label_values(&[])
391            .observe(key_sizes as f64);
392        self.store.contains_keys(keys).await
393    }
394
395    async fn read_multi_values_bytes(
396        &self,
397        keys: &[Vec<u8>],
398    ) -> Result<Vec<Option<Vec<u8>>>, Self::Error> {
399        let _latency = self
400            .counter
401            .read_multi_values_bytes_latency
402            .measure_latency();
403        self.counter
404            .read_multi_values_num_entries
405            .with_label_values(&[])
406            .observe(keys.len() as f64);
407        let key_sizes = keys.iter().map(|k| k.len()).sum::<usize>();
408        self.counter
409            .read_multi_values_key_sizes
410            .with_label_values(&[])
411            .observe(key_sizes as f64);
412        self.store.read_multi_values_bytes(keys).await
413    }
414
415    async fn find_keys_by_prefix(&self, key_prefix: &[u8]) -> Result<Vec<Vec<u8>>, Self::Error> {
416        let _latency = self.counter.find_keys_by_prefix_latency.measure_latency();
417        self.counter
418            .find_keys_by_prefix_prefix_size
419            .with_label_values(&[])
420            .observe(key_prefix.len() as f64);
421        let result = self.store.find_keys_by_prefix(key_prefix).await?;
422        let (num_keys, keys_size) = result
423            .iter()
424            .map(|key| key.len())
425            .fold((0, 0), |(count, size), len| (count + 1, size + len));
426        self.counter
427            .find_keys_by_prefix_num_keys
428            .with_label_values(&[])
429            .observe(num_keys as f64);
430        self.counter
431            .find_keys_by_prefix_keys_size
432            .with_label_values(&[])
433            .observe(keys_size as f64);
434        Ok(result)
435    }
436
437    async fn find_key_values_by_prefix(
438        &self,
439        key_prefix: &[u8],
440    ) -> Result<Vec<(Vec<u8>, Vec<u8>)>, Self::Error> {
441        let _latency = self
442            .counter
443            .find_key_values_by_prefix_latency
444            .measure_latency();
445        self.counter
446            .find_key_values_by_prefix_prefix_size
447            .with_label_values(&[])
448            .observe(key_prefix.len() as f64);
449        let result = self.store.find_key_values_by_prefix(key_prefix).await?;
450        let (num_keys, key_values_size) = result
451            .iter()
452            .map(|(key, value)| key.len() + value.len())
453            .fold((0, 0), |(count, size), len| (count + 1, size + len));
454        self.counter
455            .find_key_values_by_prefix_num_keys
456            .with_label_values(&[])
457            .observe(num_keys as f64);
458        self.counter
459            .find_key_values_by_prefix_key_values_size
460            .with_label_values(&[])
461            .observe(key_values_size as f64);
462        Ok(result)
463    }
464}
465
466impl<S> WritableKeyValueStore for MeteredStore<S>
467where
468    S: WritableKeyValueStore,
469{
470    const MAX_VALUE_SIZE: usize = S::MAX_VALUE_SIZE;
471
472    async fn write_batch(&self, batch: Batch) -> Result<(), Self::Error> {
473        let _latency = self.counter.write_batch_latency.measure_latency();
474        self.counter
475            .write_batch_size
476            .with_label_values(&[])
477            .observe(batch.size() as f64);
478        self.store.write_batch(batch).await
479    }
480
481    async fn clear_journal(&self) -> Result<(), Self::Error> {
482        let _metric = self.counter.clear_journal_latency.measure_latency();
483        self.store.clear_journal().await
484    }
485}
486
487impl<D> KeyValueDatabase for MeteredDatabase<D>
488where
489    D: KeyValueDatabase,
490{
491    type Config = D::Config;
492    type Store = MeteredStore<D::Store>;
493
494    fn get_name() -> String {
495        D::get_name()
496    }
497
498    async fn connect(config: &Self::Config, namespace: &str) -> Result<Self, Self::Error> {
499        let name = D::get_name();
500        let counter = get_counter(&name);
501        let _latency = counter.connect_latency.measure_latency();
502        let database = D::connect(config, namespace).await?;
503        let counter = get_counter(&name);
504        Ok(Self { counter, database })
505    }
506
507    fn open_shared(&self, root_key: &[u8]) -> Result<Self::Store, Self::Error> {
508        let _latency = self.counter.open_shared_latency.measure_latency();
509        let store = self.database.open_shared(root_key)?;
510        let counter = self.counter.clone();
511        Ok(MeteredStore { counter, store })
512    }
513
514    fn open_exclusive(&self, root_key: &[u8]) -> Result<Self::Store, Self::Error> {
515        let _latency = self.counter.open_exclusive_latency.measure_latency();
516        let store = self.database.open_exclusive(root_key)?;
517        let counter = self.counter.clone();
518        Ok(MeteredStore { counter, store })
519    }
520
521    async fn list_all(config: &Self::Config) -> Result<Vec<String>, Self::Error> {
522        let name = D::get_name();
523        let counter = get_counter(&name);
524        let _latency = counter.list_all_latency.measure_latency();
525        let namespaces = D::list_all(config).await?;
526        let counter = get_counter(&name);
527        counter
528            .list_all_sizes
529            .with_label_values(&[])
530            .observe(namespaces.len() as f64);
531        Ok(namespaces)
532    }
533
534    async fn list_root_keys(&self) -> Result<Vec<Vec<u8>>, Self::Error> {
535        let _latency = self.counter.list_root_keys_latency.measure_latency();
536        self.database.list_root_keys().await
537    }
538
539    async fn delete_all(config: &Self::Config) -> Result<(), Self::Error> {
540        let name = D::get_name();
541        let counter = get_counter(&name);
542        let _latency = counter.delete_all_latency.measure_latency();
543        D::delete_all(config).await
544    }
545
546    async fn exists(config: &Self::Config, namespace: &str) -> Result<bool, Self::Error> {
547        let name = D::get_name();
548        let counter = get_counter(&name);
549        let _latency = counter.exists_latency.measure_latency();
550        let result = D::exists(config, namespace).await?;
551        if result {
552            let counter = get_counter(&name);
553            counter.exists_true_cases.with_label_values(&[]).inc();
554        }
555        Ok(result)
556    }
557
558    async fn create(config: &Self::Config, namespace: &str) -> Result<(), Self::Error> {
559        let name = D::get_name();
560        let counter = get_counter(&name);
561        let _latency = counter.create_latency.measure_latency();
562        D::create(config, namespace).await
563    }
564
565    async fn delete(config: &Self::Config, namespace: &str) -> Result<(), Self::Error> {
566        let name = D::get_name();
567        let counter = get_counter(&name);
568        let _latency = counter.delete_latency.measure_latency();
569        D::delete(config, namespace).await
570    }
571}
572
573#[cfg(with_testing)]
574impl<D> TestKeyValueDatabase for MeteredDatabase<D>
575where
576    D: TestKeyValueDatabase,
577{
578    async fn new_test_config() -> Result<D::Config, Self::Error> {
579        D::new_test_config().await
580    }
581}
582
583#[cfg(with_testing)]
584impl<D: crate::backends::DatabaseBackup> crate::backends::DatabaseBackup for MeteredDatabase<D> {
585    fn backup_to(&self, dir: &std::path::Path) -> anyhow::Result<()> {
586        self.database.backup_to(dir)
587    }
588}