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-util",
|
||||
"num_cpus",
|
||||
"rustfs-common",
|
||||
"rustfs-s3-ops",
|
||||
"sysinfo",
|
||||
"thiserror 2.0.20",
|
||||
|
||||
@@ -127,6 +127,44 @@ pub struct TraceEvent {
|
||||
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 {
|
||||
pub fn new(kind: TraceKind, func: TraceFunc) -> Self {
|
||||
Self {
|
||||
@@ -174,15 +212,20 @@ impl TraceEvent {
|
||||
pub struct TraceBus {
|
||||
sender: broadcast::Sender<Arc<TraceEvent>>,
|
||||
subscriber_count: Arc<AtomicUsize>,
|
||||
telemetry_sender: broadcast::Sender<Arc<TelemetryTraceEvent>>,
|
||||
telemetry_subscriber_count: Arc<AtomicUsize>,
|
||||
}
|
||||
|
||||
impl TraceBus {
|
||||
pub fn new(capacity: usize) -> Self {
|
||||
let capacity = capacity.max(1);
|
||||
let (sender, _receiver) = broadcast::channel(capacity);
|
||||
let (telemetry_sender, _telemetry_receiver) = broadcast::channel(capacity);
|
||||
Self {
|
||||
sender,
|
||||
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()
|
||||
}
|
||||
|
||||
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 {
|
||||
@@ -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 {
|
||||
GLOBAL_TRACE_BUS.get_or_init(TraceBus::default)
|
||||
}
|
||||
@@ -252,6 +338,18 @@ pub fn trace_subscriber_count() -> usize {
|
||||
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)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
@@ -330,4 +428,30 @@ mod tests {
|
||||
.expect_err("receiver should observe lag instead of blocking publishers");
|
||||
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]
|
||||
hotpath.workspace = true
|
||||
metrics = { workspace = true }
|
||||
rustfs-common = { workspace = true }
|
||||
rustfs-s3-ops = { workspace = true }
|
||||
num_cpus = { workspace = true }
|
||||
thiserror = { workspace = true }
|
||||
|
||||
@@ -16,10 +16,14 @@
|
||||
//! Admin snapshots and metric exporters share these counters. The older
|
||||
//! 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 std::cell::Cell;
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::sync::{LazyLock, OnceLock};
|
||||
use std::time::Instant;
|
||||
|
||||
const METRIC: &str = "rustfs_s3_http_requests_total";
|
||||
const METHODS: [&str; 10] = [
|
||||
@@ -110,6 +114,7 @@ pub(crate) fn observe_s3_http_operation(op: S3Operation) {
|
||||
pub struct S3HttpRequestGuard {
|
||||
method: usize,
|
||||
operation: usize,
|
||||
telemetry_started_at: Option<Instant>,
|
||||
finished: bool,
|
||||
}
|
||||
|
||||
@@ -122,6 +127,7 @@ impl S3HttpRequestGuard {
|
||||
Self {
|
||||
method: METHODS.iter().position(|known| *known == method).unwrap_or(METHODS.len() - 1),
|
||||
operation: UNKNOWN_OPERATION,
|
||||
telemetry_started_at: (telemetry_trace_subscriber_count() != 0).then(Instant::now),
|
||||
finished: false,
|
||||
}
|
||||
}
|
||||
@@ -151,11 +157,29 @@ impl S3HttpRequestGuard {
|
||||
fn finish(&mut self, outcome: usize) {
|
||||
if !self.finished {
|
||||
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;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
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 {
|
||||
fn drop(&mut self) {
|
||||
self.finish(7);
|
||||
@@ -172,6 +196,32 @@ mod tests {
|
||||
use metrics::with_local_recorder;
|
||||
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]
|
||||
fn outcome_counters_distinguish_partial_and_complete_write_failure() {
|
||||
let counters = HttpOutcomeCounters::new();
|
||||
|
||||
@@ -27,7 +27,9 @@ use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
|
||||
use base64_simd::URL_SAFE_NO_PAD;
|
||||
use p256::ecdsa::{Signature, SigningKey, signature::Signer 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 sha2::{Digest as _, Sha256};
|
||||
use thiserror::Error;
|
||||
@@ -163,6 +165,14 @@ impl LocalTelemetryConsent {
|
||||
.filter(|remaining| !remaining.is_zero())
|
||||
.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)]
|
||||
@@ -563,12 +573,8 @@ pub async fn record_trace(
|
||||
if cancel.is_cancelled() {
|
||||
return Err(TelemetryProducerError::Cancelled);
|
||||
}
|
||||
let remaining = consent.remaining()?;
|
||||
if remaining < limits.duration {
|
||||
return Err(TelemetryProducerError::ConsentExpired);
|
||||
}
|
||||
let deadline = consent.capture_deadline(limits.duration)?;
|
||||
let _lease = acquire_telemetry_lease()?;
|
||||
let deadline = Instant::now() + limits.duration;
|
||||
let mut spans = Vec::with_capacity(limits.max_spans.min(64));
|
||||
let mut dropped_span_count = 0u64;
|
||||
let mut completion = TraceRecordCompletion::Complete;
|
||||
@@ -610,13 +616,7 @@ pub async fn record_trace(
|
||||
})
|
||||
}
|
||||
|
||||
/// Capture the process-local RustFS trace bus.
|
||||
///
|
||||
/// 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.
|
||||
/// Capture the process-local, pre-classified S3 and internode RPC trace source.
|
||||
pub async fn record_trace_bus(
|
||||
consent: LocalTelemetryConsent,
|
||||
limits: TraceRecordLimits,
|
||||
@@ -626,22 +626,78 @@ pub async fn record_trace_bus(
|
||||
if cancel.is_cancelled() {
|
||||
return Err(TelemetryProducerError::Cancelled);
|
||||
}
|
||||
let remaining = consent.remaining()?;
|
||||
if remaining < limits.duration {
|
||||
return Err(TelemetryProducerError::ConsentExpired);
|
||||
}
|
||||
let deadline = consent.capture_deadline(limits.duration)?;
|
||||
let _lease = acquire_telemetry_lease()?;
|
||||
let deadline = Instant::now() + limits.duration;
|
||||
let mut subscription = subscribe_trace_events();
|
||||
let mut subscription = subscribe_telemetry_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 {
|
||||
tokio::select! {
|
||||
biased;
|
||||
_ = 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 {
|
||||
Ok(_event) => {}
|
||||
Err(tokio::sync::broadcast::error::RecvError::Lagged(_dropped)) => {}
|
||||
Err(tokio::sync::broadcast::error::RecvError::Closed) => return Err(TelemetryProducerError::SourceUnavailable),
|
||||
Ok(event) => {
|
||||
if spans.len() == limits.max_spans {
|
||||
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_body_util::{BodyExt, Limited};
|
||||
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_io_metrics::internode_metrics::{
|
||||
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();
|
||||
if let Err(response) = verify_internode_rpc_signature(req.uri(), req.method(), req.headers()) {
|
||||
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();
|
||||
@@ -461,10 +464,20 @@ async fn handle_internode_rpc(req: Request<Incoming>) -> Response<Body> {
|
||||
started_at.elapsed(),
|
||||
);
|
||||
}
|
||||
emit_internode_rpc_telemetry(started_at, response.status());
|
||||
|
||||
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> {
|
||||
match path {
|
||||
READ_FILE_STREAM_PATH => Some(INTERNODE_OPERATION_READ_FILE_STREAM),
|
||||
@@ -1761,6 +1774,21 @@ mod tests {
|
||||
(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]
|
||||
#[serial_test::serial]
|
||||
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,
|
||||
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 tokio::sync::mpsc;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
@@ -24,6 +29,16 @@ fn consent() -> LocalTelemetryConsent {
|
||||
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 {
|
||||
let now = SystemTime::now().duration_since(UNIX_EPOCH).expect("current time").as_secs() as i64;
|
||||
let organization = "organizations/019e3ae0-0000-7000-8000-000000000001";
|
||||
@@ -242,19 +257,21 @@ fn connect_trace_record_enforces_exact_duration_and_sample_boundaries() {
|
||||
|
||||
#[tokio::test]
|
||||
#[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 {
|
||||
record_trace_bus(
|
||||
consent(),
|
||||
TraceRecordLimits {
|
||||
duration: Duration::from_millis(80),
|
||||
duration: Duration::from_millis(40),
|
||||
max_spans: 8,
|
||||
},
|
||||
&CancellationToken::new(),
|
||||
)
|
||||
.await
|
||||
});
|
||||
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||
wait_for_telemetry_subscriber().await;
|
||||
|
||||
assert!(trace_emit(|| {
|
||||
TraceEvent::new(TraceKind::Scanner, TraceFunc::ScannerHealCandidate)
|
||||
.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_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
|
||||
.await
|
||||
.expect("capture task")
|
||||
.expect_err("heal/scanner events have no frozen telemetry semantics");
|
||||
assert_eq!(error, TelemetryProducerError::SourceUnavailable);
|
||||
.expect_err("cancelled typed capture must fail");
|
||||
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]
|
||||
|
||||
Reference in New Issue
Block a user