Skip to main content

sui_rpc_store/reader/
indexes.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4//! [`RpcIndexes`] adapter — owned-object, dynamic-field, balance,
5//! package-versions, epoch-info, ledger-history, and coin-info
6//! lookups.
7//!
8//! Trait methods returning `Result` over [`typed_store_error::TypedStoreError`]
9//! wrap our [`sui_consistent_store::error::Error`] via
10//! [`TypedStoreError::RocksDBError`].
11
12use 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
34/// Local alias matching the inaccessible
35/// `sui_types::storage::read_store::PackageVersionsIterator`.
36type 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
125/// Find the object id of the first row whose Move type matches the
126/// pinned `struct_tag`. The `object_by_type` index sorts by
127/// `(type, id)`, so the first prefix-scan row IS the lowest-id
128/// match — there should only be at most one in practice for the
129/// coin-wrapper types this is used with (CoinMetadata, TreasuryCap,
130/// RegulatedCoinMetadata are all unique per coin type).
131fn 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        // When resuming, the page token carries the cursor object's full sort
168        // position -- type, balance, and id -- so the scan seeks straight to it
169        // (inclusive) and stops at the end of the prefix. No post-filtering.
170        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        // Dynamic fields are `Field<Name, Value>` objects owned (in the
207        // object-owner sense) by `parent`, so they share the `object_by_owner`
208        // index with address-owned objects. The cursor carries the field's
209        // type and id -- its full sort position -- so the scan seeks straight
210        // to it. Field objects are never coins, so the balance component is
211        // always absent.
212        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        // Coin metadata / treasury cap / regulated coin metadata
240        // are typed objects whose Move type wraps the requested
241        // `coin_type`. Discover each via an `object_by_type`
242        // prefix scan keyed on the corresponding wrapper struct
243        // tag and take the first match.
244        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        // Report the coin and address halves independently (each clamped to
279        // non-negative); the caller sums them for the total. Reporting the
280        // total as `coin_balance` would double-count the address half once
281        // the caller adds them back together.
282        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        // When resuming, the page token carries the cursor's coin type -- the
296        // index's sort key after the owner -- so the scan seeks straight to it
297        // (inclusive) and stops at the end of the owner's balances.
298        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            // Project the merged proto value back into the typed
314            // `Balance` view through the same decoder `get_balance`
315            // uses, so a malformed payload surfaces as an error here
316            // too rather than being silently read as zero.
317            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            // Report the coin and address halves independently; the caller
322            // sums them for the total (reporting the total here would
323            // double-count the address half).
324            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        // The cursor passed in by `list_package_versions` is the version of the
345        // first row of the next page (the previous page popped its
346        // `page_size + 1`th row to derive it), so seek straight to it and
347        // resume inclusively. Versions sort ascending within a package, so the
348        // seek lands on the first row with `version >= cursor`.
349        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            // Decode storage_id (32 bytes).
364            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        // Min across the registered pipelines' watermarks: every CF
378        // this reader serves has caught up to at least this
379        // checkpoint, so reads against any of them are coherent
380        // through here. Bounded to the registered set -- rather than
381        // every row in the framework CF -- so a stale watermark left
382        // behind by a pipeline that is no longer registered cannot
383        // pin the reported tip forever. `None` while any registered
384        // pipeline has no watermark yet: nothing is fully indexed,
385        // so nothing is served (fail closed).
386        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        // Only the live cohort (owned objects, types, balances); the
393        // ledger-history cohort backfills independently and is excluded so a
394        // node is healthy as soon as its live-object reads are caught up.
395        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
521/// Project a `(transaction_bitmap::Key, Protobuf<BitmapBlob>)`
522/// row into a [`LedgerBitmapBucket`], deserializing the raw
523/// RoaringBitmap bytes off the protobuf payload.
524fn 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
542/// Project an `(event_bitmap::Key, Protobuf<BitmapBlob>)` row into
543/// a [`LedgerBitmapBucket`]. Same shape as
544/// [`project_bitmap_row`] but typed against the distinct event-CF
545/// key.
546fn 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        // Buckets `[1, 3)` — only the middle bucket survives.
629        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        // Bob's bucket is not visible under Alice's dimension.
686        assert_eq!(buckets, vec![0]);
687    }
688
689    /// The indexed tip is bounded by the registered pipeline set: a
690    /// stale watermark from a pipeline that is no longer registered
691    /// (e.g. `transactions`, disabled on the embedded deployment) must
692    /// not pin the reported tip, and a registered pipeline with no
693    /// watermark yet reads as "nothing fully indexed".
694    #[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        // Nothing committed yet: nothing is fully indexed. A stale row
718        // from an unregistered pipeline does not change that.
719        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        // Only part of the registered set committed: still nothing
730        // fully indexed (fail closed).
731        stamp(crate::LIVE_COHORT, 40);
732        assert_eq!(
733            reader.get_highest_indexed_checkpoint_seq_number().unwrap(),
734            None,
735        );
736
737        // The whole registered set committed: the slowest registered
738        // pipeline bounds the tip, and the stale `transactions` row at
739        // 3 stays excluded.
740        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        // A reader for a different deployment bounds by exactly its
748        // own registered set.
749        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        // Construct a synthetic coin type and seed
843        // `object_by_type` rows for its CoinMetadata and
844        // TreasuryCap wrappers. The real on-chain pipeline writes
845        // these rows for every Move object it sees; the test
846        // bypasses pipelines and writes the rows directly.
847        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    // Pagination regression coverage for the cursor-skip iterators. The
931    // `list_*` handlers page by taking `page_size` rows and then peeking the
932    // next row as the opaque page token; a broken skip-past-cursor predicate
933    // re-yields earlier rows on every page (duplicates, and a cursor that never
934    // advances). Each test walks every page and asserts each row appears once.
935
936    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        // Bounded so a non-advancing cursor surfaces as a failed assertion
975        // rather than an infinite loop.
976        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}