Skip to main content

laminar_core/streaming/
checkpoint.rs

1//! Streaming checkpoint configuration.
2
3/// Configuration for streaming checkpoints.
4#[derive(Debug, Clone, Default)]
5pub struct StreamCheckpointConfig {
6    /// Checkpoint interval in milliseconds. `None` = manual only.
7    pub interval_ms: Option<u64>,
8    /// One end-to-end attempt deadline in milliseconds, spanning sink fencing, alignment,
9    /// capture, durable publication, and completion delivery. `None` = default (`120_000`).
10    pub timeout_ms: Option<u64>,
11    /// Directory for persisting checkpoints. `None` uses the database storage directory, then
12    /// falls back to `./data`; it never silently selects volatile checkpoint storage.
13    pub data_dir: Option<std::path::PathBuf>,
14    /// Number of predecessor checkpoints retained alongside the current recovery cut.
15    /// `None` = default (3); predecessors keep reference/delta chains resolvable.
16    pub max_retained: Option<usize>,
17    /// Maximum bytes admitted for one checkpoint across in-flight capture and
18    /// persisted/restored external state. `None` uses
19    /// `DEFAULT_MAX_CHECKPOINT_STATE_BYTES`.
20    pub max_staged_bytes: Option<u64>,
21}
22
23#[cfg(test)]
24mod tests {
25    use super::*;
26
27    #[test]
28    fn test_default_config() {
29        let config = StreamCheckpointConfig::default();
30        assert!(config.interval_ms.is_none());
31        assert!(config.timeout_ms.is_none());
32        assert!(config.data_dir.is_none());
33        assert!(config.max_retained.is_none());
34        assert!(config.max_staged_bytes.is_none());
35    }
36}