Skip to main content

tuwunel_database/engine/
cf_opts.rs

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
16/// Adjust options for the specific column by name. Provide the result of
17/// db_options() as the argument to this function and use the return value in
18/// the arguments to open the specific column.
19pub(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_preview_retention(&mut desc, config.url_preview_cache_ttl);
33	set_table_options(&mut opts, &desc, cache)?;
34
35	opts.set_min_write_buffer_number(1);
36	opts.set_max_write_buffer_number(2);
37	opts.set_write_buffer_size(desc.write_size);
38
39	opts.set_target_file_size_base(desc.file_size);
40	opts.set_target_file_size_multiplier(desc.file_shape);
41
42	opts.set_max_compaction_bytes(desc.compaction_size);
43	opts.set_level_zero_file_num_compaction_trigger(desc.level0_width);
44	opts.set_level_compaction_dynamic_level_bytes(false);
45	opts.set_ttl(desc.ttl);
46
47	opts.set_max_bytes_for_level_base(desc.level_size);
48	opts.set_max_bytes_for_level_multiplier(1.0);
49	opts.set_max_bytes_for_level_multiplier_additional(&desc.level_shape);
50
51	opts.set_disable_auto_compactions(desc.ignored);
52	opts.set_compaction_style(desc.compaction);
53	opts.set_compaction_pri(desc.compaction_pri);
54	opts.set_universal_compaction_options(&uc_options(&desc));
55	opts.set_fifo_compaction_options(&fifo_options(&desc));
56
57	let compression_shape: Vec<_> = desc
58		.compression_shape
59		.into_iter()
60		.map(|val| (val > 0).then_some(desc.compression))
61		.map(|val| val.unwrap_or(CompressionType::None))
62		.collect();
63
64	opts.set_compression_type(desc.compression);
65	opts.set_compression_per_level(compression_shape.as_slice());
66	opts.set_compression_options(-14, desc.compression_level, 0, 0); // -14 w_bits used by zlib.
67	if let Some(&bottommost_level) = desc.bottommost_level.as_ref() {
68		opts.set_bottommost_compression_type(desc.compression);
69		opts.set_bottommost_zstd_max_train_bytes(0, true);
70		opts.set_bottommost_compression_options(
71			-14, // -14 w_bits is only read by zlib.
72			bottommost_level,
73			0,
74			0,
75			true,
76		);
77	}
78
79	opts.set_options_from_string("{{arena_block_size=2097152;}}")
80		.map_err(map_err)?;
81
82	#[cfg(debug_assertions)]
83	opts.set_options_from_string(
84		"{{paranoid_checks=true;paranoid_file_checks=true;force_consistency_checks=true;\
85		 verify_sst_unique_id_in_manifest=true;}}",
86	)
87	.map_err(map_err)?;
88
89	Ok(opts)
90}
91
92fn set_table_options(opts: &mut Options, desc: &Descriptor, cache: Option<&Cache>) -> Result {
93	let mut table = table_options(desc, cache.is_some());
94
95	if let Some(cache) = cache {
96		table.set_block_cache(cache);
97	} else {
98		table.disable_cache();
99	}
100
101	let prepopulate = if desc.write_to_cache { "kFlushOnly" } else { "kDisable" };
102
103	let string = format!(
104		"{{block_based_table_factory={{num_file_reads_for_auto_readahead={0};\
105		 max_auto_readahead_size={1};initial_auto_readahead_size={2};\
106		 enable_index_compression={3};prepopulate_block_cache={4}}}}}",
107		desc.auto_readahead_thresh,
108		desc.auto_readahead_max,
109		desc.auto_readahead_init,
110		desc.compressed_index,
111		prepopulate,
112	);
113
114	opts.set_options_from_string(&string)
115		.map_err(map_err)?;
116
117	opts.set_block_based_table_factory(&table);
118
119	Ok(())
120}
121
122fn set_compression(desc: &mut Descriptor, config: &Config) {
123	// A column opting out of compression is never overridden by the config.
124	if desc.compression == CompressionType::None {
125		return;
126	}
127
128	desc.compression = match config.rocksdb_compression_algo.as_ref() {
129		| "snappy" => CompressionType::Snappy,
130		| "zlib" => CompressionType::Zlib,
131		| "bz2" => CompressionType::Bz2,
132		| "lz4" => CompressionType::Lz4,
133		| "lz4hc" => CompressionType::Lz4hc,
134		| "none" => CompressionType::None,
135		| _ => CompressionType::Zstd,
136	};
137
138	let can_override_level = config.rocksdb_compression_level == SENTINEL_COMPRESSION_LEVEL
139		&& desc.compression == CompressionType::Zstd;
140
141	if !can_override_level {
142		desc.compression_level = config.rocksdb_compression_level;
143	}
144
145	let can_override_bottom = config.rocksdb_bottommost_compression_level
146		== SENTINEL_COMPRESSION_LEVEL
147		&& desc.compression == CompressionType::Zstd;
148
149	if !can_override_bottom {
150		desc.bottommost_level = Some(config.rocksdb_bottommost_compression_level);
151	}
152
153	if !config.rocksdb_bottommost_compression {
154		desc.bottommost_level = None;
155	}
156}
157
158/// Raise the preview columns' FIFO ttl to the configured preview lifetime
159/// when it exceeds their floor.
160///
161/// A preview evicted before it expires is fetched from the origin again, and
162/// one whose lazy media mapping went with it serves an `mxc://` URI that no
163/// longer resolves, so the configured lifetime raises both columns. Only the
164/// time arm moves: `limit_size` still caps each column, and evicts oldest
165/// first once reached.
166pub(super) fn set_preview_retention(desc: &mut Descriptor, configured: u64) {
167	if matches!(desc.name, "url_preview" | "mediaid_lazy") {
168		desc.ttl = desc.ttl.max(configured);
169	}
170}
171
172fn fifo_options(desc: &Descriptor) -> FifoCompactOptions {
173	let mut opts = FifoCompactOptions::default();
174	opts.set_max_table_files_size(desc.limit_size);
175	opts.set_allow_compaction(true);
176
177	opts
178}
179
180fn uc_options(desc: &Descriptor) -> UniversalCompactOptions {
181	let mut opts = UniversalCompactOptions::default();
182	opts.set_stop_style(UniversalCompactionStopStyle::Total);
183	opts.set_min_merge_width(desc.merge_width.0);
184	opts.set_max_merge_width(desc.merge_width.1);
185	opts.set_max_size_amplification_percent(10000);
186	opts.set_compression_size_percent(-1);
187	opts.set_size_ratio(1);
188
189	opts
190}
191
192fn table_options(desc: &Descriptor, has_cache: bool) -> BlockBasedOptions {
193	let mut opts = BlockBasedOptions::default();
194
195	opts.set_block_size(desc.block_size);
196	opts.set_metadata_block_size(desc.index_size);
197
198	opts.set_cache_index_and_filter_blocks(has_cache);
199	opts.set_pin_top_level_index_and_filter(false);
200	opts.set_pin_l0_filter_and_index_blocks_in_cache(false);
201	opts.set_partition_pinning_tier(BlockBasedPinningTier::None);
202	opts.set_unpartitioned_pinning_tier(BlockBasedPinningTier::None);
203	opts.set_top_level_index_pinning_tier(BlockBasedPinningTier::None);
204
205	opts.set_partition_filters(true);
206	opts.set_use_delta_encoding(false);
207	opts.set_index_type(BlockBasedIndexType::TwoLevelIndexSearch);
208
209	opts.set_data_block_index_type(match desc.block_index_hashing {
210		| None if desc.index_size > 512 => DataBlockIndexType::BinaryAndHash,
211		| Some(enable) if enable => DataBlockIndexType::BinaryAndHash,
212		| Some(_) | None => DataBlockIndexType::BinarySearch,
213	});
214
215	opts
216}
217
218fn get_cache(ctx: &Context, desc: &Descriptor) -> Option<Cache> {
219	if desc.dropped {
220		return None;
221	}
222
223	// Capacity is configured in entries, converted to bytes against the
224	// descriptor's size hints; the lookup by column name is legacy-compat.
225	let config = &ctx.server.config;
226	let cap = match desc.name {
227		| "eventid_pduid" => Some(config.eventid_pdu_cache_capacity),
228		| "eventid_shorteventid" => Some(config.eventidshort_cache_capacity),
229		| "eventid_backoff" => Some(config.eventid_backoff_cache_capacity),
230		| "shorteventid_eventid" => Some(config.shorteventid_cache_capacity),
231		| "shortstatekey_statekey" => Some(config.shortstatekey_cache_capacity),
232		| "statekey_shortstatekey" => Some(config.statekeyshort_cache_capacity),
233		| "servernameevent_data" => Some(config.servernameevent_data_cache_capacity),
234		| "servername_destination" | "servername_override" =>
235			Some(config.resolver_cache_capacity),
236		| "servername_status" => Some(config.servername_status_cache_capacity),
237		| "mediaid_lazycontent" => Some(config.mediaid_lazycontent_cache_capacity),
238		| "pduid_pdu" | "eventid_outlierpdu" => Some(config.pdu_cache_capacity),
239		| "authchainkey_authchain" => Some(config.auth_chain_cache_capacity),
240		| _ => None,
241	}
242	.map(TryInto::try_into)
243	.transpose()
244	.expect("u32 to usize");
245
246	let ent_size: usize = desc
247		.key_size_hint
248		.unwrap_or_default()
249		.expected_add(desc.val_size_hint.unwrap_or_default());
250
251	let size = match cap {
252		| Some(cap) => cache_size(config, cap, ent_size),
253		| _ => desc.cache_size,
254	};
255
256	let shard_bits: i32 = desc
257		.cache_shards
258		.ilog2()
259		.try_into()
260		.expect("u32 to i32 conversion");
261
262	debug_assert!(shard_bits <= 10, "cache shards probably too large");
263	let mut cache_opts = LruCacheOptions::default();
264	cache_opts.set_num_shard_bits(shard_bits);
265	cache_opts.set_capacity(size);
266
267	let mut caches = ctx.col_cache.lock().expect("locked");
268	register_pool(&mut caches, desc, || Cache::new_lru_cache_opts(&cache_opts))
269}
270
271/// Returns the cache for `desc`'s pool; `build_cache` is invoked only when a
272/// new pool must be created.
273pub(crate) fn register_pool(
274	caches: &mut ColCaches,
275	desc: &Descriptor,
276	build_cache: impl FnOnce() -> Cache,
277) -> Option<Cache> {
278	match desc.cache_disp {
279		| CacheDisp::Unique if desc.cache_size == 0 => None,
280		| CacheDisp::Unique => {
281			let cache = build_cache();
282			caches.insert(desc.name, ColCache {
283				cache: cache.clone(),
284				participants: vec![desc.name],
285			});
286
287			Some(cache)
288		},
289
290		| CacheDisp::SharedWith(other) => Some(match caches.get_mut(other) {
291			| Some(pool) => {
292				pool.participants.push(desc.name);
293				pool.cache.clone()
294			},
295			| None => {
296				let new = build_cache();
297				caches.insert(desc.name, ColCache {
298					cache: new.clone(),
299					participants: vec![desc.name],
300				});
301
302				new
303			},
304		}),
305
306		| CacheDisp::Shared => {
307			let pool = caches
308				.get_mut(SHARED_POOL)
309				.expect("shared cache must already exist");
310
311			pool.participants.push(desc.name);
312			Some(pool.cache.clone())
313		},
314	}
315}
316
317pub(crate) fn cache_size(config: &Config, base_size: u32, entity_size: usize) -> usize {
318	cache_size_f64(config, f64::from(base_size), entity_size)
319}
320
321#[expect(
322	clippy::as_conversions,
323	clippy::cast_sign_loss,
324	clippy::cast_possible_truncation
325)]
326pub(crate) fn cache_size_f64(config: &Config, base_size: f64, entity_size: usize) -> usize {
327	let ents = base_size * config.cache_capacity_modifier;
328
329	(ents as usize)
330		.checked_mul(entity_size)
331		.ok_or_else(|| err!(Config("cache_capacity_modifier", "Cache size is too large.")))
332		.expect("invalid cache size")
333}