Skip to main content

aws_smithy_runtime/client/
metrics.rs

1/*
2 * Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
3 * SPDX-License-Identifier: Apache-2.0
4 */
5
6use aws_smithy_async::time::{SharedTimeSource, TimeSource};
7use aws_smithy_observability::{
8    global::get_telemetry_provider, instruments::Histogram, AttributeValue, Attributes,
9    ObservabilityError,
10};
11use aws_smithy_runtime_api::client::http::telemetry::{
12    CaptureHttpAttemptTelemetry, ConnectionAcquisitionTelemetry, HttpAttemptTelemetry,
13};
14use aws_smithy_runtime_api::client::{
15    interceptors::{
16        context::BeforeTransmitInterceptorContextMut, dyn_dispatch_hint, Intercept,
17        SharedInterceptor,
18    },
19    orchestrator::Metadata,
20    runtime_components::RuntimeComponentsBuilder,
21    runtime_plugin::RuntimePlugin,
22};
23use aws_smithy_types::config_bag::{FrozenLayer, Layer, Storable, StoreReplace};
24use aws_smithy_types::telemetry::{CapturedTelemetryAttributes, RequestedTelemetryAttributes};
25use std::{borrow::Cow, sync::Arc, time::SystemTime};
26
27/// Sets the outcome attributes (`error.type` and `http.status_code`) on `attrs` from a
28/// finalizer-phase context.
29fn add_outcome_attrs(
30    attrs: &mut Attributes,
31    context: &aws_smithy_runtime_api::client::interceptors::context::FinalizerInterceptorContextRef<
32        '_,
33    >,
34) {
35    // Coarse category only; the error is type-erased here, so the modeled name isn't reachable.
36    // Absent on success, per OTel convention.
37    if let Some(Err(err)) = context.output_or_error() {
38        let category = if err.is_timeout_error() {
39            "timeout"
40        } else if err.is_connector_error() {
41            "connector"
42        } else if err.is_response_error() {
43            "response"
44        } else if err.is_operation_error() {
45            "operation"
46        } else {
47            "other"
48        };
49        attrs.set("error.type", AttributeValue::String(category.into()));
50    }
51
52    // Raw HTTP status code, whenever a response reached us.
53    if let Some(response) = context.response() {
54        attrs.set(
55            "http.status_code",
56            AttributeValue::I64(i64::from(response.status().as_u16())),
57        );
58    }
59}
60
61/// Struct to hold metric data in the ConfigBag
62#[derive(Debug, Clone)]
63pub(crate) struct MeasurementsContainer {
64    call_start: SystemTime,
65    attempts: u32,
66    attempt_start: SystemTime,
67    attempt_capture: Option<CaptureHttpAttemptTelemetry>,
68}
69
70impl Storable for MeasurementsContainer {
71    type Storer = StoreReplace<Self>;
72}
73
74/// Instruments for recording a single operation
75#[derive(Debug, Clone)]
76pub(crate) struct OperationTelemetry {
77    pub(crate) operation_duration: Arc<dyn Histogram>,
78    pub(crate) attempt_duration: Arc<dyn Histogram>,
79    pub(crate) connection_acquisition_duration: Arc<dyn Histogram>,
80    // Body sizes are their own instruments rather than attributes on the duration histogram: body
81    // size is near-unique per call, so attaching it as a label would fragment the duration metric
82    // into one time series per byte count.
83    pub(crate) request_body_size: Arc<dyn Histogram>,
84    pub(crate) response_body_size: Arc<dyn Histogram>,
85}
86
87impl OperationTelemetry {
88    pub(crate) fn new(scope: &'static str) -> Result<Self, ObservabilityError> {
89        let meter = get_telemetry_provider()?
90            .meter_provider()
91            .get_meter(scope, None);
92
93        Ok(Self{
94            operation_duration: meter
95                .create_histogram("smithy.client.call.duration")
96                .set_units("s")
97                .set_description("Overall call duration (including retries and time to send or receive request and response body)")
98                .build(),
99            attempt_duration: meter
100                .create_histogram("smithy.client.call.attempt.duration")
101                .set_units("s")
102                .set_description("The time it takes to connect to the service, send the request, and get back HTTP status code and headers (including time queued waiting to be sent)")
103                .build(),
104            connection_acquisition_duration: meter
105                .create_histogram("smithy.client.http.connections.acquire_duration")
106                .set_units("s")
107                .set_description("Time spent acquiring the connection that accepted the request")
108                .build(),
109            request_body_size: meter
110                .create_histogram("smithy.client.call.request.size")
111                .set_units("By")
112                .set_description("Size of the transferred request body, in bytes")
113                .build(),
114            response_body_size: meter
115                .create_histogram("smithy.client.call.response.size")
116                .set_units("By")
117                .set_description("Size of the transferred response body, in bytes")
118                .build(),
119        })
120    }
121
122    fn record_connection_acquisition(&self, telemetry: &HttpAttemptTelemetry) {
123        if let Some(duration) = telemetry
124            .acquisition()
125            .and_then(ConnectionAcquisitionTelemetry::duration)
126        {
127            self.connection_acquisition_duration
128                .record(duration.as_secs_f64(), None, None);
129        }
130    }
131}
132
133impl Storable for OperationTelemetry {
134    type Storer = StoreReplace<Self>;
135}
136
137#[derive(Debug)]
138pub(crate) struct MetricsInterceptor {
139    // Holding a TimeSource here isn't ideal, but RuntimeComponents aren't available in
140    // the read_before_execution hook and that is when we need to start the timer for
141    // the operation.
142    time_source: SharedTimeSource,
143}
144
145impl MetricsInterceptor {
146    pub(crate) fn new(time_source: SharedTimeSource) -> Result<Self, ObservabilityError> {
147        Ok(MetricsInterceptor { time_source })
148    }
149
150    pub(crate) fn get_attrs_from_cfg(
151        &self,
152        cfg: &aws_smithy_types::config_bag::ConfigBag,
153    ) -> Option<Attributes> {
154        let operation_metadata = cfg.load::<Metadata>();
155
156        if let Some(md) = operation_metadata {
157            let mut attributes = Attributes::new();
158            attributes.set("rpc.service", AttributeValue::String(md.service().into()));
159            attributes.set("rpc.method", AttributeValue::String(md.name().into()));
160
161            // Merge captured input members that the customer opted in to *emit*. Capture-only
162            // members are present in the bag for in-process reads but are deliberately excluded
163            // from the metric label set.
164            if let (Some(captured), Some(requested)) = (
165                cfg.load::<CapturedTelemetryAttributes>(),
166                cfg.load::<RequestedTelemetryAttributes>(),
167            ) {
168                for (name, value) in captured.iter() {
169                    if requested.should_emit(name) {
170                        attributes.set(name, AttributeValue::String(value.into()));
171                    }
172                }
173            }
174
175            Some(attributes)
176        } else {
177            None
178        }
179    }
180
181    pub(crate) fn get_measurements_and_instruments<'a>(
182        &self,
183        cfg: &'a aws_smithy_types::config_bag::ConfigBag,
184    ) -> (&'a MeasurementsContainer, &'a OperationTelemetry) {
185        let measurements = cfg
186            .load::<MeasurementsContainer>()
187            .expect("set in `read_before_execution`");
188
189        let instruments = cfg
190            .load::<OperationTelemetry>()
191            .expect("set in RuntimePlugin");
192
193        (measurements, instruments)
194    }
195}
196
197#[dyn_dispatch_hint]
198impl Intercept for MetricsInterceptor {
199    fn name(&self) -> &'static str {
200        "MetricsInterceptor"
201    }
202
203    fn read_before_execution(
204        &self,
205        _context: &aws_smithy_runtime_api::client::interceptors::context::BeforeSerializationInterceptorContextRef<'_>,
206        cfg: &mut aws_smithy_types::config_bag::ConfigBag,
207    ) -> Result<(), aws_smithy_runtime_api::box_error::BoxError> {
208        cfg.interceptor_state().store_put(MeasurementsContainer {
209            call_start: self.time_source.now(),
210            attempts: 0,
211            attempt_start: SystemTime::UNIX_EPOCH,
212            attempt_capture: None,
213        });
214
215        Ok(())
216    }
217
218    fn read_after_execution(
219        &self,
220        context: &aws_smithy_runtime_api::client::interceptors::context::FinalizerInterceptorContextRef<'_>,
221        _runtime_components: &aws_smithy_runtime_api::client::runtime_components::RuntimeComponents,
222        cfg: &mut aws_smithy_types::config_bag::ConfigBag,
223    ) -> Result<(), aws_smithy_runtime_api::box_error::BoxError> {
224        let (measurements, instruments) = self.get_measurements_and_instruments(cfg);
225
226        let attributes = self.get_attrs_from_cfg(cfg);
227
228        if let Some(mut attrs) = attributes {
229            // The outcome is only known at the finalizer, so it is set here rather than in
230            // `get_attrs_from_cfg` (which also serves the per-attempt path).
231            add_outcome_attrs(&mut attrs, context);
232
233            // Transferred byte sizes are recorded on their own instruments by the byte
234            // interceptor (see `telemetry_bytes`), not as attributes on the duration histogram.
235
236            let call_end = self.time_source.now();
237            let call_duration = call_end.duration_since(measurements.call_start);
238            if let Ok(elapsed) = call_duration {
239                instruments
240                    .operation_duration
241                    .record(elapsed.as_secs_f64(), Some(&attrs), None);
242            }
243        }
244
245        Ok(())
246    }
247
248    fn read_before_attempt(
249        &self,
250        _context: &aws_smithy_runtime_api::client::interceptors::context::BeforeTransmitInterceptorContextRef<'_>,
251        _runtime_components: &aws_smithy_runtime_api::client::runtime_components::RuntimeComponents,
252        cfg: &mut aws_smithy_types::config_bag::ConfigBag,
253    ) -> Result<(), aws_smithy_runtime_api::box_error::BoxError> {
254        let measurements = cfg
255            .get_mut::<MeasurementsContainer>()
256            .expect("set in `read_before_execution`");
257
258        measurements.attempts += 1;
259        measurements.attempt_start = self.time_source.now();
260        measurements.attempt_capture = None;
261
262        Ok(())
263    }
264
265    fn modify_before_transmit(
266        &self,
267        context: &mut BeforeTransmitInterceptorContextMut<'_>,
268        _runtime_components: &aws_smithy_runtime_api::client::runtime_components::RuntimeComponents,
269        _cfg: &mut aws_smithy_types::config_bag::ConfigBag,
270    ) -> Result<(), aws_smithy_runtime_api::box_error::BoxError> {
271        if context
272            .request()
273            .extension::<CaptureHttpAttemptTelemetry>()
274            .is_none()
275        {
276            context
277                .request_mut()
278                .add_extension(CaptureHttpAttemptTelemetry::new());
279        }
280        Ok(())
281    }
282
283    fn read_before_transmit(
284        &self,
285        context: &aws_smithy_runtime_api::client::interceptors::context::BeforeTransmitInterceptorContextRef<'_>,
286        _runtime_components: &aws_smithy_runtime_api::client::runtime_components::RuntimeComponents,
287        cfg: &mut aws_smithy_types::config_bag::ConfigBag,
288    ) -> Result<(), aws_smithy_runtime_api::box_error::BoxError> {
289        let capture = context
290            .request()
291            .extension::<CaptureHttpAttemptTelemetry>()
292            .cloned();
293        cfg.get_mut::<MeasurementsContainer>()
294            .expect("set in `read_before_execution`")
295            .attempt_capture = capture;
296        Ok(())
297    }
298
299    fn read_after_attempt(
300        &self,
301        _context: &aws_smithy_runtime_api::client::interceptors::context::FinalizerInterceptorContextRef<'_>,
302        _runtime_components: &aws_smithy_runtime_api::client::runtime_components::RuntimeComponents,
303        cfg: &mut aws_smithy_types::config_bag::ConfigBag,
304    ) -> Result<(), aws_smithy_runtime_api::box_error::BoxError> {
305        let (attempts, attempt_start, http_telemetry) = {
306            let measurements = cfg
307                .get_mut::<MeasurementsContainer>()
308                .expect("set in `read_before_execution`");
309            (
310                measurements.attempts,
311                measurements.attempt_start,
312                measurements
313                    .attempt_capture
314                    .take()
315                    .map(|capture| capture.get()),
316            )
317        };
318        let instruments = cfg
319            .load::<OperationTelemetry>()
320            .expect("set in RuntimePlugin");
321
322        let attempt_end = self.time_source.now();
323        let attempt_duration = attempt_end.duration_since(attempt_start);
324        let attributes = self.get_attrs_from_cfg(cfg);
325
326        if let Some(http_telemetry) = http_telemetry {
327            instruments.record_connection_acquisition(&http_telemetry);
328        }
329
330        if let (Ok(elapsed), Some(mut attrs)) = (attempt_duration, attributes) {
331            attrs.set("attempt", AttributeValue::I64(attempts.into()));
332
333            instruments
334                .attempt_duration
335                .record(elapsed.as_secs_f64(), Some(&attrs), None);
336        }
337        Ok(())
338    }
339}
340
341/// Runtime plugin that adds an interceptor for collecting metrics
342#[derive(Debug, Default)]
343pub struct MetricsRuntimePlugin {
344    scope: &'static str,
345    time_source: SharedTimeSource,
346    metadata: Option<Metadata>,
347}
348
349impl MetricsRuntimePlugin {
350    /// Create a [MetricsRuntimePluginBuilder]
351    pub fn builder() -> MetricsRuntimePluginBuilder {
352        MetricsRuntimePluginBuilder::default()
353    }
354}
355
356impl RuntimePlugin for MetricsRuntimePlugin {
357    fn runtime_components(
358        &self,
359        _current_components: &RuntimeComponentsBuilder,
360    ) -> Cow<'_, RuntimeComponentsBuilder> {
361        let interceptor = MetricsInterceptor::new(self.time_source.clone());
362        if let Ok(interceptor) = interceptor {
363            Cow::Owned(
364                RuntimeComponentsBuilder::new("Metrics")
365                    .with_interceptor(SharedInterceptor::permanent(interceptor))
366                    // Counts transferred bytes into the bag for the metrics interceptor to read.
367                    .with_interceptor(SharedInterceptor::permanent(
368                        crate::client::telemetry_bytes::TelemetryBytesInterceptor,
369                    )),
370            )
371        } else {
372            Cow::Owned(RuntimeComponentsBuilder::new("Metrics"))
373        }
374    }
375
376    fn config(&self) -> Option<FrozenLayer> {
377        let instruments = OperationTelemetry::new(self.scope);
378
379        if let Ok(instruments) = instruments {
380            let mut cfg = Layer::new("Metrics");
381            cfg.store_put(instruments);
382
383            if let Some(metadata) = &self.metadata {
384                cfg.store_put(metadata.clone());
385            }
386
387            Some(cfg.freeze())
388        } else {
389            None
390        }
391    }
392}
393
394/// Builder for [MetricsRuntimePlugin]
395#[derive(Debug, Default)]
396pub struct MetricsRuntimePluginBuilder {
397    scope: Option<&'static str>,
398    time_source: Option<SharedTimeSource>,
399    metadata: Option<Metadata>,
400}
401
402impl MetricsRuntimePluginBuilder {
403    /// Set the scope for the metrics
404    pub fn with_scope(mut self, scope: &'static str) -> Self {
405        self.scope = Some(scope);
406        self
407    }
408
409    /// Set the [TimeSource] for the metrics
410    pub fn with_time_source(mut self, time_source: impl TimeSource + 'static) -> Self {
411        self.time_source = Some(SharedTimeSource::new(time_source));
412        self
413    }
414
415    /// Set the [Metadata] for the metrics.
416    ///
417    /// Note: the Metadata is optional, most operations set it themselves, but this is useful
418    /// for operations that do not, like some of the credential providers.
419    pub fn with_metadata(mut self, metadata: Metadata) -> Self {
420        self.metadata = Some(metadata);
421        self
422    }
423
424    /// Build a [MetricsRuntimePlugin]
425    pub fn build(
426        self,
427    ) -> Result<MetricsRuntimePlugin, aws_smithy_runtime_api::box_error::BoxError> {
428        if let Some(scope) = self.scope {
429            Ok(MetricsRuntimePlugin {
430                scope,
431                time_source: self.time_source.unwrap_or_default(),
432                metadata: self.metadata,
433            })
434        } else {
435            Err("Scope is required for MetricsRuntimePlugin.".into())
436        }
437    }
438}
439
440#[cfg(test)]
441mod test {
442    use super::*;
443    use aws_smithy_async::time::SystemTimeSource;
444    use aws_smithy_runtime_api::client::connection::ConnectionMetadata;
445    use aws_smithy_runtime_api::client::http::telemetry::{
446        ConnectionAcquisitionTelemetry, ConnectionUsage,
447    };
448    use aws_smithy_runtime_api::client::interceptors::context::{Input, InterceptorContext};
449    use aws_smithy_runtime_api::client::orchestrator::HttpRequest;
450    use aws_smithy_types::config_bag::ConfigBag;
451    use std::sync::Mutex;
452
453    #[derive(Debug, Default)]
454    struct RecordingHistogram {
455        records: Mutex<Vec<(f64, Option<Attributes>)>>,
456    }
457
458    impl RecordingHistogram {
459        fn only_record(&self) -> (f64, Option<Attributes>) {
460            let records = self.records.lock().unwrap();
461            assert_eq!(records.len(), 1);
462            records[0].clone()
463        }
464
465        fn record_count(&self) -> usize {
466            self.records.lock().unwrap().len()
467        }
468    }
469
470    impl Histogram for RecordingHistogram {
471        fn record(
472            &self,
473            value: f64,
474            attributes: Option<&Attributes>,
475            _context: Option<&dyn aws_smithy_observability::Context>,
476        ) {
477            self.records
478                .lock()
479                .unwrap()
480                .push((value, attributes.cloned()));
481        }
482    }
483
484    fn interceptor() -> MetricsInterceptor {
485        MetricsInterceptor::new(SharedTimeSource::new(SystemTimeSource::new())).unwrap()
486    }
487
488    fn cfg_with(layer: Layer) -> ConfigBag {
489        ConfigBag::of_layers(vec![layer])
490    }
491
492    fn string_attr<'a>(attrs: &'a Attributes, key: &str) -> Option<&'a str> {
493        match attrs.get(key) {
494            Some(AttributeValue::String(s)) => Some(s.as_str()),
495            _ => None,
496        }
497    }
498
499    #[test]
500    fn base_attrs_are_service_and_method() {
501        let mut layer = Layer::new("test");
502        layer.store_put(Metadata::new("GetObject", "S3"));
503
504        let attrs = interceptor()
505            .get_attrs_from_cfg(&cfg_with(layer))
506            .expect("metadata present");
507
508        assert_eq!(Some("S3"), string_attr(&attrs, "rpc.service"));
509        assert_eq!(Some("GetObject"), string_attr(&attrs, "rpc.method"));
510    }
511
512    #[test]
513    fn no_attrs_without_metadata() {
514        // Nothing to key the metric on, so no attributes are produced.
515        assert!(interceptor()
516            .get_attrs_from_cfg(&cfg_with(Layer::new("test")))
517            .is_none());
518    }
519
520    fn test_instruments(acquisition: Arc<RecordingHistogram>) -> OperationTelemetry {
521        let ignored: Arc<dyn Histogram> = Arc::new(RecordingHistogram::default());
522        OperationTelemetry {
523            operation_duration: ignored.clone(),
524            attempt_duration: ignored.clone(),
525            connection_acquisition_duration: acquisition,
526            request_body_size: ignored.clone(),
527            response_body_size: ignored,
528        }
529    }
530
531    fn install_measurements(cfg: &mut ConfigBag) {
532        cfg.interceptor_state().store_put(MeasurementsContainer {
533            call_start: SystemTime::UNIX_EPOCH,
534            attempts: 0,
535            attempt_start: SystemTime::UNIX_EPOCH,
536            attempt_capture: None,
537        });
538    }
539
540    #[test]
541    fn before_transmit_adopts_an_existing_attempt_capture() {
542        let runtime_components = RuntimeComponentsBuilder::for_tests().build().unwrap();
543        let mut context = InterceptorContext::new(Input::doesnt_matter());
544        context.set_request(HttpRequest::empty());
545        let mut cfg = cfg_with(Layer::new("test"));
546        install_measurements(&mut cfg);
547        let caller_capture = CaptureHttpAttemptTelemetry::new();
548        context
549            .request_mut()
550            .unwrap()
551            .add_extension(caller_capture.clone());
552
553        interceptor()
554            .modify_before_transmit(&mut (&mut context).into(), &runtime_components, &mut cfg)
555            .unwrap();
556        interceptor()
557            .read_before_transmit(&(&context).into(), &runtime_components, &mut cfg)
558            .unwrap();
559
560        let request_capture = context
561            .request()
562            .unwrap()
563            .extension::<CaptureHttpAttemptTelemetry>()
564            .expect("request capture");
565        request_capture.record_connector_call_duration(std::time::Duration::from_millis(4));
566        assert_eq!(
567            cfg.load::<MeasurementsContainer>()
568                .and_then(|measurements| measurements.attempt_capture.as_ref())
569                .expect("attempt capture")
570                .get()
571                .connector_call_duration(),
572            Some(std::time::Duration::from_millis(4))
573        );
574        assert_eq!(
575            caller_capture.get().connector_call_duration(),
576            Some(std::time::Duration::from_millis(4))
577        );
578    }
579
580    #[test]
581    fn before_transmit_adopts_the_final_replacement_capture() {
582        let runtime_components = RuntimeComponentsBuilder::for_tests().build().unwrap();
583        let mut context = InterceptorContext::new(Input::doesnt_matter());
584        context.set_request(HttpRequest::empty());
585        let mut cfg = cfg_with(Layer::new("test"));
586        install_measurements(&mut cfg);
587
588        interceptor()
589            .modify_before_transmit(&mut (&mut context).into(), &runtime_components, &mut cfg)
590            .unwrap();
591        let replacement = CaptureHttpAttemptTelemetry::new();
592        context
593            .request_mut()
594            .unwrap()
595            .add_extension(replacement.clone());
596        interceptor()
597            .read_before_transmit(&(&context).into(), &runtime_components, &mut cfg)
598            .unwrap();
599
600        replacement.record_connector_call_duration(std::time::Duration::from_millis(6));
601        assert_eq!(
602            cfg.load::<MeasurementsContainer>()
603                .and_then(|measurements| measurements.attempt_capture.as_ref())
604                .expect("attempt capture")
605                .get()
606                .connector_call_duration(),
607            Some(std::time::Duration::from_millis(6))
608        );
609    }
610
611    #[test]
612    fn attempt_capture_is_consumed_once_and_cleared_before_retry() {
613        let runtime_components = RuntimeComponentsBuilder::for_tests().build().unwrap();
614        let acquisition = Arc::new(RecordingHistogram::default());
615        let mut layer = Layer::new("test");
616        layer.store_put(Metadata::new("GetObject", "S3"));
617        layer.store_put(test_instruments(acquisition.clone()));
618        let mut cfg = cfg_with(layer);
619        install_measurements(&mut cfg);
620        let mut context = InterceptorContext::new(Input::doesnt_matter());
621        context.set_request(HttpRequest::empty());
622        let metrics = interceptor();
623
624        metrics
625            .read_before_attempt(&(&context).into(), &runtime_components, &mut cfg)
626            .unwrap();
627        metrics
628            .modify_before_transmit(&mut (&mut context).into(), &runtime_components, &mut cfg)
629            .unwrap();
630        metrics
631            .read_before_transmit(&(&context).into(), &runtime_components, &mut cfg)
632            .unwrap();
633        context
634            .request()
635            .unwrap()
636            .extension::<CaptureHttpAttemptTelemetry>()
637            .expect("request capture")
638            .record_connection_selection(
639                ConnectionAcquisitionTelemetry::new(
640                    std::time::Duration::from_millis(3),
641                    ConnectionUsage::Fresh,
642                ),
643                ConnectionMetadata::builder()
644                    .proxied(false)
645                    .poison_fn(|| {})
646                    .build(),
647            );
648        metrics
649            .read_after_attempt(&(&context).into(), &runtime_components, &mut cfg)
650            .unwrap();
651        assert_eq!(acquisition.record_count(), 1);
652
653        metrics
654            .read_before_attempt(&(&context).into(), &runtime_components, &mut cfg)
655            .unwrap();
656        metrics
657            .read_after_attempt(&(&context).into(), &runtime_components, &mut cfg)
658            .unwrap();
659        assert_eq!(
660            acquisition.record_count(),
661            1,
662            "a pre-transmit retry halt must not repeat the prior acquisition"
663        );
664    }
665
666    #[test]
667    fn records_connection_acquisition_without_attributes() {
668        let acquisition = Arc::new(RecordingHistogram::default());
669        let instruments = test_instruments(acquisition.clone());
670        let capture = CaptureHttpAttemptTelemetry::new();
671        capture.record_connector_call_duration(std::time::Duration::from_millis(7));
672        capture.record_connection_selection(
673            ConnectionAcquisitionTelemetry::new(
674                std::time::Duration::from_millis(3),
675                ConnectionUsage::Reused,
676            ),
677            ConnectionMetadata::builder()
678                .proxied(false)
679                .poison_fn(|| {})
680                .build(),
681        );
682        instruments.record_connection_acquisition(&capture.get());
683
684        let (acquisition_value, acquisition_attrs) = acquisition.only_record();
685        assert_eq!(
686            acquisition_value,
687            std::time::Duration::from_millis(3).as_secs_f64()
688        );
689        assert!(acquisition_attrs.is_none());
690    }
691
692    #[test]
693    fn emitted_members_are_merged_onto_attrs() {
694        let mut captured = CapturedTelemetryAttributes::new();
695        captured.insert("Bucket", "my-bucket");
696
697        let mut layer = Layer::new("test");
698        layer.store_put(Metadata::new("GetObject", "S3"));
699        layer.store_put(captured);
700        layer.store_put(RequestedTelemetryAttributes::new(["Bucket"]));
701
702        let attrs = interceptor()
703            .get_attrs_from_cfg(&cfg_with(layer))
704            .expect("metadata present");
705
706        // The emitted input member rides alongside the built-in rpc.* attributes.
707        assert_eq!(Some("my-bucket"), string_attr(&attrs, "Bucket"));
708        assert_eq!(Some("S3"), string_attr(&attrs, "rpc.service"));
709    }
710
711    #[test]
712    fn capture_only_members_are_not_emitted() {
713        // A value captured for in-process reads must not land on the metric.
714        let mut captured = CapturedTelemetryAttributes::new();
715        captured.insert("Prefix", "logs/");
716
717        let mut requested = RequestedTelemetryAttributes::default();
718        requested.capture_only(["Prefix"]);
719
720        let mut layer = Layer::new("test");
721        layer.store_put(Metadata::new("GetObject", "S3"));
722        layer.store_put(captured);
723        layer.store_put(requested);
724
725        let attrs = interceptor()
726            .get_attrs_from_cfg(&cfg_with(layer))
727            .expect("metadata present");
728
729        assert!(
730            attrs.get("Prefix").is_none(),
731            "capture-only member must not be emitted on the metric"
732        );
733    }
734
735    #[test]
736    fn nothing_captured_leaves_only_base_attrs() {
737        // Opt-in is off by default: an empty capture set adds nothing.
738        let mut layer = Layer::new("test");
739        layer.store_put(Metadata::new("GetObject", "S3"));
740        layer.store_put(CapturedTelemetryAttributes::new());
741
742        let attrs = interceptor()
743            .get_attrs_from_cfg(&cfg_with(layer))
744            .expect("metadata present");
745
746        assert_eq!(Some("GetObject"), string_attr(&attrs, "rpc.method"));
747        assert!(attrs.get("Bucket").is_none());
748    }
749
750    // --- add_outcome_attrs (the `status` dimension) ---
751
752    use aws_smithy_runtime_api::client::interceptors::context::{Error, Output};
753    use aws_smithy_runtime_api::client::orchestrator::OrchestratorError;
754    use aws_smithy_runtime_api::client::result::ConnectorError;
755    use aws_smithy_runtime_api::http::{Response, StatusCode};
756    use aws_smithy_types::body::SdkBody;
757
758    fn i64_attr(attrs: &Attributes, key: &str) -> Option<i64> {
759        match attrs.get(key) {
760            Some(AttributeValue::I64(v)) => Some(*v),
761            _ => None,
762        }
763    }
764
765    #[test]
766    fn outcome_on_success_has_status_code_and_no_error_type() {
767        let mut ctx = InterceptorContext::new(Input::doesnt_matter());
768        ctx.set_output_or_error(Ok(Output::doesnt_matter()));
769        ctx.set_response(Response::new(
770            StatusCode::try_from(200).unwrap(),
771            SdkBody::empty(),
772        ));
773
774        let mut attrs = Attributes::new();
775        add_outcome_attrs(&mut attrs, &(&ctx).into());
776
777        // error.type is absent on success (OTel convention); status code is present.
778        assert!(attrs.get("error.type").is_none());
779        assert_eq!(Some(200), i64_attr(&attrs, "http.status_code"));
780    }
781
782    #[test]
783    fn outcome_on_failure_sets_error_type_category() {
784        let mut ctx = InterceptorContext::new(Input::doesnt_matter());
785        ctx.set_output_or_error(Err(OrchestratorError::connector(ConnectorError::io(
786            "boom".into(),
787        ))));
788
789        let mut attrs = Attributes::new();
790        add_outcome_attrs(&mut attrs, &(&ctx).into());
791
792        // A connector error maps to the `connector` category.
793        assert_eq!(Some("connector"), string_attr(&attrs, "error.type"));
794    }
795
796    #[test]
797    fn outcome_without_response_omits_status_code() {
798        let mut ctx: InterceptorContext<Input, Output, Error> =
799            InterceptorContext::new(Input::doesnt_matter());
800        ctx.set_output_or_error(Err(OrchestratorError::connector(ConnectorError::io(
801            "boom".into(),
802        ))));
803
804        let mut attrs = Attributes::new();
805        add_outcome_attrs(&mut attrs, &(&ctx).into());
806
807        // No response reached us, so there is no HTTP status code to record.
808        assert!(attrs.get("http.status_code").is_none());
809        assert_eq!(Some("connector"), string_attr(&attrs, "error.type"));
810    }
811}