Add custom logging

This commit is contained in:
Wez Furlong
2023-03-05 20:29:14 -07:00
parent a0e73b72b1
commit 1bfa3a42af
2 changed files with 505 additions and 95 deletions
+266 -95
View File
@@ -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<String>,
/// Where to place the log files; overrides the global setting
#[serde(default)]
pub log_dir: Option<PathBuf>,
#[serde(default = "default_true")]
pub enable: bool,
/// Instead of logging the json object, format it with this
/// minijinja template
#[serde(default)]
pub template: Option<String>,
}
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<String>,
#[serde(default)]
pub per_record: HashMap<RecordType, LogRecordParams>,
}
impl LogFileParams {
@@ -109,6 +135,7 @@ pub struct Logger {
thread: Mutex<Option<JoinHandle<()>>>,
meta: Vec<String>,
headers: Vec<String>,
enabled: HashMap<RecordType, bool>,
}
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 &params.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(&params.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 &params.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<LogCommand>) {
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<OpenedFile> = None;
fn do_record(
params: &LogFileParams,
file: &mut Option<OpenedFile>,
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(&params, &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<String, Value>, HashMap<String, Value>) {
pub fn extract_fields(
&self,
msg: &Message,
) -> (HashMap<String, Value>, HashMap<String, Value>) {
let mut headers = HashMap::new();
let mut meta = HashMap::new();
if !self.headers.is_empty() {
let mut all_headers :HashMap<String, Vec<Value>> = 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<String, Vec<Value>> = 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<String>,
}
struct OpenedFile {
file: AutoFinishEncoder<'static, File>,
name: PathBuf,
written: u64,
expires: Option<Instant>,
}
struct LogThreadState {
params: LogFileParams,
receiver: Receiver<LogCommand>,
template_engine: Environment<'static>,
file_map: HashMap<FileNameKey, OpenedFile>,
}
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<Instant> {
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<Template<'a>> {
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(())
}
}
+239
View File
@@ -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": "<somespammer@example.net>",
"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": [
"<user@example.com>",
],
"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: <somespammer@example.net>
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: <Undisclosed Recipients>
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",
},
}
}
```