diff --git a/Cargo.lock b/Cargo.lock index f5afb8b..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" @@ -539,7 +548,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]] @@ -1221,7 +1230,9 @@ dependencies = [ "tokio-rustls", "toml", "tracing", + "tracing-appender", "tracing-subscriber", + "uuid", "winnow", "zlink", ] @@ -1391,7 +1402,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.52.0", + "windows-sys 0.61.2", ] [[package]] @@ -1646,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" @@ -1664,10 +1681,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]] @@ -1889,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" @@ -1921,6 +1951,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 +1971,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]] @@ -1975,6 +2018,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" @@ -2135,7 +2189,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..2b8ad81 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -26,10 +26,12 @@ 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", "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/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"; }; }; }; 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 e59933e..60eda2d 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,18 +1,18 @@ 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::{error, level_filters::LevelFilter}; -use tracing::{info, warn}; -use tracing_subscriber::EnvFilter; +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; mod acl; mod config; +mod logging; mod proxy; mod rpc; mod state; @@ -63,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. @@ -76,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. @@ -101,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, }, } @@ -108,26 +135,26 @@ 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(); - 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() + ))?; + 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 @@ -149,10 +176,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 @@ -173,10 +196,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 @@ -197,9 +216,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 @@ -255,13 +271,18 @@ 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)); 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 { @@ -288,22 +309,36 @@ 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 } => { + 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")?, ) @@ -311,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); } } diff --git a/src/proxy/context.rs b/src/proxy/context.rs index cc628ce..d1b51dd 100644 --- a/src/proxy/context.rs +++ b/src/proxy/context.rs @@ -49,13 +49,15 @@ impl From for TargetAddr { } } +#[derive(Debug)] pub struct TargetContext { pub initial_target: TargetAddr, pub resolved_target: Option, } #[derive(Debug, Clone)] -pub struct InitialRequestContext { +pub struct OwnedRequestContext { + pub trace_id: uuid::Uuid, pub client_address: SocketAddr, pub acl_ctx: crate::acl::OwnedEvaluationContext, } @@ -64,13 +66,16 @@ 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>, } -impl InitialRequestContext { +impl OwnedRequestContext { pub fn new(client_address: SocketAddr) -> Self { Self { client_address, + trace_id: uuid::Uuid::new_v4(), acl_ctx: crate::acl::OwnedEvaluationContext::empty(), } } @@ -79,6 +84,7 @@ impl InitialRequestContext { LocalRequestContext { client_address: &self.client_address, acl_ctx: self.acl_ctx.fork(), + trace_id: self.trace_id, } } } diff --git a/src/proxy/http_connect.rs b/src/proxy/http_connect.rs index 9db70fc..9ca22f8 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}, @@ -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: InitialRequestContext, + 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: InitialRequestContext, + ctx: OwnedRequestContext, stream: S, ) -> Result<(), ProxyError> { + debug!(subsystem = "proxy_access", "HTTP/2 CONNECT request"); let io = TokioIo::new(stream); // TODO: update ctx @@ -102,19 +107,21 @@ 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(); + 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 + ); } }); diff --git a/src/proxy/mod.rs b/src/proxy/mod.rs index f024d7b..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,9 +10,9 @@ use tokio_rustls::{ server::{VerifierBuilderError, WebPkiClientVerifier}, }, }; -use tracing::error; +use tracing::{Instrument, debug, 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?; @@ -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?; - - let acceptor = tls_acceptor.clone(); - let settings = settings.clone(); - let state = state.clone(); - let ctx = InitialRequestContext::new(addr); - - tokio::spawn(async move { + 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(); + + 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; } } diff --git a/src/proxy/socks5.rs b/src/proxy/socks5.rs index 2af95cb..707b6da 100644 --- a/src/proxy/socks5.rs +++ b/src/proxy/socks5.rs @@ -11,14 +11,14 @@ use tokio::{ io::{AsyncRead, AsyncWrite}, net::TcpStream, sync::RwLock, - time::timeout, + time::{Instant, timeout}, }; use tokio_rustls::TlsStream; use tracing::{debug, info, warn}; use crate::{ config::{BackendSettings, KnownBackend, Settings}, - proxy::context::{InitialRequestContext, TargetContext}, + proxy::context::{OwnedRequestContext, TargetContext}, state::State, }; @@ -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: InitialRequestContext, + 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(()) 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()