Skip to main content

servers/
metrics.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#[cfg(not(windows))]
16pub(crate) mod jemalloc;
17
18use std::task::{Context, Poll};
19use std::time::Instant;
20
21use axum::extract::{MatchedPath, Request};
22use axum::middleware::Next;
23use axum::response::IntoResponse;
24use lazy_static::lazy_static;
25use prometheus::{
26    Histogram, HistogramVec, IntCounter, IntCounterVec, IntGauge, register_histogram,
27    register_histogram_vec, register_int_counter, register_int_counter_vec, register_int_gauge,
28};
29use session::context::QueryContext;
30use tonic::body::Body;
31use tower::{Layer, Service};
32
33pub(crate) const METRIC_DB_LABEL: &str = "db";
34pub(crate) const METRIC_CODE_LABEL: &str = "code";
35pub(crate) const METRIC_TYPE_LABEL: &str = "type";
36pub(crate) const METRIC_PROTOCOL_LABEL: &str = "protocol";
37pub(crate) const METRIC_ERROR_COUNTER_LABEL_MYSQL: &str = "mysql";
38pub(crate) const METRIC_MYSQL_SUBPROTOCOL_LABEL: &str = "subprotocol";
39pub(crate) const METRIC_MYSQL_BINQUERY: &str = "binquery";
40pub(crate) const METRIC_MYSQL_TEXTQUERY: &str = "textquery";
41pub(crate) const METRIC_POSTGRES_SUBPROTOCOL_LABEL: &str = "subprotocol";
42pub(crate) const METRIC_POSTGRES_SIMPLE_QUERY: &str = "simple";
43pub(crate) const METRIC_POSTGRES_EXTENDED_QUERY: &str = "extended";
44pub(crate) const METRIC_METHOD_LABEL: &str = "method";
45pub(crate) const METRIC_PATH_LABEL: &str = "path";
46pub(crate) const METRIC_RESULT_LABEL: &str = "result";
47pub(crate) const METRIC_VERSION_LABEL: &str = "version";
48
49pub(crate) const METRIC_SUCCESS_VALUE: &str = "success";
50pub(crate) const METRIC_FAILURE_VALUE: &str = "failure";
51
52lazy_static! {
53
54    pub static ref HTTP_REQUEST_COUNTER: IntCounterVec = register_int_counter_vec!(
55        "greptime_servers_http_request_counter",
56        "servers http request counter",
57        &[METRIC_METHOD_LABEL, METRIC_PATH_LABEL, METRIC_CODE_LABEL, METRIC_DB_LABEL]
58    ).unwrap();
59
60    pub static ref METRIC_ERROR_COUNTER: IntCounterVec = register_int_counter_vec!(
61        "greptime_servers_error",
62        "servers error",
63        &[METRIC_PROTOCOL_LABEL]
64    )
65    .unwrap();
66    /// Http SQL query duration per database.
67    pub static ref METRIC_HTTP_SQL_ELAPSED: HistogramVec = register_histogram_vec!(
68        "greptime_servers_http_sql_elapsed",
69        "servers http sql elapsed",
70        &[METRIC_DB_LABEL],
71        vec![0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 60.0, 300.0]
72    )
73    .unwrap();
74    /// Http pql query duration per database.
75    pub static ref METRIC_HTTP_PROMQL_ELAPSED: HistogramVec = register_histogram_vec!(
76        "greptime_servers_http_promql_elapsed",
77        "servers http promql elapsed",
78        &[METRIC_DB_LABEL],
79        vec![0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 60.0, 300.0]
80    )
81    .unwrap();
82    /// Http logs query duration per database.
83    pub static ref METRIC_HTTP_LOGS_ELAPSED: HistogramVec = register_histogram_vec!(
84        "greptime_servers_http_logs_elapsed",
85        "servers http logs elapsed",
86        &[METRIC_DB_LABEL],
87        vec![0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 60.0, 300.0]
88    )
89    .unwrap();
90    pub static ref METRIC_AUTH_FAILURE: IntCounterVec = register_int_counter_vec!(
91        "greptime_servers_auth_failure_count",
92        "servers auth failure count",
93        &[METRIC_CODE_LABEL]
94    )
95    .unwrap();
96    /// Http influxdb write duration per database.
97    pub static ref METRIC_HTTP_INFLUXDB_WRITE_ELAPSED: HistogramVec = register_histogram_vec!(
98        "greptime_servers_http_influxdb_write_elapsed",
99        "servers http influxdb write elapsed",
100        &[METRIC_DB_LABEL],
101        vec![0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 60.0, 300.0]
102    )
103    .unwrap();
104    /// Http prometheus write duration per database and remote write protocol version.
105    pub static ref METRIC_HTTP_PROM_STORE_WRITE_ELAPSED: HistogramVec = register_histogram_vec!(
106        "greptime_servers_http_prometheus_write_elapsed",
107        "servers http prometheus write elapsed",
108        &[METRIC_DB_LABEL, METRIC_VERSION_LABEL],
109        vec![0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 60.0, 300.0]
110    )
111    .unwrap();
112    /// Prometheus remote write codec duration per protocol version.
113    pub static ref METRIC_HTTP_PROM_STORE_CODEC_ELAPSED: HistogramVec = register_histogram_vec!(
114        "greptime_servers_http_prometheus_codec_elapsed",
115        "servers http prometheus request codec duration",
116        &["type", METRIC_VERSION_LABEL],
117    )
118    .unwrap();
119    /// The samples count of Prometheus remote write per protocol version.
120    pub static ref PROM_STORE_REMOTE_WRITE_SAMPLES: IntCounterVec = register_int_counter_vec!(
121        "greptime_servers_prometheus_remote_write_samples",
122        "frontend prometheus remote write samples",
123        &[METRIC_DB_LABEL, METRIC_VERSION_LABEL]
124    )
125    .unwrap();
126    /// The native histograms count of Prometheus remote write per protocol version.
127    pub static ref PROM_STORE_REMOTE_WRITE_HISTOGRAMS: IntCounterVec = register_int_counter_vec!(
128        "greptime_servers_prometheus_remote_write_histograms",
129        "frontend prometheus remote write native histograms",
130        &[METRIC_DB_LABEL, METRIC_VERSION_LABEL]
131    )
132    .unwrap();
133    /// OTLP exponential-histogram data points rejected during conversion, by reason.
134    pub static ref OTLP_EXPONENTIAL_HISTOGRAM_REJECTED_DATA_POINTS: IntCounterVec = register_int_counter_vec!(
135        "greptime_servers_otlp_exponential_histogram_rejected_data_points_total",
136        "Number of OTLP exponential histogram data points rejected during conversion",
137        &["reason"]
138    )
139    .unwrap();
140    pub static ref PENDING_BATCHES: IntGauge = register_int_gauge!(
141        "greptime_prom_store_pending_batches",
142        "Number of pending batches waiting to be flushed"
143    )
144    .unwrap();
145    pub static ref PENDING_ROWS: IntGauge = register_int_gauge!(
146        "greptime_prom_store_pending_rows",
147        "Number of pending rows waiting to be flushed"
148    )
149    .unwrap();
150    pub static ref PENDING_WORKERS: IntGauge = register_int_gauge!(
151        "greptime_prom_store_pending_workers",
152        "Number of active pending rows batch workers"
153    )
154    .unwrap();
155    pub static ref FLUSH_TOTAL: IntCounter = register_int_counter!(
156        "greptime_prom_store_flush_total",
157        "Total number of batch flushes"
158    )
159    .unwrap();
160    pub static ref FLUSH_ROWS: Histogram = register_histogram!(
161        "greptime_prom_store_flush_rows",
162        "Number of rows per flush",
163        vec![100.0, 1000.0, 10000.0, 50000.0, 100000.0, 500000.0]
164    )
165    .unwrap();
166    pub static ref FLUSH_ELAPSED: Histogram = register_histogram!(
167        "greptime_prom_store_flush_elapsed",
168        "Elapsed time of pending rows batch flush in seconds",
169        vec![0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 60.0, 300.0]
170    )
171    .unwrap();
172    pub static ref FLUSH_DROPPED_ROWS: IntCounter = register_int_counter!(
173        "greptime_pending_rows_flush_dropped_rows",
174        "Total rows dropped due to pending rows flush failures"
175    )
176    .unwrap();
177    pub static ref FLUSH_FAILURES: IntCounter = register_int_counter!(
178        "greptime_pending_rows_flush_failures",
179        "Total pending rows flush failures"
180    )
181    .unwrap();
182    pub static ref FLOW_NOTIFICATION_DROPPED: IntCounterVec = register_int_counter_vec!(
183        "greptime_prom_store_flow_notification_dropped_total",
184        "Total flow notifications dropped by the pending rows batcher",
185        &["reason"]
186    )
187    .unwrap();
188    pub static ref PENDING_ROWS_BATCH_INGEST_STAGE_ELAPSED: HistogramVec = register_histogram_vec!(
189        "greptime_prom_store_pending_rows_batch_ingest_stage_elapsed",
190        "Elapsed time of pending rows batch ingestion stages in seconds",
191        &["stage"],
192        vec![0.0005, 0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 60.0]
193    )
194    .unwrap();
195    pub static ref PENDING_ROWS_BATCH_FLUSH_STAGE_ELAPSED: HistogramVec = register_histogram_vec!(
196        "greptime_prom_store_pending_rows_batch_flush_stage_elapsed",
197        "Elapsed time of pending rows batch flush stages in seconds",
198        &["stage"],
199        vec![0.0005, 0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 60.0]
200    )
201    .unwrap();
202    /// Http prometheus read duration per database.
203    pub static ref METRIC_HTTP_PROM_STORE_READ_ELAPSED: HistogramVec = register_histogram_vec!(
204        "greptime_servers_http_prometheus_read_elapsed",
205        "servers http prometheus read elapsed",
206        &[METRIC_DB_LABEL]
207    )
208    .unwrap();
209    /// Http prometheus endpoint query duration per database.
210    pub static ref METRIC_HTTP_PROMETHEUS_PROMQL_ELAPSED: HistogramVec = register_histogram_vec!(
211        "greptime_servers_http_prometheus_promql_elapsed",
212        "servers http prometheus promql elapsed",
213        &[METRIC_DB_LABEL, METRIC_METHOD_LABEL]
214    )
215    .unwrap();
216    pub static ref METRIC_HTTP_OPENTELEMETRY_METRICS_ELAPSED: HistogramVec =
217        register_histogram_vec!(
218            "greptime_servers_http_otlp_metrics_elapsed",
219            "servers_http_otlp_metrics_elapsed",
220            &[METRIC_DB_LABEL]
221        )
222        .unwrap();
223    pub static ref METRIC_HTTP_OPENTELEMETRY_TRACES_ELAPSED: HistogramVec =
224        register_histogram_vec!(
225            "greptime_servers_http_otlp_traces_elapsed",
226            "servers http otlp traces elapsed",
227            &[METRIC_DB_LABEL]
228        )
229        .unwrap();
230    pub static ref METRIC_HTTP_OPENTELEMETRY_LOGS_ELAPSED: HistogramVec =
231    register_histogram_vec!(
232        "greptime_servers_http_otlp_logs_elapsed",
233        "servers http otlp logs elapsed",
234        &[METRIC_DB_LABEL]
235    )
236    .unwrap();
237    pub static ref METRIC_HTTP_LOGS_INGESTION_COUNTER: IntCounterVec = register_int_counter_vec!(
238        "greptime_servers_http_logs_ingestion_counter",
239        "servers http logs ingestion counter",
240        &[METRIC_DB_LABEL]
241    )
242    .unwrap();
243    pub static ref METRIC_HTTP_LOGS_INGESTION_ELAPSED: HistogramVec =
244        register_histogram_vec!(
245            "greptime_servers_http_logs_ingestion_elapsed",
246            "servers http logs ingestion elapsed",
247            &[METRIC_DB_LABEL, METRIC_RESULT_LABEL]
248        )
249        .unwrap();
250
251    /// Count of logs ingested into Loki.
252    pub static ref METRIC_LOKI_LOGS_INGESTION_COUNTER: IntCounterVec = register_int_counter_vec!(
253        "greptime_servers_loki_logs_ingestion_counter",
254        "servers loki logs ingestion counter",
255        &[METRIC_DB_LABEL]
256    )
257    .unwrap();
258    pub static ref METRIC_LOKI_LOGS_INGESTION_ELAPSED: HistogramVec =
259        register_histogram_vec!(
260            "greptime_servers_loki_logs_ingestion_elapsed",
261            "servers loki logs ingestion elapsed",
262            &[METRIC_DB_LABEL, METRIC_RESULT_LABEL]
263        )
264        .unwrap();
265    pub static ref METRIC_ELASTICSEARCH_LOGS_INGESTION_ELAPSED: HistogramVec =
266        register_histogram_vec!(
267            "greptime_servers_elasticsearch_logs_ingestion_elapsed",
268            "servers elasticsearch logs ingestion elapsed",
269            &[METRIC_DB_LABEL]
270        )
271        .unwrap();
272
273    /// Count of documents ingested into Elasticsearch logs.
274    pub static ref METRIC_ELASTICSEARCH_LOGS_DOCS_COUNT: IntCounterVec = register_int_counter_vec!(
275        "greptime_servers_elasticsearch_logs_docs_count",
276        "servers elasticsearch ingest logs docs count",
277        &[METRIC_DB_LABEL]
278    )
279    .unwrap();
280
281    pub static ref METRIC_HTTP_LOGS_TRANSFORM_ELAPSED: HistogramVec =
282        register_histogram_vec!(
283            "greptime_servers_http_logs_transform_elapsed",
284            "servers http logs transform elapsed",
285            &[METRIC_DB_LABEL, METRIC_RESULT_LABEL]
286        )
287        .unwrap();
288    pub static ref METRIC_MYSQL_CONNECTIONS: IntGauge = register_int_gauge!(
289        "greptime_servers_mysql_connection_count",
290        "servers mysql connection count"
291    )
292    .unwrap();
293    pub static ref METRIC_MYSQL_QUERY_TIMER: HistogramVec = register_histogram_vec!(
294        "greptime_servers_mysql_query_elapsed",
295        "servers mysql query elapsed",
296        &[METRIC_MYSQL_SUBPROTOCOL_LABEL, METRIC_DB_LABEL],
297        vec![0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 60.0, 300.0]
298    )
299    .unwrap();
300    pub static ref METRIC_MYSQL_PREPARED_COUNT: IntCounterVec = register_int_counter_vec!(
301        "greptime_servers_mysql_prepared_count",
302        "servers mysql prepared count",
303        &[METRIC_DB_LABEL]
304    )
305    .unwrap();
306    pub static ref METRIC_POSTGRES_CONNECTIONS: IntGauge = register_int_gauge!(
307        "greptime_servers_postgres_connection_count",
308        "servers postgres connection count"
309    )
310    .unwrap();
311    pub static ref METRIC_POSTGRES_QUERY_TIMER: HistogramVec = register_histogram_vec!(
312        "greptime_servers_postgres_query_elapsed",
313        "servers postgres query elapsed",
314        &[METRIC_POSTGRES_SUBPROTOCOL_LABEL, METRIC_DB_LABEL],
315        vec![0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 60.0, 300.0]
316    )
317    .unwrap();
318    pub static ref METRIC_POSTGRES_PREPARED_COUNT: IntCounter = register_int_counter!(
319        "greptime_servers_postgres_prepared_count",
320        "servers postgres prepared count"
321    )
322    .unwrap();
323    pub static ref METRIC_SERVER_GRPC_DB_REQUEST_TIMER: HistogramVec = register_histogram_vec!(
324        "greptime_servers_grpc_db_request_elapsed",
325        "servers grpc db request elapsed",
326        &[METRIC_DB_LABEL, METRIC_TYPE_LABEL, METRIC_CODE_LABEL]
327    )
328    .unwrap();
329    pub static ref METRIC_SERVER_GRPC_PROM_REQUEST_TIMER: HistogramVec = register_histogram_vec!(
330        "greptime_servers_grpc_prom_request_elapsed",
331        "servers grpc prom request elapsed",
332        &[METRIC_DB_LABEL],
333        vec![0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 60.0, 300.0]
334    )
335    .unwrap();
336    pub static ref METRIC_HTTP_REQUESTS_TOTAL: IntCounterVec = register_int_counter_vec!(
337        "greptime_servers_http_requests_total",
338        "servers http requests total",
339        &[METRIC_METHOD_LABEL, METRIC_PATH_LABEL, METRIC_CODE_LABEL, METRIC_DB_LABEL]
340    )
341    .unwrap();
342    pub static ref METRIC_HTTP_REQUESTS_ELAPSED: HistogramVec = register_histogram_vec!(
343        "greptime_servers_http_requests_elapsed",
344        "servers http requests elapsed",
345        &[METRIC_METHOD_LABEL, METRIC_PATH_LABEL, METRIC_CODE_LABEL, METRIC_DB_LABEL],
346        vec![0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 60.0, 300.0]
347    )
348    .unwrap();
349    pub static ref METRIC_GRPC_REQUESTS_TOTAL: IntCounterVec = register_int_counter_vec!(
350        "greptime_servers_grpc_requests_total",
351        "servers grpc requests total",
352        &[METRIC_PATH_LABEL, METRIC_CODE_LABEL]
353    )
354    .unwrap();
355    pub static ref METRIC_GRPC_REQUESTS_ELAPSED: HistogramVec = register_histogram_vec!(
356        "greptime_servers_grpc_requests_elapsed",
357        "servers grpc requests elapsed",
358        &[METRIC_PATH_LABEL, METRIC_CODE_LABEL],
359        vec![0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 60.0, 300.0]
360    )
361    .unwrap();
362    pub static ref METRIC_JAEGER_QUERY_ELAPSED: HistogramVec = register_histogram_vec!(
363        "greptime_servers_jaeger_query_elapsed",
364        "servers jaeger query elapsed",
365        &[METRIC_DB_LABEL, METRIC_PATH_LABEL]
366    ).unwrap();
367
368    pub static ref GRPC_BULK_INSERT_ELAPSED: Histogram = register_histogram!(
369        "greptime_servers_bulk_insert_elapsed",
370        "servers handle bulk insert elapsed",
371    ).unwrap();
372
373    // Unified request memory metrics
374    /// Current memory in use by all concurrent requests (HTTP, gRPC, Flight).
375    pub static ref REQUEST_MEMORY_IN_USE: IntGauge = register_int_gauge!(
376        "greptime_servers_request_memory_in_use_bytes",
377        "bytes currently reserved for all concurrent request bodies and messages"
378    ).unwrap();
379
380    /// Maximum configured memory for all concurrent requests.
381    pub static ref REQUEST_MEMORY_LIMIT: IntGauge = register_int_gauge!(
382        "greptime_servers_request_memory_limit_bytes",
383        "maximum bytes allowed for all concurrent request bodies and messages"
384    ).unwrap();
385
386    /// Total number of requests rejected due to memory exhaustion.
387    pub static ref REQUEST_MEMORY_REJECTED: IntCounterVec = register_int_counter_vec!(
388        "greptime_servers_request_memory_rejected_total",
389        "number of requests rejected due to memory limit",
390        &["reason"]
391    ).unwrap();
392}
393
394// Based on https://github.com/hyperium/tonic/blob/master/examples/src/tower/server.rs
395// See https://github.com/hyperium/tonic/issues/242
396/// A metrics middleware.
397#[derive(Debug, Clone, Default)]
398pub(crate) struct MetricsMiddlewareLayer;
399
400impl<S> Layer<S> for MetricsMiddlewareLayer {
401    type Service = MetricsMiddleware<S>;
402
403    fn layer(&self, service: S) -> Self::Service {
404        MetricsMiddleware { inner: service }
405    }
406}
407
408#[derive(Debug, Clone)]
409pub(crate) struct MetricsMiddleware<S> {
410    inner: S,
411}
412
413impl<S> Service<http::Request<Body>> for MetricsMiddleware<S>
414where
415    S: Service<http::Request<Body>, Response = http::Response<Body>> + Clone + Send + 'static,
416    S::Future: Send + 'static,
417{
418    type Response = S::Response;
419    type Error = S::Error;
420    type Future = futures::future::BoxFuture<'static, Result<Self::Response, Self::Error>>;
421
422    fn poll_ready(&mut self, cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
423        self.inner.poll_ready(cx)
424    }
425
426    fn call(&mut self, req: http::Request<Body>) -> Self::Future {
427        // This is necessary because tonic internally uses `tower::buffer::Buffer`.
428        // See https://github.com/tower-rs/tower/issues/547#issuecomment-767629149
429        // for details on why this is necessary
430        let clone = self.inner.clone();
431        let mut inner = std::mem::replace(&mut self.inner, clone);
432
433        Box::pin(async move {
434            let start = Instant::now();
435            let path = req.uri().path().to_string();
436
437            // Do extra async work here...
438            let response = inner.call(req).await?;
439
440            let latency = start.elapsed().as_secs_f64();
441            let status = response.status().as_u16().to_string();
442
443            let labels = [path.as_str(), status.as_str()];
444            METRIC_GRPC_REQUESTS_TOTAL.with_label_values(&labels).inc();
445            METRIC_GRPC_REQUESTS_ELAPSED
446                .with_label_values(&labels)
447                .observe(latency);
448
449            Ok(response)
450        })
451    }
452}
453
454/// A middleware to record metrics for HTTP.
455// Based on https://github.com/tokio-rs/axum/blob/axum-v0.6.16/examples/prometheus-metrics/src/main.rs
456pub(crate) async fn http_metrics_layer(req: Request, next: Next) -> impl IntoResponse {
457    let start = Instant::now();
458    let path = if let Some(matched_path) = req.extensions().get::<MatchedPath>() {
459        matched_path.as_str().to_string()
460    } else {
461        req.uri().path().to_string()
462    };
463    let method = req.method().clone();
464
465    let db = req
466        .extensions()
467        .get::<QueryContext>()
468        .map(|ctx| ctx.get_db_string())
469        .unwrap_or_else(|| "unknown".to_string());
470
471    let response = next.run(req).await;
472
473    let latency = start.elapsed().as_secs_f64();
474    let status = response.status();
475    let status = status.as_str();
476    let method_str = method.as_str();
477
478    let labels = [method_str, &path, status, db.as_str()];
479    METRIC_HTTP_REQUESTS_TOTAL.with_label_values(&labels).inc();
480    METRIC_HTTP_REQUESTS_ELAPSED
481        .with_label_values(&labels)
482        .observe(latency);
483
484    response
485}