diff --git a/changelog.d/20356_aws_component_labels.fix.md b/changelog.d/20356_aws_component_labels.fix.md new file mode 100644 index 0000000000000..be9568a38124f --- /dev/null +++ b/changelog.d/20356_aws_component_labels.fix.md @@ -0,0 +1,3 @@ +AWS sinks (`aws_cloudwatch_logs`, `aws_s3`, `aws_kinesis_firehose`, `aws_kinesis_streams`, `aws_sns`, `aws_sqs`) now emit `component_sent_bytes_total` through the Driver instead of the transport layer, ensuring `component_id`, `component_kind`, `component_type`, and `region` labels are always present. + +authors: clee2691 diff --git a/lib/vector-common/src/internal_event/bytes_sent.rs b/lib/vector-common/src/internal_event/bytes_sent.rs index b363893783f43..6353cf7f025c6 100644 --- a/lib/vector-common/src/internal_event/bytes_sent.rs +++ b/lib/vector-common/src/internal_event/bytes_sent.rs @@ -8,8 +8,15 @@ use super::{ByteSize, CounterName, Protocol, SharedString}; crate::registered_event!( BytesSent { protocol: SharedString, + extra_labels: Vec<(SharedString, SharedString)>, } => { - bytes_sent: Counter = counter!(CounterName::ComponentSentBytesTotal, "protocol" => self.protocol.clone()), + bytes_sent: Counter = { + let mut labels: Vec<(String, String)> = vec![("protocol".to_string(), self.protocol.to_string())]; + for (k, v) in &self.extra_labels { + labels.push((k.to_string(), v.to_string())); + } + counter!(CounterName::ComponentSentBytesTotal, &labels) + }, protocol: SharedString = self.protocol, } @@ -23,6 +30,7 @@ impl From for BytesSent { fn from(protocol: Protocol) -> Self { Self { protocol: protocol.0, + extra_labels: vec![], } } } diff --git a/lib/vector-stream/src/driver.rs b/lib/vector-stream/src/driver.rs index b391eb90677d6..7d829eac7aa09 100644 --- a/lib/vector-stream/src/driver.rs +++ b/lib/vector-stream/src/driver.rs @@ -44,6 +44,7 @@ pub struct Driver { input: St, service: Svc, protocol: Option, + extra_labels: Vec<(SharedString, SharedString)>, } impl Driver { @@ -52,6 +53,7 @@ impl Driver { input, service, protocol: None, + extra_labels: vec![], } } @@ -64,6 +66,13 @@ impl Driver { self.protocol = Some(protocol.into()); self } + + /// Add an extra label to the `BytesSent` metric emitted by this driver. + #[must_use] + pub fn label(mut self, key: impl Into, value: impl Into) -> Self { + self.extra_labels.push((key.into(), value.into())); + self + } } impl Driver @@ -92,12 +101,18 @@ where input, mut service, protocol, + extra_labels, } = self; let batched_input = input.ready_chunks(1024); pin!(batched_input); - let bytes_sent = protocol.map(|protocol| register(BytesSent { protocol })); + let bytes_sent = protocol.map(|protocol| { + register(BytesSent { + protocol, + extra_labels, + }) + }); let events_sent = RegisteredEventCache::new(()); loop { diff --git a/src/aws/mod.rs b/src/aws/mod.rs index 4438e427cc501..88cb94b643185 100644 --- a/src/aws/mod.rs +++ b/src/aws/mod.rs @@ -183,9 +183,47 @@ pub async fn create_client( where T: ClientBuilder, { - create_client_and_region::(builder, auth, region, endpoint, proxy, tls_options, timeout) - .await - .map(|(client, _)| client) + build_client_inner::( + builder, + auth, + region, + endpoint, + proxy, + tls_options, + timeout, + true, + ) + .await + .map(|(client, _)| client) +} + +/// Like [`create_client`], but suppresses transport-level `AwsBytesSent` emission. +/// +/// Use this for sinks that report bytes through the Driver to avoid double-counting +/// `component_sent_bytes_total`. +pub async fn create_client_without_transport_metrics( + builder: &T, + auth: &AwsAuthentication, + region: Option, + endpoint: Option, + proxy: &ProxyConfig, + tls_options: Option<&TlsConfig>, + timeout: Option<&AwsTimeout>, +) -> crate::Result<(T::Client, Region)> +where + T: ClientBuilder, +{ + build_client_inner::( + builder, + auth, + region, + endpoint, + proxy, + tls_options, + timeout, + false, + ) + .await } /// Create the SDK client and resolve the region using the provided settings. @@ -198,6 +236,33 @@ pub async fn create_client_and_region( tls_options: Option<&TlsConfig>, timeout: Option<&AwsTimeout>, ) -> crate::Result<(T::Client, Region)> +where + T: ClientBuilder, +{ + build_client_inner::( + builder, + auth, + region, + endpoint, + proxy, + tls_options, + timeout, + true, + ) + .await +} + +#[allow(clippy::too_many_arguments)] +async fn build_client_inner( + builder: &T, + auth: &AwsAuthentication, + region: Option, + endpoint: Option, + proxy: &ProxyConfig, + tls_options: Option<&TlsConfig>, + timeout: Option<&AwsTimeout>, + emit_bytes_sent: bool, +) -> crate::Result<(T::Client, Region)> where T: ClientBuilder, { @@ -216,7 +281,7 @@ where let connector = AwsHttpClient { http: connector, region: region.clone(), - emit_bytes_sent: true, + emit_bytes_sent, }; // Build the configuration first. @@ -487,8 +552,7 @@ where return HttpConnectorFuture::new(self.call_inner(req)); } - let bytes_sent = Arc::new(AtomicUsize::new(0)); - + let bytes_sent = Arc::new(std::sync::atomic::AtomicUsize::new(0)); let req = req.map(|body| { let bytes_sent = Arc::clone(&bytes_sent); body.map_preserve_contents(move |body| { diff --git a/src/sinks/aws_cloudwatch_logs/config.rs b/src/sinks/aws_cloudwatch_logs/config.rs index 73c39a581a490..57cce4242c457 100644 --- a/src/sinks/aws_cloudwatch_logs/config.rs +++ b/src/sinks/aws_cloudwatch_logs/config.rs @@ -7,8 +7,12 @@ use tower::ServiceBuilder; use vector_lib::{codecs::JsonSerializerConfig, configurable::configurable_component, schema}; use vrl::value::Kind; +use aws_config::Region; + use crate::{ - aws::{AwsAuthentication, ClientBuilder, RegionOrEndpoint, create_client}, + aws::{ + AwsAuthentication, ClientBuilder, RegionOrEndpoint, create_client_without_transport_metrics, + }, codecs::{Encoder, EncodingConfig}, config::{ AcknowledgementsConfig, DataType, GenerateConfig, Input, ProxyConfig, SinkConfig, @@ -189,8 +193,11 @@ pub struct CloudwatchLogsSinkConfig { } impl CloudwatchLogsSinkConfig { - pub async fn create_client(&self, proxy: &ProxyConfig) -> crate::Result { - create_client::( + pub async fn create_client( + &self, + proxy: &ProxyConfig, + ) -> crate::Result<(CloudwatchLogsClient, Region)> { + create_client_without_transport_metrics::( &CloudwatchLogsClientBuilder {}, &self.auth, self.region.region(), @@ -218,12 +225,13 @@ impl SinkConfig for CloudwatchLogsSinkConfig { let batcher_settings = self.batch.into_batcher_settings()?; let request_settings = self.request.tower.into_settings(); - let client = self.create_client(cx.proxy()).await?; + let (client, resolved_region) = self.create_client(cx.proxy()).await?; let svc = ServiceBuilder::new() .settings(request_settings, CloudwatchRetryLogic::new()) .service(CloudwatchLogsPartitionSvc::new( self.clone(), client.clone(), + resolved_region.to_string(), )?); let transformer = self.encoding.transformer(); let serializer = self.encoding.build()?; @@ -237,7 +245,7 @@ impl SinkConfig for CloudwatchLogsSinkConfig { transformer, encoder, }, - + region: resolved_region.to_string(), service: svc, }; Ok((VectorSink::from_event_streamsink(sink), healthcheck)) diff --git a/src/sinks/aws_cloudwatch_logs/integration_tests.rs b/src/sinks/aws_cloudwatch_logs/integration_tests.rs index 0b5a00a1ff8dc..d7811f0793832 100644 --- a/src/sinks/aws_cloudwatch_logs/integration_tests.rs +++ b/src/sinks/aws_cloudwatch_logs/integration_tests.rs @@ -576,7 +576,7 @@ async fn cloudwatch_healthcheck() { confinement: Default::default(), }; - let client = config.create_client(&ProxyConfig::default()).await.unwrap(); + let (client, _resolved_region) = config.create_client(&ProxyConfig::default()).await.unwrap(); healthcheck(config, client).await.unwrap(); } diff --git a/src/sinks/aws_cloudwatch_logs/request.rs b/src/sinks/aws_cloudwatch_logs/request.rs index de3b03da25396..1aa86febed306 100644 --- a/src/sinks/aws_cloudwatch_logs/request.rs +++ b/src/sinks/aws_cloudwatch_logs/request.rs @@ -22,7 +22,10 @@ use http::{HeaderValue, header::HeaderName}; use indexmap::IndexMap; use tokio::sync::oneshot; -use crate::sinks::aws_cloudwatch_logs::{config::Retention, service::CloudwatchError}; +use crate::sinks::{ + aws_cloudwatch_logs::{config::Retention, service::CloudwatchError}, + util::EncodedLength, +}; pub struct CloudwatchFuture { client: Client, @@ -32,6 +35,8 @@ pub struct CloudwatchFuture { retention_enabled: bool, events: Vec>, token_tx: Option>>, + current_batch_bytes: usize, + accumulated_bytes_sent: usize, } struct Client { @@ -55,6 +60,10 @@ enum State { } impl CloudwatchFuture { + fn batch_message_bytes(batch: &[InputLogEvent]) -> usize { + batch.iter().map(|e| e.encoded_length()).sum() + } + /// Panics if events.is_empty() #[allow(clippy::too_many_arguments)] pub(super) fn new( @@ -82,10 +91,12 @@ impl CloudwatchFuture { tags, }; - let state = if let Some(token) = token { - State::Put(client.put_logs(Some(token), events.pop().expect("No Events to send"))) + let (state, current_batch_bytes) = if let Some(token) = token { + let batch = events.pop().expect("No Events to send"); + let bytes = Self::batch_message_bytes(&batch); + (State::Put(client.put_logs(Some(token), batch)), bytes) } else { - State::DescribeStream(client.describe_stream()) + (State::DescribeStream(client.describe_stream()), 0) }; let retention_enabled = retention.enabled; @@ -98,6 +109,8 @@ impl CloudwatchFuture { create_missing_group, create_missing_stream, retention_enabled, + current_batch_bytes, + accumulated_bytes_sent: 0, } } } @@ -145,6 +158,7 @@ impl Future for CloudwatchFuture { let token = stream.upload_sequence_token; + self.current_batch_bytes = Self::batch_message_bytes(&events); info!(message = "Putting logs.", token = ?token); self.state = State::Put(self.client.put_logs(token, events)); } else if self.create_missing_stream { @@ -210,10 +224,21 @@ impl Future for CloudwatchFuture { State::Put(fut) => { let next_token = match ready!(fut.poll_unpin(cx)) { Ok(resp) => resp.next_sequence_token, - Err(err) => return Poll::Ready(Err(CloudwatchError::Put(err))), + Err(err) => { + if self.accumulated_bytes_sent > 0 { + return Poll::Ready(Err(CloudwatchError::PutPartial { + error: err, + bytes_sent: self.accumulated_bytes_sent, + })); + } + return Poll::Ready(Err(CloudwatchError::Put(err))); + } }; + self.accumulated_bytes_sent += self.current_batch_bytes; + if let Some(events) = self.events.pop() { + self.current_batch_bytes = Self::batch_message_bytes(&events); debug!(message = "Putting logs.", next_token = ?next_token); self.state = State::Put(self.client.put_logs(next_token, events)); } else { diff --git a/src/sinks/aws_cloudwatch_logs/request_builder.rs b/src/sinks/aws_cloudwatch_logs/request_builder.rs index b5c582ec05ae5..2017e7587576f 100644 --- a/src/sinks/aws_cloudwatch_logs/request_builder.rs +++ b/src/sinks/aws_cloudwatch_logs/request_builder.rs @@ -109,8 +109,8 @@ impl CloudwatchRequestBuilder { return None; } - let bytes_len = - NonZeroUsize::new(message_bytes.len()).expect("payload should never be zero length"); + let bytes_len = NonZeroUsize::new(message_bytes.len() + BATCH_SIZE_OVERHEAD) + .expect("payload should never be zero length"); let metadata = builder.with_request_size(bytes_len); Some(CloudwatchRequest { diff --git a/src/sinks/aws_cloudwatch_logs/retry.rs b/src/sinks/aws_cloudwatch_logs/retry.rs index 11759f8f49976..07b4598725442 100644 --- a/src/sinks/aws_cloudwatch_logs/retry.rs +++ b/src/sinks/aws_cloudwatch_logs/retry.rs @@ -45,7 +45,7 @@ impl RetryLogic #[allow(clippy::cognitive_complexity)] // long, but just a hair over our limit fn is_retriable_error(&self, error: &Self::Error) -> bool { match error { - CloudwatchError::Put(err) => { + CloudwatchError::Put(err) | CloudwatchError::PutPartial { error: err, .. } => { if let SdkError::ServiceError(inner) = err { let err = inner.err(); if matches!(err, PutLogEventsError::ServiceUnavailableException(_)) { diff --git a/src/sinks/aws_cloudwatch_logs/service.rs b/src/sinks/aws_cloudwatch_logs/service.rs index cdb34702aa413..17b2c7fcafd06 100644 --- a/src/sinks/aws_cloudwatch_logs/service.rs +++ b/src/sinks/aws_cloudwatch_logs/service.rs @@ -16,7 +16,6 @@ use aws_sdk_cloudwatchlogs::{ use aws_smithy_runtime_api::client::{orchestrator::HttpResponse, result::SdkError}; use chrono::Duration; use futures::{FutureExt, future::BoxFuture}; -use futures_util::TryFutureExt; use http::{ HeaderValue, header::{HeaderName, InvalidHeaderName, InvalidHeaderValue}, @@ -33,6 +32,7 @@ use tower::{ }; use vector_lib::{ finalization::EventStatus, + internal_event::{ByteSize, BytesSent, InternalEventHandle as _, Registered, register}, request_metadata::{GroupedCountByteSize, MetaDescriptive}, stream::DriverResponse, }; @@ -66,6 +66,10 @@ type Svc = Buffer< #[derive(Debug)] pub enum CloudwatchError { Put(SdkError), + PutPartial { + error: SdkError, + bytes_sent: usize, + }, DescribeLogStreams(SdkError), CreateStream(SdkError), CreateGroup(SdkError), @@ -77,6 +81,12 @@ impl fmt::Display for CloudwatchError { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { match self { CloudwatchError::Put(error) => write!(f, "CloudwatchError::Put: {error}"), + CloudwatchError::PutPartial { error, bytes_sent } => { + write!( + f, + "CloudwatchError::PutPartial({bytes_sent} bytes sent): {error}" + ) + } CloudwatchError::DescribeLogStreams(error) => { write!(f, "CloudwatchError::DescribeLogStreams: {error}") } @@ -111,6 +121,7 @@ impl From> for CloudwatchError { #[derive(Debug)] pub struct CloudwatchResponse { events_byte_size: GroupedCountByteSize, + byte_size: usize, } impl crate::sinks::util::sink::Response for CloudwatchResponse { @@ -131,6 +142,10 @@ impl DriverResponse for CloudwatchResponse { fn events_sent(&self) -> &GroupedCountByteSize { &self.events_byte_size } + + fn bytes_sent(&self) -> Option { + Some(self.byte_size) + } } #[derive(Snafu, Debug)] @@ -145,6 +160,7 @@ impl CloudwatchLogsPartitionSvc { pub fn new( config: CloudwatchLogsSinkConfig, client: CloudwatchLogsClient, + region: String, ) -> crate::Result { let request_settings = config.request.tower.into_settings(); @@ -160,12 +176,18 @@ impl CloudwatchLogsPartitionSvc { }) .collect::, HeaderError>>()?; + let bytes_sent = register(BytesSent { + protocol: "https".into(), + extra_labels: vec![("region".into(), region.into())], + }); + Ok(Self { config, clients: HashMap::new(), request_settings, client, headers, + bytes_sent, }) } } @@ -181,6 +203,7 @@ impl Service for CloudwatchLogsPartitionSvc { fn call(&mut self, mut req: BatchCloudwatchRequest) -> Self::Future { let metadata = std::mem::take(req.metadata_mut()); + let byte_size = metadata.request_encoded_size(); let events_byte_size = metadata.into_events_estimated_json_encoded_byte_size(); let key = req.key; @@ -224,9 +247,23 @@ impl Service for CloudwatchLogsPartitionSvc { svc }; + let bytes_sent_handle = self.bytes_sent.clone(); + svc.oneshot(events) - .map_ok(move |_x| CloudwatchResponse { events_byte_size }) - .map_err(Into::into) + .map(move |result| match result { + Ok(()) => Ok(CloudwatchResponse { + events_byte_size, + byte_size, + }), + Err(box_err) => { + if let Some(CloudwatchError::PutPartial { bytes_sent, .. }) = + box_err.downcast_ref::() + { + bytes_sent_handle.emit(ByteSize(*bytes_sent)); + } + Err(box_err) + } + }) .boxed() } } @@ -372,4 +409,5 @@ pub struct CloudwatchLogsPartitionSvc { clients: HashMap, request_settings: TowerRequestSettings, client: CloudwatchLogsClient, + bytes_sent: Registered, } diff --git a/src/sinks/aws_cloudwatch_logs/sink.rs b/src/sinks/aws_cloudwatch_logs/sink.rs index 7b47a5b34ff8f..b9cb2338b04b7 100644 --- a/src/sinks/aws_cloudwatch_logs/sink.rs +++ b/src/sinks/aws_cloudwatch_logs/sink.rs @@ -25,6 +25,8 @@ use crate::{ pub struct CloudwatchSink { pub batcher_settings: BatcherSettings, pub(super) request_builder: CloudwatchRequestBuilder, + /// The AWS region string for metric labels. + pub region: String, pub service: S, } @@ -64,6 +66,8 @@ where } }) .into_driver(service) + .protocol("https") + .label("region", self.region) .run() .await } diff --git a/src/sinks/aws_kinesis/config.rs b/src/sinks/aws_kinesis/config.rs index 1f358633f36b5..2666c9b4ec105 100644 --- a/src/sinks/aws_kinesis/config.rs +++ b/src/sinks/aws_kinesis/config.rs @@ -88,6 +88,7 @@ pub fn build_sink( batch_settings: BatcherSettings, client: C, retry_logic: RT, + region: String, ) -> crate::Result where C: SendRecord + Clone + Send + Sync + 'static, @@ -101,13 +102,13 @@ where { let request_limits = config.request.into_settings(); - let region = config.region.region(); + let sdk_region = config.region.region(); let service = ServiceBuilder::new() .settings::>(request_limits, retry_logic) .service(KinesisService:: { client, stream_name: config.stream_name.clone(), - region, + region: sdk_region, _phantom_t: PhantomData, _phantom_e: PhantomData, }); @@ -127,6 +128,7 @@ where service, request_builder, partition_key_field, + region, _phantom: PhantomData, }; Ok(VectorSink::from_event_streamsink(sink)) diff --git a/src/sinks/aws_kinesis/firehose/config.rs b/src/sinks/aws_kinesis/firehose/config.rs index f48eb8ad6caaa..b4be4391c0c55 100644 --- a/src/sinks/aws_kinesis/firehose/config.rs +++ b/src/sinks/aws_kinesis/firehose/config.rs @@ -11,8 +11,10 @@ use super::{ record::{KinesisFirehoseClient, KinesisFirehoseRecord}, sink::BatchKinesisRequest, }; +use aws_config::Region; + use crate::{ - aws::{ClientBuilder, create_client, is_retriable_error}, + aws::{ClientBuilder, create_client_without_transport_metrics, is_retriable_error}, config::{AcknowledgementsConfig, GenerateConfig, Input, ProxyConfig, SinkConfig, SinkContext}, sinks::{ Healthcheck, VectorSink, @@ -101,8 +103,11 @@ impl KinesisFirehoseSinkConfig { } } - pub async fn create_client(&self, proxy: &ProxyConfig) -> crate::Result { - create_client::( + pub async fn create_client( + &self, + proxy: &ProxyConfig, + ) -> crate::Result<(KinesisClient, Region)> { + create_client_without_transport_metrics::( &KinesisFirehoseClientBuilder {}, &self.base.auth, self.base.region.region(), @@ -119,7 +124,7 @@ impl KinesisFirehoseSinkConfig { #[typetag::serde(name = "aws_kinesis_firehose")] impl SinkConfig for KinesisFirehoseSinkConfig { async fn build(&self, cx: SinkContext) -> crate::Result<(VectorSink, Healthcheck)> { - let client = self.create_client(&cx.proxy).await?; + let (client, resolved_region) = self.create_client(&cx.proxy).await?; let healthcheck = self.clone().healthcheck(client.clone()).boxed(); let batch_settings = self @@ -129,6 +134,7 @@ impl SinkConfig for KinesisFirehoseSinkConfig { .limit_max_events(MAX_PAYLOAD_EVENTS)? .into_batcher_settings()?; + let region = resolved_region.to_string(); let sink = build_sink::< KinesisFirehoseClient, KinesisRecord, @@ -143,6 +149,7 @@ impl SinkConfig for KinesisFirehoseSinkConfig { KinesisRetryLogic { retry_partial: self.base.request_retry_partial, }, + region, )?; Ok((sink, healthcheck)) diff --git a/src/sinks/aws_kinesis/firehose/record.rs b/src/sinks/aws_kinesis/firehose/record.rs index 27dc2fe53e0b5..248b254db221b 100644 --- a/src/sinks/aws_kinesis/firehose/record.rs +++ b/src/sinks/aws_kinesis/firehose/record.rs @@ -66,6 +66,7 @@ impl SendRecord for KinesisFirehoseClient { .map(|output: PutRecordBatchOutput| KinesisResponse { failure_count: output.failed_put_count() as usize, events_byte_size: CountByteSize(rec_count, JsonSize::new(total_size)).into(), + byte_size: 0, #[cfg(feature = "sinks-aws_kinesis_streams")] failed_records: vec![], // Firehose doesn't support partial failure retry }) diff --git a/src/sinks/aws_kinesis/service.rs b/src/sinks/aws_kinesis/service.rs index 7c364d4ba7d8c..a18754670224c 100644 --- a/src/sinks/aws_kinesis/service.rs +++ b/src/sinks/aws_kinesis/service.rs @@ -38,6 +38,7 @@ where pub struct KinesisResponse { pub(crate) failure_count: usize, pub(crate) events_byte_size: GroupedCountByteSize, + pub(crate) byte_size: usize, #[cfg(feature = "sinks-aws_kinesis_streams")] /// Track individual failed records for retry logic (Streams only) pub(crate) failed_records: Vec, @@ -59,6 +60,10 @@ impl DriverResponse for KinesisResponse { fn events_sent(&self) -> &GroupedCountByteSize { &self.events_byte_size } + + fn bytes_sent(&self) -> Option { + Some(self.byte_size) + } } impl Service> for KinesisService @@ -79,6 +84,11 @@ where // Emission of internal events for errors and dropped events is handled upstream by the caller. fn call(&mut self, mut requests: BatchKinesisRequest) -> Self::Future { + let byte_size: usize = requests + .events + .iter() + .map(|req| req.record.encoded_length()) + .sum(); let metadata = std::mem::take(requests.metadata_mut()); let events_byte_size = metadata.into_events_estimated_json_encoded_byte_size(); @@ -95,6 +105,7 @@ where client.send(records, stream_name).await.map(|mut r| { // augment the response r.events_byte_size = events_byte_size; + r.byte_size = byte_size; r }) }) diff --git a/src/sinks/aws_kinesis/sink.rs b/src/sinks/aws_kinesis/sink.rs index dabda85bd86d6..370bddd32ce6f 100644 --- a/src/sinks/aws_kinesis/sink.rs +++ b/src/sinks/aws_kinesis/sink.rs @@ -29,6 +29,8 @@ pub struct KinesisSink { pub service: S, pub request_builder: KinesisRequestBuilder, pub partition_key_field: Option, + /// The AWS region string for metric labels. + pub region: String, pub _phantom: PhantomData, } @@ -72,6 +74,8 @@ where BatchKinesisRequest { events, metadata } }) .into_driver(self.service) + .protocol("https") + .label("region", self.region) .run() .await } diff --git a/src/sinks/aws_kinesis/streams/config.rs b/src/sinks/aws_kinesis/streams/config.rs index 8d3a1e9770d77..79c0603de175f 100644 --- a/src/sinks/aws_kinesis/streams/config.rs +++ b/src/sinks/aws_kinesis/streams/config.rs @@ -11,8 +11,10 @@ use super::{ record::{KinesisStreamClient, KinesisStreamRecord}, sink::BatchKinesisRequest, }; +use aws_config::Region; + use crate::{ - aws::{ClientBuilder, create_client, is_retriable_error}, + aws::{ClientBuilder, create_client_without_transport_metrics, is_retriable_error}, config::{AcknowledgementsConfig, Input, ProxyConfig, SinkConfig, SinkContext}, sinks::{ Healthcheck, VectorSink, @@ -100,8 +102,11 @@ impl KinesisStreamsSinkConfig { } } - pub async fn create_client(&self, proxy: &ProxyConfig) -> crate::Result { - create_client::( + pub async fn create_client( + &self, + proxy: &ProxyConfig, + ) -> crate::Result<(KinesisClient, Region)> { + create_client_without_transport_metrics::( &KinesisClientBuilder {}, &self.base.auth, self.base.region.region(), @@ -118,7 +123,7 @@ impl KinesisStreamsSinkConfig { #[typetag::serde(name = "aws_kinesis_streams")] impl SinkConfig for KinesisStreamsSinkConfig { async fn build(&self, cx: SinkContext) -> crate::Result<(VectorSink, Healthcheck)> { - let client = self.create_client(&cx.proxy).await?; + let (client, resolved_region) = self.create_client(&cx.proxy).await?; let healthcheck = self.clone().healthcheck(client.clone()).boxed(); let batch_settings = self @@ -128,6 +133,7 @@ impl SinkConfig for KinesisStreamsSinkConfig { .limit_max_events(MAX_PAYLOAD_EVENTS)? .into_batcher_settings()?; + let region = resolved_region.to_string(); let sink = build_sink::< KinesisStreamClient, KinesisRecord, @@ -142,6 +148,7 @@ impl SinkConfig for KinesisStreamsSinkConfig { KinesisRetryLogic { retry_partial: self.base.request_retry_partial, }, + region, )?; Ok((sink, healthcheck)) diff --git a/src/sinks/aws_kinesis/streams/record.rs b/src/sinks/aws_kinesis/streams/record.rs index 420ce9b62ff3b..640d3eb1deae3 100644 --- a/src/sinks/aws_kinesis/streams/record.rs +++ b/src/sinks/aws_kinesis/streams/record.rs @@ -78,6 +78,7 @@ impl SendRecord for KinesisStreamClient { failed_records: extract_failed_records(&output), failure_count: output.failed_record_count().unwrap_or(0) as usize, events_byte_size: CountByteSize(rec_count, JsonSize::new(total_size)).into(), + byte_size: 0, }) } } diff --git a/src/sinks/aws_s3/config.rs b/src/sinks/aws_s3/config.rs index a3bdd4cd06b69..c558533fd42bd 100644 --- a/src/sinks/aws_s3/config.rs +++ b/src/sinks/aws_s3/config.rs @@ -229,9 +229,9 @@ impl GenerateConfig for S3SinkConfig { #[typetag::serde(name = "aws_s3")] impl SinkConfig for S3SinkConfig { async fn build(&self, cx: SinkContext) -> crate::Result<(VectorSink, Healthcheck)> { - let service = self.create_service(&cx.proxy).await?; + let (service, region) = self.create_service(&cx.proxy).await?; let healthcheck = self.build_healthcheck(service.client())?; - let sink = self.build_processor(service, cx)?; + let sink = self.build_processor(service, cx, region)?; Ok((sink, healthcheck)) } @@ -259,6 +259,7 @@ impl S3SinkConfig { &self, service: S3Service, cx: SinkContext, + region: String, ) -> crate::Result { // Build our S3 client/service, which is what we'll ultimately feed // requests into in order to ship files to S3. We build this here in @@ -339,7 +340,13 @@ impl S3SinkConfig { filename_tz_offset: offset, }; - let sink = S3Sink::new(service, request_options, partitioner, batch_settings); + let sink = S3Sink::new( + service, + request_options, + partitioner, + batch_settings, + region, + ); return Ok(VectorSink::from_event_streamsink(sink)); } @@ -357,7 +364,13 @@ impl S3SinkConfig { filename_tz_offset: offset, }; - let sink = S3Sink::new(service, request_options, partitioner, batch_settings); + let sink = S3Sink::new( + service, + request_options, + partitioner, + batch_settings, + region, + ); Ok(VectorSink::from_event_streamsink(sink)) } @@ -366,7 +379,7 @@ impl S3SinkConfig { s3_common::config::build_healthcheck(self.bucket.clone(), client) } - pub async fn create_service(&self, proxy: &ProxyConfig) -> crate::Result { + pub async fn create_service(&self, proxy: &ProxyConfig) -> crate::Result<(S3Service, String)> { s3_common::config::create_service( &self.region, &self.auth, diff --git a/src/sinks/aws_s3/integration_tests.rs b/src/sinks/aws_s3/integration_tests.rs index 7a6106f2da7d1..71e6171ac1fbb 100644 --- a/src/sinks/aws_s3/integration_tests.rs +++ b/src/sinks/aws_s3/integration_tests.rs @@ -68,8 +68,8 @@ async fn s3_insert_message_into_with_flat_key_prefix() { ..config(&bucket, 1000000, 5.0) }; let prefix = config.key_prefix.clone(); - let service = config.create_service(&cx.globals.proxy).await.unwrap(); - let sink = config.build_processor(service, cx).unwrap(); + let (service, region) = config.create_service(&cx.globals.proxy).await.unwrap(); + let sink = config.build_processor(service, cx, region).unwrap(); let (lines, events, receiver) = make_events_batch(100, 10); run_and_assert_sink_compliance(sink, events, &AWS_SINK_TAGS).await; @@ -104,8 +104,8 @@ async fn s3_insert_message_into_with_folder_key_prefix() { ..config(&bucket, 1000000, 5.0) }; let prefix = config.key_prefix.clone(); - let service = config.create_service(&cx.globals.proxy).await.unwrap(); - let sink = config.build_processor(service, cx).unwrap(); + let (service, region) = config.create_service(&cx.globals.proxy).await.unwrap(); + let sink = config.build_processor(service, cx, region).unwrap(); let (lines, events, receiver) = make_events_batch(100, 10); run_and_assert_sink_compliance(sink, events, &AWS_SINK_TAGS).await; @@ -146,8 +146,8 @@ async fn s3_insert_message_into_with_ssekms_key_id() { }; let prefix = config.key_prefix.clone(); - let service = config.create_service(&cx.globals.proxy).await.unwrap(); - let sink = config.build_processor(service, cx).unwrap(); + let (service, region) = config.create_service(&cx.globals.proxy).await.unwrap(); + let sink = config.build_processor(service, cx, region).unwrap(); let (lines, events, receiver) = make_events_batch(100, 10); run_and_assert_sink_compliance(sink, events, &AWS_SINK_TAGS).await; @@ -184,8 +184,8 @@ async fn s3_rotate_files_after_the_buffer_size_is_reached() { ..config(&bucket, 10, 5.0) }; let prefix = config.key_prefix.clone(); - let service = config.create_service(&cx.globals.proxy).await.unwrap(); - let sink = config.build_processor(service, cx).unwrap(); + let (service, region) = config.create_service(&cx.globals.proxy).await.unwrap(); + let sink = config.build_processor(service, cx, region).unwrap(); let (lines, _events) = random_lines_with_stream(100, 30, None); @@ -243,8 +243,8 @@ async fn s3_gzip() { }; let prefix = config.key_prefix.clone(); - let service = config.create_service(&cx.globals.proxy).await.unwrap(); - let sink = config.build_processor(service, cx).unwrap(); + let (service, region) = config.create_service(&cx.globals.proxy).await.unwrap(); + let sink = config.build_processor(service, cx, region).unwrap(); let (lines, events, receiver) = make_events_batch(100, batch_size * batch_multiplier); run_and_assert_sink_compliance(sink, events, &AWS_SINK_TAGS).await; @@ -288,8 +288,8 @@ async fn s3_zstd() { }; let prefix = config.key_prefix.clone(); - let service = config.create_service(&cx.globals.proxy).await.unwrap(); - let sink = config.build_processor(service, cx).unwrap(); + let (service, region) = config.create_service(&cx.globals.proxy).await.unwrap(); + let sink = config.build_processor(service, cx, region).unwrap(); let (lines, events, receiver) = make_events_batch(100, batch_size * batch_multiplier); run_and_assert_sink_compliance(sink, events, &AWS_SINK_TAGS).await; @@ -350,8 +350,8 @@ async fn s3_insert_message_into_object_lock() { let config = config(&bucket, 1000000, 5.0); let prefix = config.key_prefix.clone(); - let service = config.create_service(&cx.globals.proxy).await.unwrap(); - let sink = config.build_processor(service, cx).unwrap(); + let (service, region) = config.create_service(&cx.globals.proxy).await.unwrap(); + let sink = config.build_processor(service, cx, region).unwrap(); let (lines, events, receiver) = make_events_batch(100, 10); run_and_assert_sink_compliance(sink, events, &AWS_SINK_TAGS).await; @@ -383,8 +383,8 @@ async fn acknowledges_failures() { ..config(&bucket, 1, 5.0) }; let prefix = config.key_prefix.clone(); - let service = config.create_service(&cx.globals.proxy).await.unwrap(); - let sink = config.build_processor(service, cx).unwrap(); + let (service, region) = config.create_service(&cx.globals.proxy).await.unwrap(); + let sink = config.build_processor(service, cx, region).unwrap(); let (_lines, events, receiver) = make_events_batch(1, 1); run_and_assert_sink_error(sink, events, &COMPONENT_ERROR_TAGS).await; @@ -401,7 +401,7 @@ async fn s3_healthchecks() { create_bucket(&bucket, false).await; let config = config(&bucket, 1, 5.0); - let service = config + let (service, _region) = config .create_service(&ProxyConfig::from_env()) .await .unwrap(); @@ -415,7 +415,7 @@ async fn s3_healthchecks() { #[tokio::test] async fn s3_healthchecks_invalid_bucket() { let config = config("s3_healthchecks_invalid_bucket", 1, 5.0); - let service = config + let (service, _region) = config .create_service(&ProxyConfig::from_env()) .await .unwrap(); @@ -438,8 +438,8 @@ async fn s3_flush_on_exhaustion() { // batch size of ten events, timeout of ten seconds let config = config(&bucket, 10, 10.0); let prefix = config.key_prefix.clone(); - let service = config.create_service(&cx.globals.proxy).await.unwrap(); - let sink = config.build_processor(service, cx).unwrap(); + let (service, region) = config.create_service(&cx.globals.proxy).await.unwrap(); + let sink = config.build_processor(service, cx, region).unwrap(); let (lines, _events) = random_lines_with_stream(100, 2, None); // only generate two events (less than batch size) @@ -507,8 +507,8 @@ async fn s3_parquet_insert_message() { }; let prefix = config.key_prefix.clone(); - let service = config.create_service(&cx.globals.proxy).await.unwrap(); - let sink = config.build_processor(service, cx).unwrap(); + let (service, region) = config.create_service(&cx.globals.proxy).await.unwrap(); + let sink = config.build_processor(service, cx, region).unwrap(); let (batch_notifier, receiver) = BatchNotifier::new_with_receiver(); let events: Vec = (0..10) diff --git a/src/sinks/aws_s_s/service.rs b/src/sinks/aws_s_s/service.rs index 9211e7af989f3..b69e5284d1333 100644 --- a/src/sinks/aws_s_s/service.rs +++ b/src/sinks/aws_s_s/service.rs @@ -7,7 +7,7 @@ use aws_smithy_runtime_api::client::{orchestrator::HttpResponse, result::SdkErro use futures::future::BoxFuture; use tower::Service; use vector_lib::{ - ByteSizeOf, event::EventStatus, request_metadata::GroupedCountByteSize, stream::DriverResponse, + event::EventStatus, request_metadata::GroupedCountByteSize, stream::DriverResponse, }; use super::{client::Client, request_builder::SendMessageEntry}; @@ -63,7 +63,7 @@ where // Emission of internal events for errors and dropped events is handled upstream by the caller. fn call(&mut self, entry: SendMessageEntry) -> Self::Future { - let byte_size = entry.size_of(); + let byte_size = entry.metadata.request_encoded_size(); let client = self.client.clone(); Box::pin(async move { client.send_message(entry, byte_size).await }) diff --git a/src/sinks/aws_s_s/sink.rs b/src/sinks/aws_s_s/sink.rs index d4d7d19093dbb..dd4183f4c8459 100644 --- a/src/sinks/aws_s_s/sink.rs +++ b/src/sinks/aws_s_s/sink.rs @@ -10,6 +10,8 @@ where request_builder: SSRequestBuilder, service: SSService, request: TowerRequestConfig, + /// The AWS region string for metric labels. + region: String, } impl SSSink @@ -21,11 +23,13 @@ where request_builder: SSRequestBuilder, request: TowerRequestConfig, publisher: C, + region: String, ) -> crate::Result { Ok(SSSink { request_builder, service: SSService::new(publisher), request, + region, }) } @@ -48,6 +52,8 @@ where .ok() }) .into_driver(service) + .protocol("https") + .label("region", self.region) .run() .await } diff --git a/src/sinks/aws_s_s/sns/config.rs b/src/sinks/aws_s_s/sns/config.rs index 1203316813c72..17b6b3e44d0db 100644 --- a/src/sinks/aws_s_s/sns/config.rs +++ b/src/sinks/aws_s_s/sns/config.rs @@ -5,8 +5,10 @@ use super::{ BaseSSSinkConfig, SSRequestBuilder, SSSink, client::SnsMessagePublisher, message_deduplication_id, message_group_id, }; +use aws_config::Region; + use crate::{ - aws::{ClientBuilder, RegionOrEndpoint, create_client}, + aws::{ClientBuilder, RegionOrEndpoint, create_client_without_transport_metrics}, config::{ AcknowledgementsConfig, DataType, GenerateConfig, Input, ProxyConfig, SinkConfig, SinkContext, @@ -45,8 +47,11 @@ impl GenerateConfig for SnsSinkConfig { } impl SnsSinkConfig { - pub(super) async fn create_client(&self, proxy: &ProxyConfig) -> crate::Result { - create_client::( + pub(super) async fn create_client( + &self, + proxy: &ProxyConfig, + ) -> crate::Result<(SnsClient, Region)> { + create_client_without_transport_metrics::( &SnsClientBuilder {}, &self.base_config.auth, self.region.region(), @@ -66,7 +71,7 @@ impl SinkConfig for SnsSinkConfig { &self, cx: SinkContext, ) -> crate::Result<(crate::sinks::VectorSink, crate::sinks::Healthcheck)> { - let client = self.create_client(&cx.proxy).await?; + let (client, resolved_region) = self.create_client(&cx.proxy).await?; let publisher = SnsMessagePublisher::new(client.clone(), self.topic_arn.clone()); @@ -79,6 +84,7 @@ impl SinkConfig for SnsSinkConfig { let message_deduplication_id = message_deduplication_id(self.base_config.message_deduplication_id.clone()); + let region = resolved_region.to_string(); let sink = SSSink::new( SSRequestBuilder::new( message_group_id?, @@ -87,6 +93,7 @@ impl SinkConfig for SnsSinkConfig { )?, self.base_config.request, publisher, + region, )?; Ok(( crate::sinks::VectorSink::from_event_streamsink(sink), diff --git a/src/sinks/aws_s_s/sqs/config.rs b/src/sinks/aws_s_s/sqs/config.rs index 67e8f22870d5b..18386f16d8b08 100644 --- a/src/sinks/aws_s_s/sqs/config.rs +++ b/src/sinks/aws_s_s/sqs/config.rs @@ -5,8 +5,10 @@ use super::{ BaseSSSinkConfig, SSRequestBuilder, SSSink, client::SqsMessagePublisher, message_deduplication_id, message_group_id, }; +use aws_config::Region; + use crate::{ - aws::{RegionOrEndpoint, create_client}, + aws::{RegionOrEndpoint, create_client_without_transport_metrics}, common::sqs::SqsClientBuilder, config::{ AcknowledgementsConfig, DataType, GenerateConfig, Input, ProxyConfig, SinkConfig, @@ -48,8 +50,11 @@ impl GenerateConfig for SqsSinkConfig { } impl SqsSinkConfig { - pub(super) async fn create_client(&self, proxy: &ProxyConfig) -> crate::Result { - create_client::( + pub(super) async fn create_client( + &self, + proxy: &ProxyConfig, + ) -> crate::Result<(SqsClient, Region)> { + create_client_without_transport_metrics::( &SqsClientBuilder {}, &self.base_config.auth, self.region.region(), @@ -69,7 +74,7 @@ impl SinkConfig for SqsSinkConfig { &self, cx: SinkContext, ) -> crate::Result<(crate::sinks::VectorSink, crate::sinks::Healthcheck)> { - let client = self.create_client(&cx.proxy).await?; + let (client, resolved_region) = self.create_client(&cx.proxy).await?; let publisher = SqsMessagePublisher::new(client.clone(), self.queue_url.clone()); @@ -81,6 +86,7 @@ impl SinkConfig for SqsSinkConfig { let message_deduplication_id = message_deduplication_id(self.base_config.message_deduplication_id.clone()); + let region = resolved_region.to_string(); let sink = SSSink::new( SSRequestBuilder::new( message_group_id?, @@ -89,6 +95,7 @@ impl SinkConfig for SqsSinkConfig { )?, self.base_config.request, publisher, + region, )?; Ok(( crate::sinks::VectorSink::from_event_streamsink(sink), diff --git a/src/sinks/s3_common/config.rs b/src/sinks/s3_common/config.rs index f78c649c6f667..a068843516344 100644 --- a/src/sinks/s3_common/config.rs +++ b/src/sinks/s3_common/config.rs @@ -15,7 +15,10 @@ use vector_lib::configurable::configurable_component; use super::service::{S3Request, S3Response, S3Service}; use crate::{ - aws::{AwsAuthentication, RegionOrEndpoint, create_client, is_retriable_error}, + aws::{ + AwsAuthentication, RegionOrEndpoint, create_client_without_transport_metrics, + is_retriable_error, + }, common::s3::S3ClientBuilder, config::ProxyConfig, http::status, @@ -431,11 +434,11 @@ pub async fn create_service( proxy: &ProxyConfig, tls_options: Option<&TlsConfig>, force_path_style: impl Into, -) -> crate::Result { +) -> crate::Result<(S3Service, String)> { let endpoint = region.endpoint(); let region = region.region(); let force_path_style_value: bool = force_path_style.into(); - let client = create_client::( + let (client, resolved_region) = create_client_without_transport_metrics::( &S3ClientBuilder { force_path_style: Some(force_path_style_value), }, @@ -447,7 +450,7 @@ pub async fn create_service( None, ) .await?; - Ok(S3Service::new(client)) + Ok((S3Service::new(client), resolved_region.to_string())) } #[cfg(test)] diff --git a/src/sinks/s3_common/service.rs b/src/sinks/s3_common/service.rs index 08b5cad0e0016..fd1cefcd71994 100644 --- a/src/sinks/s3_common/service.rs +++ b/src/sinks/s3_common/service.rs @@ -53,6 +53,7 @@ pub struct S3Metadata { #[derive(Debug)] pub struct S3Response { events_byte_size: GroupedCountByteSize, + byte_size: usize, } impl DriverResponse for S3Response { @@ -63,6 +64,10 @@ impl DriverResponse for S3Response { fn events_sent(&self) -> &GroupedCountByteSize { &self.events_byte_size } + + fn bytes_sent(&self) -> Option { + Some(self.byte_size) + } } /// Wrapper for the AWS SDK S3 client. @@ -118,6 +123,7 @@ impl Service for S3Service { tagging.finish() }); + let byte_size = request.request_metadata.request_encoded_size(); let events_byte_size = request .request_metadata .into_events_estimated_json_encoded_byte_size(); @@ -153,7 +159,10 @@ impl Service for S3Service { key = request.metadata.s3_key ); - S3Response { events_byte_size } + S3Response { + events_byte_size, + byte_size, + } }) }) } diff --git a/src/sinks/s3_common/sink.rs b/src/sinks/s3_common/sink.rs index 5cad7ba89b6ee..af491d65a7b68 100644 --- a/src/sinks/s3_common/sink.rs +++ b/src/sinks/s3_common/sink.rs @@ -10,20 +10,25 @@ pub struct S3Sink { request_builder: RB, partitioner: P, batcher_settings: BatcherSettings, + /// The AWS region string for metric labels. + region: String, } impl S3Sink { + /// Create a new S3 sink. pub const fn new( service: Svc, request_builder: RB, partitioner: P, batcher_settings: BatcherSettings, + region: String, ) -> Self { Self { partitioner, service, request_builder, batcher_settings, + region, } } } @@ -62,6 +67,8 @@ where } }) .into_driver(self.service) + .protocol("https") + .label("region", self.region) .run() .await }