diff --git a/crates/kumod/src/logging.rs b/crates/kumod/src/logging.rs index b08af7f8..b688f7b4 100644 --- a/crates/kumod/src/logging.rs +++ b/crates/kumod/src/logging.rs @@ -1,6 +1,4 @@ use crate::mx::ResolvedAddress; -use std::collections::HashMap; -use serde_json::Value; use crate::smtp_server::RelayDisposition; use anyhow::{anyhow, Context}; use async_channel::{Receiver, Sender}; @@ -9,9 +7,12 @@ use chrono::{DateTime, Utc}; use message::rfc3464::ReportAction; use message::rfc5965::ARFReport; use message::Message; +use minijinja::{Environment, Source, Template}; use once_cell::sync::OnceCell; use rfc5321::{EnhancedStatusCode, Response}; use serde::{Deserialize, Serialize}; +use serde_json::Value; +use std::collections::HashMap; use std::fs::File; use std::io::Write; use std::net::Ipv4Addr; @@ -56,6 +57,28 @@ impl ClassifierParams { } } +#[derive(Deserialize, Clone, Debug)] +pub struct LogRecordParams { + #[serde(default)] + pub suffix: Option, + + /// Where to place the log files; overrides the global setting + #[serde(default)] + pub log_dir: Option, + + #[serde(default = "default_true")] + pub enable: bool, + + /// Instead of logging the json object, format it with this + /// minijinja template + #[serde(default)] + pub template: Option, +} + +fn default_true() -> bool { + true +} + #[derive(Deserialize, Clone, Debug)] pub struct LogFileParams { /// Where to place the log files @@ -85,6 +108,9 @@ pub struct LogFileParams { /// List of message headers to capture in the log #[serde(default)] pub headers: Vec, + + #[serde(default)] + pub per_record: HashMap, } impl LogFileParams { @@ -109,6 +135,7 @@ pub struct Logger { thread: Mutex>>, meta: Vec, headers: Vec, + enabled: HashMap, } impl Logger { @@ -117,10 +144,32 @@ impl Logger { } pub fn init(params: LogFileParams) -> anyhow::Result<()> { + let mut source = Source::new(); + + for (kind, per_rec) in ¶ms.per_record { + if let Some(template_source) = &per_rec.template { + source + .add_template(format!("{kind:?}"), template_source) + .with_context(|| { + format!( + "compiling template:\n{template_source}\nfor log record type {kind:?}" + ) + })?; + } + } + + let mut template_engine = Environment::new(); + template_engine.set_source(source); + std::fs::create_dir_all(¶ms.log_dir) .with_context(|| format!("creating log directory {}", params.log_dir.display()))?; - let headers =params.headers.clone(); + let mut enabled = HashMap::new(); + for (kind, cfg) in ¶ms.per_record { + enabled.insert(*kind, cfg.enable); + } + + let headers = params.headers.clone(); let meta = params.meta.clone(); let (sender, receiver) = async_channel::bounded(params.back_pressure); let thread = std::thread::spawn(move || { @@ -128,7 +177,15 @@ impl Logger { .enable_time() .build() .expect("create logger runtime"); - runtime.block_on(Self::logger_thread(params, receiver)); + runtime.block_on(async move { + let mut state = LogThreadState { + params, + receiver, + template_engine, + file_map: HashMap::new(), + }; + state.logger_thread().await + }); }); let logger = Self { @@ -136,6 +193,7 @@ impl Logger { thread: Mutex::new(Some(thread)), meta, headers, + enabled, }; LOGGER @@ -144,94 +202,14 @@ impl Logger { Ok(()) } - async fn logger_thread(params: LogFileParams, receiver: Receiver) { - struct OpenedFile { - file: AutoFinishEncoder<'static, File>, - name: PathBuf, - written: u64, - opened: Instant, + pub fn record_is_enabled(&self, kind: RecordType) -> bool { + if let Some(enabled) = self.enabled.get(&kind) { + return *enabled; } - - let mut file: Option = None; - - fn do_record( - params: &LogFileParams, - file: &mut Option, - mut record: JsonLogRecord, - ) -> anyhow::Result<()> { - if let Some(classifier) = CLASSIFY.get() { - record.bounce_classification = classifier.classify_response(&record.response); - } - if file.is_none() { - let now = Utc::now(); - let name = params.log_dir.join(now.format("%Y%m%d-%H%M%S").to_string()); - - let f = std::fs::OpenOptions::new() - .append(true) - .create(true) - .open(&name) - .with_context(|| format!("open log file {name:?}"))?; - - file.replace(OpenedFile { - file: Encoder::new(f, params.compression_level) - .context("set up zstd encoder")? - .auto_finish(), - name, - written: 0, - opened: Instant::now(), - }); - } - - let mut need_rotate = false; - - if let Some(file) = file.as_mut() { - let mut json = serde_json::to_string(&record).context("serializing record")?; - json.push_str("\n"); - file.file - .write_all(json.as_bytes()) - .with_context(|| format!("writing record to {}", file.name.display()))?; - file.written += json.len() as u64; - - need_rotate = file.written >= params.max_file_size; - } - - if need_rotate { - file.take(); - } - - Ok(()) - } - - loop { - let cmd = if let Some(deadline) = params - .max_segment_duration - .and_then(|duration| file.as_ref().and_then(|of| Some(of.opened + duration))) - { - tokio::select! { - cmd = receiver.recv() => cmd, - _ = tokio::time::sleep_until(deadline.into()) => { - file.take(); - continue; - } - } - } else { - receiver.recv().await - }; - let cmd = match cmd { - Ok(cmd) => cmd, - _ => return, - }; - match cmd { - LogCommand::Terminate => { - break; - } - LogCommand::Record(record) => { - if let Err(err) = do_record(¶ms, &mut file, record) { - tracing::error!("failed to log: {err:#}"); - } - } - } + if let Some(enabled) = self.enabled.get(&RecordType::Any) { + return *enabled; } + true } pub async fn log(&self, record: JsonLogRecord) -> anyhow::Result<()> { @@ -250,14 +228,20 @@ impl Logger { } } - pub fn extract_fields(&self, msg: &Message) -> (HashMap, HashMap) { + pub fn extract_fields( + &self, + msg: &Message, + ) -> (HashMap, HashMap) { let mut headers = HashMap::new(); let mut meta = HashMap::new(); if !self.headers.is_empty() { - let mut all_headers :HashMap> = HashMap::new(); - for (name, value) in msg.get_all_headers().unwrap_or_else(|_|vec![]) { - all_headers.entry(name.to_ascii_lowercase()).or_default().push(value.into()); + let mut all_headers: HashMap> = HashMap::new(); + for (name, value) in msg.get_all_headers().unwrap_or_else(|_| vec![]) { + all_headers + .entry(name.to_ascii_lowercase()) + .or_default() + .push(value.into()); } for name in &self.headers { @@ -283,7 +267,7 @@ impl Logger { } } -#[derive(Serialize, Deserialize, Debug, Copy, Clone, Eq, PartialEq)] +#[derive(Serialize, Deserialize, Debug, Copy, Clone, Eq, PartialEq, Hash)] pub enum RecordType { /// Recorded by a receiving listener Reception, @@ -300,6 +284,9 @@ pub enum RecordType { OOB, /// Contains a feedback report Feedback, + + /// Special for matching anything in the logging config + Any, } #[derive(Serialize, Deserialize, Debug)] @@ -381,6 +368,10 @@ pub async fn log_disposition(args: LogDisposition<'_>) { } } + if !logger.record_is_enabled(kind) { + return; + } + let (headers, meta) = logger.extract_fields(&msg); let record = JsonLogRecord { @@ -499,3 +490,183 @@ pub async fn log_disposition(args: LogDisposition<'_>) { } } } + +#[derive(Clone, Debug, Hash, PartialEq, Eq)] +struct FileNameKey { + log_dir: PathBuf, + suffix: Option, +} + +struct OpenedFile { + file: AutoFinishEncoder<'static, File>, + name: PathBuf, + written: u64, + expires: Option, +} + +struct LogThreadState { + params: LogFileParams, + receiver: Receiver, + template_engine: Environment<'static>, + file_map: HashMap, +} + +impl LogThreadState { + async fn logger_thread(&mut self) { + loop { + let deadline = self.get_deadline(); + + let cmd = if let Some(deadline) = deadline { + tokio::select! { + cmd = self.receiver.recv() => cmd, + _ = tokio::time::sleep_until(deadline.into()) => { + self.expire(); + continue; + } + } + } else { + self.receiver.recv().await + }; + let cmd = match cmd { + Ok(cmd) => cmd, + _ => return, + }; + match cmd { + LogCommand::Terminate => { + break; + } + LogCommand::Record(record) => { + if let Err(err) = self.do_record(record) { + tracing::error!("failed to log: {err:#}"); + }; + } + } + } + } + + fn expire(&mut self) { + let now = Instant::now(); + self.file_map.retain(|_, of| match of.expires { + Some(exp) => exp > now, + None => true, + }); + } + + fn get_deadline(&self) -> Option { + self.file_map + .values() + .filter_map(|of| of.expires.clone()) + .min() + } + + fn per_record(&self, kind: RecordType) -> Option<&LogRecordParams> { + self.params + .per_record + .get(&kind) + .or_else(|| self.params.per_record.get(&RecordType::Any)) + } + + fn resolve_template<'a>( + params: &LogFileParams, + template_engine: &'a Environment, + kind: RecordType, + ) -> Option> { + if let Some(pr) = params.per_record.get(&kind) { + if pr.template.is_some() { + let label = format!("{kind:?}"); + return template_engine.get_template(&label).ok(); + } + return None; + } + if let Some(pr) = params.per_record.get(&RecordType::Any) { + if pr.template.is_some() { + return template_engine.get_template("Any").ok(); + } + } + None + } + + fn do_record(&mut self, mut record: JsonLogRecord) -> anyhow::Result<()> { + let file_key = if let Some(per_rec) = self.per_record(record.kind) { + FileNameKey { + log_dir: per_rec + .log_dir + .as_deref() + .unwrap_or(&self.params.log_dir) + .to_path_buf(), + suffix: per_rec.suffix.clone(), + } + } else { + // Just use the global settings + FileNameKey { + log_dir: self.params.log_dir.clone(), + suffix: None, + } + }; + + if let Some(classifier) = CLASSIFY.get() { + record.bounce_classification = classifier.classify_response(&record.response); + } + + if !self.file_map.contains_key(&file_key) { + let now = Utc::now(); + + let mut base_name = now.format("%Y%m%d-%H%M%S").to_string(); + if let Some(suffix) = &file_key.suffix { + base_name.push_str(suffix); + } + + let name = file_key.log_dir.join(base_name); + + let f = std::fs::OpenOptions::new() + .append(true) + .create(true) + .open(&name) + .with_context(|| format!("open log file {name:?}"))?; + + self.file_map.insert( + file_key.clone(), + OpenedFile { + file: Encoder::new(f, self.params.compression_level) + .context("set up zstd encoder")? + .auto_finish(), + name, + written: 0, + expires: self + .params + .max_segment_duration + .map(|duration| Instant::now() + duration), + }, + ); + } + + let mut need_rotate = false; + + if let Some(file) = self.file_map.get_mut(&file_key) { + let mut record_text = Vec::new(); + + if let Some(template) = + Self::resolve_template(&self.params, &self.template_engine, record.kind) + { + template.render_to_write(&record, &mut record_text)?; + } else { + serde_json::to_writer(&mut record_text, &record).context("serializing record")?; + } + if record_text.last() != Some(&b'\n') { + record_text.push(b'\n'); + } + file.file + .write_all(&record_text) + .with_context(|| format!("writing record to {}", file.name.display()))?; + file.written += record_text.len() as u64; + + need_rotate = file.written >= self.params.max_file_size; + } + + if need_rotate { + self.file_map.remove(&file_key); + } + + Ok(()) + } +} diff --git a/docs/reference/kumo/configure_local_logs.md b/docs/reference/kumo/configure_local_logs.md index 95b93511..0238f792 100644 --- a/docs/reference/kumo/configure_local_logs.md +++ b/docs/reference/kumo/configure_local_logs.md @@ -120,3 +120,242 @@ kumo.configure_local_logs { } ``` +## per_record + +Allows configuring per-record type logging + +```lua +kumo.configure_local_logs { + per_record = { + Reception = { + -- use names like "20230306-022811_recv" for reception logs + suffix = '_recv', + }, + + Delivery = { + -- put delivery logs in a different directory + log_dir = '/var/log/kumo/delivery', + }, + + TransientFailure = { + -- Don't log transient failures + enable = false, + }, + + Bounce = { + -- Instead of logging the json record, evaluate this + -- template string and log the result. + template = [[Bounce! id={{ id }}, from={{ sender }} code={{ code }} age={{ timestamp - created }}]], + }, + + -- For any record type not explicitly listed, apply these settings. + -- This effectively turns off all other log records + Any = { + enable = false, + }, + }, +} +``` + +The [Mini Jinja](https://docs.rs/minijinja/latest/minijinja/) templating engine +is used to evalute logging templates. The full supported syntax is [documented +here](https://docs.rs/minijinja/latest/minijinja/syntax/index.html). + +The JSON log record fields shown in the section below are assigned as template +variables, so using `{{ id }}` in your log template will be substituted with +the `id` field from the log record section below. + +# Log Record + +The log record is a JSON object with the following shape: + +```json +{ + // The record type; can be one of "Reception", "Delivery", + // "Bounce", "TransientFailure", "Expiration", "AdminBounce", + // "OOB" or "Feedback" + "type": "Delivery", + + // The message spool id; corresponds to the value returned by + // message:id() + "id": "1d98076abbbc11ed940250ebf67f93bd", + + // The envelope sender + "sender": "user@sender.example.com", + + // The envelope recipient + "recipient": "user@recipient.example.com", + + // Which named queue the message was associaed with + "queue": "campaign:tenant@domain", + + // Which MX site the message was being delivered to. + // Empty string for Reception records. + "site": "source2->(alt1|alt2|alt3|alt4)?.gmail-smtp-in.l.google.com.", + + // The size of the message payload, in bytes + "size": 1047, + + // The response from the peer, if applicable + "response": { + // the SMTP status code + "code": 250, + + // The ENHANCEDSTATUSCODE portion of the response parsed + // out into individual fields. + // This one is from a "2.0.0" status code + "enhanced_code": { + "class": 2, + "subject": 0, + "detail": 0, + }, + + // the remainder of the response content + "content": "OK ids=8a5475ccbbc611eda12250ebf67f93bd", + + // the SMTP command verb to which the response was made. + // eg: "MAIL FROM", "RCPT TO" etc. "." isn't really a command + // but is used to represent the response to the final ".: + // we send to indicate the end of the message payload. + "command": "." + }, + + // Information about the peer in the communication. This is either + // the submitter or the receiver, depending on the record type + "peer_address": { + // When delivering, this is the name from the MX record. + // When receiving, this is the EHLO/HELO string sent by + // the sender + "name": "gmail-smtp-in.l.google.com.", + "addr": "142.251.2.27" + }, + + // The time at which this record was generated, expressed + // as a unix timestamp: seconds since the unix epoch + "timestamp": 1678069691, + + // The time at which the message was received, expressed + // as a unix timestamp: seconds since the unix epoch + "created": 1678069691, + + // The number of delivery attempts. + "num_attempts": 0, + + // the classification assigned by the bounce classifier, + // or Uncategorized if unknown or the classifier is not configured. + "bounce_classification": "Uncategorized", + + // The name of the egress pool used as the source for the delivery + "egress_pool": "pool0", + + // The name of the selected egress source (a member of the egress pool) + // used for the delivery + "egress_source": "source2", + + // when "type" == "Feedback", holds the parsed feedback report + "feedback_report": null, + + // holds the values of the list of meta fields from the logger + // configuration + "meta": {}, + + // holds the values of the list of message headers from the logger + // configuration + "headers": {} +} +``` + +## Record Types + +The following record types are defined: + +* `"Reception"` - logging the reception of a message via SMTP or via + the HTTP injection API +* `"Delivery"` - logging the successful delivery of a message via SMTP +* `"Bounce"` - logging a permanent failure response and end of delivery + attempts for the message. +* `"TransientFailure"` - logging a transient failure when attempting delivery +* `"Expiration"` - logged when the message exceeds the configured maximum + lifetime in the queue. +* `"AdminBounce"` - logged when an administrator uses the `/api/admin/bounce` + API to fail message(s). +* `"OOB"` - when receiving an out of band bounce with an attached RFC3464 + delivery status report, the parsed report is used to synthesize an OOB + record for each recipient in the report. +* `"Feedback"` - when receiving an ARF feedback report, instead of logging + a `"Reception"`, a `"Feedback" record is logged instead with the report + contents parsed out and made available in the `feedback_report` field. + +## Feedback Report + +ARF feedback reports are parsed into a JSON object that has the following +structure. The fields of the `feedback_report` correspond to those defined +by RFC 5965. + +See also [trace_headers](start_esmtp_listener.md#trace_headers) for information +about the `supplemental_trace` field. + +```json +{ + "type": "Feedback", + "feedback_report": { + "feedback_type": "abuse", + "user_agent": "SomeGenerator/1.0", + "version": 1, + "arrival_date": "2005-03-08T18:00:00Z", + "incidents": nil, + "original_envelope_id": nil, + "original_mail_from": "", + "reporting_mta": { + "mta_type": "dns", + "name": "mail.example.com", + }, + "source_ip": "192.0.2.1", + "authentication_results": [ + "mail.example.com; spf=fail smtp.mail=somespammer@example.com", + ], + "original_rcpto_to": [ + "", + ], + "reported_domain": [ + "example.net", + ], + "reported_uri": [ + "http://example.net/earn_money.html", + "mailto:user@example.com", + ], + "extensions": { + "removal-recipient": [ + "user@example.com", + ], + }, + + // The original message or message headers, if provided in + // the report + "original_message": "From: +Received: from mailserver.example.net (mailserver.example.net + [192.0.2.1]) by example.com with ESMTP id M63d4137594e46; + Tue, 08 Mar 2005 14:00:00 -0400 +X-KumoRef: eyJfQF8iOiJcXF8vIiwicmVjaXBpZW50IjoidGVzdEBleGFtcGxlLmNvbSJ9 +To: +Subject: Earn money +MIME-Version: 1.0 +Content-type: text/plain +Message-ID: 8787KJKJ3K4J3K4J3K4J3.mail@example.net +Date: Thu, 02 Sep 2004 12:31:03 -0500 + +Spam Spam Spam +Spam Spam Spam +Spam Spam Spam +Spam Spam Spam +", + + // if original_message is present, and a kumo-style trace + // header was decoded from it, then this holds the decoded + // trace information + "supplemental_trace": { + "recipient": "test@example.com", + }, + } +} +```