1use std::collections::hash_map::Entry;
16
17use ahash::{HashMap, HashMapExt, HashSet, HashSetExt};
18#[cfg(test)]
19use api::greptime_proto::io::prometheus::write::v2::BucketSpan;
20#[cfg(test)]
21use api::greptime_proto::io::prometheus::write::v2::histogram::{Count, ZeroCount};
22use api::greptime_proto::io::prometheus::write::v2::{
23 Exemplar, Histogram, Metadata, Sample, metadata,
24};
25#[cfg(test)]
26use api::greptime_proto::io::prometheus::write::v2::{Request, TimeSeries};
27#[cfg(test)]
28use api::v1::ColumnSchema;
29use api::v1::value::ValueData;
30use api::v1::{ColumnDataType, RowInsertRequest, Rows, SemanticType, Value};
31use bytes::{Buf, Bytes};
32use common_grpc::precision::Precision;
33#[cfg(test)]
34use common_query::native_histogram::*;
35use common_query::native_histogram::{
36 NATIVE_HISTOGRAM_FIELD, encode_native_histogram, native_histogram_column_schema,
37};
38use common_query::prelude::{greptime_native_histogram, greptime_timestamp, greptime_value};
39use pipeline::{ContextOpt, ContextReq};
40use prost::encoding::{
41 DecodeContext, WireType, decode_key, decode_varint, message, skip_field, uint32,
42};
43use prost::{DecodeError, Message};
44use snafu::{OptionExt, ResultExt, ensure};
45use table::requests::{
46 METADATA_QUALITY_DECLARED, SEMANTIC_METRIC_METADATA_QUALITY, SEMANTIC_METRIC_TYPE,
47 SEMANTIC_METRIC_UNIT,
48};
49
50use crate::error::{self, Result};
51use crate::prom_remote_write::row_builder::PromCtx;
52use crate::prom_remote_write::validation::validate_label_name;
53use crate::prom_remote_write::{REMOTE_WRITE_V2_VERSION, decompress_remote_write_body};
54#[allow(deprecated)]
55use crate::prom_store::{
56 DATABASE_LABEL, DATABASE_LABEL_ALT, METRIC_NAME_LABEL, PHYSICAL_TABLE_LABEL,
57 PHYSICAL_TABLE_LABEL_ALT, SCHEMA_LABEL,
58};
59use crate::request_memory_limiter::ServerMemoryLimiter;
60use crate::row_writer::{self, TableData};
61use crate::semantic::{
62 METRIC_TYPE_COUNTER, METRIC_TYPE_GAUGE, METRIC_TYPE_GAUGE_HISTOGRAM, METRIC_TYPE_HISTOGRAM,
63 METRIC_TYPE_INFO, METRIC_TYPE_STATESET, METRIC_TYPE_SUMMARY, SemanticIndexes,
64 openmetrics_unit_to_ucum,
65};
66
67type PromTags<'a> = Vec<(&'a str, String)>;
68type ResolvedSeriesLabels<'a> = (PromCtx, String, PromTags<'a>);
69const TIME_SERIES_LABELS_REFS_TAG: u32 = 1;
70const TIME_SERIES_SAMPLES_TAG: u32 = 2;
71const TIME_SERIES_HISTOGRAMS_TAG: u32 = 3;
72const TIME_SERIES_EXEMPLARS_TAG: u32 = 4;
73const TIME_SERIES_METADATA_TAG: u32 = 5;
74
75struct BorrowedRequest<'a> {
76 symbols: Vec<&'a str>,
77 timeseries: Vec<&'a [u8]>,
78}
79
80impl<'a> BorrowedRequest<'a> {
81 fn decode(mut buf: &'a [u8]) -> std::result::Result<Self, DecodeError> {
82 let mut symbols = Vec::new();
83 let mut timeseries = Vec::new();
84
85 while buf.has_remaining() {
86 let (tag, wire_type) = decode_key(&mut buf)?;
87 match tag {
88 4 => {
89 let value =
90 take_length_delimited(wire_type, &mut buf).map_err(|mut error| {
91 error.push("Request", "symbols");
92 error
93 })?;
94 let symbol = std::str::from_utf8(value).map_err(|_| {
95 let mut error =
96 DecodeError::new("invalid string value: data is not UTF-8 encoded");
97 error.push("Request", "symbols");
98 error
99 })?;
100 symbols.push(symbol);
101 }
102 5 => {
103 let series =
104 take_length_delimited(wire_type, &mut buf).map_err(|mut error| {
105 error.push("Request", "timeseries");
106 error
107 })?;
108 timeseries.push(series);
109 }
110 _ => skip_field(wire_type, tag, &mut buf, DecodeContext::default())?,
111 }
112 }
113
114 Ok(Self {
115 symbols,
116 timeseries,
117 })
118 }
119}
120
121pub(crate) struct RemoteWriteV2WriteRequests {
122 pub samples: ContextReq,
123 pub histograms: ContextReq,
124 pub sample_count: u64,
125 pub histogram_count: u64,
126 pub semantic_index: SemanticIndexes,
129}
130
131pub(crate) async fn decode_remote_write_v2(
132 is_zstd: bool,
133 body: Bytes,
134 limiter: &ServerMemoryLimiter,
135) -> Result<RemoteWriteV2WriteRequests> {
136 let decode_timer = crate::metrics::METRIC_HTTP_PROM_STORE_CODEC_ELAPSED
137 .with_label_values(&["decode", REMOTE_WRITE_V2_VERSION])
138 .start_timer();
139
140 let buf = decompress_remote_write_body(is_zstd, &body[..], limiter).await?;
143 drop(body);
145 let request = BorrowedRequest::decode(&buf).context(error::DecodePromRemoteRequestSnafu)?;
146 drop(decode_timer);
147
148 let _convert_timer = crate::metrics::METRIC_HTTP_PROM_STORE_CODEC_ELAPSED
149 .with_label_values(&["convert", REMOTE_WRITE_V2_VERSION])
150 .start_timer();
151 convert_remote_write_v2(request)
152}
153
154fn convert_remote_write_v2(request: BorrowedRequest<'_>) -> Result<RemoteWriteV2WriteRequests> {
155 ensure!(
156 request.symbols.first().copied() == Some(""),
157 error::InvalidPromRemoteRequestSnafu {
158 msg: "remote write v2 symbols must start with an empty string".to_string(),
159 }
160 );
161
162 let mut sample_tables = HashMap::<PromCtx, HashMap<String, TableData>>::new();
163 let mut histogram_tables = HashMap::<PromCtx, HashMap<String, TableData>>::new();
164 let mut label_names = HashSet::new();
165 let mut sample_count_total = 0;
166 let mut histogram_count_total = 0;
167 let mut labels_refs = Vec::new();
168 let mut metadata = Metadata::default();
169 let mut scratch = LeafScratch::default();
170 let mut semantic_index = SemanticIndexes::default();
171
172 for series in request.timeseries {
173 let counts = scan_series(series, &mut labels_refs, &mut metadata)
174 .context(error::DecodePromRemoteRequestSnafu)?;
175
176 if counts.samples == 0 && counts.histograms == 0 {
177 decode_series_leaves(series, None, Vec::new(), 0, &mut scratch)?;
178 continue;
179 }
180
181 let (prom_ctx, table_name, tags) =
182 resolve_series_labels(&request.symbols, &labels_refs, &mut label_names)?;
183 ensure_no_internal_histogram_labels(&tags)?;
184 record_series_metadata(
185 &mut semantic_index,
186 &request.symbols,
187 &metadata,
188 &prom_ctx,
189 &table_name,
190 )?;
191 let has_other_value_type = if counts.samples > 0 {
192 counts.histograms > 0
193 || histogram_tables
194 .get(&prom_ctx)
195 .is_some_and(|tables| tables.contains_key(&table_name))
196 } else {
197 sample_tables
198 .get(&prom_ctx)
199 .is_some_and(|tables| tables.contains_key(&table_name))
200 };
201 ensure!(
202 !has_other_value_type,
203 error::InvalidPromRemoteRequestSnafu {
204 msg: format!(
205 "remote write v2 metric `{table_name}` contains both samples and native histograms"
206 ),
207 }
208 );
209
210 let column_count =
211 tags.len()
212 .checked_add(2)
213 .with_context(|| error::InvalidPromRemoteRequestSnafu {
214 msg: "remote write v2 series has too many labels".to_string(),
215 })?;
216 let (writer, row_count) = if counts.samples > 0 {
217 (
218 SeriesWriter::Samples(get_or_create_table_data(
219 &mut sample_tables,
220 prom_ctx,
221 table_name,
222 column_count,
223 counts.samples,
224 )),
225 counts.samples,
226 )
227 } else {
228 (
229 SeriesWriter::Histograms(get_or_create_table_data(
230 &mut histogram_tables,
231 prom_ctx,
232 table_name,
233 column_count,
234 counts.histograms,
235 )),
236 counts.histograms,
237 )
238 };
239
240 decode_series_leaves(series, Some(writer), tags, row_count, &mut scratch)?;
241 sample_count_total = checked_total(sample_count_total, counts.samples, "sample")?;
242 histogram_count_total =
243 checked_total(histogram_count_total, counts.histograms, "histogram")?;
244 }
245
246 Ok(RemoteWriteV2WriteRequests {
247 samples: into_context_req(sample_tables),
248 histograms: into_context_req(histogram_tables),
249 sample_count: sample_count_total,
250 histogram_count: histogram_count_total,
251 semantic_index,
252 })
253}
254
255fn record_series_metadata(
265 index: &mut SemanticIndexes,
266 symbols: &[&str],
267 series_metadata: &Metadata,
268 prom_ctx: &PromCtx,
269 table_name: &str,
270) -> Result<()> {
271 if series_metadata.help_ref != 0 {
273 symbol_ref(symbols, series_metadata.help_ref, "metadata help")?;
274 }
275 let unit = if series_metadata.unit_ref != 0 {
276 Some(symbol_ref(
277 symbols,
278 series_metadata.unit_ref,
279 "metadata unit",
280 )?)
281 } else {
282 None
283 };
284 let metric_type = metric_type_value(series_metadata.r#type);
285 let ucum = unit.and_then(|unit| openmetrics_unit_to_ucum(unit.trim()));
286 if metric_type.is_none() && ucum.is_none() {
287 return Ok(());
288 }
289
290 let index = index.index_for(prom_ctx.schema.as_deref());
291 if let Some(metric_type) = metric_type {
292 index.record_scalar(table_name, SEMANTIC_METRIC_TYPE, metric_type);
293 index.record_scalar(
294 table_name,
295 SEMANTIC_METRIC_METADATA_QUALITY,
296 METADATA_QUALITY_DECLARED,
297 );
298 }
299 if let Some(ucum) = ucum {
300 index.record_scalar(table_name, SEMANTIC_METRIC_UNIT, ucum);
301 }
302 Ok(())
303}
304
305fn metric_type_value(wire_type: i32) -> Option<&'static str> {
308 match metadata::MetricType::try_from(wire_type).ok()? {
309 metadata::MetricType::Unspecified => None,
310 metadata::MetricType::Counter => Some(METRIC_TYPE_COUNTER),
311 metadata::MetricType::Gauge => Some(METRIC_TYPE_GAUGE),
312 metadata::MetricType::Histogram => Some(METRIC_TYPE_HISTOGRAM),
313 metadata::MetricType::Gaugehistogram => Some(METRIC_TYPE_GAUGE_HISTOGRAM),
314 metadata::MetricType::Summary => Some(METRIC_TYPE_SUMMARY),
315 metadata::MetricType::Info => Some(METRIC_TYPE_INFO),
316 metadata::MetricType::Stateset => Some(METRIC_TYPE_STATESET),
317 }
318}
319
320#[derive(Default)]
321struct SeriesCounts {
322 samples: usize,
323 histograms: usize,
324}
325
326fn scan_series(
327 mut buf: &[u8],
328 labels_refs: &mut Vec<u32>,
329 metadata: &mut Metadata,
330) -> std::result::Result<SeriesCounts, DecodeError> {
331 labels_refs.clear();
332 metadata.clear();
333 let mut counts = SeriesCounts::default();
334
335 while buf.has_remaining() {
336 let (tag, wire_type) = decode_key(&mut buf)?;
337 match tag {
338 TIME_SERIES_LABELS_REFS_TAG => {
339 uint32::merge_repeated(wire_type, labels_refs, &mut buf, DecodeContext::default())
340 .map_err(|mut error| {
341 error.push("TimeSeries", "labels_refs");
342 error
343 })?
344 }
345 TIME_SERIES_SAMPLES_TAG => {
346 take_length_delimited(wire_type, &mut buf).map_err(|mut error| {
347 error.push("TimeSeries", "samples");
348 error
349 })?;
350 counts.samples = counts.samples.checked_add(1).ok_or_else(|| {
351 DecodeError::new("remote write v2 sample count overflows usize")
352 })?;
353 }
354 TIME_SERIES_HISTOGRAMS_TAG => {
355 take_length_delimited(wire_type, &mut buf).map_err(|mut error| {
356 error.push("TimeSeries", "histograms");
357 error
358 })?;
359 counts.histograms = counts.histograms.checked_add(1).ok_or_else(|| {
360 DecodeError::new("remote write v2 histogram count overflows usize")
361 })?;
362 }
363 TIME_SERIES_EXEMPLARS_TAG => {
364 take_length_delimited(wire_type, &mut buf)?;
365 }
366 TIME_SERIES_METADATA_TAG => {
367 message::merge(wire_type, metadata, &mut buf, DecodeContext::default()).map_err(
368 |mut error| {
369 error.push("TimeSeries", "metadata");
370 error
371 },
372 )?
373 }
374 _ => skip_field(wire_type, tag, &mut buf, DecodeContext::default())?,
375 }
376 }
377
378 Ok(counts)
379}
380
381#[derive(Default)]
382struct LeafScratch {
383 sample: Sample,
384 histogram: Histogram,
385 exemplar: Exemplar,
386}
387
388enum SeriesWriter<'a> {
389 Samples(&'a mut TableData),
390 Histograms(&'a mut TableData),
391}
392
393fn decode_series_leaves(
394 mut buf: &[u8],
395 mut writer: Option<SeriesWriter<'_>>,
396 mut tags: PromTags<'_>,
397 mut rows_remaining: usize,
398 scratch: &mut LeafScratch,
399) -> Result<()> {
400 let mut sample_row_template = None;
401
402 while buf.has_remaining() {
403 let (tag, wire_type) = decode_key(&mut buf).context(error::DecodePromRemoteRequestSnafu)?;
404 match tag {
405 TIME_SERIES_LABELS_REFS_TAG => {
406 skip_field(wire_type, tag, &mut buf, DecodeContext::default())
407 .context(error::DecodePromRemoteRequestSnafu)?
408 }
409 TIME_SERIES_SAMPLES_TAG => {
410 scratch.sample.clear();
411 message::merge(
412 wire_type,
413 &mut scratch.sample,
414 &mut buf,
415 DecodeContext::default(),
416 )
417 .map_err(|mut error| {
418 error.push("TimeSeries", "samples");
419 error
420 })
421 .context(error::DecodePromRemoteRequestSnafu)?;
422 rows_remaining = rows_remaining.checked_sub(1).with_context(|| {
423 error::InvalidPromRemoteRequestSnafu {
424 msg: "remote write v2 sample count changed between scans".to_string(),
425 }
426 })?;
427 if let Some(SeriesWriter::Samples(table_data)) = &mut writer {
428 if sample_row_template.is_none() {
429 let timestamp_index = table_data.ensure_column_by_name(
430 greptime_timestamp(),
431 ColumnDataType::TimestampMillisecond,
432 SemanticType::Timestamp,
433 )?;
434 let value_index = table_data.ensure_column_by_name(
435 greptime_value(),
436 ColumnDataType::Float64,
437 SemanticType::Field,
438 )?;
439 let mut row = table_data.alloc_one_row();
440 row_writer::write_tags(
441 table_data,
442 std::mem::take(&mut tags).into_iter(),
443 &mut row,
444 )?;
445 sample_row_template = Some((row, timestamp_index, value_index));
446 }
447
448 if let Some((row, timestamp_index, value_index)) = sample_row_template.as_mut()
449 {
450 row[*timestamp_index].value_data = Some(
451 ValueData::TimestampMillisecondValue(scratch.sample.timestamp),
452 );
453 row[*value_index].value_data =
454 Some(ValueData::F64Value(scratch.sample.value));
455 let row = if rows_remaining == 0 {
456 std::mem::take(row)
457 } else {
458 row.clone()
459 };
460 table_data.add_row(row);
461 }
462 }
463 }
464 TIME_SERIES_HISTOGRAMS_TAG => {
465 scratch.histogram.clear();
466 message::merge(
467 wire_type,
468 &mut scratch.histogram,
469 &mut buf,
470 DecodeContext::default(),
471 )
472 .map_err(|mut error| {
473 error.push("TimeSeries", "histograms");
474 error
475 })
476 .context(error::DecodePromRemoteRequestSnafu)?;
477 rows_remaining = rows_remaining.checked_sub(1).with_context(|| {
478 error::InvalidPromRemoteRequestSnafu {
479 msg: "remote write v2 histogram count changed between scans".to_string(),
480 }
481 })?;
482 if let Some(SeriesWriter::Histograms(table_data)) = &mut writer {
483 if rows_remaining == 0 {
484 write_native_histogram(
485 table_data,
486 &scratch.histogram,
487 std::mem::take(&mut tags).into_iter(),
488 )?;
489 } else {
490 write_native_histogram(
491 table_data,
492 &scratch.histogram,
493 tags.iter().cloned(),
494 )?;
495 }
496 }
497 }
498 TIME_SERIES_EXEMPLARS_TAG => {
499 scratch.exemplar.clear();
500 message::merge(
501 wire_type,
502 &mut scratch.exemplar,
503 &mut buf,
504 DecodeContext::default(),
505 )
506 .map_err(|mut error| {
507 error.push("TimeSeries", "exemplars");
508 error
509 })
510 .context(error::DecodePromRemoteRequestSnafu)?;
511 }
512 TIME_SERIES_METADATA_TAG => {
513 take_length_delimited(wire_type, &mut buf)
514 .map_err(|mut error| {
515 error.push("TimeSeries", "metadata");
516 error
517 })
518 .context(error::DecodePromRemoteRequestSnafu)?;
519 }
520 _ => skip_field(wire_type, tag, &mut buf, DecodeContext::default())
521 .context(error::DecodePromRemoteRequestSnafu)?,
522 }
523 }
524
525 ensure!(
526 rows_remaining == 0,
527 error::InvalidPromRemoteRequestSnafu {
528 msg: "remote write v2 row count changed between scans".to_string(),
529 }
530 );
531
532 Ok(())
533}
534
535fn take_length_delimited<'a>(
536 wire_type: WireType,
537 buf: &mut &'a [u8],
538) -> std::result::Result<&'a [u8], DecodeError> {
539 if wire_type != WireType::LengthDelimited {
540 return Err(DecodeError::new(format!(
541 "invalid wire type: {wire_type:?} (expected LengthDelimited)"
542 )));
543 }
544
545 let len = decode_varint(buf)?;
546 let len =
547 usize::try_from(len).map_err(|_| DecodeError::new("length delimiter exceeds usize"))?;
548 if len > buf.len() {
549 return Err(DecodeError::new("buffer underflow"));
550 }
551 let (value, remaining) = buf.split_at(len);
552 *buf = remaining;
553 Ok(value)
554}
555
556fn checked_total(total: u64, count: usize, name: &str) -> Result<u64> {
557 let count =
558 u64::try_from(count)
559 .ok()
560 .with_context(|| error::InvalidPromRemoteRequestSnafu {
561 msg: format!("remote write v2 {name} count exceeds u64"),
562 })?;
563 total
564 .checked_add(count)
565 .with_context(|| error::InvalidPromRemoteRequestSnafu {
566 msg: format!("remote write v2 {name} count overflows u64"),
567 })
568}
569
570fn get_or_create_table_data(
571 tables: &mut HashMap<PromCtx, HashMap<String, TableData>>,
572 prom_ctx: PromCtx,
573 table_name: String,
574 column_count: usize,
575 row_count: usize,
576) -> &mut TableData {
577 match tables.entry(prom_ctx).or_default().entry(table_name) {
578 Entry::Occupied(entry) => {
579 let table_data = entry.into_mut();
580 table_data.reserve_rows(row_count);
581 table_data
582 }
583 Entry::Vacant(entry) => entry.insert(TableData::new(column_count, row_count)),
584 }
585}
586
587fn write_native_histogram<'a>(
588 table_data: &mut TableData,
589 histogram: &Histogram,
590 tags: impl Iterator<Item = (&'a str, String)>,
591) -> Result<()> {
592 let value = encode_native_histogram(histogram).map_err(|error| {
593 error::InvalidPromRemoteRequestSnafu {
594 msg: format!("remote write v2 {error}"),
595 }
596 .build()
597 })?;
598 let column_schema = native_histogram_column_schema().map_err(|error| {
599 error::InvalidPromRemoteRequestSnafu {
600 msg: format!("remote write v2 {error}"),
601 }
602 .build()
603 })?;
604
605 let mut row = table_data.alloc_one_row();
609 row_writer::write_ts_to_millis(
610 table_data,
611 greptime_timestamp(),
612 Some(histogram.timestamp),
613 Precision::Millisecond,
614 &mut row,
615 )?;
616 row_writer::write_by_schema(
617 table_data,
618 std::iter::once((column_schema, Some(value))),
619 &mut row,
620 )?;
621
622 row_writer::write_tags(table_data, tags, &mut row)?;
623 table_data.add_row(row);
624
625 Ok(())
626}
627
628fn ensure_no_internal_histogram_labels(tags: &PromTags<'_>) -> Result<()> {
629 for (name, _) in tags {
631 ensure!(
632 *name != greptime_native_histogram() && *name != NATIVE_HISTOGRAM_FIELD,
633 error::InvalidPromRemoteRequestSnafu {
634 msg: format!(
635 "remote write v2 label `{name}` conflicts with an internal native histogram label"
636 ),
637 }
638 );
639 }
640
641 Ok(())
642}
643
644fn resolve_series_labels<'a>(
645 symbols: &'a [&str],
646 labels_refs: &[u32],
647 label_names: &mut HashSet<&'a str>,
648) -> Result<ResolvedSeriesLabels<'a>> {
649 ensure!(
650 labels_refs.len().is_multiple_of(2),
651 error::InvalidPromRemoteRequestSnafu {
652 msg: "remote write v2 labels_refs must contain name/value pairs".to_string(),
653 }
654 );
655
656 let mut prom_ctx = PromCtx::default();
657 let mut table_name = None;
658 let mut tags = Vec::with_capacity(labels_refs.len() / 2);
659 label_names.clear();
660
661 for pair in labels_refs.chunks_exact(2) {
662 let name = symbol_ref(symbols, pair[0], "label name")?;
663 let value = symbol_ref(symbols, pair[1], "label value")?;
664 validate_label(name)?;
665 ensure!(
666 label_names.insert(name),
667 error::InvalidPromRemoteRequestSnafu {
668 msg: format!("remote write v2 label name `{name}` is repeated"),
669 }
670 );
671
672 if name == METRIC_NAME_LABEL {
673 table_name = Some(value.to_string());
674 continue;
675 }
676 if apply_remote_write_special_label(name, value, &mut prom_ctx) {
677 continue;
678 }
679
680 tags.push((name, value.to_string()));
681 }
682
683 let table_name = table_name.with_context(|| error::InvalidPromRemoteRequestSnafu {
684 msg: "missing '__name__' label in time-series".to_string(),
685 })?;
686 ensure!(
687 !table_name.is_empty(),
688 error::InvalidPromRemoteRequestSnafu {
689 msg: "remote write v2 label `__name__` value must not be empty".to_string(),
690 }
691 );
692
693 Ok((prom_ctx, table_name, tags))
694}
695
696fn validate_label(name: &str) -> Result<()> {
697 ensure!(
698 validate_label_name(name.as_bytes()),
699 error::InvalidPromRemoteRequestSnafu {
700 msg: format!("remote write v2 invalid label name `{name}`"),
701 }
702 );
703
704 Ok(())
705}
706
707fn symbol_ref<'a>(symbols: &'a [&str], idx: u32, field: &str) -> Result<&'a str> {
708 let idx = usize::try_from(idx)
709 .ok()
710 .with_context(|| error::InvalidPromRemoteRequestSnafu {
711 msg: format!("remote write v2 {field} symbol reference exceeds usize"),
712 })?;
713 symbols
714 .get(idx)
715 .copied()
716 .with_context(|| error::InvalidPromRemoteRequestSnafu {
717 msg: format!(
718 "remote write v2 {field} symbol reference {idx} is out of range, symbols len: {}",
719 symbols.len()
720 ),
721 })
722}
723
724#[allow(deprecated)]
725fn apply_remote_write_special_label(name: &str, value: &str, prom_ctx: &mut PromCtx) -> bool {
726 match name {
727 SCHEMA_LABEL => {
728 prom_ctx.schema = Some(value.to_string());
729 true
730 }
731 DATABASE_LABEL | DATABASE_LABEL_ALT => {
732 if prom_ctx.schema.is_none() {
733 prom_ctx.schema = Some(value.to_string());
734 }
735 true
736 }
737 PHYSICAL_TABLE_LABEL | PHYSICAL_TABLE_LABEL_ALT => {
738 prom_ctx.physical_table = Some(value.to_string());
739 true
740 }
741 _ => false,
742 }
743}
744
745fn into_context_req(tables: HashMap<PromCtx, HashMap<String, TableData>>) -> ContextReq {
746 let mut ctx_req = ContextReq::default();
747 for (prom_ctx, tables) in tables {
748 let mut opt = ContextOpt::default();
749 if let Some(schema) = prom_ctx.schema {
750 opt.set_schema(schema);
751 }
752 if let Some(physical_table) = prom_ctx.physical_table {
753 opt.set_physical_table(physical_table);
754 }
755
756 ctx_req.add_rows(
757 opt,
758 tables.into_iter().map(|(table_name, table_data)| {
759 table_data_to_row_insert_request(table_name, table_data)
760 }),
761 );
762 }
763 ctx_req
764}
765
766fn table_data_to_row_insert_request(table_name: String, table_data: TableData) -> RowInsertRequest {
767 let num_columns = table_data.num_columns();
768 let (schema, mut rows) = table_data.into_schema_and_rows();
769 for row in &mut rows {
770 if num_columns > row.values.len() {
771 row.values.resize(num_columns, Value { value_data: None });
772 }
773 }
774
775 RowInsertRequest {
776 table_name,
777 rows: Some(Rows { schema, rows }),
778 }
779}
780
781#[cfg(any(test, feature = "testing"))]
782pub mod test_util {
783 use api::greptime_proto::io::prometheus::write::v2::{Histogram, Request, Sample, TimeSeries};
784 use api::v1::RowInsertRequest;
785 use bytes::Bytes;
786 use prost::Message;
787 use snafu::ResultExt;
788
789 use crate::error::{self, Result};
790 use crate::prom_remote_write::decompress_remote_write_body;
791 use crate::prom_store::snappy_compress;
792 use crate::request_memory_limiter::ServerMemoryLimiter;
793
794 pub fn request_with_labels_and_samples(
795 labels: Vec<(&str, &str)>,
796 samples: Vec<Sample>,
797 ) -> Request {
798 request_with_labels(labels, samples, Vec::new())
799 }
800
801 pub fn request_with_labels_and_histograms(
802 labels: Vec<(&str, &str)>,
803 histograms: Vec<Histogram>,
804 ) -> Request {
805 request_with_labels(labels, Vec::new(), histograms)
806 }
807
808 pub fn decode_request(is_zstd: bool, body: Bytes) -> Result<Request> {
809 let buf = block_on_decode(is_zstd, body)?;
810 Request::decode(&buf[..]).context(error::DecodePromRemoteRequestSnafu)
811 }
812
813 fn block_on_decode(is_zstd: bool, body: Bytes) -> Result<Vec<u8>> {
816 tokio::runtime::Builder::new_current_thread()
817 .enable_all()
818 .build()
819 .unwrap()
820 .block_on(async move {
821 decompress_remote_write_body(is_zstd, &body[..], &ServerMemoryLimiter::default())
822 .await
823 .map(|buf| buf.to_vec())
824 })
825 }
826
827 pub fn write_requests(
828 request: Request,
829 ) -> Result<(Vec<RowInsertRequest>, Vec<RowInsertRequest>, u64, u64)> {
830 let body = Bytes::from(snappy_compress(&request.encode_to_vec())?);
831 decode_write_requests(false, body)
832 }
833
834 pub fn decode_write_requests(
835 is_zstd: bool,
836 body: Bytes,
837 ) -> Result<(Vec<RowInsertRequest>, Vec<RowInsertRequest>, u64, u64)> {
838 let requests = tokio::runtime::Builder::new_current_thread()
839 .enable_all()
840 .build()
841 .unwrap()
842 .block_on(super::decode_remote_write_v2(
843 is_zstd,
844 body,
845 &ServerMemoryLimiter::default(),
846 ))?;
847 Ok((
848 requests.samples.all_req().collect(),
849 requests.histograms.all_req().collect(),
850 requests.sample_count,
851 requests.histogram_count,
852 ))
853 }
854
855 pub fn decode_uncompressed_write_requests(
856 body: &[u8],
857 ) -> Result<(Vec<RowInsertRequest>, Vec<RowInsertRequest>, u64, u64)> {
858 let request =
859 super::BorrowedRequest::decode(body).context(error::DecodePromRemoteRequestSnafu)?;
860 let requests = super::convert_remote_write_v2(request)?;
861 Ok((
862 requests.samples.all_req().collect(),
863 requests.histograms.all_req().collect(),
864 requests.sample_count,
865 requests.histogram_count,
866 ))
867 }
868
869 pub fn histogram(timestamp: i64) -> Histogram {
870 Histogram {
871 timestamp,
872 ..Default::default()
873 }
874 }
875
876 fn request_with_labels(
877 labels: Vec<(&str, &str)>,
878 samples: Vec<Sample>,
879 histograms: Vec<Histogram>,
880 ) -> Request {
881 let mut symbols = vec!["".to_string()];
882 let mut labels_refs = Vec::with_capacity(labels.len() * 2);
883 for (name, value) in labels {
884 labels_refs.push(push_symbol(&mut symbols, name));
885 labels_refs.push(push_symbol(&mut symbols, value));
886 }
887
888 Request {
889 symbols,
890 timeseries: vec![TimeSeries {
891 labels_refs,
892 samples,
893 histograms,
894 exemplars: Vec::new(),
895 metadata: None,
896 }],
897 }
898 }
899
900 fn push_symbol(symbols: &mut Vec<String>, symbol: &str) -> u32 {
901 if let Some(idx) = symbols.iter().position(|s| s == symbol) {
902 return idx as u32;
903 }
904
905 let idx = symbols.len();
906 symbols.push(symbol.to_string());
907 idx as u32
908 }
909}
910
911#[cfg(test)]
912mod tests {
913 use std::sync::Arc;
914
915 use api::v1::value::ValueData;
916 use common_query::prelude::{greptime_timestamp, greptime_value, set_default_prefix};
917 use session::context::QueryContext;
918
919 use super::*;
920 use crate::error;
921 use crate::http::prom_store::PHYSICAL_TABLE_PARAM;
922 use crate::prom_store::{DATABASE_LABEL, PHYSICAL_TABLE_LABEL};
923
924 #[test]
925 fn test_decode_remote_write_v2_request() {
926 let request = Request {
927 symbols: vec![
928 "".to_string(),
929 "__name__".to_string(),
930 "http_requests_total".to_string(),
931 ],
932 timeseries: vec![TimeSeries {
933 labels_refs: vec![1, 2],
934 samples: vec![Sample {
935 value: 42.0,
936 timestamp: 1000,
937 start_timestamp: 0,
938 }],
939 histograms: Vec::new(),
940 exemplars: Vec::new(),
941 metadata: Some(Metadata {
942 r#type: metadata::MetricType::Counter as i32,
943 help_ref: 0,
944 unit_ref: 0,
945 }),
946 }],
947 };
948 let body =
949 Bytes::from(crate::prom_store::snappy_compress(&request.encode_to_vec()).unwrap());
950
951 let decoded = test_util::decode_request(false, body.clone()).unwrap();
952
953 assert_eq!(decoded.symbols, request.symbols);
954 assert_eq!(decoded.timeseries.len(), 1);
955 assert_eq!(decoded.timeseries[0].labels_refs, vec![1, 2]);
956 assert_eq!(decoded.timeseries[0].samples.len(), 1);
957 assert_eq!(decoded.timeseries[0].samples[0].value, 42.0);
958 assert_eq!(decoded.timeseries[0].metadata.as_ref().unwrap().r#type, 1);
959 assert_eq!(
960 decode_v2_on_test_runtime(true, body).unwrap().sample_count,
961 1
962 );
963 }
964
965 #[test]
966 fn test_fused_decoder_accepts_arbitrary_field_order_and_split_labels() {
967 let mut first_sample = Sample {
968 value: 42.0,
969 timestamp: 1000,
970 start_timestamp: 500,
971 }
972 .encode_to_vec();
973 first_sample.extend(varint_field(90, 1));
974 let second_sample = Sample {
975 value: 43.0,
976 timestamp: 2000,
977 start_timestamp: 1000,
978 }
979 .encode_to_vec();
980
981 let mut series = encoded_message_field(2, &first_sample);
982 series.extend(packed_u32_field(1, &[1]));
983 series.extend(varint_field(90, 1));
984 series.extend(encoded_message_field(2, &second_sample));
985 series.extend(varint_field(1, 2));
986
987 let mut wire = encoded_message_field(5, &series);
988 wire.extend(string_field(4, b""));
989 wire.extend(varint_field(90, 1));
990 wire.extend(string_field(4, METRIC_NAME_LABEL.as_bytes()));
991 wire.extend(encoded_message_field(5, &packed_u32_field(1, &[99])));
992 wire.extend(string_field(4, b"http_requests_total"));
993
994 let requests = decode_wire(&wire).unwrap();
995 assert_eq!(requests.sample_count, 2);
996 assert_eq!(requests.histogram_count, 0);
997 let rows = requests.samples.all_req().next().unwrap().rows.unwrap();
998 assert_eq!(rows.rows.len(), 2);
999 assert_eq!(
1000 rows.schema
1001 .iter()
1002 .map(|column| column.column_name.as_str())
1003 .collect::<Vec<_>>(),
1004 vec![greptime_timestamp(), greptime_value()]
1005 );
1006 assert_eq!(
1007 rows.rows[0].values[1].value_data,
1008 Some(ValueData::F64Value(42.0))
1009 );
1010 assert_eq!(
1011 rows.rows[1].values[1].value_data,
1012 Some(ValueData::F64Value(43.0))
1013 );
1014 }
1015
1016 #[test]
1017 fn test_fused_decoder_accepts_histogram_before_labels_and_unknown_fields() {
1018 let mut histogram = Histogram {
1019 count: Some(Count::CountInt(0)),
1020 zero_count: Some(ZeroCount::ZeroCountInt(0)),
1021 timestamp: 2000,
1022 start_timestamp: 1000,
1023 ..Default::default()
1024 }
1025 .encode_to_vec();
1026 histogram.extend(varint_field(90, 1));
1027
1028 let mut series = encoded_message_field(3, &histogram);
1029 series.extend(packed_u32_field(1, &[1, 2]));
1030 let wire = request_wire(&["", METRIC_NAME_LABEL, "metric"], &[series]);
1031
1032 let requests = decode_wire(&wire).unwrap();
1033 assert_eq!(requests.histogram_count, 1);
1034 let rows = requests.histograms.all_req().next().unwrap().rows.unwrap();
1035 assert_eq!(
1036 histogram_field_value(&rows, 0, START_TIMESTAMP_FIELD),
1037 Some(ValueData::TimestampMillisecondValue(1000))
1038 );
1039 }
1040
1041 #[test]
1042 fn test_fused_decoder_rejects_malformed_wire() {
1043 let mut wrong_request_wire = varint_field(4, 0);
1044 let invalid_utf8 = string_field(4, &[0xff]);
1045
1046 let mut truncated_request = Vec::new();
1047 prost::encoding::encode_key(5, WireType::LengthDelimited, &mut truncated_request);
1048 prost::encoding::encode_varint(2, &mut truncated_request);
1049 truncated_request.push(0);
1050
1051 let mut oversized_request = Vec::new();
1052 prost::encoding::encode_key(4, WireType::LengthDelimited, &mut oversized_request);
1053 prost::encoding::encode_varint(u64::MAX, &mut oversized_request);
1054
1055 let malformed_varint = vec![0x80; 10];
1056
1057 let mut wrong_labels_wire = Vec::new();
1058 prost::encoding::encode_key(1, WireType::SixtyFourBit, &mut wrong_labels_wire);
1059 wrong_labels_wire.extend([0; 8]);
1060 wrong_labels_wire = request_wire(&[""], &[std::mem::take(&mut wrong_labels_wire)]);
1061
1062 let wrong_series_wire = request_wire(&[""], &[varint_field(2, 0)]);
1063
1064 let mut truncated_series = Vec::new();
1065 prost::encoding::encode_key(2, WireType::LengthDelimited, &mut truncated_series);
1066 prost::encoding::encode_varint(2, &mut truncated_series);
1067 truncated_series.push(0x08);
1068 let truncated_series = request_wire(&[""], &[truncated_series]);
1069
1070 let mut oversized_series = Vec::new();
1071 prost::encoding::encode_key(3, WireType::LengthDelimited, &mut oversized_series);
1072 prost::encoding::encode_varint(u64::MAX, &mut oversized_series);
1073 let oversized_series = request_wire(&[""], &[oversized_series]);
1074
1075 let invalid_sample = series_wire(&[1, 2], 2, &varint_field(1, 1));
1076 let invalid_sample = request_wire(&["", METRIC_NAME_LABEL, "metric"], &[invalid_sample]);
1077 let invalid_histogram = series_wire(&[1, 2], 3, &varint_field(3, 1));
1078 let invalid_histogram =
1079 request_wire(&["", METRIC_NAME_LABEL, "metric"], &[invalid_histogram]);
1080
1081 for (name, wire) in [
1082 (
1083 "wrong request wire type",
1084 std::mem::take(&mut wrong_request_wire),
1085 ),
1086 ("invalid utf8", invalid_utf8),
1087 ("truncated request", truncated_request),
1088 ("oversized request length", oversized_request),
1089 ("malformed varint", malformed_varint),
1090 ("wrong labels wire type", wrong_labels_wire),
1091 ("wrong series wire type", wrong_series_wire),
1092 ("truncated series", truncated_series),
1093 ("oversized series length", oversized_series),
1094 ("invalid sample", invalid_sample),
1095 ("invalid histogram", invalid_histogram),
1096 ] {
1097 let error = decode_wire_error(&wire, name);
1098 assert!(
1099 matches!(error, error::Error::DecodePromRemoteRequest { .. }),
1100 "{name}: {error}"
1101 );
1102 }
1103 }
1104
1105 #[test]
1106 fn test_fused_decoder_rejects_malformed_ignored_messages() {
1107 for tag in [4, 5] {
1108 let mut series = series_wire(&[1, 2], 2, &Sample::default().encode_to_vec());
1109 series.extend(encoded_message_field(tag, &[0x08]));
1110 let wire = request_wire(&["", METRIC_NAME_LABEL, "metric"], &[series]);
1111
1112 let error = decode_wire_error(&wire, "malformed ignored message");
1113 assert!(matches!(
1114 error,
1115 error::Error::DecodePromRemoteRequest { .. }
1116 ));
1117 }
1118 }
1119
1120 #[test]
1121 fn test_fused_decoder_ignores_exemplar_symbol_refs() {
1122 let request = Request {
1123 symbols: vec![
1124 String::new(),
1125 METRIC_NAME_LABEL.to_string(),
1126 "metric".to_string(),
1127 ],
1128 timeseries: vec![TimeSeries {
1129 labels_refs: vec![1, 2],
1130 samples: vec![Sample::default()],
1131 exemplars: vec![Exemplar {
1132 labels_refs: vec![99],
1133 ..Default::default()
1134 }],
1135 ..Default::default()
1136 }],
1137 };
1138
1139 assert_eq!(decode_test_request(request).unwrap().sample_count, 1);
1140 }
1141
1142 #[test]
1143 fn test_fused_decoder_preserves_empty_request_and_series_behavior() {
1144 let request = Request {
1145 symbols: vec![
1146 String::new(),
1147 METRIC_NAME_LABEL.to_string(),
1148 "metric".to_string(),
1149 "job".to_string(),
1150 ],
1151 timeseries: vec![
1152 TimeSeries {
1153 labels_refs: vec![99, 99],
1154 ..Default::default()
1155 },
1156 TimeSeries {
1157 labels_refs: vec![1],
1158 ..Default::default()
1159 },
1160 TimeSeries {
1161 labels_refs: vec![3, 2, 3, 2],
1162 ..Default::default()
1163 },
1164 TimeSeries::default(),
1165 ],
1166 };
1167
1168 let requests = decode_test_request(request).unwrap();
1169 assert_eq!(requests.sample_count, 0);
1170 assert_eq!(requests.histogram_count, 0);
1171
1172 let requests = decode_test_request(Request {
1173 symbols: vec![String::new()],
1174 timeseries: Vec::new(),
1175 })
1176 .unwrap();
1177 assert_eq!(requests.sample_count, 0);
1178 assert!(decode_wire(&[]).is_err());
1179 assert!(decode_wire(&[0x0a, 0x00]).is_err());
1180 }
1181
1182 #[test]
1183 fn test_fused_decoder_rejects_duplicate_resolved_label_names() {
1184 let request = Request {
1185 symbols: vec![
1186 String::new(),
1187 METRIC_NAME_LABEL.to_string(),
1188 "metric".to_string(),
1189 "job".to_string(),
1190 "job".to_string(),
1191 "api".to_string(),
1192 "worker".to_string(),
1193 ],
1194 timeseries: vec![TimeSeries {
1195 labels_refs: vec![1, 2, 3, 5, 4, 6],
1196 samples: vec![Sample::default()],
1197 ..Default::default()
1198 }],
1199 };
1200
1201 assert_invalid(
1202 "duplicate resolved labels",
1203 request,
1204 "label name `job` is repeated",
1205 );
1206 }
1207
1208 #[test]
1209 fn test_fused_decoder_rejects_malformed_series_with_histograms() {
1210 let histogram = series_wire(
1211 &[1, 2],
1212 3,
1213 &Histogram {
1214 count: Some(Count::CountInt(0)),
1215 zero_count: Some(ZeroCount::ZeroCountInt(0)),
1216 ..Default::default()
1217 }
1218 .encode_to_vec(),
1219 );
1220 let malformed_sample = series_wire(&[1, 3], 2, &[0x08]);
1221
1222 let wire = request_wire(
1223 &["", METRIC_NAME_LABEL, "histogram", "sample"],
1224 &[histogram.clone(), malformed_sample.clone()],
1225 );
1226 let error = decode_wire_error(&wire, "histogram before malformed series");
1227 assert!(matches!(
1228 error,
1229 error::Error::DecodePromRemoteRequest { .. }
1230 ));
1231
1232 let wire = request_wire(
1233 &["", METRIC_NAME_LABEL, "histogram", "sample"],
1234 &[malformed_sample, histogram.clone()],
1235 );
1236 let error = decode_wire_error(&wire, "malformed series before histogram");
1237 assert!(matches!(
1238 error,
1239 error::Error::DecodePromRemoteRequest { .. }
1240 ));
1241
1242 let missing_name = series_wire(&[3, 4], 2, &Sample::default().encode_to_vec());
1243 let wire = request_wire(
1244 &["", METRIC_NAME_LABEL, "metric", "job", "api"],
1245 &[missing_name, histogram],
1246 );
1247 let error = decode_wire_error(&wire, "conversion error before histogram");
1248 assert!(error.to_string().contains("missing '__name__'"));
1249 }
1250
1251 #[test]
1252 fn test_into_context_req_samples() {
1253 let ctx_req = decode_test_request(test_util::request_with_labels_and_samples(
1254 vec![
1255 (METRIC_NAME_LABEL, "http_requests_total"),
1256 ("job", "api"),
1257 ("instance", "localhost:9090"),
1258 ],
1259 vec![
1260 Sample {
1261 value: 42.0,
1262 timestamp: 1000,
1263 start_timestamp: 0,
1264 },
1265 Sample {
1266 value: 43.0,
1267 timestamp: 2000,
1268 start_timestamp: 0,
1269 },
1270 ],
1271 ))
1272 .unwrap();
1273
1274 assert_eq!(ctx_req.sample_count, 2);
1275 assert_eq!(ctx_req.histogram_count, 0);
1276 assert_eq!(ctx_req.histograms.all_req().count(), 0);
1277 let mut inserts = ctx_req.samples.all_req().collect::<Vec<_>>();
1278 assert_eq!(inserts.len(), 1);
1279
1280 let request = inserts.pop().unwrap();
1281 assert_eq!(request.table_name, "http_requests_total");
1282 let rows = request.rows.unwrap();
1283 assert_eq!(rows.rows.len(), 2);
1284 assert_eq!(
1285 rows.schema
1286 .iter()
1287 .map(|col| col.column_name.as_str())
1288 .collect::<Vec<_>>(),
1289 vec![greptime_timestamp(), greptime_value(), "job", "instance"]
1290 );
1291 assert_eq!(
1292 rows.rows[0].values[0].value_data,
1293 Some(ValueData::TimestampMillisecondValue(1000))
1294 );
1295 assert_eq!(
1296 rows.rows[0].values[1].value_data,
1297 Some(ValueData::F64Value(42.0))
1298 );
1299 assert_eq!(
1300 rows.rows[0].values[2].value_data,
1301 Some(ValueData::StringValue("api".to_string()))
1302 );
1303 assert_eq!(
1304 rows.rows[0].values[3].value_data,
1305 Some(ValueData::StringValue("localhost:9090".to_string()))
1306 );
1307 assert_eq!(
1308 rows.rows[1].values[0].value_data,
1309 Some(ValueData::TimestampMillisecondValue(2000))
1310 );
1311 assert_eq!(
1312 rows.rows[1].values[1].value_data,
1313 Some(ValueData::F64Value(43.0))
1314 );
1315 assert_eq!(
1316 rows.rows[1].values[2].value_data,
1317 Some(ValueData::StringValue("api".to_string()))
1318 );
1319 assert_eq!(
1320 rows.rows[1].values[3].value_data,
1321 Some(ValueData::StringValue("localhost:9090".to_string()))
1322 );
1323 }
1324
1325 #[test]
1326 fn test_into_context_req_special_labels() {
1327 let ctx_req = decode_test_request(test_util::request_with_labels_and_samples(
1328 vec![
1329 (METRIC_NAME_LABEL, "cpu_usage"),
1330 (DATABASE_LABEL, "tenant_a"),
1331 (PHYSICAL_TABLE_LABEL, "metrics_physical"),
1332 ("job", "api"),
1333 ],
1334 vec![Sample {
1335 value: 1.0,
1336 timestamp: 1000,
1337 start_timestamp: 0,
1338 }],
1339 ))
1340 .unwrap();
1341
1342 let mut iter = ctx_req
1343 .samples
1344 .as_req_iter(Arc::new(QueryContext::with("greptime", "public")));
1345 let (ctx, reqs) = iter.next().unwrap();
1346 assert!(iter.next().is_none());
1347
1348 assert_eq!(ctx.current_schema(), "tenant_a");
1349 assert_eq!(
1350 ctx.extension(PHYSICAL_TABLE_PARAM),
1351 Some("metrics_physical")
1352 );
1353 assert_eq!(reqs.inserts.len(), 1);
1354
1355 let rows = reqs.inserts[0].rows.as_ref().unwrap();
1356 assert_eq!(
1357 rows.schema
1358 .iter()
1359 .map(|col| col.column_name.as_str())
1360 .collect::<Vec<_>>(),
1361 vec![greptime_timestamp(), greptime_value(), "job"]
1362 );
1363 }
1364
1365 #[test]
1366 fn test_into_context_req_rejects_invalid_requests() {
1367 let mut cases = Vec::new();
1368
1369 cases.push((
1370 "missing metric name",
1371 request_with_sample(vec![("job", "api")]),
1372 "missing '__name__'",
1373 ));
1374
1375 let mut request = request_with_sample(vec![(METRIC_NAME_LABEL, "metric")]);
1376 request.timeseries[0].labels_refs.push(1);
1377 cases.push((
1378 "odd label refs",
1379 request,
1380 "labels_refs must contain name/value pairs",
1381 ));
1382
1383 let mut request = request_with_sample(vec![(METRIC_NAME_LABEL, "metric")]);
1384 request.timeseries[0].labels_refs[1] = 99;
1385 cases.push((
1386 "out of range symbol ref",
1387 request,
1388 "symbol reference 99 is out of range",
1389 ));
1390
1391 let mut request = request_with_sample(vec![(METRIC_NAME_LABEL, "metric")]);
1392 request.symbols[0] = "not-empty".to_string();
1393 cases.push((
1394 "non-empty first symbol",
1395 request,
1396 "symbols must start with an empty string",
1397 ));
1398
1399 cases.push((
1400 "repeated label name",
1401 request_with_sample(vec![
1402 (METRIC_NAME_LABEL, "metric"),
1403 ("job", "api"),
1404 ("job", "worker"),
1405 ]),
1406 "label name `job` is repeated",
1407 ));
1408
1409 cases.push((
1410 "empty label name",
1411 request_with_sample(vec![(METRIC_NAME_LABEL, "metric"), ("", "api")]),
1412 "invalid label name",
1413 ));
1414
1415 cases.push((
1416 "invalid label name",
1417 request_with_sample(vec![(METRIC_NAME_LABEL, "metric"), ("has-dash", "api")]),
1418 "invalid label name",
1419 ));
1420
1421 cases.push((
1422 "dotted label name",
1423 request_with_sample(vec![(METRIC_NAME_LABEL, "metric"), ("service.name", "api")]),
1424 "invalid label name",
1425 ));
1426
1427 cases.push((
1428 "non-ascii label name",
1429 request_with_sample(vec![(METRIC_NAME_LABEL, "metric"), ("区域", "api")]),
1430 "invalid label name",
1431 ));
1432
1433 cases.push((
1434 "empty metric name",
1435 request_with_sample(vec![(METRIC_NAME_LABEL, "")]),
1436 "label `__name__` value must not be empty",
1437 ));
1438
1439 cases.push((
1440 "internal histogram label on samples",
1441 request_with_sample(vec![
1442 (METRIC_NAME_LABEL, "metric"),
1443 (greptime_native_histogram(), "user_value"),
1444 ]),
1445 "conflicts with an internal native histogram label",
1446 ));
1447
1448 cases.push((
1449 "int count with float zero count",
1450 request_with_histogram(Histogram {
1451 count: Some(Count::CountInt(1)),
1452 zero_count: Some(ZeroCount::ZeroCountFloat(0.5)),
1453 ..Default::default()
1454 }),
1455 "count and zero_count must use the same integer or float family",
1456 ));
1457
1458 cases.push((
1459 "float count with int zero count",
1460 request_with_histogram(Histogram {
1461 count: Some(Count::CountFloat(1.0)),
1462 zero_count: Some(ZeroCount::ZeroCountInt(1)),
1463 ..Default::default()
1464 }),
1465 "count and zero_count must use the same integer or float family",
1466 ));
1467
1468 cases.push((
1469 "reducible schema",
1470 request_with_histogram(Histogram {
1471 schema: 9,
1472 ..Default::default()
1473 }),
1474 "schema 9 must be reduced before ingestion",
1475 ));
1476
1477 cases.push((
1478 "unsupported schema",
1479 request_with_histogram(Histogram {
1480 schema: 53,
1481 ..Default::default()
1482 }),
1483 "schema 53 is unsupported",
1484 ));
1485
1486 cases.push((
1487 "standard schema with custom values",
1488 request_with_histogram(Histogram {
1489 schema: 1,
1490 custom_values: vec![1.0],
1491 ..Default::default()
1492 }),
1493 "standard native histogram must not use custom_values",
1494 ));
1495
1496 cases.push((
1497 "custom values with inf",
1498 request_with_histogram(Histogram {
1499 schema: CUSTOM_BUCKETS_SCHEMA,
1500 custom_values: vec![f64::INFINITY],
1501 ..Default::default()
1502 }),
1503 "custom_values must not contain +Inf or NaN",
1504 ));
1505
1506 cases.push((
1507 "custom values not sorted",
1508 request_with_histogram(Histogram {
1509 schema: CUSTOM_BUCKETS_SCHEMA,
1510 custom_values: vec![2.0, 1.0],
1511 ..Default::default()
1512 }),
1513 "custom_values must be sorted",
1514 ));
1515
1516 cases.push((
1517 "custom schema with zero bucket",
1518 request_with_histogram(Histogram {
1519 schema: CUSTOM_BUCKETS_SCHEMA,
1520 zero_count: Some(ZeroCount::ZeroCountInt(1)),
1521 ..Default::default()
1522 }),
1523 "custom native histogram must not use a zero bucket",
1524 ));
1525
1526 cases.push((
1527 "custom schema with negative buckets",
1528 request_with_histogram(Histogram {
1529 schema: CUSTOM_BUCKETS_SCHEMA,
1530 negative_spans: vec![BucketSpan {
1531 offset: -1,
1532 length: 1,
1533 }],
1534 negative_deltas: vec![1],
1535 ..Default::default()
1536 }),
1537 "custom native histogram must not use negative buckets",
1538 ));
1539
1540 cases.push((
1541 "span count mismatch",
1542 request_with_histogram(Histogram {
1543 positive_spans: vec![BucketSpan {
1544 offset: 0,
1545 length: 2,
1546 }],
1547 positive_deltas: vec![1],
1548 ..Default::default()
1549 }),
1550 "positive spans describe 2 buckets, found 1",
1551 ));
1552
1553 cases.push((
1554 "negative offset after first span",
1555 request_with_histogram(Histogram {
1556 count: Some(Count::CountInt(2)),
1557 positive_spans: vec![
1558 BucketSpan {
1559 offset: 0,
1560 length: 1,
1561 },
1562 BucketSpan {
1563 offset: -1,
1564 length: 1,
1565 },
1566 ],
1567 positive_deltas: vec![1, 0],
1568 ..Default::default()
1569 }),
1570 "positive span 2 has negative offset -1",
1571 ));
1572
1573 cases.push((
1574 "negative custom span offset",
1575 request_with_histogram(Histogram {
1576 count: Some(Count::CountInt(1)),
1577 schema: CUSTOM_BUCKETS_SCHEMA,
1578 custom_values: vec![1.0],
1579 positive_spans: vec![BucketSpan {
1580 offset: -1,
1581 length: 1,
1582 }],
1583 positive_deltas: vec![1],
1584 ..Default::default()
1585 }),
1586 "positive span 1 has negative offset -1",
1587 ));
1588
1589 cases.push((
1590 "integer bucket total mismatch",
1591 request_with_histogram(Histogram {
1592 count: Some(Count::CountInt(0)),
1593 positive_spans: vec![BucketSpan {
1594 offset: 0,
1595 length: 1,
1596 }],
1597 positive_deltas: vec![1],
1598 ..Default::default()
1599 }),
1600 "has 1 observations in buckets, count is 0",
1601 ));
1602
1603 cases.push((
1604 "negative float count",
1605 request_with_histogram(Histogram {
1606 count: Some(Count::CountFloat(-1.0)),
1607 ..Default::default()
1608 }),
1609 "float count must not be negative",
1610 ));
1611
1612 cases.push((
1613 "negative float zero count",
1614 request_with_histogram(Histogram {
1615 count: Some(Count::CountFloat(0.0)),
1616 zero_count: Some(ZeroCount::ZeroCountFloat(-1.0)),
1617 ..Default::default()
1618 }),
1619 "float zero_count must not be negative",
1620 ));
1621
1622 cases.push((
1623 "negative float bucket count",
1624 request_with_histogram(Histogram {
1625 count: Some(Count::CountFloat(0.0)),
1626 positive_spans: vec![BucketSpan {
1627 offset: 0,
1628 length: 1,
1629 }],
1630 positive_counts: vec![-1.0],
1631 ..Default::default()
1632 }),
1633 "positive bucket 1 count must not be negative",
1634 ));
1635
1636 cases.push((
1637 "custom span index out of range",
1638 request_with_histogram(Histogram {
1639 schema: CUSTOM_BUCKETS_SCHEMA,
1640 custom_values: vec![1.0],
1641 positive_spans: vec![BucketSpan {
1642 offset: 2,
1643 length: 1,
1644 }],
1645 positive_deltas: vec![1],
1646 ..Default::default()
1647 }),
1648 "positive bucket index 2 is out of range",
1649 ));
1650
1651 for (name, request, expected) in cases {
1652 assert_invalid(name, request, expected);
1653 }
1654 }
1655
1656 #[test]
1657 fn test_into_context_req_allows_nan_observations_outside_buckets() {
1658 decode_test_request(request_with_histogram(Histogram {
1659 count: Some(Count::CountInt(2)),
1660 sum: f64::NAN,
1661 positive_spans: vec![BucketSpan {
1662 offset: 0,
1663 length: 1,
1664 }],
1665 positive_deltas: vec![1],
1666 ..Default::default()
1667 }))
1668 .unwrap();
1669 }
1670
1671 #[test]
1672 fn test_into_context_req_allows_empty_label_values() {
1673 let ctx_req = decode_test_request(test_util::request_with_labels_and_samples(
1674 vec![(METRIC_NAME_LABEL, "metric"), ("job", "")],
1675 vec![Sample {
1676 value: 1.0,
1677 timestamp: 1000,
1678 start_timestamp: 0,
1679 }],
1680 ))
1681 .unwrap();
1682
1683 let rows = ctx_req.samples.all_req().next().unwrap().rows.unwrap();
1684 let job_idx = column_index(&rows.schema, "job");
1685 assert_eq!(
1686 rows.rows[0].values[job_idx].value_data,
1687 Some(ValueData::StringValue(String::new()))
1688 );
1689 }
1690
1691 #[test]
1692 fn test_into_context_req_rejects_same_metric_samples_and_histograms() {
1693 let mut request = test_util::request_with_labels_and_samples(
1694 vec![(METRIC_NAME_LABEL, "metric")],
1695 vec![Sample {
1696 value: 1.0,
1697 timestamp: 1000,
1698 start_timestamp: 0,
1699 }],
1700 );
1701 request.timeseries[0].histograms.push(Histogram::default());
1702
1703 assert_invalid(
1704 "same metric samples and histograms",
1705 request,
1706 "contains both samples and native histograms",
1707 );
1708
1709 let mut request = test_util::request_with_labels_and_samples(
1710 vec![(METRIC_NAME_LABEL, "metric")],
1711 vec![Sample {
1712 value: 1.0,
1713 timestamp: 1000,
1714 start_timestamp: 0,
1715 }],
1716 );
1717 request.timeseries.push(TimeSeries {
1718 labels_refs: request.timeseries[0].labels_refs.clone(),
1719 histograms: vec![Histogram::default()],
1720 ..Default::default()
1721 });
1722
1723 assert_invalid(
1724 "same metric samples and histograms across series",
1725 request,
1726 "contains both samples and native histograms",
1727 );
1728 }
1729
1730 #[test]
1731 fn test_into_context_req_rejects_metric_kind_conflict_across_label_sets() {
1732 let request = Request {
1733 symbols: vec![
1734 "".to_string(),
1735 METRIC_NAME_LABEL.to_string(),
1736 "metric".to_string(),
1737 "job".to_string(),
1738 "api".to_string(),
1739 "worker".to_string(),
1740 ],
1741 timeseries: vec![
1742 TimeSeries {
1743 labels_refs: vec![1, 2, 3, 4],
1744 samples: vec![Sample {
1745 value: 1.0,
1746 timestamp: 1000,
1747 start_timestamp: 0,
1748 }],
1749 ..Default::default()
1750 },
1751 TimeSeries {
1752 labels_refs: vec![1, 2, 3, 5],
1753 histograms: vec![Histogram::default()],
1754 ..Default::default()
1755 },
1756 ],
1757 };
1758
1759 assert_invalid(
1760 "same metric kind conflict across label sets",
1761 request,
1762 "contains both samples and native histograms",
1763 );
1764 }
1765
1766 #[test]
1767 fn test_into_context_req_validates_exponential_overflow_bucket_index() {
1768 for schema in [-4, 0, 8] {
1769 let max_index = exponential_overflow_bucket_index(schema).unwrap();
1770 for positive in [true, false] {
1771 let mut histogram = Histogram {
1772 schema,
1773 count: Some(Count::CountInt(1)),
1774 ..Default::default()
1775 };
1776 if positive {
1777 histogram.positive_spans = vec![BucketSpan {
1778 offset: max_index,
1779 length: 1,
1780 }];
1781 histogram.positive_deltas = vec![1];
1782 } else {
1783 histogram.negative_spans = vec![BucketSpan {
1784 offset: max_index,
1785 length: 1,
1786 }];
1787 histogram.negative_deltas = vec![1];
1788 }
1789 decode_test_request(request_with_histogram(histogram.clone())).unwrap();
1790
1791 let beyond = max_index + 1;
1792 if positive {
1793 histogram.positive_spans[0].offset = beyond;
1794 } else {
1795 histogram.negative_spans[0].offset = beyond;
1796 }
1797 assert_invalid(
1798 "exponential overflow bucket index",
1799 request_with_histogram(histogram),
1800 &format!("bucket index {beyond} is out of range"),
1801 );
1802 }
1803 }
1804 }
1805
1806 #[test]
1807 fn test_into_context_req_converts_histograms_and_ignores_exemplars() {
1808 let request = Request {
1809 symbols: vec![
1810 "".to_string(),
1811 METRIC_NAME_LABEL.to_string(),
1812 "sample_metric".to_string(),
1813 "histogram_metric".to_string(),
1814 ],
1815 timeseries: vec![
1816 TimeSeries {
1817 labels_refs: vec![1, 2],
1818 samples: vec![Sample {
1819 value: 1.0,
1820 timestamp: 1000,
1821 start_timestamp: 0,
1822 }],
1823 ..Default::default()
1824 },
1825 TimeSeries {
1826 labels_refs: vec![1, 3],
1827 histograms: vec![Histogram::default()],
1828 exemplars: vec![Exemplar::default()],
1829 ..Default::default()
1830 },
1831 ],
1832 };
1833
1834 let ctx_req = decode_test_request(request).unwrap();
1835
1836 assert_eq!(ctx_req.sample_count, 1);
1837 assert_eq!(ctx_req.histogram_count, 1);
1838 assert_eq!(ctx_req.samples.all_req().count(), 1);
1839 assert_eq!(ctx_req.histograms.all_req().count(), 1);
1840 }
1841
1842 #[test]
1843 fn test_into_context_req_converts_histogram_only_series() {
1844 let mut request =
1845 test_util::request_with_labels_and_samples(vec![(METRIC_NAME_LABEL, "metric")], vec![]);
1846 request.timeseries[0].histograms.push(Histogram::default());
1847
1848 let ctx_req = decode_test_request(request).unwrap();
1849
1850 assert_eq!(ctx_req.sample_count, 0);
1851 assert_eq!(ctx_req.histogram_count, 1);
1852 assert_eq!(ctx_req.samples.all_req().count(), 0);
1853 let mut inserts = ctx_req.histograms.all_req().collect::<Vec<_>>();
1854 assert_eq!(inserts.len(), 1);
1855
1856 let request = inserts.pop().unwrap();
1857 assert_eq!(request.table_name, "metric");
1858 let rows = request.rows.unwrap();
1859 assert_eq!(rows.rows.len(), 1);
1860 assert_eq!(
1861 rows.schema
1862 .iter()
1863 .map(|col| col.column_name.as_str())
1864 .collect::<Vec<_>>(),
1865 vec![greptime_timestamp(), greptime_native_histogram()]
1866 );
1867 assert_eq!(
1868 rows.rows[0].values[0].value_data,
1869 Some(ValueData::TimestampMillisecondValue(0))
1870 );
1871 assert_eq!(
1872 histogram_field_value(&rows, 0, SCHEMA_FIELD),
1873 Some(ValueData::I32Value(0))
1874 );
1875 assert_eq!(
1876 histogram_field_value(&rows, 0, COUNT_I64_FIELD),
1877 Some(ValueData::I64Value(0))
1878 );
1879 assert_eq!(histogram_field_value(&rows, 0, COUNT_F64_FIELD), None);
1880 }
1881
1882 #[test]
1883 fn test_into_context_req_preserves_histogram_start_timestamp() {
1884 let ctx_req = decode_test_request(test_util::request_with_labels_and_histograms(
1885 vec![(METRIC_NAME_LABEL, "metric")],
1886 vec![Histogram {
1887 timestamp: 2000,
1888 start_timestamp: 1000,
1889 ..Default::default()
1890 }],
1891 ))
1892 .unwrap();
1893
1894 let mut inserts = ctx_req.histograms.all_req().collect::<Vec<_>>();
1895 let rows = inserts.pop().unwrap().rows.unwrap();
1896
1897 assert_eq!(
1898 histogram_field_value(&rows, 0, START_TIMESTAMP_FIELD),
1899 Some(ValueData::TimestampMillisecondValue(1000))
1900 );
1901 }
1902
1903 #[test]
1904 fn test_into_context_req_preserves_exponential_zero_threshold() {
1905 for zero_threshold in [-1.0, f64::NAN] {
1906 let ctx_req = decode_test_request(request_with_histogram(Histogram {
1907 zero_threshold,
1908 ..Default::default()
1909 }))
1910 .unwrap();
1911 let rows = ctx_req.histograms.all_req().next().unwrap().rows.unwrap();
1912 let Some(ValueData::F64Value(actual)) =
1913 histogram_field_value(&rows, 0, ZERO_THRESHOLD_FIELD)
1914 else {
1915 panic!("expected zero threshold");
1916 };
1917
1918 assert_eq!(zero_threshold.to_bits(), actual.to_bits());
1919 }
1920 }
1921
1922 #[test]
1923 fn test_into_context_req_rejects_internal_histogram_labels() {
1924 let mut request = test_util::request_with_labels_and_samples(
1925 vec![
1926 (METRIC_NAME_LABEL, "metric"),
1927 (greptime_native_histogram(), "user_value"),
1928 ],
1929 vec![],
1930 );
1931 request.timeseries[0].histograms.push(Histogram::default());
1932
1933 let err = match decode_test_request(request) {
1934 Ok(_) => panic!("expected invalid request error"),
1935 Err(err) => err,
1936 };
1937 assert_eq!(
1938 err.to_string(),
1939 "Invalid prometheus remote request, msg: remote write v2 label `greptime_native_histogram` conflicts with an internal native histogram label"
1940 );
1941 }
1942
1943 #[test]
1944 fn test_rejects_legacy_histogram_label_after_prefix_change() {
1945 set_default_prefix(Some("custom")).unwrap();
1946 assert_eq!(greptime_native_histogram(), "custom_native_histogram");
1947
1948 let err = ensure_no_internal_histogram_labels(&vec![(
1949 NATIVE_HISTOGRAM_FIELD,
1950 "user_value".to_string(),
1951 )])
1952 .unwrap_err();
1953 assert!(
1954 err.to_string()
1955 .contains("conflicts with an internal native histogram label")
1956 );
1957 }
1958
1959 #[test]
1960 fn test_into_context_req_converts_int_and_float_histograms_to_one_schema() {
1961 let float_histogram = Histogram {
1962 count: Some(api::greptime_proto::io::prometheus::write::v2::histogram::Count::CountFloat(6.0)),
1963 zero_count: Some(
1964 api::greptime_proto::io::prometheus::write::v2::histogram::ZeroCount::ZeroCountFloat(
1965 0.5,
1966 ),
1967 ),
1968 positive_counts: vec![2.0, 3.5],
1969 positive_spans: vec![api::greptime_proto::io::prometheus::write::v2::BucketSpan {
1970 offset: 3,
1971 length: 2,
1972 }],
1973 timestamp: 2000,
1974 ..Default::default()
1975 };
1976 let request = Request {
1977 symbols: vec![
1978 "".to_string(),
1979 METRIC_NAME_LABEL.to_string(),
1980 "metric".to_string(),
1981 ],
1982 timeseries: vec![
1983 TimeSeries {
1984 labels_refs: vec![1, 2],
1985 histograms: vec![test_util::histogram(1000)],
1986 ..Default::default()
1987 },
1988 TimeSeries {
1989 labels_refs: vec![1, 2],
1990 histograms: vec![float_histogram],
1991 ..Default::default()
1992 },
1993 ],
1994 };
1995
1996 let ctx_req = decode_test_request(request).unwrap();
1997
1998 assert_eq!(ctx_req.histogram_count, 2);
1999 let mut inserts = ctx_req.histograms.all_req().collect::<Vec<_>>();
2000 assert_eq!(inserts.len(), 1);
2001 let rows = inserts.pop().unwrap().rows.unwrap();
2002 assert_eq!(rows.rows.len(), 2);
2003 assert_eq!(
2004 rows.schema
2005 .iter()
2006 .map(|col| col.column_name.as_str())
2007 .collect::<Vec<_>>(),
2008 vec![greptime_timestamp(), greptime_native_histogram()]
2009 );
2010
2011 assert_eq!(
2012 histogram_field_value(&rows, 0, COUNT_I64_FIELD),
2013 Some(ValueData::I64Value(0))
2014 );
2015 assert_eq!(histogram_field_value(&rows, 0, COUNT_F64_FIELD), None);
2016 assert!(matches!(
2017 histogram_field_value(&rows, 0, POSITIVE_BUCKETS_I64_FIELD),
2018 Some(ValueData::ListValue(_))
2019 ));
2020 assert!(is_empty_list(histogram_field_value(
2021 &rows,
2022 0,
2023 POSITIVE_BUCKETS_F64_FIELD
2024 )));
2025
2026 assert_eq!(histogram_field_value(&rows, 1, COUNT_I64_FIELD), None);
2027 assert_eq!(
2028 histogram_field_value(&rows, 1, COUNT_F64_FIELD),
2029 Some(ValueData::F64Value(6.0))
2030 );
2031 assert!(is_empty_list(histogram_field_value(
2032 &rows,
2033 1,
2034 POSITIVE_BUCKETS_I64_FIELD
2035 )));
2036 assert!(matches!(
2037 histogram_field_value(&rows, 1, POSITIVE_BUCKETS_F64_FIELD),
2038 Some(ValueData::ListValue(_))
2039 ));
2040 }
2041
2042 fn decode_wire(wire: &[u8]) -> Result<RemoteWriteV2WriteRequests> {
2043 let body = Bytes::from(crate::prom_store::snappy_compress(wire).unwrap());
2044 decode_v2_on_test_runtime(false, body)
2045 }
2046
2047 fn decode_v2_on_test_runtime(is_zstd: bool, body: Bytes) -> Result<RemoteWriteV2WriteRequests> {
2050 tokio::runtime::Builder::new_current_thread()
2051 .enable_all()
2052 .build()
2053 .unwrap()
2054 .block_on(decode_remote_write_v2(
2055 is_zstd,
2056 body,
2057 &ServerMemoryLimiter::default(),
2058 ))
2059 }
2060
2061 fn decode_wire_error(wire: &[u8], name: &str) -> error::Error {
2062 match decode_wire(wire) {
2063 Ok(_) => panic!("{name}: expected decoder error"),
2064 Err(error) => error,
2065 }
2066 }
2067
2068 fn request_wire(symbols: &[&str], series: &[Vec<u8>]) -> Vec<u8> {
2069 let mut wire = Vec::new();
2070 for symbol in symbols {
2071 wire.extend(string_field(4, symbol.as_bytes()));
2072 }
2073 for series in series {
2074 wire.extend(encoded_message_field(5, series));
2075 }
2076 wire
2077 }
2078
2079 fn series_wire(labels_refs: &[u32], leaf_tag: u32, leaf: &[u8]) -> Vec<u8> {
2080 let mut series = packed_u32_field(1, labels_refs);
2081 series.extend(encoded_message_field(leaf_tag, leaf));
2082 series
2083 }
2084
2085 fn string_field(tag: u32, value: &[u8]) -> Vec<u8> {
2086 encoded_message_field(tag, value)
2087 }
2088
2089 fn encoded_message_field(tag: u32, value: &[u8]) -> Vec<u8> {
2090 let mut field = Vec::new();
2091 prost::encoding::encode_key(tag, WireType::LengthDelimited, &mut field);
2092 prost::encoding::encode_varint(u64::try_from(value.len()).unwrap(), &mut field);
2093 field.extend_from_slice(value);
2094 field
2095 }
2096
2097 fn varint_field(tag: u32, value: u64) -> Vec<u8> {
2098 let mut field = Vec::new();
2099 prost::encoding::encode_key(tag, WireType::Varint, &mut field);
2100 prost::encoding::encode_varint(value, &mut field);
2101 field
2102 }
2103
2104 fn packed_u32_field(tag: u32, values: &[u32]) -> Vec<u8> {
2105 let mut packed = Vec::new();
2106 for value in values {
2107 prost::encoding::encode_varint(u64::from(*value), &mut packed);
2108 }
2109 encoded_message_field(tag, &packed)
2110 }
2111
2112 fn request_with_sample(labels: Vec<(&str, &str)>) -> Request {
2113 test_util::request_with_labels_and_samples(
2114 labels,
2115 vec![Sample {
2116 value: 1.0,
2117 timestamp: 1000,
2118 start_timestamp: 0,
2119 }],
2120 )
2121 }
2122
2123 fn request_with_histogram(histogram: Histogram) -> Request {
2124 test_util::request_with_labels_and_histograms(
2125 vec![(METRIC_NAME_LABEL, "metric")],
2126 vec![histogram],
2127 )
2128 }
2129
2130 fn decode_test_request(request: Request) -> Result<RemoteWriteV2WriteRequests> {
2131 let body =
2132 Bytes::from(crate::prom_store::snappy_compress(&request.encode_to_vec()).unwrap());
2133 decode_v2_on_test_runtime(false, body)
2134 }
2135
2136 fn assert_invalid(name: &str, request: Request, expected: &str) {
2137 let err = match decode_test_request(request) {
2138 Ok(_) => panic!("{name}: expected invalid request error"),
2139 Err(err) => err,
2140 };
2141 assert!(
2142 matches!(err, error::Error::InvalidPromRemoteRequest { .. }),
2143 "{name}: expected invalid request error, got {err}"
2144 );
2145 assert!(
2146 err.to_string().contains(expected),
2147 "{name}: expected error containing {expected:?}, got {err}"
2148 );
2149 }
2150
2151 fn column_index(schema: &[ColumnSchema], column_name: &str) -> usize {
2152 schema
2153 .iter()
2154 .position(|column| column.column_name == column_name)
2155 .unwrap()
2156 }
2157
2158 fn histogram_field_value(rows: &Rows, row_idx: usize, field_name: &str) -> Option<ValueData> {
2159 let histogram_idx = column_index(&rows.schema, greptime_native_histogram());
2160 let Some(ValueData::StructValue(histogram)) =
2161 &rows.rows[row_idx].values[histogram_idx].value_data
2162 else {
2163 panic!("expected native histogram struct value");
2164 };
2165 let field_idx = NATIVE_HISTOGRAM_FIELD_NAMES
2166 .iter()
2167 .position(|name| *name == field_name)
2168 .unwrap();
2169 histogram.items[field_idx].value_data.clone()
2170 }
2171
2172 fn is_empty_list(value: Option<ValueData>) -> bool {
2173 matches!(value, Some(ValueData::ListValue(list)) if list.items.is_empty())
2174 }
2175
2176 fn push_test_symbol(symbols: &mut Vec<String>, symbol: &str) -> u32 {
2177 if let Some(idx) = symbols.iter().position(|s| s == symbol) {
2178 return idx as u32;
2179 }
2180 let idx = symbols.len();
2181 symbols.push(symbol.to_string());
2182 idx as u32
2183 }
2184
2185 type MetadataSeries<'a> = (Vec<(&'a str, &'a str)>, i32, Option<&'a str>);
2187
2188 fn metadata_request(series: Vec<MetadataSeries<'_>>) -> Request {
2189 let mut symbols = vec!["".to_string()];
2190 let timeseries = series
2191 .into_iter()
2192 .map(|(labels, metric_type, unit)| {
2193 let mut labels_refs = Vec::with_capacity(labels.len() * 2);
2194 for (name, value) in labels {
2195 labels_refs.push(push_test_symbol(&mut symbols, name));
2196 labels_refs.push(push_test_symbol(&mut symbols, value));
2197 }
2198 let unit_ref = unit.map_or(0, |unit| push_test_symbol(&mut symbols, unit));
2199 TimeSeries {
2200 labels_refs,
2201 samples: vec![Sample {
2202 value: 1.0,
2203 timestamp: 1000,
2204 start_timestamp: 0,
2205 }],
2206 histograms: Vec::new(),
2207 exemplars: Vec::new(),
2208 metadata: Some(Metadata {
2209 r#type: metric_type,
2210 help_ref: 0,
2211 unit_ref,
2212 }),
2213 }
2214 })
2215 .collect();
2216 Request {
2217 symbols,
2218 timeseries,
2219 }
2220 }
2221
2222 type DecodedIndex = std::collections::BTreeMap<
2223 String,
2224 std::collections::BTreeMap<String, std::collections::BTreeMap<String, String>>,
2225 >;
2226
2227 fn decoded_index(request: Request) -> DecodedIndex {
2228 let encoded = decode_test_request(request)
2229 .unwrap()
2230 .semantic_index
2231 .encode("public")
2232 .expect("non-empty semantic index");
2233 serde_json::from_str(&encoded).unwrap()
2234 }
2235
2236 #[test]
2237 fn test_metadata_stamps_semantic_index() {
2238 use table::requests::METADATA_QUALITY_DECLARED;
2239
2240 let index = decoded_index(metadata_request(vec![
2241 (
2242 vec![(METRIC_NAME_LABEL, "http_requests_total")],
2243 metadata::MetricType::Counter as i32,
2244 Some("seconds"),
2245 ),
2246 (
2247 vec![(METRIC_NAME_LABEL, "queue_depth")],
2248 metadata::MetricType::Gauge as i32,
2249 Some("requests"),
2251 ),
2252 ]));
2253
2254 let tables = &index["public"];
2255 let typed = &tables["http_requests_total"];
2256 assert_eq!(typed[SEMANTIC_METRIC_TYPE], "counter");
2257 assert_eq!(
2258 typed[SEMANTIC_METRIC_METADATA_QUALITY],
2259 METADATA_QUALITY_DECLARED
2260 );
2261 assert_eq!(typed[SEMANTIC_METRIC_UNIT], "s");
2262
2263 let unitless = &tables["queue_depth"];
2264 assert_eq!(unitless[SEMANTIC_METRIC_TYPE], "gauge");
2265 assert!(!unitless.contains_key(SEMANTIC_METRIC_UNIT));
2266 }
2267
2268 #[test]
2269 fn test_metadata_unspecified_stamps_unit_but_not_type() {
2270 let index = decoded_index(metadata_request(vec![(
2271 vec![(METRIC_NAME_LABEL, "untyped_total")],
2272 metadata::MetricType::Unspecified as i32,
2273 Some("seconds"),
2274 )]));
2275 let untyped = &index["public"]["untyped_total"];
2276 assert_eq!(untyped[SEMANTIC_METRIC_UNIT], "s");
2277 assert!(!untyped.contains_key(SEMANTIC_METRIC_TYPE));
2278 assert!(!untyped.contains_key(SEMANTIC_METRIC_METADATA_QUALITY));
2279
2280 let requests = decode_test_request(metadata_request(vec![(
2281 vec![(METRIC_NAME_LABEL, "untyped_unitless_total")],
2282 metadata::MetricType::Unspecified as i32,
2283 None,
2284 )]))
2285 .unwrap();
2286 assert!(requests.semantic_index.is_empty());
2287
2288 let requests = decode_test_request(test_util::request_with_labels_and_samples(
2290 vec![(METRIC_NAME_LABEL, "bare_total")],
2291 vec![Sample {
2292 value: 1.0,
2293 timestamp: 1000,
2294 start_timestamp: 0,
2295 }],
2296 ))
2297 .unwrap();
2298 assert!(requests.semantic_index.is_empty());
2299 }
2300
2301 #[test]
2302 fn test_metadata_type_conflict_collapses_to_mixed() {
2303 let index = decoded_index(metadata_request(vec![
2304 (
2305 vec![(METRIC_NAME_LABEL, "flappy_metric"), ("job", "a")],
2306 metadata::MetricType::Counter as i32,
2307 None,
2308 ),
2309 (
2310 vec![(METRIC_NAME_LABEL, "flappy_metric"), ("job", "b")],
2311 metadata::MetricType::Gauge as i32,
2312 None,
2313 ),
2314 ]));
2315 assert_eq!(
2316 index["public"]["flappy_metric"][SEMANTIC_METRIC_TYPE],
2317 "mixed"
2318 );
2319 }
2320
2321 #[test]
2322 fn test_metadata_schema_overrides_stay_apart() {
2323 let index = decoded_index(metadata_request(vec![
2326 (
2327 vec![(METRIC_NAME_LABEL, "cpu_usage")],
2328 metadata::MetricType::Counter as i32,
2329 None,
2330 ),
2331 (
2332 vec![
2333 (METRIC_NAME_LABEL, "cpu_usage"),
2334 (DATABASE_LABEL, "tenant_b"),
2335 ],
2336 metadata::MetricType::Gauge as i32,
2337 None,
2338 ),
2339 ]));
2340 assert_eq!(
2341 index["public"]["cpu_usage"][SEMANTIC_METRIC_TYPE],
2342 "counter"
2343 );
2344 assert_eq!(
2345 index["tenant_b"]["cpu_usage"][SEMANTIC_METRIC_TYPE],
2346 "gauge"
2347 );
2348 }
2349
2350 #[test]
2351 fn test_metadata_out_of_range_refs_are_rejected() {
2352 let mut request = metadata_request(vec![(
2353 vec![(METRIC_NAME_LABEL, "broken_total")],
2354 metadata::MetricType::Counter as i32,
2355 None,
2356 )]);
2357 request.timeseries[0].metadata.as_mut().unwrap().unit_ref = 999;
2358 let err = decode_test_request(request).err().unwrap();
2359 assert!(err.to_string().contains("out of range"), "{err}");
2360
2361 let mut request = metadata_request(vec![(
2364 vec![(METRIC_NAME_LABEL, "broken_total")],
2365 metadata::MetricType::Unspecified as i32,
2366 None,
2367 )]);
2368 request.timeseries[0].metadata.as_mut().unwrap().unit_ref = 999;
2369 let err = decode_test_request(request).err().unwrap();
2370 assert!(err.to_string().contains("out of range"), "{err}");
2371
2372 let mut request = metadata_request(vec![(
2374 vec![(METRIC_NAME_LABEL, "broken_total")],
2375 metadata::MetricType::Counter as i32,
2376 None,
2377 )]);
2378 request.timeseries[0].metadata.as_mut().unwrap().help_ref = 999;
2379 let err = decode_test_request(request).err().unwrap();
2380 assert!(err.to_string().contains("out of range"), "{err}");
2381 }
2382}