vector/sinks/aws_cloudwatch_logs/
config.rs1use std::collections::HashMap;
2
3use aws_sdk_cloudwatchlogs::Client as CloudwatchLogsClient;
4use futures::FutureExt;
5use serde::{Deserialize, Deserializer, de};
6use tower::ServiceBuilder;
7use vector_lib::{codecs::JsonSerializerConfig, configurable::configurable_component, schema};
8use vrl::value::Kind;
9
10use crate::{
11 aws::{AwsAuthentication, ClientBuilder, RegionOrEndpoint, create_client},
12 codecs::{Encoder, EncodingConfig},
13 config::{
14 AcknowledgementsConfig, DataType, GenerateConfig, Input, ProxyConfig, SinkConfig,
15 SinkContext,
16 },
17 sinks::{
18 Healthcheck, VectorSink,
19 aws_cloudwatch_logs::{
20 healthcheck::healthcheck, request_builder::CloudwatchRequestBuilder,
21 retry::CloudwatchRetryLogic, service::CloudwatchLogsPartitionSvc, sink::CloudwatchSink,
22 },
23 util::{
24 BatchConfig, Compression, ServiceBuilderExt, SinkBatchSettings, http::RequestConfig,
25 },
26 },
27 template::{ConfinementConfig, Template},
28 tls::TlsConfig,
29};
30
31pub struct CloudwatchLogsClientBuilder;
32
33impl ClientBuilder for CloudwatchLogsClientBuilder {
34 type Client = aws_sdk_cloudwatchlogs::client::Client;
35
36 fn build(&self, config: &aws_types::SdkConfig) -> Self::Client {
37 aws_sdk_cloudwatchlogs::client::Client::new(config)
38 }
39}
40
41#[configurable_component]
42#[derive(Clone, Debug, Default)]
43pub struct Retention {
45 #[serde(default)]
47 pub enabled: bool,
48
49 #[serde(
51 default,
52 deserialize_with = "retention_days",
53 skip_serializing_if = "crate::serde::is_default"
54 )]
55 pub days: u32,
56}
57
58fn retention_days<'de, D>(deserializer: D) -> Result<u32, D::Error>
59where
60 D: Deserializer<'de>,
61{
62 let days: u32 = Deserialize::deserialize(deserializer)?;
63 const ALLOWED_VALUES: &[u32] = &[
64 1, 3, 5, 7, 14, 30, 60, 90, 120, 150, 180, 365, 400, 545, 731, 1096, 1827, 2192, 2557,
65 2922, 3288, 3653,
66 ];
67 if ALLOWED_VALUES.contains(&days) {
68 Ok(days)
69 } else {
70 let msg = format!("one of allowed values: {ALLOWED_VALUES:?}").to_owned();
71 let expected: &str = msg.as_str();
72 Err(de::Error::invalid_value(
73 de::Unexpected::Signed(days.into()),
74 &expected,
75 ))
76 }
77}
78
79#[configurable_component(sink(
81 "aws_cloudwatch_logs",
82 "Publish log events to AWS CloudWatch Logs."
83))]
84#[derive(Clone, Debug)]
85#[serde(deny_unknown_fields)]
86pub struct CloudwatchLogsSinkConfig {
87 #[configurable(metadata(docs::examples = "group-name"))]
91 #[configurable(metadata(docs::examples = "{{ file }}"))]
92 pub group_name: Template,
93
94 #[configurable(metadata(docs::examples = "{{ host }}"))]
102 #[configurable(metadata(docs::examples = "%Y-%m-%d"))]
103 #[configurable(metadata(docs::examples = "stream-name"))]
104 pub stream_name: Template,
105
106 #[serde(flatten)]
110 pub region: RegionOrEndpoint,
111
112 #[serde(default = "crate::serde::default_true")]
119 pub create_missing_group: bool,
120
121 #[serde(default = "crate::serde::default_true")]
125 pub create_missing_stream: bool,
126
127 #[configurable(derived)]
128 #[serde(default)]
129 pub retention: Retention,
130
131 #[configurable(derived)]
132 pub encoding: EncodingConfig,
133
134 #[configurable(derived)]
135 #[serde(default)]
136 pub compression: Compression,
137
138 #[configurable(derived)]
139 #[serde(default)]
140 pub batch: BatchConfig<CloudwatchLogsDefaultBatchSettings>,
141
142 #[configurable(derived)]
143 #[serde(default)]
144 pub request: RequestConfig,
145
146 #[configurable(derived)]
147 pub tls: Option<TlsConfig>,
148
149 #[configurable(deprecated)]
153 #[configurable(metadata(docs::hidden))]
154 pub assume_role: Option<String>,
155
156 #[configurable(derived)]
157 #[serde(default)]
158 pub auth: AwsAuthentication,
159
160 #[configurable(derived)]
161 #[serde(
162 default,
163 deserialize_with = "crate::serde::bool_or_struct",
164 skip_serializing_if = "crate::serde::is_default"
165 )]
166 pub acknowledgements: AcknowledgementsConfig,
167
168 #[configurable(derived)]
173 #[serde(default)]
174 pub kms_key: Option<String>,
175
176 #[configurable(derived)]
180 #[serde(default)]
181 #[configurable(metadata(
182 docs::additional_props_description = "A tag represented as a key-value pair"
183 ))]
184 pub tags: Option<HashMap<String, String>>,
185
186 #[configurable(derived)]
187 #[serde(flatten)]
188 pub confinement: ConfinementConfig,
189}
190
191impl CloudwatchLogsSinkConfig {
192 pub async fn create_client(&self, proxy: &ProxyConfig) -> crate::Result<CloudwatchLogsClient> {
193 create_client::<CloudwatchLogsClientBuilder>(
194 &CloudwatchLogsClientBuilder {},
195 &self.auth,
196 self.region.region(),
197 self.region.endpoint(),
198 proxy,
199 self.tls.as_ref(),
200 None,
201 )
202 .await
203 }
204}
205
206#[async_trait::async_trait]
207#[typetag::serde(name = "aws_cloudwatch_logs")]
208impl SinkConfig for CloudwatchLogsSinkConfig {
209 async fn build(&self, cx: SinkContext) -> crate::Result<(VectorSink, Healthcheck)> {
210 let group_template =
211 self.group_name
212 .clone()
213 .confine(&self.confinement, Self::NAME, "group_name")?;
214 let stream_template =
215 self.stream_name
216 .clone()
217 .confine(&self.confinement, Self::NAME, "stream_name")?;
218
219 let batcher_settings = self.batch.into_batcher_settings()?;
220 let request_settings = self.request.tower.into_settings();
221 let client = self.create_client(cx.proxy()).await?;
222 let svc = ServiceBuilder::new()
223 .settings(request_settings, CloudwatchRetryLogic::new())
224 .service(CloudwatchLogsPartitionSvc::new(
225 self.clone(),
226 client.clone(),
227 )?);
228 let transformer = self.encoding.transformer();
229 let serializer = self.encoding.build()?;
230 let encoder = Encoder::<()>::new(serializer);
231 let healthcheck = healthcheck(self.clone(), client).boxed();
232 let sink = CloudwatchSink {
233 batcher_settings,
234 request_builder: CloudwatchRequestBuilder {
235 group_template,
236 stream_template,
237 transformer,
238 encoder,
239 },
240
241 service: svc,
242 };
243
244 self.confinement.set_confinement_gauge("sink", Self::NAME);
245 Ok((VectorSink::from_event_streamsink(sink), healthcheck))
246 }
247
248 fn input(&self) -> Input {
249 let requirement =
250 schema::Requirement::empty().optional_meaning("timestamp", Kind::timestamp());
251
252 Input::new(self.encoding.config().input_type() & DataType::Log)
253 .with_schema_requirement(requirement)
254 }
255
256 fn acknowledgements(&self) -> &AcknowledgementsConfig {
257 &self.acknowledgements
258 }
259}
260
261impl GenerateConfig for CloudwatchLogsSinkConfig {
262 fn generate_config() -> toml::Value {
263 toml::Value::try_from(default_config(JsonSerializerConfig::default().into())).unwrap()
264 }
265}
266
267fn default_config(encoding: EncodingConfig) -> CloudwatchLogsSinkConfig {
268 CloudwatchLogsSinkConfig {
269 encoding,
270 group_name: Default::default(),
271 stream_name: Default::default(),
272 region: Default::default(),
273 create_missing_group: true,
274 create_missing_stream: true,
275 retention: Default::default(),
276 compression: Default::default(),
277 batch: Default::default(),
278 request: Default::default(),
279 tls: Default::default(),
280 assume_role: Default::default(),
281 auth: Default::default(),
282 acknowledgements: Default::default(),
283 kms_key: Default::default(),
284 tags: Default::default(),
285 confinement: Default::default(),
286 }
287}
288
289#[derive(Clone, Copy, Debug, Default)]
290pub struct CloudwatchLogsDefaultBatchSettings;
291
292impl SinkBatchSettings for CloudwatchLogsDefaultBatchSettings {
293 const MAX_EVENTS: Option<usize> = Some(10_000);
294 const MAX_BYTES: Option<usize> = Some(1_048_576);
295 const TIMEOUT_SECS: f64 = 1.0;
296}
297
298#[cfg(test)]
299mod tests {
300 use crate::sinks::aws_cloudwatch_logs::config::CloudwatchLogsSinkConfig;
301 use crate::template::{ConfinementConfig, Template};
302
303 #[test]
304 fn test_generate_config() {
305 crate::test_util::test_generate_config::<CloudwatchLogsSinkConfig>();
306 }
307
308 #[test]
309 fn confinement_rejects_unconfined_group_name() {
310 let template = Template::try_from("{{ group }}").unwrap();
311 let config = ConfinementConfig::default();
312 let result = template.confine(&config, "aws_cloudwatch_logs", "group_name");
313 assert!(result.is_err());
314 }
315
316 #[test]
317 fn confinement_opt_out_allows_unconfined_group_name() {
318 let template = Template::try_from("{{ group }}").unwrap();
319 let config = ConfinementConfig {
320 dangerously_allow_unconfined_template_resolution: true,
321 };
322 let result = template.confine(&config, "aws_cloudwatch_logs", "group_name");
323 assert!(result.is_ok());
324 }
325
326 #[test]
327 fn confinement_allows_prefixed_group_name() {
328 let template = Template::try_from("events-{{ env }}").unwrap();
329 let config = ConfinementConfig::default();
330 let result = template.confine(&config, "aws_cloudwatch_logs", "group_name");
331 assert!(result.is_ok());
332 }
333}