diff --git a/Cargo.lock b/Cargo.lock index 2ff0d57dc..676d0c149 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -10035,6 +10035,7 @@ dependencies = [ "metrics", "metrics-util", "num_cpus", + "rustfs-common", "rustfs-s3-ops", "sysinfo", "thiserror 2.0.20", diff --git a/crates/common/src/trace_bus.rs b/crates/common/src/trace_bus.rs index 5e6a90d60..9239fb5cd 100644 --- a/crates/common/src/trace_bus.rs +++ b/crates/common/src/trace_bus.rs @@ -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>, subscriber_count: Arc, + telemetry_sender: broadcast::Sender>, + telemetry_subscriber_count: Arc, } 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>, + subscriber_count: Arc, +} + +impl TelemetryTraceSubscription { + pub async fn recv(&mut self) -> Result, broadcast::error::RecvError> { + self.receiver.recv().await + } + + pub fn try_recv(&mut self) -> Result, 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); + } } diff --git a/crates/io-metrics/Cargo.toml b/crates/io-metrics/Cargo.toml index 0591aa4e7..06c95e3e5 100644 --- a/crates/io-metrics/Cargo.toml +++ b/crates/io-metrics/Cargo.toml @@ -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 } diff --git a/crates/io-metrics/src/s3_http_metrics.rs b/crates/io-metrics/src/s3_http_metrics.rs index 283d01a2e..f16d3682d 100644 --- a/crates/io-metrics/src/s3_http_metrics.rs +++ b/crates/io-metrics/src/s3_http_metrics.rs @@ -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, 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 { + 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(); diff --git a/rustfs/src/connect/diagnostics/trace_record.rs b/rustfs/src/connect/diagnostics/trace_record.rs index eeb382b9f..523414609 100644 --- a/rustfs/src/connect/diagnostics/trace_record.rs +++ b/rustfs/src/connect/diagnostics/trace_record.rs @@ -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 { + 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; } } } diff --git a/rustfs/src/storage/rpc/http_service.rs b/rustfs/src/storage/rpc/http_service.rs index 41b1719d7..fa902bbfd 100644 --- a/rustfs/src/storage/rpc/http_service.rs +++ b/rustfs/src/storage/rpc/http_service.rs @@ -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) -> Response { 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) -> Response { 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() { diff --git a/rustfs/tests/connect_trace_record.rs b/rustfs/tests/connect_trace_record.rs index 85bc5553b..0b0db58b3 100644 --- a/rustfs/tests/connect_trace_record.rs +++ b/rustfs/tests/connect_trace_record.rs @@ -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]