vector/sinks/aws_kinesis/firehose/
config.rs1use aws_sdk_firehose::operation::{
2 describe_delivery_stream::DescribeDeliveryStreamError, put_record_batch::PutRecordBatchError,
3};
4use aws_smithy_runtime_api::client::{orchestrator::HttpResponse, result::SdkError};
5use futures::FutureExt;
6use snafu::Snafu;
7use vector_lib::configurable::configurable_component;
8use vector_lib::stream::BatcherSettings;
9
10use super::{
11 KinesisClient, KinesisError, KinesisRecord, KinesisResponse, KinesisSinkBaseConfig, build_sink,
12 record::{KinesisFirehoseClient, KinesisFirehoseRecord},
13 sink::BatchKinesisRequest,
14};
15use crate::{
16 aws::{ClientBuilder, create_client, is_retriable_error},
17 config::{
18 AcknowledgementsConfig, DynValidatedSink, GenerateConfig, Input, ProxyConfig, SinkConfig,
19 SinkContext, ValidatedSink,
20 },
21 sinks::{
22 Healthcheck, VectorSink,
23 util::{
24 BatchConfig, SinkBatchSettings,
25 retries::{RetryAction, RetryLogic},
26 },
27 },
28};
29
30#[allow(clippy::large_enum_variant)]
31#[derive(Debug, Snafu)]
32enum HealthcheckError {
33 #[snafu(display("DescribeDeliveryStream failed: {}", source))]
34 DescribeDeliveryStreamFailed {
35 source: SdkError<DescribeDeliveryStreamError, HttpResponse>,
36 },
37 #[snafu(display("Stream name does not match, got {}, expected {}", name, stream_name))]
38 StreamNamesMismatch { name: String, stream_name: String },
39}
40
41pub struct KinesisFirehoseClientBuilder;
42
43impl ClientBuilder for KinesisFirehoseClientBuilder {
44 type Client = KinesisClient;
45
46 fn build(&self, config: &aws_types::SdkConfig) -> Self::Client {
47 Self::Client::new(config)
48 }
49}
50
51pub const MAX_PAYLOAD_SIZE: usize = 1024 * 1024 * 4;
54pub const MAX_PAYLOAD_EVENTS: usize = 500;
55
56#[derive(Clone, Copy, Debug, Default)]
57pub struct KinesisFirehoseDefaultBatchSettings;
58
59impl SinkBatchSettings for KinesisFirehoseDefaultBatchSettings {
60 const MAX_EVENTS: Option<usize> = Some(MAX_PAYLOAD_EVENTS);
61 const MAX_BYTES: Option<usize> = Some(MAX_PAYLOAD_SIZE);
62 const TIMEOUT_SECS: f64 = 1.0;
63}
64
65#[configurable_component(sink(
67 "aws_kinesis_firehose",
68 "Publish logs to AWS Kinesis Data Firehose topics."
69))]
70#[derive(Clone, Debug)]
71pub struct KinesisFirehoseSinkConfig {
72 #[serde(flatten)]
73 pub base: KinesisSinkBaseConfig,
74
75 #[configurable(derived)]
76 #[serde(default)]
77 pub batch: BatchConfig<KinesisFirehoseDefaultBatchSettings>,
78}
79
80impl KinesisFirehoseSinkConfig {
81 async fn healthcheck(self, client: KinesisClient) -> crate::Result<()> {
82 let stream_name = self.base.stream_name;
83
84 let result = client
85 .describe_delivery_stream()
86 .delivery_stream_name(stream_name.clone())
87 .set_exclusive_start_destination_id(None)
88 .limit(1)
89 .send()
90 .await;
91
92 match result {
93 Ok(resp) => {
94 let name = resp
95 .delivery_stream_description
96 .map(|x| x.delivery_stream_name)
97 .unwrap_or_default();
98 if name == stream_name {
99 Ok(())
100 } else {
101 Err(HealthcheckError::StreamNamesMismatch { name, stream_name }.into())
102 }
103 }
104 Err(source) => Err(HealthcheckError::DescribeDeliveryStreamFailed { source }.into()),
105 }
106 }
107
108 pub async fn create_client(&self, proxy: &ProxyConfig) -> crate::Result<KinesisClient> {
109 create_client::<KinesisFirehoseClientBuilder>(
110 &KinesisFirehoseClientBuilder {},
111 &self.base.auth,
112 self.base.region.region(),
113 self.base.region.endpoint(),
114 proxy,
115 self.base.tls.as_ref(),
116 None,
117 )
118 .await
119 }
120}
121
122#[async_trait::async_trait]
123#[typetag::serde(name = "aws_kinesis_firehose")]
124impl SinkConfig for KinesisFirehoseSinkConfig {
125 fn input(&self) -> Input {
126 self.base.input()
127 }
128
129 fn acknowledgements(&self) -> &AcknowledgementsConfig {
130 self.base.acknowledgements()
131 }
132
133 fn as_dyn_validated(&self) -> Option<&dyn DynValidatedSink> {
134 Some(self)
135 }
136}
137
138#[derive(Clone, Debug)]
139pub struct ValidatedKinesisFirehose {
140 batch_settings: BatcherSettings,
141}
142
143#[async_trait::async_trait]
144impl ValidatedSink for KinesisFirehoseSinkConfig {
145 type Validated = ValidatedKinesisFirehose;
146
147 fn validate(&self) -> crate::Result<ValidatedKinesisFirehose> {
148 let batch_settings = self
149 .batch
150 .validate()?
151 .limit_max_bytes(MAX_PAYLOAD_SIZE)?
152 .limit_max_events(MAX_PAYLOAD_EVENTS)?
153 .into_batcher_settings()?;
154
155 Ok(ValidatedKinesisFirehose { batch_settings })
156 }
157
158 async fn build(
159 &self,
160 validated: &ValidatedKinesisFirehose,
161 cx: SinkContext,
162 ) -> crate::Result<(VectorSink, Healthcheck)> {
163 let client = self.create_client(&cx.proxy).await?;
164 let healthcheck = self.clone().healthcheck(client.clone()).boxed();
165
166 let sink = build_sink::<
167 KinesisFirehoseClient,
168 KinesisRecord,
169 KinesisFirehoseRecord,
170 KinesisError,
171 KinesisRetryLogic,
172 >(
173 &self.base,
174 self.base.partition_key_field.clone(),
175 validated.batch_settings,
176 KinesisFirehoseClient { client },
177 KinesisRetryLogic {
178 retry_partial: self.base.request_retry_partial,
179 },
180 )?;
181
182 Ok((sink, healthcheck))
183 }
184}
185
186impl GenerateConfig for KinesisFirehoseSinkConfig {
187 fn generate_config() -> serde_json::Value {
188 serde_yaml::from_str(indoc::indoc! {
189 r#"stream_name: my-stream
190 encoding:
191 codec: json"#,
192 })
193 .unwrap()
194 }
195}
196
197#[derive(Clone, Default)]
198struct KinesisRetryLogic {
199 retry_partial: bool,
200}
201
202impl RetryLogic for KinesisRetryLogic {
203 type Error = SdkError<KinesisError, HttpResponse>;
204 type Request = BatchKinesisRequest<KinesisFirehoseRecord>;
205 type Response = KinesisResponse;
206
207 fn is_retriable_error(&self, error: &Self::Error) -> bool {
208 if let SdkError::ServiceError(inner) = error
209 && matches!(
210 inner.err(),
211 PutRecordBatchError::ServiceUnavailableException(_)
212 )
213 {
214 return true;
215 }
216 is_retriable_error(error)
217 }
218
219 fn should_retry_response(&self, response: &Self::Response) -> RetryAction<Self::Request> {
220 if response.failure_count > 0 && self.retry_partial {
221 let msg = format!("partial error count {}", response.failure_count);
222 RetryAction::Retry(msg.into())
223 } else {
224 RetryAction::Successful
225 }
226 }
227}
228
229#[cfg(test)]
230mod tests {
231 use super::*;
232
233 #[test]
234 fn validate_produces_batch_settings() {
235 let config = KinesisFirehoseSinkConfig {
236 batch: BatchConfig::<KinesisFirehoseDefaultBatchSettings>::default(),
237 base: KinesisSinkBaseConfig {
238 stream_name: String::from("test"),
239 region: crate::aws::RegionOrEndpoint::with_both(
240 "us-east-1",
241 "http://localhost:4566",
242 ),
243 encoding: vector_lib::codecs::JsonSerializerConfig::default().into(),
244 compression: crate::sinks::util::Compression::None,
245 request: Default::default(),
246 tls: None,
247 auth: Default::default(),
248 request_retry_partial: false,
249 acknowledgements: Default::default(),
250 partition_key_field: None,
251 },
252 };
253
254 let validated = config.validate().expect("validation should succeed");
255 assert_eq!(validated.batch_settings.item_limit, MAX_PAYLOAD_EVENTS);
256 assert_eq!(validated.batch_settings.size_limit, MAX_PAYLOAD_SIZE);
257 }
258}