mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-02 05:14:35 +08:00
fix(connect): accept the health-era heartbeat and measured Top windows (#8195)
Two Connect tests fail on main in every Test and Lint lane. Heartbeat: #8155 replaced "today's capabilities minus profile.memory.service@1" with "the pre-health set minus it". That drops the release that advertised health.check.service@1 but not yet in-service memory profiling, so a pending heartbeat persisted by that release now fails validation and the runtime stops. Accept that release as its own frozen list. Top disk/net: execute_*_job sets max_duration_millis to the requested duration, then the capture compared the measured sleep against it. The timer only overshoots, so a job that ran exactly as requested was rejected with LimitExceeded whenever the overshoot reached 1 ms. Report the authorized window, as Top API already does; validate_capture bounds it by the limit.
This commit is contained in:
@@ -15,8 +15,6 @@
|
||||
//! Bounded process disk-I/O window backed by RustFS's existing process sampler.
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
#[cfg(target_os = "linux")]
|
||||
use tokio::time::Instant;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
|
||||
use super::top_api::{MAX_SAFE_INTEGER, TopCaptureError, TopCaptureRequest, TopReasonCode, TopResult};
|
||||
@@ -91,12 +89,15 @@ pub async fn capture_top_disk(
|
||||
}
|
||||
let mut sampler = rustfs_io_metrics::ProcessSampler::new();
|
||||
let before = process_snapshot(&mut sampler)?;
|
||||
let started = Instant::now();
|
||||
if !request.wait_window(TOOL_ID, cancel).await? {
|
||||
return request.cancelled(TOOL_ID);
|
||||
}
|
||||
let after = process_snapshot(&mut sampler)?;
|
||||
evaluate_disk_window(request, before, after, elapsed_millis(started.elapsed()))
|
||||
// Report the authorized window, as Top API does. The timer can only
|
||||
// overshoot it, and validate_capture already bounded it by the limit;
|
||||
// measuring the sleep would reject a capture that ran as requested.
|
||||
let window_millis = u64::try_from(request.window.as_millis()).map_err(|_| TopCaptureError::Limits)?;
|
||||
evaluate_disk_window(request, before, after, window_millis)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -154,8 +155,3 @@ fn process_snapshot(sampler: &mut rustfs_io_metrics::ProcessSampler) -> Result<D
|
||||
io_count,
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
fn elapsed_millis(duration: std::time::Duration) -> u64 {
|
||||
u64::try_from(duration.as_millis()).unwrap_or(u64::MAX).max(1)
|
||||
}
|
||||
|
||||
@@ -16,7 +16,6 @@
|
||||
|
||||
use serde::Serialize;
|
||||
use sysinfo::Networks;
|
||||
use tokio::time::Instant;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
|
||||
use super::top_api::{MAX_SAFE_INTEGER, TopCaptureError, TopCaptureRequest, TopReasonCode, TopResult};
|
||||
@@ -57,10 +56,11 @@ pub async fn capture_top_net(
|
||||
if networks.is_empty() {
|
||||
return request.failed(TOOL_ID, 0, TopReasonCode::SourceUnavailable);
|
||||
}
|
||||
let started = Instant::now();
|
||||
if !request.wait_window(TOOL_ID, cancel).await? {
|
||||
return request.cancelled(TOOL_ID);
|
||||
}
|
||||
// Report the authorized window; see capture_top_disk.
|
||||
let window_millis = u64::try_from(request.window.as_millis()).map_err(|_| TopCaptureError::Limits)?;
|
||||
let observed = rustfs_obs::metrics::stats_collector::collect_host_network_stats(&mut networks);
|
||||
evaluate_network_window(
|
||||
request,
|
||||
@@ -72,7 +72,7 @@ pub async fn capture_top_net(
|
||||
received_bytes: observed.total_received,
|
||||
sent_bytes: observed.total_transmitted,
|
||||
},
|
||||
elapsed_millis(started.elapsed()),
|
||||
window_millis,
|
||||
)
|
||||
}
|
||||
|
||||
@@ -109,7 +109,3 @@ pub fn evaluate_network_window(
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
fn elapsed_millis(duration: std::time::Duration) -> u64 {
|
||||
u64::try_from(duration.as_millis()).unwrap_or(u64::MAX).max(1)
|
||||
}
|
||||
|
||||
@@ -121,6 +121,7 @@ impl PendingHeartbeat {
|
||||
.iter()
|
||||
.filter(|capability| capability.as_str() != "profile.memory.service@1")
|
||||
.eq(self.capabilities.iter())
|
||||
|| self.capabilities == pre_service_memory_heartbeat_capabilities(job_capable)
|
||||
}))
|
||||
&& self.sequence <= MAX_SEQUENCE
|
||||
&& self.coarse_node_summary.is_valid()
|
||||
@@ -367,6 +368,42 @@ fn pre_health_heartbeat_capabilities(job_capable: bool) -> Vec<String> {
|
||||
capabilities
|
||||
}
|
||||
|
||||
/// The release that advertised the health check but not yet in-service
|
||||
/// memory profiling. Frozen: do not derive it from today's capabilities.
|
||||
fn pre_service_memory_heartbeat_capabilities(job_capable: bool) -> Vec<String> {
|
||||
let mut capabilities = [
|
||||
"heartbeat",
|
||||
"diagnostics.policy.v1",
|
||||
"inventory.environment@1",
|
||||
"health.check.service@1",
|
||||
"performance.client@1",
|
||||
"performance.drive@1",
|
||||
"performance.network@1",
|
||||
"performance.object@1",
|
||||
"performance.siteReplication@1",
|
||||
"logs.capture@1",
|
||||
"profile.cpu@1",
|
||||
"profile.memory@1",
|
||||
"profile.threads@1",
|
||||
"telemetry.record@1",
|
||||
"telemetry.otlp@1",
|
||||
"telemetry.replay@1",
|
||||
"top.api@1",
|
||||
"top.disk@1",
|
||||
"top.locks@1",
|
||||
"top.net@1",
|
||||
"top.rpc@1",
|
||||
"inspect.object@1",
|
||||
]
|
||||
.into_iter()
|
||||
.map(str::to_owned)
|
||||
.collect::<Vec<_>>();
|
||||
if job_capable {
|
||||
capabilities.insert(3, "jobs".to_owned());
|
||||
}
|
||||
capabilities
|
||||
}
|
||||
|
||||
fn heartbeat_capabilities(job_capable: bool) -> Vec<String> {
|
||||
let mut capabilities = vec![
|
||||
"heartbeat".to_owned(),
|
||||
|
||||
Reference in New Issue
Block a user