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