Skip to main content

servers/prom_remote_write/
mod.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
15//! Prometheus remote write support.
16//!
17//! This module groups validation, row building, and protobuf decoding for
18//! the Prometheus remote write API.
19
20pub mod decode;
21pub(crate) mod row_builder;
22pub(crate) mod types;
23#[cfg(any(test, feature = "testing"))]
24pub mod v2;
25#[cfg(not(any(test, feature = "testing")))]
26pub(crate) mod v2;
27pub mod validation;
28
29use bytes::Bytes;
30use common_memory_manager::MemoryGuard;
31use lazy_static::lazy_static;
32use object_pool::Pool;
33use snafu::ResultExt;
34
35use crate::error;
36use crate::prom_remote_write::decode::{PromSeriesProcessor, PromWriteRequest};
37use crate::prom_remote_write::row_builder::TablesBuilder;
38use crate::prom_remote_write::validation::PromValidationMode;
39use crate::prom_store::{
40    ChargedBuffer, MAX_DECOMPRESSED_REQUEST_SIZE, snappy_decompress_limited,
41    zstd_decompress_limited,
42};
43use crate::request_memory_limiter::ServerMemoryLimiter;
44use crate::request_memory_metrics::RequestMemoryMetrics;
45
46/// Prometheus remote write protocol versions, also used as the `version` label
47/// of the remote write metrics.
48pub const REMOTE_WRITE_V1_VERSION: &str = "1.0";
49pub const REMOTE_WRITE_V2_VERSION: &str = "2.0";
50
51lazy_static! {
52    static ref PROM_WRITE_REQUEST_POOL: Pool<PromWriteRequest<'static>> =
53        Pool::new(256, PromWriteRequest::default);
54}
55
56/// Decompresses a remote-write body with the codec selected by `is_zstd`,
57/// bounded by the hard decoded-size cap and charged to the aggregate
58/// request-memory limiter.
59///
60/// Due to vmagent's limitation there is a chance that vmagent sends the
61/// content-encoding header wrong, so decoding falls back to the other codec
62/// when the first attempt fails on malformed input. Size-limit and quota
63/// errors are final: the body either expanded beyond the cap under the
64/// declared codec, or it could not be admitted, and the other codec cannot
65/// produce a different outcome.
66///
67/// see <https://github.com/VictoriaMetrics/VictoriaMetrics/issues/5301>
68/// see <https://github.com/GreptimeTeam/greptimedb/issues/3929>
69pub(crate) async fn decompress_remote_write_body(
70    is_zstd: bool,
71    body: &[u8],
72    limiter: &ServerMemoryLimiter,
73) -> crate::error::Result<ChargedBuffer> {
74    let decompress = |is_zstd: bool| async move {
75        if is_zstd {
76            zstd_decompress_limited(body, MAX_DECOMPRESSED_REQUEST_SIZE, limiter).await
77        } else {
78            snappy_decompress_limited(body, MAX_DECOMPRESSED_REQUEST_SIZE, limiter).await
79        }
80    };
81
82    match decompress(is_zstd).await {
83        Ok(buf) => Ok(buf),
84        Err(e) => {
85            if matches!(
86                e,
87                crate::error::Error::DecompressSnappyPromRemoteRequest { .. }
88                    | crate::error::Error::DecompressZstdPromRemoteRequest { .. }
89            ) {
90                decompress(!is_zstd).await
91            } else {
92                Err(e)
93            }
94        }
95    }
96}
97
98pub async fn decode_remote_write_request(
99    is_zstd: bool,
100    body: Bytes,
101    prom_validation_mode: PromValidationMode,
102    processor: &mut PromSeriesProcessor,
103    limiter: &ServerMemoryLimiter,
104) -> crate::error::Result<(
105    TablesBuilder<'static>,
106    Vec<MemoryGuard<RequestMemoryMetrics>>,
107)> {
108    let _timer = crate::metrics::METRIC_HTTP_PROM_STORE_CODEC_ELAPSED
109        .with_label_values(&["decode", REMOTE_WRITE_V1_VERSION])
110        .start_timer();
111
112    // Holds the memory permits for the decompressed bytes; the caller must
113    // keep them alive as long as the returned `TablesBuilder`, which retains
114    // the decompressed buffer as its raw data.
115    let buf = decompress_remote_write_body(is_zstd, &body[..], limiter).await?;
116    // Decompression copied the payload out, so the compressed body is no longer needed.
117    drop(body);
118    let (data, guards) = buf.into_parts();
119    let mut request = PROM_WRITE_REQUEST_POOL.pull(PromWriteRequest::default);
120
121    request
122        .decode(data, prom_validation_mode, processor)
123        .context(error::DecodePromRemoteRequestSnafu)?;
124    let tables = std::mem::take(&mut request.table_data);
125    Ok((tables, guards))
126}