Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 7 additions & 3 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -176,19 +183,16 @@ 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
treated as uncompressed and failed with a JSON parse error. Uncompressed
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
Expand Down
4 changes: 2 additions & 2 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ members = [
]

[workspace.package]
version = "0.22.0"
version = "0.23.0"
edition = "2021"
license = "MIT"
authors = ["OSO"]
Expand Down
31 changes: 31 additions & 0 deletions crates/kafka-backup-core/src/circuit_breaker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@ impl Default for CircuitBreakerConfig {
/// Circuit breaker for managing failure states
pub struct CircuitBreaker {
config: CircuitBreakerConfig,
enabled: bool,
state: Mutex<CircuitBreakerState>,
}

Expand All @@ -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,
Expand All @@ -70,15 +72,28 @@ 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
}

/// 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);

Expand All @@ -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 {
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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 {
Expand Down
115 changes: 115 additions & 0 deletions crates/kafka-backup-core/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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<std::path::PathBuf>,
Expand Down Expand Up @@ -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(),
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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
Expand Down
19 changes: 13 additions & 6 deletions crates/kafka-backup-core/src/restore/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
23 changes: 23 additions & 0 deletions docs/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down