pre-metrics

This commit is contained in:
2026-05-15 19:35:06 +02:00
parent 0b441b0759
commit 132057217d
43 changed files with 1170 additions and 170 deletions
+68
View File
@@ -0,0 +1,68 @@
use super::error::HTTPError;
use crate::models::user;
use axum::extract::FromRequestParts;
use axum::http::request::Parts;
use std::ops::Deref;
use std::time::Instant;
use uuid::Uuid;
#[derive(Clone, Debug)]
pub struct RequestContext {
pub request_id: Uuid,
pub started_at: Instant,
pub method: axum::http::Method,
pub uri: axum::http::Uri,
pub user: Option<CurrentUser>,
}
/// Représente l'utilisateur actuellement authentifié, enveloppant le modèle de base de données.
///
/// **Philosophie (vs Django) :**
/// Au lieu de passer une `request` entière, on demande `user: CurrentUser` dans la signature
/// de la vue. C'est un **contrat** : si l'utilisateur n'est pas là, la vue n'est pas appelée (401).
///
/// **Usage :**
/// ```rust
/// pub async fn ma_vue(user: CurrentUser) {
/// if user.is_superuser { ... }
/// }
/// ```
#[derive(Clone, Debug)]
pub struct CurrentUser(pub user::Model);
impl Deref for CurrentUser {
type Target = user::Model;
/// Permet d'accéder aux champs du modèle (`user.username`, etc.) directement
/// sur l'objet `CurrentUser`, sans avoir à faire `user.0.username`.
fn deref(&self) -> &Self::Target {
&self.0
}
}
/// Le "Moteur" derrière la magie des signatures de fonction d'Axum.
///
/// Implémenter ce trait permet à Axum d'extraire automatiquement l'utilisateur
/// depuis les extensions de la requête (injectées par le middleware).
impl<S> FromRequestParts<S> for CurrentUser
where
S: Send + Sync,
{
type Rejection = HTTPError;
/// Cette méthode est appelée par Axum AVANT d'exécuter votre vue.
/// 1. On cherche le `RequestContext` dans les extensions.
/// 2. On vérifie si un utilisateur y est présent.
/// 3. Si oui, on le retourne (succès).
/// 4. Si non, on retourne une erreur 401 (rejet de la requête).
async fn from_request_parts(parts: &mut Parts, _state: &S) -> Result<Self, Self::Rejection> {
// On récupère le contexte injecté par le middleware
let context = parts
.extensions
.get::<RequestContext>()
.ok_or(HTTPError::Unauthorized)?;
// On retourne l'utilisateur cloné s'il existe, sinon on rejette avec Unauthorized
context.user.clone().ok_or(HTTPError::Unauthorized)
}
}
+97
View File
@@ -0,0 +1,97 @@
use axum::http::StatusCode;
use axum::response::{IntoResponse, Response};
use axum::Json;
use sea_orm::DbErr;
use serde::Serialize;
use serde_json::json;
use std::collections::HashMap;
use thiserror::Error;
#[derive(Debug, Serialize)]
pub struct ValidationErrorResponse {
pub errors: HashMap<String, Vec<String>>,
}
#[derive(Debug, Error)]
pub enum HTTPError {
#[error("Database error: {0}")]
Database(#[from] DbErr),
#[error("Resource not found")]
NotFound,
#[error("Validation failed")]
UnprocessableEntity(HashMap<String, Vec<String>>),
#[error("Bad request: {0}")]
BadRequest(String),
#[error("Internal server error: {0}")]
InternalServerError(String),
#[error("Unauthorized")]
Unauthorized,
#[error("Forbidden")]
Forbidden,
#[error("Invalid UUID: {0}")]
UuidError(#[from] uuid::Error),
#[error("Internal server error: {0}")]
Internal(#[from] anyhow::Error),
}
// Implémentation pour Axum : transformer AppError en réponse HTTP
impl IntoResponse for HTTPError {
fn into_response(self) -> Response {
let (status, error_message) = match self {
HTTPError::Database(err) => {
eprintln!("Database error: {:?}", err);
(StatusCode::INTERNAL_SERVER_ERROR, "Database error")
}
HTTPError::NotFound => (StatusCode::NOT_FOUND, "Resource not found"),
HTTPError::Unauthorized => (StatusCode::UNAUTHORIZED, "Unauthorized"),
HTTPError::Forbidden => (StatusCode::FORBIDDEN, "Forbidden"),
HTTPError::UuidError(err) => {
return (
StatusCode::BAD_REQUEST,
Json(json!({ "error": format!("Invalid UUID: {}", err) })),
)
.into_response();
}
HTTPError::UnprocessableEntity(errors) => {
return (
StatusCode::UNPROCESSABLE_ENTITY,
Json(ValidationErrorResponse { errors }),
)
.into_response();
}
HTTPError::BadRequest(msg) => {
return (StatusCode::BAD_REQUEST, Json(json!({ "error": msg }))).into_response();
}
HTTPError::InternalServerError(msg) => {
eprintln!("Internal error: {}", msg);
return (
StatusCode::INTERNAL_SERVER_ERROR,
Json(json!({ "error": msg })),
)
.into_response();
}
HTTPError::Internal(err) => {
tracing::error!(error = ?err, "An unexpected error occurred");
(StatusCode::INTERNAL_SERVER_ERROR, "Internal server error")
}
};
(status, Json(json!({ "error": error_message }))).into_response()
}
}
impl HTTPError {
pub fn validation_error(field: &str, message: &str) -> Self {
let mut errors = HashMap::new();
errors.insert(field.to_string(), vec![message.to_string()]);
HTTPError::UnprocessableEntity(errors)
}
}
+69
View File
@@ -0,0 +1,69 @@
use axum::{extract::State, http::Request, middleware::Next, response::Response};
use std::time::Instant;
use tracing::info;
use uuid::Uuid;
use super::context::{CurrentUser, RequestContext};
use crate::auth::token::verify_jwt;
use crate::core::AppState;
pub async fn context_middleware(
State(app_state): State<AppState>,
mut req: Request<axum::body::Body>,
next: Next,
) -> Response {
let request_id = Uuid::new_v4();
let started_at = Instant::now();
// Infos "type Django request"
let method = req.method().clone();
let uri = req.uri().clone();
// Authentification par JWT
let user: Option<CurrentUser> = match req
.headers()
.get(axum::http::header::AUTHORIZATION)
.and_then(|v| v.to_str().ok())
.and_then(|auth_header| {
if auth_header.starts_with("Bearer ") {
Some(&auth_header[7..])
} else {
None
}
})
.and_then(|token| verify_jwt(token, &app_state.config.jwt.secret).ok())
{
Some(claims) => app_state
.repositories
.user
.get_by_id(claims.user_id)
.await
.ok()
.flatten()
.map(CurrentUser),
None => None,
};
let user_id = user.as_ref().map(|u| u.id);
// Injecte le contexte dans la requête (espace de stockage partagé)
// C'est ce qui permettra aux extracteurs comme 'CurrentUser' de retrouver ces données plus tard.
req.extensions_mut().insert(RequestContext {
request_id,
started_at,
method: method.clone(),
uri: uri.clone(),
user,
});
info!(
request_id = %request_id,
user_id = ?user_id,
method = %method,
uri = %uri,
"Incoming request"
);
// Passe la requête au reste de la stack
next.run(req).await
}
+10
View File
@@ -1 +1,11 @@
use crate::core::AppState;
use axum::Router;
pub mod context;
pub mod error;
pub mod metrics;
mod middleware;
pub mod server;
pub mod validation;
pub type OxRouter = Router<AppState>;
+157 -12
View File
@@ -1,23 +1,168 @@
use std::net::SocketAddr;
//! Serveur HTTP axum.
//!
//! [`HttpServer`] encapsule la configuration réseau, l'état applicatif,
//! les métriques et les layers Tower. Il est construit dans [`crate::core::App`]
//! et démarré via [`HttpServer::run`].
use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use axum::middleware as axum_middleware;
use axum::Router;
use tokio::net::TcpListener;
use tokio::sync::broadcast;
use tower_http::catch_panic::CatchPanicLayer;
use tower_http::cors::CorsLayer;
use tower_http::trace::TraceLayer;
use crate::config::AppConfig;
use crate::config::{AppConfig, NetworkConfig};
use crate::core::AppState;
use crate::http::OxRouter;
use crate::routes;
pub async fn start(config: &AppConfig) -> Result<(), Box<dyn std::error::Error>> {
let addr = SocketAddr::new(
std::net::IpAddr::V4(config.network.host),
config.network.tcp_port,
);
use super::metrics::{self, HttpMetrics};
use super::middleware::context_middleware;
let app: Router = routes::router();
// ── Erreurs ───────────────────────────────────────────────────────────────────
let listener = TcpListener::bind(addr).await?;
tracing::info!(%addr, "HTTP server listening");
/// Erreurs pouvant survenir pendant l'opération du serveur HTTP.
#[derive(Debug, thiserror::Error)]
pub enum HttpServerError {
#[error("failed to bind TCP listener to {addr}: {source}")]
Bind {
addr: SocketAddr,
#[source]
source: std::io::Error,
},
#[error("I/O error: {0}")]
Io(#[from] std::io::Error),
}
axum::serve(listener, app).await?;
// ── Struct ────────────────────────────────────────────────────────────────────
Ok(())
/// Serveur HTTP asynchrone basé sur axum.
///
/// - Injecte [`AppState`] via le middleware (les handlers l'obtiendront via
/// `State<AppState>` quand ils seront implémentés ; `.with_state()` sera
/// ajouté en conséquence).
/// - Applique [`context_middleware`] (authentification JWT + contexte de
/// requête) sur l'ensemble du router.
/// - Empile les layers Tower : [`CatchPanicLayer`], [`CorsLayer`] (permissif),
/// [`TraceLayer`].
/// - Collecte des métriques via [`HttpMetrics`] et les reporte périodiquement.
/// - Supporte un shutdown gracieux via un [`broadcast::Sender`].
///
/// # Exemple
/// ```no_run
/// use oxspeak_server_lib::config::AppConfig;
/// use oxspeak_server_lib::core::{App, AppState};
/// use oxspeak_server_lib::http::server::HttpServer;
///
/// #[tokio::main]
/// async fn main() {
/// let config = AppConfig::load().unwrap();
/// // AppState construit via App::build(config)
/// }
/// ```
pub struct HttpServer {
bind_addr: SocketAddr,
app_state: AppState,
metrics: Arc<HttpMetrics>,
shutdown_rx: broadcast::Receiver<()>,
}
impl HttpServer {
/// Construit un [`HttpServer`] depuis la configuration et l'état applicatif.
///
/// Retourne le serveur et un [`broadcast::Sender`] pour déclencher le
/// shutdown gracieux depuis l'extérieur.
pub fn new(
network_config: &NetworkConfig,
app_state: AppState,
) -> (Self, broadcast::Sender<()>) {
let bind_addr = SocketAddr::new(network_config.host.into(), network_config.tcp_port);
let metrics = HttpMetrics::new();
let (shutdown_tx, shutdown_rx) = broadcast::channel(1);
(
Self {
bind_addr,
app_state,
metrics,
shutdown_rx,
},
shutdown_tx,
)
}
/// Retourne une référence aux métriques du serveur.
pub fn metrics(&self) -> &Arc<HttpMetrics> {
&self.metrics
}
/// Bind le listener TCP et démarre la boucle de service.
///
/// La future se résout lorsqu'un signal de shutdown est reçu ou qu'une
/// erreur I/O fatale survient.
pub async fn run(mut self) -> Result<(), HttpServerError> {
// Lance le reporter de métriques toutes les 30 secondes
metrics::spawn_reporter(self.metrics.clone(), Duration::from_secs(30));
let metrics = self.metrics.clone();
let app_state = self.app_state.clone();
// Construit le router avec state + middleware + layers Tower.
//
// Ordre d'application en axum : chaque `.layer()` enveloppe le service
// courant — le dernier appel devient la couche la plus externe.
//
// CatchPanicLayer (outermost — intercepte tout panic en dessous)
// └─ CorsLayer (répond aux preflight avant auth)
// └─ TraceLayer (span tracing sur la durée totale)
// └─ context_middleware (JWT + RequestContext + métriques)
// └─ routes (handlers)
let app: Router = routes::router()
.with_state(app_state.clone())
// Innermost : auth JWT, injection contexte, métriques
.layer(axum_middleware::from_fn_with_state(
app_state.clone(),
move |state, req, next| {
let metrics = metrics.clone();
async move {
metrics.inc_request();
let started = std::time::Instant::now();
let response = context_middleware(state, req, next).await;
let latency = started.elapsed();
metrics.record_response(response.status().as_u16(), latency);
response
}
},
))
// Spans tracing par requête (method, uri, status, latency)
.layer(TraceLayer::new_for_http())
// CORS permissif (à affiner en production)
.layer(CorsLayer::permissive())
// Outermost : intercepte les panics et retourne une 500 propre
.layer(CatchPanicLayer::new());
let listener =
TcpListener::bind(self.bind_addr)
.await
.map_err(|source| HttpServerError::Bind {
addr: self.bind_addr,
source,
})?;
tracing::info!(addr = %self.bind_addr, "HTTP server listening");
axum::serve(listener, app)
.with_graceful_shutdown(async move {
let _ = self.shutdown_rx.recv().await;
tracing::info!("HTTP server shutting down");
})
.await?;
Ok(())
}
}
+63
View File
@@ -0,0 +1,63 @@
use axum::{
extract::{FromRequest, Request},
http::StatusCode,
response::{IntoResponse, Response},
Json,
};
use serde::de::DeserializeOwned;
use serde::Serialize;
use std::collections::HashMap;
use validator::Validate;
#[derive(Debug, Serialize)]
pub struct ValidationErrorResponse {
pub errors: HashMap<String, Vec<String>>,
}
/// Extracteur qui valide le JSON entrant à l'aide de la crate `validator`.
/// Si la validation échoue (JSON malformé ou règles métier violées),
/// il retourne une erreur 400 avec le détail par champ.
pub struct ValidatedJson<T>(pub T);
impl<S, T> FromRequest<S> for ValidatedJson<T>
where
S: Send + Sync,
T: DeserializeOwned + Validate + 'static,
{
type Rejection = Response;
async fn from_request(req: Request, state: &S) -> Result<Self, Self::Rejection> {
// 1. Tente d'extraire le JSON.
// Si le JSON est syntaxiquement invalide, on retourne l'erreur d'Axum directement.
let Json(value) = Json::<T>::from_request(req, state)
.await
.map_err(|rejection| rejection.into_response())?;
// 2. Exécute la validation définie par #[derive(Validate)] sur le DTO.
value.validate().map_err(|errors| {
let mut error_map = HashMap::new();
for (field, field_errors) in errors.field_errors() {
let messages: Vec<String> = field_errors
.iter()
.map(|e| {
e.message
.as_ref()
.map(|m| m.to_string())
.unwrap_or_else(|| e.code.to_string())
})
.collect();
error_map.insert(field.to_string(), messages);
}
// Retourne un format JSON structuré : {"errors": {"field": ["msg"]}}
(
StatusCode::BAD_REQUEST,
Json(ValidationErrorResponse { errors: error_map }),
)
.into_response()
})?;
Ok(ValidatedJson(value))
}
}