Skip to main content

servers/prom_remote_write/
v2.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::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    /// Per-table semantic metadata from the series' inline `Metadata`, folded
127    /// into table options at auto-create time.
128    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    // Holds the memory permits for the decompressed bytes until the protobuf
141    // decoding and conversion below are finished.
142    let buf = decompress_remote_write_body(is_zstd, &body[..], limiter).await?;
143    // Decompression copied the payload out, so the compressed body is no longer needed.
144    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
255/// Stamps the series' inline metadata for the written table's auto-create: an
256/// explicit metric type upgrades the table's metadata quality to `declared`,
257/// `UNSPECIFIED` series keep the request-level `inferred` stamp, and units are
258/// canonicalised from OpenMetrics words to UCUM. Type and unit stamp
259/// independently, as OpenMetrics defines them. Help text is not persisted.
260///
261/// Every non-zero symbol reference is validated up front, independent of what
262/// ends up persisted: the spec requires all references to point into the
263/// symbol table.
264fn 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    // Symbol 0 is the mandatory empty string: no help / no unit.
272    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
305/// The `greptime.semantic.metric.type` value for a wire metric type; `None`
306/// for `UNSPECIFIED` (nothing was declared) and out-of-range values.
307fn 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    // Persist both int and float families into the logical table schema. Only one
606    // family is populated per row; the other is written as NULL so PromQL can
607    // infer the original histogram flavor without a separate type column.
608    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    // The histogram field column is generated from the protobuf payload.
630    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    /// Runs the charged decompression on a throwaway runtime so plain `#[test]`
814    /// callers can stay synchronous.
815    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    /// Runs the async (charged) v2 decoder on a throwaway runtime so plain
2048    /// `#[test]` callers can stay synchronous.
2049    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    /// One sample-carrying test series: `(labels, metric_type, unit)`.
2186    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                // Outside the OpenMetrics base set: dropped, not passed through.
2250                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        // No metadata at all behaves the same.
2289        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        // The same metric name written into two schemas by one request must not
2324        // collapse each other's metadata.
2325        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        // unit_ref must be validated even when the type is UNSPECIFIED
2362        // (which persists nothing).
2363        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        // help_ref is validated although help is never persisted.
2373        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}