1use 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 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
146pub 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
156pub 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
166pub 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 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 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
281async 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 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 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 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 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 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 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 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
705pub 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 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 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 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 let mut row = Vec::with_capacity(schema_len);
793 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 for (k, v) in labels {
823 if let Some(index) = column_indexer.get(&k) {
824 row[*index] = GreptimeValue {
827 value_data: Some(ValueData::StringValue(v)),
828 };
829 } else {
830 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 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 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 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 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 let too_short = r#"}"#;
1118 let re = parse_loki_labels(too_short);
1119 assert!(matches!(re.err().unwrap(), InvalidLokiLabels { .. }));
1120
1121 let missing_start = r#"job="foobar"}"#;
1123 let re = parse_loki_labels(missing_start);
1124 assert!(matches!(re.err().unwrap(), InvalidLokiLabels { .. }));
1125
1126 let missing_end = r#"{job="foobar""#;
1128 let re = parse_loki_labels(missing_end);
1129 assert!(matches!(re.err().unwrap(), InvalidLokiLabels { .. }));
1130
1131 let missing_equal = r#"{job"foobar"}"#;
1133 let re = parse_loki_labels(missing_equal);
1134 assert!(matches!(re.err().unwrap(), InvalidLokiLabels { .. }));
1135
1136 let missing_quote = r#"{job=foobar}"#;
1138 let re = parse_loki_labels(missing_quote);
1139 assert!(matches!(re.err().unwrap(), InvalidLokiLabels { .. }));
1140
1141 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}