1use super::sink::DorisSink;
4
5use crate::{
6 codecs::EncodingConfigWithFraming,
7 http::{Auth, HttpClient},
8 sinks::{
9 doris::{
10 client::DorisSinkClient, common::DorisCommon, health::DorisHealthLogic,
11 retry::DorisRetryLogic, service::DorisService,
12 },
13 prelude::*,
14 util::{RealtimeSizeBasedDefaultBatchSettings, UriSerde, service::HealthConfig},
15 },
16 template::ConfinementConfig,
17};
18use futures;
19use futures_util::TryFutureExt;
20use std::collections::HashMap;
21use std::sync::Arc;
22
23#[configurable_component(sink("doris", "Deliver log data to an Apache Doris database."))]
25#[derive(Clone, Debug)]
26#[serde(deny_unknown_fields)]
27pub struct DorisConfig {
28 #[serde(default)]
33 #[configurable(metadata(docs::examples = "http://127.0.0.1:8030"))]
34 pub endpoints: Vec<UriSerde>,
35
36 #[configurable(metadata(docs::examples = "mydatabase"))]
38 pub database: Template,
39
40 #[configurable(metadata(docs::examples = "mytable"))]
42 pub table: Template,
43
44 #[configurable(metadata(docs::examples = "vector"))]
47 #[serde(default = "default_label_prefix")]
48 pub label_prefix: String,
49
50 #[serde(default, skip_serializing_if = "crate::serde::is_default")]
52 pub log_request: bool,
53
54 #[serde(default)]
65 #[configurable(metadata(docs::additional_props_description = "An HTTP header value."))]
66 pub headers: HashMap<String, String>,
67
68 #[serde(flatten)]
69 pub encoding: EncodingConfigWithFraming,
70
71 #[serde(default)]
73 pub compression: Compression,
74
75 #[serde(default = "default_max_retries")]
77 pub max_retries: isize,
78
79 #[configurable(derived)]
80 #[serde(default)]
81 pub batch: BatchConfig<RealtimeSizeBasedDefaultBatchSettings>,
82
83 #[configurable(derived)]
84 pub auth: Option<Auth>,
85
86 #[serde(default)]
87 #[configurable(derived)]
88 pub request: TowerRequestConfig,
89
90 #[configurable(derived)]
91 pub tls: Option<TlsConfig>,
92
93 #[serde(default)]
95 #[configurable(derived)]
96 #[serde(rename = "distribution")]
97 pub endpoint_health: Option<HealthConfig>,
98
99 #[configurable(derived)]
100 #[serde(
101 default,
102 deserialize_with = "crate::serde::bool_or_struct",
103 skip_serializing_if = "crate::serde::is_default"
104 )]
105 pub acknowledgements: AcknowledgementsConfig,
106
107 #[configurable(derived)]
108 #[serde(flatten)]
109 pub confinement: ConfinementConfig,
110}
111
112fn default_label_prefix() -> String {
113 "vector".to_string()
114}
115
116const fn default_max_retries() -> isize {
117 -1
118}
119
120impl Default for DorisConfig {
121 fn default() -> Self {
122 Self {
123 endpoints: Vec::new(),
124 database: Template::try_from("").unwrap(),
125 table: Template::try_from("").unwrap(),
126 label_prefix: default_label_prefix(),
127 log_request: false,
128 headers: HashMap::new(),
129 encoding: (
130 Some(vector_lib::codecs::encoding::FramingConfig::NewlineDelimited),
131 vector_lib::codecs::JsonSerializerConfig::default(),
132 )
133 .into(),
134 compression: Compression::default(),
135 max_retries: default_max_retries(),
136 batch: BatchConfig::default(),
137 auth: None,
138 request: TowerRequestConfig::default(),
139 tls: None,
140 endpoint_health: None,
141 acknowledgements: AcknowledgementsConfig::default(),
142 confinement: ConfinementConfig::default(),
143 }
144 }
145}
146
147impl_generate_config_from_default!(DorisConfig);
148
149#[async_trait::async_trait]
150#[typetag::serde(name = "doris")]
151impl SinkConfig for DorisConfig {
152 async fn build(&self, cx: SinkContext) -> crate::Result<(VectorSink, Healthcheck)> {
153 let endpoints = self.endpoints.clone();
154
155 if endpoints.is_empty() {
156 return Err("No endpoints configured.'.".into());
157 }
158 let commons = DorisCommon::parse_many(self).await?;
159 let common = commons[0].clone();
160
161 let client = HttpClient::new(common.tls_settings.clone(), &cx.proxy)?;
162
163 let request_settings = self.request.into_settings();
165
166 let health_config = self.endpoint_health.clone().unwrap_or_default();
167
168 let services_futures = commons
169 .iter()
170 .map(|common| {
171 let client_clone = client.clone();
172 let compression = self.compression;
173 let label_prefix = self.label_prefix.clone();
174 let headers = self.headers.clone();
175 let log_request = self.log_request;
176 let base_url = common.base_url.clone();
177 let auth = common.auth.clone();
178
179 async move {
180 let endpoint = base_url.to_string();
181
182 let doris_client = DorisSinkClient::new(
183 client_clone,
184 base_url,
185 auth,
186 compression,
187 label_prefix,
188 headers,
189 )
190 .await;
191
192 let doris_client_safe = doris_client.into_thread_safe();
193
194 let service = DorisService::new(doris_client_safe, log_request);
195
196 Ok::<_, crate::Error>((endpoint, service))
197 }
198 })
199 .collect::<Vec<_>>();
200
201 let services_results = futures::future::join_all(services_futures).await;
203
204 let services = services_results
206 .into_iter()
207 .filter_map(Result::ok)
208 .collect::<Vec<_>>();
209
210 let service = request_settings.distributed_service(
211 DorisRetryLogic {},
212 services,
213 health_config,
214 DorisHealthLogic,
215 1, );
217
218 let sink = DorisSink::new(service, self, &common)?;
220
221 let sink = VectorSink::from_event_streamsink(sink);
222
223 let healthcheck_doris_client = {
225 let doris_client = DorisSinkClient::new(
226 client.clone(),
227 common.base_url.clone(),
228 common.auth.clone(),
229 self.compression,
230 self.label_prefix.clone(),
231 self.headers.clone(),
232 )
233 .await;
234 doris_client.into_thread_safe()
235 };
236
237 let healthcheck = futures::future::select_ok(commons.into_iter().map(move |common| {
239 let client = Arc::clone(&healthcheck_doris_client);
240 async move { common.healthcheck(client).await }.boxed()
241 }))
242 .map_ok(|((), _)| ())
243 .boxed();
244
245 self.confinement.set_confinement_gauge("sink", Self::NAME);
246 Ok((sink, healthcheck))
247 }
248
249 fn input(&self) -> Input {
250 Input::log()
251 }
252
253 fn acknowledgements(&self) -> &AcknowledgementsConfig {
254 &self.acknowledgements
255 }
256}
257
258#[cfg(test)]
259mod tests {
260 use super::*;
261
262 #[test]
263 fn generate_config() {
264 crate::test_util::test_generate_config::<DorisConfig>();
265 }
266
267 #[test]
268 fn test_default_values() {
269 assert_eq!(default_label_prefix(), "vector");
270 assert_eq!(default_max_retries(), -1);
271 }
272
273 #[test]
274 fn confinement_rejects_unconfined_database_template() {
275 let template = Template::try_from("{{ tenant }}").unwrap();
276 let config = ConfinementConfig::default();
277 let result = template.confine(&config, "doris", "database");
278 assert!(
279 result.is_err(),
280 "bare template with no literal prefix must be rejected"
281 );
282 }
283
284 #[test]
285 fn confinement_opt_out_allows_unconfined_database_template() {
286 let template = Template::try_from("{{ tenant }}").unwrap();
287 let config = ConfinementConfig {
288 dangerously_allow_unconfined_template_resolution: true,
289 };
290 let result = template.confine(&config, "doris", "database");
291 assert!(result.is_ok(), "opt-out must allow bare template");
292 }
293
294 #[test]
295 fn confinement_blocks_dotdot_escape_at_render() {
296 use crate::event::LogEvent;
297 use vrl::event_path;
298 let template = Template::try_from("mydb_{{ tenant }}").unwrap();
299 let config = ConfinementConfig::default();
300 let confined = template.confine(&config, "doris", "database").unwrap();
301 let mut event = LogEvent::default();
302 event.insert(event_path!("tenant"), "/../evil");
303 let result = confined.render_string(&crate::event::Event::Log(event));
304 assert!(
305 result.is_err(),
306 "dotdot escape in rendered value must be rejected by prefix check"
307 );
308 }
309}