Skip to main content

tuwunel_database/engine/
events.rs

1use rocksdb::{
2	Env,
3	event_listener::{
4		CompactionJobInfo, DBBackgroundErrorReason, DBWriteStallCondition, EventListener,
5		FlushJobInfo, IngestionInfo, MemTableInfo, MutableStatus, SubcompactionJobInfo,
6		WriteStallInfo,
7	},
8};
9use tuwunel_core::{Config, debug, debug::INFO_SPAN_LEVEL, debug_info, error, info, warn};
10
11/// Bridges RocksDB background events into structured tracing.
12///
13/// The listener records stalls, compactions, flushes, memtable seals, and
14/// external file ingestion without retaining per-event state.
15pub(super) struct Events;
16
17impl Events {
18	pub(super) fn new(_config: &Config, _env: &Env) -> Self { Self {} }
19}
20
21impl EventListener for Events {
22	#[tracing::instrument(name = "error", level = "error", skip_all)]
23	fn on_background_error(&self, reason: DBBackgroundErrorReason, status: MutableStatus) {
24		error!(
25			?reason,
26			severity = ?status.severity(),
27			error = ?status.result().as_ref().err(),
28			"Critical RocksDB Error",
29		);
30	}
31
32	#[tracing::instrument(name = "stall", level = "warn", skip_all)]
33	fn on_stall_conditions_changed(&self, info: &WriteStallInfo) {
34		let col = info.cf_name();
35		let col = col
36			.as_deref()
37			.map(str::from_utf8)
38			.expect("column has a name")
39			.expect("column name is valid utf8");
40
41		let prev = info.prev();
42		match info.cur() {
43			| DBWriteStallCondition::KStopped => {
44				error!(?col, ?prev, "Database Stalled");
45			},
46			| DBWriteStallCondition::KDelayed if prev == DBWriteStallCondition::KStopped => {
47				warn!(?col, ?prev, "Database Stall Recovering");
48			},
49			| DBWriteStallCondition::KDelayed => {
50				warn!(?col, ?prev, "Database Stalling");
51			},
52			| DBWriteStallCondition::KNormal
53				if prev == DBWriteStallCondition::KStopped
54					|| prev == DBWriteStallCondition::KDelayed =>
55			{
56				info!(?col, ?prev, "Database Stall Recovered");
57			},
58			| DBWriteStallCondition::KNormal => {
59				debug!(?col, ?prev, "Database Normal");
60			},
61		}
62	}
63
64	#[tracing::instrument(
65		name = "compaction",
66		level = INFO_SPAN_LEVEL,
67		skip_all,
68	)]
69	fn on_compaction_begin(&self, info: &CompactionJobInfo) {
70		let col = info.cf_name();
71		let col = col
72			.as_deref()
73			.map(str::from_utf8)
74			.expect("column has a name")
75			.expect("column name is valid utf8");
76
77		let level = (info.base_input_level(), info.output_level());
78		let records = (info.input_records(), info.output_records());
79		let bytes = (info.total_input_bytes(), info.total_output_bytes());
80		let files = (
81			info.input_file_count(),
82			info.output_file_count(),
83			info.num_input_files_at_output_level(),
84		);
85
86		debug!(
87			status = ?info.status(),
88			?level,
89			?files,
90			?records,
91			?bytes,
92			micros = info.elapsed_micros(),
93			errs = info.num_corrupt_keys(),
94			reason = ?info.compaction_reason(),
95			?col,
96			"Compaction Starting",
97		);
98	}
99
100	#[tracing::instrument(
101		name = "compaction",
102		level = INFO_SPAN_LEVEL,
103		skip_all,
104	)]
105	fn on_compaction_completed(&self, info: &CompactionJobInfo) {
106		let col = info.cf_name();
107		let col = col
108			.as_deref()
109			.map(str::from_utf8)
110			.expect("column has a name")
111			.expect("column name is valid utf8");
112
113		let level = (info.base_input_level(), info.output_level());
114		let records = (info.input_records(), info.output_records());
115		let bytes = (info.total_input_bytes(), info.total_output_bytes());
116		let files = (
117			info.input_file_count(),
118			info.output_file_count(),
119			info.num_input_files_at_output_level(),
120		);
121
122		debug_info!(
123			status = ?info.status(),
124			?level,
125			?files,
126			?records,
127			?bytes,
128			micros = info.elapsed_micros(),
129			errs = info.num_corrupt_keys(),
130			reason = ?info.compaction_reason(),
131			?col,
132			"Compaction Complete",
133		);
134	}
135
136	#[tracing::instrument(name = "compaction", level = "debug", skip_all)]
137	fn on_subcompaction_begin(&self, info: &SubcompactionJobInfo) {
138		let col = info.cf_name();
139		let col = col
140			.as_deref()
141			.map(str::from_utf8)
142			.expect("column has a name")
143			.expect("column name is valid utf8");
144
145		let level = (info.base_input_level(), info.output_level());
146
147		debug!(
148			status = ?info.status(),
149			?level,
150			tid = info.thread_id(),
151			reason = ?info.compaction_reason(),
152			?col,
153			"Compaction Starting",
154		);
155	}
156
157	#[tracing::instrument(name = "compaction", level = "debug", skip_all)]
158	fn on_subcompaction_completed(&self, info: &SubcompactionJobInfo) {
159		let col = info.cf_name();
160		let col = col
161			.as_deref()
162			.map(str::from_utf8)
163			.expect("column has a name")
164			.expect("column name is valid utf8");
165
166		let level = (info.base_input_level(), info.output_level());
167
168		debug!(
169			status = ?info.status(),
170			?level,
171			tid = info.thread_id(),
172			reason = ?info.compaction_reason(),
173			?col,
174			"Compaction Complete",
175		);
176	}
177
178	#[tracing::instrument(
179		name = "flush",
180		level = INFO_SPAN_LEVEL,
181		skip_all,
182	)]
183	fn on_flush_begin(&self, info: &FlushJobInfo) {
184		let col = info.cf_name();
185		let col = col
186			.as_deref()
187			.map(str::from_utf8)
188			.expect("column has a name")
189			.expect("column name is valid utf8");
190
191		debug!(
192			seq_start = info.smallest_seqno(),
193			seq_end = info.largest_seqno(),
194			slow = info.triggered_writes_slowdown(),
195			stop = info.triggered_writes_stop(),
196			reason = ?info.flush_reason(),
197			?col,
198			"Flush Starting",
199		);
200	}
201
202	#[tracing::instrument(
203		name = "flush",
204		level = INFO_SPAN_LEVEL,
205		skip_all,
206	)]
207	fn on_flush_completed(&self, info: &FlushJobInfo) {
208		let col = info.cf_name();
209		let col = col
210			.as_deref()
211			.map(str::from_utf8)
212			.expect("column has a name")
213			.expect("column name is valid utf8");
214
215		debug_info!(
216			seq_start = info.smallest_seqno(),
217			seq_end = info.largest_seqno(),
218			slow = info.triggered_writes_slowdown(),
219			stop = info.triggered_writes_stop(),
220			reason = ?info.flush_reason(),
221			?col,
222			"Flush Complete",
223		);
224	}
225
226	#[tracing::instrument(
227		name = "memtable",
228		level = INFO_SPAN_LEVEL,
229		skip_all,
230	)]
231	fn on_memtable_sealed(&self, info: &MemTableInfo) {
232		let col = info.cf_name();
233		let col = col
234			.as_deref()
235			.map(str::from_utf8)
236			.expect("column has a name")
237			.expect("column name is valid utf8");
238
239		debug_info!(
240			seq_first = info.first_seqno(),
241			seq_early = info.earliest_seqno(),
242			ents = info.num_entries(),
243			dels = info.num_deletes(),
244			?col,
245			"Buffer Filled",
246		);
247	}
248
249	fn on_external_file_ingested(&self, _info: &IngestionInfo) {
250		unimplemented!();
251	}
252}