1#[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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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#[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 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 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
454pub(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}