From 6bcfac479f01e3c1d776623e30c97db6878427ae Mon Sep 17 00:00:00 2001 From: Ryan Lahfa Date: Wed, 10 Jun 2026 13:42:04 +0200 Subject: [PATCH 01/13] tracing: add json support Signed-off-by: Ryan Lahfa --- Cargo.lock | 23 ++++++++++++++++++----- Cargo.toml | 2 +- 2 files changed, 19 insertions(+), 6 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index f5afb8b..9f156ea 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -539,7 +539,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -1391,7 +1391,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -1664,10 +1664,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd" dependencies = [ "fastrand", - "getrandom 0.3.4", + "getrandom 0.4.2", "once_cell", "rustix", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -1921,6 +1921,16 @@ dependencies = [ "tracing-core", ] +[[package]] +name = "tracing-serde" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "704b1aeb7be0d0a84fc9828cae51dab5970fee5088f83d1dd7ee6f6246fc6ff1" +dependencies = [ + "serde", + "tracing-core", +] + [[package]] name = "tracing-subscriber" version = "0.3.23" @@ -1931,12 +1941,15 @@ dependencies = [ "nu-ansi-term", "once_cell", "regex-automata", + "serde", + "serde_json", "sharded-slab", "smallvec", "thread_local", "tracing", "tracing-core", "tracing-log", + "tracing-serde", ] [[package]] @@ -2135,7 +2148,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 771b09e..8780243 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -26,7 +26,7 @@ tokio = { version = "^1", features = ["full"] } tokio-rustls = "^0.26" toml = "^1" tracing = "^0.1" -tracing-subscriber = { version = "^0.3", features = ["env-filter"] } +tracing-subscriber = { version = "^0.3", features = ["env-filter", "json"] } winnow = "^1" zlink = "^0.5" rustls-pki-types = "^1" From 5722dc35e36e3261224c06bd5600ad1482681690 Mon Sep 17 00:00:00 2001 From: Ryan Lahfa Date: Wed, 10 Jun 2026 10:27:05 +0200 Subject: [PATCH 02/13] proxy/context: introduce trace IDs in contexts This way, it is possible to print it to users and be able to reconcile a very specific request. Signed-off-by: Ryan Lahfa --- Cargo.lock | 12 ++++++++++++ Cargo.toml | 1 + src/proxy/context.rs | 5 +++++ 3 files changed, 18 insertions(+) diff --git a/Cargo.lock b/Cargo.lock index 9f156ea..04045b1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1222,6 +1222,7 @@ dependencies = [ "toml", "tracing", "tracing-subscriber", + "uuid", "winnow", "zlink", ] @@ -1988,6 +1989,17 @@ version = "0.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821" +[[package]] +name = "uuid" +version = "1.23.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "144d6b123cef80b301b8f72a9e2ca4370ddec21950d0a103dd22c437006d2db7" +dependencies = [ + "getrandom 0.4.2", + "js-sys", + "wasm-bindgen", +] + [[package]] name = "valuable" version = "0.1.1" diff --git a/Cargo.toml b/Cargo.toml index 8780243..79f6376 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -30,6 +30,7 @@ tracing-subscriber = { version = "^0.3", features = ["env-filter", "json"] } winnow = "^1" zlink = "^0.5" rustls-pki-types = "^1" +uuid = { version = "1.23.3", features = ["v4"] } [dev-dependencies] criterion = { version = "0.8.2", features = ["html_reports"] } diff --git a/src/proxy/context.rs b/src/proxy/context.rs index cc628ce..4a1ee8a 100644 --- a/src/proxy/context.rs +++ b/src/proxy/context.rs @@ -56,6 +56,7 @@ pub struct TargetContext { #[derive(Debug, Clone)] pub struct InitialRequestContext { + pub trace_id: uuid::Uuid, pub client_address: SocketAddr, pub acl_ctx: crate::acl::OwnedEvaluationContext, } @@ -64,6 +65,8 @@ pub struct InitialRequestContext { pub struct LocalRequestContext<'s> { #[allow(dead_code)] pub client_address: &'s SocketAddr, + #[allow(dead_code)] + pub trace_id: uuid::Uuid, pub acl_ctx: crate::acl::EvaluationContext<'s>, } @@ -71,6 +74,7 @@ impl InitialRequestContext { pub fn new(client_address: SocketAddr) -> Self { Self { client_address, + trace_id: uuid::Uuid::new_v4(), acl_ctx: crate::acl::OwnedEvaluationContext::empty(), } } @@ -79,6 +83,7 @@ impl InitialRequestContext { LocalRequestContext { client_address: &self.client_address, acl_ctx: self.acl_ctx.fork(), + trace_id: self.trace_id, } } } From 67e69dfd85d616255508f3e36a4a4503ee34c550 Mon Sep 17 00:00:00 2001 From: Ryan Lahfa Date: Wed, 10 Jun 2026 10:27:55 +0200 Subject: [PATCH 03/13] proxy/context: rename Initial into Owned As it can be used in final contexts as well after an async move, making it non-initial. Signed-off-by: Ryan Lahfa --- src/proxy/context.rs | 4 ++-- src/proxy/http_connect.rs | 10 +++++----- src/proxy/mod.rs | 8 ++++---- src/proxy/socks5.rs | 4 ++-- 4 files changed, 13 insertions(+), 13 deletions(-) diff --git a/src/proxy/context.rs b/src/proxy/context.rs index 4a1ee8a..2139806 100644 --- a/src/proxy/context.rs +++ b/src/proxy/context.rs @@ -55,7 +55,7 @@ pub struct TargetContext { } #[derive(Debug, Clone)] -pub struct InitialRequestContext { +pub struct OwnedRequestContext { pub trace_id: uuid::Uuid, pub client_address: SocketAddr, pub acl_ctx: crate::acl::OwnedEvaluationContext, @@ -70,7 +70,7 @@ pub struct LocalRequestContext<'s> { pub acl_ctx: crate::acl::EvaluationContext<'s>, } -impl InitialRequestContext { +impl OwnedRequestContext { pub fn new(client_address: SocketAddr) -> Self { Self { client_address, diff --git a/src/proxy/http_connect.rs b/src/proxy/http_connect.rs index 9db70fc..423f93d 100644 --- a/src/proxy/http_connect.rs +++ b/src/proxy/http_connect.rs @@ -1,5 +1,5 @@ use crate::config::KnownBackend; -use crate::proxy::context::InitialRequestContext; +use crate::proxy::context::OwnedRequestContext; use crate::proxy::protocol_detect::{ALPN_H2, ALPN_HTTP1_1}; use crate::{ config::{BackendSettings, Settings}, @@ -42,7 +42,7 @@ enum InboundHttpProtocol { pub async fn serve_http1_connect( settings: Arc, state: Arc>, - ctx: InitialRequestContext, + ctx: OwnedRequestContext, stream: S, ) -> Result<(), ProxyError> { let io = TokioIo::new(stream); @@ -72,7 +72,7 @@ pub async fn serve_http1_connect( settings: Arc, state: Arc>, - ctx: InitialRequestContext, + ctx: OwnedRequestContext, stream: S, ) -> Result<(), ProxyError> { let io = TokioIo::new(stream); @@ -104,12 +104,12 @@ pub async fn serve_http2_connect, - ctx: InitialRequestContext, + initial_ctx: OwnedRequestContext, settings: Arc, state: Arc>, inbound_protocol: InboundHttpProtocol, ) -> Result>, hyper::Error> { - let mut ctx = ctx.as_local(); + let mut ctx = initial_ctx.as_local(); if req.method() != Method::CONNECT { debug!( diff --git a/src/proxy/mod.rs b/src/proxy/mod.rs index f024d7b..5d9ff7d 100644 --- a/src/proxy/mod.rs +++ b/src/proxy/mod.rs @@ -12,7 +12,7 @@ use tokio_rustls::{ }; use tracing::error; -use crate::{config::Settings, proxy::context::InitialRequestContext, state::State}; +use crate::{config::Settings, proxy::context::OwnedRequestContext, state::State}; mod client_tls; mod context; @@ -36,7 +36,7 @@ enum ProxyError { async fn serve_authenticated_proxy( settings: Arc, state: Arc>, - ctx: InitialRequestContext, + ctx: OwnedRequestContext, stream: TlsStream, ) -> anyhow::Result<()> { // TODO: extract context @@ -58,7 +58,7 @@ async fn serve_authenticated_proxy( async fn serve_unauthenticated_proxy( settings: Arc, state: Arc>, - ctx: InitialRequestContext, + ctx: OwnedRequestContext, stream: tokio::net::TcpStream, ) -> anyhow::Result<()> { let (proto, stream) = detect_protocol(InboundStream::TcpStream(stream)).await?; @@ -127,7 +127,7 @@ pub async fn start( let acceptor = tls_acceptor.clone(); let settings = settings.clone(); let state = state.clone(); - let ctx = InitialRequestContext::new(addr); + let ctx = OwnedRequestContext::new(addr); tokio::spawn(async move { match detect_tls(&socket).await { diff --git a/src/proxy/socks5.rs b/src/proxy/socks5.rs index 2af95cb..613ff4c 100644 --- a/src/proxy/socks5.rs +++ b/src/proxy/socks5.rs @@ -18,7 +18,7 @@ use tracing::{debug, info, warn}; use crate::{ config::{BackendSettings, KnownBackend, Settings}, - proxy::context::{InitialRequestContext, TargetContext}, + proxy::context::{OwnedRequestContext, TargetContext}, state::State, }; @@ -80,7 +80,7 @@ pub async fn route_to_backend( pub async fn serve_socks5( opts: Arc, state: Arc>, - ctx: InitialRequestContext, + ctx: OwnedRequestContext, socket: S, ) -> Result<(), SocksError> { let mut ctx = ctx.as_local(); From 3aa1b5800c40253711cc2f0185eb59a666e1ee00 Mon Sep 17 00:00:00 2001 From: Ryan Lahfa Date: Tue, 9 Jun 2026 23:47:21 +0200 Subject: [PATCH 04/13] rpc/cli: separate logging Signed-off-by: Ryan Lahfa --- src/logging.rs | 12 ++++++++++++ src/main.rs | 13 +++---------- 2 files changed, 15 insertions(+), 10 deletions(-) create mode 100644 src/logging.rs diff --git a/src/logging.rs b/src/logging.rs new file mode 100644 index 0000000..ca42adf --- /dev/null +++ b/src/logging.rs @@ -0,0 +1,12 @@ +use tracing::level_filters::LevelFilter; +use tracing_subscriber::EnvFilter; + +pub fn init() { + tracing_subscriber::fmt::fmt() + .with_env_filter( + EnvFilter::builder() + .with_default_directive(LevelFilter::INFO.into()) + .from_env_lossy(), + ) + .init(); +} diff --git a/src/main.rs b/src/main.rs index e59933e..f85922d 100644 --- a/src/main.rs +++ b/src/main.rs @@ -4,15 +4,14 @@ use std::net::SocketAddr; use std::os::fd::FromRawFd; use std::{path::PathBuf, sync::Arc}; use tokio::sync::RwLock; -use tracing::{error, level_filters::LevelFilter}; -use tracing::{info, warn}; -use tracing_subscriber::EnvFilter; +use tracing::{error, info, warn}; use crate::rpc::fr_gouv_portail_control::{DynamicBackendSpec, GetCurrentBackendOutput}; use crate::systemd::sd_notify_ready; mod acl; mod config; +mod logging; mod proxy; mod rpc; mod state; @@ -108,13 +107,7 @@ enum Commands { async fn main() -> Result<()> { let cli = Cli::parse(); - tracing_subscriber::fmt::fmt() - .with_env_filter( - EnvFilter::builder() - .with_default_directive(LevelFilter::INFO.into()) - .from_env_lossy(), - ) - .init(); + logging::init(); match cli.command { Commands::Rpc { From c365f0a82c4cca5dca505d866a6f1f7ee95cf41d Mon Sep 17 00:00:00 2001 From: Ryan Lahfa Date: Mon, 15 Jun 2026 17:23:17 +0200 Subject: [PATCH 05/13] rpc/main: move the connection at the start of the match block Signed-off-by: Ryan Lahfa --- src/main.rs | 20 +++++--------------- 1 file changed, 5 insertions(+), 15 deletions(-) diff --git a/src/main.rs b/src/main.rs index f85922d..dbcbc49 100644 --- a/src/main.rs +++ b/src/main.rs @@ -115,12 +115,13 @@ async fn main() -> Result<()> { json, command, } => { + let mut connection = zlink::unix::connect(&rpc_socket).await.context(format!( + "Opening the RPC socket at path '{}'", + rpc_socket.display() + ))?; + match command { RpcCommands::PrintCurrentBackend => { - let mut connection = zlink::unix::connect(&rpc_socket).await.context( - format!("Opening the RPC socket at path '{}'", rpc_socket.display()), - )?; - let cur_backend = connection .get_current_backend() .await @@ -142,10 +143,6 @@ async fn main() -> Result<()> { } RpcCommands::SetDefaultBackend { backend_id } => { - let mut connection = zlink::unix::connect(&rpc_socket).await.context( - format!("Opening the RPC socket at path '{}'", rpc_socket.display()), - )?; - connection .set_default_backend(Some(&backend_id)) .await @@ -166,10 +163,6 @@ async fn main() -> Result<()> { } RpcCommands::UnsetDefaultBackend => { - let mut connection = zlink::unix::connect(&rpc_socket).await.context( - format!("Opening the RPC socket at path '{}'", rpc_socket.display()), - )?; - connection .set_default_backend(None) .await @@ -190,9 +183,6 @@ async fn main() -> Result<()> { } RpcCommands::ListBackends => { - let mut connection = zlink::unix::connect(&rpc_socket).await.context( - format!("Opening the RPC socket at path '{}'", rpc_socket.display()), - )?; let backends = connection .list_backends() .await From 2f9f93166d77761b7ccccb1af811ce5a978f1b02 Mon Sep 17 00:00:00 2001 From: Ryan Lahfa Date: Wed, 10 Jun 2026 00:27:40 +0200 Subject: [PATCH 06/13] daemon: log more startup and sync readiness Signed-off-by: Ryan Lahfa --- src/main.rs | 25 ++++++++++++++++--------- 1 file changed, 16 insertions(+), 9 deletions(-) diff --git a/src/main.rs b/src/main.rs index dbcbc49..8fb8122 100644 --- a/src/main.rs +++ b/src/main.rs @@ -4,7 +4,7 @@ use std::net::SocketAddr; use std::os::fd::FromRawFd; use std::{path::PathBuf, sync::Arc}; use tokio::sync::RwLock; -use tracing::{error, info, warn}; +use tracing::{debug, error, info, warn}; use crate::rpc::fr_gouv_portail_control::{DynamicBackendSpec, GetCurrentBackendOutput}; use crate::systemd::sd_notify_ready; @@ -239,12 +239,15 @@ async fn main() -> Result<()> { bind_proxy_address, bind_rpc_socket, } => { + info!("Reading Portail settings from '{}'", config.display()); let settings: Arc = Arc::new(config::init(&config)); let state: Arc> = Arc::new(RwLock::new( state::init(&settings).context("While initializing application state")?, )); + debug!("Loaded Portail settings and state"); + let fds_named = systemd::listen_fds_named(); let tcp_listener = if let Some(proxy_address) = bind_proxy_address { @@ -271,20 +274,24 @@ async fn main() -> Result<()> { tokio::net::UnixListener::from_std(std)? }; - info!("starting services"); + info!("Starting services"); + + let (proxy_fut, rpc_fut) = ( + proxy::start(settings.clone(), state.clone(), tcp_listener), + rpc::start(settings.clone(), state.clone(), rpc_listener), + ); if let Err(e) = sd_notify_ready() { - warn!("failed to notify systemd about readiness: {e}"); + warn!("Failed to notify systemd about readiness: {e}"); } else { - info!("notified systemd about readiness"); + debug!("Notified systemd about readiness"); } - tokio::try_join!( - proxy::start(settings.clone(), state.clone(), tcp_listener), - rpc::start(settings.clone(), state.clone(), rpc_listener), - )?; + info!("Services are ready."); + + tokio::try_join!(proxy_fut, rpc_fut)?; - info!("exiting"); + info!("Exiting..."); } Commands::CheckACLSyntax { config, acl_file } => { let contents = String::from_utf8_lossy( From bfe7336f1b90742cfeb51fd8cdf2193668d29168 Mon Sep 17 00:00:00 2001 From: Ryan Lahfa Date: Mon, 15 Jun 2026 18:41:16 +0200 Subject: [PATCH 07/13] proxy/context: make TargetContext derive Debug Signed-off-by: Ryan Lahfa --- src/proxy/context.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/src/proxy/context.rs b/src/proxy/context.rs index 2139806..d1b51dd 100644 --- a/src/proxy/context.rs +++ b/src/proxy/context.rs @@ -49,6 +49,7 @@ impl From for TargetAddr { } } +#[derive(Debug)] pub struct TargetContext { pub initial_target: TargetAddr, pub resolved_target: Option, From b2d84cd0eec7cfacbabe405c3710cca590c5678f Mon Sep 17 00:00:00 2001 From: Ryan Lahfa Date: Mon, 15 Jun 2026 18:40:52 +0200 Subject: [PATCH 08/13] proxy/socks5: add structured logging With timings Signed-off-by: Ryan Lahfa --- src/proxy/socks5.rs | 116 ++++++++++++++++++++++++++++++++++++-------- 1 file changed, 95 insertions(+), 21 deletions(-) diff --git a/src/proxy/socks5.rs b/src/proxy/socks5.rs index 613ff4c..707b6da 100644 --- a/src/proxy/socks5.rs +++ b/src/proxy/socks5.rs @@ -11,7 +11,7 @@ use tokio::{ io::{AsyncRead, AsyncWrite}, net::TcpStream, sync::RwLock, - time::timeout, + time::{Instant, timeout}, }; use tokio_rustls::TlsStream; use tracing::{debug, info, warn}; @@ -53,6 +53,7 @@ pub async fn connect_to_backend( )) } else { debug!( + subsystem = "proxy_access", "Backend is not identity-aware, establishing a plain SOCKS5 connection to the backend" ); Ok(OutboundSock5Stream::Plain( @@ -77,37 +78,56 @@ pub async fn route_to_backend( Ok(()) } +#[tracing::instrument(skip_all, fields(trace_id = %ctx.trace_id, client_address = %ctx.client_address, subsystem = "proxy_access"))] pub async fn serve_socks5( opts: Arc, state: Arc>, ctx: OwnedRequestContext, socket: S, ) -> Result<(), SocksError> { - let mut ctx = ctx.as_local(); let should_resolve_dns: bool = state.read().await.default_backend.is_none(); + let start = Instant::now(); let (proto, cmd, target_addr) = Socks5ServerProtocol::accept_no_auth(socket) .await? .read_command() .await?; - ctx.acl_ctx.insert( - "proxy.protocol", - crate::acl::ast::ConcreteOperand::String("socks5"), + debug!( + subsystem = "proxy_access", + target_addr = %target_addr, + duration_ms = %start.elapsed().as_millis(), + "SOCKS5 target address obtained" ); + let start = Instant::now(); let mut target_context = TargetContext { initial_target: target_addr.clone().into(), resolved_target: None, }; + let target_addr = if should_resolve_dns { target_addr.resolve_dns().await? } else { target_addr }; + debug!( + subsystem = "proxy_access", + target_addr = %target_addr, + duration_ms = %start.elapsed().as_millis(), + "SOCKS5 resolved target address obtained" + ); + + let start = Instant::now(); + let (host, port) = target_context.initial_target.clone().into_string_and_port(); + let mut ctx = ctx.as_local(); + ctx.acl_ctx.insert( + "proxy.protocol", + crate::acl::ast::ConcreteOperand::String("socks5"), + ); ctx.acl_ctx .insert("host", crate::acl::ast::ConcreteOperand::String(&host)); ctx.acl_ctx.insert( @@ -124,7 +144,11 @@ pub async fn serve_socks5( }; if cmd != Socks5Command::TCPConnect && cmd != Socks5Command::UDPAssociate { - debug!("Unsupported SOCKS5 command received, terminating connection"); + info!( + subsystem = "proxy_errors", + command = ?cmd, + "Unsupported SOCKS5 command received, terminating connection" + ); proto.reply_error(&ReplyError::CommandNotSupported).await?; return Err(ReplyError::CommandNotSupported.into()); } @@ -143,6 +167,12 @@ pub async fn serve_socks5( ); } + debug!( + subsystem = "proxy_access", + command = ?cmd, + "SOCKS5 command allowed" + ); + let mut backends: Vec = Vec::with_capacity(1); let acl = &state.read().await.acl_rules; let backend_specs = &state.read().await.backends; @@ -153,8 +183,8 @@ pub async fn serve_socks5( Err(failure) => { proto.reply_error(&ReplyError::GeneralFailure).await?; warn!( - "Failed to evaluate routes for a request: {} (Context: {:#?})", - failure, ctx + subsystem = "proxy_errors", + "Failed to evaluate routes for a request: {} (Context: {:#?})", failure, ctx ); return Ok(()); } @@ -197,8 +227,8 @@ pub async fn serve_socks5( Err(failure) => { proto.reply_error(&ReplyError::GeneralFailure).await?; warn!( - "Failed to evaluate a request: {} (Context: {:#?})", - failure, ctx + subsystem = "proxy_errors", + "Failed to evaluate a request: {} (Context: {:#?})", failure, ctx ); return Ok(()); } @@ -207,14 +237,22 @@ pub async fn serve_socks5( match assessment.action { // FIXME: render the deny template if there's one. crate::acl::Action::Deny(_explain_template) => { - info!("Request to {0} is blocked", &target_context.initial_target); + info!( + subsystem = "proxy_access", + target_context = ?target_context, + duration_us = start.elapsed().as_micros(), + "SOCKS5 request blocked due to ACL" + ); proto.reply_error(&ReplyError::ConnectionNotAllowed).await?; return Ok(()); } crate::acl::Action::Redirect(target) => { info!( - "Request to {} redirected to {}", - &target_context.initial_target, target + subsystem = "proxy_access", + target_context = ?target_context, + redirected_to = %target, + duration_us = start.elapsed().as_micros(), + "SOCKS5 request redirected due to ACL" ); final_addr = TargetAddr::Domain( target @@ -230,8 +268,21 @@ pub async fn serve_socks5( _ => {} } + info!( + subsystem = "proxy_access", + duration_us = start.elapsed().as_micros(), + "SOCKS5 request allowed due to ACL" + ); + + let start = Instant::now(); // Either, we route to another backend or we do the SOCKS5 proxying ourselves. while let Some(backend) = backends.pop() { + debug!( + subsystem = "proxy_access", + backend = ?backend, + duration_ms = start.elapsed().as_millis(), + "Backend selected for connection routing" + ); match backend { BackendSettings::UnresolvedBackend => { // TODO: keep the IDs to print them here. @@ -242,11 +293,6 @@ pub async fn serve_socks5( return Ok(()); } BackendSettings::KnownBackend(backend) => { - debug!( - "Backend {} selected for routing the connection", - &backend.target_address - ); - match timeout( opts.request_timeout, connect_to_backend(&backend, &final_addr, state.clone()), @@ -254,15 +300,33 @@ pub async fn serve_socks5( .await { Ok(Ok(stream)) => { + debug!( + subsystem = "proxy_access", + backend = ?backend, + duration_ms = start.elapsed().as_millis(), + "Connection to upstream backend successful" + ); + + let start = Instant::now(); + route_to_backend(stream, proto).await?; + + debug!( + subsystem = "proxy_access", + duration_ms = start.elapsed().as_millis(), + "SOCKS5 request finished" + ); + return Ok(()); } Ok(Err(err)) => { debug!( - "Backend {} failed to route the request: {err}, trying the next one", - &backend.target_address + subsystem = "proxy_errors", + backend = ?backend, + "Backend failed to route the request: {err}, trying the next one", ); + continue; } @@ -280,7 +344,12 @@ pub async fn serve_socks5( // If we get there, this means that we did not have any backend at all. { - debug!("No backend, terminating the connection ourself"); + debug!( + subsystem = "proxy_access", + duration_ms = start.elapsed().as_millis(), + "No backend, terminating the connection ourself" + ); + let start = Instant::now(); match (cmd, opts.public_address) { (Socks5Command::TCPConnect, _) => { fast_socks5::server::run_tcp_proxy( @@ -302,6 +371,11 @@ pub async fn serve_socks5( return Err(ReplyError::CommandNotSupported.into()); } } + debug!( + subsystem = "proxy_access", + duration_ms = start.elapsed().as_millis(), + "SOCKS5 request finished" + ); } Ok(()) From a01c6021996da78081ae284ca3817ab04f51b3dc Mon Sep 17 00:00:00 2001 From: Ryan Lahfa Date: Thu, 18 Jun 2026 14:59:34 +0200 Subject: [PATCH 09/13] proxy/http: add structured logging With timings --- src/proxy/http_connect.rs | 166 +++++++++++++++++++++++++++++--------- 1 file changed, 129 insertions(+), 37 deletions(-) diff --git a/src/proxy/http_connect.rs b/src/proxy/http_connect.rs index 423f93d..9ca22f8 100644 --- a/src/proxy/http_connect.rs +++ b/src/proxy/http_connect.rs @@ -21,8 +21,8 @@ use std::sync::Arc; use tokio::io::{AsyncRead, AsyncWrite}; use tokio::net::TcpStream; use tokio::sync::RwLock; -use tokio::time::timeout; -use tracing::{debug, error, info, warn}; +use tokio::time::{Instant, timeout}; +use tracing::{Instrument, debug, error, info, warn}; /// This is a workaround for the restriction `only auto traits can be used as additional traits in a trait object` trait OutboundStreamIo: AsyncRead + AsyncWrite {} @@ -30,6 +30,7 @@ impl OutboundStreamIo for T {} type OutboundStream = Box; +#[derive(Debug)] enum InboundHttpProtocol { Http1, Http2, @@ -39,12 +40,14 @@ enum InboundHttpProtocol { /// - https://docs.rs/hyper/latest/hyper/upgrade/index.html /// - https://github.com/hyperium/hyper/blob/master/examples/http_proxy.rs /// - https://github.com/hyperium/hyper/blob/master/examples/upgrades.rs +#[tracing::instrument(skip_all, fields(trace_id = %ctx.trace_id, client_address = %ctx.client_address, subsystem = "proxy_access"))] pub async fn serve_http1_connect( settings: Arc, state: Arc>, ctx: OwnedRequestContext, stream: S, ) -> Result<(), ProxyError> { + debug!(subsystem = "proxy_access", "HTTP/1.1 CONNECT request"); let io = TokioIo::new(stream); // TODO: update ctx @@ -69,12 +72,14 @@ pub async fn serve_http1_connect( settings: Arc, state: Arc>, ctx: OwnedRequestContext, stream: S, ) -> Result<(), ProxyError> { + debug!(subsystem = "proxy_access", "HTTP/2 CONNECT request"); let io = TokioIo::new(stream); // TODO: update ctx @@ -102,6 +107,7 @@ pub async fn serve_http2_connect, initial_ctx: OwnedRequestContext, @@ -110,11 +116,12 @@ async fn handle_http_request( inbound_protocol: InboundHttpProtocol, ) -> Result>, hyper::Error> { let mut ctx = initial_ctx.as_local(); + let start = Instant::now(); if req.method() != Method::CONNECT { debug!( - "Unsupported HTTP method `{}` received, terminating connection", - req.method() + subsystem = "proxy_errors", + "Unsupported HTTP method received, terminating connection", ); // TODO: we might want to handle this case in the future let mut resp = Response::new(empty_body()); @@ -128,7 +135,10 @@ async fn handle_http_request( ); let Some(target_authority) = req.uri().authority() else { - debug!("Invalid authority in CONNECT URI, terminating connection"); + debug!( + subsystem = "proxy_errors", + "Invalid authority in CONNECT URI, terminating connection" + ); let mut resp = Response::new(empty_body()); *resp.status_mut() = StatusCode::BAD_REQUEST; return Ok(resp); @@ -138,6 +148,8 @@ async fn handle_http_request( let mut final_address = target_address.clone(); let mut backends: Vec = Vec::with_capacity(1); let backend_specs = &state.read().await.backends; + debug!(subsystem = "proxy_access", final_address = %final_address, + "HTTP CONNECT request"); // We evaluate first whether we are allowed then we evaluate routes. let acl = &state.read().await.acl_rules; @@ -181,8 +193,8 @@ async fn handle_http_request( Err(failure) => { let mut resp = Response::new(empty_body()); warn!( - "Failed to evaluate a request: {} (Context: {:#?})", - failure, ctx + subsystem = "proxy_errors", + "Failed to evaluate a request: {} (Context: {:#?})", failure, ctx ); *resp.status_mut() = StatusCode::INTERNAL_SERVER_ERROR; return Ok(resp); @@ -226,8 +238,8 @@ async fn handle_http_request( Err(failure) => { let mut resp = Response::new(empty_body()); warn!( - "Failed to evaluate a request: {} (Context: {:#?})", - failure, ctx + subsystem = "proxy_errors", + "Failed to evaluate a request: {} (Context: {:#?})", failure, ctx ); *resp.status_mut() = StatusCode::INTERNAL_SERVER_ERROR; return Ok(resp); @@ -237,16 +249,22 @@ async fn handle_http_request( match assessment.action { // FIXME: render the deny template if there's one. crate::acl::Action::Deny(_explain_template) => { - info!("Request to {} is blocked", target_authority.host()); + info!( + subsystem = "proxy_access", + duration_us = start.elapsed().as_micros(), + "Request denied by ACL", + ); let mut resp = Response::new(empty_body()); *resp.status_mut() = StatusCode::FORBIDDEN; return Ok(resp); } crate::acl::Action::Redirect(target) => { info!( - "Request to {} redirected to {}", - target_authority.host(), - target + subsystem = "proxy_access", + original_host = %target_authority.host(), + redirected_to = %target, + duration_us = start.elapsed().as_micros(), + "Request redirected by ACL", ); final_address = target.to_string(); } @@ -254,8 +272,21 @@ async fn handle_http_request( _ => {} } + info!( + subsystem = "proxy_access", + duration_us = start.elapsed().as_micros(), + "Request allowed by ACL", + ); + let mut stream: Option = None; + let start = Instant::now(); for backend in backends { + debug!( + subsystem = "proxy_access", + address = %final_address, + backend = ?backend, + "Backend selected for HTTP CONNECT" + ); match backend { BackendSettings::UnresolvedBackend => { // TODO: keep the IDs to print them here. @@ -267,10 +298,6 @@ async fn handle_http_request( return Ok(resp); } BackendSettings::KnownBackend(backend) => { - debug!( - "Backend {} selected for HTTP CONNECT to {}", - backend.target_address, final_address - ); match timeout( settings.request_timeout, connect_to_http_proxy_backend( @@ -283,19 +310,33 @@ async fn handle_http_request( .await { Ok(Ok(upstream)) => { + debug!( + subsystem = "proxy_access", + address = %final_address, + backend = ?backend, + duration_ms = start.elapsed().as_millis(), + "Stream established to upstream backend" + ); + stream = Some(upstream); break; } Ok(Err(err)) => { debug!( - "Backend {} failed for HTTP CONNECT: {}, trying next", - backend.target_address, err + subsystem = "proxy_access", + backend = ?backend, + duration_ms = start.elapsed().as_millis(), + "Backend failed for HTTP CONNECT: {}, trying next", + err ); } Err(_) => { debug!( - "Backend {} timed out for HTTP CONNECT after {:?}, trying next", - backend.target_address, settings.request_timeout + subsystem = "proxy_access", + backend = ?backend, + duration_ms = start.elapsed().as_millis(), + configured_timeout_ms = settings.request_timeout.as_millis(), + "Backend timed out for HTTP CONNECT, trying next", ); } } @@ -304,36 +345,80 @@ async fn handle_http_request( } if stream.is_none() { debug!( - "No backend, establishing a direct connection to `{}`", - final_address + subsystem = "proxy_access", + address = %final_address, + duration_ms = start.elapsed().as_millis(), + "No backend, establishing a direct connection to the target" ); + + let start = Instant::now(); match timeout(settings.request_timeout, TcpStream::connect(&final_address)).await { - Ok(Ok(socket)) => stream = Some(Box::new(socket)), - Ok(Err(e)) => warn!("Direct connection to `{}` failed: {}", final_address, e), + Ok(Ok(socket)) => { + debug!( + subsystem = "proxy_access", + address = %final_address, + duration_ms = start.elapsed().as_millis(), + "Stream directly established to final address (local exit)" + ); + stream = Some(Box::new(socket)) + } + Ok(Err(e)) => warn!( + subsystem = "proxy_errors", + address = %final_address, + duration_ms = start.elapsed().as_millis(), + "Direct connection failed: {}", + e), Err(_) => warn!( - "Direct connection to `{}` timed out after {:?}", - final_address, settings.request_timeout + subsystem = "proxy_errors", + address = %final_address, + duration_ms = start.elapsed().as_millis(), + configured_timeout_ms = settings.request_timeout.as_millis(), + "Direct connection failed due to timeout", ), } } let Some(mut stream) = stream else { + warn!( + subsystem = "proxy_errors", + address = %final_address, + "No outbound stream could be established" + ); + let mut resp = Response::new(empty_body()); *resp.status_mut() = StatusCode::BAD_GATEWAY; return Ok(resp); }; - tokio::task::spawn(async move { - match hyper::upgrade::on(req).await { - Ok(upgraded) => { - let mut client = TokioIo::new(upgraded); - if let Err(e) = tokio::io::copy_bidirectional(&mut client, &mut *stream).await { - error!("CONNECT tunnel error: {}", e); + tokio::task::spawn( + async move { + let start = Instant::now(); + match hyper::upgrade::on(req).await { + Ok(upgraded) => { + let mut client = TokioIo::new(upgraded); + match tokio::io::copy_bidirectional(&mut client, &mut *stream).await { + Err(e) => error!( + subsystem = "proxy_errors", + duration_ms = start.elapsed().as_millis(), + "CONNECT tunnel error: {}", + e + ), + Ok((n_bytes_sent, n_bytes_recv)) => { + info!( + subsystem = "proxy_access", + n_bytes_sent = %n_bytes_sent, + n_bytes_recv = %n_bytes_recv, + duration_ms = start.elapsed().as_millis(), + "CONNECT tunnel finished successfully" + ); + } + } } + Err(e) => error!(subsystem = "proxy_errors", "CONNECT upgrade error: {}", e), } - Err(e) => debug!("CONNECT upgrade error: {}", e), } - }); + .in_current_span(), + ); // Connection: keep-alive is // - the default for HTTP/1.1 @@ -368,6 +453,7 @@ async fn connect_to_http_proxy_backend( let (stream, use_http2): (OutboundStream, bool) = if backend.identity_aware { debug!( + subsystem = "proxy_access", "Backend is identity-aware, establishing a TLS connection to {}", backend.target_address ); @@ -415,7 +501,10 @@ async fn connect_to_http_proxy_backend( .map_err(io::Error::other)?; tokio::spawn(async move { if let Err(e) = conn.await { - warn!("Cannot HTTP CONNECT to upstream: {}", e); + warn!( + subsystem = "proxy_errors", + "Cannot HTTP CONNECT to upstream: {}", e + ); } }); @@ -448,7 +537,10 @@ async fn connect_to_http_proxy_backend( // https://docs.rs/hyper/latest/hyper/client/conn/http1/struct.Builder.html#method.handshake tokio::spawn(async move { if let Err(e) = conn.with_upgrades().await { - warn!("Cannot HTTP CONNECT to upstream: {}", e); + warn!( + subsystem = "proxy_errors", + "Cannot HTTP CONNECT to upstream: {}", e + ); } }); From 42e50203133dec15b20d1f87644f264adddf1ade Mon Sep 17 00:00:00 2001 From: Ryan Lahfa Date: Mon, 15 Jun 2026 18:41:42 +0200 Subject: [PATCH 10/13] proxy/accept: add structured logging Signed-off-by: Ryan Lahfa --- src/proxy/mod.rs | 73 +++++++++++++++++++++++++++++++++++------------- 1 file changed, 53 insertions(+), 20 deletions(-) diff --git a/src/proxy/mod.rs b/src/proxy/mod.rs index 5d9ff7d..76316d4 100644 --- a/src/proxy/mod.rs +++ b/src/proxy/mod.rs @@ -2,7 +2,7 @@ use anyhow::bail; use fast_socks5::SocksError; use std::sync::Arc; use thiserror::Error; -use tokio::sync::RwLock; +use tokio::{net::TcpStream, sync::RwLock}; use tokio_rustls::{ TlsAcceptor, TlsStream, rustls::{ @@ -10,7 +10,7 @@ use tokio_rustls::{ server::{VerifierBuilderError, WebPkiClientVerifier}, }, }; -use tracing::error; +use tracing::{Instrument, debug, error}; use crate::{config::Settings, proxy::context::OwnedRequestContext, state::State}; @@ -114,27 +114,32 @@ async fn build_tls_acceptor( } } -pub async fn start( +#[tracing::instrument(skip_all, fields(trace_id = %ctx.trace_id, client_address = %ctx.client_address, subsystem = "proxy_access"))] +pub async fn accept_client( settings: Arc, state: Arc>, - listener: tokio::net::TcpListener, -) -> anyhow::Result<()> { - let tls_acceptor: Option = build_tls_acceptor(&settings, state.clone()).await?; - - loop { - let (socket, addr) = listener.accept().await?; + socket: TcpStream, + tls_acceptor: Option, + ctx: OwnedRequestContext, +) { + debug!(subsystem = "proxy_access", "Accepting a proxy connection"); - let acceptor = tls_acceptor.clone(); - let settings = settings.clone(); - let state = state.clone(); - let ctx = OwnedRequestContext::new(addr); + let acceptor = tls_acceptor.clone(); + let settings = settings.clone(); + let state = state.clone(); - tokio::spawn(async move { + tokio::spawn( + async move { match detect_tls(&socket).await { Ok(true) => { + debug!(subsystem = "proxy_access", "TLS detected"); if let Some(acceptor) = acceptor { match acceptor.accept(socket).await { Ok(tls_stream) => { + debug!( + subsystem = "proxy_access", + "Authenticated TLS stream (client certificates)" + ); if let Err(e) = serve_authenticated_proxy( settings, state, @@ -143,33 +148,61 @@ pub async fn start( ) .await { - error!("TLS proxy error from {addr}: {e:?}"); + error!(subsystem = "proxy_errors", "TLS proxy error: {e:?}"); } } Err(e) => { - error!("TLS handshake failed from {addr}: {e:?}"); + error!(subsystem = "proxy_errors", "TLS handshake failed: {e:?}"); } } } else { error!( - "TLS received from {addr}: but no TLS configuration set in the proxy" + subsystem = "proxy_errors", + "TLS received but no TLS configuration set in the proxy" ); } } Ok(false) => { + debug!( + subsystem = "proxy_access", + "No TLS detected, serving unauthenticated requests", + ); if let Err(e) = serve_unauthenticated_proxy(settings, state, ctx, socket).await { - error!("Proxy error from {addr}: {e:?}"); + error!(subsystem = "proxy_errors", "Proxy error: {e:?}"); } } Err(err) => { error!( - "While detecting the header for TLS from {addr}, error occurred: {err:?}" + subsystem = "proxy_errors", + "While detecting the header for TLS, error occurred: {err:?}" ); } } - }); + } + .in_current_span(), + ); +} + +pub async fn start( + settings: Arc, + state: Arc>, + listener: tokio::net::TcpListener, +) -> anyhow::Result<()> { + let tls_acceptor: Option = build_tls_acceptor(&settings, state.clone()).await?; + + loop { + let (socket, addr) = listener.accept().await?; + let ctx = OwnedRequestContext::new(addr); + accept_client( + settings.clone(), + state.clone(), + socket, + tls_acceptor.clone(), + ctx, + ) + .await; } } From d10955f7d773cfc6649bba0c49dcce791f49d338 Mon Sep 17 00:00:00 2001 From: Ryan Lahfa Date: Mon, 15 Jun 2026 18:13:58 +0200 Subject: [PATCH 11/13] rpc: add structured logging Known to be broken: spans for connection-passing function are dropped. I assume this is a macro problem with zlink. Signed-off-by: Ryan Lahfa --- src/rpc/mod.rs | 60 ++++++++++++++++++++++++++++++++++++-------------- 1 file changed, 43 insertions(+), 17 deletions(-) diff --git a/src/rpc/mod.rs b/src/rpc/mod.rs index 50c421f..cb09958 100644 --- a/src/rpc/mod.rs +++ b/src/rpc/mod.rs @@ -8,7 +8,7 @@ use crate::{ }; use std::{collections::HashSet, ffi::CStr, sync::Arc}; use tokio::{net::UnixListener, sync::RwLock}; -use tracing::info; +use tracing::{debug, info, warn}; use zlink::{ Server, connection::{Gid, Socket, socket::FetchPeerCredentials}, @@ -102,22 +102,27 @@ impl Control { groups.push(creds.unix_primary_group_id()); let groups = resolve_numeric_groups_to_names(groups); + let authorized_groups = authorized_groups + .into_iter() + .map(|s| s.as_ref().to_owned()) + .collect::>(); + + // TODO: resolve numeric UIDs into usernames proper for better logs. if creds.unix_user_id().is_root() || authorized_groups - .into_iter() - .any(|trusted| groups.contains(trusted.as_ref())) + .iter() + .any(|trusted| groups.contains(trusted)) { - info!( - "Privileged RPC allowed from user '{}'", - creds.unix_user_id() - ); + info!(user = %creds.unix_user_id(), subsystem = "rpc", "Privileged RPC allowed"); Ok(()) } else { - tracing::warn!( - "Privileged RPC call attempt from user '{}' (groups: '{:?}') refused", - creds.unix_user_id(), - groups + warn!( + user = %creds.unix_user_id(), + groups = ?groups, + authorized_groups = ?authorized_groups, + subsystem = "rpc", + "Privileged RPC call attempt refused", ); Err(ControlError::PermissionDenied) @@ -136,6 +141,7 @@ impl Control where Sock::ReadHalf: FetchPeerCredentials, { + #[tracing::instrument(skip_all, fields(subsystem = "rpc", backend_id = %backend_id))] async fn set_default_backend( &mut self, backend_id: Option<&str>, @@ -148,20 +154,24 @@ where if let Some(backend_id) = backend_id { if !state.backends.contains_key(backend_id) { + info!(target_backend = %backend_id, "Backend not found in state"); return Err(ControlError::BackendNotFound { provided_backend: backend_id.to_string(), available_backends: state.backends.keys().cloned().collect(), }); } + info!(previous_backend = ?state.default_backend, new_backend = %backend_id, "Default backend changed"); state.default_backend = Some(backend_id.to_owned()); } else { + info!(previous_backend = ?state.default_backend, "Default backend unset"); state.default_backend = None; } Ok(()) } + #[tracing::instrument(skip_all, fields(subsystem = "rpc", backend_id = %backend_id, backend_spec = %backend_spec))] async fn update_dynamic_backend( &mut self, backend_id: &str, @@ -176,6 +186,10 @@ where self.settings.backends.get(backend_id), Some(BackendSettings::KnownBackend(_)) ) { + warn!( + backend_target = %backend_id, + "Attempt to change a non-dynamic backend, rejected" + ); return Err(ControlError::ImmutableBackend); } @@ -185,22 +199,32 @@ where match state.backends.get_mut(backend_id) { Some(backend) => { info!( - "Changed backend `{}` from `{:?}` to `{:?}`", - backend_id, *backend, new_backend + changed_backend = %backend_id, + old_spec = ?*backend, + new_spec = ?new_backend, + "Dynamic backend specification changed" ); *backend = BackendSettings::KnownBackend(new_backend); Ok(()) } - None => Err(ControlError::BackendNotFound { - provided_backend: backend_id.to_string(), - available_backends: state.backends.keys().cloned().collect(), - }), + None => { + info!( + target_backend = %backend_id, + "Backend not found" + ); + Err(ControlError::BackendNotFound { + provided_backend: backend_id.to_string(), + available_backends: state.backends.keys().cloned().collect(), + }) + } } } + #[tracing::instrument(skip_all, fields(subsystem = "rpc"))] async fn get_current_backend(&mut self) -> GetCurrentBackendOutput { + debug!("Current backend read"); GetCurrentBackendOutput { backend_id: self .state @@ -212,9 +236,11 @@ where } } + #[tracing::instrument(skip_all, fields(subsystem = "rpc"))] async fn list_backends(&mut self) -> ListBackendsOutput { let cur_backend = self.state.read().await.default_backend.clone(); let backends = &self.state.read().await.backends; + debug!("Backend list read"); ListBackendsOutput { backends: backends .iter() From c40d7262b2cb1d60d666a598bc63826ba4a702b5 Mon Sep 17 00:00:00 2001 From: Ryan Lahfa Date: Wed, 10 Jun 2026 00:27:49 +0200 Subject: [PATCH 12/13] logging: add structured logging configuration MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit This adds a bunch of structures to configure log subsystems according to targets and filter them adequately. The logging mechanism supports routing to multiple outputs and different formats such as JSON, compact, pretty or the normal one (called full). Instead of letting the user configure everything, we offer 3 presets: - development — everything on stdout & stderr in the pretty format - systemd — some files from LogsDirectory=portail and stderr for system & errors, ideally, via journald later on, with json in some places - container — stdout & stderr in the json format. Let's not offer any customization for now as this complicates (for no good reason) the user's job to configure the system. Log shippers loves structured JSON, this is our default format. Our traces are updated to be routed properly and enriched with more spans now. Logging has been made non-blocking using tracing_appender which uses a thread pool to dispatch logs to a queue and let the worker depop them, this should ensure that logging does not add any meaningful overhead during proxying. This needs to be benchmarked and measured wrt to latency (with runtime deactivation, compile-time deactivation, etc.). Signed-off-by: Ryan Lahfa --- Cargo.lock | 29 ++++ Cargo.toml | 3 +- src/logging.rs | 12 -- src/logging/mod.rs | 404 +++++++++++++++++++++++++++++++++++++++++++++ src/main.rs | 60 ++++++- 5 files changed, 488 insertions(+), 20 deletions(-) delete mode 100644 src/logging.rs create mode 100644 src/logging/mod.rs diff --git a/Cargo.lock b/Cargo.lock index 04045b1..3cb1797 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -427,6 +427,15 @@ dependencies = [ "itertools", ] +[[package]] +name = "crossbeam-channel" +version = "0.5.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "82b8f8f868b36967f9606790d1903570de9ceaf870a7bf9fbbd3016d636a2cb2" +dependencies = [ + "crossbeam-utils", +] + [[package]] name = "crossbeam-deque" version = "0.8.6" @@ -1221,6 +1230,7 @@ dependencies = [ "tokio-rustls", "toml", "tracing", + "tracing-appender", "tracing-subscriber", "uuid", "winnow", @@ -1647,6 +1657,12 @@ version = "2.6.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "13c2bddecc57b384dee18652358fb23172facb8a2c51ccc10d74c157bdea3292" +[[package]] +name = "symlink" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a7973cce6668464ea31f176d85b13c7ab3bba2cb3b77a2ed26abd7801688010a" + [[package]] name = "syn" version = "2.0.117" @@ -1890,6 +1906,19 @@ dependencies = [ "tracing-core", ] +[[package]] +name = "tracing-appender" +version = "0.2.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "050686193eb999b4bb3bc2acfa891a13da00f79734704c4b8b4ef1a10b368a3c" +dependencies = [ + "crossbeam-channel", + "symlink", + "thiserror 2.0.18", + "time", + "tracing-subscriber", +] + [[package]] name = "tracing-attributes" version = "0.1.31" diff --git a/Cargo.toml b/Cargo.toml index 79f6376..2b8ad81 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -26,11 +26,12 @@ tokio = { version = "^1", features = ["full"] } tokio-rustls = "^0.26" toml = "^1" tracing = "^0.1" -tracing-subscriber = { version = "^0.3", features = ["env-filter", "json"] } +tracing-subscriber = { version = "^0.3", features = ["env-filter", "json", "std"] } winnow = "^1" zlink = "^0.5" rustls-pki-types = "^1" uuid = { version = "1.23.3", features = ["v4"] } +tracing-appender = "0.2.5" [dev-dependencies] criterion = { version = "0.8.2", features = ["html_reports"] } diff --git a/src/logging.rs b/src/logging.rs deleted file mode 100644 index ca42adf..0000000 --- a/src/logging.rs +++ /dev/null @@ -1,12 +0,0 @@ -use tracing::level_filters::LevelFilter; -use tracing_subscriber::EnvFilter; - -pub fn init() { - tracing_subscriber::fmt::fmt() - .with_env_filter( - EnvFilter::builder() - .with_default_directive(LevelFilter::INFO.into()) - .from_env_lossy(), - ) - .init(); -} diff --git a/src/logging/mod.rs b/src/logging/mod.rs new file mode 100644 index 0000000..b2999c4 --- /dev/null +++ b/src/logging/mod.rs @@ -0,0 +1,404 @@ +use std::{collections::HashMap, fs::OpenOptions, path::PathBuf}; + +use anyhow::Context; +use clap::ValueEnum; +use tracing::level_filters::LevelFilter; +use tracing_appender::non_blocking::WorkerGuard; +use tracing_subscriber::{ + EnvFilter, Layer, Registry, + filter::FilterExt, + layer::{Filter, SubscriberExt}, + util::SubscriberInitExt, +}; + +#[derive(Debug, Clone)] +pub enum LogFormat { + #[allow(dead_code)] + Full, + Compact, + Pretty, + Json, +} + +#[derive(Debug, Clone)] +pub struct LogRoute { + format: LogFormat, + output: LogOutput, +} + +#[derive(Debug, Clone)] +pub struct LogConfig { + routes: Vec, +} + +#[derive(Debug, Clone, PartialEq, Hash)] +pub enum LogOutput { + Stdout, + Stderr, + File(PathBuf), + #[allow(dead_code)] + Journald, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub enum Subsystem { + /// All log entries related to the proxy (access logs) + ProxyAccess, + /// All log entries related to the proxy (error logs) + ProxyErrors, + /// All log entries related to the RPC + Rpc, + /// All log entries related to everything else + System, +} + +impl std::str::FromStr for Subsystem { + type Err = (); + + fn from_str(s: &str) -> Result { + match s.to_lowercase().as_str() { + "proxy_access" => Ok(Self::ProxyAccess), + "proxy_errors" => Ok(Self::ProxyErrors), + "rpc" => Ok(Self::Rpc), + _ => Ok(Self::System), + } + } +} + +#[derive(Clone)] +pub struct SubsystemFilter { + destination: Subsystem, +} + +impl SubsystemFilter { + pub fn new(destination: Subsystem) -> Self { + Self { destination } + } +} + +impl Filter for SubsystemFilter { + fn enabled( + &self, + _metadata: &tracing::Metadata<'_>, + _ctx: &tracing_subscriber::layer::Context<'_, S>, + ) -> bool { + true + } + + fn event_enabled( + &self, + event: &tracing::Event<'_>, + _ctx: &tracing_subscriber::layer::Context<'_, S>, + ) -> bool { + let mut visitor = SubsystemVisitor::default(); + event.record(&mut visitor); + + visitor.subsystem.unwrap_or(Subsystem::System) == self.destination + } +} + +#[derive(Default)] +struct SubsystemVisitor { + pub subsystem: Option, +} + +const SUBSYSTEM_ROUTING_FIELD: &str = "subsystem"; + +impl tracing::field::Visit for SubsystemVisitor { + fn record_str(&mut self, field: &tracing::field::Field, value: &str) { + if field.name() == SUBSYSTEM_ROUTING_FIELD { + self.subsystem = value.parse().ok(); + } + } + + fn record_debug(&mut self, _field: &tracing::field::Field, _value: &dyn std::fmt::Debug) {} +} + +#[derive(ValueEnum, Debug, Clone, Copy)] +pub enum LogPreset { + /// Logs are all compact, printed on stdout (traces) and stderr (errors) + Cli, + /// Logs are all compact, printed on stdout (traces) and stderr (errors) with JSON output + Scripting, + /// All logs will be printed nicely on stdout (traces) and stderr (errors) + Development, + /// Errors will be printed nicely on stderr but traces and errors will be routed into files as + /// well in /var/log/portail in JSON format + Systemd, + /// All traces (errors included) will be printed in JSON format on stdout and stderr + Container, +} + +fn preset_to_config(preset: LogPreset) -> HashMap { + match preset { + LogPreset::Cli => { + let stdout = LogRoute { + format: LogFormat::Compact, + output: LogOutput::Stdout, + }; + let stderr = LogRoute { + format: LogFormat::Compact, + output: LogOutput::Stderr, + }; + + [ + ( + Subsystem::ProxyAccess, + LogConfig { + routes: vec![stdout.clone()], + }, + ), + ( + Subsystem::ProxyErrors, + LogConfig { + routes: vec![stderr], + }, + ), + ( + Subsystem::Rpc, + LogConfig { + routes: vec![stdout.clone()], + }, + ), + ( + Subsystem::System, + LogConfig { + routes: vec![stdout], + }, + ), + ] + .into() + } + LogPreset::Scripting => { + let stdout = LogRoute { + format: LogFormat::Json, + output: LogOutput::Stdout, + }; + let stderr = LogRoute { + format: LogFormat::Json, + output: LogOutput::Stderr, + }; + + [ + ( + Subsystem::ProxyAccess, + LogConfig { + routes: vec![stdout.clone()], + }, + ), + ( + Subsystem::ProxyErrors, + LogConfig { + routes: vec![stderr], + }, + ), + ( + Subsystem::Rpc, + LogConfig { + routes: vec![stdout.clone()], + }, + ), + ( + Subsystem::System, + LogConfig { + routes: vec![stdout], + }, + ), + ] + .into() + } + LogPreset::Development => { + let stdout = LogRoute { + format: LogFormat::Pretty, + output: LogOutput::Stdout, + }; + let stderr = LogRoute { + format: LogFormat::Pretty, + output: LogOutput::Stderr, + }; + + [ + ( + Subsystem::ProxyAccess, + LogConfig { + routes: vec![stdout.clone()], + }, + ), + ( + Subsystem::ProxyErrors, + LogConfig { + routes: vec![stderr], + }, + ), + ( + Subsystem::Rpc, + LogConfig { + routes: vec![stdout.clone()], + }, + ), + ( + Subsystem::System, + LogConfig { + routes: vec![stdout], + }, + ), + ] + .into() + } + LogPreset::Systemd => { + let paccess_file = LogRoute { + format: LogFormat::Json, + output: LogOutput::File("/var/log/portail/access.log".into()), + }; + + let perror_file = LogRoute { + format: LogFormat::Json, + output: LogOutput::File("/var/log/portail/error.log".into()), + }; + + let rpc_file = LogRoute { + format: LogFormat::Json, + output: LogOutput::File("/var/log/portail/rpc.log".into()), + }; + + let stderr = LogRoute { + format: LogFormat::Pretty, + output: LogOutput::Stderr, + }; + + [ + ( + Subsystem::ProxyAccess, + LogConfig { + routes: vec![paccess_file], + }, + ), + ( + Subsystem::ProxyErrors, + LogConfig { + routes: vec![stderr.clone(), perror_file], + }, + ), + ( + Subsystem::Rpc, + LogConfig { + routes: vec![rpc_file], + }, + ), + ( + Subsystem::System, + LogConfig { + routes: vec![stderr], + }, + ), + ] + .into() + } + LogPreset::Container => { + let stderr = LogRoute { + format: LogFormat::Json, + output: LogOutput::Stderr, + }; + + let stdout = LogRoute { + format: LogFormat::Json, + output: LogOutput::Stdout, + }; + + [ + ( + Subsystem::ProxyAccess, + LogConfig { + routes: vec![stdout.clone()], + }, + ), + ( + Subsystem::ProxyErrors, + LogConfig { + routes: vec![stderr.clone()], + }, + ), + ( + Subsystem::Rpc, + LogConfig { + routes: vec![stdout], + }, + ), + ( + Subsystem::System, + LogConfig { + routes: vec![stderr], + }, + ), + ] + .into() + } + } +} + +fn create_layer_for_route( + subsystem: Subsystem, + route: &LogRoute, +) -> Result<(Box + Send + Sync>, WorkerGuard), std::io::Error> { + let (writer, guard) = match &route.output { + LogOutput::Stdout => tracing_appender::non_blocking(std::io::stdout()), + LogOutput::Stderr => tracing_appender::non_blocking(std::io::stderr()), + LogOutput::File(path) => { + let file = OpenOptions::new().append(true).create(true).open(path)?; + tracing_appender::non_blocking(file) + } + LogOutput::Journald => unimplemented!(), + }; + + let fmt_layer = tracing_subscriber::fmt::layer() + .with_target(true) + .with_thread_ids(true) + .with_file(true) + .with_line_number(true); + + let fmt_layer = match route.format { + LogFormat::Full => fmt_layer.with_writer(writer).boxed(), + LogFormat::Pretty => fmt_layer.with_writer(writer).pretty().boxed(), + LogFormat::Compact => fmt_layer.with_writer(writer).compact().boxed(), + LogFormat::Json => fmt_layer.with_writer(writer).json().boxed(), + }; + + Ok(( + fmt_layer + .with_filter( + SubsystemFilter::new(subsystem).and( + EnvFilter::builder() + .with_default_directive(LevelFilter::INFO.into()) + .from_env_lossy(), + ), + ) + .boxed(), + guard, + )) +} + +/// Guards for each worker that drains the logging queue. +/// When they are dropped, the corresponding log outputs are closed. +pub struct LogGuard { + guards: Vec, +} + +pub fn init(preset: LogPreset) -> anyhow::Result { + let config = preset_to_config(preset); + let mut layers = Vec::new(); + let mut guards = LogGuard { guards: Vec::new() }; + + // For each subsystem, create one layer per route tagged by the subsystem filter. + for (subsystem, log_config) in config { + for route in log_config.routes { + let (layer, guard) = create_layer_for_route(subsystem, &route) + .context("Creating a layer for a log route")?; + layers.push(layer); + guards.guards.push(guard); + } + } + + tracing_subscriber::registry().with(layers).init(); + + Ok(guards) +} diff --git a/src/main.rs b/src/main.rs index 8fb8122..60eda2d 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,11 +1,12 @@ use anyhow::{Context, Result}; -use clap::{Parser, Subcommand}; +use clap::{Parser, Subcommand, ValueEnum}; use std::net::SocketAddr; use std::os::fd::FromRawFd; use std::{path::PathBuf, sync::Arc}; use tokio::sync::RwLock; use tracing::{debug, error, info, warn}; +use crate::logging::LogPreset; use crate::rpc::fr_gouv_portail_control::{DynamicBackendSpec, GetCurrentBackendOutput}; use crate::systemd::sd_notify_ready; @@ -62,6 +63,27 @@ enum RpcCommands { }, } +#[derive(ValueEnum, Debug, Clone, Copy)] +pub enum UserLogPreset { + /// All logs will be printed nicely on stdout (traces) and stderr (errors) + Development, + /// Errors will be printed nicely on stderr but traces and errors will be routed into files as + /// well in /var/log/portail in JSON format + Systemd, + /// All traces (errors included) will be printed in JSON format on stdout and stderr + Container, +} + +impl From for LogPreset { + fn from(value: UserLogPreset) -> Self { + match value { + UserLogPreset::Development => Self::Development, + UserLogPreset::Systemd => Self::Systemd, + UserLogPreset::Container => Self::Container, + } + } +} + #[derive(Subcommand)] enum Commands { /// Run the portail daemon. @@ -75,6 +97,9 @@ enum Commands { #[arg(long, value_name = "FILE")] /// Path where to create the RPC socket if the daemon must create it itself bind_rpc_socket: Option, + /// Preset for logging, defaults to development + #[arg(long, default_value = "development")] + log_preset: UserLogPreset, }, /// Run RPC commands to the daemon. @@ -100,6 +125,9 @@ enum Commands { config: PathBuf, /// Path to the ACL file acl_file: PathBuf, + /// Whether to provide JSON output for scripting. + #[arg(long, value_name = "BOOLEAN", default_value_t = false)] + json: bool, }, } @@ -107,14 +135,19 @@ enum Commands { async fn main() -> Result<()> { let cli = Cli::parse(); - logging::init(); - match cli.command { Commands::Rpc { rpc_socket, json, command, } => { + let preset = if json { + LogPreset::Scripting + } else { + LogPreset::Cli + }; + let _guards = logging::init(preset).expect("Failed to initialize logging"); + let mut connection = zlink::unix::connect(&rpc_socket).await.context(format!( "Opening the RPC socket at path '{}'", rpc_socket.display() @@ -238,7 +271,9 @@ async fn main() -> Result<()> { config, bind_proxy_address, bind_rpc_socket, + log_preset, } => { + let _guards = logging::init(log_preset.into()).expect("Failed to initialize logging"); info!("Reading Portail settings from '{}'", config.display()); let settings: Arc = Arc::new(config::init(&config)); @@ -293,7 +328,17 @@ async fn main() -> Result<()> { info!("Exiting..."); } - Commands::CheckACLSyntax { config, acl_file } => { + Commands::CheckACLSyntax { + config, + json, + acl_file, + } => { + let preset = if json { + LogPreset::Scripting + } else { + LogPreset::Cli + }; + let guards = logging::init(preset).expect("Failed to initialize logging"); let contents = String::from_utf8_lossy( &std::fs::read(&acl_file).context("while reading ACL file")?, ) @@ -301,12 +346,13 @@ async fn main() -> Result<()> { let settings: Arc = Arc::new(config::init(&config)); match acl::load_rules_from_str(contents.as_str(), &settings) { Ok(rules) => info!( - "Parsed {} ACL policies and {} routes successfully", - rules.hir.policies.len(), - rules.hir.routes.len() + n_policies = %rules.hir.policies.len(), + n_routes = %rules.hir.routes.len(), + "Parsed ACL policies and routes successfully", ), Err(err) => { error!("Error while parsing the ACL rules:\n{err}"); + drop(guards); std::process::exit(1); } } From 4877751c822f79b4ae15ee3248359fb7a97729ef Mon Sep 17 00:00:00 2001 From: Ryan Lahfa Date: Mon, 15 Jun 2026 18:30:24 +0200 Subject: [PATCH 13/13] nix/module: add logging directives This makes full use of the logging work. Signed-off-by: Ryan Lahfa --- nix/module.nix | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/nix/module.nix b/nix/module.nix index 900d0df..0bb0f28 100644 --- a/nix/module.nix +++ b/nix/module.nix @@ -133,7 +133,7 @@ in # Use "notify-reload" when https://github.com/cloud-gouv/portail/issues/9 is done. Type = "notify"; NotifyAccess = "main"; - ExecStart = "${cfg.package}/bin/portail daemon --config ${configFile}"; + ExecStart = "${cfg.package}/bin/portail daemon --log-preset systemd --config ${configFile}"; # Enable when https://github.com/cloud-gouv/portail/issues/10 is done. # FileDescriptorStoreMax = 1000; @@ -154,6 +154,7 @@ in RuntimeDirectory = "portail"; StateDirectory = "portail"; + LogsDirectory = "portail"; }; }; };