Skip to main content

servers/
elasticsearch.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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
52// The headers for every response of Elasticsearch API.
53static 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
66// The fake version of Elasticsearch and used for `_version` API.
67const ELASTICSEARCH_VERSION: &str = "8.16.0";
68
69// Return fake response for Elasticsearch ping request.
70#[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// Return fake response for Elasticsearch license request.
81// Reference: https://www.elastic.co/guide/en/elasticsearch/reference/current/get-license.html.
82#[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/// Process `_bulk` API requests. Only support to create logs.
96/// Reference: https://www.elastic.co/guide/en/elasticsearch/reference/current/docs-bulk.html#docs-bulk-api-request.
97#[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/// Process `/${index}/_bulk` API requests. Only support to create logs.
110/// Reference: https://www.elastic.co/guide/en/elasticsearch/reference/current/docs-bulk.html#docs-bulk-api-request.
111#[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    // The `schema` is already set in the query_ctx in auth process.
139    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    // Record the ingestion time histogram.
146    let _timer = METRIC_ELASTICSEARCH_LOGS_INGESTION_ELAPSED
147        .with_label_values(&[&db])
148        .start_timer();
149
150    // If pipeline_name is not provided, use the internal pipeline.
151    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    // Read the ndjson payload and convert it to a vector of Value.
159    let requests = match parse_bulk_request(&payload, &index, &params.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            // should be unreachable
180            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    // Record the number of documents ingested.
218    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
234// It will generate the following response when write _bulk request to GreptimeDB successfully:
235// {
236//     "took": 1000,
237//     "errors": false,
238//     "items": [
239//         { "create": { "status": 201 } },
240//         { "create": { "status": 201 } },
241//         ...
242//     ]
243// }
244// If the status code is not 201, it will generate the following response:
245// {
246//     "took": 1000,
247//     "errors": true,
248//     "items": [
249//         { "create": { "status": 400, "error": { "type": "illegal_argument_exception", "reason": "<error_reason>" } } }
250//     ]
251// }
252fn 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
280/// Returns the headers for every response of Elasticsearch API.
281pub fn elasticsearch_headers() -> HeaderMap {
282    ELASTICSEARCH_HEADERS.clone()
283}
284
285// Parse the Elasticsearch bulk request and convert it to multiple LogIngestRequests.
286// The input will be Elasticsearch bulk request in NDJSON format.
287// For example, the input will be like this:
288// { "index" : { "_index" : "test", "_id" : "1" } }
289// { "field1" : "value1" }
290// { "index" : { "_index" : "test", "_id" : "2" } }
291// { "field2" : "value2" }
292fn parse_bulk_request(
293    input: &str,
294    index_from_url: &Option<String>,
295    msg_field: &Option<String>,
296) -> ServersResult<Vec<PipelineIngestRequest>> {
297    // Read the ndjson payload and convert it to `Vec<Value>`. Return error if the input is not a valid JSON.
298    let values: Vec<VrlValue> = Deserializer::from_str(input)
299        .into_iter::<VrlValue>()
300        .collect::<Result<_, _>>()
301        .context(ParseJsonSnafu)?;
302
303    // Check if the input is empty.
304    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    // Read the ndjson payload and convert it to a (index, value) vector.
315    // For Elasticsearch post `_bulk` API, each chunk contains two objects:
316    //   1. The first object is the command, it should be `create` or `index`.
317    //   2. The second object is the document data.
318    while let Some(cmd) = values.next() {
319        // NOTE: Although the native Elasticsearch API supports upsert in `index` command, we don't support change any data in `index` command and it's same as `create` command.
320        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        // Read the second object to get the document data. Stop the loop if there is no document.
336        if let Some(document) = values.next() {
337            // If the msg_field is provided, fetch the value of the field from the document data.
338            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
367// Get the index from the command. We will take index as the table name in GreptimeDB.
368fn 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        // If the `_index` exists, it should be a string.
377        InvalidElasticsearchInputSnafu {
378            reason: "index is not a string in bulk request",
379        }
380        .fail()
381    }
382}
383
384// If the msg_field is provided, fetch the value of the field from the document data.
385// For example, if the `msg_field` is `message`, and the document data is `{"message":"hello"}`, the log value will be Value::String("hello").
386fn 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                    // If the message is not a valid JSON, return a map with the original message key and value.
397                    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            // If the message is not a string, just use the original message as the log value.
407            _ => message,
408        }
409    } else {
410        // If the msg_field is not found, just use the original message as the log value.
411        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            // Normal case.
423            (
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            // Case with index.
448            (
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            // Case with index.
473            (
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            // Case with incomplete bulk request.
498            (
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            // Specify the `data` field as the message field and the value is a JSON string.
516            (
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            // Simulate the log data from Logstash.
541            (
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            // With invalid bulk request.
566            (
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}