vector/sinks/http/
sink.rs1use 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 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 .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 .request_builder(
54 default_request_builder_concurrency_limit(),
55 self.request_builder,
56 )
57 .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 .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}