mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-18 00:25:31 +00:00
feat(connect): add typed telemetry trace source (#7775)
This commit is contained in:
Generated
+1
@@ -10035,6 +10035,7 @@ dependencies = [
|
|||||||
"metrics",
|
"metrics",
|
||||||
"metrics-util",
|
"metrics-util",
|
||||||
"num_cpus",
|
"num_cpus",
|
||||||
|
"rustfs-common",
|
||||||
"rustfs-s3-ops",
|
"rustfs-s3-ops",
|
||||||
"sysinfo",
|
"sysinfo",
|
||||||
"thiserror 2.0.20",
|
"thiserror 2.0.20",
|
||||||
|
|||||||
@@ -127,6 +127,44 @@ pub struct TraceEvent {
|
|||||||
pub attrs: SmallVec<[TraceAttr; TRACE_ATTR_INLINE_CAPACITY]>,
|
pub attrs: SmallVec<[TraceAttr; TRACE_ATTR_INLINE_CAPACITY]>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A telemetry operation whose published shape cannot retain request data.
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub enum TelemetryTraceOperation {
|
||||||
|
GetObject,
|
||||||
|
PutObject,
|
||||||
|
HeadObject,
|
||||||
|
ListObjects,
|
||||||
|
InternalRpc,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub enum TelemetryTraceStatus {
|
||||||
|
Ok,
|
||||||
|
Error,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A pre-classified S3 or internode RPC observation.
|
||||||
|
///
|
||||||
|
/// This event deliberately has no request path, headers, object identity,
|
||||||
|
/// payload, or free-form attributes. Producers must classify the operation
|
||||||
|
/// and outcome before publishing it.
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub struct TelemetryTraceEvent {
|
||||||
|
pub operation: TelemetryTraceOperation,
|
||||||
|
pub duration: Duration,
|
||||||
|
pub status: TelemetryTraceStatus,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl TelemetryTraceEvent {
|
||||||
|
pub const fn new(operation: TelemetryTraceOperation, duration: Duration, status: TelemetryTraceStatus) -> Self {
|
||||||
|
Self {
|
||||||
|
operation,
|
||||||
|
duration,
|
||||||
|
status,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
impl TraceEvent {
|
impl TraceEvent {
|
||||||
pub fn new(kind: TraceKind, func: TraceFunc) -> Self {
|
pub fn new(kind: TraceKind, func: TraceFunc) -> Self {
|
||||||
Self {
|
Self {
|
||||||
@@ -174,15 +212,20 @@ impl TraceEvent {
|
|||||||
pub struct TraceBus {
|
pub struct TraceBus {
|
||||||
sender: broadcast::Sender<Arc<TraceEvent>>,
|
sender: broadcast::Sender<Arc<TraceEvent>>,
|
||||||
subscriber_count: Arc<AtomicUsize>,
|
subscriber_count: Arc<AtomicUsize>,
|
||||||
|
telemetry_sender: broadcast::Sender<Arc<TelemetryTraceEvent>>,
|
||||||
|
telemetry_subscriber_count: Arc<AtomicUsize>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl TraceBus {
|
impl TraceBus {
|
||||||
pub fn new(capacity: usize) -> Self {
|
pub fn new(capacity: usize) -> Self {
|
||||||
let capacity = capacity.max(1);
|
let capacity = capacity.max(1);
|
||||||
let (sender, _receiver) = broadcast::channel(capacity);
|
let (sender, _receiver) = broadcast::channel(capacity);
|
||||||
|
let (telemetry_sender, _telemetry_receiver) = broadcast::channel(capacity);
|
||||||
Self {
|
Self {
|
||||||
sender,
|
sender,
|
||||||
subscriber_count: Arc::new(AtomicUsize::new(0)),
|
subscriber_count: Arc::new(AtomicUsize::new(0)),
|
||||||
|
telemetry_sender,
|
||||||
|
telemetry_subscriber_count: Arc::new(AtomicUsize::new(0)),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -206,6 +249,27 @@ impl TraceBus {
|
|||||||
|
|
||||||
self.sender.send(Arc::new(build())).is_ok()
|
self.sender.send(Arc::new(build())).is_ok()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn telemetry_subscriber_count(&self) -> usize {
|
||||||
|
self.telemetry_subscriber_count.load(Ordering::Acquire)
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn subscribe_telemetry(&self) -> TelemetryTraceSubscription {
|
||||||
|
let receiver = self.telemetry_sender.subscribe();
|
||||||
|
self.telemetry_subscriber_count.fetch_add(1, Ordering::AcqRel);
|
||||||
|
TelemetryTraceSubscription {
|
||||||
|
receiver,
|
||||||
|
subscriber_count: Arc::clone(&self.telemetry_subscriber_count),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn emit_telemetry(&self, build: impl FnOnce() -> TelemetryTraceEvent) -> bool {
|
||||||
|
if self.telemetry_subscriber_count() == 0 {
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
|
||||||
|
self.telemetry_sender.send(Arc::new(build())).is_ok()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Default for TraceBus {
|
impl Default for TraceBus {
|
||||||
@@ -236,6 +300,28 @@ impl Drop for TraceSubscription {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug)]
|
||||||
|
pub struct TelemetryTraceSubscription {
|
||||||
|
receiver: broadcast::Receiver<Arc<TelemetryTraceEvent>>,
|
||||||
|
subscriber_count: Arc<AtomicUsize>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl TelemetryTraceSubscription {
|
||||||
|
pub async fn recv(&mut self) -> Result<Arc<TelemetryTraceEvent>, broadcast::error::RecvError> {
|
||||||
|
self.receiver.recv().await
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn try_recv(&mut self) -> Result<Arc<TelemetryTraceEvent>, broadcast::error::TryRecvError> {
|
||||||
|
self.receiver.try_recv()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Drop for TelemetryTraceSubscription {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
self.subscriber_count.fetch_sub(1, Ordering::AcqRel);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
pub fn global_trace_bus() -> &'static TraceBus {
|
pub fn global_trace_bus() -> &'static TraceBus {
|
||||||
GLOBAL_TRACE_BUS.get_or_init(TraceBus::default)
|
GLOBAL_TRACE_BUS.get_or_init(TraceBus::default)
|
||||||
}
|
}
|
||||||
@@ -252,6 +338,18 @@ pub fn trace_subscriber_count() -> usize {
|
|||||||
global_trace_bus().subscriber_count()
|
global_trace_bus().subscriber_count()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn subscribe_telemetry_trace_events() -> TelemetryTraceSubscription {
|
||||||
|
global_trace_bus().subscribe_telemetry()
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn telemetry_trace_emit(build: impl FnOnce() -> TelemetryTraceEvent) -> bool {
|
||||||
|
global_trace_bus().emit_telemetry(build)
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn telemetry_trace_subscriber_count() -> usize {
|
||||||
|
global_trace_bus().telemetry_subscriber_count()
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
@@ -330,4 +428,30 @@ mod tests {
|
|||||||
.expect_err("receiver should observe lag instead of blocking publishers");
|
.expect_err("receiver should observe lag instead of blocking publishers");
|
||||||
assert!(matches!(err, broadcast::error::RecvError::Lagged(_)));
|
assert!(matches!(err, broadcast::error::RecvError::Lagged(_)));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn telemetry_subscription_exposes_only_classified_fields() {
|
||||||
|
let bus = TraceBus::new(4);
|
||||||
|
let _unrelated_subscription = bus.subscribe();
|
||||||
|
let built = AtomicUsize::new(0);
|
||||||
|
assert!(!bus.emit_telemetry(|| {
|
||||||
|
built.fetch_add(1, Ordering::Relaxed);
|
||||||
|
TelemetryTraceEvent::new(TelemetryTraceOperation::GetObject, Duration::ZERO, TelemetryTraceStatus::Ok)
|
||||||
|
}));
|
||||||
|
assert_eq!(built.load(Ordering::Relaxed), 0);
|
||||||
|
|
||||||
|
let mut subscription = bus.subscribe_telemetry();
|
||||||
|
assert_eq!(bus.telemetry_subscriber_count(), 1);
|
||||||
|
|
||||||
|
assert!(bus.emit_telemetry(|| {
|
||||||
|
TelemetryTraceEvent::new(TelemetryTraceOperation::GetObject, Duration::from_micros(7), TelemetryTraceStatus::Ok)
|
||||||
|
}));
|
||||||
|
|
||||||
|
let event = subscription.recv().await.expect("classified telemetry event");
|
||||||
|
assert_eq!(event.operation, TelemetryTraceOperation::GetObject);
|
||||||
|
assert_eq!(event.duration, Duration::from_micros(7));
|
||||||
|
assert_eq!(event.status, TelemetryTraceStatus::Ok);
|
||||||
|
drop(subscription);
|
||||||
|
assert_eq!(bus.telemetry_subscriber_count(), 0);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -49,6 +49,7 @@ hotpath-cpu = [
|
|||||||
[dependencies]
|
[dependencies]
|
||||||
hotpath.workspace = true
|
hotpath.workspace = true
|
||||||
metrics = { workspace = true }
|
metrics = { workspace = true }
|
||||||
|
rustfs-common = { workspace = true }
|
||||||
rustfs-s3-ops = { workspace = true }
|
rustfs-s3-ops = { workspace = true }
|
||||||
num_cpus = { workspace = true }
|
num_cpus = { workspace = true }
|
||||||
thiserror = { workspace = true }
|
thiserror = { workspace = true }
|
||||||
|
|||||||
@@ -16,10 +16,14 @@
|
|||||||
//! Admin snapshots and metric exporters share these counters. The older
|
//! Admin snapshots and metric exporters share these counters. The older
|
||||||
//! operation counter counts handler entries and is not an HTTP denominator.
|
//! operation counter counts handler entries and is not an HTTP denominator.
|
||||||
|
|
||||||
|
use rustfs_common::trace_bus::{
|
||||||
|
TelemetryTraceEvent, TelemetryTraceOperation, TelemetryTraceStatus, telemetry_trace_emit, telemetry_trace_subscriber_count,
|
||||||
|
};
|
||||||
use rustfs_s3_ops::S3Operation;
|
use rustfs_s3_ops::S3Operation;
|
||||||
use std::cell::Cell;
|
use std::cell::Cell;
|
||||||
use std::sync::atomic::{AtomicU64, Ordering};
|
use std::sync::atomic::{AtomicU64, Ordering};
|
||||||
use std::sync::{LazyLock, OnceLock};
|
use std::sync::{LazyLock, OnceLock};
|
||||||
|
use std::time::Instant;
|
||||||
|
|
||||||
const METRIC: &str = "rustfs_s3_http_requests_total";
|
const METRIC: &str = "rustfs_s3_http_requests_total";
|
||||||
const METHODS: [&str; 10] = [
|
const METHODS: [&str; 10] = [
|
||||||
@@ -110,6 +114,7 @@ pub(crate) fn observe_s3_http_operation(op: S3Operation) {
|
|||||||
pub struct S3HttpRequestGuard {
|
pub struct S3HttpRequestGuard {
|
||||||
method: usize,
|
method: usize,
|
||||||
operation: usize,
|
operation: usize,
|
||||||
|
telemetry_started_at: Option<Instant>,
|
||||||
finished: bool,
|
finished: bool,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -122,6 +127,7 @@ impl S3HttpRequestGuard {
|
|||||||
Self {
|
Self {
|
||||||
method: METHODS.iter().position(|known| *known == method).unwrap_or(METHODS.len() - 1),
|
method: METHODS.iter().position(|known| *known == method).unwrap_or(METHODS.len() - 1),
|
||||||
operation: UNKNOWN_OPERATION,
|
operation: UNKNOWN_OPERATION,
|
||||||
|
telemetry_started_at: (telemetry_trace_subscriber_count() != 0).then(Instant::now),
|
||||||
finished: false,
|
finished: false,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -151,11 +157,29 @@ impl S3HttpRequestGuard {
|
|||||||
fn finish(&mut self, outcome: usize) {
|
fn finish(&mut self, outcome: usize) {
|
||||||
if !self.finished {
|
if !self.finished {
|
||||||
COUNTERS.record(self.method, self.operation, outcome);
|
COUNTERS.record(self.method, self.operation, outcome);
|
||||||
|
if let Some((started_at, operation)) = self.telemetry_started_at.take().zip(telemetry_operation(self.operation)) {
|
||||||
|
let status = if outcome == 1 {
|
||||||
|
TelemetryTraceStatus::Ok
|
||||||
|
} else {
|
||||||
|
TelemetryTraceStatus::Error
|
||||||
|
};
|
||||||
|
telemetry_trace_emit(|| TelemetryTraceEvent::new(operation, started_at.elapsed(), status));
|
||||||
|
}
|
||||||
self.finished = true;
|
self.finished = true;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn telemetry_operation(index: usize) -> Option<TelemetryTraceOperation> {
|
||||||
|
match S3Operation::ALL.get(index)? {
|
||||||
|
S3Operation::GetObject => Some(TelemetryTraceOperation::GetObject),
|
||||||
|
S3Operation::PutObject => Some(TelemetryTraceOperation::PutObject),
|
||||||
|
S3Operation::HeadObject => Some(TelemetryTraceOperation::HeadObject),
|
||||||
|
S3Operation::ListObjects | S3Operation::ListObjectsV2 => Some(TelemetryTraceOperation::ListObjects),
|
||||||
|
_ => None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
impl Drop for S3HttpRequestGuard {
|
impl Drop for S3HttpRequestGuard {
|
||||||
fn drop(&mut self) {
|
fn drop(&mut self) {
|
||||||
self.finish(7);
|
self.finish(7);
|
||||||
@@ -172,6 +196,32 @@ mod tests {
|
|||||||
use metrics::with_local_recorder;
|
use metrics::with_local_recorder;
|
||||||
use metrics_util::debugging::DebuggingRecorder;
|
use metrics_util::debugging::DebuggingRecorder;
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn telemetry_adapter_accepts_only_the_frozen_s3_operations() {
|
||||||
|
assert_eq!(
|
||||||
|
telemetry_operation(S3Operation::GetObject.metric_index()),
|
||||||
|
Some(TelemetryTraceOperation::GetObject)
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
telemetry_operation(S3Operation::PutObject.metric_index()),
|
||||||
|
Some(TelemetryTraceOperation::PutObject)
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
telemetry_operation(S3Operation::HeadObject.metric_index()),
|
||||||
|
Some(TelemetryTraceOperation::HeadObject)
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
telemetry_operation(S3Operation::ListObjects.metric_index()),
|
||||||
|
Some(TelemetryTraceOperation::ListObjects)
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
telemetry_operation(S3Operation::ListObjectsV2.metric_index()),
|
||||||
|
Some(TelemetryTraceOperation::ListObjects)
|
||||||
|
);
|
||||||
|
assert_eq!(telemetry_operation(S3Operation::DeleteObject.metric_index()), None);
|
||||||
|
assert_eq!(telemetry_operation(UNKNOWN_OPERATION), None);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn outcome_counters_distinguish_partial_and_complete_write_failure() {
|
fn outcome_counters_distinguish_partial_and_complete_write_failure() {
|
||||||
let counters = HttpOutcomeCounters::new();
|
let counters = HttpOutcomeCounters::new();
|
||||||
|
|||||||
@@ -27,7 +27,9 @@ use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
|
|||||||
use base64_simd::URL_SAFE_NO_PAD;
|
use base64_simd::URL_SAFE_NO_PAD;
|
||||||
use p256::ecdsa::{Signature, SigningKey, signature::Signer as _};
|
use p256::ecdsa::{Signature, SigningKey, signature::Signer as _};
|
||||||
use p256::pkcs8::DecodePrivateKey as _;
|
use p256::pkcs8::DecodePrivateKey as _;
|
||||||
use rustfs_common::trace_bus::subscribe_trace_events;
|
use rustfs_common::trace_bus::{
|
||||||
|
TelemetryTraceEvent, TelemetryTraceOperation, TelemetryTraceStatus, subscribe_telemetry_trace_events,
|
||||||
|
};
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use sha2::{Digest as _, Sha256};
|
use sha2::{Digest as _, Sha256};
|
||||||
use thiserror::Error;
|
use thiserror::Error;
|
||||||
@@ -163,6 +165,14 @@ impl LocalTelemetryConsent {
|
|||||||
.filter(|remaining| !remaining.is_zero())
|
.filter(|remaining| !remaining.is_zero())
|
||||||
.ok_or(TelemetryProducerError::ConsentExpired)
|
.ok_or(TelemetryProducerError::ConsentExpired)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn capture_deadline(self, duration: Duration) -> Result<Instant, TelemetryProducerError> {
|
||||||
|
let deadline = Instant::now() + duration;
|
||||||
|
if deadline > self.expires_at {
|
||||||
|
return Err(TelemetryProducerError::ConsentExpired);
|
||||||
|
}
|
||||||
|
Ok(deadline)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Clone, Debug, Eq, PartialEq)]
|
#[derive(Clone, Debug, Eq, PartialEq)]
|
||||||
@@ -563,12 +573,8 @@ pub async fn record_trace(
|
|||||||
if cancel.is_cancelled() {
|
if cancel.is_cancelled() {
|
||||||
return Err(TelemetryProducerError::Cancelled);
|
return Err(TelemetryProducerError::Cancelled);
|
||||||
}
|
}
|
||||||
let remaining = consent.remaining()?;
|
let deadline = consent.capture_deadline(limits.duration)?;
|
||||||
if remaining < limits.duration {
|
|
||||||
return Err(TelemetryProducerError::ConsentExpired);
|
|
||||||
}
|
|
||||||
let _lease = acquire_telemetry_lease()?;
|
let _lease = acquire_telemetry_lease()?;
|
||||||
let deadline = Instant::now() + limits.duration;
|
|
||||||
let mut spans = Vec::with_capacity(limits.max_spans.min(64));
|
let mut spans = Vec::with_capacity(limits.max_spans.min(64));
|
||||||
let mut dropped_span_count = 0u64;
|
let mut dropped_span_count = 0u64;
|
||||||
let mut completion = TraceRecordCompletion::Complete;
|
let mut completion = TraceRecordCompletion::Complete;
|
||||||
@@ -610,13 +616,7 @@ pub async fn record_trace(
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Capture the process-local RustFS trace bus.
|
/// Capture the process-local, pre-classified S3 and internode RPC trace source.
|
||||||
///
|
|
||||||
/// The current bus exposes only heal and scanner operations, none of which has
|
|
||||||
/// the frozen GET/PUT/HEAD/LIST/INTERNAL_RPC semantics. The adapter consumes
|
|
||||||
/// that real source but refuses to infer an operation or status from raw
|
|
||||||
/// fields, so it remains unavailable until RustFS publishes an approved typed
|
|
||||||
/// event.
|
|
||||||
pub async fn record_trace_bus(
|
pub async fn record_trace_bus(
|
||||||
consent: LocalTelemetryConsent,
|
consent: LocalTelemetryConsent,
|
||||||
limits: TraceRecordLimits,
|
limits: TraceRecordLimits,
|
||||||
@@ -626,22 +626,78 @@ pub async fn record_trace_bus(
|
|||||||
if cancel.is_cancelled() {
|
if cancel.is_cancelled() {
|
||||||
return Err(TelemetryProducerError::Cancelled);
|
return Err(TelemetryProducerError::Cancelled);
|
||||||
}
|
}
|
||||||
let remaining = consent.remaining()?;
|
let deadline = consent.capture_deadline(limits.duration)?;
|
||||||
if remaining < limits.duration {
|
|
||||||
return Err(TelemetryProducerError::ConsentExpired);
|
|
||||||
}
|
|
||||||
let _lease = acquire_telemetry_lease()?;
|
let _lease = acquire_telemetry_lease()?;
|
||||||
let deadline = Instant::now() + limits.duration;
|
let mut subscription = subscribe_telemetry_trace_events();
|
||||||
let mut subscription = subscribe_trace_events();
|
let mut spans = Vec::with_capacity(limits.max_spans.min(64));
|
||||||
|
let mut dropped_span_count = 0u64;
|
||||||
|
let mut completion = TraceRecordCompletion::Complete;
|
||||||
loop {
|
loop {
|
||||||
tokio::select! {
|
tokio::select! {
|
||||||
biased;
|
biased;
|
||||||
_ = cancel.cancelled() => return Err(TelemetryProducerError::Cancelled),
|
_ = cancel.cancelled() => return Err(TelemetryProducerError::Cancelled),
|
||||||
_ = tokio::time::sleep_until(deadline.into()) => return Err(TelemetryProducerError::SourceUnavailable),
|
_ = tokio::time::sleep_until(deadline.into()) => break,
|
||||||
received = subscription.recv() => match received {
|
received = subscription.recv() => match received {
|
||||||
Ok(_event) => {}
|
Ok(event) => {
|
||||||
Err(tokio::sync::broadcast::error::RecvError::Lagged(_dropped)) => {}
|
if spans.len() == limits.max_spans {
|
||||||
Err(tokio::sync::broadcast::error::RecvError::Closed) => return Err(TelemetryProducerError::SourceUnavailable),
|
completion = TraceRecordCompletion::LimitExceeded;
|
||||||
|
dropped_span_count = dropped_span_count.saturating_add(1).min(MAX_SAFE_INTEGER);
|
||||||
|
drain_telemetry_drops(&mut subscription, &mut dropped_span_count);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
spans.push(to_contract_span(observe_trace_bus_event(&event))?);
|
||||||
|
}
|
||||||
|
Err(tokio::sync::broadcast::error::RecvError::Lagged(dropped)) => {
|
||||||
|
completion = TraceRecordCompletion::LimitExceeded;
|
||||||
|
dropped_span_count = dropped_span_count.saturating_add(dropped).min(MAX_SAFE_INTEGER);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
Err(tokio::sync::broadcast::error::RecvError::Closed) => {
|
||||||
|
if spans.is_empty() {
|
||||||
|
return Err(TelemetryProducerError::SourceUnavailable);
|
||||||
|
}
|
||||||
|
completion = TraceRecordCompletion::SourceUnavailable;
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let result = RecordedTrace {
|
||||||
|
spans,
|
||||||
|
dropped_span_count,
|
||||||
|
};
|
||||||
|
ensure_result_size(&result)?;
|
||||||
|
Ok(TraceRecordCapture {
|
||||||
|
data: result,
|
||||||
|
completion,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
fn observe_trace_bus_event(event: &TelemetryTraceEvent) -> ObservedTelemetrySpan {
|
||||||
|
let operation = match event.operation {
|
||||||
|
TelemetryTraceOperation::GetObject => TelemetryOperation::GetObject,
|
||||||
|
TelemetryTraceOperation::PutObject => TelemetryOperation::PutObject,
|
||||||
|
TelemetryTraceOperation::HeadObject => TelemetryOperation::HeadObject,
|
||||||
|
TelemetryTraceOperation::ListObjects => TelemetryOperation::ListObjects,
|
||||||
|
TelemetryTraceOperation::InternalRpc => TelemetryOperation::InternalRpc,
|
||||||
|
};
|
||||||
|
let status = match event.status {
|
||||||
|
TelemetryTraceStatus::Ok => TelemetrySpanStatus::Ok,
|
||||||
|
TelemetryTraceStatus::Error => TelemetrySpanStatus::Error,
|
||||||
|
};
|
||||||
|
ObservedTelemetrySpan::new(operation, event.duration, status)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn drain_telemetry_drops(subscription: &mut rustfs_common::trace_bus::TelemetryTraceSubscription, dropped_span_count: &mut u64) {
|
||||||
|
loop {
|
||||||
|
match subscription.try_recv() {
|
||||||
|
Ok(_) => *dropped_span_count = dropped_span_count.saturating_add(1).min(MAX_SAFE_INTEGER),
|
||||||
|
Err(tokio::sync::broadcast::error::TryRecvError::Lagged(dropped)) => {
|
||||||
|
*dropped_span_count = dropped_span_count.saturating_add(dropped).min(MAX_SAFE_INTEGER);
|
||||||
|
}
|
||||||
|
Err(tokio::sync::broadcast::error::TryRecvError::Empty | tokio::sync::broadcast::error::TryRecvError::Closed) => {
|
||||||
|
break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -36,6 +36,7 @@ use futures_util::{Stream, StreamExt, TryStreamExt, stream};
|
|||||||
use http::{HeaderMap, HeaderValue, Method, Request, Response, StatusCode, Uri};
|
use http::{HeaderMap, HeaderValue, Method, Request, Response, StatusCode, Uri};
|
||||||
use http_body_util::{BodyExt, Limited};
|
use http_body_util::{BodyExt, Limited};
|
||||||
use hyper::body::Incoming;
|
use hyper::body::Incoming;
|
||||||
|
use rustfs_common::trace_bus::{TelemetryTraceEvent, TelemetryTraceOperation, TelemetryTraceStatus, telemetry_trace_emit};
|
||||||
use rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE;
|
use rustfs_config::MAX_ADMIN_REQUEST_BODY_SIZE;
|
||||||
use rustfs_io_metrics::internode_metrics::{
|
use rustfs_io_metrics::internode_metrics::{
|
||||||
INTERNODE_OPERATION_NS_SCANNER, INTERNODE_OPERATION_PUT_FILE_CAPABILITY, INTERNODE_OPERATION_PUT_FILE_STREAM,
|
INTERNODE_OPERATION_NS_SCANNER, INTERNODE_OPERATION_PUT_FILE_CAPABILITY, INTERNODE_OPERATION_PUT_FILE_STREAM,
|
||||||
@@ -415,7 +416,9 @@ async fn handle_internode_rpc(req: Request<Incoming>) -> Response<Body> {
|
|||||||
let started_at = Instant::now();
|
let started_at = Instant::now();
|
||||||
if let Err(response) = verify_internode_rpc_signature(req.uri(), req.method(), req.headers()) {
|
if let Err(response) = verify_internode_rpc_signature(req.uri(), req.method(), req.headers()) {
|
||||||
record_internode_rpc_error(operation);
|
record_internode_rpc_error(operation);
|
||||||
return *response;
|
let response = *response;
|
||||||
|
emit_internode_rpc_telemetry(started_at, response.status());
|
||||||
|
return response;
|
||||||
}
|
}
|
||||||
|
|
||||||
let method = req.method().clone();
|
let method = req.method().clone();
|
||||||
@@ -461,10 +464,20 @@ async fn handle_internode_rpc(req: Request<Incoming>) -> Response<Body> {
|
|||||||
started_at.elapsed(),
|
started_at.elapsed(),
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
emit_internode_rpc_telemetry(started_at, response.status());
|
||||||
|
|
||||||
response
|
response
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn emit_internode_rpc_telemetry(started_at: Instant, status: StatusCode) {
|
||||||
|
let status = if status.is_success() {
|
||||||
|
TelemetryTraceStatus::Ok
|
||||||
|
} else {
|
||||||
|
TelemetryTraceStatus::Error
|
||||||
|
};
|
||||||
|
telemetry_trace_emit(|| TelemetryTraceEvent::new(TelemetryTraceOperation::InternalRpc, started_at.elapsed(), status));
|
||||||
|
}
|
||||||
|
|
||||||
fn internode_http_operation(path: &str) -> Option<&'static str> {
|
fn internode_http_operation(path: &str) -> Option<&'static str> {
|
||||||
match path {
|
match path {
|
||||||
READ_FILE_STREAM_PATH => Some(INTERNODE_OPERATION_READ_FILE_STREAM),
|
READ_FILE_STREAM_PATH => Some(INTERNODE_OPERATION_READ_FILE_STREAM),
|
||||||
@@ -1761,6 +1774,21 @@ mod tests {
|
|||||||
(disk, dir)
|
(disk, dir)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial_test::serial]
|
||||||
|
async fn internode_rpc_telemetry_publishes_only_the_classified_outcome() {
|
||||||
|
let mut subscription = rustfs_common::trace_bus::subscribe_telemetry_trace_events();
|
||||||
|
|
||||||
|
super::emit_internode_rpc_telemetry(std::time::Instant::now(), StatusCode::FORBIDDEN);
|
||||||
|
|
||||||
|
let event = tokio::time::timeout(Duration::from_secs(1), subscription.recv())
|
||||||
|
.await
|
||||||
|
.expect("telemetry event should arrive")
|
||||||
|
.expect("telemetry source should remain open");
|
||||||
|
assert_eq!(event.operation, rustfs_common::trace_bus::TelemetryTraceOperation::InternalRpc);
|
||||||
|
assert_eq!(event.status, rustfs_common::trace_bus::TelemetryTraceStatus::Error);
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial_test::serial]
|
#[serial_test::serial]
|
||||||
async fn authenticated_put_route_checks_server_epoch_before_disk_lookup() {
|
async fn authenticated_put_route_checks_server_epoch_before_disk_lookup() {
|
||||||
|
|||||||
@@ -14,7 +14,12 @@ use rustfs::connect::{
|
|||||||
TelemetrySpanStatus, TelemetryTool, TraceRecordCompletion, TraceRecordLimits, encode_signed_telemetry_export,
|
TelemetrySpanStatus, TelemetryTool, TraceRecordCompletion, TraceRecordLimits, encode_signed_telemetry_export,
|
||||||
record_diagnostic_result, record_trace, record_trace_bus, save_signed_telemetry_export,
|
record_diagnostic_result, record_trace, record_trace_bus, save_signed_telemetry_export,
|
||||||
};
|
};
|
||||||
use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, trace_emit};
|
use rustfs_common::trace_bus::{
|
||||||
|
TelemetryTraceEvent, TelemetryTraceOperation, TelemetryTraceStatus, TraceEvent, TraceFunc, TraceKind, subscribe_trace_events,
|
||||||
|
telemetry_trace_emit, telemetry_trace_subscriber_count, trace_emit,
|
||||||
|
};
|
||||||
|
use rustfs_io_metrics::{record_s3_op, s3_http_metrics::S3HttpRequestGuard};
|
||||||
|
use rustfs_s3_ops::S3Operation;
|
||||||
use sha2::{Digest as _, Sha256};
|
use sha2::{Digest as _, Sha256};
|
||||||
use tokio::sync::mpsc;
|
use tokio::sync::mpsc;
|
||||||
use tokio_util::sync::CancellationToken;
|
use tokio_util::sync::CancellationToken;
|
||||||
@@ -24,6 +29,16 @@ fn consent() -> LocalTelemetryConsent {
|
|||||||
LocalTelemetryConsent::new(Instant::now() + Duration::from_secs(5)).expect("future consent")
|
LocalTelemetryConsent::new(Instant::now() + Duration::from_secs(5)).expect("future consent")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn wait_for_telemetry_subscriber() {
|
||||||
|
tokio::time::timeout(Duration::from_secs(1), async {
|
||||||
|
while telemetry_trace_subscriber_count() == 0 {
|
||||||
|
tokio::task::yield_now().await;
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.expect("telemetry subscriber should become ready");
|
||||||
|
}
|
||||||
|
|
||||||
fn artifact_request() -> TelemetryArtifactRequest {
|
fn artifact_request() -> TelemetryArtifactRequest {
|
||||||
let now = SystemTime::now().duration_since(UNIX_EPOCH).expect("current time").as_secs() as i64;
|
let now = SystemTime::now().duration_since(UNIX_EPOCH).expect("current time").as_secs() as i64;
|
||||||
let organization = "organizations/019e3ae0-0000-7000-8000-000000000001";
|
let organization = "organizations/019e3ae0-0000-7000-8000-000000000001";
|
||||||
@@ -242,19 +257,21 @@ fn connect_trace_record_enforces_exact_duration_and_sample_boundaries() {
|
|||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial]
|
#[serial]
|
||||||
async fn connect_trace_record_refuses_to_relabel_the_real_heal_bus_as_frozen_telemetry() {
|
async fn connect_trace_record_uses_classified_s3_and_rpc_events_only() {
|
||||||
|
let _unrelated_subscription = subscribe_trace_events();
|
||||||
let task = tokio::spawn(async {
|
let task = tokio::spawn(async {
|
||||||
record_trace_bus(
|
record_trace_bus(
|
||||||
consent(),
|
consent(),
|
||||||
TraceRecordLimits {
|
TraceRecordLimits {
|
||||||
duration: Duration::from_millis(80),
|
duration: Duration::from_millis(40),
|
||||||
max_spans: 8,
|
max_spans: 8,
|
||||||
},
|
},
|
||||||
&CancellationToken::new(),
|
&CancellationToken::new(),
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
});
|
});
|
||||||
tokio::time::sleep(Duration::from_millis(10)).await;
|
wait_for_telemetry_subscriber().await;
|
||||||
|
|
||||||
assert!(trace_emit(|| {
|
assert!(trace_emit(|| {
|
||||||
TraceEvent::new(TraceKind::Scanner, TraceFunc::ScannerHealCandidate)
|
TraceEvent::new(TraceKind::Scanner, TraceFunc::ScannerHealCandidate)
|
||||||
.with_bucket("SYNTHETIC_SECRET_BUCKET")
|
.with_bucket("SYNTHETIC_SECRET_BUCKET")
|
||||||
@@ -262,11 +279,121 @@ async fn connect_trace_record_refuses_to_relabel_the_real_heal_bus_as_frozen_tel
|
|||||||
.with_duration(Duration::from_micros(41))
|
.with_duration(Duration::from_micros(41))
|
||||||
.with_attr("error", "SYNTHETIC_SECRET_ERROR")
|
.with_attr("error", "SYNTHETIC_SECRET_ERROR")
|
||||||
}));
|
}));
|
||||||
|
|
||||||
|
let mut s3_request = S3HttpRequestGuard::new("GET");
|
||||||
|
s3_request.in_scope(|| record_s3_op(S3Operation::GetObject));
|
||||||
|
tokio::time::sleep(Duration::from_millis(1)).await;
|
||||||
|
s3_request.response(200);
|
||||||
|
assert!(telemetry_trace_emit(|| {
|
||||||
|
TelemetryTraceEvent::new(
|
||||||
|
TelemetryTraceOperation::InternalRpc,
|
||||||
|
Duration::from_micros(41),
|
||||||
|
TelemetryTraceStatus::Error,
|
||||||
|
)
|
||||||
|
}));
|
||||||
|
|
||||||
|
let record = task.await.expect("capture task").expect("typed capture should succeed");
|
||||||
|
assert_eq!(record.completion, TraceRecordCompletion::Complete);
|
||||||
|
assert_eq!(record.data.spans.len(), 2);
|
||||||
|
assert_eq!(record.data.spans[0].operation, TelemetryOperation::GetObject);
|
||||||
|
assert_eq!(record.data.spans[0].status, TelemetrySpanStatus::Ok);
|
||||||
|
assert_eq!(record.data.spans[1].operation, TelemetryOperation::InternalRpc);
|
||||||
|
assert_eq!(record.data.spans[1].duration_micros, 41);
|
||||||
|
assert_eq!(record.data.spans[1].status, TelemetrySpanStatus::Error);
|
||||||
|
|
||||||
|
let json = serde_json::to_string(&record.data).expect("serialize typed telemetry");
|
||||||
|
for forbidden in ["SYNTHETIC_SECRET_BUCKET", "private/object", "SYNTHETIC_SECRET_ERROR"] {
|
||||||
|
assert!(!json.contains(forbidden));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial]
|
||||||
|
async fn connect_trace_record_typed_bus_stops_at_the_span_limit() {
|
||||||
|
let task = tokio::spawn(async {
|
||||||
|
record_trace_bus(
|
||||||
|
consent(),
|
||||||
|
TraceRecordLimits {
|
||||||
|
duration: Duration::from_secs(1),
|
||||||
|
max_spans: 1,
|
||||||
|
},
|
||||||
|
&CancellationToken::new(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
});
|
||||||
|
wait_for_telemetry_subscriber().await;
|
||||||
|
|
||||||
|
for _ in 0..3 {
|
||||||
|
assert!(telemetry_trace_emit(|| {
|
||||||
|
TelemetryTraceEvent::new(TelemetryTraceOperation::InternalRpc, Duration::from_micros(1), TelemetryTraceStatus::Ok)
|
||||||
|
}));
|
||||||
|
}
|
||||||
|
|
||||||
|
let record = task
|
||||||
|
.await
|
||||||
|
.expect("capture task")
|
||||||
|
.expect("bounded typed capture should finish");
|
||||||
|
assert_eq!(record.completion, TraceRecordCompletion::LimitExceeded);
|
||||||
|
assert_eq!(record.data.spans.len(), 1);
|
||||||
|
assert_eq!(record.data.dropped_span_count, 2);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial]
|
||||||
|
async fn connect_trace_record_typed_bus_honors_expiry_and_stop() {
|
||||||
|
let short_consent = LocalTelemetryConsent::new(Instant::now() + Duration::from_millis(10)).expect("short consent");
|
||||||
|
let error = record_trace_bus(
|
||||||
|
short_consent,
|
||||||
|
TraceRecordLimits {
|
||||||
|
duration: Duration::from_secs(1),
|
||||||
|
max_spans: 1,
|
||||||
|
},
|
||||||
|
&CancellationToken::new(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect_err("capture must fit inside the consent window");
|
||||||
|
assert_eq!(error, TelemetryProducerError::ConsentExpired);
|
||||||
|
|
||||||
|
let cancel = CancellationToken::new();
|
||||||
|
let task_cancel = cancel.clone();
|
||||||
|
let task = tokio::spawn(async move {
|
||||||
|
record_trace_bus(
|
||||||
|
consent(),
|
||||||
|
TraceRecordLimits {
|
||||||
|
duration: Duration::from_secs(1),
|
||||||
|
max_spans: 1,
|
||||||
|
},
|
||||||
|
&task_cancel,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
});
|
||||||
|
wait_for_telemetry_subscriber().await;
|
||||||
|
cancel.cancel();
|
||||||
|
|
||||||
let error = task
|
let error = task
|
||||||
.await
|
.await
|
||||||
.expect("capture task")
|
.expect("capture task")
|
||||||
.expect_err("heal/scanner events have no frozen telemetry semantics");
|
.expect_err("cancelled typed capture must fail");
|
||||||
assert_eq!(error, TelemetryProducerError::SourceUnavailable);
|
assert_eq!(error, TelemetryProducerError::Cancelled);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial]
|
||||||
|
async fn connect_trace_record_completes_an_idle_typed_window_without_inventing_spans() {
|
||||||
|
let record = record_trace_bus(
|
||||||
|
consent(),
|
||||||
|
TraceRecordLimits {
|
||||||
|
duration: Duration::from_millis(5),
|
||||||
|
max_spans: 8,
|
||||||
|
},
|
||||||
|
&CancellationToken::new(),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("idle typed source remains available");
|
||||||
|
|
||||||
|
assert_eq!(record.completion, TraceRecordCompletion::Complete);
|
||||||
|
assert!(record.data.spans.is_empty());
|
||||||
|
assert_eq!(record.data.dropped_span_count, 0);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
|
|||||||
Reference in New Issue
Block a user