micromegas/servers/
firehose.rs1use super::firehose_common::{firehose_auth_middleware, firehose_response, request_id_from};
22use super::ingestion_limits::apply_ingestion_body_limits;
23use axum::Extension;
24use axum::Router;
25use axum::http::{HeaderMap, StatusCode};
26use axum::middleware;
27use axum::response::Response;
28use axum::routing::post;
29use micromegas_auth::types::AuthProvider;
30use micromegas_ingestion::web_ingestion_service::WebIngestionService;
31use micromegas_otel_ingestion::{Signal, handler};
32use micromegas_tracing::prelude::*;
33use std::sync::Arc;
34
35async fn firehose_handler(
36 Extension(service): Extension<Arc<WebIngestionService>>,
37 headers: HeaderMap,
38 body: bytes::Bytes,
39) -> Response {
40 let mut request_id = request_id_from(&headers);
41 let envelope = match handler::decode_firehose_envelope(&body, Signal::Metrics) {
42 Ok(e) => e,
43 Err(err) => {
44 error!("firehose decode error (request_id={request_id}): {err}");
45 let status = StatusCode::from_u16(err.http_status())
46 .unwrap_or(StatusCode::INTERNAL_SERVER_ERROR);
47 return firehose_response(status, &request_id, Some(&err.public_message()));
48 }
49 };
50 if request_id.is_empty() {
51 request_id = envelope.request_id.clone(); }
53 match handler::ingest_firehose_metrics(service, envelope.records).await {
54 Ok(()) => firehose_response(StatusCode::OK, &request_id, None),
55 Err(err) => {
56 error!("firehose ingest error (request_id={request_id}): {err}");
57 let status = StatusCode::from_u16(err.http_status())
58 .unwrap_or(StatusCode::INTERNAL_SERVER_ERROR);
59 firehose_response(status, &request_id, Some(&err.public_message()))
60 }
61 }
62}
63
64pub fn firehose_router(
72 service: Arc<WebIngestionService>,
73 auth_provider: Option<Arc<dyn AuthProvider>>,
74) -> Router {
75 let mut router = Router::new()
76 .route(
77 "/ingestion/otlp/v1/metrics/firehose",
78 post(firehose_handler),
79 )
80 .layer(Extension(service));
81 if let Some(provider) = auth_provider {
82 router = router.layer(middleware::from_fn(move |req, next| {
83 firehose_auth_middleware(provider.clone(), req, next)
84 }));
85 }
86 apply_ingestion_body_limits(router)
87}