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