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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
29pub struct ObserverConfig {
30 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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
75pub enum OutputStage {
76 Serialize,
77 Record,
78 Newline,
79 Flush,
80}
81
82#[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
121pub 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#[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
308pub 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
319pub 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
333pub fn global() -> Option<&'static Observer> {
335 GLOBAL.get()
336}
337
338impl Observer {
339 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 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 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}