Skip to main content

vector/sinks/http/
sink.rs

1//! Implementation of the `http` sink.
2
3use std::collections::BTreeMap;
4
5use super::{batch::HttpBatchSizer, request_builder::HttpRequestBuilder};
6use crate::sinks::{prelude::*, util::http::HttpRequest};
7
8pub(super) struct HttpSink<S> {
9    service: S,
10    uri: Template,
11    headers: BTreeMap<String, Template>,
12    batch_settings: BatcherSettings,
13    request_builder: HttpRequestBuilder,
14}
15
16impl<S> HttpSink<S>
17where
18    S: Service<HttpRequest<PartitionKey>> + Send + 'static,
19    S::Future: Send + 'static,
20    S::Response: DriverResponse + Send + 'static,
21    S::Error: std::fmt::Debug + Into<crate::Error> + Send,
22{
23    /// Creates a new `HttpSink`.
24    pub(super) const fn new(
25        service: S,
26        uri: Template,
27        headers: BTreeMap<String, Template>,
28        batch_settings: BatcherSettings,
29        request_builder: HttpRequestBuilder,
30    ) -> Self {
31        Self {
32            service,
33            uri,
34            headers,
35            batch_settings,
36            request_builder,
37        }
38    }
39
40    async fn run_inner(self: Box<Self>, input: BoxStream<'_, Event>) -> Result<(), ()> {
41        let batch_sizer = HttpBatchSizer {
42            encoder: self.request_builder.encoder.encoder.clone(),
43        };
44        input
45            // Batch the input stream with size calculation based on the configured codec
46            .batched_partitioned(
47                KeyPartitioner::new(self.uri, self.headers),
48                self.batch_settings.timeout,
49                |_| self.batch_settings.as_item_size_config(batch_sizer.clone()),
50            )
51            .filter_map(|(key, batch)| async move { key.map(move |k| (k, batch)) })
52            // Build requests with default concurrency limit.
53            .request_builder(
54                default_request_builder_concurrency_limit(),
55                self.request_builder,
56            )
57            // Filter out any errors that occurred in the request building.
58            .filter_map(|request| async move {
59                match request {
60                    Err(error) => {
61                        emit!(SinkRequestBuildError { error });
62                        None
63                    }
64                    Ok(req) => Some(req),
65                }
66            })
67            // Generate the driver that will send requests and handle retries,
68            // event finalization, and logging/internal metric reporting.
69            .into_driver(self.service)
70            .run()
71            .await
72    }
73}
74
75#[async_trait::async_trait]
76impl<S> StreamSink<Event> for HttpSink<S>
77where
78    S: Service<HttpRequest<PartitionKey>> + Send + 'static,
79    S::Future: Send + 'static,
80    S::Response: DriverResponse + Send + 'static,
81    S::Error: std::fmt::Debug + Into<crate::Error> + Send,
82{
83    async fn run(
84        self: Box<Self>,
85        input: futures_util::stream::BoxStream<'_, Event>,
86    ) -> Result<(), ()> {
87        self.run_inner(input).await
88    }
89}
90
91#[derive(Eq, PartialEq, Clone, Debug, Hash)]
92pub struct PartitionKey {
93    pub uri: String,
94    pub headers: BTreeMap<String, String>,
95}
96
97struct KeyPartitioner {
98    uri: Template,
99    headers: BTreeMap<String, Template>,
100}
101
102impl KeyPartitioner {
103    const fn new(uri: Template, headers: BTreeMap<String, Template>) -> Self {
104        Self { uri, headers }
105    }
106}
107
108impl Partitioner for KeyPartitioner {
109    type Item = Event;
110    type Key = Option<PartitionKey>;
111
112    fn partition(&self, event: &Event) -> Self::Key {
113        let uri = self
114            .uri
115            .render_string(event)
116            .map_err(|error| {
117                emit!(TemplateRenderingError {
118                    error,
119                    field: Some("uri"),
120                    drop_event: true,
121                });
122            })
123            .ok()?;
124
125        let mut headers = BTreeMap::new();
126        for (name, template) in &self.headers {
127            let value = template
128                .render_string(event)
129                .map_err(|error| {
130                    emit!(TemplateRenderingError {
131                        error,
132                        field: Some("headers"),
133                        drop_event: true,
134                    });
135                })
136                .ok()?;
137            headers.insert(name.clone(), value);
138        }
139
140        Some(PartitionKey { uri, headers })
141    }
142}