From 2bdcd9cb4ef3ef84223cfe155991f58d3a3df3e8 Mon Sep 17 00:00:00 2001 From: Wez Furlong Date: Thu, 4 May 2023 15:22:43 -0700 Subject: [PATCH] Allow configuring multiple instances of local logging --- crates/kumod/src/logging.rs | 50 +++++++++++++++++++++---------------- docs/changelog/main.md | 3 +++ 2 files changed, 31 insertions(+), 22 deletions(-) diff --git a/crates/kumod/src/logging.rs b/crates/kumod/src/logging.rs index d0973d15..5103414b 100644 --- a/crates/kumod/src/logging.rs +++ b/crates/kumod/src/logging.rs @@ -7,7 +7,7 @@ use kumo_log_types::rfc3464::ReportAction; pub use kumo_log_types::*; use message::Message; use minijinja::{Environment, Source, Template}; -use once_cell::sync::OnceCell; +use once_cell::sync::{Lazy, OnceCell}; use rfc5321::{EnhancedStatusCode, Response}; use serde::Deserialize; use serde_json::Value; @@ -16,12 +16,12 @@ use std::fs::File; use std::io::Write; use std::net::Ipv4Addr; use std::path::PathBuf; -use std::sync::Mutex; +use std::sync::{Arc, Mutex}; use std::thread::JoinHandle; use std::time::{Duration, Instant}; use zstd::stream::write::Encoder; -static LOGGER: OnceCell = OnceCell::new(); +static LOGGER: Lazy>>> = Lazy::new(|| Mutex::new(vec![])); static CLASSIFY: OnceCell = OnceCell::new(); #[derive(Deserialize, Clone, Debug)] @@ -139,8 +139,8 @@ pub struct Logger { } impl Logger { - pub fn get() -> Option<&'static Logger> { - LOGGER.get() + fn get_loggers() -> Vec> { + LOGGER.lock().unwrap().iter().map(Arc::clone).collect() } pub fn init(params: LogFileParams) -> anyhow::Result<()> { @@ -200,9 +200,7 @@ impl Logger { enabled, }; - LOGGER - .set(logger) - .map_err(|_| anyhow::anyhow!("logger already initialized"))?; + LOGGER.lock().unwrap().push(Arc::new(logger)); Ok(()) } @@ -221,7 +219,8 @@ impl Logger { } pub async fn signal_shutdown() { - if let Some(logger) = Self::get() { + let loggers = Self::get_loggers(); + for logger in loggers.iter() { logger.sender.send(LogCommand::Terminate).await.ok(); logger .thread @@ -296,22 +295,29 @@ pub async fn log_disposition(args: LogDisposition<'_>) { relay_disposition, } = args; - if let Some(logger) = Logger::get() { - let mut feedback_report = None; + let loggers = Logger::get_loggers(); + if loggers.is_empty() { + return; + } - msg.load_meta_if_needed().await.ok(); + let mut feedback_report = None; - if kind == RecordType::Reception { - if let Some(RelayDisposition { log_arf: true, .. }) = relay_disposition { - if let Ok(Some(report)) = msg.parse_rfc5965() { - feedback_report.replace(report); - kind = RecordType::Feedback; - } + msg.load_meta_if_needed().await.ok(); + + if kind == RecordType::Reception { + if let Some(RelayDisposition { log_arf: true, .. }) = relay_disposition { + if let Ok(Some(report)) = msg.parse_rfc5965() { + feedback_report.replace(report); + kind = RecordType::Feedback; } } + } + let now = Utc::now(); + + for logger in loggers.iter() { if !logger.record_is_enabled(kind) { - return; + continue; } let (headers, meta) = logger.extract_fields(&msg).await; @@ -333,14 +339,14 @@ pub async fn log_disposition(args: LogDisposition<'_>) { .unwrap_or_else(|err| format!("{err:#}")), site: site.to_string(), peer_address: peer_address.cloned(), - response, - timestamp: Utc::now(), + response: response.clone(), + timestamp: now, created: msg.id().created(), num_attempts: msg.get_num_attempts(), egress_pool: egress_pool.map(|s| s.to_string()), egress_source: egress_source.map(|s| s.to_string()), bounce_classification: BounceClass::Uncategorized, - feedback_report, + feedback_report: feedback_report.clone(), headers, meta, }; diff --git a/docs/changelog/main.md b/docs/changelog/main.md index 094c14a9..4ea625b3 100644 --- a/docs/changelog/main.md +++ b/docs/changelog/main.md @@ -2,3 +2,6 @@ * Expose ready queue size to metrics. #30 * Fixed IPv6 lookups for domains without MX records +* [kumo.configure_local_logs](../reference/kumo/configure_local_logs.md) can now be called + multiple times to configure multiple different logging locations and + configurations.