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