servers/prom_remote_write/
mod.rs1pub 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
46pub 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
56pub(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 let buf = decompress_remote_write_body(is_zstd, &body[..], limiter).await?;
116 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}