1use std::{str::FromStr, sync::Arc, time::Duration};
6
7use anyhow;
8use async_trait::async_trait;
9use bytes::Bytes;
10use futures::stream::{self, StreamExt};
11use iota_sdk_types::{
12 Address, CheckpointDigest, ObjectId, TransactionDigest, TransactionEffects, TransactionEvents,
13 Version, checkpoint::CheckpointContents,
14};
15use iota_types::{
16 effects::TransactionEffectsAPI,
17 error::{IotaError, IotaResult},
18 messages_checkpoint::{CertifiedCheckpointSummary, CheckpointSequenceNumber},
19 object::Object,
20 storage::ObjectKey,
21 transaction::TransactionEnvelope,
22};
23use moka::sync::{Cache as MokaCache, CacheBuilder as MokaCacheBuilder};
24use reqwest::{
25 Client, Url,
26 header::{CONTENT_LENGTH, HeaderValue},
27};
28use serde::{Deserialize, Serialize};
29use tap::TapFallible;
30use tracing::{error, info, instrument, trace, warn};
31
32use crate::{
33 key_value_store::{
34 KVStoreTransactionData, TransactionKeyValueStore, TransactionKeyValueStoreTrait,
35 },
36 key_value_store_metrics::KeyValueStoreMetrics,
37};
38
39pub struct HttpKVStore {
40 base_url: Url,
41 client: Client,
42 cache: MokaCache<Url, Bytes>,
43 metrics: Arc<KeyValueStoreMetrics>,
44}
45
46pub fn encode_digest<T: AsRef<[u8]>>(digest: &T) -> String {
47 base64_url::encode(digest)
48}
49
50#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
52pub enum TaggedKey {
53 CheckpointSequenceNumber(CheckpointSequenceNumber),
54}
55
56pub fn encoded_tagged_key(key: &TaggedKey) -> String {
57 let bytes = bcs::to_bytes(key).expect("failed to serialize key");
58 base64_url::encode(&bytes)
59}
60
61pub fn encode_object_key(object_key: &ObjectKey) -> String {
62 let bytes = bcs::to_bytes(object_key).expect("failed to serialize object key");
63 base64_url::encode(&bytes)
64}
65
66trait IntoIotaResult<T> {
67 fn into_iota_result(self) -> IotaResult<T>;
68}
69
70impl<T, E> IntoIotaResult<T> for Result<T, E>
71where
72 E: std::error::Error,
73{
74 fn into_iota_result(self) -> IotaResult<T> {
75 self.map_err(|e| IotaError::Storage(e.to_string()))
76 }
77}
78
79#[derive(
82 Clone,
83 Copy,
84 Debug,
85 PartialEq,
86 Eq,
87 PartialOrd,
88 Ord,
89 Hash,
90 Deserialize,
91 strum::EnumString,
92 strum::Display,
93)]
94pub enum ItemType {
95 #[strum(serialize = "tx")]
96 #[serde(rename = "tx")]
97 Transaction,
98 #[strum(serialize = "fx")]
99 #[serde(rename = "fx")]
100 TransactionEffects,
101 #[strum(serialize = "cc")]
102 #[serde(rename = "cc")]
103 CheckpointContents,
104 #[strum(serialize = "cs")]
105 #[serde(rename = "cs")]
106 CheckpointSummary,
107 #[strum(serialize = "tx2c")]
108 #[serde(rename = "tx2c")]
109 TransactionToCheckpoint,
110 #[strum(serialize = "ob")]
111 #[serde(rename = "ob")]
112 Object,
113 #[strum(serialize = "evtx")]
114 #[serde(rename = "evtx")]
115 EventTransactionDigest,
116 #[strum(serialize = "txa")]
117 #[serde(rename = "txa")]
118 TransactionDigestsByAddress,
119}
120
121#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
122pub enum Key {
123 Transaction(TransactionDigest),
124 TransactionEffects(TransactionDigest),
125 CheckpointContents(CheckpointSequenceNumber),
126 CheckpointSummary(CheckpointSequenceNumber),
127 CheckpointSummaryByDigest(CheckpointDigest),
128 TransactionToCheckpoint(TransactionDigest),
129 ObjectKey(ObjectKey),
130 EventsByTransactionDigest(TransactionDigest),
131 TransactionDigestsByAddress(Address),
132}
133
134impl Key {
135 pub fn new(item_type: &str, encoded_key: &str) -> anyhow::Result<Self> {
155 let item_type =
156 ItemType::from_str(item_type).map_err(|e| anyhow::anyhow!("invalid item type: {e}"))?;
157 let decoded_key = base64_url::decode(encoded_key)
158 .map_err(|err| anyhow::anyhow!("invalid base64 url string: {err}"))?;
159
160 match item_type {
161 ItemType::Transaction => Ok(Key::Transaction(TransactionDigest::from_bytes(
162 decoded_key.as_slice(),
163 )?)),
164 ItemType::TransactionEffects => Ok(Key::TransactionEffects(
165 TransactionDigest::from_bytes(decoded_key.as_slice())?,
166 )),
167 ItemType::CheckpointContents => {
168 let tagged_key = bcs::from_bytes(&decoded_key).map_err(|err| {
169 anyhow::anyhow!("failed to deserialize checkpoint sequence number: {err}")
170 })?;
171 match tagged_key {
172 TaggedKey::CheckpointSequenceNumber(seq) => Ok(Key::CheckpointContents(seq)),
173 }
174 }
175 ItemType::CheckpointSummary => {
176 match CheckpointDigest::from_bytes(decoded_key.clone()) {
178 Err(_) => {
179 let tagged_key = bcs::from_bytes(&decoded_key).map_err(|err| {
180 anyhow::anyhow!(
181 "failed to deserialize checkpoint sequence number: {err}"
182 )
183 })?;
184 match tagged_key {
185 TaggedKey::CheckpointSequenceNumber(seq) => {
186 Ok(Key::CheckpointSummary(seq))
187 }
188 }
189 }
190 Ok(cs_digest) => Ok(Key::CheckpointSummaryByDigest(cs_digest)),
191 }
192 }
193 ItemType::TransactionToCheckpoint => Ok(Key::TransactionToCheckpoint(
194 TransactionDigest::from_bytes(decoded_key.as_slice())?,
195 )),
196 ItemType::Object => {
197 let object_key: ObjectKey = bcs::from_bytes(&decoded_key)
198 .map_err(|err| anyhow::anyhow!("failed to deserialize object key: {err}"))?;
199
200 Ok(Key::ObjectKey(ObjectKey(object_key.0, object_key.1)))
201 }
202 ItemType::EventTransactionDigest => Ok(Key::EventsByTransactionDigest(
203 TransactionDigest::from_bytes(decoded_key.as_slice())?,
204 )),
205 ItemType::TransactionDigestsByAddress => Ok(Key::TransactionDigestsByAddress(
206 Address::from_bytes(decoded_key.as_slice())?,
207 )),
208 }
209 }
210
211 pub fn item_type(&self) -> ItemType {
230 match self {
231 Key::Transaction(_) => ItemType::Transaction,
232 Key::TransactionEffects(_) => ItemType::TransactionEffects,
233 Key::CheckpointContents(_) => ItemType::CheckpointContents,
234 Key::CheckpointSummary(_) | Key::CheckpointSummaryByDigest(_) => {
235 ItemType::CheckpointSummary
236 }
237 Key::TransactionToCheckpoint(_) => ItemType::TransactionToCheckpoint,
238 Key::ObjectKey(_) => ItemType::Object,
239 Key::EventsByTransactionDigest(_) => ItemType::EventTransactionDigest,
240 Key::TransactionDigestsByAddress(_) => ItemType::TransactionDigestsByAddress,
241 }
242 }
243
244 pub fn to_path_elements(&self) -> (ItemType, String) {
275 let encoded_key_digest = match self {
276 Key::Transaction(digest) => encode_digest(digest),
277 Key::TransactionEffects(digest) => encode_digest(digest),
278 Key::CheckpointContents(seq) => {
279 encoded_tagged_key(&TaggedKey::CheckpointSequenceNumber(*seq))
280 }
281 Key::CheckpointSummary(seq) => {
282 encoded_tagged_key(&TaggedKey::CheckpointSequenceNumber(*seq))
283 }
284 Key::CheckpointSummaryByDigest(digest) => encode_digest(digest),
285 Key::TransactionToCheckpoint(digest) => encode_digest(digest),
286 Key::ObjectKey(object_key) => encode_object_key(object_key),
287 Key::EventsByTransactionDigest(digest) => encode_digest(digest),
288 Key::TransactionDigestsByAddress(address) => encode_digest(address),
291 };
292
293 (self.item_type(), encoded_key_digest)
294 }
295}
296
297#[derive(Clone, Debug)]
298enum Value {
299 Tx(Box<TransactionEnvelope>),
300 Fx(Box<TransactionEffects>),
301 Events(Box<TransactionEvents>),
302 CheckpointContents(Box<CheckpointContents>),
303 CheckpointSummary(Box<CertifiedCheckpointSummary>),
304 TxToCheckpoint(CheckpointSequenceNumber),
305}
306
307impl HttpKVStore {
308 pub fn new_kv(
309 base_url: &str,
310 cache_size: u64,
311 metrics: Arc<KeyValueStoreMetrics>,
312 ) -> IotaResult<TransactionKeyValueStore> {
313 let inner = Arc::new(Self::new(base_url, cache_size, metrics.clone())?);
314 Ok(TransactionKeyValueStore::new("http", metrics, inner))
315 }
316
317 pub fn new(
318 base_url: &str,
319 cache_size: u64,
320 metrics: Arc<KeyValueStoreMetrics>,
321 ) -> IotaResult<Self> {
322 info!("creating HttpKVStore with base_url: {}", base_url);
323
324 let client = Client::builder().http2_prior_knowledge().build().unwrap();
325
326 let base_url = if base_url.ends_with('/') {
327 base_url.to_string()
328 } else {
329 format!("{base_url}/")
330 };
331
332 let base_url = Url::parse(&base_url).into_iota_result()?;
333
334 let cache = MokaCacheBuilder::new(cache_size)
335 .time_to_idle(Duration::from_secs(600))
336 .build();
337
338 Ok(Self {
339 base_url,
340 client,
341 cache,
342 metrics,
343 })
344 }
345
346 fn get_url(&self, key: &Key) -> IotaResult<Url> {
347 let (item_type, digest) = key.to_path_elements();
348 let joined = self
349 .base_url
350 .join(&format!("{item_type}/{digest}"))
351 .into_iota_result()?;
352 Url::from_str(joined.as_str()).into_iota_result()
353 }
354
355 async fn multi_fetch(&self, uris: Vec<Key>) -> Vec<IotaResult<Option<Bytes>>> {
356 let uris_vec = uris.to_vec();
357 let fetches = stream::iter(uris_vec.into_iter().map(|url| self.fetch(url)));
358 fetches.buffered(uris.len()).collect::<Vec<_>>().await
359 }
360
361 async fn fetch(&self, key: Key) -> IotaResult<Option<Bytes>> {
362 let url = self.get_url(&key)?;
363
364 trace!("fetching url: {}", url);
365
366 if let Some(res) = self.cache.get(&url) {
367 trace!("found cached data for url: {}, len: {:?}", url, res.len());
368 self.metrics
369 .key_value_store_num_fetches_success
370 .with_label_values(&["http_cache", "url"])
371 .inc();
372 return Ok(Some(res));
373 }
374
375 self.metrics
376 .key_value_store_num_fetches_not_found
377 .with_label_values(&["http_cache", "url"])
378 .inc();
379
380 let resp = self
381 .client
382 .get(url.clone())
383 .send()
384 .await
385 .into_iota_result()?;
386 trace!(
387 "got response {} for url: {}, len: {:?}",
388 url,
389 resp.status(),
390 resp.headers()
391 .get(CONTENT_LENGTH)
392 .unwrap_or(&HeaderValue::from_static("0"))
393 );
394 if resp.status().is_success() {
396 let bytes = resp.bytes().await.into_iota_result()?;
397 self.cache.insert(url, bytes.clone());
398
399 Ok(Some(bytes))
400 } else {
401 Ok(None)
402 }
403 }
404}
405
406fn deser<K, T>(key: &K, bytes: &[u8]) -> Option<T>
407where
408 K: std::fmt::Debug,
409 T: for<'de> Deserialize<'de>,
410{
411 bcs::from_bytes(bytes)
412 .tap_err(|e| warn!("Error deserializing data for key {:?}: {:?}", key, e))
413 .ok()
414}
415
416fn map_fetch<'a, K>(fetch: (&'a IotaResult<Option<Bytes>>, &'a K)) -> Option<(&'a Bytes, &'a K)>
417where
418 K: std::fmt::Debug,
419{
420 let (fetch, key) = fetch;
421 match fetch {
422 Ok(Some(bytes)) => Some((bytes, key)),
423 Ok(None) => None,
424 Err(err) => {
425 warn!("Error fetching key: {:?}, error: {:?}", key, err);
426 None
427 }
428 }
429}
430
431fn multi_split_slice<'a, T>(slice: &'a [T], lengths: &'a [usize]) -> Vec<&'a [T]> {
432 let mut start = 0;
433 lengths
434 .iter()
435 .map(|length| {
436 let end = start + length;
437 let result = &slice[start..end];
438 start = end;
439 result
440 })
441 .collect()
442}
443
444fn deser_check_digest<T, D>(
445 digest: &D,
446 bytes: &Bytes,
447 get_expected_digest: impl FnOnce(&T) -> D,
448) -> Option<T>
449where
450 D: std::fmt::Debug + PartialEq,
451 T: for<'de> Deserialize<'de>,
452{
453 deser(digest, bytes).and_then(|o: T| {
454 let expected_digest = get_expected_digest(&o);
455 if expected_digest == *digest {
456 Some(o)
457 } else {
458 error!(
459 "Digest mismatch - expected: {:?}, got: {:?}",
460 digest, expected_digest,
461 );
462 None
463 }
464 })
465}
466
467#[async_trait]
468impl TransactionKeyValueStoreTrait for HttpKVStore {
469 #[instrument(level = "trace", skip_all)]
470 async fn multi_get(
471 &self,
472 transaction_keys: &[TransactionDigest],
473 effects_keys: &[TransactionDigest],
474 ) -> IotaResult<KVStoreTransactionData> {
475 let num_txns = transaction_keys.len();
476 let num_effects = effects_keys.len();
477
478 let keys = transaction_keys
479 .iter()
480 .map(|tx| Key::Transaction(*tx))
481 .chain(effects_keys.iter().map(|fx| Key::TransactionEffects(*fx)))
482 .collect::<Vec<_>>();
483
484 let fetches = self.multi_fetch(keys).await;
485 let txn_slice = fetches[..num_txns].to_vec();
486 let fx_slice = fetches[num_txns..num_txns + num_effects].to_vec();
487
488 let txn_results = txn_slice
489 .iter()
490 .take(num_txns)
491 .zip(transaction_keys.iter())
492 .map(map_fetch)
493 .map(|maybe_bytes| {
494 maybe_bytes.and_then(|(bytes, digest)| {
495 deser_check_digest(digest, bytes, |tx: &TransactionEnvelope| *tx.digest())
496 })
497 })
498 .collect::<Vec<_>>();
499
500 let fx_results = fx_slice
501 .iter()
502 .take(num_effects)
503 .zip(effects_keys.iter())
504 .map(map_fetch)
505 .map(|maybe_bytes| {
506 maybe_bytes.and_then(|(bytes, digest)| {
507 deser_check_digest(digest, bytes, |fx: &TransactionEffects| {
508 *fx.transaction_digest()
509 })
510 })
511 })
512 .collect::<Vec<_>>();
513
514 Ok((txn_results, fx_results))
515 }
516
517 #[instrument(level = "trace", skip_all)]
518 async fn multi_get_checkpoints(
519 &self,
520 checkpoint_summaries: &[CheckpointSequenceNumber],
521 checkpoint_contents: &[CheckpointSequenceNumber],
522 checkpoint_summaries_by_digest: &[CheckpointDigest],
523 ) -> IotaResult<(
524 Vec<Option<CertifiedCheckpointSummary>>,
525 Vec<Option<CheckpointContents>>,
526 Vec<Option<CertifiedCheckpointSummary>>,
527 )> {
528 let keys = checkpoint_summaries
529 .iter()
530 .map(|cp| Key::CheckpointSummary(*cp))
531 .chain(
532 checkpoint_contents
533 .iter()
534 .map(|cp| Key::CheckpointContents(*cp)),
535 )
536 .chain(
537 checkpoint_summaries_by_digest
538 .iter()
539 .map(|cp| Key::CheckpointSummaryByDigest(*cp)),
540 )
541 .collect::<Vec<_>>();
542
543 let summaries_len = checkpoint_summaries.len();
544 let contents_len = checkpoint_contents.len();
545 let summaries_by_digest_len = checkpoint_summaries_by_digest.len();
546
547 let fetches = self.multi_fetch(keys).await;
548
549 let input_slices = [summaries_len, contents_len, summaries_by_digest_len];
550
551 let result_slices = multi_split_slice(&fetches, &input_slices);
552
553 let summaries_results = result_slices[0]
554 .iter()
555 .zip(checkpoint_summaries.iter())
556 .map(map_fetch)
557 .map(|maybe_bytes| {
558 maybe_bytes
559 .and_then(|(bytes, seq)| deser::<_, CertifiedCheckpointSummary>(seq, bytes))
560 })
561 .collect::<Vec<_>>();
562
563 let contents_results = result_slices[1]
564 .iter()
565 .zip(checkpoint_contents.iter())
566 .map(map_fetch)
567 .map(|maybe_bytes| {
568 maybe_bytes.and_then(|(bytes, seq)| deser::<_, CheckpointContents>(seq, bytes))
569 })
570 .collect::<Vec<_>>();
571
572 let summaries_by_digest_results = result_slices[2]
573 .iter()
574 .zip(checkpoint_summaries_by_digest.iter())
575 .map(map_fetch)
576 .map(|maybe_bytes| {
577 maybe_bytes.and_then(|(bytes, digest)| {
578 deser_check_digest(digest, bytes, |s: &CertifiedCheckpointSummary| *s.digest())
579 })
580 })
581 .collect::<Vec<_>>();
582
583 Ok((
584 summaries_results,
585 contents_results,
586 summaries_by_digest_results,
587 ))
588 }
589
590 #[instrument(level = "trace", skip_all)]
591 async fn get_transaction_perpetual_checkpoint(
592 &self,
593 digest: TransactionDigest,
594 ) -> IotaResult<Option<CheckpointSequenceNumber>> {
595 let key = Key::TransactionToCheckpoint(digest);
596 self.fetch(key).await.map(|maybe| {
597 maybe.and_then(|bytes| deser::<_, CheckpointSequenceNumber>(&key, bytes.as_ref()))
598 })
599 }
600
601 #[instrument(level = "trace", skip_all)]
602 async fn get_object(
603 &self,
604 object_id: ObjectId,
605 version: Version,
606 ) -> IotaResult<Option<Object>> {
607 let key = Key::ObjectKey(ObjectKey(object_id, version));
608 self.fetch(key)
609 .await
610 .map(|maybe| maybe.and_then(|bytes| deser::<_, Object>(&key, bytes.as_ref())))
611 }
612
613 #[instrument(level = "trace", skip_all)]
614 async fn multi_get_objects(
615 &self,
616 object_keys: &[ObjectKey],
617 ) -> IotaResult<Vec<Option<Object>>> {
618 let keys = object_keys
619 .iter()
620 .map(|key| Key::ObjectKey(*key))
621 .collect::<Vec<_>>();
622
623 let fetches = self.multi_fetch(keys).await;
624
625 let results = fetches
626 .iter()
627 .zip(object_keys.iter())
628 .map(map_fetch)
629 .map(|maybe_bytes| maybe_bytes.and_then(|(bytes, key)| deser::<_, Object>(&key, bytes)))
630 .collect::<Vec<_>>();
631
632 Ok(results)
633 }
634
635 #[instrument(level = "trace", skip_all)]
636 async fn multi_get_transactions_perpetual_checkpoints(
637 &self,
638 digests: &[TransactionDigest],
639 ) -> IotaResult<Vec<Option<CheckpointSequenceNumber>>> {
640 let keys = digests
641 .iter()
642 .map(|digest| Key::TransactionToCheckpoint(*digest))
643 .collect::<Vec<_>>();
644
645 let fetches = self.multi_fetch(keys).await;
646
647 let results = fetches
648 .iter()
649 .zip(digests.iter())
650 .map(map_fetch)
651 .map(|maybe_bytes| {
652 maybe_bytes
653 .and_then(|(bytes, key)| deser::<_, CheckpointSequenceNumber>(&key, bytes))
654 })
655 .collect::<Vec<_>>();
656
657 Ok(results)
658 }
659
660 #[instrument(level = "trace", skip_all)]
661 async fn multi_get_events_by_tx_digests(
662 &self,
663 digests: &[TransactionDigest],
664 ) -> IotaResult<Vec<Option<TransactionEvents>>> {
665 let keys = digests
666 .iter()
667 .map(|digest| Key::EventsByTransactionDigest(*digest))
668 .collect::<Vec<_>>();
669 Ok(self
670 .multi_fetch(keys)
671 .await
672 .iter()
673 .zip(digests.iter())
674 .map(map_fetch)
675 .map(|maybe_bytes| {
676 maybe_bytes.and_then(|(bytes, key)| deser::<_, TransactionEvents>(&key, bytes))
677 })
678 .collect::<Vec<_>>())
679 }
680}