1use std::ops::Bound;
13
14use move_core_types::language_storage::StructTag;
15use sui_consistent_store::Encode;
16use sui_consistent_store::reader::Reader;
17use sui_types::base_types::ObjectID;
18use sui_types::base_types::SuiAddress;
19use sui_types::messages_checkpoint::CheckpointSequenceNumber;
20use sui_types::storage::BalanceInfo;
21use sui_types::storage::BalanceIterator;
22use sui_types::storage::CoinInfo;
23use sui_types::storage::DynamicFieldIteratorItem;
24use sui_types::storage::DynamicFieldKey;
25use sui_types::storage::EpochInfo;
26use sui_types::storage::LedgerBitmapBucket;
27use sui_types::storage::LedgerBitmapBucketIter;
28use sui_types::storage::LedgerBitmapBucketIterator;
29use sui_types::storage::LedgerTxSeqDigest;
30use sui_types::storage::LedgerTxSeqDigestIterator;
31use sui_types::storage::OwnedObjectInfo;
32use sui_types::storage::RpcIndexes;
33
34type PackageVersionsIterator<'a> =
37 Box<dyn Iterator<Item = Result<(u64, ObjectID), TypedStoreError>> + 'a>;
38use sui_types::storage::error::Result as StorageResult;
39use typed_store_error::TypedStoreError;
40
41use crate::reader::RpcStoreReader;
42use crate::schema::type_filter::TypeFilter;
43
44enum BoundedBitmapIter<'d, K, V> {
45 Fwd(sui_consistent_store::Iter<'d, K, V>),
46 Rev(sui_consistent_store::RevIter<'d, K, V>),
47}
48
49impl<K, V> BoundedBitmapIter<'_, K, V> {
50 fn seek(&mut self, probe: impl AsRef<[u8]>) {
51 match self {
52 Self::Fwd(iter) => iter.seek(probe),
53 Self::Rev(iter) => iter.seek(probe),
54 }
55 }
56}
57
58struct TransactionBitmapBucketIter<'d> {
59 inner: BoundedBitmapIter<
60 'd,
61 crate::schema::transaction_bitmap::Key,
62 crate::schema::transaction_bitmap::Value,
63 >,
64 dimension_key: Vec<u8>,
65}
66
67impl Iterator for TransactionBitmapBucketIter<'_> {
68 type Item = Result<LedgerBitmapBucket, TypedStoreError>;
69
70 fn next(&mut self) -> Option<Self::Item> {
71 let row = match &mut self.inner {
72 BoundedBitmapIter::Fwd(iter) => iter.next(),
73 BoundedBitmapIter::Rev(iter) => iter.next(),
74 };
75 row.map(project_bitmap_row)
76 }
77}
78
79impl LedgerBitmapBucketIter for TransactionBitmapBucketIter<'_> {
80 fn seek_bucket(&mut self, bucket_id: u64) {
81 let probe = crate::schema::transaction_bitmap::Key {
82 dimension_key: self.dimension_key.clone(),
83 bucket: bucket_id,
84 }
85 .encode()
86 .expect("bitmap key encodes infallibly: raw bytes + u64 BE");
87 self.inner.seek(probe);
88 }
89}
90
91struct EventBitmapBucketIter<'d> {
92 inner:
93 BoundedBitmapIter<'d, crate::schema::event_bitmap::Key, crate::schema::event_bitmap::Value>,
94 dimension_key: Vec<u8>,
95}
96
97impl Iterator for EventBitmapBucketIter<'_> {
98 type Item = Result<LedgerBitmapBucket, TypedStoreError>;
99
100 fn next(&mut self) -> Option<Self::Item> {
101 let row = match &mut self.inner {
102 BoundedBitmapIter::Fwd(iter) => iter.next(),
103 BoundedBitmapIter::Rev(iter) => iter.next(),
104 };
105 row.map(project_event_bitmap_row)
106 }
107}
108
109impl LedgerBitmapBucketIter for EventBitmapBucketIter<'_> {
110 fn seek_bucket(&mut self, bucket_id: u64) {
111 let probe = crate::schema::event_bitmap::Key {
112 dimension_key: self.dimension_key.clone(),
113 bucket: bucket_id,
114 }
115 .encode()
116 .expect("bitmap key encodes infallibly: raw bytes + u64 BE");
117 self.inner.seek(probe);
118 }
119}
120
121fn to_typed_store_err(e: sui_consistent_store::error::Error) -> TypedStoreError {
122 TypedStoreError::RocksDBError(format!("{e:#}"))
123}
124
125fn first_object_of_type<R: Reader + Send + Sync>(
132 reader: &RpcStoreReader<R>,
133 struct_tag: move_core_types::language_storage::StructTag,
134) -> StorageResult<Option<ObjectID>> {
135 let filter = TypeFilter::Type(struct_tag);
136 let mut iter = reader
137 .schema()
138 .iter_objects_of_type(&filter)
139 .map_err(sui_types::storage::error::Error::custom)?;
140 match iter.next() {
141 Some(Ok((key, _value))) => Ok(Some(key.object_id)),
142 Some(Err(e)) => Err(sui_types::storage::error::Error::custom(e)),
143 None => Ok(None),
144 }
145}
146
147impl<R: Reader + Send + Sync> RpcIndexes for RpcStoreReader<R> {
148 fn get_epoch_info(
149 &self,
150 epoch: sui_types::committee::EpochId,
151 ) -> StorageResult<Option<EpochInfo>> {
152 self.schema()
153 .get_epoch(epoch)
154 .map_err(sui_types::storage::error::Error::custom)
155 }
156
157 fn owned_objects_iter(
158 &self,
159 owner: SuiAddress,
160 object_type: Option<StructTag>,
161 cursor: Option<OwnedObjectInfo>,
162 ) -> StorageResult<Box<dyn Iterator<Item = Result<OwnedObjectInfo, TypedStoreError>> + '_>>
163 {
164 use crate::schema::object_by_owner::{Key, OwnerKind};
165 let map = &self.schema().object_by_owner;
166 let kind = OwnerKind::AddressOwner(owner);
167 let from = cursor.map(|c| Key {
171 kind,
172 type_: c.object_type,
173 inverted_balance: c.balance.map(|b| !b),
174 object_id: c.object_id,
175 });
176 let iter = match (object_type, &from) {
177 (Some(struct_tag), Some(from)) => {
178 map.iter_prefix_from(&(kind, TypeFilter::Type(struct_tag)), from)
179 }
180 (Some(struct_tag), None) => map.iter_prefix(&(kind, TypeFilter::Type(struct_tag))),
181 (None, Some(from)) => map.iter_prefix_from(&kind, from),
182 (None, None) => map.iter_prefix(&kind),
183 }
184 .map_err(sui_types::storage::error::Error::custom)?;
185
186 let mapped = iter.map(move |row| {
187 let (key, value) = row.map_err(to_typed_store_err)?;
188 Ok(OwnedObjectInfo {
189 owner,
190 object_type: key.type_,
191 balance: key.inverted_balance.map(|b| !b),
192 object_id: key.object_id,
193 version: sui_types::base_types::SequenceNumber::from_u64(value.0),
194 })
195 });
196
197 Ok(Box::new(mapped))
198 }
199
200 fn dynamic_field_iter(
201 &self,
202 parent: ObjectID,
203 cursor: Option<DynamicFieldKey>,
204 ) -> StorageResult<Box<dyn Iterator<Item = DynamicFieldIteratorItem> + '_>> {
205 use crate::schema::object_by_owner::{Key, OwnerKind};
206 let map = &self.schema().object_by_owner;
213 let kind = OwnerKind::ObjectOwner(parent.into());
214 let from = cursor.map(|c| Key {
215 kind,
216 type_: c.object_type,
217 inverted_balance: None,
218 object_id: c.field_id,
219 });
220 let iter = match &from {
221 Some(from) => map.iter_prefix_from(&kind, from),
222 None => map.iter_prefix(&kind),
223 }
224 .map_err(sui_types::storage::error::Error::custom)?;
225
226 let mapped = iter.map(move |row| {
227 let (key, _value) = row.map_err(to_typed_store_err)?;
228 Ok(DynamicFieldKey {
229 parent,
230 field_id: key.object_id,
231 object_type: key.type_,
232 })
233 });
234
235 Ok(Box::new(mapped))
236 }
237
238 fn get_coin_info(&self, coin_type: &StructTag) -> StorageResult<Option<CoinInfo>> {
239 let coin_metadata_object_id = first_object_of_type(
245 self,
246 sui_types::coin::CoinMetadata::type_(coin_type.clone()),
247 )?;
248 let treasury_object_id =
249 first_object_of_type(self, sui_types::coin::TreasuryCap::type_(coin_type.clone()))?;
250 let regulated_coin_metadata_object_id = first_object_of_type(
251 self,
252 sui_types::coin::RegulatedCoinMetadata::type_(coin_type.clone()),
253 )?;
254
255 if coin_metadata_object_id.is_none()
256 && treasury_object_id.is_none()
257 && regulated_coin_metadata_object_id.is_none()
258 {
259 return Ok(None);
260 }
261
262 Ok(Some(CoinInfo {
263 coin_metadata_object_id,
264 treasury_object_id,
265 regulated_coin_metadata_object_id,
266 }))
267 }
268
269 fn get_balance(
270 &self,
271 owner: &SuiAddress,
272 coin_type: &StructTag,
273 ) -> StorageResult<Option<BalanceInfo>> {
274 let balance = self
275 .schema()
276 .get_balance(*owner, coin_type.clone().into())
277 .map_err(sui_types::storage::error::Error::custom)?;
278 Ok(balance.map(|b| BalanceInfo {
283 coin_balance: b.coin.clamp(0, u64::MAX as i128) as u64,
284 address_balance: b.address.clamp(0, u64::MAX as i128) as u64,
285 }))
286 }
287
288 fn balance_iter(
289 &self,
290 owner: &SuiAddress,
291 cursor: Option<(SuiAddress, StructTag)>,
292 ) -> StorageResult<BalanceIterator<'_>> {
293 use crate::schema::balance::{Key, OwnerPrefix};
294 let map = &self.schema().balance;
295 let from = cursor.map(|(_, tag)| Key {
299 owner: *owner,
300 coin_type: move_core_types::language_storage::TypeTag::Struct(Box::new(tag)),
301 });
302 let iter = match &from {
303 Some(from) => map.iter_prefix_from(&OwnerPrefix(*owner), from),
304 None => map.iter_prefix(&OwnerPrefix(*owner)),
305 }
306 .map_err(sui_types::storage::error::Error::custom)?;
307
308 let mapped = iter.filter_map(move |row| {
309 let (key, value) = match row {
310 Ok(pair) => pair,
311 Err(e) => return Some(Err(sui_types::storage::error::Error::custom(e))),
312 };
313 let balance = match crate::schema::balance::Balance::from_delta(&value.into_inner()) {
318 Ok(b) => b,
319 Err(e) => return Some(Err(sui_types::storage::error::Error::custom(e))),
320 };
321 let info = BalanceInfo {
325 coin_balance: balance.coin.clamp(0, u64::MAX as i128) as u64,
326 address_balance: balance.address.clamp(0, u64::MAX as i128) as u64,
327 };
328 let struct_tag = match key.coin_type {
329 move_core_types::language_storage::TypeTag::Struct(b) => *b,
330 _ => return None,
331 };
332 Some(Ok((struct_tag, info)))
333 });
334
335 Ok(Box::new(mapped))
336 }
337
338 fn package_versions_iter(
339 &self,
340 original_id: ObjectID,
341 cursor: Option<u64>,
342 ) -> StorageResult<PackageVersionsIterator<'_>> {
343 use crate::schema::package_versions::{Key, OriginalIdPrefix};
344 let map = &self.schema().package_versions;
350 let iter = match cursor {
351 Some(version) => map.iter_prefix_from(
352 &OriginalIdPrefix(original_id),
353 &Key {
354 original_id,
355 version,
356 },
357 ),
358 None => map.iter_prefix(&OriginalIdPrefix(original_id)),
359 }
360 .map_err(sui_types::storage::error::Error::custom)?;
361 let mapped = iter.map(move |row| {
362 let (key, value) = row.map_err(to_typed_store_err)?;
363 let storage_id_bytes: [u8; 32] = (&value.into_inner().storage_id[..])
365 .try_into()
366 .map_err(|_| {
367 TypedStoreError::SerializationError("package_versions storage_id length".into())
368 })?;
369 Ok((key.version, ObjectID::new(storage_id_bytes)))
370 });
371 Ok(Box::new(mapped))
372 }
373
374 fn get_highest_indexed_checkpoint_seq_number(
375 &self,
376 ) -> StorageResult<Option<CheckpointSequenceNumber>> {
377 self.min_committed(self.pipelines.iter().copied())
387 }
388
389 fn get_highest_live_indexed_checkpoint_seq_number(
390 &self,
391 ) -> StorageResult<Option<CheckpointSequenceNumber>> {
392 self.highest_live_committed_checkpoint()
396 }
397
398 fn ledger_tx_seq_digest(&self, tx_seq: u64) -> StorageResult<Option<LedgerTxSeqDigest>> {
399 let meta = self
400 .schema()
401 .get_tx_metadata_by_seq(tx_seq)
402 .map_err(sui_types::storage::error::Error::custom)?;
403 Ok(meta.map(|m| LedgerTxSeqDigest {
404 tx_sequence_number: tx_seq,
405 digest: m.digest,
406 event_count: m.event_count,
407 tx_offset: m.ckpt_position,
408 checkpoint_number: m.checkpoint_seq,
409 }))
410 }
411
412 fn ledger_tx_seq_digest_iter(
413 &self,
414 start: u64,
415 end_exclusive: u64,
416 descending: bool,
417 ) -> StorageResult<LedgerTxSeqDigestIterator<'_>> {
418 use crate::schema::primitives::U64Be;
419 let range = (
420 Bound::Included(U64Be(start)),
421 Bound::Excluded(U64Be(end_exclusive)),
422 );
423 let map = &self.schema().tx_metadata_by_seq;
424 let project = move |row: Result<
425 (U64Be, crate::schema::tx_metadata_by_seq::Value),
426 sui_consistent_store::error::Error,
427 >| {
428 let (U64Be(seq), value) = row.map_err(to_typed_store_err)?;
429 let stored = value.into_inner();
430 let digest_bytes: [u8; 32] = (&stored.digest[..]).try_into().map_err(|_| {
431 TypedStoreError::SerializationError("tx_metadata digest length".into())
432 })?;
433 Ok(LedgerTxSeqDigest {
434 tx_sequence_number: seq,
435 digest: sui_types::digests::TransactionDigest::new(digest_bytes),
436 event_count: stored.event_count,
437 tx_offset: stored.ckpt_position,
438 checkpoint_number: stored.checkpoint_seq,
439 })
440 };
441 if descending {
442 let iter = map
443 .iter_rev(range)
444 .map_err(sui_types::storage::error::Error::custom)?;
445 Ok(Box::new(iter.map(project)))
446 } else {
447 let iter = map
448 .iter(range)
449 .map_err(sui_types::storage::error::Error::custom)?;
450 Ok(Box::new(iter.map(project)))
451 }
452 }
453
454 fn transaction_bitmap_bucket_iter(
455 &self,
456 dimension_key: Vec<u8>,
457 start_bucket: u64,
458 end_bucket_exclusive: u64,
459 descending: bool,
460 ) -> StorageResult<LedgerBitmapBucketIterator<'_>> {
461 let map = &self.schema().transaction_bitmap;
462 let lower = crate::schema::transaction_bitmap::Key {
463 dimension_key: dimension_key.clone(),
464 bucket: start_bucket,
465 };
466 let upper = crate::schema::transaction_bitmap::Key {
467 dimension_key: dimension_key.clone(),
468 bucket: end_bucket_exclusive,
469 };
470 let inner = if descending {
471 BoundedBitmapIter::Rev(
472 map.iter_rev(lower..upper)
473 .map_err(sui_types::storage::error::Error::custom)?,
474 )
475 } else {
476 BoundedBitmapIter::Fwd(
477 map.iter(lower..upper)
478 .map_err(sui_types::storage::error::Error::custom)?,
479 )
480 };
481 Ok(Box::new(TransactionBitmapBucketIter {
482 inner,
483 dimension_key,
484 }))
485 }
486
487 fn event_bitmap_bucket_iter(
488 &self,
489 dimension_key: Vec<u8>,
490 start_bucket: u64,
491 end_bucket_exclusive: u64,
492 descending: bool,
493 ) -> StorageResult<LedgerBitmapBucketIterator<'_>> {
494 let map = &self.schema().event_bitmap;
495 let lower = crate::schema::event_bitmap::Key {
496 dimension_key: dimension_key.clone(),
497 bucket: start_bucket,
498 };
499 let upper = crate::schema::event_bitmap::Key {
500 dimension_key: dimension_key.clone(),
501 bucket: end_bucket_exclusive,
502 };
503 let inner = if descending {
504 BoundedBitmapIter::Rev(
505 map.iter_rev(lower..upper)
506 .map_err(sui_types::storage::error::Error::custom)?,
507 )
508 } else {
509 BoundedBitmapIter::Fwd(
510 map.iter(lower..upper)
511 .map_err(sui_types::storage::error::Error::custom)?,
512 )
513 };
514 Ok(Box::new(EventBitmapBucketIter {
515 inner,
516 dimension_key,
517 }))
518 }
519}
520
521fn project_bitmap_row(
525 row: Result<
526 (
527 crate::schema::transaction_bitmap::Key,
528 crate::schema::transaction_bitmap::Value,
529 ),
530 sui_consistent_store::error::Error,
531 >,
532) -> Result<sui_types::storage::LedgerBitmapBucket, TypedStoreError> {
533 let (key, value) = row.map_err(to_typed_store_err)?;
534 let bitmap = roaring::RoaringBitmap::deserialize_from(value.into_inner().data.as_ref())
535 .map_err(|e| TypedStoreError::SerializationError(format!("RoaringBitmap: {e}")))?;
536 Ok(sui_types::storage::LedgerBitmapBucket {
537 bucket_id: key.bucket,
538 bitmap,
539 })
540}
541
542fn project_event_bitmap_row(
547 row: Result<
548 (
549 crate::schema::event_bitmap::Key,
550 crate::schema::event_bitmap::Value,
551 ),
552 sui_consistent_store::error::Error,
553 >,
554) -> Result<sui_types::storage::LedgerBitmapBucket, TypedStoreError> {
555 let (key, value) = row.map_err(to_typed_store_err)?;
556 let bitmap = roaring::RoaringBitmap::deserialize_from(value.into_inner().data.as_ref())
557 .map_err(|e| TypedStoreError::SerializationError(format!("RoaringBitmap: {e}")))?;
558 Ok(sui_types::storage::LedgerBitmapBucket {
559 bucket_id: key.bucket,
560 bitmap,
561 })
562}
563
564#[cfg(test)]
565mod tests {
566 use std::sync::Arc;
567
568 use sui_consistent_store::Db;
569 use sui_consistent_store::DbOptions;
570 use sui_types::base_types::ObjectID;
571 use sui_types::storage::RpcIndexes;
572
573 use crate::RpcStoreSchema;
574 use crate::reader::RpcStoreReader;
575 use crate::schema::transaction_bitmap;
576
577 fn setup() -> (tempfile::TempDir, Db, RpcStoreReader) {
578 let dir = tempfile::tempdir().unwrap();
579 let (db, schema) = Db::open::<RpcStoreSchema>(dir.path(), DbOptions::default()).unwrap();
580 let reader = RpcStoreReader::new(db.clone(), Arc::new(schema));
581 (dir, db, reader)
582 }
583
584 #[test]
585 fn transaction_bitmap_bucket_iter_walks_range_ascending() {
586 let (_dir, db, reader) = setup();
587 let dim = b"sender:alice".to_vec();
588
589 let mut batch = db.batch();
590 for tx_seq in [
591 1u64,
592 transaction_bitmap::TX_BUCKET_SIZE + 5,
593 3 * transaction_bitmap::TX_BUCKET_SIZE + 9,
594 ] {
595 let (k, v) = transaction_bitmap::store_match(dim.clone(), tx_seq);
596 batch
597 .merge(&reader.schema().transaction_bitmap, &k, &v)
598 .unwrap();
599 }
600 batch.commit().unwrap();
601
602 let buckets: Vec<u64> = reader
603 .transaction_bitmap_bucket_iter(dim.clone(), 0, 5, false)
604 .unwrap()
605 .map(|res| res.unwrap().bucket_id)
606 .collect();
607 assert_eq!(buckets, vec![0, 1, 3]);
608 }
609
610 #[test]
611 fn transaction_bitmap_bucket_iter_respects_bucket_range_bounds() {
612 let (_dir, db, reader) = setup();
613 let dim = b"sender:alice".to_vec();
614
615 let mut batch = db.batch();
616 for tx_seq in [
617 1u64,
618 transaction_bitmap::TX_BUCKET_SIZE + 5,
619 3 * transaction_bitmap::TX_BUCKET_SIZE + 9,
620 ] {
621 let (k, v) = transaction_bitmap::store_match(dim.clone(), tx_seq);
622 batch
623 .merge(&reader.schema().transaction_bitmap, &k, &v)
624 .unwrap();
625 }
626 batch.commit().unwrap();
627
628 let buckets: Vec<u64> = reader
630 .transaction_bitmap_bucket_iter(dim.clone(), 1, 3, false)
631 .unwrap()
632 .map(|res| res.unwrap().bucket_id)
633 .collect();
634 assert_eq!(buckets, vec![1]);
635 }
636
637 #[test]
638 fn transaction_bitmap_bucket_iter_descending_reverses_order() {
639 let (_dir, db, reader) = setup();
640 let dim = b"sender:alice".to_vec();
641
642 let mut batch = db.batch();
643 for tx_seq in [
644 1u64,
645 transaction_bitmap::TX_BUCKET_SIZE + 5,
646 3 * transaction_bitmap::TX_BUCKET_SIZE + 9,
647 ] {
648 let (k, v) = transaction_bitmap::store_match(dim.clone(), tx_seq);
649 batch
650 .merge(&reader.schema().transaction_bitmap, &k, &v)
651 .unwrap();
652 }
653 batch.commit().unwrap();
654
655 let buckets: Vec<u64> = reader
656 .transaction_bitmap_bucket_iter(dim, 0, 5, true)
657 .unwrap()
658 .map(|res| res.unwrap().bucket_id)
659 .collect();
660 assert_eq!(buckets, vec![3, 1, 0]);
661 }
662
663 #[test]
664 fn transaction_bitmap_bucket_iter_isolates_dimension() {
665 let (_dir, db, reader) = setup();
666 let alice = b"sender:alice".to_vec();
667 let bob = b"sender:bob".to_vec();
668
669 let mut batch = db.batch();
670 let (k_a, v_a) = transaction_bitmap::store_match(alice.clone(), 1);
671 let (k_b, v_b) = transaction_bitmap::store_match(bob, 1);
672 batch
673 .merge(&reader.schema().transaction_bitmap, &k_a, &v_a)
674 .unwrap();
675 batch
676 .merge(&reader.schema().transaction_bitmap, &k_b, &v_b)
677 .unwrap();
678 batch.commit().unwrap();
679
680 let buckets: Vec<u64> = reader
681 .transaction_bitmap_bucket_iter(alice, 0, 5, false)
682 .unwrap()
683 .map(|res| res.unwrap().bucket_id)
684 .collect();
685 assert_eq!(buckets, vec![0]);
687 }
688
689 #[test]
695 fn highest_indexed_bounds_to_registered_pipelines() {
696 use sui_consistent_store::FrameworkSchema;
697 use sui_consistent_store::PipelineTaskKey;
698 use sui_consistent_store::Watermark;
699
700 let (_dir, db, reader) = setup();
701
702 let framework = FrameworkSchema::new(db.clone());
703 let stamp = |names: &[&str], hi: u64| {
704 let mut batch = db.batch();
705 for name in names {
706 batch
707 .put(
708 &framework.watermarks,
709 &PipelineTaskKey::new(*name),
710 &Watermark::for_checkpoint(hi),
711 )
712 .unwrap();
713 }
714 batch.commit().unwrap();
715 };
716
717 assert_eq!(
720 reader.get_highest_indexed_checkpoint_seq_number().unwrap(),
721 None,
722 );
723 stamp(&["transactions"], 3);
724 assert_eq!(
725 reader.get_highest_indexed_checkpoint_seq_number().unwrap(),
726 None,
727 );
728
729 stamp(crate::LIVE_COHORT, 40);
732 assert_eq!(
733 reader.get_highest_indexed_checkpoint_seq_number().unwrap(),
734 None,
735 );
736
737 stamp(crate::HISTORY_COHORT, 40);
741 stamp(&[crate::HISTORY_COHORT[0]], 25);
742 assert_eq!(
743 reader.get_highest_indexed_checkpoint_seq_number().unwrap(),
744 Some(25),
745 );
746
747 let custom = reader.clone().with_pipelines(["transactions"]);
750 assert_eq!(
751 custom.get_highest_indexed_checkpoint_seq_number().unwrap(),
752 Some(3),
753 );
754 }
755
756 #[test]
757 fn transaction_bitmap_bucket_iter_seeks_within_range() {
758 let (_dir, db, reader) = setup();
759 let dimension_key = b"sender:alice".to_vec();
760 let mut batch = db.batch();
761 for bucket in [0, 3, 7] {
762 let (key, value) = transaction_bitmap::store_match(
763 dimension_key.clone(),
764 bucket * transaction_bitmap::TX_BUCKET_SIZE + 1,
765 );
766 batch
767 .merge(&reader.schema().transaction_bitmap, &key, &value)
768 .unwrap();
769 }
770 batch.commit().unwrap();
771
772 let mut iter = reader
773 .transaction_bitmap_bucket_iter(dimension_key, 0, 8, false)
774 .unwrap();
775 iter.seek_bucket(2);
776 assert_eq!(iter.next().unwrap().unwrap().bucket_id, 3);
777 iter.seek_bucket(7);
778 assert_eq!(iter.next().unwrap().unwrap().bucket_id, 7);
779 iter.seek_bucket(9);
780 assert!(iter.next().is_none());
781 }
782
783 #[test]
784 fn transaction_bitmap_bucket_iter_seeks_descending() {
785 let (_dir, db, reader) = setup();
786 let dimension_key = b"sender:alice".to_vec();
787 let mut batch = db.batch();
788 for bucket in [0, 3, 7] {
789 let (key, value) = transaction_bitmap::store_match(
790 dimension_key.clone(),
791 bucket * transaction_bitmap::TX_BUCKET_SIZE + 1,
792 );
793 batch
794 .merge(&reader.schema().transaction_bitmap, &key, &value)
795 .unwrap();
796 }
797 batch.commit().unwrap();
798
799 let mut iter = reader
800 .transaction_bitmap_bucket_iter(dimension_key, 0, 8, true)
801 .unwrap();
802 iter.seek_bucket(5);
803 assert_eq!(iter.next().unwrap().unwrap().bucket_id, 3);
804 iter.seek_bucket(0);
805 assert_eq!(iter.next().unwrap().unwrap().bucket_id, 0);
806 }
807
808 #[test]
809 fn bitmap_bucket_iter_seek_respects_dimension_isolation() {
810 let (_dir, db, reader) = setup();
811 let alice = b"sender:alice".to_vec();
812 let bob = b"sender:bob".to_vec();
813 let mut batch = db.batch();
814 for (dimension_key, bucket) in [(&alice, 0), (&alice, 7), (&bob, 3), (&bob, 7)] {
815 let (key, value) = transaction_bitmap::store_match(
816 dimension_key.clone(),
817 bucket * transaction_bitmap::TX_BUCKET_SIZE + 1,
818 );
819 batch
820 .merge(&reader.schema().transaction_bitmap, &key, &value)
821 .unwrap();
822 }
823 batch.commit().unwrap();
824
825 let mut iter = reader
826 .transaction_bitmap_bucket_iter(alice, 0, 8, false)
827 .unwrap();
828 iter.seek_bucket(3);
829 assert_eq!(iter.next().unwrap().unwrap().bucket_id, 7);
830 assert!(iter.next().is_none());
831 }
832
833 #[test]
834 fn get_coin_info_finds_metadata_and_treasury_objects() {
835 use move_core_types::language_storage::StructTag;
836
837 use crate::schema::object_by_type;
838 use crate::schema::primitives::U64Varint;
839
840 let (_dir, db, reader) = setup();
841
842 let coin_type = StructTag {
848 address: move_core_types::account_address::AccountAddress::new([2u8; 32]),
849 module: move_core_types::identifier::Identifier::new("sui").unwrap(),
850 name: move_core_types::identifier::Identifier::new("SUI").unwrap(),
851 type_params: vec![],
852 };
853 let metadata_type = sui_types::coin::CoinMetadata::type_(coin_type.clone());
854 let treasury_type = sui_types::coin::TreasuryCap::type_(coin_type.clone());
855
856 let metadata_object_id = ObjectID::from_single_byte(0xA1);
857 let treasury_object_id = ObjectID::from_single_byte(0xA2);
858
859 let mut batch = db.batch();
860 batch
861 .put(
862 &reader.schema().object_by_type,
863 &object_by_type::Key {
864 type_: metadata_type,
865 object_id: metadata_object_id,
866 },
867 &U64Varint(1),
868 )
869 .unwrap();
870 batch
871 .put(
872 &reader.schema().object_by_type,
873 &object_by_type::Key {
874 type_: treasury_type,
875 object_id: treasury_object_id,
876 },
877 &U64Varint(1),
878 )
879 .unwrap();
880 batch.commit().unwrap();
881
882 let info = reader
883 .get_coin_info(&coin_type)
884 .unwrap()
885 .expect("coin info present");
886 assert_eq!(info.coin_metadata_object_id, Some(metadata_object_id));
887 assert_eq!(info.treasury_object_id, Some(treasury_object_id));
888 assert_eq!(info.regulated_coin_metadata_object_id, None);
889 }
890
891 #[test]
892 fn get_coin_info_returns_none_when_no_wrappers_indexed() {
893 use move_core_types::language_storage::StructTag;
894
895 let (_dir, _db, reader) = setup();
896
897 let coin_type = StructTag {
898 address: move_core_types::account_address::AccountAddress::new([3u8; 32]),
899 module: move_core_types::identifier::Identifier::new("custom").unwrap(),
900 name: move_core_types::identifier::Identifier::new("COIN").unwrap(),
901 type_params: vec![],
902 };
903 assert!(reader.get_coin_info(&coin_type).unwrap().is_none());
904 }
905
906 #[test]
907 fn transaction_bitmap_bucket_iter_returns_decoded_bitmap() {
908 let (_dir, db, reader) = setup();
909 let dim = b"sender:alice".to_vec();
910
911 let mut batch = db.batch();
912 for tx_seq in [1u64, 17, 256] {
913 let (k, v) = transaction_bitmap::store_match(dim.clone(), tx_seq);
914 batch
915 .merge(&reader.schema().transaction_bitmap, &k, &v)
916 .unwrap();
917 }
918 batch.commit().unwrap();
919
920 let first = reader
921 .transaction_bitmap_bucket_iter(dim, 0, 1, false)
922 .unwrap()
923 .next()
924 .unwrap()
925 .unwrap();
926 let bits: Vec<u32> = first.bitmap.iter().collect();
927 assert_eq!(bits, vec![1, 17, 256]);
928 }
929
930 fn a_type() -> move_core_types::language_storage::StructTag {
937 move_core_types::language_storage::StructTag {
938 address: move_core_types::account_address::AccountAddress::TWO,
939 module: move_core_types::identifier::Identifier::new("m").unwrap(),
940 name: move_core_types::identifier::Identifier::new("T").unwrap(),
941 type_params: vec![],
942 }
943 }
944
945 #[test]
946 fn owned_objects_iter_paginates_each_object_once() {
947 use crate::schema::object_by_owner::{Key, OwnerKind};
948 use crate::schema::primitives::U64Varint;
949 use sui_types::base_types::SuiAddress;
950
951 let (_dir, db, reader) = setup();
952 let owner = SuiAddress::ZERO;
953 let ids: Vec<ObjectID> = (1u8..=5).map(ObjectID::from_single_byte).collect();
954
955 let mut batch = db.batch();
956 for id in &ids {
957 batch
958 .put(
959 &reader.schema().object_by_owner,
960 &Key {
961 kind: OwnerKind::AddressOwner(owner),
962 type_: a_type(),
963 inverted_balance: None,
964 object_id: *id,
965 },
966 &U64Varint(1),
967 )
968 .unwrap();
969 }
970 batch.commit().unwrap();
971
972 let mut seen = Vec::new();
973 let mut cursor = None;
974 for _ in 0..ids.len() + 2 {
977 let mut iter = reader
978 .owned_objects_iter(owner, None, cursor.take())
979 .unwrap();
980 for _ in 0..2 {
981 match iter.next() {
982 Some(res) => seen.push(res.unwrap().object_id),
983 None => break,
984 }
985 }
986 cursor = iter.next().transpose().unwrap();
987 if cursor.is_none() {
988 break;
989 }
990 }
991
992 assert_eq!(seen.len(), ids.len(), "each object yielded exactly once");
993 let mut unique = seen.clone();
994 unique.sort();
995 unique.dedup();
996 assert_eq!(unique, ids, "every object, no gaps or duplicates");
997 }
998
999 #[test]
1000 fn dynamic_field_iter_paginates_each_field_once() {
1001 use crate::schema::object_by_owner::{Key, OwnerKind};
1002 use crate::schema::primitives::U64Varint;
1003
1004 let (_dir, db, reader) = setup();
1005 let parent = ObjectID::from_single_byte(0xAA);
1006 let fields: Vec<ObjectID> = (1u8..=5).map(ObjectID::from_single_byte).collect();
1007
1008 let mut batch = db.batch();
1009 for id in &fields {
1010 batch
1011 .put(
1012 &reader.schema().object_by_owner,
1013 &Key {
1014 kind: OwnerKind::ObjectOwner(parent.into()),
1015 type_: a_type(),
1016 inverted_balance: None,
1017 object_id: *id,
1018 },
1019 &U64Varint(1),
1020 )
1021 .unwrap();
1022 }
1023 batch.commit().unwrap();
1024
1025 let mut seen = Vec::new();
1026 let mut cursor = None;
1027 for _ in 0..fields.len() + 2 {
1028 let mut iter = reader.dynamic_field_iter(parent, cursor.take()).unwrap();
1029 for _ in 0..2 {
1030 match iter.next() {
1031 Some(res) => seen.push(res.unwrap().field_id),
1032 None => break,
1033 }
1034 }
1035 cursor = iter.next().transpose().unwrap();
1036 if cursor.is_none() {
1037 break;
1038 }
1039 }
1040
1041 assert_eq!(seen.len(), fields.len(), "each field yielded exactly once");
1042 let mut unique = seen.clone();
1043 unique.sort();
1044 unique.dedup();
1045 assert_eq!(unique, fields, "every field, no gaps or duplicates");
1046 }
1047
1048 #[test]
1049 fn balance_iter_paginates_each_coin_type_once() {
1050 use crate::schema::balance::coin_delta;
1051 use move_core_types::language_storage::{StructTag, TypeTag};
1052 use sui_types::base_types::SuiAddress;
1053
1054 let coin_type = |name: &str| -> TypeTag {
1055 TypeTag::Struct(Box::new(StructTag {
1056 address: move_core_types::account_address::AccountAddress::TWO,
1057 module: move_core_types::identifier::Identifier::new("coin").unwrap(),
1058 name: move_core_types::identifier::Identifier::new(name).unwrap(),
1059 type_params: vec![],
1060 }))
1061 };
1062
1063 let (_dir, db, reader) = setup();
1064 let owner = SuiAddress::ZERO;
1065 let names = ["c0", "c1", "c2", "c3", "c4"];
1066
1067 let mut batch = db.batch();
1068 for name in names {
1069 let (k, v) = coin_delta(owner, coin_type(name), 100);
1070 batch.merge(&reader.schema().balance, &k, &v).unwrap();
1071 }
1072 batch.commit().unwrap();
1073
1074 let mut seen: Vec<StructTag> = Vec::new();
1075 let mut cursor: Option<(SuiAddress, StructTag)> = None;
1076 for _ in 0..names.len() + 2 {
1077 let mut iter = reader.balance_iter(&owner, cursor.take()).unwrap();
1078 for _ in 0..2 {
1079 match iter.next() {
1080 Some(res) => seen.push(res.unwrap().0),
1081 None => break,
1082 }
1083 }
1084 cursor = iter
1085 .next()
1086 .transpose()
1087 .unwrap()
1088 .map(|(tag, _)| (owner, tag));
1089 if cursor.is_none() {
1090 break;
1091 }
1092 }
1093
1094 assert_eq!(
1095 seen.len(),
1096 names.len(),
1097 "each coin type yielded exactly once"
1098 );
1099 let mut unique = seen.clone();
1100 unique.sort();
1101 unique.dedup();
1102 assert_eq!(unique.len(), names.len(), "no duplicate coin types");
1103 }
1104
1105 #[test]
1106 fn package_versions_iter_paginates_each_version_once() {
1107 let (_dir, db, reader) = setup();
1108 let original = ObjectID::from_single_byte(0xBB);
1109 let versions: Vec<u64> = (1..=5).collect();
1110
1111 let mut batch = db.batch();
1112 for &version in &versions {
1113 let (k, v) = crate::schema::package_versions::store(
1114 original,
1115 version,
1116 ObjectID::from_single_byte(version as u8),
1117 version,
1118 );
1119 batch
1120 .put(&reader.schema().package_versions, &k, &v)
1121 .unwrap();
1122 }
1123 batch.commit().unwrap();
1124
1125 let mut seen: Vec<u64> = Vec::new();
1126 let mut cursor: Option<u64> = None;
1127 for _ in 0..versions.len() + 2 {
1128 let mut iter = reader
1129 .package_versions_iter(original, cursor.take())
1130 .unwrap();
1131 for _ in 0..2 {
1132 match iter.next() {
1133 Some(res) => seen.push(res.unwrap().0),
1134 None => break,
1135 }
1136 }
1137 cursor = iter.next().transpose().unwrap().map(|(v, _)| v);
1138 if cursor.is_none() {
1139 break;
1140 }
1141 }
1142
1143 assert_eq!(
1144 seen.len(),
1145 versions.len(),
1146 "each version yielded exactly once"
1147 );
1148 let mut unique = seen.clone();
1149 unique.sort();
1150 unique.dedup();
1151 assert_eq!(unique, versions, "every version, no gaps or duplicates");
1152 }
1153}