Skip to main content

saddle_observability/
logger.rs

1use std::{
2    error::Error,
3    fmt,
4    future::Future,
5    io::{self, Write},
6    pin::Pin,
7    sync::{
8        Arc, Mutex, OnceLock,
9        atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering},
10        mpsc::{self, Receiver, SyncSender, TrySendError},
11    },
12    task::{Context, Poll, Waker},
13    thread::{self, JoinHandle},
14    time::{SystemTime, UNIX_EPOCH},
15};
16
17use saddle_core::{ComponentLifecycle, ErrorKind, LifecycleFuture, SaddleError, SpanId, TraceId};
18use serde::Serialize;
19
20use crate::{
21    calendar_writer::{CalendarFileWriter, FileLoggingConfig},
22    event::EventLevel,
23};
24
25static GLOBAL: OnceLock<Observer> = OnceLock::new();
26
27/// Configuration for the bounded non-blocking logger.
28#[derive(Clone, Copy, Debug, Eq, PartialEq)]
29pub struct ObserverConfig {
30    /// Maximum records waiting for the output worker. New records are dropped
31    /// when the queue is full rather than blocking an async request.
32    pub queue_capacity: usize,
33}
34
35impl Default for ObserverConfig {
36    fn default() -> Self {
37        Self {
38            queue_capacity: 8_192,
39        }
40    }
41}
42
43#[derive(Debug)]
44pub enum InitError {
45    EmptyQueue,
46    RandomSource,
47    Output(io::Error),
48    Spawn(io::Error),
49}
50
51impl fmt::Display for InitError {
52    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
53        match self {
54            Self::EmptyQueue => formatter.write_str("log queue capacity must be greater than zero"),
55            Self::RandomSource => {
56                formatter.write_str("operating system random source is unavailable")
57            }
58            Self::Output(error) => write!(formatter, "failed to initialize log output: {error}"),
59            Self::Spawn(error) => write!(formatter, "failed to start log writer: {error}"),
60        }
61    }
62}
63
64impl Error for InitError {
65    fn source(&self) -> Option<&(dyn Error + 'static)> {
66        match self {
67            Self::Spawn(error) | Self::Output(error) => Some(error),
68            Self::EmptyQueue | Self::RandomSource => None,
69        }
70    }
71}
72
73/// The output operation whose first failure made log delivery unhealthy.
74#[derive(Clone, Copy, Debug, Eq, PartialEq)]
75pub enum OutputStage {
76    Serialize,
77    Record,
78    Newline,
79    Flush,
80}
81
82/// A flush or managed shutdown result that is safe to expose to framework code.
83#[derive(Clone, Debug, Eq, PartialEq)]
84pub enum FlushError {
85    OutputFailed(OutputStage),
86    DroppedEvents(u64),
87    OutputFailedAndDropped {
88        stage: OutputStage,
89        dropped_events: u64,
90    },
91    WriterStopped,
92    AlreadyShuttingDown,
93    CoordinatorUnavailable,
94    WorkerPanicked,
95}
96
97impl fmt::Display for FlushError {
98    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
99        match self {
100            Self::OutputFailed(stage) => write!(formatter, "log output failed during {stage:?}"),
101            Self::DroppedEvents(count) => write!(formatter, "{count} log events were dropped"),
102            Self::OutputFailedAndDropped {
103                stage,
104                dropped_events,
105            } => write!(
106                formatter,
107                "log output failed during {stage:?} and {dropped_events} events were dropped"
108            ),
109            Self::WriterStopped => formatter.write_str("log writer is not running"),
110            Self::AlreadyShuttingDown => formatter.write_str("log writer shutdown already started"),
111            Self::CoordinatorUnavailable => {
112                formatter.write_str("failed to start log lifecycle coordinator")
113            }
114            Self::WorkerPanicked => formatter.write_str("log writer thread panicked"),
115        }
116    }
117}
118
119impl Error for FlushError {}
120
121/// A non-blocking future returned by flush and shutdown lifecycle operations.
122pub struct FlushFuture {
123    completion: Arc<Completion>,
124}
125
126impl Future for FlushFuture {
127    type Output = Result<(), FlushError>;
128
129    fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
130        if let Some(result) = self.completion.result.lock().unwrap().take() {
131            return Poll::Ready(result);
132        }
133        *self.completion.waker.lock().unwrap() = Some(context.waker().clone());
134        if let Some(result) = self.completion.result.lock().unwrap().take() {
135            Poll::Ready(result)
136        } else {
137            Poll::Pending
138        }
139    }
140}
141
142struct Completion {
143    result: Mutex<Option<Result<(), FlushError>>>,
144    waker: Mutex<Option<Waker>>,
145}
146
147impl Completion {
148    fn pending() -> Arc<Self> {
149        Arc::new(Self {
150            result: Mutex::new(None),
151            waker: Mutex::new(None),
152        })
153    }
154
155    fn ready(result: Result<(), FlushError>) -> Arc<Self> {
156        Arc::new(Self {
157            result: Mutex::new(Some(result)),
158            waker: Mutex::new(None),
159        })
160    }
161
162    fn complete(&self, result: Result<(), FlushError>) {
163        *self.result.lock().unwrap() = Some(result);
164        if let Some(waker) = self.waker.lock().unwrap().take() {
165            waker.wake();
166        }
167    }
168}
169
170/// Process-level structured logger and trace correlator.
171///
172/// Cloning this value is cheap. Emission uses `try_send`; output and JSON
173/// serialization happen on the dedicated writer thread.
174#[derive(Clone)]
175pub struct Observer {
176    pub(crate) inner: Arc<Inner>,
177}
178
179pub(crate) struct Inner {
180    sender: SyncSender<Command>,
181    dropped: AtomicU64,
182    accepting: AtomicBool,
183    emitting: AtomicUsize,
184    shutdown_started: AtomicBool,
185    worker: Mutex<Option<JoinHandle<()>>>,
186    ids: IdGenerator,
187    pub(crate) metrics: crate::metrics::Metrics,
188}
189
190struct IdGenerator {
191    trace_high: u64,
192    trace_seed: u64,
193    trace_counter: AtomicU64,
194    span_seed: u64,
195    span_counter: AtomicU64,
196}
197
198impl IdGenerator {
199    fn from_seed(seed: [u8; 24]) -> Self {
200        let mut high = u64::from_be_bytes(seed[0..8].try_into().unwrap());
201        if high == 0 {
202            high = 1;
203        }
204        Self {
205            trace_high: high,
206            trace_seed: u64::from_be_bytes(seed[8..16].try_into().unwrap()),
207            trace_counter: AtomicU64::new(0),
208            span_seed: u64::from_be_bytes(seed[16..24].try_into().unwrap()),
209            span_counter: AtomicU64::new(0),
210        }
211    }
212
213    fn trace_id(&self) -> TraceId {
214        let low = self
215            .trace_seed
216            .wrapping_add(self.trace_counter.fetch_add(1, Ordering::Relaxed));
217        TraceId::from_u128((u128::from(self.trace_high) << 64) | u128::from(low))
218    }
219
220    fn span_id(&self) -> SpanId {
221        loop {
222            let value = self
223                .span_seed
224                .wrapping_add(self.span_counter.fetch_add(1, Ordering::Relaxed));
225            if value != 0 {
226                return SpanId::from_u64(value);
227            }
228        }
229    }
230}
231
232enum Command {
233    Record(LogRecord),
234    Flush {
235        unreported_dropped: u64,
236        response: mpsc::Sender<WorkerStatus>,
237    },
238    Shutdown {
239        unreported_dropped: u64,
240        response: mpsc::Sender<WorkerStatus>,
241    },
242}
243
244#[derive(Clone, Copy, Debug, Default)]
245struct WorkerStatus {
246    failure: Option<OutputStage>,
247    unreported_dropped: u64,
248}
249
250impl WorkerStatus {
251    fn into_result(self) -> Result<(), FlushError> {
252        match (self.failure, self.unreported_dropped) {
253            (None, 0) => Ok(()),
254            (Some(stage), 0) => Err(FlushError::OutputFailed(stage)),
255            (None, dropped_events) => Err(FlushError::DroppedEvents(dropped_events)),
256            (Some(stage), dropped_events) => Err(FlushError::OutputFailedAndDropped {
257                stage,
258                dropped_events,
259            }),
260        }
261    }
262}
263
264#[derive(Serialize)]
265pub(crate) struct LogRecord {
266    pub timestamp_unix_ms: u128,
267    pub level: EventLevel,
268    pub event: &'static str,
269    #[serde(skip_serializing_if = "Option::is_none")]
270    pub trace_id: Option<String>,
271    #[serde(skip_serializing_if = "Option::is_none")]
272    pub span: Option<String>,
273    #[serde(skip_serializing_if = "Option::is_none")]
274    pub span_id: Option<String>,
275    #[serde(skip_serializing_if = "Option::is_none")]
276    pub parent: Option<String>,
277    #[serde(skip_serializing_if = "Option::is_none")]
278    pub parent_span_id: Option<String>,
279    #[serde(flatten)]
280    pub data: serde_json::Map<String, serde_json::Value>,
281    #[serde(skip_serializing_if = "is_zero")]
282    pub dropped_events: u64,
283}
284
285const fn is_zero(value: &u64) -> bool {
286    *value == 0
287}
288
289impl LogRecord {
290    pub(crate) fn new(level: EventLevel, event: &'static str) -> Self {
291        Self {
292            timestamp_unix_ms: SystemTime::now()
293                .duration_since(UNIX_EPOCH)
294                .unwrap_or_default()
295                .as_millis(),
296            level,
297            event,
298            trace_id: None,
299            span: None,
300            span_id: None,
301            parent: None,
302            parent_span_id: None,
303            data: serde_json::Map::new(),
304            dropped_events: 0,
305        }
306    }
307}
308
309/// Initializes stdout logging once. Later calls return the first observer.
310pub fn init(config: ObserverConfig) -> Result<&'static Observer, InitError> {
311    if let Some(observer) = GLOBAL.get() {
312        return Ok(observer);
313    }
314
315    let observer = Observer::with_writer(config, io::stdout())?;
316    let _ = GLOBAL.set(observer);
317    Ok(GLOBAL.get().expect("global observer was initialized"))
318}
319
320/// Initializes the process logger with the sole calendar-rotated file sink.
321pub fn init_file(
322    config: ObserverConfig,
323    file: FileLoggingConfig,
324) -> Result<&'static Observer, InitError> {
325    if let Some(observer) = GLOBAL.get() {
326        return Ok(observer);
327    }
328    let writer = CalendarFileWriter::open(file).map_err(InitError::Output)?;
329    let observer = Observer::with_writer(config, writer)?;
330    let _ = GLOBAL.set(observer);
331    Ok(GLOBAL.get().expect("global observer was initialized"))
332}
333
334/// Returns the initialized process observer, if application startup installed it.
335pub fn global() -> Option<&'static Observer> {
336    GLOBAL.get()
337}
338
339impl Observer {
340    /// Creates an observer with a framework-owned writer.
341    ///
342    /// Application startup should normally use [`init`]. This constructor is
343    /// useful for embedding and deterministic tests.
344    pub fn with_writer(
345        config: ObserverConfig,
346        writer: impl Write + Send + 'static,
347    ) -> Result<Self, InitError> {
348        let mut seed = [0_u8; 24];
349        getrandom::fill(&mut seed).map_err(|_| InitError::RandomSource)?;
350        Self::with_writer_and_seed(config, writer, Ok(seed))
351    }
352
353    fn with_writer_and_seed(
354        config: ObserverConfig,
355        writer: impl Write + Send + 'static,
356        seed: Result<[u8; 24], InitError>,
357    ) -> Result<Self, InitError> {
358        if config.queue_capacity == 0 {
359            return Err(InitError::EmptyQueue);
360        }
361        let ids = IdGenerator::from_seed(seed?);
362        let (sender, receiver) = mpsc::sync_channel(config.queue_capacity);
363        let worker = thread::Builder::new()
364            .name("saddle-log-writer".to_owned())
365            .spawn(move || write_records(receiver, writer))
366            .map_err(InitError::Spawn)?;
367        Ok(Self {
368            inner: Arc::new(Inner {
369                sender,
370                dropped: AtomicU64::new(0),
371                accepting: AtomicBool::new(true),
372                emitting: AtomicUsize::new(0),
373                shutdown_started: AtomicBool::new(false),
374                worker: Mutex::new(Some(worker)),
375                ids,
376                metrics: crate::metrics::Metrics::default(),
377            }),
378        })
379    }
380
381    pub(crate) fn new_trace_id(&self) -> TraceId {
382        self.inner.ids.trace_id()
383    }
384
385    pub(crate) fn new_span_id(&self) -> SpanId {
386        self.inner.ids.span_id()
387    }
388
389    pub(crate) fn emit(&self, mut record: LogRecord) {
390        if !self.inner.accepting.load(Ordering::Acquire) {
391            return;
392        }
393        self.inner.emitting.fetch_add(1, Ordering::AcqRel);
394        if !self.inner.accepting.load(Ordering::Acquire) {
395            self.inner.emitting.fetch_sub(1, Ordering::AcqRel);
396            return;
397        }
398
399        record.dropped_events = self.inner.dropped.swap(0, Ordering::Relaxed);
400        match self.inner.sender.try_send(Command::Record(record)) {
401            Ok(()) => {}
402            Err(TrySendError::Full(Command::Record(record))) => {
403                self.inner.metrics.logger_dropped(1);
404                self.inner
405                    .dropped
406                    .fetch_add(record.dropped_events.saturating_add(1), Ordering::Relaxed);
407            }
408            Err(TrySendError::Disconnected(_)) => {
409                self.inner.metrics.logger_dropped(1);
410                self.inner.metrics.logger_output(true);
411                self.inner.dropped.fetch_add(1, Ordering::Relaxed);
412            }
413            Err(TrySendError::Full(_)) => unreachable!("emit only sends records"),
414        }
415        self.inner.emitting.fetch_sub(1, Ordering::Release);
416    }
417
418    /// Returns a fixed-size, allocation-free copy of the current process metrics.
419    pub fn metrics_snapshot(&self) -> crate::MetricsSnapshot {
420        self.inner.metrics.snapshot()
421    }
422
423    /// Flushes records accepted before the lifecycle coordinator reaches the
424    /// writer. The returned future never blocks the async executor thread.
425    pub fn flush(&self) -> FlushFuture {
426        if !self.inner.accepting.load(Ordering::Acquire) {
427            return FlushFuture {
428                completion: Completion::ready(Err(FlushError::WriterStopped)),
429            };
430        }
431        let completion = Completion::pending();
432        let future = FlushFuture {
433            completion: completion.clone(),
434        };
435        let inner = self.inner.clone();
436        if thread::Builder::new()
437            .name("saddle-log-flush".to_owned())
438            .spawn(move || coordinate_flush(inner, completion.clone()))
439            .is_err()
440        {
441            future
442                .completion
443                .complete(Err(FlushError::CoordinatorUnavailable));
444        }
445        future
446    }
447
448    /// Stops admission, drains accepted records, flushes output and joins the
449    /// writer without blocking the async executor thread.
450    pub fn shutdown(&self) -> FlushFuture {
451        if self.inner.shutdown_started.swap(true, Ordering::AcqRel) {
452            return FlushFuture {
453                completion: Completion::ready(Err(FlushError::AlreadyShuttingDown)),
454            };
455        }
456        self.inner.accepting.store(false, Ordering::Release);
457        let completion = Completion::pending();
458        let future = FlushFuture {
459            completion: completion.clone(),
460        };
461        let inner = self.inner.clone();
462        if thread::Builder::new()
463            .name("saddle-log-shutdown".to_owned())
464            .spawn(move || coordinate_shutdown(inner, completion.clone()))
465            .is_err()
466        {
467            self.inner.accepting.store(true, Ordering::Release);
468            self.inner.shutdown_started.store(false, Ordering::Release);
469            future
470                .completion
471                .complete(Err(FlushError::CoordinatorUnavailable));
472        }
473        future
474    }
475
476    pub fn dropped_events(&self) -> u64 {
477        self.inner.dropped.load(Ordering::Relaxed)
478    }
479}
480
481impl ComponentLifecycle for Observer {
482    fn name(&self) -> &'static str {
483        "observability"
484    }
485
486    fn start(&self) -> LifecycleFuture<'_> {
487        Box::pin(async { Ok(()) })
488    }
489
490    fn shutdown(&self) -> LifecycleFuture<'_> {
491        let shutdown = Observer::shutdown(self);
492        Box::pin(async move {
493            shutdown.await.map_err(|_| {
494                SaddleError::new(
495                    ErrorKind::Infrastructure,
496                    "observability.shutdown_failed",
497                    "structured log shutdown failed",
498                )
499            })
500        })
501    }
502}
503
504fn coordinate_flush(inner: Arc<Inner>, completion: Arc<Completion>) {
505    let dropped = inner.dropped.swap(0, Ordering::AcqRel);
506    let (response, receiver) = mpsc::channel();
507    let result = if inner
508        .sender
509        .send(Command::Flush {
510            unreported_dropped: dropped,
511            response,
512        })
513        .is_err()
514    {
515        Err(FlushError::WriterStopped)
516    } else {
517        receiver
518            .recv()
519            .map_err(|_| FlushError::WriterStopped)
520            .and_then(WorkerStatus::into_result)
521    };
522    completion.complete(result);
523}
524
525fn coordinate_shutdown(inner: Arc<Inner>, completion: Arc<Completion>) {
526    while inner.emitting.load(Ordering::Acquire) != 0 {
527        thread::yield_now();
528    }
529    let dropped = inner.dropped.swap(0, Ordering::AcqRel);
530    let (response, receiver) = mpsc::channel();
531    let mut result = if inner
532        .sender
533        .send(Command::Shutdown {
534            unreported_dropped: dropped,
535            response,
536        })
537        .is_err()
538    {
539        Err(FlushError::WriterStopped)
540    } else {
541        receiver
542            .recv()
543            .map_err(|_| FlushError::WriterStopped)
544            .and_then(WorkerStatus::into_result)
545    };
546
547    if let Some(worker) = inner.worker.lock().unwrap().take() {
548        if worker.join().is_err() {
549            result = Err(FlushError::WorkerPanicked);
550        }
551    }
552    completion.complete(result);
553}
554
555fn write_records(receiver: Receiver<Command>, mut writer: impl Write) {
556    let mut status = WorkerStatus::default();
557    while let Ok(command) = receiver.recv() {
558        match command {
559            Command::Record(record) => write_record(&mut writer, record, &mut status),
560            Command::Flush {
561                unreported_dropped,
562                response,
563            } => {
564                status.unreported_dropped =
565                    status.unreported_dropped.saturating_add(unreported_dropped);
566                flush_writer(&mut writer, &mut status);
567                let _ = response.send(status);
568            }
569            Command::Shutdown {
570                unreported_dropped,
571                response,
572            } => {
573                status.unreported_dropped =
574                    status.unreported_dropped.saturating_add(unreported_dropped);
575                flush_writer(&mut writer, &mut status);
576                let _ = response.send(status);
577                break;
578            }
579        }
580    }
581}
582
583fn write_record(writer: &mut impl Write, record: LogRecord, status: &mut WorkerStatus) {
584    if status.failure.is_some() {
585        status.unreported_dropped = status
586            .unreported_dropped
587            .saturating_add(record.dropped_events);
588        return;
589    }
590
591    let dropped = record.dropped_events;
592    let bytes = match serde_json::to_vec(&record) {
593        Ok(bytes) => bytes,
594        Err(_) => {
595            status.failure = Some(OutputStage::Serialize);
596            status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
597            return;
598        }
599    };
600    if writer.write_all(&bytes).is_err() {
601        status.failure = Some(OutputStage::Record);
602        status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
603        return;
604    }
605    if writer.write_all(b"\n").is_err() {
606        status.failure = Some(OutputStage::Newline);
607        status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
608    }
609}
610
611fn flush_writer(writer: &mut impl Write, status: &mut WorkerStatus) {
612    if writer.flush().is_err() && status.failure.is_none() {
613        status.failure = Some(OutputStage::Flush);
614    }
615}
616
617#[cfg(test)]
618mod tests {
619    use std::{
620        sync::{Condvar, mpsc},
621        task::{Wake, Waker},
622    };
623
624    use super::*;
625
626    struct ThreadWaker(thread::Thread);
627
628    impl Wake for ThreadWaker {
629        fn wake(self: Arc<Self>) {
630            self.0.unpark();
631        }
632    }
633
634    fn block_on<T>(future: impl Future<Output = T>) -> T {
635        let waker = Waker::from(Arc::new(ThreadWaker(thread::current())));
636        let mut context = Context::from_waker(&waker);
637        let mut future = std::pin::pin!(future);
638        loop {
639            match future.as_mut().poll(&mut context) {
640                Poll::Ready(output) => return output,
641                Poll::Pending => thread::park(),
642            }
643        }
644    }
645
646    struct BlockingWriter {
647        entered: Option<mpsc::Sender<()>>,
648        release: Arc<(Mutex<bool>, Condvar)>,
649        fail_after_release: bool,
650    }
651
652    impl Write for BlockingWriter {
653        fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
654            if let Some(entered) = self.entered.take() {
655                let _ = entered.send(());
656                let (lock, condition) = &*self.release;
657                let mut released = lock.lock().unwrap();
658                while !*released {
659                    released = condition.wait(released).unwrap();
660                }
661                if self.fail_after_release {
662                    return Err(io::Error::other("injected writer failure"));
663                }
664            }
665            Ok(bytes.len())
666        }
667
668        fn flush(&mut self) -> io::Result<()> {
669            Ok(())
670        }
671    }
672
673    struct FailingWriter {
674        fail_write: Option<usize>,
675        writes: usize,
676        fail_flush: bool,
677    }
678
679    impl Write for FailingWriter {
680        fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
681            self.writes += 1;
682            if self.fail_write == Some(self.writes) {
683                Err(io::Error::other("injected writer failure"))
684            } else {
685                Ok(bytes.len())
686            }
687        }
688
689        fn flush(&mut self) -> io::Result<()> {
690            if self.fail_flush {
691                Err(io::Error::other("injected flush failure"))
692            } else {
693                Ok(())
694            }
695        }
696    }
697
698    fn observer(writer: impl Write + Send + 'static, capacity: usize) -> Observer {
699        Observer::with_writer_and_seed(
700            ObserverConfig {
701                queue_capacity: capacity,
702            },
703            writer,
704            Ok([7; 24]),
705        )
706        .unwrap()
707    }
708
709    #[test]
710    fn random_source_failure_is_reported_during_initialization() {
711        let result = Observer::with_writer_and_seed(
712            ObserverConfig::default(),
713            io::sink(),
714            Err(InitError::RandomSource),
715        );
716        assert!(matches!(result, Err(InitError::RandomSource)));
717    }
718
719    #[test]
720    fn initialized_generator_produces_non_zero_unique_ids_without_more_io() {
721        let observer = observer(io::sink(), 8);
722        let first_trace = observer.new_trace_id();
723        let second_trace = observer.new_trace_id();
724        let first_span = observer.new_span_id();
725        let second_span = observer.new_span_id();
726        assert_ne!(first_trace.as_u128(), 0);
727        assert_ne!(first_trace, second_trace);
728        assert_ne!(first_span.as_u64(), 0);
729        assert_ne!(first_span, second_span);
730        block_on(observer.shutdown()).unwrap();
731    }
732
733    #[test]
734    fn observer_is_a_managed_component() {
735        let observer = observer(io::sink(), 8);
736        assert_eq!(ComponentLifecycle::name(&observer), "observability");
737        block_on(ComponentLifecycle::start(&observer)).unwrap();
738        block_on(ComponentLifecycle::shutdown(&observer)).unwrap();
739    }
740
741    #[test]
742    fn full_queue_is_reported_by_shutdown_without_waiting_in_emit() {
743        let (entered_sender, entered_receiver) = mpsc::channel();
744        let release = Arc::new((Mutex::new(false), Condvar::new()));
745        let writer = BlockingWriter {
746            entered: Some(entered_sender),
747            release: release.clone(),
748            fail_after_release: false,
749        };
750        let observer = observer(writer, 1);
751
752        observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
753        entered_receiver.recv().unwrap();
754        observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
755        observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
756        assert_eq!(observer.dropped_events(), 1);
757
758        let shutdown = observer.shutdown();
759        let (lock, condition) = &*release;
760        *lock.lock().unwrap() = true;
761        condition.notify_one();
762        assert_eq!(block_on(shutdown), Err(FlushError::DroppedEvents(1)));
763    }
764
765    #[test]
766    fn shutdown_reports_writer_failure_and_unreported_drops_together() {
767        let (entered_sender, entered_receiver) = mpsc::channel();
768        let release = Arc::new((Mutex::new(false), Condvar::new()));
769        let writer = BlockingWriter {
770            entered: Some(entered_sender),
771            release: release.clone(),
772            fail_after_release: true,
773        };
774        let observer = observer(writer, 1);
775
776        observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
777        entered_receiver.recv().unwrap();
778        observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
779        observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
780        let shutdown = observer.shutdown();
781
782        let (lock, condition) = &*release;
783        *lock.lock().unwrap() = true;
784        condition.notify_one();
785        assert_eq!(
786            block_on(shutdown),
787            Err(FlushError::OutputFailedAndDropped {
788                stage: OutputStage::Record,
789                dropped_events: 1,
790            })
791        );
792    }
793
794    #[test]
795    fn record_write_failure_persists_until_shutdown() {
796        let observer = observer(
797            FailingWriter {
798                fail_write: Some(1),
799                writes: 0,
800                fail_flush: false,
801            },
802            8,
803        );
804        observer.emit(LogRecord::new(EventLevel::Info, "test.record"));
805        assert_eq!(
806            block_on(observer.shutdown()),
807            Err(FlushError::OutputFailed(OutputStage::Record))
808        );
809    }
810
811    #[test]
812    fn newline_failure_persists_until_flush() {
813        let observer = observer(
814            FailingWriter {
815                fail_write: Some(2),
816                writes: 0,
817                fail_flush: false,
818            },
819            8,
820        );
821        observer.emit(LogRecord::new(EventLevel::Info, "test.newline"));
822        assert_eq!(
823            block_on(observer.flush()),
824            Err(FlushError::OutputFailed(OutputStage::Newline))
825        );
826        let _ = block_on(observer.shutdown());
827    }
828
829    #[test]
830    fn flush_failure_is_reported() {
831        let observer = observer(
832            FailingWriter {
833                fail_write: None,
834                writes: 0,
835                fail_flush: true,
836            },
837            8,
838        );
839        assert_eq!(
840            block_on(observer.flush()),
841            Err(FlushError::OutputFailed(OutputStage::Flush))
842        );
843        assert_eq!(
844            block_on(observer.shutdown()),
845            Err(FlushError::OutputFailed(OutputStage::Flush))
846        );
847    }
848}