diff --git a/CHANGELOG.md b/CHANGELOG.md index 25f2671..8d69499 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,13 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [0.23.0] - 2026-09-24 + +### Added +- Configurable `backup.circuit_breaker` and `restore.circuit_breaker` settings, + including an `enabled` switch, failure threshold, reset timeout, and success + threshold. Disabled breakers remain closed and do not block operations. + ## [0.22.0] - 2026-09-07 ### Added @@ -176,7 +183,6 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [0.19.2] - 2026-08-30 -### Fixed - Restore can read **compressed legacy JSON segments** again. `read_segment` passed only the bare extension (`zst`) to `detect_from_extension`, which matches on `.zst` / `.lz4`, so every compressed pre-binary-format segment was @@ -184,11 +190,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 legacy segments were unaffected. Covered by a unit test per codec. ## [0.19.1] - 2026-08-29 - ### Fixed - `kafka_backup_snapshot_records_target` and `kafka_backup_snapshot_records_remaining` now describe the work of the - current run. Snapshot mode (`stop_at_current_offsets`) sized both gauges from the whole captured offset range (`latest - earliest` summed over partitions) and only subtracted a partition's checkpointed prefix once that partition's task started, so an incremental run over a large archive began diff --git a/Cargo.lock b/Cargo.lock index 43926e0..5067c33 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1802,7 +1802,7 @@ dependencies = [ [[package]] name = "kafka-backup-cli" -version = "0.22.0" +version = "0.23.0" dependencies = [ "anyhow", "bytes", @@ -1823,7 +1823,7 @@ dependencies = [ [[package]] name = "kafka-backup-core" -version = "0.22.0" +version = "0.23.0" dependencies = [ "anyhow", "async-trait", diff --git a/Cargo.toml b/Cargo.toml index 4664b10..6d22a77 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -6,7 +6,7 @@ members = [ ] [workspace.package] -version = "0.22.0" +version = "0.23.0" edition = "2021" license = "MIT" authors = ["OSO"] diff --git a/crates/kafka-backup-core/src/circuit_breaker.rs b/crates/kafka-backup-core/src/circuit_breaker.rs index 4cea7db..08b5cc8 100644 --- a/crates/kafka-backup-core/src/circuit_breaker.rs +++ b/crates/kafka-backup-core/src/circuit_breaker.rs @@ -45,6 +45,7 @@ impl Default for CircuitBreakerConfig { /// Circuit breaker for managing failure states pub struct CircuitBreaker { config: CircuitBreakerConfig, + enabled: bool, state: Mutex, } @@ -61,6 +62,7 @@ impl CircuitBreaker { info!("Created circuit breaker: {}", config.name); Self { config, + enabled: true, state: Mutex::new(CircuitBreakerState { state: CircuitState::Closed, failure_count: 0, @@ -70,8 +72,18 @@ impl CircuitBreaker { } } + /// Create a circuit breaker that never blocks operations. + pub fn disabled() -> Self { + let mut breaker = Self::new(CircuitBreakerConfig::default()); + breaker.enabled = false; + breaker + } + /// Get the current circuit state pub fn state(&self) -> CircuitState { + if !self.enabled { + return CircuitState::Closed; + } let mut state = self.state.lock(); self.maybe_transition_to_half_open(&mut state); state.state @@ -79,6 +91,9 @@ impl CircuitBreaker { /// Check if the circuit allows the operation pub fn is_allowed(&self) -> bool { + if !self.enabled { + return true; + } let mut state = self.state.lock(); self.maybe_transition_to_half_open(&mut state); @@ -91,6 +106,9 @@ impl CircuitBreaker { /// Record a successful operation pub fn record_success(&self) { + if !self.enabled { + return; + } let mut state = self.state.lock(); match state.state { @@ -125,6 +143,9 @@ impl CircuitBreaker { /// Record a failed operation pub fn record_failure(&self) { + if !self.enabled { + return; + } let mut state = self.state.lock(); state.failure_count += 1; @@ -242,6 +263,16 @@ mod tests { assert!(cb.is_allowed()); } + #[test] + fn disabled_circuit_breaker_allows_failures() { + let cb = CircuitBreaker::disabled(); + + cb.record_failure(); + + assert_eq!(cb.state(), CircuitState::Closed); + assert!(cb.is_allowed()); + } + #[test] fn test_circuit_breaker_opens_after_failures() { let cb = CircuitBreaker::new(CircuitBreakerConfig { diff --git a/crates/kafka-backup-core/src/config.rs b/crates/kafka-backup-core/src/config.rs index c838f37..4a90758 100644 --- a/crates/kafka-backup-core/src/config.rs +++ b/crates/kafka-backup-core/src/config.rs @@ -776,6 +776,58 @@ pub enum ExistingTopicConfigPolicy { Fail, } +/// Circuit breaker settings for restore Kafka operations. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct RestoreCircuitBreakerConfig { + /// Whether the restore circuit breaker should block operations. + #[serde(default = "default_restore_circuit_breaker_enabled")] + pub enabled: bool, + + /// Number of consecutive failures before opening the circuit. + #[serde(default = "default_restore_circuit_breaker_failure_threshold")] + pub failure_threshold: u32, + + /// Time to wait before probing a previously failing circuit. + #[serde(default = "default_restore_circuit_breaker_reset_timeout_ms")] + pub reset_timeout_ms: u64, + + /// Number of successful probes required to close the circuit. + #[serde(default = "default_restore_circuit_breaker_success_threshold")] + pub success_threshold: u32, +} + +impl Default for RestoreCircuitBreakerConfig { + fn default() -> Self { + Self { + enabled: default_restore_circuit_breaker_enabled(), + failure_threshold: default_restore_circuit_breaker_failure_threshold(), + reset_timeout_ms: default_restore_circuit_breaker_reset_timeout_ms(), + success_threshold: default_restore_circuit_breaker_success_threshold(), + } + } +} + +impl RestoreCircuitBreakerConfig { + fn validate(&self) -> crate::Result<()> { + if self.failure_threshold == 0 { + return Err(crate::Error::Config( + "restore.circuit_breaker.failure_threshold must be > 0".to_string(), + )); + } + if self.reset_timeout_ms == 0 { + return Err(crate::Error::Config( + "restore.circuit_breaker.reset_timeout_ms must be > 0".to_string(), + )); + } + if self.success_threshold == 0 { + return Err(crate::Error::Config( + "restore.circuit_breaker.success_threshold must be > 0".to_string(), + )); + } + Ok(()) + } +} + /// Restore-specific options #[derive(Debug, Clone, Serialize, Deserialize)] pub struct RestoreOptions { @@ -867,6 +919,10 @@ pub struct RestoreOptions { #[serde(default = "default_produce_timeout_ms")] pub produce_timeout_ms: i32, + /// Kafka circuit breaker settings for restore operations. + #[serde(default)] + pub circuit_breaker: RestoreCircuitBreakerConfig, + /// Checkpoint state file path for resumable restores #[serde(default)] pub checkpoint_state: Option, @@ -1047,6 +1103,7 @@ impl Default for RestoreOptions { produce_batch_size: default_produce_batch_size(), produce_acks: default_produce_acks(), produce_timeout_ms: default_produce_timeout_ms(), + circuit_breaker: RestoreCircuitBreakerConfig::default(), checkpoint_state: None, checkpoint_interval_secs: default_restore_checkpoint_interval_secs(), consumer_groups: Vec::new(), @@ -1087,6 +1144,22 @@ fn default_produce_timeout_ms() -> i32 { 30_000 // 30 seconds, matching the previous hardcoded value } +fn default_restore_circuit_breaker_enabled() -> bool { + true +} + +fn default_restore_circuit_breaker_failure_threshold() -> u32 { + 5 +} + +fn default_restore_circuit_breaker_reset_timeout_ms() -> u64 { + 30_000 +} + +fn default_restore_circuit_breaker_success_threshold() -> u32 { + 2 +} + /// Public accessor for integration tests that assert the default acks value. pub fn default_produce_acks_pub() -> i16 { default_produce_acks() @@ -1309,6 +1382,8 @@ impl RestoreOptions { )); } + self.circuit_breaker.validate()?; + // Validate consumer group offset reset if self.reset_consumer_offsets && self.consumer_groups.is_empty() @@ -1354,6 +1429,46 @@ mod tests { assert_eq!(defaults.segment_max_records, None); } + #[test] + fn restore_circuit_breaker_options_parse_and_default() { + let defaults: RestoreOptions = serde_yaml::from_str("{}").unwrap(); + assert_eq!( + defaults.circuit_breaker, + RestoreCircuitBreakerConfig::default() + ); + + let configured: RestoreOptions = serde_yaml::from_str( + r#" +circuit_breaker: + enabled: false + failure_threshold: 15 + reset_timeout_ms: 2000 + success_threshold: 1 +"#, + ) + .unwrap(); + + assert_eq!(configured.circuit_breaker.enabled, false); + assert_eq!(configured.circuit_breaker.failure_threshold, 15); + assert_eq!(configured.circuit_breaker.reset_timeout_ms, 2000); + assert_eq!(configured.circuit_breaker.success_threshold, 1); + } + + #[test] + fn restore_circuit_breaker_options_reject_zero_values() { + let mut options = RestoreOptions::default(); + options.circuit_breaker.failure_threshold = 0; + assert!(options.validate().is_err()); + + options.circuit_breaker = RestoreCircuitBreakerConfig::default(); + options.circuit_breaker.reset_timeout_ms = 0; + assert!(options.validate().is_err()); + + options.circuit_breaker = RestoreCircuitBreakerConfig::default(); + options.circuit_breaker.success_threshold = 0; + assert!(options.validate().is_err()); + } + #[test] fn pre_v0_16_config_parses_unchanged_with_no_warnings() { // A config written for older versions (no v0.16.0 keys) must parse diff --git a/crates/kafka-backup-core/src/restore/engine.rs b/crates/kafka-backup-core/src/restore/engine.rs index faabf46..9a82738 100644 --- a/crates/kafka-backup-core/src/restore/engine.rs +++ b/crates/kafka-backup-core/src/restore/engine.rs @@ -320,12 +320,19 @@ impl RestoreEngine { health.register_component("storage"); // Initialize circuit breakers - let kafka_circuit_breaker = Arc::new(CircuitBreaker::new(CircuitBreakerConfig { - failure_threshold: 5, - reset_timeout: Duration::from_secs(30), - success_threshold: 2, - name: "kafka".to_string(), - })); + let restore_options = config.restore.as_ref().cloned().unwrap_or_default(); + let kafka_circuit_breaker = if restore_options.circuit_breaker.enabled { + Arc::new(CircuitBreaker::new(CircuitBreakerConfig { + failure_threshold: restore_options.circuit_breaker.failure_threshold, + reset_timeout: Duration::from_millis( + restore_options.circuit_breaker.reset_timeout_ms, + ), + success_threshold: restore_options.circuit_breaker.success_threshold, + name: "kafka".to_string(), + })) + } else { + Arc::new(CircuitBreaker::disabled()) + }; let storage_circuit_breaker = Arc::new(CircuitBreaker::new(CircuitBreakerConfig { failure_threshold: 3, diff --git a/docs/configuration.md b/docs/configuration.md index 655e76b..3cd5206 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -835,6 +835,29 @@ restore: rate_limit_records_per_sec: 10000 ``` +#### Circuit Breaker + +The Kafka circuit breaker protects restore operations from cascading failures. +Its defaults preserve the existing behavior. Set `enabled: false` for a +checkpointed restore when isolated transient failures should not pause the +whole restore. + +| Option | Type | Required | Default | Description | +|--------|------|----------|---------|-------------| +| `circuit_breaker.enabled` | bool | No | `true` | Enable Kafka circuit-breaker blocking | +| `circuit_breaker.failure_threshold` | int | No | `5` | Failures before opening the circuit | +| `circuit_breaker.reset_timeout_ms` | int | No | `30000` | Delay before a half-open probe | +| `circuit_breaker.success_threshold` | int | No | `2` | Successful probes before closing the circuit | + +```yaml +restore: + circuit_breaker: + enabled: true + failure_threshold: 15 + reset_timeout_ms: 2000 + success_threshold: 1 +``` + ### Resumable Restores | Option | Type | Required | Default | Description |