1use rocksdb::{
2 BlockBasedIndexType, BlockBasedOptions, BlockBasedPinningTier, Cache,
3 DBCompressionType as CompressionType, DataBlockIndexType, FifoCompactOptions,
4 LruCacheOptions, Options, UniversalCompactOptions, UniversalCompactionStopStyle,
5};
6use tuwunel_core::{Config, Result, err, utils::math::Expected};
7
8use super::{
9 context::{ColCache, ColCaches, SHARED_POOL},
10 descriptor::{CacheDisp, Descriptor},
11};
12use crate::{Context, util::map_err};
13
14pub(super) const SENTINEL_COMPRESSION_LEVEL: i32 = 32767;
15
16pub(crate) fn cf_options(ctx: &Context, opts: Options, desc: &Descriptor) -> Result<Options> {
20 let cache = get_cache(ctx, desc);
21 let config = &ctx.server.config;
22 descriptor_cf_options(opts, *desc, config, cache.as_ref())
23}
24
25fn descriptor_cf_options(
26 mut opts: Options,
27 mut desc: Descriptor,
28 config: &Config,
29 cache: Option<&Cache>,
30) -> Result<Options> {
31 set_compression(&mut desc, config);
32 set_table_options(&mut opts, &desc, cache)?;
33
34 opts.set_min_write_buffer_number(1);
35 opts.set_max_write_buffer_number(2);
36 opts.set_write_buffer_size(desc.write_size);
37
38 opts.set_target_file_size_base(desc.file_size);
39 opts.set_target_file_size_multiplier(desc.file_shape);
40
41 opts.set_level_zero_file_num_compaction_trigger(desc.level0_width);
42 opts.set_level_compaction_dynamic_level_bytes(false);
43 opts.set_ttl(desc.ttl);
44
45 opts.set_max_bytes_for_level_base(desc.level_size);
46 opts.set_max_bytes_for_level_multiplier(1.0);
47 opts.set_max_bytes_for_level_multiplier_additional(&desc.level_shape);
48
49 opts.set_disable_auto_compactions(desc.ignored);
50 opts.set_compaction_style(desc.compaction);
51 opts.set_compaction_pri(desc.compaction_pri);
52 opts.set_universal_compaction_options(&uc_options(&desc));
53 opts.set_fifo_compaction_options(&fifo_options(&desc));
54
55 let compression_shape: Vec<_> = desc
56 .compression_shape
57 .into_iter()
58 .map(|val| (val > 0).then_some(desc.compression))
59 .map(|val| val.unwrap_or(CompressionType::None))
60 .collect();
61
62 opts.set_compression_type(desc.compression);
63 opts.set_compression_per_level(compression_shape.as_slice());
64 opts.set_compression_options(-14, desc.compression_level, 0, 0); if let Some(&bottommost_level) = desc.bottommost_level.as_ref() {
66 opts.set_bottommost_compression_type(desc.compression);
67 opts.set_bottommost_zstd_max_train_bytes(0, true);
68 opts.set_bottommost_compression_options(
69 -14, bottommost_level,
71 0,
72 0,
73 true,
74 );
75 }
76
77 opts.set_options_from_string("{{arena_block_size=2097152;}}")
78 .map_err(map_err)?;
79
80 #[cfg(debug_assertions)]
81 opts.set_options_from_string(
82 "{{paranoid_checks=true;paranoid_file_checks=true;force_consistency_checks=true;\
83 verify_sst_unique_id_in_manifest=true;}}",
84 )
85 .map_err(map_err)?;
86
87 Ok(opts)
88}
89
90fn set_table_options(opts: &mut Options, desc: &Descriptor, cache: Option<&Cache>) -> Result {
91 let mut table = table_options(desc, cache.is_some());
92
93 if let Some(cache) = cache {
94 table.set_block_cache(cache);
95 } else {
96 table.disable_cache();
97 }
98
99 let prepopulate = if desc.write_to_cache { "kFlushOnly" } else { "kDisable" };
100
101 let string = format!(
102 "{{block_based_table_factory={{num_file_reads_for_auto_readahead={0};\
103 max_auto_readahead_size={1};initial_auto_readahead_size={2};\
104 enable_index_compression={3};prepopulate_block_cache={4}}}}}",
105 desc.auto_readahead_thresh,
106 desc.auto_readahead_max,
107 desc.auto_readahead_init,
108 desc.compressed_index,
109 prepopulate,
110 );
111
112 opts.set_options_from_string(&string)
113 .map_err(map_err)?;
114
115 opts.set_block_based_table_factory(&table);
116
117 Ok(())
118}
119
120fn set_compression(desc: &mut Descriptor, config: &Config) {
121 if desc.compression == CompressionType::None {
123 return;
124 }
125
126 desc.compression = match config.rocksdb_compression_algo.as_ref() {
127 | "snappy" => CompressionType::Snappy,
128 | "zlib" => CompressionType::Zlib,
129 | "bz2" => CompressionType::Bz2,
130 | "lz4" => CompressionType::Lz4,
131 | "lz4hc" => CompressionType::Lz4hc,
132 | "none" => CompressionType::None,
133 | _ => CompressionType::Zstd,
134 };
135
136 let can_override_level = config.rocksdb_compression_level == SENTINEL_COMPRESSION_LEVEL
137 && desc.compression == CompressionType::Zstd;
138
139 if !can_override_level {
140 desc.compression_level = config.rocksdb_compression_level;
141 }
142
143 let can_override_bottom = config.rocksdb_bottommost_compression_level
144 == SENTINEL_COMPRESSION_LEVEL
145 && desc.compression == CompressionType::Zstd;
146
147 if !can_override_bottom {
148 desc.bottommost_level = Some(config.rocksdb_bottommost_compression_level);
149 }
150
151 if !config.rocksdb_bottommost_compression {
152 desc.bottommost_level = None;
153 }
154}
155
156fn fifo_options(desc: &Descriptor) -> FifoCompactOptions {
157 let mut opts = FifoCompactOptions::default();
158 opts.set_max_table_files_size(desc.limit_size);
159 opts.set_allow_compaction(true);
160
161 opts
162}
163
164fn uc_options(desc: &Descriptor) -> UniversalCompactOptions {
165 let mut opts = UniversalCompactOptions::default();
166 opts.set_stop_style(UniversalCompactionStopStyle::Total);
167 opts.set_min_merge_width(desc.merge_width.0);
168 opts.set_max_merge_width(desc.merge_width.1);
169 opts.set_max_size_amplification_percent(10000);
170 opts.set_compression_size_percent(-1);
171 opts.set_size_ratio(1);
172
173 opts
174}
175
176fn table_options(desc: &Descriptor, has_cache: bool) -> BlockBasedOptions {
177 let mut opts = BlockBasedOptions::default();
178
179 opts.set_block_size(desc.block_size);
180 opts.set_metadata_block_size(desc.index_size);
181
182 opts.set_cache_index_and_filter_blocks(has_cache);
183 opts.set_pin_top_level_index_and_filter(false);
184 opts.set_pin_l0_filter_and_index_blocks_in_cache(false);
185 opts.set_partition_pinning_tier(BlockBasedPinningTier::None);
186 opts.set_unpartitioned_pinning_tier(BlockBasedPinningTier::None);
187 opts.set_top_level_index_pinning_tier(BlockBasedPinningTier::None);
188
189 opts.set_partition_filters(true);
190 opts.set_use_delta_encoding(false);
191 opts.set_index_type(BlockBasedIndexType::TwoLevelIndexSearch);
192
193 opts.set_data_block_index_type(match desc.block_index_hashing {
194 | None if desc.index_size > 512 => DataBlockIndexType::BinaryAndHash,
195 | Some(enable) if enable => DataBlockIndexType::BinaryAndHash,
196 | Some(_) | None => DataBlockIndexType::BinarySearch,
197 });
198
199 opts
200}
201
202fn get_cache(ctx: &Context, desc: &Descriptor) -> Option<Cache> {
203 if desc.dropped {
204 return None;
205 }
206
207 let config = &ctx.server.config;
210 let cap = match desc.name {
211 | "eventid_pduid" => Some(config.eventid_pdu_cache_capacity),
212 | "eventid_shorteventid" => Some(config.eventidshort_cache_capacity),
213 | "eventid_backoff" => Some(config.eventid_backoff_cache_capacity),
214 | "shorteventid_eventid" => Some(config.shorteventid_cache_capacity),
215 | "shortstatekey_statekey" => Some(config.shortstatekey_cache_capacity),
216 | "statekey_shortstatekey" => Some(config.statekeyshort_cache_capacity),
217 | "servernameevent_data" => Some(config.servernameevent_data_cache_capacity),
218 | "servername_destination" | "servername_override" =>
219 Some(config.resolver_cache_capacity),
220 | "servername_status" => Some(config.servername_status_cache_capacity),
221 | "mediaid_lazycontent" => Some(config.mediaid_lazycontent_cache_capacity),
222 | "pduid_pdu" | "eventid_outlierpdu" => Some(config.pdu_cache_capacity),
223 | "authchainkey_authchain" => Some(config.auth_chain_cache_capacity),
224 | _ => None,
225 }
226 .map(TryInto::try_into)
227 .transpose()
228 .expect("u32 to usize");
229
230 let ent_size: usize = desc
231 .key_size_hint
232 .unwrap_or_default()
233 .expected_add(desc.val_size_hint.unwrap_or_default());
234
235 let size = match cap {
236 | Some(cap) => cache_size(config, cap, ent_size),
237 | _ => desc.cache_size,
238 };
239
240 let shard_bits: i32 = desc
241 .cache_shards
242 .ilog2()
243 .try_into()
244 .expect("u32 to i32 conversion");
245
246 debug_assert!(shard_bits <= 10, "cache shards probably too large");
247 let mut cache_opts = LruCacheOptions::default();
248 cache_opts.set_num_shard_bits(shard_bits);
249 cache_opts.set_capacity(size);
250
251 let mut caches = ctx.col_cache.lock().expect("locked");
252 register_pool(&mut caches, desc, || Cache::new_lru_cache_opts(&cache_opts))
253}
254
255pub(crate) fn register_pool(
258 caches: &mut ColCaches,
259 desc: &Descriptor,
260 build_cache: impl FnOnce() -> Cache,
261) -> Option<Cache> {
262 match desc.cache_disp {
263 | CacheDisp::Unique if desc.cache_size == 0 => None,
264 | CacheDisp::Unique => {
265 let cache = build_cache();
266 caches.insert(desc.name, ColCache {
267 cache: cache.clone(),
268 participants: vec![desc.name],
269 });
270
271 Some(cache)
272 },
273
274 | CacheDisp::SharedWith(other) => Some(match caches.get_mut(other) {
275 | Some(pool) => {
276 pool.participants.push(desc.name);
277 pool.cache.clone()
278 },
279 | None => {
280 let new = build_cache();
281 caches.insert(desc.name, ColCache {
282 cache: new.clone(),
283 participants: vec![desc.name],
284 });
285
286 new
287 },
288 }),
289
290 | CacheDisp::Shared => {
291 let pool = caches
292 .get_mut(SHARED_POOL)
293 .expect("shared cache must already exist");
294
295 pool.participants.push(desc.name);
296 Some(pool.cache.clone())
297 },
298 }
299}
300
301pub(crate) fn cache_size(config: &Config, base_size: u32, entity_size: usize) -> usize {
302 cache_size_f64(config, f64::from(base_size), entity_size)
303}
304
305#[expect(
306 clippy::as_conversions,
307 clippy::cast_sign_loss,
308 clippy::cast_possible_truncation
309)]
310pub(crate) fn cache_size_f64(config: &Config, base_size: f64, entity_size: usize) -> usize {
311 let ents = base_size * config.cache_capacity_modifier;
312
313 (ents as usize)
314 .checked_mul(entity_size)
315 .ok_or_else(|| err!(Config("cache_capacity_modifier", "Cache size is too large.")))
316 .expect("invalid cache size")
317}