feat(connect): add typed telemetry trace source (#7775)

This commit is contained in:
Chris
2026-09-14 02:46:22 +08:00
committed by GitHub
parent 44c49733fd
commit dc700d3662
7 changed files with 417 additions and 30 deletions
Generated
+1
View File
@@ -10035,6 +10035,7 @@ dependencies = [
"metrics",
"metrics-util",
"num_cpus",
"rustfs-common",
"rustfs-s3-ops",
"sysinfo",
"thiserror 2.0.20",
+124
View File
@@ -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);
}
}
+1
View File
@@ -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 }
+50
View File
@@ -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();
+79 -23
View File
@@ -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;
}
}
}
+29 -1
View File
@@ -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() {
+133 -6
View File
@@ -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]