Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
293 changes: 256 additions & 37 deletions src/logger.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,8 @@
mod file;
mod filter;

use std::io::{self, Write};
use std::fmt::Write as _;
use std::io::{self, Write as IoWrite};
use std::sync::Mutex;
use std::time::{SystemTime, UNIX_EPOCH};

Expand All @@ -10,23 +11,26 @@ use log::{Log, Metadata, Record};
use self::file::{DailyFileWriter, resolve_log_target};
use self::filter::LogFilter;

const INITIAL_RECORD_CAPACITY: usize = 256;
const MAX_RETAINED_RECORD_CAPACITY: usize = 64 * 1024;

pub(crate) fn init(log_to_file: bool) {
let filter = LogFilter::from_env();
let max_level = filter.max_level();
let writer: Box<dyn Write + Send> = if log_to_file {
let sink = if log_to_file {
let target = resolve_log_target();
let file_writer = DailyFileWriter::new(target);
Box::new(TeeWriter::new(
Box::new(io::stderr()),
Box::new(file_writer),
))
LogSink::tee(Box::new(io::stderr()), Box::new(file_writer))
} else {
Box::new(io::stderr())
LogSink::single(Box::new(io::stderr()))
};

let logger = SimpleLogger {
filter,
writer: Mutex::new(writer),
state: Mutex::new(LoggerState {
sink,
record_buffer: String::with_capacity(INITIAL_RECORD_CAPACITY),
}),
include_timestamp: log_to_file,
};

Expand All @@ -37,10 +41,44 @@ pub(crate) fn init(log_to_file: bool) {

struct SimpleLogger {
filter: LogFilter,
writer: Mutex<Box<dyn Write + Send>>,
state: Mutex<LoggerState>,
include_timestamp: bool,
}

struct LoggerState {
sink: LogSink,
record_buffer: String,
}

impl LoggerState {
fn write_record(&mut self, record: &Record<'_>, include_timestamp: bool) {
self.record_buffer.clear();
if include_timestamp {
let _ = writeln!(
self.record_buffer,
"{} {} {}: {}",
timestamp_millis(),
record.level(),
record.target(),
record.args()
);
} else {
let _ = writeln!(
self.record_buffer,
"{} {}: {}",
record.level(),
record.target(),
record.args()
);
}

self.sink.write_record(self.record_buffer.as_bytes());
if self.record_buffer.capacity() > MAX_RETAINED_RECORD_CAPACITY {
self.record_buffer = String::with_capacity(INITIAL_RECORD_CAPACITY);
}
}
}

impl Log for SimpleLogger {
fn enabled(&self, metadata: &Metadata<'_>) -> bool {
self.filter.enabled(metadata.target(), metadata.level())
Expand All @@ -51,26 +89,17 @@ impl Log for SimpleLogger {
return;
}

let mut writer = match self.writer.lock() {
Ok(writer) => writer,
let mut state = match self.state.lock() {
Ok(state) => state,
Err(poisoned) => poisoned.into_inner(),
};

if self.include_timestamp {
let _ = write!(writer, "{} ", timestamp_millis());
}
let _ = writeln!(
writer,
"{} {}: {}",
record.level(),
record.target(),
record.args()
);
state.write_record(record, self.include_timestamp);
}

fn flush(&self) {
if let Ok(mut writer) = self.writer.lock() {
let _ = writer.flush();
if let Ok(mut state) = self.state.lock() {
state.sink.flush();
}
}
}
Expand All @@ -82,26 +111,216 @@ fn timestamp_millis() -> u128 {
.unwrap_or(0)
}

struct TeeWriter {
left: Box<dyn Write + Send>,
right: Box<dyn Write + Send>,
enum LogSink {
Single(Box<dyn IoWrite + Send>),
Tee {
primary: Box<dyn IoWrite + Send>,
secondary: Box<dyn IoWrite + Send>,
},
}

impl TeeWriter {
fn new(left: Box<dyn Write + Send>, right: Box<dyn Write + Send>) -> Self {
Self { left, right }
impl LogSink {
fn single(writer: Box<dyn IoWrite + Send>) -> Self {
Self::Single(writer)
}

fn tee(primary: Box<dyn IoWrite + Send>, secondary: Box<dyn IoWrite + Send>) -> Self {
Self::Tee { primary, secondary }
}

fn write_record(&mut self, record: &[u8]) {
match self {
Self::Single(writer) => {
let _ = writer.write_all(record);
}
Self::Tee { primary, secondary } => {
let _ = primary.write_all(record);
let _ = secondary.write_all(record);
}
}
}

fn flush(&mut self) {
match self {
Self::Single(writer) => {
let _ = writer.flush();
}
Self::Tee { primary, secondary } => {
let _ = primary.flush();
let _ = secondary.flush();
}
}
}
}

impl Write for TeeWriter {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
self.left.write_all(buf)?;
self.right.write_all(buf)?;
Ok(buf.len())
#[cfg(test)]
mod tests {
use super::*;
use log::Level;
use std::sync::{Arc, Mutex};

#[derive(Clone, Default)]
struct SharedWriter {
state: Arc<Mutex<SharedWriterState>>,
}

#[derive(Default)]
struct SharedWriterState {
bytes: Vec<u8>,
flushes: usize,
}

impl SharedWriter {
fn bytes(&self) -> Vec<u8> {
self.state
.lock()
.expect("shared writer state")
.bytes
.clone()
}

fn flushes(&self) -> usize {
self.state.lock().expect("shared writer state").flushes
}
}

impl IoWrite for SharedWriter {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
let mut state = self.state.lock().expect("shared writer state");
state.bytes.extend_from_slice(buf);
Ok(buf.len())
}

fn flush(&mut self) -> io::Result<()> {
self.state.lock().expect("shared writer state").flushes += 1;
Ok(())
}
}

struct FailingWriter;

impl IoWrite for FailingWriter {
fn write(&mut self, _buf: &[u8]) -> io::Result<usize> {
Err(io::Error::new(io::ErrorKind::BrokenPipe, "sink closed"))
}

fn flush(&mut self) -> io::Result<()> {
Err(io::Error::new(io::ErrorKind::BrokenPipe, "sink closed"))
}
}

struct InterruptingWriteAll;

impl IoWrite for InterruptingWriteAll {
fn write(&mut self, _buf: &[u8]) -> io::Result<usize> {
Err(io::Error::new(
io::ErrorKind::Interrupted,
"sink interrupted",
))
}

fn write_all(&mut self, _buf: &[u8]) -> io::Result<()> {
Err(io::Error::new(
io::ErrorKind::Interrupted,
"sink interrupted",
))
}

fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}

fn warning_record() -> Record<'static> {
Record::builder()
.args(format_args!("sink failure record"))
.level(Level::Warn)
.target("wayscriber::daemon")
.build()
}

#[test]
fn failed_primary_sink_does_not_starve_the_secondary_sink() {
let secondary = SharedWriter::default();
let mut sink = LogSink::tee(Box::new(FailingWriter), Box::new(secondary.clone()));

sink.write_record(b"WARN wayscriber::daemon: sink failure record\n");
assert_eq!(
secondary.bytes(),
b"WARN wayscriber::daemon: sink failure record\n"
);
}

#[test]
fn failed_secondary_sink_does_not_starve_the_primary_sink() {
let primary = SharedWriter::default();
let mut sink = LogSink::tee(Box::new(primary.clone()), Box::new(FailingWriter));

sink.write_record(b"WARN wayscriber::daemon: sink failure record\n");
assert_eq!(
primary.bytes(),
b"WARN wayscriber::daemon: sink failure record\n"
);
}

fn flush(&mut self) -> io::Result<()> {
self.left.flush()?;
self.right.flush()
#[test]
fn interrupted_sink_does_not_duplicate_the_healthy_sink_record() {
let secondary = SharedWriter::default();
let mut sink = LogSink::tee(Box::new(InterruptingWriteAll), Box::new(secondary.clone()));

sink.write_record(b"WARN wayscriber::daemon: sink failure record\n");

assert_eq!(
secondary.bytes(),
b"WARN wayscriber::daemon: sink failure record\n"
);
assert_eq!(
secondary
.bytes()
.windows(b"sink failure record".len())
.filter(|window| *window == b"sink failure record")
.count(),
1,
"the healthy sink receives the record exactly once"
);
}

#[test]
fn failed_flush_still_flushes_the_other_sink() {
let secondary = SharedWriter::default();
let mut sink = LogSink::tee(Box::new(FailingWriter), Box::new(secondary.clone()));

sink.flush();
assert_eq!(secondary.flushes(), 1);
}

#[test]
fn file_record_keeps_timestamp_and_message_in_one_line() {
let output = SharedWriter::default();
let mut state = LoggerState {
sink: LogSink::single(Box::new(output.clone())),
record_buffer: String::with_capacity(INITIAL_RECORD_CAPACITY),
};

state.write_record(&warning_record(), true);

let line = String::from_utf8(output.bytes()).expect("valid log output");
let (timestamp, message) = line.split_once(' ').expect("timestamp separator");

assert!(timestamp.parse::<u128>().is_ok());
assert_eq!(message, "WARN wayscriber::daemon: sink failure record\n");
}

#[test]
fn oversized_record_buffer_is_not_retained() {
let output = SharedWriter::default();
let mut state = LoggerState {
sink: LogSink::single(Box::new(output)),
record_buffer: String::with_capacity(MAX_RETAINED_RECORD_CAPACITY + 1),
};

state.write_record(&warning_record(), false);

assert!(state.record_buffer.capacity() <= MAX_RETAINED_RECORD_CAPACITY);
}
}
Loading