1use self::health::{Health, Policy};
2use crate::{
3 Spool, SpoolBackpressureTimeout, SpoolCallerDeadlineExceeded, SpoolEntry, SpoolId,
4 SpoolUnhealthyError,
5};
6use anyhow::Context;
7use async_trait::async_trait;
8use chrono::{DateTime, Utc};
9use flume::Sender;
10use kumo_prometheus::declare_metric;
11use rocksdb::perf::get_memory_usage_stats;
12use rocksdb::properties::{
13 ACTUAL_DELAYED_WRITE_RATE, BACKGROUND_ERRORS, COMPACTION_PENDING,
14 ESTIMATE_PENDING_COMPACTION_BYTES, IS_WRITE_STOPPED, NUM_RUNNING_COMPACTIONS,
15};
16use rocksdb::{
17 BottommostLevelCompaction, CompactOptions, DBCompressionType, ErrorKind, IteratorMode,
18 LogLevel, Options, WaitForCompactOptions, WriteBatch, WriteOptions, DB,
19};
20use serde::{Deserialize, Serialize};
21use std::path::{Path, PathBuf};
22use std::sync::{Arc, Weak};
23use std::time::{Duration, Instant};
24use tokio::runtime::Handle;
25use tokio::sync::{OwnedSemaphorePermit, Semaphore};
26use tokio::time::{sleep, timeout_at};
27
28mod health;
29
30#[derive(Serialize, Deserialize, Debug)]
31pub struct RocksSpoolParams {
32 pub increase_parallelism: Option<i32>,
33
34 pub optimize_level_style_compaction: Option<usize>,
35 pub optimize_universal_style_compaction: Option<usize>,
36 #[serde(default)]
37 pub paranoid_checks: bool,
38 #[serde(default)]
39 pub compression_type: DBCompressionTypeDef,
40
41 pub compaction_readahead_size: Option<usize>,
45
46 #[serde(default)]
47 pub level_compaction_dynamic_level_bytes: bool,
48
49 #[serde(default)]
50 pub max_open_files: Option<usize>,
51
52 #[serde(default)]
62 pub write_buffer_size: Option<usize>,
63
64 #[serde(default)]
71 pub level0_stop_writes_trigger: Option<i32>,
72
73 #[serde(default)]
74 pub log_level: LogLevelDef,
75
76 #[serde(default)]
79 pub memtable_huge_page_size: Option<usize>,
80
81 #[serde(
82 with = "duration_serde",
83 default = "RocksSpoolParams::default_log_file_time_to_roll"
84 )]
85 pub log_file_time_to_roll: Duration,
86
87 #[serde(
88 with = "duration_serde",
89 default = "RocksSpoolParams::default_obsolete_files_period"
90 )]
91 pub obsolete_files_period: Duration,
92
93 #[serde(default)]
94 pub limit_concurrent_stores: Option<usize>,
95 #[serde(default)]
96 pub limit_concurrent_loads: Option<usize>,
97 #[serde(default)]
98 pub limit_concurrent_removes: Option<usize>,
99
100 #[serde(
109 with = "duration_serde",
110 default = "RocksSpoolParams::default_store_deadline"
111 )]
112 pub store_deadline: Duration,
113
114 #[serde(
122 with = "duration_serde",
123 default = "RocksSpoolParams::default_error_latch_duration"
124 )]
125 pub error_latch_duration: Duration,
126
127 #[serde(
133 with = "duration_serde",
134 default = "RocksSpoolParams::default_error_unlatch_duration"
135 )]
136 pub error_unlatch_duration: Duration,
137
138 #[serde(default = "RocksSpoolParams::default_allow_error_unlatch")]
143 pub allow_error_unlatch: bool,
144}
145
146impl Default for RocksSpoolParams {
147 fn default() -> Self {
148 Self {
149 increase_parallelism: None,
150 optimize_level_style_compaction: None,
151 optimize_universal_style_compaction: None,
152 paranoid_checks: false,
153 compression_type: DBCompressionTypeDef::default(),
154 compaction_readahead_size: None,
155 level_compaction_dynamic_level_bytes: false,
156 max_open_files: None,
157 write_buffer_size: None,
158 level0_stop_writes_trigger: None,
159 log_level: LogLevelDef::default(),
160 memtable_huge_page_size: None,
161 log_file_time_to_roll: Self::default_log_file_time_to_roll(),
162 obsolete_files_period: Self::default_obsolete_files_period(),
163 limit_concurrent_stores: None,
164 limit_concurrent_loads: None,
165 limit_concurrent_removes: None,
166 store_deadline: Self::default_store_deadline(),
167 error_latch_duration: Self::default_error_latch_duration(),
168 error_unlatch_duration: Self::default_error_unlatch_duration(),
169 allow_error_unlatch: Self::default_allow_error_unlatch(),
170 }
171 }
172}
173
174impl RocksSpoolParams {
175 fn default_log_file_time_to_roll() -> Duration {
176 Duration::from_secs(86400)
177 }
178
179 fn default_obsolete_files_period() -> Duration {
180 Duration::from_secs(6 * 60 * 60)
181 }
182
183 fn default_store_deadline() -> Duration {
184 Duration::from_secs(30)
185 }
186
187 fn default_error_latch_duration() -> Duration {
188 Duration::from_secs(15)
189 }
190
191 fn default_error_unlatch_duration() -> Duration {
192 Duration::from_secs(5 * 60)
193 }
194
195 fn default_allow_error_unlatch() -> bool {
196 true
197 }
198}
199
200#[derive(Serialize, Deserialize, Debug)]
201pub enum DBCompressionTypeDef {
202 None,
203 Snappy,
204 Zlib,
205 Bz2,
206 Lz4,
207 Lz4hc,
208 Zstd,
209}
210
211impl From<DBCompressionTypeDef> for DBCompressionType {
212 fn from(val: DBCompressionTypeDef) -> Self {
213 match val {
214 DBCompressionTypeDef::None => DBCompressionType::None,
215 DBCompressionTypeDef::Snappy => DBCompressionType::Snappy,
216 DBCompressionTypeDef::Zlib => DBCompressionType::Zlib,
217 DBCompressionTypeDef::Bz2 => DBCompressionType::Bz2,
218 DBCompressionTypeDef::Lz4 => DBCompressionType::Lz4,
219 DBCompressionTypeDef::Lz4hc => DBCompressionType::Lz4hc,
220 DBCompressionTypeDef::Zstd => DBCompressionType::Zstd,
221 }
222 }
223}
224
225impl Default for DBCompressionTypeDef {
226 fn default() -> Self {
227 Self::Snappy
228 }
229}
230
231#[derive(Serialize, Deserialize, Debug)]
232pub enum LogLevelDef {
233 Debug,
234 Info,
235 Warn,
236 Error,
237 Fatal,
238 Header,
239}
240
241impl Default for LogLevelDef {
242 fn default() -> Self {
243 Self::Info
244 }
245}
246
247impl From<LogLevelDef> for LogLevel {
248 fn from(val: LogLevelDef) -> Self {
249 match val {
250 LogLevelDef::Debug => LogLevel::Debug,
251 LogLevelDef::Info => LogLevel::Info,
252 LogLevelDef::Warn => LogLevel::Warn,
253 LogLevelDef::Error => LogLevel::Error,
254 LogLevelDef::Fatal => LogLevel::Fatal,
255 LogLevelDef::Header => LogLevel::Header,
256 }
257 }
258}
259
260pub struct RocksSpool {
261 db: Arc<DB>,
262 runtime: Handle,
263 limit_concurrent_stores: Option<Arc<Semaphore>>,
264 limit_concurrent_loads: Option<Arc<Semaphore>>,
265 limit_concurrent_removes: Option<Arc<Semaphore>>,
266 health: Arc<Health>,
269 store_deadline: Duration,
270}
271
272const BACKOFF_INITIAL: Duration = Duration::from_micros(500);
276const BACKOFF_MAX: Duration = Duration::from_millis(50);
280
281fn select_effective_deadline(
286 caller_deadline: Option<Instant>,
287 spool_deadline: Instant,
288) -> (Instant, bool) {
289 match caller_deadline {
290 Some(c) if c < spool_deadline => (c, true),
291 _ => (spool_deadline, false),
292 }
293}
294
295impl RocksSpool {
296 fn timeout_err(&self, caller_wins: bool) -> anyhow::Error {
300 if caller_wins {
301 SpoolCallerDeadlineExceeded.into()
302 } else {
303 SpoolBackpressureTimeout {
304 deadline: self.store_deadline,
305 }
306 .into()
307 }
308 }
309
310 fn note_backpressure_timeout(&self) {
317 if self.health.record_foreground_error(false) {
318 tracing::error!(
319 "rocksdb at {}: a write timed out waiting for RocksDB to accept it. \
320 This is the first sign of trouble in a new incident; it does not by \
321 itself mean the spool has stopped accepting writes yet, but the \
322 load-shedding gate will close {:?} from now, at which point ingress \
323 will start rejecting traffic -- this happens even if no further \
324 errors occur. Investigate disk I/O and space on this host, and check \
325 the RocksDB LOG file in the spool directory for background errors.",
326 self.db.path().display(),
327 self.health.latch_duration(),
328 );
329 }
330 }
331
332 async fn acquire_store_permit(
336 &self,
337 permits: Option<Arc<Semaphore>>,
338 effective_deadline: Instant,
339 caller_wins: bool,
340 ) -> anyhow::Result<Option<OwnedSemaphorePermit>> {
341 let Some(s) = permits else {
342 return Ok(None);
343 };
344 match timeout_at(effective_deadline.into(), s.acquire_owned()).await {
345 Ok(r) => Ok(Some(r?)),
346 Err(_) => {
347 self.note_backpressure_timeout();
348 Err(self.timeout_err(caller_wins))
349 }
350 }
351 }
352
353 async fn write_with_backpressure(
359 &self,
360 opts: WriteOptions,
361 caller_deadline: Option<Instant>,
362 permits: Option<Arc<Semaphore>>,
363 apply: impl Fn(&mut WriteBatch),
364 ) -> anyhow::Result<()> {
365 if self.health.is_active() {
375 return Err(SpoolUnhealthyError.into());
376 }
377
378 let mut batch = WriteBatch::default();
379 apply(&mut batch);
380 match self.db.write_opt(batch, &opts) {
381 Ok(()) => return Ok(()),
382 Err(err) if err.kind() == ErrorKind::Incomplete => {}
383 Err(err) => {
384 record_foreground_error(&self.health, self.db.path(), &err);
385 return Err(err.into());
386 }
387 }
388
389 let spool_deadline = Instant::now() + self.store_deadline;
390 let (effective_deadline, caller_wins) =
391 select_effective_deadline(caller_deadline, spool_deadline);
392
393 let _permit = self
394 .acquire_store_permit(permits, effective_deadline, caller_wins)
395 .await?;
396
397 let mut backoff = BACKOFF_INITIAL;
398 loop {
399 if self.health.is_active() {
400 return Err(SpoolUnhealthyError.into());
401 }
402 if Instant::now() >= effective_deadline {
403 self.note_backpressure_timeout();
404 return Err(self.timeout_err(caller_wins));
405 }
406 sleep(backoff).await;
407 backoff = (backoff * 2).min(BACKOFF_MAX);
408
409 let mut batch = WriteBatch::default();
410 apply(&mut batch);
411 match self.db.write_opt(batch, &opts) {
412 Ok(()) => return Ok(()),
413 Err(err) if err.kind() == ErrorKind::Incomplete => continue,
414 Err(err) => {
415 record_foreground_error(&self.health, self.db.path(), &err);
416 return Err(err.into());
417 }
418 }
419 }
420 }
421
422 pub fn new(
423 path: &Path,
424 flush: bool,
425 params: Option<RocksSpoolParams>,
426 runtime: Handle,
427 ) -> anyhow::Result<Self> {
428 let mut opts = Options::default();
429 opts.set_use_fsync(flush);
430 opts.create_if_missing(true);
431 opts.set_keep_log_file_num(10);
433
434 let p = params.unwrap_or_default();
435 let policy = Policy {
436 latch_duration: p.error_latch_duration,
437 unlatch_duration: p.error_unlatch_duration,
438 allow_unlatch: p.allow_error_unlatch,
439 };
440 policy.validate()?;
441 if let Some(i) = p.increase_parallelism {
442 opts.increase_parallelism(i);
443 }
444 if let Some(i) = p.optimize_level_style_compaction {
445 opts.optimize_level_style_compaction(i);
446 }
447 if let Some(i) = p.optimize_universal_style_compaction {
448 opts.optimize_universal_style_compaction(i);
449 }
450 if let Some(i) = p.compaction_readahead_size {
451 opts.set_compaction_readahead_size(i);
452 }
453 if let Some(i) = p.max_open_files {
454 opts.set_max_open_files(i as _);
455 }
456 if let Some(i) = p.write_buffer_size {
457 opts.set_write_buffer_size(i);
458 }
459 if let Some(i) = p.level0_stop_writes_trigger {
460 opts.set_level_zero_stop_writes_trigger(i);
461 }
462 if let Some(i) = p.memtable_huge_page_size {
463 opts.set_memtable_huge_page_size(i);
464 }
465 opts.set_paranoid_checks(p.paranoid_checks);
466 opts.set_level_compaction_dynamic_level_bytes(p.level_compaction_dynamic_level_bytes);
467 opts.set_compression_type(p.compression_type.into());
468 opts.set_log_level(p.log_level.into());
469 opts.set_log_file_time_to_roll(p.log_file_time_to_roll.as_secs() as usize);
470 opts.set_delete_obsolete_files_period_micros(p.obsolete_files_period.as_micros() as u64);
471
472 let limit_concurrent_stores = p
473 .limit_concurrent_stores
474 .map(|n| Arc::new(Semaphore::new(n)));
475 let limit_concurrent_loads = p
476 .limit_concurrent_loads
477 .map(|n| Arc::new(Semaphore::new(n)));
478 let limit_concurrent_removes = p
479 .limit_concurrent_removes
480 .map(|n| Arc::new(Semaphore::new(n)));
481
482 #[cfg(unix)]
487 {
488 use std::os::unix::fs::DirBuilderExt;
489 std::fs::DirBuilder::new()
490 .recursive(true)
491 .mode(0o755)
492 .create(path)
493 .with_context(|| format!("creating spool directory {}", path.display()))?;
494 }
495 #[cfg(not(unix))]
496 std::fs::create_dir_all(path)
497 .with_context(|| format!("creating spool directory {}", path.display()))?;
498
499 dir_probe::probe_directory(path)
503 .with_context(|| format!("spool directory {} is not usable", path.display()))?;
504
505 let db = Arc::new(DB::open(&opts, path)?);
506 let health = Arc::new(Health::new(0, policy, format!("{}", path.display())));
510 let store_deadline = p.store_deadline;
511
512 tokio::spawn(metrics_monitor(
513 Arc::downgrade(&db),
514 Arc::downgrade(&health),
515 format!("{}", path.display()),
516 ));
517
518 Ok(Self {
519 db,
520 runtime,
521 limit_concurrent_stores,
522 limit_concurrent_loads,
523 limit_concurrent_removes,
524 health,
525 store_deadline,
526 })
527 }
528}
529
530#[async_trait]
531impl Spool for RocksSpool {
532 async fn load(&self, id: SpoolId) -> anyhow::Result<Vec<u8>> {
533 let permit = match self.limit_concurrent_loads.clone() {
534 Some(s) => Some(s.acquire_owned().await?),
535 None => None,
536 };
537 let db = self.db.clone();
538 let health = self.health.clone();
539 let db_path: PathBuf = self.db.path().to_owned();
540 tokio::task::Builder::new()
541 .name("rocksdb load")
542 .spawn_blocking_on(
543 move || {
544 let result = match db.get(id.as_bytes()) {
545 Ok(Some(v)) => v,
546 Ok(None) => {
547 drop(permit);
548 anyhow::bail!("no such key {id}");
549 }
550 Err(err) => {
551 record_foreground_error(&health, &db_path, &err);
555 drop(permit);
556 return Err(err.into());
557 }
558 };
559 drop(permit);
560 Ok(result)
561 },
562 &self.runtime,
563 )?
564 .await?
565 }
566
567 async fn store(
568 &self,
569 id: SpoolId,
570 data: Arc<Box<[u8]>>,
571 force_sync: bool,
572 deadline: Option<Instant>,
573 ) -> anyhow::Result<()> {
574 let mut opts = WriteOptions::default();
575 opts.set_sync(force_sync);
576 opts.set_no_slowdown(true);
577
578 self.write_with_backpressure(
579 opts,
580 deadline,
581 self.limit_concurrent_stores.clone(),
582 |batch| batch.put(id.as_bytes(), &*data),
583 )
584 .await
585 }
586
587 async fn remove(&self, id: SpoolId) -> anyhow::Result<()> {
588 let mut opts = WriteOptions::default();
589 opts.set_no_slowdown(true);
590
591 self.write_with_backpressure(opts, None, self.limit_concurrent_removes.clone(), |batch| {
592 batch.delete(id.as_bytes())
593 })
594 .await
595 }
596
597 async fn cleanup(&self) -> anyhow::Result<()> {
598 Ok(())
599 }
600
601 async fn compact(&self) -> anyhow::Result<()> {
602 let db = self.db.clone();
603 tokio::task::spawn_blocking(move || -> anyhow::Result<()> {
604 db.flush()?;
605 let mut opts = CompactOptions::default();
610 opts.set_bottommost_level_compaction(BottommostLevelCompaction::Force);
611 opts.set_exclusive_manual_compaction(true);
612 db.compact_range_opt::<&[u8], &[u8]>(None, None, &opts);
613 let wait_opts = WaitForCompactOptions::default();
617 db.wait_for_compact(&wait_opts)?;
618 Ok(())
619 })
620 .await?
621 }
622
623 async fn shutdown(&self) -> anyhow::Result<()> {
624 let db = self.db.clone();
625 tokio::task::spawn_blocking(move || db.cancel_all_background_work(true)).await?;
626 Ok(())
627 }
628
629 fn unhealthy_reason(&self) -> Option<&'static str> {
630 if self.health.is_active() {
631 Some("the spool is not accepting writes")
632 } else {
633 None
634 }
635 }
636
637 async fn advise_low_memory(&self) -> anyhow::Result<isize> {
638 let db = self.db.clone();
639 tokio::task::spawn_blocking(move || {
640 let usage_before = match get_memory_usage_stats(Some(&[&db]), None) {
641 Ok(stats) => {
642 let stats: Stats = stats.into();
643 tracing::debug!("pre-flush: {stats:#?}");
644 stats.total()
645 }
646 Err(err) => {
647 tracing::error!("error getting stats: {err:#}");
648 0
649 }
650 };
651
652 if let Err(err) = db.flush() {
653 tracing::error!("error flushing memory: {err:#}");
654 }
655
656 let usage_after = match get_memory_usage_stats(Some(&[&db]), None) {
657 Ok(stats) => {
658 let stats: Stats = stats.into();
659 tracing::debug!("post-flush: {stats:#?}");
660 stats.total()
661 }
662 Err(err) => {
663 tracing::error!("error getting stats: {err:#}");
664 0
665 }
666 };
667
668 Ok(usage_before - usage_after)
669 })
670 .await?
671 }
672
673 fn enumerate(
674 &self,
675 sender: Sender<SpoolEntry>,
676 start_time: DateTime<Utc>,
677 ) -> anyhow::Result<()> {
678 let db = Arc::clone(&self.db);
679 let health = self.health.clone();
680 let db_path: PathBuf = self.db.path().to_owned();
681 tokio::task::Builder::new()
682 .name("rocksdb enumerate")
683 .spawn_blocking_on(
684 move || {
685 let iter = db.iterator(IteratorMode::Start);
686 for entry in iter {
687 let (key, value) = match entry {
688 Ok(e) => e,
689 Err(err) => {
690 record_foreground_error(&health, &db_path, &err);
699 return Err(err.into());
700 }
701 };
702 let id = SpoolId::from_slice(&key)
703 .ok_or_else(|| anyhow::anyhow!("invalid spool id {key:?}"))?;
704
705 if id.created() >= start_time {
706 continue;
710 }
711
712 sender
713 .send(SpoolEntry::Item {
714 id,
715 data: value.to_vec(),
716 })
717 .map_err(|err| {
718 anyhow::anyhow!("failed to send SpoolEntry for {id}: {err:#}")
719 })?;
720 }
721 Ok::<(), anyhow::Error>(())
722 },
723 &self.runtime,
724 )?;
725 Ok(())
726 }
727}
728
729#[cfg(test)]
730mod test {
731 use super::*;
732 use k9::assert_equal;
733
734 #[test]
738 fn effective_deadline_prefers_the_sooner_caller_deadline() {
739 let base = Instant::now();
740 let spool = base + Duration::from_secs(10);
741
742 let sooner = base + Duration::from_secs(1);
743 assert_equal!(
744 select_effective_deadline(Some(sooner), spool),
745 (sooner, true)
746 );
747
748 let later = base + Duration::from_secs(20);
749 assert_equal!(
750 select_effective_deadline(Some(later), spool),
751 (spool, false)
752 );
753
754 assert_equal!(select_effective_deadline(None, spool), (spool, false));
755 }
756
757 #[tokio::test]
758 async fn rocks_spool() -> anyhow::Result<()> {
759 let location = tempfile::tempdir()?;
760 let spool = RocksSpool::new(location.path(), false, None, Handle::current())?;
761
762 {
763 let id1 = SpoolId::new();
764
765 assert_eq!(
767 format!("{:#}", spool.load(id1).await.unwrap_err()),
768 format!("no such key {id1}")
769 );
770 }
771
772 let mut ids = vec![];
774 for i in 0..100 {
775 let id = SpoolId::new();
776 spool
777 .store(
778 id,
779 Arc::new(format!("I am {i}").as_bytes().to_vec().into_boxed_slice()),
780 false,
781 None,
782 )
783 .await?;
784 ids.push(id);
785 }
786
787 for (i, &id) in ids.iter().enumerate() {
789 let data = spool.load(id).await?;
790 let text = String::from_utf8(data)?;
791 assert_eq!(text, format!("I am {i}"));
792 }
793
794 {
795 let (tx, rx) = flume::bounded(32);
797 spool.enumerate(tx, Utc::now())?;
798 let mut count = 0;
799
800 while let Ok(item) = rx.recv_async().await {
801 match item {
802 SpoolEntry::Item { id, data } => {
803 let i = ids
804 .iter()
805 .position(|&item| item == id)
806 .ok_or_else(|| anyhow::anyhow!("{id} not found in ids!"))?;
807
808 let text = String::from_utf8(data)?;
809 assert_eq!(text, format!("I am {i}"));
810
811 spool.remove(id).await?;
812 assert_eq!(
814 format!("{:#}", spool.load(id).await.unwrap_err()),
815 format!("no such key {id}")
816 );
817 count += 1;
818 }
819 SpoolEntry::Corrupt { id, error } => {
820 anyhow::bail!("Corrupt: {id}: {error}");
821 }
822 }
823 }
824
825 assert_eq!(count, 100);
826 }
827
828 for _ in 0..2 {
834 let (tx, rx) = flume::bounded(32);
836 spool.enumerate(tx, Utc::now())?;
837 let mut unexpected = vec![];
838
839 while let Ok(item) = rx.recv_async().await {
840 match item {
841 SpoolEntry::Item { id, .. } | SpoolEntry::Corrupt { id, .. } => {
842 unexpected.push(id)
843 }
844 }
845 }
846
847 assert_eq!(unexpected.len(), 0);
848 }
849
850 Ok(())
851 }
852}
853
854#[allow(unused)]
856#[derive(Debug)]
857struct Stats {
858 pub mem_table_total: u64,
859 pub mem_table_unflushed: u64,
860 pub mem_table_readers_total: u64,
861 pub cache_total: u64,
862}
863
864impl Stats {
865 fn total(&self) -> isize {
866 (self.mem_table_total + self.mem_table_readers_total + self.cache_total) as isize
867 }
868}
869
870impl From<rocksdb::perf::MemoryUsageStats> for Stats {
871 fn from(s: rocksdb::perf::MemoryUsageStats) -> Self {
872 Self {
873 mem_table_total: s.mem_table_total,
874 mem_table_unflushed: s.mem_table_unflushed,
875 mem_table_readers_total: s.mem_table_readers_total,
876 cache_total: s.cache_total,
877 }
878 }
879}
880
881fn property_u64(db: &DB, name: &rocksdb::properties::PropName) -> u64 {
886 db.property_int_value(name).ok().flatten().unwrap_or(0)
887}
888
889fn is_definitively_bad(err: &rocksdb::Error) -> bool {
892 matches!(err.kind(), ErrorKind::Corruption | ErrorKind::IOError)
893}
894
895fn record_foreground_error(health: &Health, path: &Path, err: &rocksdb::Error) {
899 let fatal = is_definitively_bad(err);
900 if health.record_foreground_error(fatal) {
901 if fatal {
902 tracing::error!(
903 "rocksdb at {}: a store, load, or enumeration hit a {:?} error: {}. \
904 This usually means a missing or corrupt SST file. The load-shedding gate is \
905 latching immediately: ingress will reject traffic and this spool will not accept \
906 further writes until an operator investigates the RocksDB LOG file in the spool \
907 directory, repairs or restores the affected files, and either waits for automatic \
908 recovery (if allow_error_unlatch is enabled) or restarts the process.",
909 path.display(),
910 err.kind(),
911 err.as_ref(),
912 );
913 } else {
914 tracing::error!(
915 "rocksdb at {}: a store, load, or enumeration returned an error: {}. This is the \
916 first sign of trouble in a new incident; it does not by itself mean the spool has \
917 stopped accepting writes yet, but the load-shedding gate will close {:?} from \
918 now, at which point ingress will start rejecting traffic -- this happens even if \
919 no further errors occur. Check the RocksDB LOG file in the spool directory for \
920 the underlying cause.",
921 path.display(),
922 err.as_ref(),
923 health.latch_duration(),
924 );
925 }
926 }
927}
928
929declare_metric! {
930static MEM_TABLE_TOTAL: IntGaugeVec(
935 "rocks_spool_mem_table_total",
936 &["path"]
937 );
938}
939
940declare_metric! {
941static MEM_TABLE_UNFLUSHED: IntGaugeVec(
946 "rocks_spool_mem_table_unflushed",
947 &["path"]
948 );
949}
950
951declare_metric! {
952static MEM_TABLE_READERS_TOTAL: IntGaugeVec(
957 "rocks_spool_mem_table_readers_total",
958 &["path"]
959 );
960}
961
962declare_metric! {
963static CACHE_TOTAL: IntGaugeVec(
968 "rocks_spool_cache_total",
969 &["path"]
970 );
971}
972
973declare_metric! {
974static BACKGROUND_ERRORS_METRIC: IntGaugeVec(
995 "rocks_spool_background_errors",
996 &["path"]
997 );
998}
999
1000declare_metric! {
1001static WRITE_STOPPED: IntGaugeVec(
1013 "rocks_spool_write_stopped",
1014 &["path"]
1015 );
1016}
1017
1018declare_metric! {
1019static LOAD_SHED_ACTIVE: IntGaugeVec(
1052 "rocks_spool_load_shed_active",
1053 &["path"]
1054 );
1055}
1056
1057declare_metric! {
1058static NUM_RUNNING_COMPACTIONS_METRIC: IntGaugeVec(
1071 "rocks_spool_num_running_compactions",
1072 &["path"]
1073 );
1074}
1075
1076declare_metric! {
1077static COMPACTION_PENDING_METRIC: IntGaugeVec(
1087 "rocks_spool_compaction_pending",
1088 &["path"]
1089 );
1090}
1091
1092declare_metric! {
1093static ESTIMATE_PENDING_COMPACTION_BYTES_METRIC: IntGaugeVec(
1109 "rocks_spool_estimate_pending_compaction_bytes",
1110 &["path"]
1111 );
1112}
1113
1114declare_metric! {
1115static ACTUAL_DELAYED_WRITE_RATE_METRIC: IntGaugeVec(
1129 "rocks_spool_actual_delayed_write_rate",
1130 &["path"]
1131 );
1132}
1133
1134async fn metrics_monitor(db: Weak<DB>, health: Weak<Health>, path: String) {
1135 let mem_table_total = MEM_TABLE_TOTAL
1136 .get_metric_with_label_values(&[path.as_str()])
1137 .unwrap();
1138 let mem_table_unflushed = MEM_TABLE_UNFLUSHED
1139 .get_metric_with_label_values(&[path.as_str()])
1140 .unwrap();
1141 let mem_table_readers_total = MEM_TABLE_READERS_TOTAL
1142 .get_metric_with_label_values(&[path.as_str()])
1143 .unwrap();
1144 let cache_total = CACHE_TOTAL
1145 .get_metric_with_label_values(&[path.as_str()])
1146 .unwrap();
1147 let background_errors = BACKGROUND_ERRORS_METRIC
1148 .get_metric_with_label_values(&[path.as_str()])
1149 .unwrap();
1150 let write_stopped = WRITE_STOPPED
1151 .get_metric_with_label_values(&[path.as_str()])
1152 .unwrap();
1153 let load_shed_active = LOAD_SHED_ACTIVE
1154 .get_metric_with_label_values(&[path.as_str()])
1155 .unwrap();
1156 let num_running_compactions = NUM_RUNNING_COMPACTIONS_METRIC
1157 .get_metric_with_label_values(&[path.as_str()])
1158 .unwrap();
1159 let compaction_pending = COMPACTION_PENDING_METRIC
1160 .get_metric_with_label_values(&[path.as_str()])
1161 .unwrap();
1162 let estimate_pending_compaction_bytes = ESTIMATE_PENDING_COMPACTION_BYTES_METRIC
1163 .get_metric_with_label_values(&[path.as_str()])
1164 .unwrap();
1165 let actual_delayed_write_rate = ACTUAL_DELAYED_WRITE_RATE_METRIC
1166 .get_metric_with_label_values(&[path.as_str()])
1167 .unwrap();
1168
1169 loop {
1170 match db.upgrade() {
1171 Some(db) => {
1172 match get_memory_usage_stats(Some(&[&db]), None) {
1173 Ok(stats) => {
1174 mem_table_total.set(stats.mem_table_total as i64);
1175 mem_table_unflushed.set(stats.mem_table_unflushed as i64);
1176 mem_table_readers_total.set(stats.mem_table_readers_total as i64);
1177 cache_total.set(stats.cache_total as i64);
1178 }
1179 Err(err) => {
1180 tracing::error!("error getting stats: {err:#}");
1181 }
1182 };
1183
1184 let bg = property_u64(&db, BACKGROUND_ERRORS);
1185 let stopped = property_u64(&db, IS_WRITE_STOPPED);
1186 let compaction_pending_now = property_u64(&db, COMPACTION_PENDING);
1187 let num_running = property_u64(&db, NUM_RUNNING_COMPACTIONS);
1188 background_errors.set(bg as i64);
1189 write_stopped.set(stopped as i64);
1190 num_running_compactions.set(num_running as i64);
1191 compaction_pending.set(compaction_pending_now as i64);
1192 estimate_pending_compaction_bytes
1193 .set(property_u64(&db, ESTIMATE_PENDING_COMPACTION_BYTES) as i64);
1194 actual_delayed_write_rate.set(property_u64(&db, ACTUAL_DELAYED_WRITE_RATE) as i64);
1195
1196 let Some(health) = health.upgrade() else {
1197 return;
1198 };
1199 health.sample_background_errors(bg);
1202 load_shed_active.set(if health.is_active() { 1 } else { 0 });
1203 }
1204 None => {
1205 return;
1207 }
1208 }
1209 tokio::time::sleep(tokio::time::Duration::from_secs(5)).await;
1210 }
1211}