1use 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
27fn add_outcome_attrs(
30 attrs: &mut Attributes,
31 context: &aws_smithy_runtime_api::client::interceptors::context::FinalizerInterceptorContextRef<
32 '_,
33 >,
34) {
35 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 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#[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#[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 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 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 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 add_outcome_attrs(&mut attrs, context);
232
233 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#[derive(Debug, Default)]
343pub struct MetricsRuntimePlugin {
344 scope: &'static str,
345 time_source: SharedTimeSource,
346 metadata: Option<Metadata>,
347}
348
349impl MetricsRuntimePlugin {
350 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 .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#[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 pub fn with_scope(mut self, scope: &'static str) -> Self {
405 self.scope = Some(scope);
406 self
407 }
408
409 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 pub fn with_metadata(mut self, metadata: Metadata) -> Self {
420 self.metadata = Some(metadata);
421 self
422 }
423
424 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 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 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 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 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 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 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 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 assert!(attrs.get("http.status_code").is_none());
809 assert_eq!(Some("connector"), string_attr(&attrs, "error.type"));
810 }
811}