use std::{
error::Error,
fmt,
future::Future,
io::{self, Write},
pin::Pin,
sync::{
Arc, Mutex, OnceLock,
atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering},
mpsc::{self, Receiver, SyncSender, TrySendError},
},
task::{Context, Poll, Waker},
thread::{self, JoinHandle},
time::{SystemTime, UNIX_EPOCH},
};
use saddle_core::{ComponentLifecycle, ErrorKind, LifecycleFuture, SaddleError, SpanId, TraceId};
use serde::Serialize;
use crate::event::EventLevel;
static GLOBAL: OnceLock<Observer> = OnceLock::new();
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct ObserverConfig {
pub queue_capacity: usize,
}
impl Default for ObserverConfig {
fn default() -> Self {
Self {
queue_capacity: 8_192,
}
}
}
#[derive(Debug)]
pub enum InitError {
EmptyQueue,
RandomSource,
Spawn(io::Error),
}
impl fmt::Display for InitError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::EmptyQueue => formatter.write_str("log queue capacity must be greater than zero"),
Self::RandomSource => {
formatter.write_str("operating system random source is unavailable")
}
Self::Spawn(error) => write!(formatter, "failed to start log writer: {error}"),
}
}
}
impl Error for InitError {
fn source(&self) -> Option<&(dyn Error + 'static)> {
match self {
Self::Spawn(error) => Some(error),
Self::EmptyQueue | Self::RandomSource => None,
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum OutputStage {
Serialize,
Record,
Newline,
Flush,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum FlushError {
OutputFailed(OutputStage),
DroppedEvents(u64),
OutputFailedAndDropped {
stage: OutputStage,
dropped_events: u64,
},
WriterStopped,
AlreadyShuttingDown,
CoordinatorUnavailable,
WorkerPanicked,
}
impl fmt::Display for FlushError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::OutputFailed(stage) => write!(formatter, "log output failed during {stage:?}"),
Self::DroppedEvents(count) => write!(formatter, "{count} log events were dropped"),
Self::OutputFailedAndDropped {
stage,
dropped_events,
} => write!(
formatter,
"log output failed during {stage:?} and {dropped_events} events were dropped"
),
Self::WriterStopped => formatter.write_str("log writer is not running"),
Self::AlreadyShuttingDown => formatter.write_str("log writer shutdown already started"),
Self::CoordinatorUnavailable => {
formatter.write_str("failed to start log lifecycle coordinator")
}
Self::WorkerPanicked => formatter.write_str("log writer thread panicked"),
}
}
}
impl Error for FlushError {}
pub struct FlushFuture {
completion: Arc<Completion>,
}
impl Future for FlushFuture {
type Output = Result<(), FlushError>;
fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
if let Some(result) = self.completion.result.lock().unwrap().take() {
return Poll::Ready(result);
}
*self.completion.waker.lock().unwrap() = Some(context.waker().clone());
if let Some(result) = self.completion.result.lock().unwrap().take() {
Poll::Ready(result)
} else {
Poll::Pending
}
}
}
struct Completion {
result: Mutex<Option<Result<(), FlushError>>>,
waker: Mutex<Option<Waker>>,
}
impl Completion {
fn pending() -> Arc<Self> {
Arc::new(Self {
result: Mutex::new(None),
waker: Mutex::new(None),
})
}
fn ready(result: Result<(), FlushError>) -> Arc<Self> {
Arc::new(Self {
result: Mutex::new(Some(result)),
waker: Mutex::new(None),
})
}
fn complete(&self, result: Result<(), FlushError>) {
*self.result.lock().unwrap() = Some(result);
if let Some(waker) = self.waker.lock().unwrap().take() {
waker.wake();
}
}
}
#[derive(Clone)]
pub struct Observer {
pub(crate) inner: Arc<Inner>,
}
pub(crate) struct Inner {
sender: SyncSender<Command>,
dropped: AtomicU64,
accepting: AtomicBool,
emitting: AtomicUsize,
shutdown_started: AtomicBool,
worker: Mutex<Option<JoinHandle<()>>>,
ids: IdGenerator,
}
struct IdGenerator {
trace_high: u64,
trace_seed: u64,
trace_counter: AtomicU64,
span_seed: u64,
span_counter: AtomicU64,
}
impl IdGenerator {
fn from_seed(seed: [u8; 24]) -> Self {
let mut high = u64::from_be_bytes(seed[0..8].try_into().unwrap());
if high == 0 {
high = 1;
}
Self {
trace_high: high,
trace_seed: u64::from_be_bytes(seed[8..16].try_into().unwrap()),
trace_counter: AtomicU64::new(0),
span_seed: u64::from_be_bytes(seed[16..24].try_into().unwrap()),
span_counter: AtomicU64::new(0),
}
}
fn trace_id(&self) -> TraceId {
let low = self
.trace_seed
.wrapping_add(self.trace_counter.fetch_add(1, Ordering::Relaxed));
TraceId::from_u128((u128::from(self.trace_high) << 64) | u128::from(low))
}
fn span_id(&self) -> SpanId {
loop {
let value = self
.span_seed
.wrapping_add(self.span_counter.fetch_add(1, Ordering::Relaxed));
if value != 0 {
return SpanId::from_u64(value);
}
}
}
}
enum Command {
Record(LogRecord),
Flush {
unreported_dropped: u64,
response: mpsc::Sender<WorkerStatus>,
},
Shutdown {
unreported_dropped: u64,
response: mpsc::Sender<WorkerStatus>,
},
}
#[derive(Clone, Copy, Debug, Default)]
struct WorkerStatus {
failure: Option<OutputStage>,
unreported_dropped: u64,
}
impl WorkerStatus {
fn into_result(self) -> Result<(), FlushError> {
match (self.failure, self.unreported_dropped) {
(None, 0) => Ok(()),
(Some(stage), 0) => Err(FlushError::OutputFailed(stage)),
(None, dropped_events) => Err(FlushError::DroppedEvents(dropped_events)),
(Some(stage), dropped_events) => Err(FlushError::OutputFailedAndDropped {
stage,
dropped_events,
}),
}
}
}
#[derive(Serialize)]
pub(crate) struct LogRecord {
pub timestamp_unix_ms: u128,
pub level: EventLevel,
pub event: &'static str,
#[serde(skip_serializing_if = "Option::is_none")]
pub trace_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub span_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub parent_span_id: Option<String>,
#[serde(flatten)]
pub data: serde_json::Map<String, serde_json::Value>,
#[serde(skip_serializing_if = "is_zero")]
pub dropped_events: u64,
}
const fn is_zero(value: &u64) -> bool {
*value == 0
}
impl LogRecord {
pub(crate) fn new(level: EventLevel, event: &'static str) -> Self {
Self {
timestamp_unix_ms: SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis(),
level,
event,
trace_id: None,
span_id: None,
parent_span_id: None,
data: serde_json::Map::new(),
dropped_events: 0,
}
}
}
pub fn init(config: ObserverConfig) -> Result<&'static Observer, InitError> {
if let Some(observer) = GLOBAL.get() {
return Ok(observer);
}
let observer = Observer::with_writer(config, io::stdout())?;
let _ = GLOBAL.set(observer);
Ok(GLOBAL.get().expect("global observer was initialized"))
}
pub fn global() -> Option<&'static Observer> {
GLOBAL.get()
}
impl Observer {
pub fn with_writer(
config: ObserverConfig,
writer: impl Write + Send + 'static,
) -> Result<Self, InitError> {
let mut seed = [0_u8; 24];
getrandom::fill(&mut seed).map_err(|_| InitError::RandomSource)?;
Self::with_writer_and_seed(config, writer, Ok(seed))
}
fn with_writer_and_seed(
config: ObserverConfig,
writer: impl Write + Send + 'static,
seed: Result<[u8; 24], InitError>,
) -> Result<Self, InitError> {
if config.queue_capacity == 0 {
return Err(InitError::EmptyQueue);
}
let ids = IdGenerator::from_seed(seed?);
let (sender, receiver) = mpsc::sync_channel(config.queue_capacity);
let worker = thread::Builder::new()
.name("saddle-log-writer".to_owned())
.spawn(move || write_records(receiver, writer))
.map_err(InitError::Spawn)?;
Ok(Self {
inner: Arc::new(Inner {
sender,
dropped: AtomicU64::new(0),
accepting: AtomicBool::new(true),
emitting: AtomicUsize::new(0),
shutdown_started: AtomicBool::new(false),
worker: Mutex::new(Some(worker)),
ids,
}),
})
}
pub(crate) fn new_trace_id(&self) -> TraceId {
self.inner.ids.trace_id()
}
pub(crate) fn new_span_id(&self) -> SpanId {
self.inner.ids.span_id()
}
pub(crate) fn emit(&self, mut record: LogRecord) {
if !self.inner.accepting.load(Ordering::Acquire) {
return;
}
self.inner.emitting.fetch_add(1, Ordering::AcqRel);
if !self.inner.accepting.load(Ordering::Acquire) {
self.inner.emitting.fetch_sub(1, Ordering::AcqRel);
return;
}
record.dropped_events = self.inner.dropped.swap(0, Ordering::Relaxed);
match self.inner.sender.try_send(Command::Record(record)) {
Ok(()) => {}
Err(TrySendError::Full(Command::Record(record))) => {
self.inner
.dropped
.fetch_add(record.dropped_events.saturating_add(1), Ordering::Relaxed);
}
Err(TrySendError::Disconnected(_)) => {
self.inner.dropped.fetch_add(1, Ordering::Relaxed);
}
Err(TrySendError::Full(_)) => unreachable!("emit only sends records"),
}
self.inner.emitting.fetch_sub(1, Ordering::Release);
}
pub fn flush(&self) -> FlushFuture {
if !self.inner.accepting.load(Ordering::Acquire) {
return FlushFuture {
completion: Completion::ready(Err(FlushError::WriterStopped)),
};
}
let completion = Completion::pending();
let future = FlushFuture {
completion: completion.clone(),
};
let inner = self.inner.clone();
if thread::Builder::new()
.name("saddle-log-flush".to_owned())
.spawn(move || coordinate_flush(inner, completion.clone()))
.is_err()
{
future
.completion
.complete(Err(FlushError::CoordinatorUnavailable));
}
future
}
pub fn shutdown(&self) -> FlushFuture {
if self.inner.shutdown_started.swap(true, Ordering::AcqRel) {
return FlushFuture {
completion: Completion::ready(Err(FlushError::AlreadyShuttingDown)),
};
}
self.inner.accepting.store(false, Ordering::Release);
let completion = Completion::pending();
let future = FlushFuture {
completion: completion.clone(),
};
let inner = self.inner.clone();
if thread::Builder::new()
.name("saddle-log-shutdown".to_owned())
.spawn(move || coordinate_shutdown(inner, completion.clone()))
.is_err()
{
self.inner.accepting.store(true, Ordering::Release);
self.inner.shutdown_started.store(false, Ordering::Release);
future
.completion
.complete(Err(FlushError::CoordinatorUnavailable));
}
future
}
pub fn dropped_events(&self) -> u64 {
self.inner.dropped.load(Ordering::Relaxed)
}
}
impl ComponentLifecycle for Observer {
fn name(&self) -> &'static str {
"observability"
}
fn start(&self) -> LifecycleFuture<'_> {
Box::pin(async { Ok(()) })
}
fn shutdown(&self) -> LifecycleFuture<'_> {
let shutdown = Observer::shutdown(self);
Box::pin(async move {
shutdown.await.map_err(|_| {
SaddleError::new(
ErrorKind::Infrastructure,
"observability.shutdown_failed",
"structured log shutdown failed",
)
})
})
}
}
fn coordinate_flush(inner: Arc<Inner>, completion: Arc<Completion>) {
let dropped = inner.dropped.swap(0, Ordering::AcqRel);
let (response, receiver) = mpsc::channel();
let result = if inner
.sender
.send(Command::Flush {
unreported_dropped: dropped,
response,
})
.is_err()
{
Err(FlushError::WriterStopped)
} else {
receiver
.recv()
.map_err(|_| FlushError::WriterStopped)
.and_then(WorkerStatus::into_result)
};
completion.complete(result);
}
fn coordinate_shutdown(inner: Arc<Inner>, completion: Arc<Completion>) {
while inner.emitting.load(Ordering::Acquire) != 0 {
thread::yield_now();
}
let dropped = inner.dropped.swap(0, Ordering::AcqRel);
let (response, receiver) = mpsc::channel();
let mut result = if inner
.sender
.send(Command::Shutdown {
unreported_dropped: dropped,
response,
})
.is_err()
{
Err(FlushError::WriterStopped)
} else {
receiver
.recv()
.map_err(|_| FlushError::WriterStopped)
.and_then(WorkerStatus::into_result)
};
if let Some(worker) = inner.worker.lock().unwrap().take() {
if worker.join().is_err() {
result = Err(FlushError::WorkerPanicked);
}
}
completion.complete(result);
}
fn write_records(receiver: Receiver<Command>, mut writer: impl Write) {
let mut status = WorkerStatus::default();
while let Ok(command) = receiver.recv() {
match command {
Command::Record(record) => write_record(&mut writer, record, &mut status),
Command::Flush {
unreported_dropped,
response,
} => {
status.unreported_dropped =
status.unreported_dropped.saturating_add(unreported_dropped);
flush_writer(&mut writer, &mut status);
let _ = response.send(status);
}
Command::Shutdown {
unreported_dropped,
response,
} => {
status.unreported_dropped =
status.unreported_dropped.saturating_add(unreported_dropped);
flush_writer(&mut writer, &mut status);
let _ = response.send(status);
break;
}
}
}
}
fn write_record(writer: &mut impl Write, record: LogRecord, status: &mut WorkerStatus) {
if status.failure.is_some() {
status.unreported_dropped = status
.unreported_dropped
.saturating_add(record.dropped_events);
return;
}
let dropped = record.dropped_events;
let bytes = match serde_json::to_vec(&record) {
Ok(bytes) => bytes,
Err(_) => {
status.failure = Some(OutputStage::Serialize);
status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
return;
}
};
if writer.write_all(&bytes).is_err() {
status.failure = Some(OutputStage::Record);
status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
return;
}
if writer.write_all(b"\n").is_err() {
status.failure = Some(OutputStage::Newline);
status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
}
}
fn flush_writer(writer: &mut impl Write, status: &mut WorkerStatus) {
if writer.flush().is_err() && status.failure.is_none() {
status.failure = Some(OutputStage::Flush);
}
}
#[cfg(test)]
mod tests {
use std::{
sync::{Condvar, mpsc},
task::{Wake, Waker},
};
use super::*;
struct ThreadWaker(thread::Thread);
impl Wake for ThreadWaker {
fn wake(self: Arc<Self>) {
self.0.unpark();
}
}
fn block_on<T>(future: impl Future<Output = T>) -> T {
let waker = Waker::from(Arc::new(ThreadWaker(thread::current())));
let mut context = Context::from_waker(&waker);
let mut future = std::pin::pin!(future);
loop {
match future.as_mut().poll(&mut context) {
Poll::Ready(output) => return output,
Poll::Pending => thread::park(),
}
}
}
struct BlockingWriter {
entered: Option<mpsc::Sender<()>>,
release: Arc<(Mutex<bool>, Condvar)>,
fail_after_release: bool,
}
impl Write for BlockingWriter {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
if let Some(entered) = self.entered.take() {
let _ = entered.send(());
let (lock, condition) = &*self.release;
let mut released = lock.lock().unwrap();
while !*released {
released = condition.wait(released).unwrap();
}
if self.fail_after_release {
return Err(io::Error::other("injected writer failure"));
}
}
Ok(bytes.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
struct FailingWriter {
fail_write: Option<usize>,
writes: usize,
fail_flush: bool,
}
impl Write for FailingWriter {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
self.writes += 1;
if self.fail_write == Some(self.writes) {
Err(io::Error::other("injected writer failure"))
} else {
Ok(bytes.len())
}
}
fn flush(&mut self) -> io::Result<()> {
if self.fail_flush {
Err(io::Error::other("injected flush failure"))
} else {
Ok(())
}
}
}
fn observer(writer: impl Write + Send + 'static, capacity: usize) -> Observer {
Observer::with_writer_and_seed(
ObserverConfig {
queue_capacity: capacity,
},
writer,
Ok([7; 24]),
)
.unwrap()
}
#[test]
fn random_source_failure_is_reported_during_initialization() {
let result = Observer::with_writer_and_seed(
ObserverConfig::default(),
io::sink(),
Err(InitError::RandomSource),
);
assert!(matches!(result, Err(InitError::RandomSource)));
}
#[test]
fn initialized_generator_produces_non_zero_unique_ids_without_more_io() {
let observer = observer(io::sink(), 8);
let first_trace = observer.new_trace_id();
let second_trace = observer.new_trace_id();
let first_span = observer.new_span_id();
let second_span = observer.new_span_id();
assert_ne!(first_trace.as_u128(), 0);
assert_ne!(first_trace, second_trace);
assert_ne!(first_span.as_u64(), 0);
assert_ne!(first_span, second_span);
block_on(observer.shutdown()).unwrap();
}
#[test]
fn observer_is_a_managed_component() {
let observer = observer(io::sink(), 8);
assert_eq!(ComponentLifecycle::name(&observer), "observability");
block_on(ComponentLifecycle::start(&observer)).unwrap();
block_on(ComponentLifecycle::shutdown(&observer)).unwrap();
}
#[test]
fn full_queue_is_reported_by_shutdown_without_waiting_in_emit() {
let (entered_sender, entered_receiver) = mpsc::channel();
let release = Arc::new((Mutex::new(false), Condvar::new()));
let writer = BlockingWriter {
entered: Some(entered_sender),
release: release.clone(),
fail_after_release: false,
};
let observer = observer(writer, 1);
observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
entered_receiver.recv().unwrap();
observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
assert_eq!(observer.dropped_events(), 1);
let shutdown = observer.shutdown();
let (lock, condition) = &*release;
*lock.lock().unwrap() = true;
condition.notify_one();
assert_eq!(block_on(shutdown), Err(FlushError::DroppedEvents(1)));
}
#[test]
fn shutdown_reports_writer_failure_and_unreported_drops_together() {
let (entered_sender, entered_receiver) = mpsc::channel();
let release = Arc::new((Mutex::new(false), Condvar::new()));
let writer = BlockingWriter {
entered: Some(entered_sender),
release: release.clone(),
fail_after_release: true,
};
let observer = observer(writer, 1);
observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
entered_receiver.recv().unwrap();
observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
let shutdown = observer.shutdown();
let (lock, condition) = &*release;
*lock.lock().unwrap() = true;
condition.notify_one();
assert_eq!(
block_on(shutdown),
Err(FlushError::OutputFailedAndDropped {
stage: OutputStage::Record,
dropped_events: 1,
})
);
}
#[test]
fn record_write_failure_persists_until_shutdown() {
let observer = observer(
FailingWriter {
fail_write: Some(1),
writes: 0,
fail_flush: false,
},
8,
);
observer.emit(LogRecord::new(EventLevel::Info, "test.record"));
assert_eq!(
block_on(observer.shutdown()),
Err(FlushError::OutputFailed(OutputStage::Record))
);
}
#[test]
fn newline_failure_persists_until_flush() {
let observer = observer(
FailingWriter {
fail_write: Some(2),
writes: 0,
fail_flush: false,
},
8,
);
observer.emit(LogRecord::new(EventLevel::Info, "test.newline"));
assert_eq!(
block_on(observer.flush()),
Err(FlushError::OutputFailed(OutputStage::Newline))
);
let _ = block_on(observer.shutdown());
}
#[test]
fn flush_failure_is_reported() {
let observer = observer(
FailingWriter {
fail_write: None,
writes: 0,
fail_flush: true,
},
8,
);
assert_eq!(
block_on(observer.flush()),
Err(FlushError::OutputFailed(OutputStage::Flush))
);
assert_eq!(
block_on(observer.shutdown()),
Err(FlushError::OutputFailed(OutputStage::Flush))
);
}
}