Skip to main content

sui_rpc_store/schema/
mod.rs

1// Copyright (c) Mysten Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Column-family layout for `sui-rpc-store`.
5//!
6//! Each CF lives in its own submodule that declares:
7//!
8//! - `NAME` — the on-disk column-family name.
9//! - `Key` — the key type with `Encode` / `Decode` pinning its
10//!   on-disk layout.
11//! - `Value` — the value type, typically `Protobuf<…>`.
12//! - `options(resolver)` — per-CF `rocksdb::Options`, obtained from the
13//!   [`sui_consistent_store::CfOptionsResolver`] with the CF's merge
14//!   operator and compaction filter (if any) layered on top.
15//!
16//! [`RpcStoreSchema`] aggregates these into the schema passed to
17//! [`sui_consistent_store::Db::open`]. Keys reused across multiple
18//! CFs live in [`primitives`].
19
20pub mod balance;
21pub mod checkpoint_contents;
22pub mod checkpoint_seq_by_digest;
23pub mod checkpoint_summary;
24pub mod effects;
25pub mod epochs;
26pub mod event_bitmap;
27pub mod events;
28pub mod object_by_owner;
29pub mod object_by_type;
30pub mod object_version_by_checkpoint;
31pub mod objects;
32pub mod package_versions;
33pub mod primitives;
34pub mod pruning_watermark;
35pub mod transaction_bitmap;
36pub mod transactions;
37pub mod tx_metadata_by_seq;
38pub mod tx_seq_by_digest;
39pub mod type_filter;
40
41use std::collections::BTreeMap;
42use std::path::Path;
43use std::sync::Arc;
44use std::sync::atomic::AtomicU64;
45use std::sync::atomic::Ordering;
46use sui_consistent_store::CfDescriptor;
47use sui_consistent_store::CfOptionsResolver;
48use sui_consistent_store::CfTuning;
49use sui_consistent_store::Compression;
50use sui_consistent_store::Db;
51use sui_consistent_store::DbMap;
52use sui_consistent_store::DbWideConfig;
53use sui_consistent_store::RocksDbConfig;
54use sui_consistent_store::Schema;
55use sui_consistent_store::SchemaAtSnapshot;
56use sui_consistent_store::Snapshot;
57use sui_consistent_store::WriteStallConfig;
58use sui_consistent_store::error::OpenError;
59use sui_consistent_store::reader::Reader;
60
61/// Typed handles to every CF in the `sui-rpc-store` layout.
62pub struct RpcStoreSchema<R: Reader = Db> {
63    /// Exclusive pruning floor in transaction-sequence (`tx_seq`)
64    /// space. Bitmap compaction filters may remove buckets containing
65    /// only transaction sequences strictly below this value.
66    tx_seq_pruning_floor: Arc<AtomicU64>,
67
68    /// Per-epoch metadata: protocol version, gas price, start and
69    /// end timestamps, and the epoch's final checkpoint.
70    pub epochs: DbMap<epochs::Key, epochs::Value, R>,
71
72    /// Signed checkpoint headers. The lightweight metadata served
73    /// by most "fetch a checkpoint" requests; the heavier contents
74    /// list lives in a separate CF.
75    pub checkpoint_summary: DbMap<checkpoint_summary::Key, checkpoint_summary::Value, R>,
76
77    /// The ordered list of executed transaction digests in each
78    /// checkpoint.
79    pub checkpoint_contents: DbMap<checkpoint_contents::Key, checkpoint_contents::Value, R>,
80
81    /// Resolves a checkpoint digest to its sequence number, which
82    /// is then the key for every other checkpoint-keyed CF.
83    pub checkpoint_seq_by_digest:
84        DbMap<checkpoint_seq_by_digest::Key, checkpoint_seq_by_digest::Value, R>,
85
86    /// Signed transactions, keyed by their assigned tx_seq.
87    pub transactions: DbMap<transactions::Key, transactions::Value, R>,
88
89    /// Resolves a transaction digest to its assigned tx_seq.
90    pub tx_seq_by_digest: DbMap<tx_seq_by_digest::Key, tx_seq_by_digest::Value, R>,
91
92    /// Per-transaction metadata: digest, the containing
93    /// checkpoint, position within that checkpoint, event count,
94    /// and the checkpoint's timestamp.
95    pub tx_metadata_by_seq: DbMap<tx_metadata_by_seq::Key, tx_metadata_by_seq::Value, R>,
96
97    /// The effects produced by each transaction, together with the
98    /// set of objects loaded but unchanged during execution.
99    pub effects: DbMap<effects::Key, effects::Value, R>,
100
101    /// The events emitted by each transaction.
102    pub events: DbMap<events::Key, events::Value, R>,
103
104    /// Every version of every object that has ever existed. A
105    /// prefix scan on the object id walks all versions in ascending
106    /// order; a reverse prefix scan resolves the latest version (the
107    /// greatest `(id, version)` row), the way the validator perpetual
108    /// store serves "latest object" reads.
109    pub objects: DbMap<objects::Key, objects::Value, R>,
110
111    /// An object's version as of a checkpoint: keyed by
112    /// `(object id, checkpoint)`, a reverse prefix scan resolves the
113    /// version live at the end of the most recent checkpoint, at or
114    /// before the one queried, in which the object changed. Backs
115    /// checkpoint-pinned historical reads that the version-keyed
116    /// `objects` CF cannot answer.
117    pub object_version_by_checkpoint:
118        DbMap<object_version_by_checkpoint::Key, object_version_by_checkpoint::Value, R>,
119
120    /// Supports listing an owner's objects, optionally filtered by
121    /// Move type. Coin-like objects sort richest-first within
122    /// each `(owner, type)` group so paginating valuable holdings
123    /// is a forward prefix scan.
124    pub object_by_owner: DbMap<object_by_owner::Key, object_by_owner::Value, R>,
125
126    /// Supports listing every live object of a given Move type,
127    /// regardless of owner.
128    pub object_by_type: DbMap<object_by_type::Key, object_by_type::Value, R>,
129
130    /// Tracks an account's balance per coin type, combining the
131    /// coin-derived component (sum of owned `Coin<T>` balances)
132    /// and the accumulator-derived component into a single row
133    /// merged from independent indexer pipelines.
134    pub balance: DbMap<balance::Key, balance::Value, R>,
135
136    /// Tracks every published version of a Move package and the
137    /// storage id under which each version lives.
138    pub package_versions: DbMap<package_versions::Key, package_versions::Value, R>,
139
140    /// Inverted bitmap index over transaction-sequence space,
141    /// supporting filtered transaction queries by indexed fields
142    /// such as sender, called function, or input/changed object.
143    pub transaction_bitmap: DbMap<transaction_bitmap::Key, transaction_bitmap::Value, R>,
144
145    /// Inverted bitmap index over packed event-sequence space,
146    /// supporting filtered event queries by event type, emitting
147    /// module, sender, and similar indexed fields.
148    pub event_bitmap: DbMap<event_bitmap::Key, event_bitmap::Value, R>,
149
150    // --- Bookkeeping ---
151    /// Singleton holding the lowest still-available `tx_seq`,
152    /// `checkpoint_seq`, and object version. Drives compaction
153    /// filters and feeds `available_range` responses.
154    pub pruning_watermark: DbMap<pruning_watermark::Key, pruning_watermark::Value, R>,
155}
156
157impl Schema for RpcStoreSchema {
158    fn open(
159        path: &Path,
160        opts: &CfOptionsResolver,
161        snapshot_capacity: usize,
162    ) -> Result<(Db, Self), OpenError> {
163        let tx_seq_pruning_floor = Arc::new(AtomicU64::new(0));
164        let db = Db::open_cfs(
165            path,
166            opts,
167            snapshot_capacity,
168            vec![
169                CfDescriptor::new(epochs::NAME, epochs::options(opts)),
170                CfDescriptor::new(checkpoint_summary::NAME, checkpoint_summary::options(opts)),
171                CfDescriptor::new(
172                    checkpoint_contents::NAME,
173                    checkpoint_contents::options(opts),
174                ),
175                CfDescriptor::new(
176                    checkpoint_seq_by_digest::NAME,
177                    checkpoint_seq_by_digest::options(opts),
178                ),
179                CfDescriptor::new(transactions::NAME, transactions::options(opts)),
180                CfDescriptor::new(tx_seq_by_digest::NAME, tx_seq_by_digest::options(opts)),
181                CfDescriptor::new(tx_metadata_by_seq::NAME, tx_metadata_by_seq::options(opts)),
182                CfDescriptor::new(effects::NAME, effects::options(opts)),
183                CfDescriptor::new(events::NAME, events::options(opts)),
184                CfDescriptor::new(objects::NAME, objects::options(opts)),
185                CfDescriptor::new(
186                    object_version_by_checkpoint::NAME,
187                    object_version_by_checkpoint::options(opts),
188                ),
189                CfDescriptor::new(object_by_owner::NAME, object_by_owner::options(opts)),
190                CfDescriptor::new(object_by_type::NAME, object_by_type::options(opts)),
191                CfDescriptor::new(balance::NAME, balance::options(opts)),
192                CfDescriptor::new(package_versions::NAME, package_versions::options(opts)),
193                CfDescriptor::new(
194                    transaction_bitmap::NAME,
195                    transaction_bitmap::options(opts, tx_seq_pruning_floor.clone()),
196                ),
197                CfDescriptor::new(
198                    event_bitmap::NAME,
199                    event_bitmap::options(opts, tx_seq_pruning_floor.clone()),
200                ),
201                CfDescriptor::new(pruning_watermark::NAME, pruning_watermark::options(opts)),
202            ],
203        )?;
204
205        let schema = Self {
206            tx_seq_pruning_floor,
207            epochs: DbMap::new(db.clone(), epochs::NAME)?,
208            checkpoint_summary: DbMap::new(db.clone(), checkpoint_summary::NAME)?,
209            checkpoint_contents: DbMap::new(db.clone(), checkpoint_contents::NAME)?,
210            checkpoint_seq_by_digest: DbMap::new(db.clone(), checkpoint_seq_by_digest::NAME)?,
211            transactions: DbMap::new(db.clone(), transactions::NAME)?,
212            tx_seq_by_digest: DbMap::new(db.clone(), tx_seq_by_digest::NAME)?,
213            tx_metadata_by_seq: DbMap::new(db.clone(), tx_metadata_by_seq::NAME)?,
214            effects: DbMap::new(db.clone(), effects::NAME)?,
215            events: DbMap::new(db.clone(), events::NAME)?,
216            objects: DbMap::new(db.clone(), objects::NAME)?,
217            object_version_by_checkpoint: DbMap::new(
218                db.clone(),
219                object_version_by_checkpoint::NAME,
220            )?,
221            object_by_owner: DbMap::new(db.clone(), object_by_owner::NAME)?,
222            object_by_type: DbMap::new(db.clone(), object_by_type::NAME)?,
223            balance: DbMap::new(db.clone(), balance::NAME)?,
224            package_versions: DbMap::new(db.clone(), package_versions::NAME)?,
225            transaction_bitmap: DbMap::new(db.clone(), transaction_bitmap::NAME)?,
226            event_bitmap: DbMap::new(db.clone(), event_bitmap::NAME)?,
227            pruning_watermark: DbMap::new(db.clone(), pruning_watermark::NAME)?,
228        };
229
230        if let Some(watermarks) = schema.get_pruning_watermarks().map_err(|error| {
231            OpenError::with_source("read persisted RPC pruning watermark", error)
232        })? {
233            schema
234                .tx_seq_pruning_floor
235                .store(watermarks.tx_seq_lo, Ordering::Relaxed);
236        }
237
238        Ok((db, schema))
239    }
240}
241
242impl SchemaAtSnapshot for RpcStoreSchema {
243    type At = RpcStoreSchema<Snapshot>;
244    fn at(&self, snap: &Snapshot) -> Self::At {
245        RpcStoreSchema {
246            tx_seq_pruning_floor: self.tx_seq_pruning_floor.clone(),
247            epochs: self.epochs.at(snap),
248            checkpoint_summary: self.checkpoint_summary.at(snap),
249            checkpoint_contents: self.checkpoint_contents.at(snap),
250            checkpoint_seq_by_digest: self.checkpoint_seq_by_digest.at(snap),
251            transactions: self.transactions.at(snap),
252            tx_seq_by_digest: self.tx_seq_by_digest.at(snap),
253            tx_metadata_by_seq: self.tx_metadata_by_seq.at(snap),
254            effects: self.effects.at(snap),
255            events: self.events.at(snap),
256            objects: self.objects.at(snap),
257            object_version_by_checkpoint: self.object_version_by_checkpoint.at(snap),
258            object_by_owner: self.object_by_owner.at(snap),
259            object_by_type: self.object_by_type.at(snap),
260            balance: self.balance.at(snap),
261            package_versions: self.package_versions.at(snap),
262            transaction_bitmap: self.transaction_bitmap.at(snap),
263            event_bitmap: self.event_bitmap.at(snap),
264            pruning_watermark: self.pruning_watermark.at(snap),
265        }
266    }
267}
268
269/// The tuned [`RocksDbConfig`] this crate ships as its baseline for
270/// the `sui-rpc-store` column families.
271///
272/// Bitmap CFs use a 7-day periodic-compaction interval by default.
273///
274/// The defaults port the production-proven settings from `typed_store`
275/// and bake in a "no write stalls, generous compaction parallelism"
276/// policy: the pending-compaction stall limits are disabled and the L0
277/// triggers raised so neither the bulk restore nor steady-state indexing
278/// throttles on compaction debt, while the L0 stop trigger still bounds
279/// a runaway backlog.
280pub fn default_rocksdb_config() -> RocksDbConfig {
281    let write_stall = WriteStallConfig {
282        soft_pending_compaction_bytes_limit_mb: Some(0),
283        hard_pending_compaction_bytes_limit_mb: Some(0),
284        level0_file_num_compaction_trigger: Some(4),
285        level0_slowdown_writes_trigger: Some(512),
286        level0_stop_writes_trigger: Some(1024),
287    };
288
289    let default_cf = CfTuning {
290        write_buffer_size_mb: Some(64),
291        max_write_buffer_number: Some(6),
292        compression: Some(Compression::Lz4),
293        bottommost_compression: Some(Compression::Zstd),
294        block_size_kb: Some(16),
295        bloom_filter_bits: None,
296        memtable_prefix_bloom_ratio: Some(0.02),
297        target_file_size_mb: Some(128),
298        periodic_compaction_seconds: None,
299        write_stall,
300    };
301
302    let mut column_family = BTreeMap::new();
303
304    // Point-lookup CFs: a whole-key bloom filter lets reads skip SSTs
305    // that cannot contain the requested key.
306    let point_lookup = CfTuning {
307        bloom_filter_bits: Some(10.0),
308        ..Default::default()
309    };
310    for name in [tx_seq_by_digest::NAME, checkpoint_seq_by_digest::NAME] {
311        column_family.insert(name.to_string(), point_lookup.clone());
312    }
313
314    // Bitmap CFs accumulate large merge-blob values; a bigger memtable
315    // amortizes the merge-and-flush churn. Periodic compaction ensures
316    // old merge operands are materialized and later filtered.
317    let bitmap = CfTuning {
318        write_buffer_size_mb: Some(256),
319        periodic_compaction_seconds: Some(7 * 86_400),
320        ..Default::default()
321    };
322    for name in [transaction_bitmap::NAME, event_bitmap::NAME] {
323        column_family.insert(name.to_string(), bitmap.clone());
324    }
325
326    RocksDbConfig {
327        db: DbWideConfig {
328            parallelism: Some(8),
329            max_background_jobs: None,
330            // RocksDB's default (`-1`) keeps every SST open, which
331            // exhausts the process file-descriptor budget on a
332            // large DB (a formal-snapshot restore writes thousands
333            // of SSTs and fails with "Too many open files"). Mirror
334            // `typed_store::default_db_options`: raise the fd limit
335            // toward the hard cap and bound the table cache to an
336            // eighth of it. `None` on platforms without the syscall
337            // (e.g. Windows), leaving the RocksDB default.
338            max_open_files: fdlimit::raise_fd_limit()
339                .map(|limit| (limit / 8).try_into().unwrap_or(i32::MAX)),
340            db_write_buffer_size_mb: Some(1024),
341            max_total_wal_size_mb: Some(1024),
342            enable_pipelined_write: Some(true),
343            table_cache_num_shard_bits: Some(10),
344            block_cache_size_mb: Some(1024),
345        },
346        default_cf,
347        column_family,
348    }
349}
350
351#[cfg(test)]
352mod tests {
353    use sui_consistent_store::Db;
354    use sui_consistent_store::DbOptions;
355    use sui_types::base_types::ObjectID;
356    use sui_types::base_types::SequenceNumber;
357
358    use super::*;
359
360    #[test]
361    fn opens_with_all_cfs() {
362        let dir = tempfile::tempdir().unwrap();
363        let (_db, schema) = Db::open::<RpcStoreSchema>(dir.path(), DbOptions::default()).unwrap();
364        // Empty database — every typed handle is constructed; a
365        // miss on any of them returns None instead of an open-time
366        // missing-CF error.
367        assert!(
368            schema
369                .objects
370                .get(&objects::Key {
371                    id: ObjectID::ZERO,
372                    version: SequenceNumber::from_u64(0),
373                })
374                .unwrap()
375                .is_none()
376        );
377        assert!(
378            schema
379                .pruning_watermark
380                .get(&primitives::UnitKey)
381                .unwrap()
382                .is_none()
383        );
384    }
385
386    #[test]
387    fn default_config_validates() {
388        default_rocksdb_config()
389            .validate()
390            .expect("shipped default config must be internally consistent");
391    }
392
393    #[test]
394    fn default_config_opens_every_cf() {
395        // Exercises the full resolve path: validation, the shared
396        // block cache, and every CF's merge operator / compaction
397        // filter layered on the tuned per-CF options.
398        let dir = tempfile::tempdir().unwrap();
399        let opts = DbOptions {
400            rocksdb: default_rocksdb_config(),
401            snapshot_capacity: 32,
402        };
403        let (_db, schema) = Db::open::<RpcStoreSchema>(dir.path(), opts).unwrap();
404        assert!(
405            schema
406                .pruning_watermark
407                .get(&primitives::UnitKey)
408                .unwrap()
409                .is_none()
410        );
411    }
412
413    #[test]
414    fn default_config_sets_per_cf_deviations() {
415        let cfg = default_rocksdb_config();
416        // Point-lookup CFs get a bloom filter.
417        assert_eq!(
418            cfg.column_family[tx_seq_by_digest::NAME].bloom_filter_bits,
419            Some(10.0)
420        );
421        assert_eq!(
422            cfg.column_family[checkpoint_seq_by_digest::NAME].bloom_filter_bits,
423            Some(10.0)
424        );
425        // Bitmap CFs get a larger write buffer.
426        assert_eq!(
427            cfg.column_family[transaction_bitmap::NAME].write_buffer_size_mb,
428            Some(256)
429        );
430        assert_eq!(
431            cfg.column_family[transaction_bitmap::NAME].periodic_compaction_seconds,
432            Some(604_800)
433        );
434        assert_eq!(
435            cfg.column_family[event_bitmap::NAME].periodic_compaction_seconds,
436            Some(604_800)
437        );
438        assert_eq!(cfg.default_cf.periodic_compaction_seconds, None);
439        for (name, tuning) in &cfg.column_family {
440            if tuning.periodic_compaction_seconds.is_some() {
441                assert!(
442                    [transaction_bitmap::NAME, event_bitmap::NAME].contains(&name.as_str()),
443                    "periodic compaction unexpectedly configured for {name}",
444                );
445            }
446        }
447    }
448}