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