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_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
297pub 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
308pub fn global() -> Option<&'static Observer> {
310 GLOBAL.get()
311}
312
313impl Observer {
314 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 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 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}