#![allow(dead_code)]
use std::error::Error;
use std::fmt;
use std::io::{self, BufWriter, Stdout, Write};
use std::sync::{Mutex, PoisonError};
use crate::decode::LogRecord;
use crate::formatter::Formatter;
use crate::level::LogLevel;
#[non_exhaustive]
pub enum SinkError {
Io(io::Error),
Other(Box<dyn Error + Send + Sync + 'static>),
}
impl SinkError {
pub fn other(e: impl Error + Send + Sync + 'static) -> Self {
Self::Other(Box::new(e))
}
}
impl fmt::Display for SinkError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Io(e) => write!(f, "I/O error: {e}"),
Self::Other(e) => fmt::Display::fmt(e, f),
}
}
}
impl fmt::Debug for SinkError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Io(e) => f.debug_tuple("Io").field(e).finish(),
Self::Other(e) => f.debug_tuple("Other").field(e).finish(),
}
}
}
impl Error for SinkError {
fn source(&self) -> Option<&(dyn Error + 'static)> {
match self {
Self::Io(e) => Some(e),
Self::Other(e) => Some(e.as_ref()),
}
}
}
impl From<io::Error> for SinkError {
fn from(e: io::Error) -> Self {
Self::Io(e)
}
}
pub trait Sink: Send + Sync {
fn write_record(&self, record: &LogRecord) -> Result<(), SinkError>;
fn flush(&self) -> Result<(), SinkError>;
fn level(&self) -> LogLevel;
}
struct ConsoleState<W: Write> {
writer: W,
scratch: String,
}
pub struct ConsoleSink<F: Formatter, W: Write = BufWriter<Stdout>> {
formatter: F,
level: LogLevel,
state: Mutex<ConsoleState<W>>,
}
impl<F: Formatter> ConsoleSink<F> {
#[expect(
clippy::use_self,
reason = "Self here is ConsoleSink<F> but the return type is \
ConsoleSink<F, BufWriter<Stdout>>; they differ in W"
)]
pub fn new(formatter: F, level: LogLevel) -> ConsoleSink<F, BufWriter<Stdout>> {
ConsoleSink::with_writer(formatter, level, BufWriter::new(io::stdout()))
}
}
impl<F: Formatter, W: Write> ConsoleSink<F, W> {
pub const fn with_writer(formatter: F, level: LogLevel, writer: W) -> Self {
Self {
formatter,
level,
state: Mutex::new(ConsoleState {
writer,
scratch: String::new(),
}),
}
}
}
impl<F: Formatter> ConsoleSink<F, Vec<u8>> {
pub fn captured_output(&self) -> Vec<u8> {
self.state
.lock()
.unwrap_or_else(PoisonError::into_inner)
.writer
.clone()
}
}
impl<F: Formatter, W: Write + Send> Sink for ConsoleSink<F, W> {
#[cfg_attr(feature = "rtsan", rtsan_standalone::blocking)]
#[expect(
clippy::significant_drop_tightening,
reason = "the lock must cover format + write_all so concurrent \
ConsoleSinks don't interleave bytes mid-line"
)]
fn write_record(&self, record: &LogRecord) -> Result<(), SinkError> {
let mut guard = self.state.lock().unwrap_or_else(PoisonError::into_inner);
let ConsoleState { writer, scratch } = &mut *guard;
scratch.clear();
self.formatter.format(record, scratch);
writer.write_all(scratch.as_bytes())?;
writer.write_all(b"\n")?;
Ok(())
}
#[cfg_attr(feature = "rtsan", rtsan_standalone::blocking)]
fn flush(&self) -> Result<(), SinkError> {
let mut guard = self.state.lock().unwrap_or_else(PoisonError::into_inner);
guard.writer.flush().map_err(SinkError::Io)
}
fn level(&self) -> LogLevel {
self.level
}
}
pub struct NullSink {
level: LogLevel,
}
impl NullSink {
#[must_use]
pub const fn new(level: LogLevel) -> Self {
Self { level }
}
}
impl Sink for NullSink {
fn write_record(&self, _record: &LogRecord) -> Result<(), SinkError> {
Ok(())
}
fn flush(&self) -> Result<(), SinkError> {
Ok(())
}
fn level(&self) -> LogLevel {
self.level
}
}
#[cfg(test)]
mod tests {
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};
use super::*;
use crate::decode::{DecodedArg, LogRecord};
use crate::formatter::PatternFormatter;
use crate::metadata::LogMetadata;
static META: LogMetadata = LogMetadata {
level: LogLevel::Info,
fmt_str: "x={}",
file: "f.rs",
line: 1,
module_path: "test",
arg_count: 1,
};
struct CountingSink {
level: LogLevel,
records: AtomicUsize,
flushes: AtomicUsize,
seen_levels: Mutex<Vec<LogLevel>>,
}
impl CountingSink {
fn new(level: LogLevel) -> Self {
Self {
level,
records: AtomicUsize::new(0),
flushes: AtomicUsize::new(0),
seen_levels: Mutex::new(Vec::new()),
}
}
}
impl Sink for CountingSink {
fn write_record(&self, record: &LogRecord) -> Result<(), SinkError> {
self.records.fetch_add(1, Ordering::Relaxed);
self.seen_levels
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(record.metadata.level);
Ok(())
}
fn flush(&self) -> Result<(), SinkError> {
self.flushes.fetch_add(1, Ordering::Relaxed);
Ok(())
}
fn level(&self) -> LogLevel {
self.level
}
}
fn make_record() -> LogRecord {
LogRecord {
timestamp_ns: 0,
logger_name: "test".to_owned(),
metadata: &META,
args: vec![DecodedArg::U32(7)],
}
}
#[test]
fn sink_trait_is_dyn_compatible() {
let arc: std::sync::Arc<dyn Sink> = std::sync::Arc::new(CountingSink::new(LogLevel::Info));
assert_eq!(arc.level(), LogLevel::Info);
}
#[test]
fn sink_trait_bounds_are_send_and_sync() {
const fn assert_send_sync<T: Send + Sync + ?Sized>() {}
assert_send_sync::<dyn Sink>();
}
#[test]
fn console_sink_level_round_trips_each_variant() {
for level in [
LogLevel::Trace,
LogLevel::Debug,
LogLevel::Info,
LogLevel::Warning,
LogLevel::Error,
] {
let sink = ConsoleSink::new(PatternFormatter::default(), level);
assert_eq!(sink.level(), level);
}
}
#[test]
fn console_sink_is_send_and_sync() {
const fn assert_send_sync<T: Send + Sync>() {}
assert_send_sync::<ConsoleSink<PatternFormatter>>();
}
#[test]
fn console_sink_arc_coerces_to_arc_dyn_sink() {
let concrete: std::sync::Arc<ConsoleSink<PatternFormatter>> = std::sync::Arc::new(
ConsoleSink::new(PatternFormatter::default(), LogLevel::Info),
);
let erased: std::sync::Arc<dyn Sink> = concrete;
assert_eq!(erased.level(), LogLevel::Info);
}
fn make_vec_sink() -> ConsoleSink<PatternFormatter, Vec<u8>> {
ConsoleSink::with_writer(PatternFormatter::default(), LogLevel::Info, Vec::new())
}
fn captured(sink: ConsoleSink<PatternFormatter, Vec<u8>>) -> String {
let bytes = sink
.state
.into_inner()
.unwrap_or_else(PoisonError::into_inner)
.writer;
String::from_utf8(bytes).expect("sink output is valid UTF-8")
}
#[test]
fn console_sink_write_record_appends_newline() {
let sink = make_vec_sink();
sink.write_record(&make_record()).unwrap();
let out = captured(sink);
assert!(
out.ends_with('\n'),
"expected trailing newline, got: {out:?}"
);
}
#[test]
fn console_sink_write_record_contains_formatted_arg() {
let sink = make_vec_sink();
sink.write_record(&make_record()).unwrap();
let out = captured(sink);
assert!(
out.contains("x=7"),
"expected 'x=7' in output, got: {out:?}"
);
}
#[test]
fn console_sink_write_record_accumulates_lines() {
let sink = make_vec_sink();
sink.write_record(&make_record()).unwrap();
sink.write_record(&make_record()).unwrap();
let out = captured(sink);
assert_eq!(
out,
"[INFO 0.000] f.rs:1 x=7\n\
[INFO 0.000] f.rs:1 x=7\n",
);
}
#[test]
fn console_sink_flush_succeeds_on_vec_writer() {
let sink = make_vec_sink();
sink.write_record(&make_record()).unwrap();
sink.flush().unwrap();
}
#[test]
fn sink_trait_dispatch_drives_implementation() {
let sink = CountingSink::new(LogLevel::Warning);
let dynamic: &dyn Sink = &sink;
assert_eq!(dynamic.level(), LogLevel::Warning);
let record = make_record();
dynamic.write_record(&record).unwrap();
dynamic.write_record(&record).unwrap();
dynamic.flush().unwrap();
assert_eq!(sink.records.load(Ordering::Relaxed), 2);
assert_eq!(sink.flushes.load(Ordering::Relaxed), 1);
assert_eq!(
sink.seen_levels
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.as_slice(),
&[LogLevel::Info, LogLevel::Info],
);
}
#[test]
fn log_record_logger_name_is_accessible_from_sink() {
let record = LogRecord {
timestamp_ns: 0,
logger_name: "payments".to_owned(),
metadata: &META,
args: vec![],
};
assert_eq!(record.logger_name, "payments");
}
}