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}
188
189struct IdGenerator {
190    trace_high: u64,
191    trace_seed: u64,
192    trace_counter: AtomicU64,
193    span_seed: u64,
194    span_counter: AtomicU64,
195}
196
197impl IdGenerator {
198    fn from_seed(seed: [u8; 24]) -> Self {
199        let mut high = u64::from_be_bytes(seed[0..8].try_into().unwrap());
200        if high == 0 {
201            high = 1;
202        }
203        Self {
204            trace_high: high,
205            trace_seed: u64::from_be_bytes(seed[8..16].try_into().unwrap()),
206            trace_counter: AtomicU64::new(0),
207            span_seed: u64::from_be_bytes(seed[16..24].try_into().unwrap()),
208            span_counter: AtomicU64::new(0),
209        }
210    }
211
212    fn trace_id(&self) -> TraceId {
213        let low = self
214            .trace_seed
215            .wrapping_add(self.trace_counter.fetch_add(1, Ordering::Relaxed));
216        TraceId::from_u128((u128::from(self.trace_high) << 64) | u128::from(low))
217    }
218
219    fn span_id(&self) -> SpanId {
220        loop {
221            let value = self
222                .span_seed
223                .wrapping_add(self.span_counter.fetch_add(1, Ordering::Relaxed));
224            if value != 0 {
225                return SpanId::from_u64(value);
226            }
227        }
228    }
229}
230
231enum Command {
232    Record(LogRecord),
233    Flush {
234        unreported_dropped: u64,
235        response: mpsc::Sender<WorkerStatus>,
236    },
237    Shutdown {
238        unreported_dropped: u64,
239        response: mpsc::Sender<WorkerStatus>,
240    },
241}
242
243#[derive(Clone, Copy, Debug, Default)]
244struct WorkerStatus {
245    failure: Option<OutputStage>,
246    unreported_dropped: u64,
247}
248
249impl WorkerStatus {
250    fn into_result(self) -> Result<(), FlushError> {
251        match (self.failure, self.unreported_dropped) {
252            (None, 0) => Ok(()),
253            (Some(stage), 0) => Err(FlushError::OutputFailed(stage)),
254            (None, dropped_events) => Err(FlushError::DroppedEvents(dropped_events)),
255            (Some(stage), dropped_events) => Err(FlushError::OutputFailedAndDropped {
256                stage,
257                dropped_events,
258            }),
259        }
260    }
261}
262
263#[derive(Serialize)]
264pub(crate) struct LogRecord {
265    pub timestamp_unix_ms: u128,
266    pub level: EventLevel,
267    pub event: &'static str,
268    #[serde(skip_serializing_if = "Option::is_none")]
269    pub trace_id: Option<String>,
270    #[serde(skip_serializing_if = "Option::is_none")]
271    pub span: Option<String>,
272    #[serde(skip_serializing_if = "Option::is_none")]
273    pub span_id: Option<String>,
274    #[serde(skip_serializing_if = "Option::is_none")]
275    pub parent: Option<String>,
276    #[serde(skip_serializing_if = "Option::is_none")]
277    pub parent_span_id: Option<String>,
278    #[serde(flatten)]
279    pub data: serde_json::Map<String, serde_json::Value>,
280    #[serde(skip_serializing_if = "is_zero")]
281    pub dropped_events: u64,
282}
283
284const fn is_zero(value: &u64) -> bool {
285    *value == 0
286}
287
288impl LogRecord {
289    pub(crate) fn new(level: EventLevel, event: &'static str) -> Self {
290        Self {
291            timestamp_unix_ms: SystemTime::now()
292                .duration_since(UNIX_EPOCH)
293                .unwrap_or_default()
294                .as_millis(),
295            level,
296            event,
297            trace_id: None,
298            span: None,
299            span_id: None,
300            parent: None,
301            parent_span_id: None,
302            data: serde_json::Map::new(),
303            dropped_events: 0,
304        }
305    }
306}
307
308/// Initializes stdout logging once. Later calls return the first observer.
309pub fn init(config: ObserverConfig) -> Result<&'static Observer, InitError> {
310    if let Some(observer) = GLOBAL.get() {
311        return Ok(observer);
312    }
313
314    let observer = Observer::with_writer(config, io::stdout())?;
315    let _ = GLOBAL.set(observer);
316    Ok(GLOBAL.get().expect("global observer was initialized"))
317}
318
319/// Initializes the process logger with the sole calendar-rotated file sink.
320pub fn init_file(
321    config: ObserverConfig,
322    file: FileLoggingConfig,
323) -> Result<&'static Observer, InitError> {
324    if let Some(observer) = GLOBAL.get() {
325        return Ok(observer);
326    }
327    let writer = CalendarFileWriter::open(file).map_err(InitError::Output)?;
328    let observer = Observer::with_writer(config, writer)?;
329    let _ = GLOBAL.set(observer);
330    Ok(GLOBAL.get().expect("global observer was initialized"))
331}
332
333/// Returns the initialized process observer, if application startup installed it.
334pub fn global() -> Option<&'static Observer> {
335    GLOBAL.get()
336}
337
338impl Observer {
339    /// Creates an observer with a framework-owned writer.
340    ///
341    /// Application startup should normally use [`init`]. This constructor is
342    /// useful for embedding and deterministic tests.
343    pub fn with_writer(
344        config: ObserverConfig,
345        writer: impl Write + Send + 'static,
346    ) -> Result<Self, InitError> {
347        let mut seed = [0_u8; 24];
348        getrandom::fill(&mut seed).map_err(|_| InitError::RandomSource)?;
349        Self::with_writer_and_seed(config, writer, Ok(seed))
350    }
351
352    fn with_writer_and_seed(
353        config: ObserverConfig,
354        writer: impl Write + Send + 'static,
355        seed: Result<[u8; 24], InitError>,
356    ) -> Result<Self, InitError> {
357        if config.queue_capacity == 0 {
358            return Err(InitError::EmptyQueue);
359        }
360        let ids = IdGenerator::from_seed(seed?);
361        let (sender, receiver) = mpsc::sync_channel(config.queue_capacity);
362        let worker = thread::Builder::new()
363            .name("saddle-log-writer".to_owned())
364            .spawn(move || write_records(receiver, writer))
365            .map_err(InitError::Spawn)?;
366        Ok(Self {
367            inner: Arc::new(Inner {
368                sender,
369                dropped: AtomicU64::new(0),
370                accepting: AtomicBool::new(true),
371                emitting: AtomicUsize::new(0),
372                shutdown_started: AtomicBool::new(false),
373                worker: Mutex::new(Some(worker)),
374                ids,
375            }),
376        })
377    }
378
379    pub(crate) fn new_trace_id(&self) -> TraceId {
380        self.inner.ids.trace_id()
381    }
382
383    pub(crate) fn new_span_id(&self) -> SpanId {
384        self.inner.ids.span_id()
385    }
386
387    pub(crate) fn emit(&self, mut record: LogRecord) {
388        if !self.inner.accepting.load(Ordering::Acquire) {
389            return;
390        }
391        self.inner.emitting.fetch_add(1, Ordering::AcqRel);
392        if !self.inner.accepting.load(Ordering::Acquire) {
393            self.inner.emitting.fetch_sub(1, Ordering::AcqRel);
394            return;
395        }
396
397        record.dropped_events = self.inner.dropped.swap(0, Ordering::Relaxed);
398        match self.inner.sender.try_send(Command::Record(record)) {
399            Ok(()) => {}
400            Err(TrySendError::Full(Command::Record(record))) => {
401                self.inner
402                    .dropped
403                    .fetch_add(record.dropped_events.saturating_add(1), Ordering::Relaxed);
404            }
405            Err(TrySendError::Disconnected(_)) => {
406                self.inner.dropped.fetch_add(1, Ordering::Relaxed);
407            }
408            Err(TrySendError::Full(_)) => unreachable!("emit only sends records"),
409        }
410        self.inner.emitting.fetch_sub(1, Ordering::Release);
411    }
412
413    /// Flushes records accepted before the lifecycle coordinator reaches the
414    /// writer. The returned future never blocks the async executor thread.
415    pub fn flush(&self) -> FlushFuture {
416        if !self.inner.accepting.load(Ordering::Acquire) {
417            return FlushFuture {
418                completion: Completion::ready(Err(FlushError::WriterStopped)),
419            };
420        }
421        let completion = Completion::pending();
422        let future = FlushFuture {
423            completion: completion.clone(),
424        };
425        let inner = self.inner.clone();
426        if thread::Builder::new()
427            .name("saddle-log-flush".to_owned())
428            .spawn(move || coordinate_flush(inner, completion.clone()))
429            .is_err()
430        {
431            future
432                .completion
433                .complete(Err(FlushError::CoordinatorUnavailable));
434        }
435        future
436    }
437
438    /// Stops admission, drains accepted records, flushes output and joins the
439    /// writer without blocking the async executor thread.
440    pub fn shutdown(&self) -> FlushFuture {
441        if self.inner.shutdown_started.swap(true, Ordering::AcqRel) {
442            return FlushFuture {
443                completion: Completion::ready(Err(FlushError::AlreadyShuttingDown)),
444            };
445        }
446        self.inner.accepting.store(false, Ordering::Release);
447        let completion = Completion::pending();
448        let future = FlushFuture {
449            completion: completion.clone(),
450        };
451        let inner = self.inner.clone();
452        if thread::Builder::new()
453            .name("saddle-log-shutdown".to_owned())
454            .spawn(move || coordinate_shutdown(inner, completion.clone()))
455            .is_err()
456        {
457            self.inner.accepting.store(true, Ordering::Release);
458            self.inner.shutdown_started.store(false, Ordering::Release);
459            future
460                .completion
461                .complete(Err(FlushError::CoordinatorUnavailable));
462        }
463        future
464    }
465
466    pub fn dropped_events(&self) -> u64 {
467        self.inner.dropped.load(Ordering::Relaxed)
468    }
469}
470
471impl ComponentLifecycle for Observer {
472    fn name(&self) -> &'static str {
473        "observability"
474    }
475
476    fn start(&self) -> LifecycleFuture<'_> {
477        Box::pin(async { Ok(()) })
478    }
479
480    fn shutdown(&self) -> LifecycleFuture<'_> {
481        let shutdown = Observer::shutdown(self);
482        Box::pin(async move {
483            shutdown.await.map_err(|_| {
484                SaddleError::new(
485                    ErrorKind::Infrastructure,
486                    "observability.shutdown_failed",
487                    "structured log shutdown failed",
488                )
489            })
490        })
491    }
492}
493
494fn coordinate_flush(inner: Arc<Inner>, completion: Arc<Completion>) {
495    let dropped = inner.dropped.swap(0, Ordering::AcqRel);
496    let (response, receiver) = mpsc::channel();
497    let result = if inner
498        .sender
499        .send(Command::Flush {
500            unreported_dropped: dropped,
501            response,
502        })
503        .is_err()
504    {
505        Err(FlushError::WriterStopped)
506    } else {
507        receiver
508            .recv()
509            .map_err(|_| FlushError::WriterStopped)
510            .and_then(WorkerStatus::into_result)
511    };
512    completion.complete(result);
513}
514
515fn coordinate_shutdown(inner: Arc<Inner>, completion: Arc<Completion>) {
516    while inner.emitting.load(Ordering::Acquire) != 0 {
517        thread::yield_now();
518    }
519    let dropped = inner.dropped.swap(0, Ordering::AcqRel);
520    let (response, receiver) = mpsc::channel();
521    let mut result = if inner
522        .sender
523        .send(Command::Shutdown {
524            unreported_dropped: dropped,
525            response,
526        })
527        .is_err()
528    {
529        Err(FlushError::WriterStopped)
530    } else {
531        receiver
532            .recv()
533            .map_err(|_| FlushError::WriterStopped)
534            .and_then(WorkerStatus::into_result)
535    };
536
537    if let Some(worker) = inner.worker.lock().unwrap().take() {
538        if worker.join().is_err() {
539            result = Err(FlushError::WorkerPanicked);
540        }
541    }
542    completion.complete(result);
543}
544
545fn write_records(receiver: Receiver<Command>, mut writer: impl Write) {
546    let mut status = WorkerStatus::default();
547    while let Ok(command) = receiver.recv() {
548        match command {
549            Command::Record(record) => write_record(&mut writer, record, &mut status),
550            Command::Flush {
551                unreported_dropped,
552                response,
553            } => {
554                status.unreported_dropped =
555                    status.unreported_dropped.saturating_add(unreported_dropped);
556                flush_writer(&mut writer, &mut status);
557                let _ = response.send(status);
558            }
559            Command::Shutdown {
560                unreported_dropped,
561                response,
562            } => {
563                status.unreported_dropped =
564                    status.unreported_dropped.saturating_add(unreported_dropped);
565                flush_writer(&mut writer, &mut status);
566                let _ = response.send(status);
567                break;
568            }
569        }
570    }
571}
572
573fn write_record(writer: &mut impl Write, record: LogRecord, status: &mut WorkerStatus) {
574    if status.failure.is_some() {
575        status.unreported_dropped = status
576            .unreported_dropped
577            .saturating_add(record.dropped_events);
578        return;
579    }
580
581    let dropped = record.dropped_events;
582    let bytes = match serde_json::to_vec(&record) {
583        Ok(bytes) => bytes,
584        Err(_) => {
585            status.failure = Some(OutputStage::Serialize);
586            status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
587            return;
588        }
589    };
590    if writer.write_all(&bytes).is_err() {
591        status.failure = Some(OutputStage::Record);
592        status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
593        return;
594    }
595    if writer.write_all(b"\n").is_err() {
596        status.failure = Some(OutputStage::Newline);
597        status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
598    }
599}
600
601fn flush_writer(writer: &mut impl Write, status: &mut WorkerStatus) {
602    if writer.flush().is_err() && status.failure.is_none() {
603        status.failure = Some(OutputStage::Flush);
604    }
605}
606
607#[cfg(test)]
608mod tests {
609    use std::{
610        sync::{Condvar, mpsc},
611        task::{Wake, Waker},
612    };
613
614    use super::*;
615
616    struct ThreadWaker(thread::Thread);
617
618    impl Wake for ThreadWaker {
619        fn wake(self: Arc<Self>) {
620            self.0.unpark();
621        }
622    }
623
624    fn block_on<T>(future: impl Future<Output = T>) -> T {
625        let waker = Waker::from(Arc::new(ThreadWaker(thread::current())));
626        let mut context = Context::from_waker(&waker);
627        let mut future = std::pin::pin!(future);
628        loop {
629            match future.as_mut().poll(&mut context) {
630                Poll::Ready(output) => return output,
631                Poll::Pending => thread::park(),
632            }
633        }
634    }
635
636    struct BlockingWriter {
637        entered: Option<mpsc::Sender<()>>,
638        release: Arc<(Mutex<bool>, Condvar)>,
639        fail_after_release: bool,
640    }
641
642    impl Write for BlockingWriter {
643        fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
644            if let Some(entered) = self.entered.take() {
645                let _ = entered.send(());
646                let (lock, condition) = &*self.release;
647                let mut released = lock.lock().unwrap();
648                while !*released {
649                    released = condition.wait(released).unwrap();
650                }
651                if self.fail_after_release {
652                    return Err(io::Error::other("injected writer failure"));
653                }
654            }
655            Ok(bytes.len())
656        }
657
658        fn flush(&mut self) -> io::Result<()> {
659            Ok(())
660        }
661    }
662
663    struct FailingWriter {
664        fail_write: Option<usize>,
665        writes: usize,
666        fail_flush: bool,
667    }
668
669    impl Write for FailingWriter {
670        fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
671            self.writes += 1;
672            if self.fail_write == Some(self.writes) {
673                Err(io::Error::other("injected writer failure"))
674            } else {
675                Ok(bytes.len())
676            }
677        }
678
679        fn flush(&mut self) -> io::Result<()> {
680            if self.fail_flush {
681                Err(io::Error::other("injected flush failure"))
682            } else {
683                Ok(())
684            }
685        }
686    }
687
688    fn observer(writer: impl Write + Send + 'static, capacity: usize) -> Observer {
689        Observer::with_writer_and_seed(
690            ObserverConfig {
691                queue_capacity: capacity,
692            },
693            writer,
694            Ok([7; 24]),
695        )
696        .unwrap()
697    }
698
699    #[test]
700    fn random_source_failure_is_reported_during_initialization() {
701        let result = Observer::with_writer_and_seed(
702            ObserverConfig::default(),
703            io::sink(),
704            Err(InitError::RandomSource),
705        );
706        assert!(matches!(result, Err(InitError::RandomSource)));
707    }
708
709    #[test]
710    fn initialized_generator_produces_non_zero_unique_ids_without_more_io() {
711        let observer = observer(io::sink(), 8);
712        let first_trace = observer.new_trace_id();
713        let second_trace = observer.new_trace_id();
714        let first_span = observer.new_span_id();
715        let second_span = observer.new_span_id();
716        assert_ne!(first_trace.as_u128(), 0);
717        assert_ne!(first_trace, second_trace);
718        assert_ne!(first_span.as_u64(), 0);
719        assert_ne!(first_span, second_span);
720        block_on(observer.shutdown()).unwrap();
721    }
722
723    #[test]
724    fn observer_is_a_managed_component() {
725        let observer = observer(io::sink(), 8);
726        assert_eq!(ComponentLifecycle::name(&observer), "observability");
727        block_on(ComponentLifecycle::start(&observer)).unwrap();
728        block_on(ComponentLifecycle::shutdown(&observer)).unwrap();
729    }
730
731    #[test]
732    fn full_queue_is_reported_by_shutdown_without_waiting_in_emit() {
733        let (entered_sender, entered_receiver) = mpsc::channel();
734        let release = Arc::new((Mutex::new(false), Condvar::new()));
735        let writer = BlockingWriter {
736            entered: Some(entered_sender),
737            release: release.clone(),
738            fail_after_release: false,
739        };
740        let observer = observer(writer, 1);
741
742        observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
743        entered_receiver.recv().unwrap();
744        observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
745        observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
746        assert_eq!(observer.dropped_events(), 1);
747
748        let shutdown = observer.shutdown();
749        let (lock, condition) = &*release;
750        *lock.lock().unwrap() = true;
751        condition.notify_one();
752        assert_eq!(block_on(shutdown), Err(FlushError::DroppedEvents(1)));
753    }
754
755    #[test]
756    fn shutdown_reports_writer_failure_and_unreported_drops_together() {
757        let (entered_sender, entered_receiver) = mpsc::channel();
758        let release = Arc::new((Mutex::new(false), Condvar::new()));
759        let writer = BlockingWriter {
760            entered: Some(entered_sender),
761            release: release.clone(),
762            fail_after_release: true,
763        };
764        let observer = observer(writer, 1);
765
766        observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
767        entered_receiver.recv().unwrap();
768        observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
769        observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
770        let shutdown = observer.shutdown();
771
772        let (lock, condition) = &*release;
773        *lock.lock().unwrap() = true;
774        condition.notify_one();
775        assert_eq!(
776            block_on(shutdown),
777            Err(FlushError::OutputFailedAndDropped {
778                stage: OutputStage::Record,
779                dropped_events: 1,
780            })
781        );
782    }
783
784    #[test]
785    fn record_write_failure_persists_until_shutdown() {
786        let observer = observer(
787            FailingWriter {
788                fail_write: Some(1),
789                writes: 0,
790                fail_flush: false,
791            },
792            8,
793        );
794        observer.emit(LogRecord::new(EventLevel::Info, "test.record"));
795        assert_eq!(
796            block_on(observer.shutdown()),
797            Err(FlushError::OutputFailed(OutputStage::Record))
798        );
799    }
800
801    #[test]
802    fn newline_failure_persists_until_flush() {
803        let observer = observer(
804            FailingWriter {
805                fail_write: Some(2),
806                writes: 0,
807                fail_flush: false,
808            },
809            8,
810        );
811        observer.emit(LogRecord::new(EventLevel::Info, "test.newline"));
812        assert_eq!(
813            block_on(observer.flush()),
814            Err(FlushError::OutputFailed(OutputStage::Newline))
815        );
816        let _ = block_on(observer.shutdown());
817    }
818
819    #[test]
820    fn flush_failure_is_reported() {
821        let observer = observer(
822            FailingWriter {
823                fail_write: None,
824                writes: 0,
825                fail_flush: true,
826            },
827            8,
828        );
829        assert_eq!(
830            block_on(observer.flush()),
831            Err(FlushError::OutputFailed(OutputStage::Flush))
832        );
833        assert_eq!(
834            block_on(observer.shutdown()),
835            Err(FlushError::OutputFailed(OutputStage::Flush))
836        );
837    }
838}