Skip to main content

iota_storage/
key_value_store.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5//! Immutable key/value store trait for storing/retrieving transactions,
6//! effects, and events to/from a scalable.
7
8use std::{sync::Arc, time::Instant};
9
10use async_trait::async_trait;
11use iota_sdk_types::{
12    CheckpointDigest, ObjectId, TransactionDigest, TransactionEffects, TransactionEvents, Version,
13    checkpoint::CheckpointContents,
14};
15use iota_types::{
16    base_types::VersionNumber,
17    error::{IotaError, IotaResult, UserInputError},
18    messages_checkpoint::{CertifiedCheckpointSummary, CheckpointSequenceNumber},
19    object::Object,
20    storage::ObjectKey,
21    transaction::TransactionEnvelope,
22};
23use tracing::instrument;
24
25use crate::key_value_store_metrics::KeyValueStoreMetrics;
26
27pub type KVStoreTransactionData = (
28    Vec<Option<TransactionEnvelope>>,
29    Vec<Option<TransactionEffects>>,
30);
31
32pub type KVStoreCheckpointData = (
33    Vec<Option<CertifiedCheckpointSummary>>,
34    Vec<Option<CheckpointContents>>,
35    Vec<Option<CertifiedCheckpointSummary>>,
36);
37
38pub struct TransactionKeyValueStore {
39    store_name: &'static str,
40    metrics: Arc<KeyValueStoreMetrics>,
41    inner: Arc<dyn TransactionKeyValueStoreTrait + Send + Sync>,
42}
43
44impl TransactionKeyValueStore {
45    pub fn new(
46        store_name: &'static str,
47        metrics: Arc<KeyValueStoreMetrics>,
48        inner: Arc<dyn TransactionKeyValueStoreTrait + Send + Sync>,
49    ) -> Self {
50        Self {
51            store_name,
52            metrics,
53            inner,
54        }
55    }
56
57    /// Generic multi_get, allows implementors to get heterogenous values with a
58    /// single round trip.
59    pub async fn multi_get(
60        &self,
61        transaction_keys: &[TransactionDigest],
62        effects_keys: &[TransactionDigest],
63    ) -> IotaResult<KVStoreTransactionData> {
64        let start = Instant::now();
65        let res = self.inner.multi_get(transaction_keys, effects_keys).await;
66        let elapsed = start.elapsed();
67
68        let num_txns = transaction_keys.len() as u64;
69        let num_effects = effects_keys.len() as u64;
70        let total_keys = num_txns + num_effects;
71
72        self.metrics
73            .key_value_store_num_fetches_latency_ms
74            .with_label_values(&[self.store_name, "tx"])
75            .observe(elapsed.as_millis() as f64);
76        self.metrics
77            .key_value_store_num_fetches_batch_size
78            .with_label_values(&[self.store_name, "tx"])
79            .observe(total_keys as f64);
80
81        if let Ok((transactions, effects)) = &res {
82            let txns_not_found = transactions.iter().filter(|v| v.is_none()).count() as u64;
83            let effects_not_found = effects.iter().filter(|v| v.is_none()).count() as u64;
84
85            if num_txns > 0 {
86                self.metrics
87                    .key_value_store_num_fetches_success
88                    .with_label_values(&[self.store_name, "tx"])
89                    .inc_by(num_txns);
90            }
91            if num_effects > 0 {
92                self.metrics
93                    .key_value_store_num_fetches_success
94                    .with_label_values(&[self.store_name, "fx"])
95                    .inc_by(num_effects);
96            }
97
98            if txns_not_found > 0 {
99                self.metrics
100                    .key_value_store_num_fetches_not_found
101                    .with_label_values(&[self.store_name, "tx"])
102                    .inc_by(txns_not_found);
103            }
104            if effects_not_found > 0 {
105                self.metrics
106                    .key_value_store_num_fetches_not_found
107                    .with_label_values(&[self.store_name, "fx"])
108                    .inc_by(effects_not_found);
109            }
110        } else {
111            self.metrics
112                .key_value_store_num_fetches_error
113                .with_label_values(&[self.store_name, "tx"])
114                .inc_by(num_txns);
115            self.metrics
116                .key_value_store_num_fetches_error
117                .with_label_values(&[self.store_name, "fx"])
118                .inc_by(num_effects);
119        }
120
121        res
122    }
123
124    pub async fn multi_get_checkpoints(
125        &self,
126        checkpoint_summaries: &[CheckpointSequenceNumber],
127        checkpoint_contents: &[CheckpointSequenceNumber],
128        checkpoint_summaries_by_digest: &[CheckpointDigest],
129    ) -> IotaResult<(
130        Vec<Option<CertifiedCheckpointSummary>>,
131        Vec<Option<CheckpointContents>>,
132        Vec<Option<CertifiedCheckpointSummary>>,
133    )> {
134        let start = Instant::now();
135        let res = self
136            .inner
137            .multi_get_checkpoints(
138                checkpoint_summaries,
139                checkpoint_contents,
140                checkpoint_summaries_by_digest,
141            )
142            .await;
143        let elapsed = start.elapsed();
144
145        let num_summaries =
146            checkpoint_summaries.len() as u64 + checkpoint_summaries_by_digest.len() as u64;
147        let num_contents = checkpoint_contents.len() as u64;
148
149        self.metrics
150            .key_value_store_num_fetches_latency_ms
151            .with_label_values(&[self.store_name, "checkpoint"])
152            .observe(elapsed.as_millis() as f64);
153        self.metrics
154            .key_value_store_num_fetches_batch_size
155            .with_label_values(&[self.store_name, "checkpoint_summary"])
156            .observe(num_summaries as f64);
157        self.metrics
158            .key_value_store_num_fetches_batch_size
159            .with_label_values(&[self.store_name, "checkpoint_content"])
160            .observe(num_contents as f64);
161
162        if let Ok((summaries, contents, summaries_by_digest)) = &res {
163            let summaries_not_found = summaries.iter().filter(|v| v.is_none()).count() as u64
164                + summaries_by_digest.iter().filter(|v| v.is_none()).count() as u64;
165            let contents_not_found = contents.iter().filter(|v| v.is_none()).count() as u64;
166
167            if num_summaries > 0 {
168                self.metrics
169                    .key_value_store_num_fetches_success
170                    .with_label_values(&[self.store_name, "ckpt_summary"])
171                    .inc_by(num_summaries);
172            }
173            if num_contents > 0 {
174                self.metrics
175                    .key_value_store_num_fetches_success
176                    .with_label_values(&[self.store_name, "ckpt_contents"])
177                    .inc_by(num_contents);
178            }
179
180            if summaries_not_found > 0 {
181                self.metrics
182                    .key_value_store_num_fetches_not_found
183                    .with_label_values(&[self.store_name, "ckpt_summary"])
184                    .inc_by(summaries_not_found);
185            }
186            if contents_not_found > 0 {
187                self.metrics
188                    .key_value_store_num_fetches_not_found
189                    .with_label_values(&[self.store_name, "ckpt_contents"])
190                    .inc_by(contents_not_found);
191            }
192        } else {
193            self.metrics
194                .key_value_store_num_fetches_error
195                .with_label_values(&[self.store_name, "ckpt_summary"])
196                .inc_by(num_summaries);
197            self.metrics
198                .key_value_store_num_fetches_error
199                .with_label_values(&[self.store_name, "ckpt_contents"])
200                .inc_by(num_contents);
201        }
202
203        res
204    }
205
206    pub async fn multi_get_checkpoints_summaries(
207        &self,
208        keys: &[CheckpointSequenceNumber],
209    ) -> IotaResult<Vec<Option<CertifiedCheckpointSummary>>> {
210        self.multi_get_checkpoints(keys, &[], &[])
211            .await
212            .map(|(summaries, _, _)| summaries)
213    }
214
215    pub async fn multi_get_checkpoints_contents(
216        &self,
217        keys: &[CheckpointSequenceNumber],
218    ) -> IotaResult<Vec<Option<CheckpointContents>>> {
219        self.multi_get_checkpoints(&[], keys, &[])
220            .await
221            .map(|(_, contents, _)| contents)
222    }
223
224    pub async fn multi_get_checkpoints_summaries_by_digest(
225        &self,
226        keys: &[CheckpointDigest],
227    ) -> IotaResult<Vec<Option<CertifiedCheckpointSummary>>> {
228        self.multi_get_checkpoints(&[], &[], keys)
229            .await
230            .map(|(_, _, summaries)| summaries)
231    }
232
233    pub async fn multi_get_tx(
234        &self,
235        keys: &[TransactionDigest],
236    ) -> IotaResult<Vec<Option<TransactionEnvelope>>> {
237        self.multi_get(keys, &[]).await.map(|(txns, _)| txns)
238    }
239
240    pub async fn multi_get_fx_by_tx_digest(
241        &self,
242        keys: &[TransactionDigest],
243    ) -> IotaResult<Vec<Option<TransactionEffects>>> {
244        self.multi_get(&[], keys).await.map(|(_, fx)| fx)
245    }
246
247    /// Convenience method for fetching single digest, and returning an error if
248    /// it's not found. Prefer using multi_get_tx whenever possible.
249    pub async fn get_tx(&self, digest: TransactionDigest) -> IotaResult<TransactionEnvelope> {
250        self.multi_get_tx(&[digest])
251            .await?
252            .into_iter()
253            .next()
254            .flatten()
255            .ok_or(IotaError::TransactionNotFound { digest })
256    }
257
258    /// Convenience method for fetching single digest, and returning an error if
259    /// it's not found. Prefer using multi_get_fx_by_tx_digest whenever
260    /// possible.
261    pub async fn get_fx_by_tx_digest(
262        &self,
263        digest: TransactionDigest,
264    ) -> IotaResult<TransactionEffects> {
265        self.multi_get_fx_by_tx_digest(&[digest])
266            .await?
267            .into_iter()
268            .next()
269            .flatten()
270            .ok_or(IotaError::TransactionNotFound { digest })
271    }
272
273    /// Convenience method for fetching single checkpoint, and returning an
274    /// error if it's not found. Prefer using
275    /// multi_get_checkpoints_summaries whenever possible.
276    pub async fn get_checkpoint_summary(
277        &self,
278        checkpoint: CheckpointSequenceNumber,
279    ) -> IotaResult<CertifiedCheckpointSummary> {
280        self.multi_get_checkpoints_summaries(&[checkpoint])
281            .await?
282            .into_iter()
283            .next()
284            .flatten()
285            .ok_or(IotaError::UserInput {
286                error: UserInputError::VerifiedCheckpointNotFound(checkpoint),
287            })
288    }
289
290    /// Convenience method for fetching single checkpoint, and returning an
291    /// error if it's not found. Prefer using multi_get_checkpoints_contents
292    /// whenever possible.
293    pub async fn get_checkpoint_contents(
294        &self,
295        checkpoint: CheckpointSequenceNumber,
296    ) -> IotaResult<CheckpointContents> {
297        self.multi_get_checkpoints_contents(&[checkpoint])
298            .await?
299            .into_iter()
300            .next()
301            .flatten()
302            .ok_or(IotaError::UserInput {
303                error: UserInputError::VerifiedCheckpointNotFound(checkpoint),
304            })
305    }
306
307    /// Convenience method for fetching single checkpoint, and returning an
308    /// error if it's not found. Prefer using
309    /// multi_get_checkpoints_summaries_by_digest whenever possible.
310    pub async fn get_checkpoint_summary_by_digest(
311        &self,
312        digest: CheckpointDigest,
313    ) -> IotaResult<CertifiedCheckpointSummary> {
314        self.multi_get_checkpoints_summaries_by_digest(&[digest])
315            .await?
316            .into_iter()
317            .next()
318            .flatten()
319            .ok_or(IotaError::UserInput {
320                error: UserInputError::VerifiedCheckpointDigestNotFound(format!("{digest}")),
321            })
322    }
323
324    pub async fn get_transaction_perpetual_checkpoint(
325        &self,
326        digest: TransactionDigest,
327    ) -> IotaResult<Option<CheckpointSequenceNumber>> {
328        self.inner
329            .get_transaction_perpetual_checkpoint(digest)
330            .await
331    }
332
333    pub async fn get_object(
334        &self,
335        object_id: ObjectId,
336        version: VersionNumber,
337    ) -> IotaResult<Option<Object>> {
338        self.inner.get_object(object_id, version).await
339    }
340
341    pub async fn multi_get_objects(
342        &self,
343        object_keys: &[ObjectKey],
344    ) -> IotaResult<Vec<Option<Object>>> {
345        self.inner.multi_get_objects(object_keys).await
346    }
347
348    pub async fn multi_get_transactions_perpetual_checkpoints(
349        &self,
350        digests: &[TransactionDigest],
351    ) -> IotaResult<Vec<Option<CheckpointSequenceNumber>>> {
352        self.inner
353            .multi_get_transactions_perpetual_checkpoints(digests)
354            .await
355    }
356
357    pub async fn multi_get_events_by_tx_digests(
358        &self,
359        digests: &[TransactionDigest],
360    ) -> IotaResult<Vec<Option<TransactionEvents>>> {
361        self.inner.multi_get_events_by_tx_digests(digests).await
362    }
363}
364
365/// Immutable key/value store trait for storing/retrieving transactions,
366/// effects, and events. Only defines multi_get/multi_put methods to discourage
367/// single key/value operations.
368#[async_trait]
369pub trait TransactionKeyValueStoreTrait {
370    /// Generic multi_get, allows implementors to get heterogenous values with a
371    /// single round trip.
372    async fn multi_get(
373        &self,
374        transaction_keys: &[TransactionDigest],
375        effects_keys: &[TransactionDigest],
376    ) -> IotaResult<KVStoreTransactionData>;
377
378    /// Generic multi_get to allow implementors to get heterogenous values with
379    /// a single round trip.
380    async fn multi_get_checkpoints(
381        &self,
382        checkpoint_summaries: &[CheckpointSequenceNumber],
383        checkpoint_contents: &[CheckpointSequenceNumber],
384        checkpoint_summaries_by_digest: &[CheckpointDigest],
385    ) -> IotaResult<KVStoreCheckpointData>;
386
387    async fn get_transaction_perpetual_checkpoint(
388        &self,
389        digest: TransactionDigest,
390    ) -> IotaResult<Option<CheckpointSequenceNumber>>;
391
392    async fn get_object(&self, object_id: ObjectId, version: Version)
393    -> IotaResult<Option<Object>>;
394
395    async fn multi_get_objects(&self, object_keys: &[ObjectKey])
396    -> IotaResult<Vec<Option<Object>>>;
397
398    async fn multi_get_transactions_perpetual_checkpoints(
399        &self,
400        digests: &[TransactionDigest],
401    ) -> IotaResult<Vec<Option<CheckpointSequenceNumber>>>;
402
403    async fn multi_get_events_by_tx_digests(
404        &self,
405        digests: &[TransactionDigest],
406    ) -> IotaResult<Vec<Option<TransactionEvents>>>;
407}
408
409/// A TransactionKeyValueStoreTrait that falls back to a secondary store for any
410/// key for which the primary store returns None.
411///
412/// Will be used to check the local rocksdb store, before falling back to a
413/// remote scalable store.
414pub struct FallbackTransactionKVStore {
415    primary: TransactionKeyValueStore,
416    fallback: TransactionKeyValueStore,
417}
418
419impl FallbackTransactionKVStore {
420    pub fn new_kv(
421        primary: TransactionKeyValueStore,
422        fallback: TransactionKeyValueStore,
423        metrics: Arc<KeyValueStoreMetrics>,
424        label: &'static str,
425    ) -> TransactionKeyValueStore {
426        let store = Arc::new(Self { primary, fallback });
427        TransactionKeyValueStore::new(label, metrics, store)
428    }
429}
430
431#[async_trait]
432impl TransactionKeyValueStoreTrait for FallbackTransactionKVStore {
433    #[instrument(level = "trace", skip_all)]
434    async fn multi_get(
435        &self,
436        transaction_keys: &[TransactionDigest],
437        effects_keys: &[TransactionDigest],
438    ) -> IotaResult<KVStoreTransactionData> {
439        let (mut transactions, mut effects) = self
440            .primary
441            .multi_get(transaction_keys, effects_keys)
442            .await?;
443
444        let (fallback_transaction_keys, indices_transactions) =
445            find_fallback(&transactions, transaction_keys);
446        let (fallback_effects_keys, indices_effects) = find_fallback(&effects, effects_keys);
447
448        if fallback_transaction_keys.is_empty() && fallback_effects_keys.is_empty() {
449            return Ok((transactions, effects));
450        }
451
452        let (fallback_transactions, fallback_effects) = self
453            .fallback
454            .multi_get(&fallback_transaction_keys, &fallback_effects_keys)
455            .await?;
456
457        merge_res(
458            &mut transactions,
459            fallback_transactions,
460            &indices_transactions,
461        );
462        merge_res(&mut effects, fallback_effects, &indices_effects);
463
464        Ok((transactions, effects))
465    }
466
467    #[instrument(level = "trace", skip_all)]
468    async fn multi_get_checkpoints(
469        &self,
470        checkpoint_summaries: &[CheckpointSequenceNumber],
471        checkpoint_contents: &[CheckpointSequenceNumber],
472        checkpoint_summaries_by_digest: &[CheckpointDigest],
473    ) -> IotaResult<(
474        Vec<Option<CertifiedCheckpointSummary>>,
475        Vec<Option<CheckpointContents>>,
476        Vec<Option<CertifiedCheckpointSummary>>,
477    )> {
478        let (mut summaries, mut contents, mut summaries_by_digest) = self
479            .primary
480            .multi_get_checkpoints(
481                checkpoint_summaries,
482                checkpoint_contents,
483                checkpoint_summaries_by_digest,
484            )
485            .await?;
486
487        let (fallback_summaries, indices_summaries) =
488            find_fallback(&summaries, checkpoint_summaries);
489        let (fallback_contents, indices_contents) = find_fallback(&contents, checkpoint_contents);
490        let (fallback_summaries_by_digest, indices_summaries_by_digest) =
491            find_fallback(&summaries_by_digest, checkpoint_summaries_by_digest);
492
493        if fallback_summaries.is_empty()
494            && fallback_contents.is_empty()
495            && fallback_summaries_by_digest.is_empty()
496        {
497            return Ok((summaries, contents, summaries_by_digest));
498        }
499
500        let (fallback_summaries, fallback_contents, fallback_summaries_by_digest) = self
501            .fallback
502            .multi_get_checkpoints(
503                &fallback_summaries,
504                &fallback_contents,
505                &fallback_summaries_by_digest,
506            )
507            .await?;
508
509        merge_res(&mut summaries, fallback_summaries, &indices_summaries);
510        merge_res(&mut contents, fallback_contents, &indices_contents);
511        merge_res(
512            &mut summaries_by_digest,
513            fallback_summaries_by_digest,
514            &indices_summaries_by_digest,
515        );
516
517        Ok((summaries, contents, summaries_by_digest))
518    }
519
520    #[instrument(level = "trace", skip_all)]
521    async fn get_transaction_perpetual_checkpoint(
522        &self,
523        digest: TransactionDigest,
524    ) -> IotaResult<Option<CheckpointSequenceNumber>> {
525        let mut res = self
526            .primary
527            .get_transaction_perpetual_checkpoint(digest)
528            .await?;
529        if res.is_none() {
530            res = self
531                .fallback
532                .get_transaction_perpetual_checkpoint(digest)
533                .await?;
534        }
535        Ok(res)
536    }
537
538    #[instrument(level = "trace", skip_all)]
539    async fn get_object(
540        &self,
541        object_id: ObjectId,
542        version: Version,
543    ) -> IotaResult<Option<Object>> {
544        let mut res = self.primary.get_object(object_id, version).await?;
545        if res.is_none() {
546            res = self.fallback.get_object(object_id, version).await?;
547        }
548        Ok(res)
549    }
550
551    #[instrument(level = "trace", skip_all)]
552    async fn multi_get_objects(
553        &self,
554        object_keys: &[ObjectKey],
555    ) -> IotaResult<Vec<Option<Object>>> {
556        let mut res = self.primary.multi_get_objects(object_keys).await?;
557
558        let (fallback, indices) = find_fallback(&res, object_keys);
559
560        if fallback.is_empty() {
561            return Ok(res);
562        }
563
564        let secondary_res = self.fallback.multi_get_objects(&fallback).await?;
565
566        merge_res(&mut res, secondary_res, &indices);
567
568        Ok(res)
569    }
570
571    #[instrument(level = "trace", skip_all)]
572    async fn multi_get_transactions_perpetual_checkpoints(
573        &self,
574        digests: &[TransactionDigest],
575    ) -> IotaResult<Vec<Option<CheckpointSequenceNumber>>> {
576        let mut res = self
577            .primary
578            .multi_get_transactions_perpetual_checkpoints(digests)
579            .await?;
580
581        let (fallback, indices) = find_fallback(&res, digests);
582
583        if fallback.is_empty() {
584            return Ok(res);
585        }
586
587        let secondary_res = self
588            .fallback
589            .multi_get_transactions_perpetual_checkpoints(&fallback)
590            .await?;
591
592        merge_res(&mut res, secondary_res, &indices);
593
594        Ok(res)
595    }
596
597    #[instrument(level = "trace", skip_all)]
598    async fn multi_get_events_by_tx_digests(
599        &self,
600        digests: &[TransactionDigest],
601    ) -> IotaResult<Vec<Option<TransactionEvents>>> {
602        let mut res = self.primary.multi_get_events_by_tx_digests(digests).await?;
603        let (fallback, indices) = find_fallback(&res, digests);
604        if fallback.is_empty() {
605            return Ok(res);
606        }
607        let secondary_res = self
608            .fallback
609            .multi_get_events_by_tx_digests(&fallback)
610            .await?;
611        merge_res(&mut res, secondary_res, &indices);
612        Ok(res)
613    }
614}
615
616fn find_fallback<T, K: Clone>(values: &[Option<T>], keys: &[K]) -> (Vec<K>, Vec<usize>) {
617    let num_nones = values.iter().filter(|v| v.is_none()).count();
618    let mut fallback_keys = Vec::with_capacity(num_nones);
619    let mut fallback_indices = Vec::with_capacity(num_nones);
620    for (i, value) in values.iter().enumerate() {
621        if value.is_none() {
622            fallback_keys.push(keys[i].clone());
623            fallback_indices.push(i);
624        }
625    }
626    (fallback_keys, fallback_indices)
627}
628
629fn merge_res<T>(values: &mut [Option<T>], fallback_values: Vec<Option<T>>, indices: &[usize]) {
630    for (&index, fallback_value) in indices.iter().zip(fallback_values) {
631        values[index] = fallback_value;
632    }
633}