1pub 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
61pub struct RpcStoreSchema<R: Reader = Db> {
63 tx_seq_pruning_floor: Arc<AtomicU64>,
67
68 pub epochs: DbMap<epochs::Key, epochs::Value, R>,
71
72 pub checkpoint_summary: DbMap<checkpoint_summary::Key, checkpoint_summary::Value, R>,
76
77 pub checkpoint_contents: DbMap<checkpoint_contents::Key, checkpoint_contents::Value, R>,
80
81 pub checkpoint_seq_by_digest:
84 DbMap<checkpoint_seq_by_digest::Key, checkpoint_seq_by_digest::Value, R>,
85
86 pub transactions: DbMap<transactions::Key, transactions::Value, R>,
88
89 pub tx_seq_by_digest: DbMap<tx_seq_by_digest::Key, tx_seq_by_digest::Value, R>,
91
92 pub tx_metadata_by_seq: DbMap<tx_metadata_by_seq::Key, tx_metadata_by_seq::Value, R>,
96
97 pub effects: DbMap<effects::Key, effects::Value, R>,
100
101 pub events: DbMap<events::Key, events::Value, R>,
103
104 pub objects: DbMap<objects::Key, objects::Value, R>,
110
111 pub object_version_by_checkpoint:
118 DbMap<object_version_by_checkpoint::Key, object_version_by_checkpoint::Value, R>,
119
120 pub object_by_owner: DbMap<object_by_owner::Key, object_by_owner::Value, R>,
125
126 pub object_by_type: DbMap<object_by_type::Key, object_by_type::Value, R>,
129
130 pub balance: DbMap<balance::Key, balance::Value, R>,
135
136 pub package_versions: DbMap<package_versions::Key, package_versions::Value, R>,
139
140 pub transaction_bitmap: DbMap<transaction_bitmap::Key, transaction_bitmap::Value, R>,
144
145 pub event_bitmap: DbMap<event_bitmap::Key, event_bitmap::Value, R>,
149
150 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
269pub 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 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 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 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 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 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 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 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}