Skip to main content

servers/http/
loki.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, VecDeque};
16use std::sync::Arc;
17use std::time::Instant;
18
19use api::v1::value::ValueData;
20use api::v1::{
21    ColumnDataType, ColumnDataTypeExtension, ColumnSchema, JsonTypeExtension, Row,
22    RowInsertRequest, Rows, SemanticType, Value as GreptimeValue,
23};
24use axum::Extension;
25use axum::extract::State;
26use axum_extra::TypedHeader;
27use bytes::Bytes;
28use chrono::DateTime;
29use common_query::prelude::greptime_timestamp;
30use common_telemetry::{error, warn};
31use headers::ContentType;
32use jsonb::Value;
33use lazy_static::lazy_static;
34use loki_proto::logproto::LabelPairAdapter;
35use loki_proto::prost_types::Timestamp as LokiTimestamp;
36use pipeline::util::to_pipeline_version;
37use pipeline::{
38    ContextReq, GreptimePipelineParams, PipelineContext, PipelineDefinition, SchemaInfo,
39};
40use prost::Message;
41use quoted_string::test_utils::TestSpec;
42use session::context::{Channel, QueryContext, QueryContextRef};
43use snafu::{OptionExt, ResultExt, ensure};
44use snap::raw::Decoder;
45use table::requests::{SEMANTIC_SIGNAL_TYPE, SEMANTIC_SOURCE, SIGNAL_TYPE_LOG, SOURCE_LOKI};
46use vrl::value::{KeyString, Value as VrlValue};
47
48use crate::error::{
49    DecodeLokiRequestSnafu, DecompressSnappyLokiRequestSnafu, InvalidLokiLabelsSnafu,
50    InvalidLokiPayloadSnafu, ParseJsonSnafu, PipelineSnafu, Result, UnsupportedContentTypeSnafu,
51};
52use crate::http::HttpResponse;
53use crate::http::event::{
54    JSON_CONTENT_TYPE, LogState, PB_CONTENT_TYPE, PipelineIngestRequest, execute_log_context_req,
55};
56use crate::http::extractor::{LogTableName, PipelineInfo};
57use crate::metrics::{METRIC_LOKI_LOGS_INGESTION_COUNTER, METRIC_LOKI_LOGS_INGESTION_ELAPSED};
58use crate::pipeline::run_pipeline;
59use crate::query_handler::PipelineHandlerRef;
60
61const LOKI_TABLE_NAME: &str = "loki_logs";
62const LOKI_LINE_COLUMN: &str = "line";
63const LOKI_STRUCTURED_METADATA_COLUMN: &str = "structured_metadata";
64
65const LOKI_LINE_COLUMN_NAME: &str = "loki_line";
66
67const LOKI_PIPELINE_METADATA_PREFIX: &str = "loki_metadata_";
68const LOKI_PIPELINE_LABEL_PREFIX: &str = "loki_label_";
69
70const STREAMS_KEY: &str = "streams";
71const LABEL_KEY: &str = "stream";
72const LINES_KEY: &str = "values";
73
74lazy_static! {
75    static ref LOKI_INIT_SCHEMAS: Vec<ColumnSchema> = vec![
76        ColumnSchema {
77            column_name: greptime_timestamp().to_string(),
78            datatype: ColumnDataType::TimestampNanosecond.into(),
79            semantic_type: SemanticType::Timestamp.into(),
80            datatype_extension: None,
81            options: None,
82        },
83        ColumnSchema {
84            column_name: LOKI_LINE_COLUMN.to_string(),
85            datatype: ColumnDataType::String.into(),
86            semantic_type: SemanticType::Field.into(),
87            datatype_extension: None,
88            options: None,
89        },
90        ColumnSchema {
91            column_name: LOKI_STRUCTURED_METADATA_COLUMN.to_string(),
92            datatype: ColumnDataType::Binary.into(),
93            semantic_type: SemanticType::Field.into(),
94            datatype_extension: Some(ColumnDataTypeExtension {
95                type_ext: Some(api::v1::column_data_type_extension::TypeExt::JsonType(
96                    JsonTypeExtension::JsonBinary.into()
97                ))
98            }),
99            options: None,
100        }
101    ];
102}
103
104#[axum_macros::debug_handler]
105pub async fn loki_ingest(
106    State(log_state): State<LogState>,
107    Extension(mut ctx): Extension<QueryContext>,
108    TypedHeader(content_type): TypedHeader<ContentType>,
109    LogTableName(table_name): LogTableName,
110    pipeline_info: PipelineInfo,
111    Extension(memory_limiter): Extension<crate::request_memory_limiter::ServerMemoryLimiter>,
112    bytes: Bytes,
113) -> Result<HttpResponse> {
114    ctx.set_channel(Channel::Loki);
115    ctx.set_extension(SEMANTIC_SIGNAL_TYPE, SIGNAL_TYPE_LOG);
116    ctx.set_extension(SEMANTIC_SOURCE, SOURCE_LOKI);
117    let ctx = Arc::new(ctx);
118    let table_name = table_name.unwrap_or_else(|| LOKI_TABLE_NAME.to_string());
119    let handler = log_state.log_handler;
120    // Preserve the old elapsed metric boundary: it includes parsing, optional
121    // pipeline execution, and insertion.
122    let exec_timer = Instant::now();
123
124    let ctx_req = build_loki_context_req(
125        &handler,
126        content_type,
127        table_name,
128        pipeline_info,
129        bytes,
130        &ctx,
131        &memory_limiter,
132    )
133    .await?;
134
135    execute_log_context_req(
136        handler,
137        ctx_req,
138        ctx,
139        exec_timer,
140        &METRIC_LOKI_LOGS_INGESTION_COUNTER,
141        &METRIC_LOKI_LOGS_INGESTION_ELAPSED,
142    )
143    .await
144}
145
146/// This is the holder of the loki lines parsed from json or protobuf.
147/// The generic here is either [VrlValue] or [Vec<LabelPairAdapter>].
148/// Depending on the target destination, this can be converted to [LokiRawItem] or [LokiPipeline].
149pub struct LokiMiddleItem<T> {
150    pub ts: i64,
151    pub line: String,
152    pub structured_metadata: Option<T>,
153    pub labels: Option<BTreeMap<String, String>>,
154}
155
156/// This is the line item for the Loki raw ingestion.
157/// We'll persist the line in its whole, set labels into tags,
158/// and structured metadata into a big JSON.
159pub struct LokiRawItem {
160    pub ts: i64,
161    pub line: String,
162    pub structured_metadata: Vec<u8>,
163    pub labels: Option<BTreeMap<String, String>>,
164}
165
166/// This is the line item prepared for the pipeline engine.
167pub struct LokiPipeline {
168    pub map: VrlValue,
169}
170
171struct LokiPipelineContextReq {
172    content_type: ContentType,
173    table_name: String,
174    pipeline_name: String,
175    pipeline_version: Option<String>,
176    pipeline_params: GreptimePipelineParams,
177    bytes: Bytes,
178}
179
180async fn build_loki_context_req(
181    handler: &PipelineHandlerRef,
182    content_type: ContentType,
183    table_name: String,
184    pipeline_info: PipelineInfo,
185    bytes: Bytes,
186    ctx: &QueryContextRef,
187    memory_limiter: &crate::request_memory_limiter::ServerMemoryLimiter,
188) -> Result<ContextReq> {
189    // A pipeline header switches Loki into the generic pipeline path; without
190    // it, Loki writes directly to the target log table.
191    match pipeline_info.pipeline_name {
192        Some(pipeline_name) => {
193            let pipeline_req = LokiPipelineContextReq {
194                content_type,
195                table_name,
196                pipeline_name,
197                pipeline_version: pipeline_info.pipeline_version,
198                pipeline_params: pipeline_info.pipeline_params,
199                bytes,
200            };
201            build_loki_pipeline_context_req(handler, pipeline_req, ctx, memory_limiter).await
202        }
203        None => {
204            let req =
205                build_loki_raw_insert_request(content_type, table_name, bytes, memory_limiter)
206                    .await?;
207            Ok(ContextReq::default_opt_with_reqs(vec![req]))
208        }
209    }
210}
211
212async fn build_loki_pipeline_context_req(
213    handler: &PipelineHandlerRef,
214    pipeline_req: LokiPipelineContextReq,
215    ctx: &QueryContextRef,
216    memory_limiter: &crate::request_memory_limiter::ServerMemoryLimiter,
217) -> Result<ContextReq> {
218    let LokiPipelineContextReq {
219        content_type,
220        table_name,
221        pipeline_name,
222        pipeline_version,
223        pipeline_params,
224        bytes,
225    } = pipeline_req;
226
227    let version = to_pipeline_version(pipeline_version.as_deref()).context(PipelineSnafu)?;
228    let def =
229        PipelineDefinition::from_name(&pipeline_name, version, None).context(PipelineSnafu)?;
230    let pipeline_ctx = PipelineContext::new(&def, &pipeline_params, Channel::Loki);
231
232    let values = extract_item::<LokiPipeline>(content_type, bytes, memory_limiter)
233        .await?
234        .map(|item| item.map)
235        .collect::<Vec<_>>();
236
237    let req = PipelineIngestRequest {
238        table: table_name,
239        values,
240    };
241
242    run_pipeline(handler, &pipeline_ctx, req, ctx, true).await
243}
244
245async fn build_loki_raw_insert_request(
246    content_type: ContentType,
247    table_name: String,
248    bytes: Bytes,
249    memory_limiter: &crate::request_memory_limiter::ServerMemoryLimiter,
250) -> Result<RowInsertRequest> {
251    let mut schema_info = SchemaInfo::from_schema_list(LOKI_INIT_SCHEMAS.clone());
252    let mut rows = Vec::with_capacity(256);
253    let loki_rows = extract_item::<LokiRawItem>(content_type, bytes, memory_limiter).await?;
254    for loki_row in loki_rows {
255        let mut row = init_row(
256            schema_info.schema.len(),
257            loki_row.ts,
258            loki_row.line,
259            loki_row.structured_metadata,
260        );
261        process_labels(&mut schema_info, &mut row, loki_row.labels);
262        rows.push(row);
263    }
264
265    let schemas = schema_info.column_schemas()?;
266    // Labels can introduce new tag columns after earlier rows were built.
267    for row in rows.iter_mut() {
268        row.resize(schemas.len(), GreptimeValue::default());
269    }
270    let rows = Rows {
271        rows: rows.into_iter().map(|values| Row { values }).collect(),
272        schema: schemas,
273    };
274
275    Ok(RowInsertRequest {
276        table_name,
277        rows: Some(rows),
278    })
279}
280
281/// Extract Loki entries from the supported wire format into the caller's
282/// destination type.
283///
284/// JSON push bodies become `LokiMiddleItem<VrlValue>`, protobuf push bodies
285/// become `LokiMiddleItem<Vec<LabelPairAdapter>>`, and the generic `Into<T>`
286/// conversion selects either direct-write `LokiRawItem` or pipeline `LokiPipeline`.
287async fn extract_item<T>(
288    content_type: ContentType,
289    bytes: Bytes,
290    memory_limiter: &crate::request_memory_limiter::ServerMemoryLimiter,
291) -> Result<Box<dyn Iterator<Item = T>>>
292where
293    LokiMiddleItem<VrlValue>: Into<T>,
294    LokiMiddleItem<Vec<LabelPairAdapter>>: Into<T>,
295{
296    match content_type {
297        x if x == *JSON_CONTENT_TYPE => Ok(Box::new(
298            LokiJsonParser::from_bytes(bytes)?.flat_map(|item| item.into_iter().map(|i| i.into())),
299        )),
300        x if x == *PB_CONTENT_TYPE => Ok(Box::new(
301            LokiPbParser::from_bytes(bytes, memory_limiter)
302                .await?
303                .flat_map(|item| item.into_iter().map(|i| i.into())),
304        )),
305        _ => UnsupportedContentTypeSnafu { content_type }.fail(),
306    }
307}
308
309struct LokiJsonParser {
310    pub streams: VecDeque<VrlValue>,
311}
312
313impl LokiJsonParser {
314    pub fn from_bytes(bytes: Bytes) -> Result<Self> {
315        let payload: VrlValue = serde_json::from_slice(bytes.as_ref()).context(ParseJsonSnafu)?;
316
317        let VrlValue::Object(mut map) = payload else {
318            return InvalidLokiPayloadSnafu {
319                msg: "payload is not an object",
320            }
321            .fail();
322        };
323
324        let streams = map.remove(STREAMS_KEY).context(InvalidLokiPayloadSnafu {
325            msg: "missing streams",
326        })?;
327
328        let VrlValue::Array(streams) = streams else {
329            return InvalidLokiPayloadSnafu {
330                msg: "streams is not an array",
331            }
332            .fail();
333        };
334
335        Ok(Self {
336            streams: streams.into(),
337        })
338    }
339}
340
341impl Iterator for LokiJsonParser {
342    type Item = JsonStreamItem;
343
344    fn next(&mut self) -> Option<Self::Item> {
345        while let Some(stream) = self.streams.pop_front() {
346            // get lines from the map
347            let VrlValue::Object(mut map) = stream else {
348                warn!("stream is not an object, {:?}", stream);
349                continue;
350            };
351            let Some(lines) = map.remove(LINES_KEY) else {
352                warn!("missing lines on stream, {:?}", map);
353                continue;
354            };
355            let VrlValue::Array(lines) = lines else {
356                warn!("lines is not an array, {:?}", lines);
357                continue;
358            };
359
360            // get labels
361            let labels = map
362                .remove(LABEL_KEY)
363                .and_then(|m| match m {
364                    VrlValue::Object(labels) => Some(labels),
365                    _ => None,
366                })
367                .map(|m| {
368                    m.into_iter()
369                        .filter_map(|(k, v)| match v {
370                            VrlValue::Bytes(v) => {
371                                Some((k.into(), String::from_utf8_lossy(&v).to_string()))
372                            }
373                            _ => None,
374                        })
375                        .collect::<BTreeMap<String, String>>()
376                });
377
378            return Some(JsonStreamItem {
379                lines: lines.into(),
380                labels,
381            });
382        }
383        None
384    }
385}
386
387struct JsonStreamItem {
388    pub lines: VecDeque<VrlValue>,
389    pub labels: Option<BTreeMap<String, String>>,
390}
391
392impl Iterator for JsonStreamItem {
393    type Item = LokiMiddleItem<VrlValue>;
394
395    fn next(&mut self) -> Option<Self::Item> {
396        while let Some(line) = self.lines.pop_front() {
397            let VrlValue::Array(line) = line else {
398                warn!("line is not an array, {:?}", line);
399                continue;
400            };
401            if line.len() < 2 {
402                warn!("line is too short, {:?}", line);
403                continue;
404            }
405            let mut line: VecDeque<VrlValue> = line.into();
406
407            // get ts
408            let ts = line.pop_front().and_then(|ts| match ts {
409                VrlValue::Bytes(ts) => String::from_utf8_lossy(&ts).parse::<i64>().ok(),
410                _ => {
411                    warn!("missing or invalid timestamp, {:?}", ts);
412                    None
413                }
414            });
415            let Some(ts) = ts else {
416                continue;
417            };
418
419            let line_text = line.pop_front().and_then(|l| match l {
420                VrlValue::Bytes(l) => Some(String::from_utf8_lossy(&l).to_string()),
421                _ => {
422                    warn!("missing or invalid line, {:?}", l);
423                    None
424                }
425            });
426            let Some(line_text) = line_text else {
427                continue;
428            };
429
430            let structured_metadata = line.pop_front();
431
432            return Some(LokiMiddleItem {
433                ts,
434                line: line_text,
435                structured_metadata,
436                labels: self.labels.clone(),
437            });
438        }
439        None
440    }
441}
442
443type LokiPipelineMap = BTreeMap<KeyString, VrlValue>;
444
445fn vrl_metadata_to_jsonb(structured_metadata: Option<VrlValue>) -> Vec<u8> {
446    // JSON push structured metadata arrives as a VRL object.
447    let structured_metadata = structured_metadata
448        .and_then(|metadata| match metadata {
449            VrlValue::Object(metadata) => Some(metadata),
450            _ => None,
451        })
452        .map(|metadata| {
453            metadata
454                .into_iter()
455                .filter_map(|(key, value)| match value {
456                    VrlValue::Bytes(bytes) => Some((
457                        key.into(),
458                        Value::String(String::from_utf8_lossy(&bytes).to_string().into()),
459                    )),
460                    _ => None,
461                })
462                .collect::<BTreeMap<String, Value>>()
463        })
464        .unwrap_or_default();
465
466    Value::Object(structured_metadata).to_vec()
467}
468
469fn label_pair_metadata_to_jsonb(structured_metadata: Option<Vec<LabelPairAdapter>>) -> Vec<u8> {
470    // Protobuf push structured metadata arrives as Loki label pairs.
471    let structured_metadata = structured_metadata
472        .unwrap_or_default()
473        .into_iter()
474        .map(|metadata| (metadata.name, Value::String(metadata.value.into())))
475        .collect::<BTreeMap<String, Value>>();
476
477    Value::Object(structured_metadata).to_vec()
478}
479
480fn new_loki_pipeline_map(ts: i64, line: String) -> LokiPipelineMap {
481    let mut map = BTreeMap::new();
482    map.insert(
483        KeyString::from(greptime_timestamp()),
484        VrlValue::Timestamp(DateTime::from_timestamp_nanos(ts)),
485    );
486    map.insert(
487        KeyString::from(LOKI_LINE_COLUMN_NAME),
488        VrlValue::Bytes(line.into()),
489    );
490    map
491}
492
493fn append_vrl_pipeline_metadata(map: &mut LokiPipelineMap, structured_metadata: Option<VrlValue>) {
494    if let Some(VrlValue::Object(metadata)) = structured_metadata {
495        for (key, value) in metadata {
496            map.insert(
497                KeyString::from(format!("{}{}", LOKI_PIPELINE_METADATA_PREFIX, key)),
498                value,
499            );
500        }
501    }
502}
503
504fn append_label_pair_pipeline_metadata(
505    map: &mut LokiPipelineMap,
506    structured_metadata: Option<Vec<LabelPairAdapter>>,
507) {
508    for metadata in structured_metadata.unwrap_or_default() {
509        map.insert(
510            KeyString::from(format!(
511                "{}{}",
512                LOKI_PIPELINE_METADATA_PREFIX, metadata.name
513            )),
514            VrlValue::Bytes(metadata.value.into()),
515        );
516    }
517}
518
519fn append_pipeline_labels(map: &mut LokiPipelineMap, labels: Option<BTreeMap<String, String>>) {
520    if let Some(labels) = labels {
521        for (key, value) in labels {
522            map.insert(
523                KeyString::from(format!("{}{}", LOKI_PIPELINE_LABEL_PREFIX, key)),
524                VrlValue::Bytes(value.into()),
525            );
526        }
527    }
528}
529
530impl From<LokiMiddleItem<VrlValue>> for LokiRawItem {
531    fn from(val: LokiMiddleItem<VrlValue>) -> Self {
532        let LokiMiddleItem {
533            ts,
534            line,
535            structured_metadata,
536            labels,
537        } = val;
538
539        LokiRawItem {
540            ts,
541            line,
542            structured_metadata: vrl_metadata_to_jsonb(structured_metadata),
543            labels,
544        }
545    }
546}
547
548impl From<LokiMiddleItem<VrlValue>> for LokiPipeline {
549    fn from(value: LokiMiddleItem<VrlValue>) -> Self {
550        let LokiMiddleItem {
551            ts,
552            line,
553            structured_metadata,
554            labels,
555        } = value;
556
557        let mut map = new_loki_pipeline_map(ts, line);
558        append_vrl_pipeline_metadata(&mut map, structured_metadata);
559        append_pipeline_labels(&mut map, labels);
560
561        LokiPipeline {
562            map: VrlValue::Object(map),
563        }
564    }
565}
566
567pub struct LokiPbParser {
568    pub streams: VecDeque<loki_proto::logproto::StreamAdapter>,
569}
570
571impl LokiPbParser {
572    pub async fn from_bytes(
573        bytes: Bytes,
574        memory_limiter: &crate::request_memory_limiter::ServerMemoryLimiter,
575    ) -> Result<Self> {
576        // `decompressed` carries the memory permits for the decoded bytes and
577        // keeps them alive until the protobuf decoding below is finished.
578        let decompressed = snappy_decompress_loki_request(&bytes, memory_limiter).await?;
579        let req = loki_proto::logproto::PushRequest::decode(&decompressed[..])
580            .context(DecodeLokiRequestSnafu)?;
581
582        Ok(Self {
583            streams: req.streams.into(),
584        })
585    }
586}
587
588async fn snappy_decompress_loki_request(
589    buf: &[u8],
590    limiter: &crate::request_memory_limiter::ServerMemoryLimiter,
591) -> Result<crate::prom_store::ChargedBuffer> {
592    // Loki's protobuf push body is Snappy-compressed independent of HTTP
593    // content-encoding, so keep this decode step explicit.
594    //
595    // The decoded size is validated before the output buffer is allocated:
596    // the raw Snappy header carries the decoded length, and without this cap
597    // a tiny body could declare multiple GiB of output. The declared size is
598    // also charged to the aggregate request-memory limiter before the
599    // allocation, so concurrent Loki pushes cannot bypass the quota. The
600    // returned buffer carries the permits, keeping the reservation alive
601    // until the caller is done with the decoded bytes.
602    let decoded_len = snap::raw::decompress_len(buf).context(DecompressSnappyLokiRequestSnafu)?;
603    ensure!(
604        decoded_len <= crate::prom_store::MAX_DECOMPRESSED_REQUEST_SIZE,
605        crate::error::DecompressedBodyTooLargeSnafu {
606            size: decoded_len as u64,
607            limit: crate::prom_store::MAX_DECOMPRESSED_REQUEST_SIZE as u64,
608        }
609    );
610    let guard = limiter.acquire(decoded_len as u64).await?;
611    let mut decoder = Decoder::new();
612    let data = decoder
613        .decompress_vec(buf)
614        .context(DecompressSnappyLokiRequestSnafu)?;
615    Ok(crate::prom_store::ChargedBuffer::new(data, vec![guard]))
616}
617
618impl Iterator for LokiPbParser {
619    type Item = PbStreamItem;
620
621    fn next(&mut self) -> Option<Self::Item> {
622        let stream = self.streams.pop_front()?;
623
624        let labels = parse_loki_labels(&stream.labels)
625            .inspect_err(|e| {
626                error!(e; "failed to parse loki labels, {:?}", stream.labels);
627            })
628            .ok();
629
630        Some(PbStreamItem {
631            entries: stream.entries.into(),
632            labels,
633        })
634    }
635}
636
637pub struct PbStreamItem {
638    pub entries: VecDeque<loki_proto::logproto::EntryAdapter>,
639    pub labels: Option<BTreeMap<String, String>>,
640}
641
642impl Iterator for PbStreamItem {
643    type Item = LokiMiddleItem<Vec<LabelPairAdapter>>;
644
645    fn next(&mut self) -> Option<Self::Item> {
646        while let Some(entry) = self.entries.pop_front() {
647            let ts = if let Some(ts) = entry.timestamp {
648                ts
649            } else {
650                warn!("missing timestamp, {:?}", entry);
651                continue;
652            };
653            let line = entry.line;
654
655            let structured_metadata = entry.structured_metadata;
656
657            return Some(LokiMiddleItem {
658                ts: prost_ts_to_nano(&ts),
659                line,
660                structured_metadata: Some(structured_metadata),
661                labels: self.labels.clone(),
662            });
663        }
664        None
665    }
666}
667
668impl From<LokiMiddleItem<Vec<LabelPairAdapter>>> for LokiRawItem {
669    fn from(val: LokiMiddleItem<Vec<LabelPairAdapter>>) -> Self {
670        let LokiMiddleItem {
671            ts,
672            line,
673            structured_metadata,
674            labels,
675        } = val;
676
677        LokiRawItem {
678            ts,
679            line,
680            structured_metadata: label_pair_metadata_to_jsonb(structured_metadata),
681            labels,
682        }
683    }
684}
685
686impl From<LokiMiddleItem<Vec<LabelPairAdapter>>> for LokiPipeline {
687    fn from(value: LokiMiddleItem<Vec<LabelPairAdapter>>) -> Self {
688        let LokiMiddleItem {
689            ts,
690            line,
691            structured_metadata,
692            labels,
693        } = value;
694
695        let mut map = new_loki_pipeline_map(ts, line);
696        append_label_pair_pipeline_metadata(&mut map, structured_metadata);
697        append_pipeline_labels(&mut map, labels);
698
699        LokiPipeline {
700            map: VrlValue::Object(map),
701        }
702    }
703}
704
705/// since we're hand-parsing the labels, if any error is encountered, we'll just skip the label
706/// note: pub here for bench usage
707/// ref:
708/// 1. encoding: https://github.com/grafana/alloy/blob/be34410b9e841cc0c37c153f9550d9086a304bca/internal/component/common/loki/client/batch.go#L114-L145
709/// 2. test data: https://github.com/grafana/loki/blob/a24ef7b206e0ca63ee74ca6ecb0a09b745cd2258/pkg/push/types_test.go
710pub fn parse_loki_labels(labels: &str) -> Result<BTreeMap<String, String>> {
711    let mut labels = labels.trim();
712    ensure!(
713        labels.len() >= 2,
714        InvalidLokiLabelsSnafu {
715            msg: "labels string too short"
716        }
717    );
718    ensure!(
719        labels.starts_with("{"),
720        InvalidLokiLabelsSnafu {
721            msg: "missing `{` at the beginning"
722        }
723    );
724    ensure!(
725        labels.ends_with("}"),
726        InvalidLokiLabelsSnafu {
727            msg: "missing `}` at the end"
728        }
729    );
730
731    let mut result = BTreeMap::new();
732    labels = &labels[1..labels.len() - 1];
733
734    while !labels.is_empty() {
735        // parse key
736        let first_index = labels.find("=").with_context(|| InvalidLokiLabelsSnafu {
737            msg: format!("missing `=` near: {}", labels),
738        })?;
739        let key = &labels[..first_index];
740        labels = &labels[first_index + 1..];
741
742        // parse value
743        let qs = quoted_string::parse::<TestSpec>(labels)
744            .map_err(|e| {
745                InvalidLokiLabelsSnafu {
746                    msg: format!(
747                        "failed to parse quoted string near: {}, reason: {}",
748                        labels, e.1
749                    ),
750                }
751                .build()
752            })?
753            .quoted_string;
754
755        labels = &labels[qs.len()..];
756
757        let value = quoted_string::to_content::<TestSpec>(qs).map_err(|e| {
758            InvalidLokiLabelsSnafu {
759                msg: format!("failed to unquote the string: {}, reason: {}", qs, e),
760            }
761            .build()
762        })?;
763
764        // insert key and value
765        result.insert(key.to_string(), value.to_string());
766
767        if labels.is_empty() {
768            break;
769        }
770        ensure!(
771            labels.starts_with(","),
772            InvalidLokiLabelsSnafu { msg: "missing `,`" }
773        );
774        labels = labels[1..].trim_start();
775    }
776
777    Ok(result)
778}
779
780#[inline]
781fn prost_ts_to_nano(ts: &LokiTimestamp) -> i64 {
782    ts.seconds * 1_000_000_000 + ts.nanos as i64
783}
784
785fn init_row(
786    schema_len: usize,
787    ts: i64,
788    line: String,
789    structured_metadata: Vec<u8>,
790) -> Vec<GreptimeValue> {
791    // create and init row
792    let mut row = Vec::with_capacity(schema_len);
793    // set ts and line
794    row.push(GreptimeValue {
795        value_data: Some(ValueData::TimestampNanosecondValue(ts)),
796    });
797    row.push(GreptimeValue {
798        value_data: Some(ValueData::StringValue(line)),
799    });
800    row.push(GreptimeValue {
801        value_data: Some(ValueData::BinaryValue(structured_metadata)),
802    });
803    for _ in 0..(schema_len - 3) {
804        row.push(GreptimeValue { value_data: None });
805    }
806    row
807}
808
809fn process_labels(
810    schema_info: &mut SchemaInfo,
811    row: &mut Vec<GreptimeValue>,
812    labels: Option<BTreeMap<String, String>>,
813) {
814    let Some(labels) = labels else {
815        return;
816    };
817
818    let column_indexer = &mut schema_info.index;
819    let schemas = &mut schema_info.schema;
820
821    // insert labels
822    for (k, v) in labels {
823        if let Some(index) = column_indexer.get(&k) {
824            // exist in schema
825            // insert value using index
826            row[*index] = GreptimeValue {
827                value_data: Some(ValueData::StringValue(v)),
828            };
829        } else {
830            // not exist
831            // add schema and append to values
832            schemas.push(
833                ColumnSchema {
834                    column_name: k.clone(),
835                    datatype: ColumnDataType::String.into(),
836                    semantic_type: SemanticType::Tag.into(),
837                    datatype_extension: None,
838                    options: None,
839                }
840                .into(),
841            );
842            column_indexer.insert(k, schemas.len() - 1);
843
844            row.push(GreptimeValue {
845                value_data: Some(ValueData::StringValue(v)),
846            });
847        }
848    }
849}
850
851#[cfg(test)]
852mod tests {
853    use std::collections::BTreeMap;
854
855    use bytes::Bytes;
856    use loki_proto::logproto::{EntryAdapter, PushRequest, StreamAdapter};
857    use loki_proto::prost_types::Timestamp;
858    use prost::Message;
859
860    use super::*;
861    use crate::error::Error::{DecompressSnappyLokiRequest, InvalidLokiLabels};
862    use crate::prom_store::snappy_compress;
863    use crate::request_memory_limiter::ServerMemoryLimiter;
864
865    const JSON_PAYLOAD: &[u8] = br#"{
866        "streams": [
867            {
868                "stream": {
869                    "job": "api",
870                    "namespace": "prod"
871                },
872                "values": [
873                    ["1731748568804293888", "line one", {"trace_id": "abc"}]
874                ]
875            },
876            {
877                "stream": {
878                    "job": "worker",
879                    "pod": "worker-0"
880                },
881                "values": [
882                    ["1731748568804293889", "line two"]
883                ]
884            }
885        ]
886    }"#;
887
888    fn row_string_value(row: &Row, index: usize) -> Option<&str> {
889        match row.values[index].value_data.as_ref() {
890            Some(ValueData::StringValue(value)) => Some(value.as_str()),
891            _ => None,
892        }
893    }
894
895    fn pipeline_bytes_value(map: &BTreeMap<KeyString, VrlValue>, key: &str) -> Option<String> {
896        match map.get(&KeyString::from(key))? {
897            VrlValue::Bytes(value) => Some(String::from_utf8_lossy(value.as_ref()).to_string()),
898            _ => None,
899        }
900    }
901
902    #[test]
903    fn test_ts_to_nano() {
904        // ts = 1731748568804293888
905        // seconds = 1731748568
906        // nano = 804293888
907        let ts = Timestamp {
908            seconds: 1731748568,
909            nanos: 804293888,
910        };
911        assert_eq!(prost_ts_to_nano(&ts), 1731748568804293888);
912    }
913
914    #[tokio::test]
915    async fn test_json_direct_ingest_builds_schema_and_pads_rows() {
916        let request = build_loki_raw_insert_request(
917            JSON_CONTENT_TYPE.clone(),
918            "custom_loki".to_string(),
919            Bytes::from_static(JSON_PAYLOAD),
920            &ServerMemoryLimiter::default(),
921        )
922        .await
923        .unwrap();
924
925        assert_eq!(request.table_name, "custom_loki");
926        let rows = request.rows.unwrap();
927        let column_names = rows
928            .schema
929            .iter()
930            .map(|schema| schema.column_name.as_str())
931            .collect::<Vec<_>>();
932        assert_eq!(
933            column_names,
934            vec![
935                greptime_timestamp(),
936                LOKI_LINE_COLUMN,
937                LOKI_STRUCTURED_METADATA_COLUMN,
938                "job",
939                "namespace",
940                "pod",
941            ]
942        );
943        assert_eq!(rows.schema[3].semantic_type, SemanticType::Tag as i32);
944        assert_eq!(rows.schema[4].semantic_type, SemanticType::Tag as i32);
945        assert_eq!(rows.schema[5].semantic_type, SemanticType::Tag as i32);
946        assert_eq!(rows.rows.len(), 2);
947
948        let first = &rows.rows[0];
949        assert_eq!(first.values.len(), rows.schema.len());
950        assert_eq!(row_string_value(first, 1), Some("line one"));
951        assert_eq!(row_string_value(first, 3), Some("api"));
952        assert_eq!(row_string_value(first, 4), Some("prod"));
953        assert!(first.values[5].value_data.is_none());
954
955        let second = &rows.rows[1];
956        assert_eq!(second.values.len(), rows.schema.len());
957        assert_eq!(row_string_value(second, 1), Some("line two"));
958        assert_eq!(row_string_value(second, 3), Some("worker"));
959        assert!(second.values[4].value_data.is_none());
960        assert_eq!(row_string_value(second, 5), Some("worker-0"));
961    }
962
963    #[tokio::test]
964    async fn test_json_pipeline_conversion_names_loki_fields() {
965        let items = extract_item::<LokiPipeline>(
966            JSON_CONTENT_TYPE.clone(),
967            Bytes::from_static(JSON_PAYLOAD),
968            &ServerMemoryLimiter::default(),
969        )
970        .await
971        .unwrap()
972        .collect::<Vec<_>>();
973
974        assert_eq!(items.len(), 2);
975        let VrlValue::Object(map) = &items[0].map else {
976            panic!("expected pipeline object");
977        };
978        assert!(matches!(
979            map.get(&KeyString::from(greptime_timestamp())),
980            Some(VrlValue::Timestamp(_))
981        ));
982        assert_eq!(
983            pipeline_bytes_value(map, LOKI_LINE_COLUMN_NAME),
984            Some("line one".to_string())
985        );
986        assert_eq!(
987            pipeline_bytes_value(map, "loki_label_job"),
988            Some("api".to_string())
989        );
990        assert_eq!(
991            pipeline_bytes_value(map, "loki_label_namespace"),
992            Some("prod".to_string())
993        );
994        assert_eq!(
995            pipeline_bytes_value(map, "loki_metadata_trace_id"),
996            Some("abc".to_string())
997        );
998    }
999
1000    #[tokio::test]
1001    async fn test_protobuf_parser_decodes_snappy_push_request() {
1002        let request = PushRequest {
1003            streams: vec![StreamAdapter {
1004                labels: r#"{job="api"}"#.to_string(),
1005                entries: vec![EntryAdapter {
1006                    timestamp: Some(Timestamp {
1007                        seconds: 1731748568,
1008                        nanos: 804293888,
1009                    }),
1010                    line: "line one".to_string(),
1011                    structured_metadata: vec![LabelPairAdapter {
1012                        name: "trace_id".to_string(),
1013                        value: "abc".to_string(),
1014                    }],
1015                    ..Default::default()
1016                }],
1017                ..Default::default()
1018            }],
1019        };
1020        let bytes = snappy_compress(&request.encode_to_vec()).unwrap();
1021
1022        let items = extract_item::<LokiRawItem>(
1023            PB_CONTENT_TYPE.clone(),
1024            Bytes::from(bytes),
1025            &ServerMemoryLimiter::default(),
1026        )
1027        .await
1028        .unwrap()
1029        .collect::<Vec<_>>();
1030
1031        assert_eq!(items.len(), 1);
1032        assert_eq!(items[0].ts, 1731748568804293888);
1033        assert_eq!(items[0].line, "line one");
1034        assert_eq!(
1035            items[0].labels.as_ref().unwrap().get("job"),
1036            Some(&"api".to_string())
1037        );
1038        assert!(!items[0].structured_metadata.is_empty());
1039    }
1040
1041    #[tokio::test]
1042    async fn test_protobuf_parser_rejects_invalid_snappy_payload() {
1043        let err = match LokiPbParser::from_bytes(
1044            Bytes::from_static(b"not-snappy"),
1045            &ServerMemoryLimiter::default(),
1046        )
1047        .await
1048        {
1049            Ok(_) => panic!("expected invalid snappy payload to fail"),
1050            Err(err) => err,
1051        };
1052
1053        assert!(matches!(err, DecompressSnappyLokiRequest { .. }));
1054    }
1055
1056    #[tokio::test]
1057    async fn test_protobuf_push_charges_decoded_size() {
1058        use common_memory_manager::OnExhaustedPolicy;
1059
1060        // A valid snappy payload whose decoded size exceeds the 1 KiB quota.
1061        let request = loki_proto::logproto::PushRequest {
1062            streams: vec![loki_proto::logproto::StreamAdapter {
1063                labels: "{job=\"quota\"}".to_string(),
1064                entries: std::iter::repeat_with(|| loki_proto::logproto::EntryAdapter {
1065                    timestamp: Some(Timestamp {
1066                        seconds: 1,
1067                        nanos: 0,
1068                    }),
1069                    line: "x".repeat(256),
1070                    ..Default::default()
1071                })
1072                .take(8)
1073                .collect(),
1074                ..Default::default()
1075            }],
1076        };
1077        let bytes = Bytes::from(snappy_compress(&request.encode_to_vec()).unwrap());
1078
1079        let limiter = ServerMemoryLimiter::new(1024, OnExhaustedPolicy::Fail);
1080        let err = match LokiPbParser::from_bytes(bytes.clone(), &limiter).await {
1081            Ok(_) => panic!("expected quota exhaustion to fail"),
1082            Err(err) => err,
1083        };
1084        assert!(matches!(
1085            err,
1086            crate::error::Error::MemoryLimitExceeded { .. }
1087        ));
1088        assert_eq!(0, limiter.used_bytes(), "failed charge must be released");
1089
1090        // The same payload fits an adequately sized quota.
1091        let limiter = ServerMemoryLimiter::new(64 * 1024, OnExhaustedPolicy::Fail);
1092        let parser = LokiPbParser::from_bytes(bytes, &limiter).await.unwrap();
1093        assert_eq!(1, parser.streams.len());
1094        assert_eq!(
1095            0,
1096            limiter.used_bytes(),
1097            "guards must be released after decode"
1098        );
1099    }
1100
1101    #[test]
1102    fn test_parse_loki_labels() {
1103        let mut expected = BTreeMap::new();
1104        expected.insert("job".to_string(), "foobar".to_string());
1105        expected.insert("cluster".to_string(), "foo-central1".to_string());
1106        expected.insert("namespace".to_string(), "bar".to_string());
1107        expected.insert("container_name".to_string(), "buzz".to_string());
1108
1109        // perfect case
1110        let valid_labels =
1111            r#"{job="foobar", cluster="foo-central1", namespace="bar", container_name="buzz"}"#;
1112        let re = parse_loki_labels(valid_labels);
1113        assert!(re.is_ok());
1114        assert_eq!(re.unwrap(), expected);
1115
1116        // too short
1117        let too_short = r#"}"#;
1118        let re = parse_loki_labels(too_short);
1119        assert!(matches!(re.err().unwrap(), InvalidLokiLabels { .. }));
1120
1121        // missing start
1122        let missing_start = r#"job="foobar"}"#;
1123        let re = parse_loki_labels(missing_start);
1124        assert!(matches!(re.err().unwrap(), InvalidLokiLabels { .. }));
1125
1126        // missing start
1127        let missing_end = r#"{job="foobar""#;
1128        let re = parse_loki_labels(missing_end);
1129        assert!(matches!(re.err().unwrap(), InvalidLokiLabels { .. }));
1130
1131        // missing equal
1132        let missing_equal = r#"{job"foobar"}"#;
1133        let re = parse_loki_labels(missing_equal);
1134        assert!(matches!(re.err().unwrap(), InvalidLokiLabels { .. }));
1135
1136        // missing quote
1137        let missing_quote = r#"{job=foobar}"#;
1138        let re = parse_loki_labels(missing_quote);
1139        assert!(matches!(re.err().unwrap(), InvalidLokiLabels { .. }));
1140
1141        // missing comma
1142        let missing_comma = r#"{job="foobar" cluster="foo-central1"}"#;
1143        let re = parse_loki_labels(missing_comma);
1144        assert!(matches!(re.err().unwrap(), InvalidLokiLabels { .. }));
1145    }
1146}