Skip to main content

vector/sinks/aws_kinesis/firehose/
config.rs

1use 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
51// AWS Kinesis Firehose API accepts payloads up to 4MB or 500 events
52// https://docs.aws.amazon.com/firehose/latest/dev/limits.html
53pub 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/// Configuration for the `aws_kinesis_firehose` sink.
66#[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}