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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
26pub struct ObserverConfig {
27 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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
70pub enum OutputStage {
71 Serialize,
72 Record,
73 Newline,
74 Flush,
75}
76
77#[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
116pub 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#[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
303pub 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
314pub fn global() -> Option<&'static Observer> {
316 GLOBAL.get()
317}
318
319impl Observer {
320 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 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 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}