1use 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)]
26pub 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
69static 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 pub fn new(name: &str) -> Self {
92 let var_name = name.replace(' ', "_");
94 let title_name = name.to_case(Case::Snake);
95
96 let latency_buckets = exponential_bucket_latencies(10000.0);
98 let size_buckets =
100 Some(exponential_buckets(1.0, 4.0, 12).expect("Size buckets creation should not fail"));
101 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#[derive(Clone)]
310pub struct MeteredDatabase<D> {
311 counter: Arc<KeyValueStoreMetrics>,
313 database: D,
315}
316
317#[derive(Clone)]
319pub struct MeteredStore<S> {
320 counter: Arc<KeyValueStoreMetrics>,
322 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}