From efcd960b65bf37968f24296bc1a342df7e63cd4b Mon Sep 17 00:00:00 2001 From: houseme Date: Wed, 26 Aug 2026 21:31:17 +0800 Subject: [PATCH] feat(startup): expose resync reconcile observability (#6667) Co-authored-by: heihutu --- rustfs/src/startup_bucket_metadata.rs | 133 ++++++++++++++++++++++++++ rustfs/tests/connect_registration.rs | 16 +++- 2 files changed, 146 insertions(+), 3 deletions(-) diff --git a/rustfs/src/startup_bucket_metadata.rs b/rustfs/src/startup_bucket_metadata.rs index 33e90cb1b..dec3f50c1 100644 --- a/rustfs/src/startup_bucket_metadata.rs +++ b/rustfs/src/startup_bucket_metadata.rs @@ -20,12 +20,31 @@ use crate::storage_api::startup::bucket_metadata::{ use std::{ io::{Error as IoError, Result as IoResult}, sync::Arc, + time::{Duration, Instant}, }; use tokio_util::sync::CancellationToken; +const EVENT_REPLICATION_RESYNC_STARTUP_BACKGROUND_CANCELED: &str = "replication_resync_startup_background_canceled"; +const EVENT_REPLICATION_RESYNC_STARTUP_BACKGROUND_COMPLETED: &str = "replication_resync_startup_background_completed"; const EVENT_REPLICATION_RESYNC_STARTUP_BACKGROUND_FAILED: &str = "replication_resync_startup_background_failed"; +const EVENT_REPLICATION_RESYNC_STARTUP_BACKGROUND_STARTED: &str = "replication_resync_startup_background_started"; const LOG_COMPONENT_STARTUP_BUCKET_METADATA: &str = "startup_bucket_metadata"; const LOG_SUBSYSTEM_REPLICATION: &str = "replication"; +const METRIC_REPLICATION_RESYNC_STARTUP_BACKGROUND_DURATION_SECONDS: &str = + "rustfs_replication_resync_startup_background_duration_seconds"; +const METRIC_REPLICATION_RESYNC_STARTUP_BACKGROUND_EVENTS_TOTAL: &str = + "rustfs_replication_resync_startup_background_events_total"; +const METRIC_REPLICATION_RESYNC_STARTUP_BACKGROUND_STATUS: &str = "rustfs_replication_resync_startup_background_status"; +const STARTUP_BACKGROUND_MODE_EMBEDDED: &str = "embedded"; +const STARTUP_BACKGROUND_MODE_SERVER: &str = "server"; +const STARTUP_BACKGROUND_OUTCOME_CANCELED: &str = "canceled"; +const STARTUP_BACKGROUND_OUTCOME_FAILED: &str = "failed"; +const STARTUP_BACKGROUND_OUTCOME_STARTED: &str = "started"; +const STARTUP_BACKGROUND_OUTCOME_SUCCEEDED: &str = "succeeded"; +const STARTUP_BACKGROUND_STATUS_FAILED: f64 = 0.0; +const STARTUP_BACKGROUND_STATUS_SUCCEEDED: f64 = 1.0; +const STARTUP_BACKGROUND_STATUS_RUNNING: f64 = 2.0; +const STARTUP_BACKGROUND_STATUS_CANCELED: f64 = 3.0; pub(crate) async fn init_embedded_bucket_metadata_runtime(store: Arc, ctx: &CancellationToken) -> IoResult> { let buckets_list = store @@ -68,23 +87,127 @@ pub(crate) async fn init_bucket_metadata_runtime(store: Arc, ctx: Cance fn spawn_bucket_resync_startup_reconcile(buckets: Vec, ctx: CancellationToken, init_resync_after_reconcile: bool) { tokio::spawn(async move { + describe_bucket_resync_startup_background_metrics(); + let bucket_count = buckets.len(); + let mode = bucket_resync_startup_background_mode(init_resync_after_reconcile); + let started = Instant::now(); + + record_bucket_resync_startup_background_started(mode); + tracing::info!( + event = EVENT_REPLICATION_RESYNC_STARTUP_BACKGROUND_STARTED, + component = LOG_COMPONENT_STARTUP_BUCKET_METADATA, + subsystem = LOG_SUBSYSTEM_REPLICATION, + state = STARTUP_BACKGROUND_OUTCOME_STARTED, + mode, + bucket_count, + init_resync_after_reconcile, + "Bucket metadata startup resync reconcile started in background" + ); + if let Err(error) = run_bucket_resync_startup_reconcile(buckets, ctx, init_resync_after_reconcile).await { if !report_bucket_resync_startup_background_error(&error) { + record_bucket_resync_startup_background_finished(mode, STARTUP_BACKGROUND_OUTCOME_CANCELED, started.elapsed()); + tracing::debug!( + event = EVENT_REPLICATION_RESYNC_STARTUP_BACKGROUND_CANCELED, + component = LOG_COMPONENT_STARTUP_BUCKET_METADATA, + subsystem = LOG_SUBSYSTEM_REPLICATION, + result = STARTUP_BACKGROUND_OUTCOME_CANCELED, + mode, + bucket_count, + init_resync_after_reconcile, + duration_ms = started.elapsed().as_millis() as u64, + "Bucket metadata startup resync reconcile canceled during shutdown" + ); return; } + record_bucket_resync_startup_background_finished(mode, STARTUP_BACKGROUND_OUTCOME_FAILED, started.elapsed()); tracing::error!( event = EVENT_REPLICATION_RESYNC_STARTUP_BACKGROUND_FAILED, component = LOG_COMPONENT_STARTUP_BUCKET_METADATA, subsystem = LOG_SUBSYSTEM_REPLICATION, result = "failed", + mode, + bucket_count, init_resync_after_reconcile, + duration_ms = started.elapsed().as_millis() as u64, error = %error, "Bucket metadata startup resync reconcile failed in background" ); + return; } + + record_bucket_resync_startup_background_finished(mode, STARTUP_BACKGROUND_OUTCOME_SUCCEEDED, started.elapsed()); + tracing::info!( + event = EVENT_REPLICATION_RESYNC_STARTUP_BACKGROUND_COMPLETED, + component = LOG_COMPONENT_STARTUP_BUCKET_METADATA, + subsystem = LOG_SUBSYSTEM_REPLICATION, + result = "ok", + mode, + bucket_count, + init_resync_after_reconcile, + duration_ms = started.elapsed().as_millis() as u64, + "Bucket metadata startup resync reconcile completed in background" + ); }); } +fn describe_bucket_resync_startup_background_metrics() { + static DESCRIBE: std::sync::Once = std::sync::Once::new(); + DESCRIBE.call_once(|| { + metrics::describe_counter!( + METRIC_REPLICATION_RESYNC_STARTUP_BACKGROUND_EVENTS_TOTAL, + "Bucket metadata startup resync background task events, by fixed mode and outcome" + ); + metrics::describe_histogram!( + METRIC_REPLICATION_RESYNC_STARTUP_BACKGROUND_DURATION_SECONDS, + "Bucket metadata startup resync background task duration in seconds, by fixed mode and outcome" + ); + metrics::describe_gauge!( + METRIC_REPLICATION_RESYNC_STARTUP_BACKGROUND_STATUS, + "Latest bucket metadata startup resync background task status by fixed mode: 0=failed, 1=succeeded, 2=running, 3=canceled" + ); + }); +} + +fn bucket_resync_startup_background_mode(init_resync_after_reconcile: bool) -> &'static str { + if init_resync_after_reconcile { + STARTUP_BACKGROUND_MODE_SERVER + } else { + STARTUP_BACKGROUND_MODE_EMBEDDED + } +} + +fn record_bucket_resync_startup_background_started(mode: &'static str) { + metrics::counter!( + METRIC_REPLICATION_RESYNC_STARTUP_BACKGROUND_EVENTS_TOTAL, + "mode" => mode, + "outcome" => STARTUP_BACKGROUND_OUTCOME_STARTED + ) + .increment(1); + metrics::gauge!(METRIC_REPLICATION_RESYNC_STARTUP_BACKGROUND_STATUS, "mode" => mode).set(STARTUP_BACKGROUND_STATUS_RUNNING); +} + +fn record_bucket_resync_startup_background_finished(mode: &'static str, outcome: &'static str, duration: Duration) { + let status = match outcome { + STARTUP_BACKGROUND_OUTCOME_SUCCEEDED => STARTUP_BACKGROUND_STATUS_SUCCEEDED, + STARTUP_BACKGROUND_OUTCOME_CANCELED => STARTUP_BACKGROUND_STATUS_CANCELED, + _ => STARTUP_BACKGROUND_STATUS_FAILED, + }; + metrics::counter!( + METRIC_REPLICATION_RESYNC_STARTUP_BACKGROUND_EVENTS_TOTAL, + "mode" => mode, + "outcome" => outcome + ) + .increment(1); + metrics::histogram!( + METRIC_REPLICATION_RESYNC_STARTUP_BACKGROUND_DURATION_SECONDS, + "mode" => mode, + "outcome" => outcome + ) + .record(duration.as_secs_f64()); + metrics::gauge!(METRIC_REPLICATION_RESYNC_STARTUP_BACKGROUND_STATUS, "mode" => mode).set(status); +} + async fn run_bucket_resync_startup_reconcile( buckets: Vec, ctx: CancellationToken, @@ -117,4 +240,14 @@ mod tests { "replication pool is not initialized" ))); } + + #[test] + fn startup_resync_background_observability_uses_fixed_modes_and_statuses() { + assert_eq!(bucket_resync_startup_background_mode(true), STARTUP_BACKGROUND_MODE_SERVER); + assert_eq!(bucket_resync_startup_background_mode(false), STARTUP_BACKGROUND_MODE_EMBEDDED); + assert_eq!(STARTUP_BACKGROUND_STATUS_FAILED, 0.0); + assert_eq!(STARTUP_BACKGROUND_STATUS_SUCCEEDED, 1.0); + assert_eq!(STARTUP_BACKGROUND_STATUS_RUNNING, 2.0); + assert_eq!(STARTUP_BACKGROUND_STATUS_CANCELED, 3.0); + } } diff --git a/rustfs/tests/connect_registration.rs b/rustfs/tests/connect_registration.rs index c6d5e8c61..2a52a9551 100644 --- a/rustfs/tests/connect_registration.rs +++ b/rustfs/tests/connect_registration.rs @@ -46,7 +46,7 @@ use sha2::{Digest as _, Sha256}; use time::OffsetDateTime; use time::format_description::well_known::Rfc3339; use tokio::net::TcpListener; -use tokio::sync::watch; +use tokio::sync::{Notify, watch}; use tokio_rustls::TlsAcceptor; use tokio_util::sync::CancellationToken; @@ -204,6 +204,7 @@ struct TestServer { seen: Arc>>, paths: Arc>>, client_certificates: Arc>>>, + request_notify: Arc, task: tokio::task::JoinHandle<()>, } @@ -225,9 +226,11 @@ async fn server_with_client_auth(pki: &TestPki, replies: Vec, require_cli let seen = Arc::new(Mutex::new(Vec::new())); let paths = Arc::new(Mutex::new(Vec::new())); let client_certificates = Arc::new(Mutex::new(Vec::new())); + let request_notify = Arc::new(Notify::new()); let captured = seen.clone(); let captured_paths = paths.clone(); let captured_certificates = client_certificates.clone(); + let captured_notify = request_notify.clone(); let task = tokio::spawn(async move { loop { let Ok((stream, _)) = listener.accept().await else { @@ -238,6 +241,7 @@ async fn server_with_client_auth(pki: &TestPki, replies: Vec, require_cli let seen = captured.clone(); let paths = captured_paths.clone(); let client_certificates = captured_certificates.clone(); + let request_notify = captured_notify.clone(); tokio::spawn(async move { let Ok(stream) = acceptor.accept(stream).await else { return; @@ -253,6 +257,7 @@ async fn server_with_client_auth(pki: &TestPki, replies: Vec, require_cli let seen = seen.clone(); let paths = paths.clone(); let client_certificates = client_certificates.clone(); + let request_notify = request_notify.clone(); let client_certificate = client_certificate.clone(); async move { let path = request.uri().path().to_owned(); @@ -265,6 +270,7 @@ async fn server_with_client_auth(pki: &TestPki, replies: Vec, require_cli .push(client_certificate); paths.lock().expect("paths lock").push(path); seen.lock().expect("seen lock").push(value.clone()); + request_notify.notify_waiters(); match reply { Reply::Json(status, value) => Ok::<_, hyper::Error>( Response::builder() @@ -344,6 +350,7 @@ async fn server_with_client_auth(pki: &TestPki, replies: Vec, require_cli seen, paths, client_certificates, + request_notify, task, } } @@ -545,9 +552,12 @@ fn heartbeat_response(server_time: &str) -> Value { } async fn wait_for_requests(server: &TestServer, count: usize) { - tokio::time::timeout(Duration::from_secs(3), async { + tokio::time::timeout(Duration::from_secs(10), async { while server.paths.lock().expect("paths lock").len() < count { - tokio::task::yield_now().await; + tokio::select! { + _ = server.request_notify.notified() => {} + _ = tokio::time::sleep(Duration::from_millis(5)) => {} + } } }) .await