1use crate::config::{ElasticsearchAuth, ElasticsearchSourceConfig};
4use async_trait::async_trait;
5use faucet_core::util::{DEFAULT_ERROR_BODY_MAX_LEN, check_http_response};
6use faucet_core::{AuthSpec, FaucetError, SharedAuthProvider, Stream, StreamPage};
7use reqwest::Client;
8use serde_json::{Value, json};
9use std::pin::Pin;
10
11pub(crate) const NO_BATCHING_SEARCH_SIZE: usize = 10_000;
15
16pub struct ElasticsearchSource {
18 config: ElasticsearchSourceConfig,
19 client: Client,
20 auth_provider: Option<SharedAuthProvider>,
24}
25
26impl ElasticsearchSource {
27 pub fn new(config: ElasticsearchSourceConfig) -> Result<Self, FaucetError> {
31 config.validate()?;
32 Ok(Self {
33 config,
34 client: Client::new(),
35 auth_provider: None,
36 })
37 }
38
39 pub fn with_auth_provider(mut self, provider: SharedAuthProvider) -> Self {
45 self.auth_provider = Some(provider);
46 self
47 }
48
49 async fn resolve_auth(&self) -> Result<ElasticsearchAuth, FaucetError> {
56 if let Some(p) = &self.auth_provider {
57 return faucet_common_elasticsearch::credential_to_auth(p.credential().await?);
58 }
59 match &self.config.auth {
60 AuthSpec::Inline(a) => Ok(a.clone()),
61 AuthSpec::Reference(r) => Err(FaucetError::Auth(format!(
62 "auth references provider '{}' but no provider was supplied",
63 r.name
64 ))),
65 }
66 }
67
68 fn apply_auth_value(
70 req: reqwest::RequestBuilder,
71 auth: &ElasticsearchAuth,
72 ) -> reqwest::RequestBuilder {
73 match auth {
74 ElasticsearchAuth::None => req,
75 ElasticsearchAuth::Basic { username, password } => {
76 req.basic_auth(username, Some(password))
77 }
78 ElasticsearchAuth::Bearer { token } => req.bearer_auth(token),
79 ElasticsearchAuth::ApiKey { key } => {
80 req.header("Authorization", format!("ApiKey {key}"))
81 }
82 }
83 }
84
85 fn extract_hits(body: &Value) -> Vec<Value> {
87 body.get("hits")
88 .and_then(|h| h.get("hits"))
89 .and_then(|h| h.as_array())
90 .map(|hits| {
91 hits.iter()
92 .filter_map(|hit| hit.get("_source").cloned())
93 .collect()
94 })
95 .unwrap_or_default()
96 }
97
98 fn extract_scroll_id(body: &Value) -> Option<String> {
100 body.get("_scroll_id")
101 .and_then(|v| v.as_str())
102 .map(|s| s.to_string())
103 }
104
105 async fn clear_scroll(&self, scroll_id: &str) {
107 let url = format!("{}/_search/scroll", self.config.base_url);
108 let req = self
109 .client
110 .delete(&url)
111 .json(&json!({"scroll_id": scroll_id}));
112 let auth = match self.resolve_auth().await {
113 Ok(a) => a,
114 Err(e) => {
115 tracing::warn!(error = %e, "failed to resolve auth for scroll cleanup");
116 return;
117 }
118 };
119 let req = Self::apply_auth_value(req, &auth);
120
121 if let Err(e) = req.send().await {
122 tracing::warn!(error = %e, "failed to clear Elasticsearch scroll context");
123 }
124 }
125
126 fn resolve_index_and_query(
129 &self,
130 context: &std::collections::HashMap<String, Value>,
131 ) -> Result<(String, Value), FaucetError> {
132 let index = if context.is_empty() {
133 self.config.index.clone()
134 } else {
135 faucet_core::util::substitute_context(&self.config.index, context)
136 };
137 let query = if context.is_empty() {
138 self.config.query.clone()
139 } else {
140 let s = serde_json::to_string(&self.config.query)
141 .map_err(|e| FaucetError::Config(format!("failed to serialize query: {e}")))?;
142 let s = faucet_core::util::substitute_context_json(&s, context);
143 serde_json::from_str(&s).map_err(|e| {
144 FaucetError::Config(format!("failed to parse substituted query: {e}"))
145 })?
146 };
147 Ok((index, query))
148 }
149}
150
151#[async_trait]
152impl faucet_core::Source for ElasticsearchSource {
153 async fn fetch_with_context(
154 &self,
155 context: &std::collections::HashMap<String, serde_json::Value>,
156 ) -> Result<Vec<Value>, FaucetError> {
157 let (index, query) = self.resolve_index_and_query(context)?;
158 let auth = self.resolve_auth().await?;
160
161 let mut all_records = Vec::new();
162
163 let page_size = if self.config.batch_size == 0 {
167 NO_BATCHING_SEARCH_SIZE
168 } else {
169 self.config.batch_size
170 };
171
172 let url = format!(
174 "{}/{}/_search?scroll={}&size={}",
175 self.config.base_url, index, self.config.scroll_timeout, page_size
176 );
177 let req = self.client.post(&url).json(&json!({"query": query}));
178 let req = Self::apply_auth_value(req, &auth);
179
180 let resp = req.send().await?;
181 let resp = check_http_response(resp, DEFAULT_ERROR_BODY_MAX_LEN).await?;
182 let body: Value = resp.json().await?;
183
184 let mut records = Self::extract_hits(&body);
185 let mut scroll_id = Self::extract_scroll_id(&body);
186 let mut pages_fetched: usize = 1;
187
188 tracing::debug!(
189 records = records.len(),
190 page = pages_fetched,
191 "Elasticsearch initial search"
192 );
193
194 all_records.append(&mut records);
195
196 while let Some(ref sid) = scroll_id {
198 if let Some(max) = self.config.max_pages
200 && pages_fetched >= max
201 {
202 tracing::debug!(max_pages = max, "max_pages reached, stopping scroll");
203 break;
204 }
205
206 let scroll_url = format!("{}/_search/scroll", self.config.base_url);
207 let req = self.client.post(&scroll_url).json(&json!({
208 "scroll": self.config.scroll_timeout,
209 "scroll_id": sid,
210 }));
211 let req = Self::apply_auth_value(req, &auth);
212
213 let resp = req.send().await?;
214 let resp = check_http_response(resp, DEFAULT_ERROR_BODY_MAX_LEN).await?;
215 let body: Value = resp.json().await?;
216
217 let mut page_records = Self::extract_hits(&body);
218 pages_fetched += 1;
219
220 tracing::debug!(
221 records = page_records.len(),
222 page = pages_fetched,
223 "Elasticsearch scroll page"
224 );
225
226 if page_records.is_empty() {
228 break;
229 }
230
231 scroll_id = Self::extract_scroll_id(&body);
233 all_records.append(&mut page_records);
234 }
235
236 if let Some(ref sid) = scroll_id {
238 self.clear_scroll(sid).await;
239 }
240
241 tracing::debug!(
242 total_records = all_records.len(),
243 pages = pages_fetched,
244 "Elasticsearch fetch complete"
245 );
246
247 Ok(all_records)
248 }
249
250 fn stream_pages<'a>(
275 &'a self,
276 context: &'a std::collections::HashMap<String, Value>,
277 _batch_size: usize,
278 ) -> Pin<Box<dyn Stream<Item = Result<StreamPage, FaucetError>> + Send + 'a>> {
279 let batch_size = self.config.batch_size;
280
281 Box::pin(async_stream::try_stream! {
282 let (index, query) = self.resolve_index_and_query(context)?;
283 let auth = self.resolve_auth().await?;
285
286 if batch_size == 0 {
288 let url = format!(
289 "{}/{}/_search?size={}",
290 self.config.base_url, index, NO_BATCHING_SEARCH_SIZE
291 );
292 let req = self.client.post(&url).json(&json!({"query": query}));
293 let req = Self::apply_auth_value(req, &auth);
294 let resp = req.send().await?;
295 let resp = check_http_response(resp, DEFAULT_ERROR_BODY_MAX_LEN).await?;
296 let body: Value = resp.json().await?;
297 let records = Self::extract_hits(&body);
298 tracing::info!(
299 docs = records.len(),
300 batch_size = 0,
301 "Elasticsearch source stream complete (no-batching path)",
302 );
303 yield StreamPage { records, bookmark: None };
304 return;
305 }
306
307 let mut guard = ScrollGuard::new(
312 self.config.base_url.clone(),
313 self.client.clone(),
314 auth.clone(),
315 );
316
317 let url = format!(
318 "{}/{}/_search?scroll={}&size={}",
319 self.config.base_url, index, self.config.scroll_timeout, batch_size
320 );
321 let req = self.client.post(&url).json(&json!({"query": query}));
322 let req = Self::apply_auth_value(req, &auth);
323 let resp = req.send().await?;
324 let resp = check_http_response(resp, DEFAULT_ERROR_BODY_MAX_LEN).await?;
325 let body: Value = resp.json().await?;
326
327 let records = Self::extract_hits(&body);
328 guard.update(Self::extract_scroll_id(&body));
329 let mut pages_emitted: usize = 0;
330 let mut total = records.len();
331
332 pages_emitted += 1;
335 let is_final = records.is_empty()
336 || guard.scroll_id().is_none()
337 || matches!(self.config.max_pages, Some(max) if pages_emitted >= max);
338 yield StreamPage { records, bookmark: None };
339 if is_final {
340 guard.disarm_if_done();
341 tracing::info!(
342 docs = total,
343 pages = pages_emitted,
344 batch_size,
345 "Elasticsearch source stream complete",
346 );
347 return;
348 }
349
350 while let Some(sid) = guard.scroll_id().map(|s| s.to_string()) {
352 let scroll_url = format!("{}/_search/scroll", self.config.base_url);
353 let req = self.client.post(&scroll_url).json(&json!({
354 "scroll": self.config.scroll_timeout,
355 "scroll_id": sid,
356 }));
357 let req = Self::apply_auth_value(req, &auth);
358 let resp = req.send().await?;
359 let resp = check_http_response(resp, DEFAULT_ERROR_BODY_MAX_LEN).await?;
360 let body: Value = resp.json().await?;
361
362 let records = Self::extract_hits(&body);
363 guard.update(Self::extract_scroll_id(&body));
364 pages_emitted += 1;
365 total += records.len();
366
367 let is_empty = records.is_empty();
368 let hit_cap = matches!(self.config.max_pages, Some(max) if pages_emitted >= max);
369
370 if is_empty {
371 break;
374 }
375
376 yield StreamPage { records, bookmark: None };
377
378 if hit_cap {
379 tracing::debug!(
380 max_pages = self.config.max_pages.unwrap_or(0),
381 "max_pages reached, stopping scroll"
382 );
383 break;
384 }
385 }
386
387 tracing::info!(
388 docs = total,
389 pages = pages_emitted,
390 batch_size,
391 "Elasticsearch source stream complete",
392 );
393
394 guard.disarm_if_done();
396 })
397 }
398
399 fn connector_name(&self) -> &'static str {
400 "elasticsearch"
401 }
402
403 fn config_schema(&self) -> serde_json::Value {
404 serde_json::to_value(faucet_core::schema_for!(ElasticsearchSourceConfig))
405 .expect("schema serialization")
406 }
407
408 fn dataset_uri(&self) -> String {
409 format!(
410 "{}/{}",
411 faucet_core::redact_uri_credentials(&self.config.base_url),
412 self.config.index
413 )
414 }
415
416 fn supports_discover(&self) -> bool {
417 true
418 }
419
420 async fn discover(&self) -> Result<Vec<faucet_core::DatasetDescriptor>, FaucetError> {
425 let auth = self.resolve_auth().await?;
427
428 let url = format!(
429 "{}/_cat/indices?format=json&h=index,docs.count",
430 self.config.base_url
431 );
432 let req = Self::apply_auth_value(self.client.get(&url), &auth);
433 let resp = req.send().await.map_err(|e| {
434 FaucetError::Source(format!("elasticsearch: catalog discovery failed: {e}"))
435 })?;
436 let resp = check_http_response(resp, DEFAULT_ERROR_BODY_MAX_LEN).await?;
437 let cat: Value = resp.json().await.map_err(|e| {
438 FaucetError::Source(format!("elasticsearch: catalog discovery failed: {e}"))
439 })?;
440
441 let entries = parse_cat_indices(&cat);
442 let mut datasets = Vec::with_capacity(entries.len());
443 for (index, doc_count) in entries {
444 let url = format!("{}/{}/_mapping", self.config.base_url, index);
445 let req = Self::apply_auth_value(self.client.get(&url), &auth);
446 let resp = req.send().await.map_err(|e| {
447 FaucetError::Source(format!(
448 "elasticsearch: catalog discovery failed (mapping for {index:?}): {e}"
449 ))
450 })?;
451 let resp = check_http_response(resp, DEFAULT_ERROR_BODY_MAX_LEN).await?;
452 let body: Value = resp.json().await.map_err(|e| {
453 FaucetError::Source(format!(
454 "elasticsearch: catalog discovery failed (mapping for {index:?}): {e}"
455 ))
456 })?;
457 datasets.push(descriptor_for_index(&index, doc_count, &body));
458 }
459 Ok(datasets)
460 }
461}
462
463fn es_type_to_json_type(es_type: &str) -> &'static str {
468 match es_type {
469 "long" | "integer" | "short" | "byte" => "integer",
470 "double" | "float" | "half_float" | "scaled_float" => "number",
471 "boolean" => "boolean",
472 "object" | "nested" => "object",
473 _ => "string",
474 }
475}
476
477fn mapping_to_schema(mappings: &Value) -> Value {
483 let mut properties = serde_json::Map::new();
484 if let Some(fields) = mappings.get("properties").and_then(Value::as_object) {
485 for (name, spec) in fields {
486 let ty = match spec.get("type").and_then(Value::as_str) {
487 Some(t) => es_type_to_json_type(t),
488 None => "object",
489 };
490 properties.insert(name.clone(), json!({ "type": ty }));
491 }
492 }
493 json!({ "type": "object", "properties": Value::Object(properties) })
494}
495
496fn parse_cat_indices(cat: &Value) -> Vec<(String, Option<u64>)> {
501 let mut entries: Vec<(String, Option<u64>)> = cat
502 .as_array()
503 .map(|rows| {
504 rows.iter()
505 .filter_map(|row| {
506 let index = row.get("index")?.as_str()?;
507 if index.starts_with('.') {
508 return None;
509 }
510 let doc_count = row.get("docs.count").and_then(|v| {
511 v.as_str()
512 .and_then(|s| s.parse::<u64>().ok())
513 .or_else(|| v.as_u64())
514 });
515 Some((index.to_string(), doc_count))
516 })
517 .collect()
518 })
519 .unwrap_or_default();
520 entries.sort_by(|a, b| a.0.cmp(&b.0));
521 entries
522}
523
524fn descriptor_for_index(
530 index: &str,
531 doc_count: Option<u64>,
532 mapping_body: &Value,
533) -> faucet_core::DatasetDescriptor {
534 let empty = json!({});
535 let mappings = mapping_body
536 .get(index)
537 .or_else(|| mapping_body.as_object().and_then(|o| o.values().next()))
538 .and_then(|entry| entry.get("mappings"))
539 .unwrap_or(&empty);
540 let mut descriptor =
541 faucet_core::DatasetDescriptor::new(index, "index", json!({ "index": index }))
542 .with_schema(mapping_to_schema(mappings));
543 if let Some(rows) = doc_count {
544 descriptor = descriptor.with_estimated_rows(rows);
545 }
546 descriptor
547}
548
549struct ScrollGuard {
554 base_url: String,
555 client: Client,
556 auth: ElasticsearchAuth,
557 scroll_id: Option<String>,
558}
559
560impl ScrollGuard {
561 fn new(base_url: String, client: Client, auth: ElasticsearchAuth) -> Self {
562 Self {
563 base_url,
564 client,
565 auth,
566 scroll_id: None,
567 }
568 }
569
570 fn scroll_id(&self) -> Option<&str> {
571 self.scroll_id.as_deref()
572 }
573
574 fn update(&mut self, new_id: Option<String>) {
575 if let Some(id) = new_id {
576 self.scroll_id = Some(id);
577 }
578 }
579
580 fn disarm_if_done(&mut self) {
583 if let Some(sid) = self.scroll_id.take() {
584 let base_url = self.base_url.clone();
585 let auth = self.auth.clone();
586 let client = self.client.clone();
587 tokio::spawn(async move {
588 let url = format!("{base_url}/_search/scroll");
589 let req = client.delete(&url).json(&json!({"scroll_id": sid}));
590 let req = apply_auth_to(req, &auth);
591 if let Err(e) = req.send().await {
592 tracing::warn!(error = %e, "failed to clear Elasticsearch scroll context");
593 }
594 });
595 }
596 }
597}
598
599impl Drop for ScrollGuard {
600 fn drop(&mut self) {
601 if let Some(sid) = self.scroll_id.take() {
602 let base_url = self.base_url.clone();
605 let auth = self.auth.clone();
606 let client = self.client.clone();
607 tokio::spawn(async move {
608 let url = format!("{base_url}/_search/scroll");
609 let req = client.delete(&url).json(&json!({"scroll_id": sid}));
610 let req = apply_auth_to(req, &auth);
611 if let Err(e) = req.send().await {
612 tracing::warn!(
613 error = %e,
614 "failed to clear Elasticsearch scroll context (drop path)",
615 );
616 }
617 });
618 }
619 }
620}
621
622fn apply_auth_to(
625 req: reqwest::RequestBuilder,
626 auth: &ElasticsearchAuth,
627) -> reqwest::RequestBuilder {
628 match auth {
629 ElasticsearchAuth::None => req,
630 ElasticsearchAuth::Basic { username, password } => req.basic_auth(username, Some(password)),
631 ElasticsearchAuth::Bearer { token } => req.bearer_auth(token),
632 ElasticsearchAuth::ApiKey { key } => req.header("Authorization", format!("ApiKey {key}")),
633 }
634}
635
636#[cfg(test)]
637mod tests {
638 use super::*;
639 use faucet_core::Source;
640
641 #[test]
642 fn new_rejects_out_of_range_batch_size() {
643 let mut config = ElasticsearchSourceConfig::new("http://localhost:9200", "idx");
644 config.batch_size = faucet_core::MAX_BATCH_SIZE + 1;
645 match ElasticsearchSource::new(config) {
646 Err(FaucetError::Config(m)) => assert!(m.contains("batch_size"), "got: {m}"),
647 _ => panic!("expected a batch_size Config error"),
648 }
649 }
650
651 #[test]
652 fn dataset_uri_returns_base_url_slash_index() {
653 let config = ElasticsearchSourceConfig::new("http://localhost:9200", "my_index");
654 let source = ElasticsearchSource::new(config).unwrap();
655 assert_eq!(source.dataset_uri(), "http://localhost:9200/my_index");
656 }
657
658 #[test]
659 fn dataset_uri_strips_credentials() {
660 let config =
661 ElasticsearchSourceConfig::new("http://user:secret@es.example.com:9200", "logs");
662 let source = ElasticsearchSource::new(config).unwrap();
663 assert_eq!(source.dataset_uri(), "http://es.example.com:9200/logs");
664 }
665
666 #[test]
669 fn es_types_map_to_json_types() {
670 for (es, want) in [
671 ("long", "integer"),
672 ("integer", "integer"),
673 ("short", "integer"),
674 ("byte", "integer"),
675 ("double", "number"),
676 ("float", "number"),
677 ("half_float", "number"),
678 ("scaled_float", "number"),
679 ("boolean", "boolean"),
680 ("object", "object"),
681 ("nested", "object"),
682 ("text", "string"),
683 ("keyword", "string"),
684 ("date", "string"),
685 ("ip", "string"),
686 ("geo_point", "string"),
687 ] {
688 assert_eq!(es_type_to_json_type(es), want, "for ES type {es:?}");
689 }
690 }
691
692 #[test]
693 fn mapping_to_schema_covers_scalar_object_and_nested_fields() {
694 let mappings = json!({
697 "properties": {
698 "id": {"type": "long"},
699 "total": {"type": "scaled_float", "scaling_factor": 100},
700 "note": {"type": "text", "fields": {"keyword": {"type": "keyword"}}},
701 "active": {"type": "boolean"},
702 "customer": {"properties": {"name": {"type": "text"}}},
703 "meta": {"type": "nested", "properties": {"k": {"type": "keyword"}}},
704 }
705 });
706 let schema = mapping_to_schema(&mappings);
707 assert_eq!(schema["type"], "object");
708 let props = &schema["properties"];
709 assert_eq!(props["id"]["type"], "integer");
710 assert_eq!(props["total"]["type"], "number");
711 assert_eq!(props["note"]["type"], "string");
712 assert_eq!(props["active"]["type"], "boolean");
713 assert_eq!(
714 props["customer"]["type"], "object",
715 "type-less field with nested properties is an object"
716 );
717 assert_eq!(props["meta"]["type"], "object");
718 }
719
720 #[test]
721 fn mapping_to_schema_empty_mappings_yield_empty_properties() {
722 assert_eq!(
723 mapping_to_schema(&json!({})),
724 json!({"type": "object", "properties": {}})
725 );
726 assert_eq!(
727 mapping_to_schema(&Value::Null),
728 json!({"type": "object", "properties": {}})
729 );
730 }
731
732 #[test]
733 fn parse_cat_indices_skips_system_and_parses_counts() {
734 let cat = json!([
735 {"index": "orders", "docs.count": "1200"},
736 {"index": ".kibana_1", "docs.count": "3"},
737 {"index": "logs", "docs.count": "n/a"},
738 {"index": "metrics"},
739 {"index": "numeric", "docs.count": 7},
740 ]);
741 let entries = parse_cat_indices(&cat);
742 assert_eq!(
743 entries,
744 vec![
745 ("logs".to_string(), None),
746 ("metrics".to_string(), None),
747 ("numeric".to_string(), Some(7)),
748 ("orders".to_string(), Some(1200)),
749 ],
750 "system index skipped, unparsable/missing counts → None, sorted"
751 );
752 }
753
754 #[test]
755 fn parse_cat_indices_non_array_is_empty() {
756 assert!(parse_cat_indices(&json!({"error": "nope"})).is_empty());
757 assert!(parse_cat_indices(&Value::Null).is_empty());
758 }
759
760 #[test]
761 fn descriptor_for_index_builds_full_descriptor() {
762 let body = json!({
763 "orders": {"mappings": {"properties": {"id": {"type": "long"}}}}
764 });
765 let d = descriptor_for_index("orders", Some(1200), &body);
766 assert_eq!(d.name, "orders");
767 assert_eq!(d.kind, "index");
768 assert_eq!(d.config_patch, json!({"index": "orders"}));
769 assert_eq!(d.estimated_rows, Some(1200));
770 let schema = d.schema.as_ref().expect("schema");
771 assert_eq!(schema["properties"]["id"]["type"], "integer");
772 }
773
774 #[test]
775 fn descriptor_for_index_falls_back_to_first_mapping_entry() {
776 let body = json!({
779 "orders-000001": {"mappings": {"properties": {"id": {"type": "long"}}}}
780 });
781 let d = descriptor_for_index("orders", None, &body);
782 assert_eq!(d.estimated_rows, None);
783 let schema = d.schema.as_ref().expect("schema");
784 assert_eq!(schema["properties"]["id"]["type"], "integer");
785 }
786
787 #[test]
788 fn descriptor_for_index_missing_mappings_yields_empty_schema() {
789 let d = descriptor_for_index("orders", Some(0), &json!({}));
790 assert_eq!(
791 d.schema,
792 Some(json!({"type": "object", "properties": {}})),
793 "no mappings → empty object schema, never a panic"
794 );
795 }
796
797 #[test]
798 fn source_advertises_discover() {
799 let config = ElasticsearchSourceConfig::new("http://localhost:9200", "idx");
800 let source = ElasticsearchSource::new(config).unwrap();
801 assert!(source.supports_discover());
802 }
803}