spool/
rocks.rs

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    /// If non-zero, we perform bigger reads when doing compaction. If you’re running RocksDB on
42    /// spinning disks, you should set this to at least 2MB. That way RocksDB’s compaction is doing
43    /// sequential instead of random reads
44    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    /// Size in bytes of the rocksdb memtable that buffers writes before
53    /// being flushed to disk as a new SST file.
54    ///
55    /// Smaller values produce smaller, more frequent SST files and
56    /// trigger compactions sooner -- useful in test setups that need
57    /// to force the storage through its full write/compact lifecycle
58    /// quickly.  Larger values amortize compaction overhead but
59    /// increase memory use and recovery time after restart.  Leave
60    /// unset to use the rocksdb default.
61    #[serde(default)]
62    pub write_buffer_size: Option<usize>,
63
64    /// Number of level-0 SST files at which rocksdb will stop
65    /// accepting writes.  Lower values transition the database into
66    /// the write-stopped state more quickly when background
67    /// compaction cannot keep up, which is useful for tests that
68    /// need to deterministically observe that condition.  Leave
69    /// unset to use the rocksdb default.
70    #[serde(default)]
71    pub level0_stop_writes_trigger: Option<i32>,
72
73    #[serde(default)]
74    pub log_level: LogLevelDef,
75
76    /// See:
77    /// <https://docs.rs/rocksdb/latest/rocksdb/struct.Options.html#method.set_memtable_huge_page_size>
78    #[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    /// Upper bound on the wait that `store()` and `remove()` will
101    /// tolerate when rocksdb is applying backpressure.  Callers may
102    /// provide a shorter deadline (typically derived from an SMTP
103    /// client's idle timeout); the effective deadline is the minimum
104    /// of the two.  Going longer than the caller-provided value risks
105    /// the client timing out and retrying, which would produce
106    /// duplicate deliveries -- this option therefore only narrows the
107    /// effective deadline, it never extends it.
108    #[serde(
109        with = "duration_serde",
110        default = "RocksSpoolParams::default_store_deadline"
111    )]
112    pub store_deadline: Duration,
113
114    /// Delay after an error before pausing writes, specified as a duration
115    /// string such as `"15s"` or `"2m"`.
116    /// Even an isolated error starts this delay. A read, write, or enumeration
117    /// returning Corruption or IOError bypasses the delay and immediately
118    /// pauses writes until an automatic retry or process restart. Errors in
119    /// background flushes and compactions start the delay when the monitor
120    /// observes their count increasing.
121    #[serde(
122        with = "duration_serde",
123        default = "RocksSpoolParams::default_error_latch_duration"
124    )]
125    pub error_latch_duration: Duration,
126
127    /// The duration to pause before an automatic retry, specified as a duration
128    /// string such as `"5m"` or `"30s"`. The timer starts when the gate latches
129    /// and restarts on each error observation. Must be nonzero when
130    /// `allow_error_unlatch` is true. Longer pauses allow more time for
131    /// operator inspection.
132    #[serde(
133        with = "duration_serde",
134        default = "RocksSpoolParams::default_error_unlatch_duration"
135    )]
136    pub error_unlatch_duration: Duration,
137
138    /// Whether writes resume automatically after `error_unlatch_duration`
139    /// elapses with the gate latched and error counts unchanged. Defaults to
140    /// true. Set to false to keep writes paused until an operator restarts
141    /// the process after inspecting the database.
142    #[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    /// Error observations and the gate shared by foreground operations
267    /// and the monitor.
268    health: Arc<Health>,
269    store_deadline: Duration,
270}
271
272/// Initial sleep interval for the `store`/`remove` backpressure loop.
273/// Chosen low enough that brief, sub-millisecond memtable backpressure
274/// is caught with negligible added latency on the slow path.
275const BACKOFF_INITIAL: Duration = Duration::from_micros(500);
276/// Upper bound on the backpressure loop sleep interval.  Keeps the
277/// load-shedding gate observable within a bounded window even during
278/// a long wedge.
279const BACKOFF_MAX: Duration = Duration::from_millis(50);
280
281/// Selects the binding deadline for a backpressured write: the caller's
282/// deadline when it falls before the `store_deadline` of the spool, otherwise
283/// the spool deadline. Returns that deadline along with whether it was the
284/// caller's.
285fn 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    /// Builds the typed error for a backpressure deadline: the
297    /// caller-provided deadline when `caller_wins` is true, otherwise the
298    /// spool's own `store_deadline`.
299    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    /// Routes a backpressure timeout to the delayed latch path rather than
311    /// latching now: a load spike can exhaust the deadline on an intact
312    /// database, and the delay waits out that transient before latching. The
313    /// log fires only on the incident-starting transition (bounded to once per
314    /// incident, since a latched gate rejects writes before they reach this
315    /// path), not per timed-out write.
316    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    /// Acquires a concurrency permit under the effective deadline, or reports
333    /// the timeout as a backpressure incident and returns the typed error.
334    /// Returns `None` when `permits` is `None`.
335    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    /// Writes a RocksDB batch, retrying with exponential backoff until the
354    /// write succeeds, the effective deadline is reached, the load-shedding
355    /// gate latches, or RocksDB returns a non-`Incomplete` error. The write is
356    /// atomic, cancellable, and doesn't hold a blocking-pool worker while it
357    /// waits.
358    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        // Gate at the top so that the load-shedding mirror affects
366        // every store, not just those that happen to hit backpressure.
367        // A relaxed atomic load is essentially free compared to the
368        // rocksdb FFI write below; this preserves the healthy hot
369        // path's latency profile while giving the gate consistent
370        // semantics across the in-flight call sites that aren't
371        // covered by the per-connection ingress checks (notably,
372        // already-established SMTP connections doing new
373        // transactions).
374        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        // The default is 1000, which is a bit high
432        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        // Ensure the directory exists so we can probe it before opening.
483        // Create it the way RocksDB would (mkdir 0755, subject to umask)
484        // so our pre-creation is indistinguishable from letting RocksDB
485        // create it; RocksDB has no option that influences this mode.
486        #[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        // Catch a split real/effective identity meeting an over-restrictive
500        // directory before RocksDB silently corrupts itself and later fails
501        // with an opaque "wal_dir contains existing log file" error.
502        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        // Baseline at zero: every error on this DB instance, including one
507        // from startup work performed by DB::open itself, can start an
508        // incident.
509        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                            // Count read failures explicitly, since
552                            // background-error sampling cannot detect a failed
553                            // read.
554                            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            // Force bottommost-level compaction so the entire keyspace
606            // is rewritten; without this, single-level layouts cause
607            // the call to be a no-op even when there are missing files
608            // that we'd want to surface as errors.
609            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            // compact_range itself does not return errors -- wait_for_compact
614            // does, and is what surfaces background failures (e.g. a
615            // missing SST encountered during compaction) to the caller.
616            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                                // Iterator errors typically indicate a
691                                // missing or corrupt SST file
692                                // discovered while walking the
693                                // keyspace.  Feed into the foreground
694                                // error machinery so the gate latches
695                                // (immediately for IOError /
696                                // Corruption) and abort the
697                                // enumeration.
698                                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                            // Entries created since we started must have
707                            // landed there after we started and are thus
708                            // not eligible for discovery via enumeration
709                            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    /// The caller's deadline binds only when it is sooner than the spool's own
735    /// `store_deadline`. The returned flag indicates whether the caller's
736    /// deadline was selected.
737    #[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            // Can't load an entry that doesn't exist
766            assert_eq!(
767                format!("{:#}", spool.load(id1).await.unwrap_err()),
768                format!("no such key {id1}")
769            );
770        }
771
772        // Insert some entries
773        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        // Verify that we can load those entries
788        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            // Verify that we can enumerate them
796            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                        // Can't load an entry that we just removed
813                        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        // Now that we've removed the files, try enumerating again.
829        // We expect to receive no entries.
830        // Do it a couple of times to verify that none of the cleanup
831        // stuff that happens in enumerate breaks the directory
832        // structure
833        for _ in 0..2 {
834            // Verify that we can enumerate them
835            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/// The rocksdb type doesn't impl Debug, so we get to do it
855#[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
881/// Read an integer-valued rocksdb property, returning 0 if the property
882/// is missing or cannot be parsed.  Used for hot-path checks and metrics
883/// gathering; callers that want to distinguish "missing" from "zero"
884/// should call `property_int_value` directly.
885fn property_u64(db: &DB, name: &rocksdb::properties::PropName) -> u64 {
886    db.property_int_value(name).ok().flatten().unwrap_or(0)
887}
888
889/// Returns true for errors that require an immediate write pause to protect
890/// stored data from further damage.
891fn is_definitively_bad(err: &rocksdb::Error) -> bool {
892    matches!(err.kind(), ErrorKind::Corruption | ErrorKind::IOError)
893}
894
895/// Records a RocksDB error and pauses writes immediately for Corruption
896/// and IOError. Logs the first error after startup or a retry, and any
897/// error that changes the spool from accepting writes to refusing them.
898fn 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! {
930/// Approximate memory usage (bytes) of all the mem-tables.
931///
932/// This may be useful when understanding the memory usage of
933/// the system.
934static MEM_TABLE_TOTAL: IntGaugeVec(
935        "rocks_spool_mem_table_total",
936        &["path"]
937    );
938}
939
940declare_metric! {
941/// Approximate memory usage (bytes) of un-flushed mem-tables.
942///
943/// This may be useful when understanding the memory usage of
944/// the system.
945static MEM_TABLE_UNFLUSHED: IntGaugeVec(
946        "rocks_spool_mem_table_unflushed",
947        &["path"]
948    );
949}
950
951declare_metric! {
952/// Approximate memory usage (bytes) of all the table readers.
953///
954/// This may be useful when understanding the memory usage of
955/// the system.
956static MEM_TABLE_READERS_TOTAL: IntGaugeVec(
957        "rocks_spool_mem_table_readers_total",
958        &["path"]
959    );
960}
961
962declare_metric! {
963/// Approximate memory (bytes) usage by cache.
964///
965/// This may be useful when understanding the memory usage of
966/// the system.
967static CACHE_TOTAL: IntGaugeVec(
968        "rocks_spool_cache_total",
969        &["path"]
970    );
971}
972
973declare_metric! {
974/// Accumulated count of background errors encountered by the rocksdb
975/// instance (failed flushes or compactions, typically caused by I/O
976/// errors such as missing or corrupt SST files, ENOSPC, or permission
977/// problems).
978///
979/// {{since('2026.09.22-a276d4a8')}}
980///
981/// This counter is **monotonic** for the lifetime of the process: it
982/// does not decrease when rocksdb auto-resumes from transient errors
983/// such as a brief ENOSPC.  A non-zero value therefore does not
984/// necessarily mean the database is currently wedged; it means at
985/// least one background error has occurred since the process started.
986///
987/// For SRE monitoring, alert on the **rate of change** (e.g.
988/// `increase(rocks_spool_background_errors[5m]) > 0`) to catch new
989/// occurrences.  For the actionable "the database is wedged right
990/// now and we are shedding load" signal, page on
991/// `rocks_spool_load_shed_active` instead, which combines this
992/// counter, foreground read/write errors, and rocksdb error
993/// severity into a single latched indicator.
994static BACKGROUND_ERRORS_METRIC: IntGaugeVec(
995        "rocks_spool_background_errors",
996        &["path"]
997    );
998}
999
1000declare_metric! {
1001/// Set to 1 when the rocksdb instance is currently refusing writes
1002/// at the WriteController layer (memtable count or L0 file count
1003/// reached the stop threshold), 0 otherwise.
1004///
1005/// {{since('2026.09.22-a276d4a8')}}
1006///
1007/// This reflects rocksdb's own `is-write-stopped` property and
1008/// indicates backpressure rather than a fatal background error.
1009/// Healthy databases under bursty load may briefly report 1 here.
1010/// For the "the database is wedged due to a background error"
1011/// signal, see `rocks_spool_load_shed_active` instead.
1012static WRITE_STOPPED: IntGaugeVec(
1013        "rocks_spool_write_stopped",
1014        &["path"]
1015    );
1016}
1017
1018declare_metric! {
1019/// Set to 1 while this spool refuses writes, or 0 otherwise. When set,
1020/// SMTP and HTTP ingress reject traffic, and store/remove operations
1021/// return an error immediately.
1022///
1023/// {{since('2026.09.22-a276d4a8')}}
1024///
1025/// A foreground operation returning `Corruption` or `IOError` immediately
1026/// latches the gate, causing subsequent writes to return errors. These failures
1027/// include missing and corrupt SST files.
1028///
1029/// Newly observed background errors, other foreground errors, and timeouts
1030/// while waiting for RocksDB to accept a write start the `error_latch_duration`
1031/// delay (default 15 seconds). Even an isolated error causes a latch after
1032/// this delay.
1033///
1034/// With `allow_error_unlatch = true` (the default), writes resume after
1035/// `error_unlatch_duration` (default 5 minutes) has elapsed since the later
1036/// of the latch time and the most recent error observation. If the database
1037/// remains damaged, another error can latch the gate again. Set
1038/// `allow_error_unlatch = false` to keep writes paused until an operator
1039/// inspects the database and restarts the process.
1040///
1041/// The monitor checks background-error growth and applies the latch and
1042/// retry timers, then sleeps for 5 seconds. Later errors are handled by
1043/// the next iteration.
1044///
1045/// Automatic retries accept this sampling delay: writes may resume between a
1046/// background error and its observation. Disable `allow_error_unlatch` to
1047/// keep writes paused across that window.
1048///
1049/// If writes remain paused, inspect `rocks_spool_background_errors` and the
1050/// RocksDB LOG to identify the storage failure.
1051static LOAD_SHED_ACTIVE: IntGaugeVec(
1052        "rocks_spool_load_shed_active",
1053        &["path"]
1054    );
1055}
1056
1057declare_metric! {
1058/// Number of background compactions currently running for this
1059/// rocksdb instance.
1060///
1061/// {{since('2026.09.22-a276d4a8')}}
1062///
1063/// In a healthy, actively-written spool this is typically non-zero
1064/// in bursts.  A value persistently stuck at 0 while
1065/// `rocks_spool_compaction_pending` or
1066/// `rocks_spool_estimate_pending_compaction_bytes` is growing is a
1067/// strong indicator that the background worker is wedged --
1068/// cross-reference `rocks_spool_write_stopped` and
1069/// `rocks_spool_background_errors`.
1070static NUM_RUNNING_COMPACTIONS_METRIC: IntGaugeVec(
1071        "rocks_spool_num_running_compactions",
1072        &["path"]
1073    );
1074}
1075
1076declare_metric! {
1077/// Set to 1 when at least one compaction is pending for this rocksdb
1078/// instance, 0 otherwise.
1079///
1080/// {{since('2026.09.22-a276d4a8')}}
1081///
1082/// Brief flapping is normal under write load.  A value of 1 that
1083/// persists alongside `rocks_spool_num_running_compactions == 0` is
1084/// suspicious and suggests the compaction worker is not making
1085/// progress.
1086static COMPACTION_PENDING_METRIC: IntGaugeVec(
1087        "rocks_spool_compaction_pending",
1088        &["path"]
1089    );
1090}
1091
1092declare_metric! {
1093/// Estimated total bytes that compaction needs to rewrite to bring
1094/// all levels back under their target sizes.
1095///
1096/// {{since('2026.09.22-a276d4a8')}}
1097///
1098/// This is a backlog indicator.  Steady-state values depend heavily
1099/// on write rate, compression, and the configured compaction style,
1100/// so absolute thresholds should be derived from each deployment's
1101/// baseline.  Unbounded growth over a multi-hour window indicates
1102/// that compaction cannot keep up with the write rate, which
1103/// eventually leads to write slowdown
1104/// (`rocks_spool_actual_delayed_write_rate` becomes non-zero) and
1105/// then to write stop (`rocks_spool_write_stopped` becomes 1).
1106///
1107/// Only meaningful for level-style compaction.
1108static ESTIMATE_PENDING_COMPACTION_BYTES_METRIC: IntGaugeVec(
1109        "rocks_spool_estimate_pending_compaction_bytes",
1110        &["path"]
1111    );
1112}
1113
1114declare_metric! {
1115/// Current delayed write rate (bytes/second) applied by rocksdb to
1116/// throttle foreground writers.  0 means no slowdown is in effect.
1117///
1118/// {{since('2026.09.22-a276d4a8')}}
1119///
1120/// A non-zero value means rocksdb is intentionally slowing writers
1121/// down because compaction or flush is falling behind.  This is the
1122/// early-warning signal that precedes a full write stop: if this
1123/// remains non-zero for an extended period, investigate the
1124/// compaction backlog
1125/// (`rocks_spool_estimate_pending_compaction_bytes`) and underlying
1126/// disk throughput before the database transitions to
1127/// `rocks_spool_write_stopped == 1`.
1128static 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                // Apply each background sample independently. New errors
1200                // observed after reopening start another latch delay.
1201                health.sample_background_errors(bg);
1202                load_shed_active.set(if health.is_active() { 1 } else { 0 });
1203            }
1204            None => {
1205                // Dead
1206                return;
1207            }
1208        }
1209        tokio::time::sleep(tokio::time::Duration::from_secs(5)).await;
1210    }
1211}