Skip to content

Commit cd83afc

Browse files
enhancement(aws sinks): use validated config pattern
1 parent 44b30fd commit cd83afc

14 files changed

Lines changed: 648 additions & 200 deletions

File tree

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,4 @@
1+
The `ssekms_key_id` option in the `aws_s3` sink now respects the configured timezone when the
2+
value is a template containing time components, matching the existing behavior of `key_prefix`.
3+
4+
authors: thomasqueirozb

‎src/sinks/aws_cloudwatch_logs/config.rs‎

Lines changed: 127 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -1,18 +1,22 @@
1-
use std::collections::HashMap;
1+
use std::collections::{BTreeMap, HashMap};
22

33
use aws_sdk_cloudwatchlogs::Client as CloudwatchLogsClient;
44
use futures::FutureExt;
5+
use http::HeaderValue;
56
use serde::{Deserialize, Deserializer, de};
67
use tower::ServiceBuilder;
7-
use vector_lib::{codecs::JsonSerializerConfig, configurable::configurable_component, schema};
8+
use vector_lib::{
9+
codecs::JsonSerializerConfig, configurable::configurable_component, schema,
10+
stream::BatcherSettings,
11+
};
812
use vrl::value::Kind;
913

1014
use crate::{
1115
aws::{AwsAuthentication, ClientBuilder, RegionOrEndpoint, create_client},
1216
codecs::{Encoder, EncodingConfig},
1317
config::{
14-
AcknowledgementsConfig, DataType, GenerateConfig, Input, ProxyConfig, SinkConfig,
15-
SinkContext,
18+
AcknowledgementsConfig, DataType, DynValidatedSink, GenerateConfig, Input, ProxyConfig,
19+
SinkConfig, SinkContext, ValidatedSink,
1620
},
1721
sinks::{
1822
Healthcheck, VectorSink,
@@ -21,10 +25,11 @@ use crate::{
2125
retry::CloudwatchRetryLogic, service::CloudwatchLogsPartitionSvc, sink::CloudwatchSink,
2226
},
2327
util::{
24-
BatchConfig, Compression, ServiceBuilderExt, SinkBatchSettings, http::RequestConfig,
28+
BatchConfig, Compression, ServiceBuilderExt, SinkBatchSettings,
29+
http::{OrderedHeaderName, RequestConfig, validate_headers},
2530
},
2631
},
27-
template::{ConfinementConfig, Template},
32+
template::{ConfinedTemplate, ConfinementConfig, Template},
2833
tls::TlsConfig,
2934
};
3035

