1use std::collections::BTreeMap;
16use std::sync::Arc;
17use std::time::Instant;
18
19use axum::Extension;
20use axum::extract::{Path, Query, State};
21use axum::http::{HeaderMap, HeaderName, HeaderValue, StatusCode};
22use axum::response::IntoResponse;
23use axum_extra::TypedHeader;
24use common_error::ext::ErrorExt;
25use common_telemetry::{debug, error};
26use headers::ContentType;
27use once_cell::sync::Lazy;
28use pipeline::{
29 GREPTIME_INTERNAL_IDENTITY_PIPELINE_NAME, GreptimePipelineParams, PipelineDefinition,
30};
31use serde_json::{Deserializer, Value, json};
32use session::context::{Channel, QueryContext};
33use snafu::{ResultExt, ensure};
34use table::requests::{
35 SEMANTIC_SIGNAL_TYPE, SEMANTIC_SOURCE, SIGNAL_TYPE_LOG, SOURCE_ELASTICSEARCH,
36};
37use vrl::value::Value as VrlValue;
38
39use crate::error::{
40 InvalidElasticsearchInputSnafu, ParseJsonSnafu, Result as ServersResult,
41 status_code_to_http_status,
42};
43use crate::http::event::{
44 LogIngesterQueryParams, LogState, PipelineIngestRequest,
45 extract_pipeline_params_map_from_headers, ingest_logs_inner,
46};
47use crate::http::header::constants::GREPTIME_PIPELINE_NAME_HEADER_NAME;
48use crate::metrics::{
49 METRIC_ELASTICSEARCH_LOGS_DOCS_COUNT, METRIC_ELASTICSEARCH_LOGS_INGESTION_ELAPSED,
50};
51
52static ELASTICSEARCH_HEADERS: Lazy<HeaderMap> = Lazy::new(|| {
54 HeaderMap::from_iter([
55 (
56 axum::http::header::CONTENT_TYPE,
57 HeaderValue::from_static("application/json"),
58 ),
59 (
60 HeaderName::from_static("x-elastic-product"),
61 HeaderValue::from_static("Elasticsearch"),
62 ),
63 ])
64});
65
66const ELASTICSEARCH_VERSION: &str = "8.16.0";
68
69#[axum_macros::debug_handler]
71pub async fn handle_get_version() -> impl IntoResponse {
72 let body = serde_json::json!({
73 "version": {
74 "number": ELASTICSEARCH_VERSION
75 }
76 });
77 (StatusCode::OK, elasticsearch_headers(), axum::Json(body))
78}
79
80#[axum_macros::debug_handler]
83pub async fn handle_get_license() -> impl IntoResponse {
84 let body = serde_json::json!({
85 "license": {
86 "uid": "cbff45e7-c553-41f7-ae4f-9205eabd80xx",
87 "type": "oss",
88 "status": "active",
89 "expiry_date_in_millis": 4891198687000_i64,
90 }
91 });
92 (StatusCode::OK, elasticsearch_headers(), axum::Json(body))
93}
94
95#[axum_macros::debug_handler]
98pub async fn handle_bulk_api(
99 State(log_state): State<LogState>,
100 Query(params): Query<LogIngesterQueryParams>,
101 Extension(query_ctx): Extension<QueryContext>,
102 TypedHeader(_content_type): TypedHeader<ContentType>,
103 headers: HeaderMap,
104 payload: String,
105) -> impl IntoResponse {
106 do_handle_bulk_api(log_state, None, params, query_ctx, headers, payload).await
107}
108
109#[axum_macros::debug_handler]
112pub async fn handle_bulk_api_with_index(
113 State(log_state): State<LogState>,
114 Path(index): Path<String>,
115 Query(params): Query<LogIngesterQueryParams>,
116 Extension(query_ctx): Extension<QueryContext>,
117 TypedHeader(_content_type): TypedHeader<ContentType>,
118 headers: HeaderMap,
119 payload: String,
120) -> impl IntoResponse {
121 do_handle_bulk_api(log_state, Some(index), params, query_ctx, headers, payload).await
122}
123
124async fn do_handle_bulk_api(
125 log_state: LogState,
126 index: Option<String>,
127 params: LogIngesterQueryParams,
128 mut query_ctx: QueryContext,
129 headers: HeaderMap,
130 payload: String,
131) -> impl IntoResponse {
132 let start = Instant::now();
133 debug!(
134 "Received bulk request, params: {:?}, payload: {:?}",
135 params, payload
136 );
137
138 query_ctx.set_channel(Channel::Elasticsearch);
140 query_ctx.set_extension(SEMANTIC_SIGNAL_TYPE, SIGNAL_TYPE_LOG);
141 query_ctx.set_extension(SEMANTIC_SOURCE, SOURCE_ELASTICSEARCH);
142
143 let db = query_ctx.current_schema();
144
145 let _timer = METRIC_ELASTICSEARCH_LOGS_INGESTION_ELAPSED
147 .with_label_values(&[&db])
148 .start_timer();
149
150 let pipeline_name = params.pipeline_name.as_deref().unwrap_or_else(|| {
152 headers
153 .get(GREPTIME_PIPELINE_NAME_HEADER_NAME)
154 .and_then(|v| v.to_str().ok())
155 .unwrap_or(GREPTIME_INTERNAL_IDENTITY_PIPELINE_NAME)
156 });
157
158 let requests = match parse_bulk_request(&payload, &index, ¶ms.msg_field) {
160 Ok(requests) => requests,
161 Err(e) => {
162 return (
163 StatusCode::BAD_REQUEST,
164 elasticsearch_headers(),
165 axum::Json(write_bulk_response(
166 start.elapsed().as_millis() as i64,
167 0,
168 StatusCode::BAD_REQUEST.as_u16() as u32,
169 e.to_string().as_str(),
170 )),
171 );
172 }
173 };
174 let log_num = requests.len();
175
176 let pipeline = match PipelineDefinition::from_name(pipeline_name, None, None) {
177 Ok(pipeline) => pipeline,
178 Err(e) => {
179 error!(e; "Failed to ingest logs");
181 return (
182 status_code_to_http_status(&e.status_code()),
183 elasticsearch_headers(),
184 axum::Json(write_bulk_response(
185 start.elapsed().as_millis() as i64,
186 0,
187 e.status_code() as u32,
188 e.to_string().as_str(),
189 )),
190 );
191 }
192 };
193 let pipeline_params =
194 GreptimePipelineParams::from_map(extract_pipeline_params_map_from_headers(&headers));
195 if let Err(e) = ingest_logs_inner(
196 log_state.log_handler,
197 pipeline,
198 requests,
199 Arc::new(query_ctx),
200 pipeline_params,
201 )
202 .await
203 {
204 error!(e; "Failed to ingest logs");
205 return (
206 status_code_to_http_status(&e.status_code()),
207 elasticsearch_headers(),
208 axum::Json(write_bulk_response(
209 start.elapsed().as_millis() as i64,
210 0,
211 e.status_code() as u32,
212 e.to_string().as_str(),
213 )),
214 );
215 }
216
217 METRIC_ELASTICSEARCH_LOGS_DOCS_COUNT
219 .with_label_values(&[&db])
220 .inc_by(log_num as u64);
221
222 (
223 StatusCode::OK,
224 elasticsearch_headers(),
225 axum::Json(write_bulk_response(
226 start.elapsed().as_millis() as i64,
227 log_num,
228 StatusCode::CREATED.as_u16() as u32,
229 "",
230 )),
231 )
232}
233
234fn write_bulk_response(took_ms: i64, n: usize, status_code: u32, error_reason: &str) -> Value {
253 if error_reason.is_empty() {
254 let items: Vec<Value> = (0..n)
255 .map(|_| {
256 json!({
257 "create": {
258 "status": status_code
259 }
260 })
261 })
262 .collect();
263 let mut response = json!({
264 "took": took_ms,
265 "errors": false,
266 });
267 response["items"] = Value::Array(items);
268 response
269 } else {
270 json!({
271 "took": took_ms,
272 "errors": true,
273 "items": [
274 { "create": { "status": status_code, "error": { "type": "illegal_argument_exception", "reason": error_reason } } }
275 ]
276 })
277 }
278}
279
280pub fn elasticsearch_headers() -> HeaderMap {
282 ELASTICSEARCH_HEADERS.clone()
283}
284
285fn parse_bulk_request(
293 input: &str,
294 index_from_url: &Option<String>,
295 msg_field: &Option<String>,
296) -> ServersResult<Vec<PipelineIngestRequest>> {
297 let values: Vec<VrlValue> = Deserializer::from_str(input)
299 .into_iter::<VrlValue>()
300 .collect::<Result<_, _>>()
301 .context(ParseJsonSnafu)?;
302
303 ensure!(
305 !values.is_empty(),
306 InvalidElasticsearchInputSnafu {
307 reason: "empty bulk request".to_string(),
308 }
309 );
310
311 let mut requests: Vec<PipelineIngestRequest> = Vec::with_capacity(values.len() / 2);
312 let mut values = values.into_iter();
313
314 while let Some(cmd) = values.next() {
319 let mut cmd = cmd.into_object();
321 let index = if let Some(cmd) = cmd.as_mut().and_then(|c| c.remove("create")) {
322 get_index_from_cmd(cmd)?
323 } else if let Some(cmd) = cmd.as_mut().and_then(|c| c.remove("index")) {
324 get_index_from_cmd(cmd)?
325 } else {
326 return InvalidElasticsearchInputSnafu {
327 reason: format!(
328 "invalid bulk request, expected 'create' or 'index' but got {:?}",
329 cmd
330 ),
331 }
332 .fail();
333 };
334
335 if let Some(document) = values.next() {
337 let log_value = if let Some(msg_field) = msg_field {
339 get_log_value_from_msg_field(document, msg_field)
340 } else {
341 document
342 };
343
344 ensure!(
345 index.is_some() || index_from_url.is_some(),
346 InvalidElasticsearchInputSnafu {
347 reason: "missing index in bulk request".to_string(),
348 }
349 );
350
351 requests.push(PipelineIngestRequest {
352 table: index.unwrap_or_else(|| index_from_url.as_ref().unwrap().clone()),
353 values: vec![log_value],
354 });
355 }
356 }
357
358 debug!(
359 "Received {} log ingest requests: {:?}",
360 requests.len(),
361 requests
362 );
363
364 Ok(requests)
365}
366
367fn get_index_from_cmd(v: VrlValue) -> ServersResult<Option<String>> {
369 let Some(index) = v.into_object().and_then(|mut m| m.remove("_index")) else {
370 return Ok(None);
371 };
372
373 if let VrlValue::Bytes(index) = index {
374 Ok(Some(String::from_utf8_lossy(&index).to_string()))
375 } else {
376 InvalidElasticsearchInputSnafu {
378 reason: "index is not a string in bulk request",
379 }
380 .fail()
381 }
382}
383
384fn get_log_value_from_msg_field(v: VrlValue, msg_field: &str) -> VrlValue {
387 let VrlValue::Object(mut m) = v else {
388 return v;
389 };
390
391 if let Some(message) = m.remove(msg_field) {
392 match message {
393 VrlValue::Bytes(bytes) => {
394 match serde_json::from_slice::<VrlValue>(&bytes) {
395 Ok(v) => v,
396 Err(_) => {
398 let map = BTreeMap::from([(
399 msg_field.to_string().into(),
400 VrlValue::Bytes(bytes),
401 )]);
402 VrlValue::Object(map)
403 }
404 }
405 }
406 _ => message,
408 }
409 } else {
410 VrlValue::Object(m)
412 }
413}
414
415#[cfg(test)]
416mod tests {
417 use super::*;
418
419 #[test]
420 fn test_parse_bulk_request() {
421 let test_cases = vec![
422 (
424 r#"
425 {"create":{"_index":"test","_id":"1"}}
426 {"foo1":"foo1_value", "bar1":"bar1_value"}
427 {"create":{"_index":"test","_id":"2"}}
428 {"foo2":"foo2_value","bar2":"bar2_value"}
429 "#,
430 None,
431 None,
432 Ok(vec![
433 PipelineIngestRequest {
434 table: "test".to_string(),
435 values: vec![
436 json!({"foo1": "foo1_value", "bar1": "bar1_value"}).into(),
437 ],
438 },
439 PipelineIngestRequest {
440 table: "test".to_string(),
441 values: vec![
442 json!({"foo2": "foo2_value", "bar2": "bar2_value"}).into(),
443 ],
444 },
445 ]),
446 ),
447 (
449 r#"
450 {"create":{"_index":"test","_id":"1"}}
451 {"foo1":"foo1_value", "bar1":"bar1_value"}
452 {"create":{"_index":"logs","_id":"2"}}
453 {"foo2":"foo2_value","bar2":"bar2_value"}
454 "#,
455 Some("logs".to_string()),
456 None,
457 Ok(vec![
458 PipelineIngestRequest {
459 table: "test".to_string(),
460 values: vec![
461 json!({"foo1": "foo1_value", "bar1": "bar1_value"}).into(),
462 ],
463 },
464 PipelineIngestRequest {
465 table: "logs".to_string(),
466 values: vec![
467 json!({"foo2": "foo2_value", "bar2": "bar2_value"}).into(),
468 ],
469 },
470 ]),
471 ),
472 (
474 r#"
475 {"create":{"_index":"test","_id":"1"}}
476 {"foo1":"foo1_value", "bar1":"bar1_value"}
477 {"create":{"_index":"logs","_id":"2"}}
478 {"foo2":"foo2_value","bar2":"bar2_value"}
479 "#,
480 Some("logs".to_string()),
481 None,
482 Ok(vec![
483 PipelineIngestRequest {
484 table: "test".to_string(),
485 values: vec![
486 json!({"foo1": "foo1_value", "bar1": "bar1_value"}).into(),
487 ],
488 },
489 PipelineIngestRequest {
490 table: "logs".to_string(),
491 values: vec![
492 json!({"foo2": "foo2_value", "bar2": "bar2_value"}).into(),
493 ],
494 },
495 ]),
496 ),
497 (
499 r#"
500 {"create":{"_index":"test","_id":"1"}}
501 {"foo1":"foo1_value", "bar1":"bar1_value"}
502 {"create":{"_index":"logs","_id":"2"}}
503 "#,
504 Some("logs".to_string()),
505 None,
506 Ok(vec![
507 PipelineIngestRequest {
508 table: "test".to_string(),
509 values: vec![
510 json!({"foo1": "foo1_value", "bar1": "bar1_value"}).into(),
511 ],
512 },
513 ]),
514 ),
515 (
517 r#"
518 {"create":{"_index":"test","_id":"1"}}
519 {"data":"{\"foo1\":\"foo1_value\", \"bar1\":\"bar1_value\"}", "not_data":"not_data_value"}
520 {"create":{"_index":"test","_id":"2"}}
521 {"data":"{\"foo2\":\"foo2_value\", \"bar2\":\"bar2_value\"}", "not_data":"not_data_value"}
522 "#,
523 None,
524 Some("data".to_string()),
525 Ok(vec![
526 PipelineIngestRequest {
527 table: "test".to_string(),
528 values: vec![
529 json!({"foo1": "foo1_value", "bar1": "bar1_value"}).into(),
530 ],
531 },
532 PipelineIngestRequest {
533 table: "test".to_string(),
534 values: vec![
535 json!({"foo2": "foo2_value", "bar2": "bar2_value"}).into(),
536 ],
537 },
538 ]),
539 ),
540 (
542 r#"
543 {"create":{"_id":null,"_index":"logs-generic-default","routing":null}}
544 {"message":"172.16.0.1 - - [25/May/2024:20:19:37 +0000] \"GET /contact HTTP/1.1\" 404 162 \"-\" \"Mozilla/5.0 (iPhone; CPU iPhone OS 14_0 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/14.0 Mobile/15E148 Safari/604.1\"","@timestamp":"2025-01-04T04:32:13.868962186Z","event":{"original":"172.16.0.1 - - [25/May/2024:20:19:37 +0000] \"GET /contact HTTP/1.1\" 404 162 \"-\" \"Mozilla/5.0 (iPhone; CPU iPhone OS 14_0 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/14.0 Mobile/15E148 Safari/604.1\""},"host":{"name":"orbstack"},"log":{"file":{"path":"/var/log/nginx/access.log"}},"@version":"1","data_stream":{"type":"logs","dataset":"generic","namespace":"default"}}
545 {"create":{"_id":null,"_index":"logs-generic-default","routing":null}}
546 {"message":"10.0.0.1 - - [25/May/2024:20:18:37 +0000] \"GET /images/logo.png HTTP/1.1\" 304 0 \"-\" \"Mozilla/5.0 (X11; Ubuntu; Linux x86_64; rv:89.0) Gecko/20100101 Firefox/89.0\"","@timestamp":"2025-01-04T04:32:13.868723810Z","event":{"original":"10.0.0.1 - - [25/May/2024:20:18:37 +0000] \"GET /images/logo.png HTTP/1.1\" 304 0 \"-\" \"Mozilla/5.0 (X11; Ubuntu; Linux x86_64; rv:89.0) Gecko/20100101 Firefox/89.0\""},"host":{"name":"orbstack"},"log":{"file":{"path":"/var/log/nginx/access.log"}},"@version":"1","data_stream":{"type":"logs","dataset":"generic","namespace":"default"}}
547 "#,
548 None,
549 Some("message".to_string()),
550 Ok(vec![
551 PipelineIngestRequest {
552 table: "logs-generic-default".to_string(),
553 values: vec![
554 json!({"message": "172.16.0.1 - - [25/May/2024:20:19:37 +0000] \"GET /contact HTTP/1.1\" 404 162 \"-\" \"Mozilla/5.0 (iPhone; CPU iPhone OS 14_0 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/14.0 Mobile/15E148 Safari/604.1\""}).into(),
555 ],
556 },
557 PipelineIngestRequest {
558 table: "logs-generic-default".to_string(),
559 values: vec![
560 json!({"message": "10.0.0.1 - - [25/May/2024:20:18:37 +0000] \"GET /images/logo.png HTTP/1.1\" 304 0 \"-\" \"Mozilla/5.0 (X11; Ubuntu; Linux x86_64; rv:89.0) Gecko/20100101 Firefox/89.0\""}).into(),
561 ],
562 },
563 ]),
564 ),
565 (
567 r#"
568 { "not_create_or_index" : { "_index" : "test", "_id" : "1" } }
569 { "foo1" : "foo1_value", "bar1" : "bar1_value" }
570 "#,
571 None,
572 None,
573 Err(InvalidElasticsearchInputSnafu {
574 reason: "it's a invalid bulk request".to_string(),
575 }),
576 ),
577 ];
578
579 for (input, index, msg_field, expected) in test_cases {
580 let requests = parse_bulk_request(input, &index, &msg_field);
581 if let Ok(expected) = expected {
582 assert_eq!(requests.unwrap(), expected);
583 } else {
584 assert!(requests.is_err());
585 }
586 }
587 }
588}