diff --git a/rustfs/src/admin/handlers/config_admin.rs b/rustfs/src/admin/handlers/config_admin.rs index a43ccc0d0..c9b6bad04 100644 --- a/rustfs/src/admin/handlers/config_admin.rs +++ b/rustfs/src/admin/handlers/config_admin.rs @@ -96,6 +96,7 @@ const JSON_CONTENT_TYPE: &str = "application/json"; const TEXT_CONTENT_TYPE: &str = "text/plain; charset=utf-8"; const CONFIG_HISTORY_PREFIX: &str = "config/history"; const CONFIG_HISTORY_SUFFIX: &str = ".kv"; +const CONFIG_HISTORY_SNAPSHOT_HEADER: &str = "# rustfs-config-history: full-snapshot-v1"; const CONFIG_APPLIED_HEADER: &str = "x-rustfs-config-applied"; const CONFIG_APPLIED_COMPAT_HEADER: &str = "x-minio-config-applied"; const CONFIG_APPLIED_TRUE: &str = "true"; @@ -833,6 +834,64 @@ async fn save_server_config_history(data: &[u8]) -> S3Result { Ok(restore_id) } +/// Encode the complete persisted state that existed immediately before a +/// configuration mutation. Set, delete, full import, and restore all write +/// this record while holding the server-config write lock, before persisting +/// their replacement state. The version marker deliberately makes older +/// directive-only entries distinguishable so restore can reject them safely. +fn encode_config_history_snapshot(config: &ServerConfig) -> Vec { + let rendered = render_full_config(config); + let mut snapshot = Vec::with_capacity(CONFIG_HISTORY_SNAPSHOT_HEADER.len() + 1 + rendered.len()); + snapshot.extend_from_slice(CONFIG_HISTORY_SNAPSHOT_HEADER.as_bytes()); + snapshot.push(b'\n'); + snapshot.extend_from_slice(&rendered); + snapshot +} + +fn decode_config_history_snapshot(data: &[u8]) -> S3Result { + let text = std::str::from_utf8(data).map_err(|_| s3_error!(InvalidRequest, "config history entry is not valid UTF-8"))?; + let Some(snapshot) = text.strip_prefix(CONFIG_HISTORY_SNAPSHOT_HEADER) else { + return Err(s3_error!( + InvalidRequest, + "legacy config history entry cannot be restored safely; create a new history entry before retrying" + )); + }; + let snapshot = snapshot + .strip_prefix('\n') + .ok_or_else(|| s3_error!(InvalidRequest, "config history snapshot header is malformed"))?; + let directives = parse_config_directives(snapshot, false)?; + validate_config_directives(&directives)?; + + let mut config = ServerConfig::new(); + apply_set_directives(&mut config, &directives)?; + Ok(config) +} + +fn validate_restore_rollback_generation(current: &ServerConfig, restored: &ServerConfig) -> S3Result<()> { + if current == restored { + Ok(()) + } else { + Err(s3_error!( + InvalidRequest, + "config restore rollback was skipped because a concurrent configuration change committed" + )) + } +} + +async fn save_server_config_history_snapshot(config: &ServerConfig) -> S3Result { + save_server_config_history(&encode_config_history_snapshot(config)).await +} + +async fn cleanup_failed_config_history_snapshot(restore_id: &str) { + if let Err(err) = delete_server_config_history(restore_id).await { + warn!( + restore_id, + error = %err, + "Failed to remove config history snapshot for an uncommitted mutation" + ); + } +} + async fn read_server_config_history(restore_id: &str) -> S3Result> { let store = object_store()?; read_admin_config(store, &history_object_name(restore_id)) @@ -1761,20 +1820,25 @@ impl Operation for SetConfigKVHandler { let (config, storage_class_applied, notify_transition) = with_admin_server_config_write_lock(config_store, move || async move { let sub_system = transaction_sub_system.as_deref(); - let mut config = load_server_config_from_store_locked().await?; + let previous_config = load_server_config_from_store_locked().await?; + let mut config = previous_config.clone(); apply_set_directives(&mut config, &directives)?; let prepared = prepare_server_config(&config, sub_system).await?; - save_server_config_history(&body).await?; + let history_restore_id = save_server_config_history_snapshot(&previous_config).await?; + if let Err(err) = save_server_config_to_store_locked(&config).await { + cleanup_failed_config_history_snapshot(&history_restore_id).await; + return Err(err); + } if sub_system == Some(STORAGE_CLASS_SUB_SYS) { - commit_prepared_config( - config.clone(), - prepared, - save_server_config_to_store_locked(&config), - publish_prepared_config_snapshots, - ) - .await?; + publish_prepared_config_snapshots(config.clone(), prepared).map_err(|err| { + s3_error!( + InternalError, + "config persisted but runtime publish failed; recovery snapshot restoreId={}: {}", + history_restore_id, + err + ) + })?; } else { - save_server_config_to_store_locked(&config).await?; publish_server_config(config.clone()); } let notify_transition = publish_notify_config_intent(&config, sub_system); @@ -1812,20 +1876,25 @@ impl Operation for DelConfigKVHandler { let (config, storage_class_applied, notify_transition) = with_admin_server_config_write_lock(config_store, move || async move { let sub_system = transaction_sub_system.as_deref(); - let mut config = load_server_config_from_store_locked().await?; + let previous_config = load_server_config_from_store_locked().await?; + let mut config = previous_config.clone(); apply_delete_directives(&mut config, &directives); let prepared = prepare_server_config(&config, sub_system).await?; - save_server_config_history(&body).await?; + let history_restore_id = save_server_config_history_snapshot(&previous_config).await?; + if let Err(err) = save_server_config_to_store_locked(&config).await { + cleanup_failed_config_history_snapshot(&history_restore_id).await; + return Err(err); + } if sub_system == Some(STORAGE_CLASS_SUB_SYS) { - commit_prepared_config( - config.clone(), - prepared, - save_server_config_to_store_locked(&config), - publish_prepared_config_snapshots, - ) - .await?; + publish_prepared_config_snapshots(config.clone(), prepared).map_err(|err| { + s3_error!( + InternalError, + "config persisted but runtime publish failed; recovery snapshot restoreId={}: {}", + history_restore_id, + err + ) + })?; } else { - save_server_config_to_store_locked(&config).await?; publish_server_config(config.clone()); } let notify_transition = publish_notify_config_intent(&config, sub_system); @@ -1913,33 +1982,81 @@ impl Operation for RestoreConfigHistoryKVHandler { .get("restoreId") .ok_or_else(|| s3_error!(InvalidRequest, "missing restoreId query parameter"))?; let history = read_server_config_history(restore_id).await?; - let directives = parse_config_directives(std::str::from_utf8(&history).map_err(ApiError::other)?, false)?; - if directives.is_empty() { - return Err(s3_error!(InvalidRequest, "history entry is empty")); - } - validate_config_directives(&directives)?; - - let mut config = ServerConfig::new(); - apply_set_directives(&mut config, &directives)?; + let config = decode_config_history_snapshot(&history)?; let prepared = prepare_server_config(&config, None).await?; let config_store = object_store()?; supervise_admin_mutation("config mutation", async move { preflight_notify_config_intent(None).await?; let persisted_config = config.clone(); - let notify_transition = with_admin_server_config_write_lock(config_store, move || async move { + let restored_config = config.clone(); + let rollback_store = config_store.clone(); + let (previous_config, recovery_restore_id, notify_transition) = + with_admin_server_config_write_lock(config_store, move || async move { + let previous_config = load_server_config_from_store_locked().await?; + let recovery_restore_id = save_server_config_history_snapshot(&previous_config).await?; + if let Err(err) = save_server_config_to_store_locked(&persisted_config).await { + cleanup_failed_config_history_snapshot(&recovery_restore_id).await; + return Err(err); + } + publish_prepared_config_snapshots(persisted_config.clone(), prepared).map_err(|err| { + s3_error!( + InternalError, + "config restore persisted but runtime publish failed; recovery snapshot restoreId={}: {}", + recovery_restore_id, + err + ) + })?; + let notify_transition = publish_notify_config_intent(&persisted_config, None); + Ok::<_, S3Error>((previous_config, recovery_restore_id, notify_transition)) + }) + .await + .map_err(|err| s3_error!(InternalError, "failed to lock server config restore: {}", err))??; + + let Err(restore_error) = reconcile_full_config(config, notify_transition).await else { + return Ok(()); + }; + + let rollback_config = previous_config.clone(); + let rollback_expected_config = restored_config.clone(); + let rollback_transition = with_admin_server_config_write_lock(rollback_store, move || async move { + let current_config = load_server_config_from_store_locked().await?; + validate_restore_rollback_generation(¤t_config, &rollback_expected_config)?; + let rollback_prepared = prepare_server_config(&rollback_config, None).await?; commit_prepared_config( - persisted_config.clone(), - prepared, - save_server_config_to_store_locked(&persisted_config), + rollback_config.clone(), + rollback_prepared, + save_server_config_to_store_locked(&rollback_config), publish_prepared_config_snapshots, ) .await?; - Ok::<_, S3Error>(publish_notify_config_intent(&persisted_config, None)) + Ok::<_, S3Error>(publish_notify_config_intent(&rollback_config, None)) }) .await - .map_err(|err| s3_error!(InternalError, "failed to lock server config restore: {}", err))??; + .map_err(|err| s3_error!(InternalError, "failed to lock server config rollback: {}", err))?; - reconcile_full_config(config, notify_transition).await + let rollback_transition = match rollback_transition { + Ok(transition) => transition, + Err(rollback_error) => { + return Err(s3_error!( + InternalError, + "config restore failed: {}; automatic rollback failed: {}; recovery snapshot restoreId={}", + restore_error, + rollback_error, + recovery_restore_id + )); + } + }; + if let Err(rollback_error) = reconcile_full_config(previous_config, rollback_transition).await { + return Err(s3_error!( + InternalError, + "config restore failed: {}; persisted rollback did not converge: {}; recovery snapshot restoreId={}", + restore_error, + rollback_error, + recovery_restore_id + )); + } + cleanup_failed_config_history_snapshot(&recovery_restore_id).await; + Err(s3_error!(InternalError, "config restore failed and was rolled back: {}", restore_error)) }) .await?; @@ -1979,16 +2096,22 @@ impl Operation for SetConfigHandler { let config_store = object_store()?; supervise_admin_mutation("config mutation", async move { preflight_notify_config_intent(None).await?; - save_server_config_history(&body).await?; let persisted_config = config.clone(); let notify_transition = with_admin_server_config_write_lock(config_store, move || async move { - commit_prepared_config( - persisted_config.clone(), - prepared, - save_server_config_to_store_locked(&persisted_config), - publish_prepared_config_snapshots, - ) - .await?; + let previous_config = load_server_config_from_store_locked().await?; + let history_restore_id = save_server_config_history_snapshot(&previous_config).await?; + if let Err(err) = save_server_config_to_store_locked(&persisted_config).await { + cleanup_failed_config_history_snapshot(&history_restore_id).await; + return Err(err); + } + publish_prepared_config_snapshots(persisted_config.clone(), prepared).map_err(|err| { + s3_error!( + InternalError, + "full config persisted but runtime publish failed; recovery snapshot restoreId={}: {}", + history_restore_id, + err + ) + })?; Ok::<_, S3Error>(publish_notify_config_intent(&persisted_config, None)) }) .await @@ -2484,6 +2607,110 @@ notify_webhook:beta endpoint="http://beta.example""#, assert!(history_restore_id_from_name("config/history/restore-123.txt").is_none()); } + #[test] + fn history_snapshot_round_trip_preserves_unrelated_subsystems_targets_and_secrets() { + crate::admin::storage_api::config::init_admin_config_defaults(); + let mut original = ServerConfig::new(); + apply_set_directives( + &mut original, + &parse_config_directives( + r#"storage_class standard="EC:4" rrs="EC:2" +identity_openid client_id="client" client_secret="unrelated-secret" +notify_webhook:primary endpoint="https://hooks.example" auth_token="target-secret""#, + false, + ) + .expect("parse fixture"), + ) + .expect("apply fixture"); + + let encoded = encode_config_history_snapshot(&original); + let restored = decode_config_history_snapshot(&encoded).expect("decode snapshot"); + + assert_eq!(render_full_config(&restored), render_full_config(&original)); + let rendered = String::from_utf8(render_full_config(&restored)).expect("utf8"); + assert!(rendered.contains(r#"client_secret="unrelated-secret""#)); + assert!(rendered.contains(r#"notify_webhook:primary"#)); + assert!(rendered.contains(r#"auth_token="target-secret""#)); + } + + #[test] + fn history_snapshot_supports_complete_default_config() { + crate::admin::storage_api::config::init_admin_config_defaults(); + let original = ServerConfig::new(); + + let encoded = encode_config_history_snapshot(&original); + let restored = decode_config_history_snapshot(&encoded).expect("decode default snapshot"); + + assert_eq!(render_full_config(&restored), render_full_config(&original)); + } + + #[test] + fn pre_change_snapshot_restores_single_key_and_deleted_target() { + crate::admin::storage_api::config::init_admin_config_defaults(); + let mut before = ServerConfig::new(); + apply_set_directives( + &mut before, + &parse_config_directives( + r#"identity_openid client_id="before-client" client_secret="preserved-secret" +notify_webhook:secondary endpoint="https://secondary.example" auth_token="secondary-secret""#, + false, + ) + .expect("parse before fixture"), + ) + .expect("apply before fixture"); + let snapshot = encode_config_history_snapshot(&before); + + let mut after = before.clone(); + apply_set_directives( + &mut after, + &parse_config_directives(r#"identity_openid client_id="after-client""#, false).expect("parse set"), + ) + .expect("apply set"); + apply_delete_directives( + &mut after, + &parse_config_directives("notify_webhook:secondary", true).expect("parse target delete"), + ); + + let restored = decode_config_history_snapshot(&snapshot).expect("restore pre-change snapshot"); + assert_eq!(render_full_config(&restored), render_full_config(&before)); + assert_ne!(render_full_config(&restored), render_full_config(&after)); + } + + #[test] + fn restore_rollback_generation_rejects_concurrent_config_change() { + crate::admin::storage_api::config::init_admin_config_defaults(); + let restored = ServerConfig::new(); + let mut concurrent = restored.clone(); + apply_set_directives( + &mut concurrent, + &parse_config_directives(r#"identity_openid client_id="concurrent-client""#, false).expect("parse concurrent"), + ) + .expect("apply concurrent"); + + validate_restore_rollback_generation(&restored, &restored).expect("unchanged restore generation"); + let error = validate_restore_rollback_generation(&concurrent, &restored) + .expect_err("concurrent change must fence automatic rollback"); + assert_eq!(error.code(), &S3ErrorCode::InvalidRequest); + assert!(error.to_string().contains("concurrent configuration change")); + } + + #[test] + fn legacy_directive_history_is_rejected_without_being_applied_as_a_snapshot() { + let error = decode_config_history_snapshot(b"identity_openid client_id=\"legacy\"") + .expect_err("legacy directive history must fail closed"); + + assert_eq!(error.code(), &S3ErrorCode::InvalidRequest); + assert!(error.to_string().contains("legacy config history entry")); + } + + #[test] + fn malformed_versioned_history_is_rejected() { + let error = decode_config_history_snapshot(b"# rustfs-config-history: full-snapshot-v1\nnot_real key=\"value\"") + .expect_err("malformed snapshot must fail closed"); + + assert_eq!(error.code(), &S3ErrorCode::InvalidRequest); + } + #[test] fn trim_history_entries_keeps_most_recent_count_in_time_order() { let entries = vec![