1use 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 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 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 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 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 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 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#[async_trait]
369pub trait TransactionKeyValueStoreTrait {
370 async fn multi_get(
373 &self,
374 transaction_keys: &[TransactionDigest],
375 effects_keys: &[TransactionDigest],
376 ) -> IotaResult<KVStoreTransactionData>;
377
378 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
409pub 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}