@@ -88,7 +93,7 @@ pub struct CloudwatchLogsSinkConfig {
8893
///
8994
/// [group_name]: https://docs.aws.amazon.com/AmazonCloudWatch/latest/logs/Working-with-log-groups-and-streams.html
9095
#[configurable(metadata(docs::examples = "group-name"))]
91-
#[configurable(metadata(docs::examples = "{{ file }}"))]
96+
#[configurable(metadata(docs::examples = "group-{{ file }}"))]
9297
pub group_name: Template,
9398

9499
/// The [stream name][stream_name] of the target CloudWatch Logs stream.
@@ -98,7 +103,7 @@ pub struct CloudwatchLogsSinkConfig {
98103
/// unique per instance.
99104
///
100105
/// [stream_name]: https://docs.aws.amazon.com/AmazonCloudWatch/latest/logs/Working-with-log-groups-and-streams.html
101-
#[configurable(metadata(docs::examples = "{{ host }}"))]
106+
#[configurable(metadata(docs::examples = "stream-{{ host }}"))]
102107
#[configurable(metadata(docs::examples = "%Y-%m-%d"))]
103108
#[configurable(metadata(docs::examples = "stream-name"))]
104109
pub stream_name: Template,
@@ -206,7 +211,40 @@ impl CloudwatchLogsSinkConfig {
206211
#[async_trait::async_trait]
207212
#[typetag::serde(name = "aws_cloudwatch_logs")]
208213
impl SinkConfig for CloudwatchLogsSinkConfig {
209-
async fn build(&self, cx: SinkContext) -> crate::Result<(VectorSink, Healthcheck)> {
214+
fn confinement_config(&self) -> Option<&crate::template::ConfinementConfig> {
215+
Some(&self.confinement)
216+
}
217+
218+
fn input(&self) -> Input {
219+
let requirement =
220+
schema::Requirement::empty().optional_meaning("timestamp", Kind::timestamp());
221+
222+
Input::new(self.encoding.config().input_type() & DataType::Log)
223+
.with_schema_requirement(requirement)
224+
}
225+
226+
fn acknowledgements(&self) -> &AcknowledgementsConfig {
227+
&self.acknowledgements
228+
}
229+
230+
fn as_dyn_validated(&self) -> Option<&dyn DynValidatedSink> {
231+
Some(self)
232+
}
233+
}
234+
235+
#[derive(Clone, Debug)]
236+
pub struct ValidatedCloudwatchLogs {
237+
group_template: ConfinedTemplate,
238+
stream_template: ConfinedTemplate,
239+
batcher_settings: BatcherSettings,
240+
headers: BTreeMap<OrderedHeaderName, HeaderValue>,
241+
}
242+
243+
#[async_trait::async_trait]
244+
impl ValidatedSink for CloudwatchLogsSinkConfig {
245+
type Validated = ValidatedCloudwatchLogs;
246+
247+
fn validate(&self) -> crate::Result<ValidatedCloudwatchLogs> {
210248
let group_template =
211249
self.group_name
212250
.clone()
@@ -215,16 +253,37 @@ impl SinkConfig for CloudwatchLogsSinkConfig {
215253
self.stream_name
216254
.clone()
217255
.confine(&self.confinement, Self::NAME, "stream_name")?;
218-
219256
let batcher_settings = self.batch.into_batcher_settings()?;
257+
let headers = validate_headers(&self.request.headers)?;
258+
259+
Ok(ValidatedCloudwatchLogs {
260+
group_template,
261+
stream_template,
262+
batcher_settings,
263+
headers,
264+
})
265+
}
266+
267+
async fn build(
268+
&self,
269+
validated: &ValidatedCloudwatchLogs,
270+
cx: SinkContext,
271+
) -> crate::Result<(VectorSink, Healthcheck)> {
272+
let ValidatedCloudwatchLogs {
273+
group_template,
274+
stream_template,
275+
batcher_settings,
276+
headers,
277+
} = validated.clone();
220278
let request_settings = self.request.tower.into_settings();
221279
let client = self.create_client(cx.proxy()).await?;
222280
let svc = ServiceBuilder::new()
223281
.settings(request_settings, CloudwatchRetryLogic::new())
224282
.service(CloudwatchLogsPartitionSvc::new(
225283
self.clone(),
226284
client.clone(),
227-
)?);
285+
headers.clone(),
286+
));
228287
let transformer = self.encoding.transformer();
229288
let serializer = self.encoding.build()?;
230289
let encoder = Encoder::<()>::new(serializer);
@@ -242,22 +301,6 @@ impl SinkConfig for CloudwatchLogsSinkConfig {
242301
};
243302
Ok((VectorSink::from_event_streamsink(sink), healthcheck))
244303
}
245-
246-
fn confinement_config(&self) -> Option<&crate::template::ConfinementConfig> {
247-
Some(&self.confinement)
248-
}
249-
250-
fn input(&self) -> Input {
251-
let requirement =
252-
schema::Requirement::empty().optional_meaning("timestamp", Kind::timestamp());
253-
254-
Input::new(self.encoding.config().input_type() & DataType::Log)
255-
.with_schema_requirement(requirement)
256-
}
257-
258-
fn acknowledgements(&self) -> &AcknowledgementsConfig {
259-
&self.acknowledgements
260-
}
261304
}
262305

263306
impl GenerateConfig for CloudwatchLogsSinkConfig {
@@ -299,14 +342,71 @@ impl SinkBatchSettings for CloudwatchLogsDefaultBatchSettings {
299342

300343
#[cfg(test)]
301344
mod tests {
345+
use crate::config::ValidatedSink;
302346
use crate::sinks::aws_cloudwatch_logs::config::CloudwatchLogsSinkConfig;
303347
use crate::template::{ConfinementConfig, Template};
348+
use vector_lib::codecs::JsonSerializerConfig;
349+
350+
#[test]
351+
fn prepares_valid_config() {
352+
let mut config = super::default_config(JsonSerializerConfig::default().into());
353+
config.group_name = "group-{{ file }}".try_into().unwrap();
354+
config.stream_name = "stream".try_into().unwrap();
355+
356+
let validated = config.validate().expect("preparation should succeed");
357+
assert_eq!(validated.group_template.to_string(), "group-{{ file }}");
358+
assert_eq!(validated.stream_template.to_string(), "stream");
359+
assert_eq!(validated.batcher_settings.item_limit, 10_000);
360+
}
304361

305362
#[test]
306363
fn test_generate_config() {
307364
crate::test_util::test_generate_config::<CloudwatchLogsSinkConfig>();
308365
}
309366

367+
#[test]
368+
fn validate_rejects_invalid_header_name() {
369+
let mut config = super::default_config(JsonSerializerConfig::default().into());
370+
config.group_name = "group".try_into().unwrap();
371+
config.stream_name = "stream".try_into().unwrap();
372+
config
373+
.request
374+
.headers
375+
.insert("invalid header name".to_string(), "value".to_string());
376+
377+
assert!(config.validate().is_err());
378+
}
379+
380+
#[test]
381+
fn validate_rejects_invalid_header_value() {
382+
let mut config = super::default_config(JsonSerializerConfig::default().into());
383+
config.group_name = "group".try_into().unwrap();
384+
config.stream_name = "stream".try_into().unwrap();
385+
config.request.headers.insert(
386+
"valid-header".to_string(),
387+
"value\nwith newline".to_string(),
388+
);
389+
390+
assert!(config.validate().is_err());
391+
}
392+
393+
#[test]
394+
fn validate_retains_valid_headers() {
395+
let mut config = super::default_config(JsonSerializerConfig::default().into());
396+
config.group_name = "group".try_into().unwrap();
397+
config.stream_name = "stream".try_into().unwrap();
398+
config
399+
.request
400+
.headers
401+
.insert("x-custom-header".to_string(), "custom-value".to_string());
402+
403+
let validated = config.validate().expect("preparation should succeed");
404+
assert_eq!(validated.headers.len(), 1);
405+
let (name, value) = validated.headers.iter().next().unwrap();
406+
assert_eq!(name.inner(), "x-custom-header");
407+
assert_eq!(value, "custom-value");
408+
}
409+
310410
#[test]
311411
fn confinement_rejects_unconfined_group_name() {
312412
let template = Template::try_from("{{ group }}").unwrap();

‎src/sinks/aws_cloudwatch_logs/request.rs‎

Lines changed: 10 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
use std::{
2-
collections::HashMap,
2+
collections::{BTreeMap, HashMap},
33
future::Future,
44
pin::Pin,
55
task::{Context, Poll, ready},
@@ -18,11 +18,13 @@ use aws_sdk_cloudwatchlogs::{
1818
};
1919
use aws_smithy_runtime_api::client::{orchestrator::HttpResponse, result::SdkError};
2020
use futures::{FutureExt, future::BoxFuture};
21-
use http::{HeaderValue, header::HeaderName};
22-
use indexmap::IndexMap;
21+
use http::HeaderValue;
2322
use tokio::sync::oneshot;
2423

25-
use crate::sinks::aws_cloudwatch_logs::{config::Retention, service::CloudwatchError};
24+
use crate::sinks::{
25+
aws_cloudwatch_logs::{config::Retention, service::CloudwatchError},
26+
util::http::OrderedHeaderName,
27+
};
2628

2729
pub struct CloudwatchFuture {
2830
client: Client,
@@ -38,7 +40,7 @@ struct Client {
3840
client: CloudwatchLogsClient,
3941
stream_name: String,
4042
group_name: String,
41-
headers: IndexMap<HeaderName, HeaderValue>,
43+
headers: BTreeMap<OrderedHeaderName, HeaderValue>,
4244
retention_days: u32,
4345
kms_key: Option<String>,
4446
tags: Option<HashMap<String, String>>,
@@ -59,7 +61,7 @@ impl CloudwatchFuture {
5961
#[allow(clippy::too_many_arguments)]
6062
pub(super) fn new(
6163
client: CloudwatchLogsClient,
62-
headers: IndexMap<HeaderName, HeaderValue>,
64+
headers: BTreeMap<OrderedHeaderName, HeaderValue>,
6365
stream_name: String,
6466
group_name: String,
6567
create_missing_group: bool,
@@ -267,7 +269,8 @@ impl Client {
267269
.customize()
268270
.mutate_request(move |req| {
269271
for (header, value) in headers.iter() {
270-
req.headers_mut().insert(header.clone(), value.clone());
272+
req.headers_mut()
273+
.insert(header.inner().clone(), value.clone());
271274
}
272275
})
273276
.send()

‎src/sinks/aws_cloudwatch_logs/service.rs‎

Lines changed: 12 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
use std::{
2-
collections::HashMap,
2+
collections::{BTreeMap, HashMap},
33
fmt,
44
task::{Context, Poll, ready},
55
};
@@ -17,12 +17,7 @@ use aws_smithy_runtime_api::client::{orchestrator::HttpResponse, result::SdkErro
1717
use chrono::Duration;
1818
use futures::{FutureExt, future::BoxFuture};
1919
use futures_util::TryFutureExt;
20-
use http::{
21-
HeaderValue,
22-
header::{HeaderName, InvalidHeaderName, InvalidHeaderValue},
23-
};
24-
use indexmap::IndexMap;
25-
use snafu::{ResultExt, Snafu};
20+
use http::HeaderValue;
2621
use tokio::sync::oneshot;
2722
use tower::{
2823
Service, ServiceBuilder, ServiceExt,
@@ -45,7 +40,9 @@ use crate::sinks::{
4540
retry::CloudwatchRetryLogic,
4641
sink::BatchCloudwatchRequest,
4742
},
48-
util::{EncodedLength, TowerRequestSettings, retries::FibonacciRetryPolicy},
43+
util::{
44+
EncodedLength, TowerRequestSettings, http::OrderedHeaderName, retries::FibonacciRetryPolicy,
45+
},
4946
};
5047

5148
type Svc = Buffer<
@@ -133,40 +130,21 @@ impl DriverResponse for CloudwatchResponse {
133130
}
134131
}
135132

136-
#[derive(Snafu, Debug)]
137-
enum HeaderError {
138-
#[snafu(display("invalid header name {source}"))]
139-
InvalidName { source: InvalidHeaderName },
140-
#[snafu(display("invalid header value {source}"))]
141-
InvalidValue { source: InvalidHeaderValue },
142-
}
143-
144133
impl CloudwatchLogsPartitionSvc {
145134
pub fn new(
146135
config: CloudwatchLogsSinkConfig,
147136
client: CloudwatchLogsClient,
148-
) -> crate::Result<Self> {
137+
headers: BTreeMap<OrderedHeaderName, HeaderValue>,
138+
) -> Self {
149139
let request_settings = config.request.tower.into_settings();
150140

151-
let headers = config
152-
.request
153-
.headers
154-
.iter()
155-
.map(|(name, value)| {
156-
Ok((
157-
HeaderName::from_bytes(name.as_bytes()).context(InvalidNameSnafu {})?,
158-
HeaderValue::from_str(value.as_str()).context(InvalidValueSnafu {})?,
159-
))
160-
})
161-
.collect::<Result<IndexMap<_, _>, HeaderError>>()?;
162-
163-
Ok(Self {
141+
Self {
164142
config,
165143
clients: HashMap::new(),
166144
request_settings,
167145
client,
168146
headers,
169-
})
147+
}
170148
}
171149
}
172150

@@ -236,7 +214,7 @@ impl CloudwatchLogsSvc {
236214
config: CloudwatchLogsSinkConfig,
237215
key: &CloudwatchKey,
238216
client: CloudwatchLogsClient,
239-
headers: IndexMap<HeaderName, HeaderValue>,
217+
headers: BTreeMap<OrderedHeaderName, HeaderValue>,
240218
) -> Self {
241219
let group_name = key.group.clone();
242220
let stream_name = key.stream.clone();
@@ -347,7 +325,7 @@ impl Service<Vec<InputLogEvent>> for CloudwatchLogsSvc {
347325

348326
pub struct CloudwatchLogsSvc {
349327
client: CloudwatchLogsClient,
350-
headers: IndexMap<HeaderName, HeaderValue>,
328+
headers: BTreeMap<OrderedHeaderName, HeaderValue>,
351329
stream_name: String,
352330
group_name: String,
353331
create_missing_group: bool,
@@ -368,7 +346,7 @@ impl EncodedLength for InputLogEvent {
368346
#[derive(Clone)]
369347
pub struct CloudwatchLogsPartitionSvc {
370348
config: CloudwatchLogsSinkConfig,
371-
headers: IndexMap<HeaderName, HeaderValue>,
349+
headers: BTreeMap<OrderedHeaderName, HeaderValue>,
372350
clients: HashMap<CloudwatchKey, Svc>,
373351
request_settings: TowerRequestSettings,
374352
client: CloudwatchLogsClient,

0 commit comments

Comments
 (0)