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::{
21    calendar_writer::{CalendarFileWriter, FileLoggingConfig},
22    event::EventLevel,
23};
24
25static GLOBAL: OnceLock<Observer> = OnceLock::new();
26
27/// Configuration for the bounded non-blocking logger.
28#[derive(Clone, Copy, Debug, Eq, PartialEq)]
29pub struct ObserverConfig {
30    /// Maximum records waiting for the output worker. New records are dropped
31    /// when the queue is full rather than blocking an async request.
32    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/// The output operation whose first failure made log delivery unhealthy.
74#[derive(Clone, Copy, Debug, Eq, PartialEq)]
75pub enum OutputStage {
76    Serialize,
77    Record,
78    Newline,
79    Flush,
80}
81
82/// A flush or managed shutdown result that is safe to expose to framework code.
83#[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
121/// A non-blocking future returned by flush and shutdown lifecycle operations.
122pub 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/// Process-level structured logger and trace correlator.
171///
172/// Cloning this value is cheap. Emission uses `try_send`; output and JSON
173/// serialization happen on the dedicated writer thread.
174#[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    // Source-encoded, bounded and self-contained. Never contains a request root.
235    RootRecord {
236        bytes: Vec<u8>,
237        dropped: u64,
238    },
239    Flush {
240        unreported_dropped: u64,
241        response: mpsc::Sender<WorkerStatus>,
242    },
243    Shutdown {
244        unreported_dropped: u64,
245        response: mpsc::Sender<WorkerStatus>,
246    },
247}
248pub(crate) fn root_queue_layout() -> std::alloc::Layout {
249    std::alloc::Layout::new::<Command>()
250}
251
252#[derive(Clone, Copy, Debug, Default)]
253struct WorkerStatus {
254    failure: Option<OutputStage>,
255    unreported_dropped: u64,
256}
257
258impl WorkerStatus {
259    fn into_result(self) -> Result<(), FlushError> {
260        match (self.failure, self.unreported_dropped) {
261            (None, 0) => Ok(()),
262            (Some(stage), 0) => Err(FlushError::OutputFailed(stage)),
263            (None, dropped_events) => Err(FlushError::DroppedEvents(dropped_events)),
264            (Some(stage), dropped_events) => Err(FlushError::OutputFailedAndDropped {
265                stage,
266                dropped_events,
267            }),
268        }
269    }
270}
271
272#[derive(Serialize)]
273pub(crate) struct LogRecord {
274    pub timestamp_unix_ms: u128,
275    pub level: EventLevel,
276    pub event: &'static str,
277    #[serde(skip_serializing_if = "Option::is_none")]
278    pub trace_id: Option<String>,
279    #[serde(skip_serializing_if = "Option::is_none")]
280    pub span: Option<String>,
281    #[serde(skip_serializing_if = "Option::is_none")]
282    pub span_id: Option<String>,
283    #[serde(skip_serializing_if = "Option::is_none")]
284    pub parent: Option<String>,
285    #[serde(skip_serializing_if = "Option::is_none")]
286    pub parent_span_id: Option<String>,
287    #[serde(flatten)]
288    pub data: serde_json::Map<String, serde_json::Value>,
289    #[serde(skip_serializing_if = "is_zero")]
290    pub dropped_events: u64,
291}
292
293const fn is_zero(value: &u64) -> bool {
294    *value == 0
295}
296
297impl LogRecord {
298    pub(crate) fn new(level: EventLevel, event: &'static str) -> Self {
299        Self {
300            timestamp_unix_ms: SystemTime::now()
301                .duration_since(UNIX_EPOCH)
302                .unwrap_or_default()
303                .as_millis(),
304            level,
305            event,
306            trace_id: None,
307            span: None,
308            span_id: None,
309            parent: None,
310            parent_span_id: None,
311            data: serde_json::Map::new(),
312            dropped_events: 0,
313        }
314    }
315}
316
317/// Initializes stdout logging once. Later calls return the first observer.
318pub fn init(config: ObserverConfig) -> Result<&'static Observer, InitError> {
319    if let Some(observer) = GLOBAL.get() {
320        return Ok(observer);
321    }
322
323    let observer = Observer::with_writer(config, io::stdout())?;
324    let _ = GLOBAL.set(observer);
325    Ok(GLOBAL.get().expect("global observer was initialized"))
326}
327
328/// Initializes the process logger with the sole calendar-rotated file sink.
329pub fn init_file(
330    config: ObserverConfig,
331    file: FileLoggingConfig,
332) -> Result<&'static Observer, InitError> {
333    if let Some(observer) = GLOBAL.get() {
334        return Ok(observer);
335    }
336    let writer = CalendarFileWriter::open(file).map_err(InitError::Output)?;
337    let observer = Observer::with_writer(config, writer)?;
338    let _ = GLOBAL.set(observer);
339    Ok(GLOBAL.get().expect("global observer was initialized"))
340}
341
342/// Returns the initialized process observer, if application startup installed it.
343pub fn global() -> Option<&'static Observer> {
344    GLOBAL.get()
345}
346
347impl Observer {
348    /// Creates an observer with a framework-owned writer.
349    ///
350    /// Application startup should normally use [`init`]. This constructor is
351    /// useful for embedding and deterministic tests.
352    pub fn with_writer(
353        config: ObserverConfig,
354        writer: impl Write + Send + 'static,
355    ) -> Result<Self, InitError> {
356        let mut seed = [0_u8; 24];
357        getrandom::fill(&mut seed).map_err(|_| InitError::RandomSource)?;
358        Self::with_writer_and_seed(config, writer, Ok(seed))
359    }
360
361    fn with_writer_and_seed(
362        config: ObserverConfig,
363        writer: impl Write + Send + 'static,
364        seed: Result<[u8; 24], InitError>,
365    ) -> Result<Self, InitError> {
366        if config.queue_capacity == 0 {
367            return Err(InitError::EmptyQueue);
368        }
369        let ids = IdGenerator::from_seed(seed?);
370        let (sender, receiver) = mpsc::sync_channel(config.queue_capacity);
371        let worker = thread::Builder::new()
372            .name("saddle-log-writer".to_owned())
373            .spawn(move || write_records(receiver, writer))
374            .map_err(InitError::Spawn)?;
375        Ok(Self {
376            inner: Arc::new(Inner {
377                sender,
378                dropped: AtomicU64::new(0),
379                accepting: AtomicBool::new(true),
380                emitting: AtomicUsize::new(0),
381                shutdown_started: AtomicBool::new(false),
382                worker: Mutex::new(Some(worker)),
383                ids,
384                metrics: crate::metrics::Metrics::default(),
385            }),
386        })
387    }
388
389    pub(crate) fn new_trace_id(&self) -> TraceId {
390        self.inner.ids.trace_id()
391    }
392
393    pub(crate) fn new_span_id(&self) -> SpanId {
394        self.inner.ids.span_id()
395    }
396
397    pub(crate) fn emit(&self, mut record: LogRecord) {
398        if !self.inner.accepting.load(Ordering::Acquire) {
399            return;
400        }
401        self.inner.emitting.fetch_add(1, Ordering::AcqRel);
402        if !self.inner.accepting.load(Ordering::Acquire) {
403            self.inner.emitting.fetch_sub(1, Ordering::AcqRel);
404            return;
405        }
406
407        record.dropped_events = self.inner.dropped.swap(0, Ordering::Relaxed);
408        match self.inner.sender.try_send(Command::Record(record)) {
409            Ok(()) => {}
410            Err(TrySendError::Full(Command::Record(record))) => {
411                self.inner.metrics.logger_dropped(1);
412                self.inner
413                    .dropped
414                    .fetch_add(record.dropped_events.saturating_add(1), Ordering::Relaxed);
415            }
416            Err(TrySendError::Disconnected(_)) => {
417                self.inner.metrics.logger_dropped(1);
418                self.inner.metrics.logger_output(true);
419                self.inner.dropped.fetch_add(1, Ordering::Relaxed);
420            }
421            Err(TrySendError::Full(_)) => unreachable!("emit only sends records"),
422        }
423        self.inner.emitting.fetch_sub(1, Ordering::Release);
424    }
425
426    pub(crate) fn emit_root_record<T: Serialize>(&self, record: &T) -> crate::DiagnosticSubmission {
427        use crate::DiagnosticSubmission as Submission;
428        self.inner.emitting.fetch_add(1, Ordering::AcqRel);
429        let result = (|| {
430            if !self.inner.accepting.load(Ordering::Acquire) {
431                return Submission::Closed;
432            }
433            let mut frame = crate::diagnostic::FixedDiagnosticBytes {
434                bytes: [0; 8192],
435                len: 0,
436            };
437            if serde_json::to_writer(&mut frame, record).is_err() {
438                return Submission::EncodingFailed;
439            }
440            let dropped = self.inner.dropped.swap(0, Ordering::Relaxed);
441            // Only initialized bytes move into the existing logger domain. This
442            // allocation and source frame peak are separate R0 layout inputs.
443            let bytes = frame.bytes[..frame.len].to_vec();
444            match self
445                .inner
446                .sender
447                .try_send(Command::RootRecord { bytes, dropped })
448            {
449                Ok(()) => Submission::Enqueued,
450                Err(TrySendError::Full(Command::RootRecord { dropped, .. })) => {
451                    self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
452                    Submission::Full
453                }
454                Err(TrySendError::Disconnected(_)) => {
455                    self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
456                    self.inner.metrics.logger_output(true);
457                    Submission::Closed
458                }
459                Err(TrySendError::Full(_)) => unreachable!("root record submission"),
460            }
461        })();
462        if result != Submission::Enqueued {
463            self.inner.dropped.fetch_add(1, Ordering::Relaxed);
464            self.inner.metrics.logger_dropped(1);
465        }
466        self.inner.emitting.fetch_sub(1, Ordering::Release);
467        result
468    }
469
470    /// Returns a fixed-size, allocation-free copy of the current process metrics.
471    pub fn metrics_snapshot(&self) -> crate::MetricsSnapshot {
472        self.inner.metrics.snapshot()
473    }
474
475    /// Flushes records accepted before the lifecycle coordinator reaches the
476    /// writer. The returned future never blocks the async executor thread.
477    pub fn flush(&self) -> FlushFuture {
478        if !self.inner.accepting.load(Ordering::Acquire) {
479            return FlushFuture {
480                completion: Completion::ready(Err(FlushError::WriterStopped)),
481            };
482        }
483        let completion = Completion::pending();
484        let future = FlushFuture {
485            completion: completion.clone(),
486        };
487        let inner = self.inner.clone();
488        if thread::Builder::new()
489            .name("saddle-log-flush".to_owned())
490            .spawn(move || coordinate_flush(inner, completion.clone()))
491            .is_err()
492        {
493            future
494                .completion
495                .complete(Err(FlushError::CoordinatorUnavailable));
496        }
497        future
498    }
499
500    /// Stops admission, drains accepted records, flushes output and joins the
501    /// writer without blocking the async executor thread.
502    pub fn shutdown(&self) -> FlushFuture {
503        if self.inner.shutdown_started.swap(true, Ordering::AcqRel) {
504            return FlushFuture {
505                completion: Completion::ready(Err(FlushError::AlreadyShuttingDown)),
506            };
507        }
508        self.inner.accepting.store(false, Ordering::Release);
509        let completion = Completion::pending();
510        let future = FlushFuture {
511            completion: completion.clone(),
512        };
513        let inner = self.inner.clone();
514        if thread::Builder::new()
515            .name("saddle-log-shutdown".to_owned())
516            .spawn(move || coordinate_shutdown(inner, completion.clone()))
517            .is_err()
518        {
519            self.inner.accepting.store(true, Ordering::Release);
520            self.inner.shutdown_started.store(false, Ordering::Release);
521            future
522                .completion
523                .complete(Err(FlushError::CoordinatorUnavailable));
524        }
525        future
526    }
527
528    pub fn dropped_events(&self) -> u64 {
529        self.inner.dropped.load(Ordering::Relaxed)
530    }
531}
532
533impl ComponentLifecycle for Observer {
534    fn name(&self) -> &'static str {
535        "observability"
536    }
537
538    fn start(&self) -> LifecycleFuture<'_> {
539        Box::pin(async { Ok(()) })
540    }
541
542    fn shutdown(&self) -> LifecycleFuture<'_> {
543        let shutdown = Observer::shutdown(self);
544        Box::pin(async move {
545            shutdown.await.map_err(|_| {
546                SaddleError::new(
547                    ErrorKind::Infrastructure,
548                    "observability.shutdown_failed",
549                    "structured log shutdown failed",
550                )
551            })
552        })
553    }
554}
555
556fn coordinate_flush(inner: Arc<Inner>, completion: Arc<Completion>) {
557    let dropped = inner.dropped.swap(0, Ordering::AcqRel);
558    let (response, receiver) = mpsc::channel();
559    let result = if inner
560        .sender
561        .send(Command::Flush {
562            unreported_dropped: dropped,
563            response,
564        })
565        .is_err()
566    {
567        Err(FlushError::WriterStopped)
568    } else {
569        receiver
570            .recv()
571            .map_err(|_| FlushError::WriterStopped)
572            .and_then(WorkerStatus::into_result)
573    };
574    completion.complete(result);
575}
576
577fn coordinate_shutdown(inner: Arc<Inner>, completion: Arc<Completion>) {
578    while inner.emitting.load(Ordering::Acquire) != 0 {
579        thread::yield_now();
580    }
581    let dropped = inner.dropped.swap(0, Ordering::AcqRel);
582    let (response, receiver) = mpsc::channel();
583    let mut result = if inner
584        .sender
585        .send(Command::Shutdown {
586            unreported_dropped: dropped,
587            response,
588        })
589        .is_err()
590    {
591        Err(FlushError::WriterStopped)
592    } else {
593        receiver
594            .recv()
595            .map_err(|_| FlushError::WriterStopped)
596            .and_then(WorkerStatus::into_result)
597    };
598
599    if let Some(worker) = inner.worker.lock().unwrap().take() {
600        if worker.join().is_err() {
601            result = Err(FlushError::WorkerPanicked);
602        }
603    }
604    completion.complete(result);
605}
606
607fn write_records(receiver: Receiver<Command>, mut writer: impl Write) {
608    let mut status = WorkerStatus::default();
609    while let Ok(command) = receiver.recv() {
610        match command {
611            Command::Record(record) => write_record(&mut writer, record, &mut status),
612            Command::RootRecord { bytes, dropped } => {
613                status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
614                if status.failure.is_none() {
615                    if writer.write_all(&bytes).is_err() {
616                        status.failure = Some(OutputStage::Record);
617                    } else if writer.write_all(b"\n").is_err() {
618                        status.failure = Some(OutputStage::Newline);
619                    }
620                }
621            }
622            Command::Flush {
623                unreported_dropped,
624                response,
625            } => {
626                status.unreported_dropped =
627                    status.unreported_dropped.saturating_add(unreported_dropped);
628                flush_writer(&mut writer, &mut status);
629                let _ = response.send(status);
630            }
631            Command::Shutdown {
632                unreported_dropped,
633                response,
634            } => {
635                status.unreported_dropped =
636                    status.unreported_dropped.saturating_add(unreported_dropped);
637                flush_writer(&mut writer, &mut status);
638                let _ = response.send(status);
639                break;
640            }
641        }
642    }
643}
644
645fn write_record(writer: &mut impl Write, record: LogRecord, status: &mut WorkerStatus) {
646    if status.failure.is_some() {
647        status.unreported_dropped = status
648            .unreported_dropped
649            .saturating_add(record.dropped_events);
650        return;
651    }
652
653    let dropped = record.dropped_events;
654    let bytes = match serde_json::to_vec(&record) {
655        Ok(bytes) => bytes,
656        Err(_) => {
657            status.failure = Some(OutputStage::Serialize);
658            status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
659            return;
660        }
661    };
662    if writer.write_all(&bytes).is_err() {
663        status.failure = Some(OutputStage::Record);
664        status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
665        return;
666    }
667    if writer.write_all(b"\n").is_err() {
668        status.failure = Some(OutputStage::Newline);
669        status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
670    }
671}
672
673fn flush_writer(writer: &mut impl Write, status: &mut WorkerStatus) {
674    if writer.flush().is_err() && status.failure.is_none() {
675        status.failure = Some(OutputStage::Flush);
676    }
677}
678
679#[cfg(test)]
680mod tests {
681    use std::{
682        sync::{Condvar, mpsc},
683        task::{Wake, Waker},
684    };
685
686    use super::*;
687
688    struct ThreadWaker(thread::Thread);
689
690    impl Wake for ThreadWaker {
691        fn wake(self: Arc<Self>) {
692            self.0.unpark();
693        }
694    }
695
696    fn block_on<T>(future: impl Future<Output = T>) -> T {
697        let waker = Waker::from(Arc::new(ThreadWaker(thread::current())));
698        let mut context = Context::from_waker(&waker);
699        let mut future = std::pin::pin!(future);
700        loop {
701            match future.as_mut().poll(&mut context) {
702                Poll::Ready(output) => return output,
703                Poll::Pending => thread::park(),
704            }
705        }
706    }
707
708    struct BlockingWriter {
709        entered: Option<mpsc::Sender<()>>,
710        release: Arc<(Mutex<bool>, Condvar)>,
711        fail_after_release: bool,
712    }
713
714    impl Write for BlockingWriter {
715        fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
716            if let Some(entered) = self.entered.take() {
717                let _ = entered.send(());
718                let (lock, condition) = &*self.release;
719                let mut released = lock.lock().unwrap();
720                while !*released {
721                    released = condition.wait(released).unwrap();
722                }
723                if self.fail_after_release {
724                    return Err(io::Error::other("injected writer failure"));
725                }
726            }
727            Ok(bytes.len())
728        }
729
730        fn flush(&mut self) -> io::Result<()> {
731            Ok(())
732        }
733    }
734
735    struct FailingWriter {
736        fail_write: Option<usize>,
737        writes: usize,
738        fail_flush: bool,
739    }
740
741    impl Write for FailingWriter {
742        fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
743            self.writes += 1;
744            if self.fail_write == Some(self.writes) {
745                Err(io::Error::other("injected writer failure"))
746            } else {
747                Ok(bytes.len())
748            }
749        }
750
751        fn flush(&mut self) -> io::Result<()> {
752            if self.fail_flush {
753                Err(io::Error::other("injected flush failure"))
754            } else {
755                Ok(())
756            }
757        }
758    }
759
760    fn observer(writer: impl Write + Send + 'static, capacity: usize) -> Observer {
761        Observer::with_writer_and_seed(
762            ObserverConfig {
763                queue_capacity: capacity,
764            },
765            writer,
766            Ok([7; 24]),
767        )
768        .unwrap()
769    }
770
771    #[test]
772    fn root_contract_ordinary_uses_same_view_and_bounded_queue() {
773        use crate::{
774            DiagnosticSubmission, RootDiagnosticScope, RootOutcomeFacts, RootRequestEvent,
775        };
776        use saddle_core::{
777            ContextFact, ContextLabel, RequestLocalFacts, RequestRootPublisher, RequestViewPhase,
778        };
779        let root = RequestRootPublisher::create(
780            ContextLabel::checked("app").unwrap(),
781            ContextFact::NotEstablished,
782        )
783        .unwrap();
784        let view = root
785            .reference()
786            .view(RequestLocalFacts::new(RequestViewPhase::Handler));
787        let expected = serde_json::to_value(&view).unwrap();
788        let mut frame = crate::diagnostic::FixedDiagnosticBytes {
789            bytes: [0; 8192],
790            len: 0,
791        };
792        serde_json::to_writer(&mut frame, &view).unwrap();
793        let ordinary_context: serde_json::Value =
794            serde_json::from_slice(&frame.bytes[..frame.len]).unwrap();
795        assert_eq!(ordinary_context, expected);
796        let (entered, wait) = mpsc::channel();
797        let release = Arc::new((Mutex::new(false), Condvar::new()));
798        struct Captured {
799            writer: BlockingWriter,
800            bytes: Arc<Mutex<Vec<u8>>>,
801        }
802        impl Write for Captured {
803            fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
804                let n = self.writer.write(bytes)?;
805                self.bytes.lock().unwrap().extend_from_slice(&bytes[..n]);
806                Ok(n)
807            }
808            fn flush(&mut self) -> io::Result<()> {
809                self.writer.flush()
810            }
811        }
812        struct Unblock(Arc<(Mutex<bool>, Condvar)>);
813        impl Drop for Unblock {
814            fn drop(&mut self) {
815                *self.0.0.lock().unwrap() = true;
816                self.0.1.notify_all();
817            }
818        }
819        let unblock = Unblock(Arc::clone(&release));
820        let bytes = Arc::new(Mutex::new(Vec::new()));
821        let logger = observer(
822            Captured {
823                writer: BlockingWriter {
824                    entered: Some(entered),
825                    release: Arc::clone(&release),
826                    fail_after_release: false,
827                },
828                bytes: Arc::clone(&bytes),
829            },
830            1,
831        );
832        let scope = RootDiagnosticScope::new(&view, None);
833        assert_eq!(
834            scope.ordinary(
835                &logger,
836                RootRequestEvent::Handler,
837                RootOutcomeFacts::default()
838            ),
839            DiagnosticSubmission::Enqueued
840        );
841        wait.recv_timeout(std::time::Duration::from_secs(5))
842            .unwrap();
843        assert_eq!(
844            scope.ordinary(
845                &logger,
846                RootRequestEvent::Response,
847                RootOutcomeFacts::default()
848            ),
849            DiagnosticSubmission::Enqueued
850        );
851        assert_eq!(
852            scope.ordinary(
853                &logger,
854                RootRequestEvent::Response,
855                RootOutcomeFacts::default()
856            ),
857            DiagnosticSubmission::Full
858        );
859        drop((view, root));
860        drop(unblock);
861        assert!(matches!(
862            block_on(logger.shutdown()),
863            Err(FlushError::DroppedEvents(1))
864        ));
865        let bytes = bytes.lock().unwrap();
866        let rows: Vec<serde_json::Value> = serde_json::Deserializer::from_slice(&bytes)
867            .into_iter()
868            .map(Result::unwrap)
869            .collect();
870        assert_eq!(rows.len(), 2);
871        assert!(rows.iter().all(|row| row["context"] == expected));
872    }
873
874    #[test]
875    fn root_contract_encoded_ordinary_readback_and_failure_use_original_writer() {
876        let record = serde_json::json!({"event":"request_stage", "context":{"local_request":5}});
877        let (tx, rx) = mpsc::sync_channel(1);
878        tx.try_send(Command::RootRecord {
879            bytes: serde_json::to_vec(&record).unwrap(),
880            dropped: 0,
881        })
882        .unwrap_or_else(|_| panic!("empty queue"));
883        drop(tx);
884        let mut bytes = Vec::new();
885        write_records(rx, &mut bytes);
886        assert_eq!(
887            serde_json::from_slice::<serde_json::Value>(&bytes).unwrap(),
888            record
889        );
890        let logger = observer(
891            FailingWriter {
892                fail_write: Some(1),
893                writes: 0,
894                fail_flush: false,
895            },
896            1,
897        );
898        assert_eq!(
899            logger.emit_root_record(&record),
900            crate::DiagnosticSubmission::Enqueued
901        );
902        assert!(matches!(
903            block_on(logger.shutdown()),
904            Err(FlushError::OutputFailed(OutputStage::Record))
905        ));
906    }
907
908    #[test]
909    fn random_source_failure_is_reported_during_initialization() {
910        let result = Observer::with_writer_and_seed(
911            ObserverConfig::default(),
912            io::sink(),
913            Err(InitError::RandomSource),
914        );
915        assert!(matches!(result, Err(InitError::RandomSource)));
916    }
917
918    #[test]
919    fn initialized_generator_produces_non_zero_unique_ids_without_more_io() {
920        let observer = observer(io::sink(), 8);
921        let first_trace = observer.new_trace_id();
922        let second_trace = observer.new_trace_id();
923        let first_span = observer.new_span_id();
924        let second_span = observer.new_span_id();
925        assert_ne!(first_trace.as_u128(), 0);
926        assert_ne!(first_trace, second_trace);
927        assert_ne!(first_span.as_u64(), 0);
928        assert_ne!(first_span, second_span);
929        block_on(observer.shutdown()).unwrap();
930    }
931
932    #[test]
933    fn observer_is_a_managed_component() {
934        let observer = observer(io::sink(), 8);
935        assert_eq!(ComponentLifecycle::name(&observer), "observability");
936        block_on(ComponentLifecycle::start(&observer)).unwrap();
937        block_on(ComponentLifecycle::shutdown(&observer)).unwrap();
938    }
939
940    #[test]
941    fn full_queue_is_reported_by_shutdown_without_waiting_in_emit() {
942        let (entered_sender, entered_receiver) = mpsc::channel();
943        let release = Arc::new((Mutex::new(false), Condvar::new()));
944        let writer = BlockingWriter {
945            entered: Some(entered_sender),
946            release: release.clone(),
947            fail_after_release: false,
948        };
949        let observer = observer(writer, 1);
950
951        observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
952        entered_receiver.recv().unwrap();
953        observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
954        observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
955        assert_eq!(observer.dropped_events(), 1);
956
957        let shutdown = observer.shutdown();
958        let (lock, condition) = &*release;
959        *lock.lock().unwrap() = true;
960        condition.notify_one();
961        assert_eq!(block_on(shutdown), Err(FlushError::DroppedEvents(1)));
962    }
963
964    #[test]
965    fn shutdown_reports_writer_failure_and_unreported_drops_together() {
966        let (entered_sender, entered_receiver) = mpsc::channel();
967        let release = Arc::new((Mutex::new(false), Condvar::new()));
968        let writer = BlockingWriter {
969            entered: Some(entered_sender),
970            release: release.clone(),
971            fail_after_release: true,
972        };
973        let observer = observer(writer, 1);
974
975        observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
976        entered_receiver.recv().unwrap();
977        observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
978        observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
979        let shutdown = observer.shutdown();
980
981        let (lock, condition) = &*release;
982        *lock.lock().unwrap() = true;
983        condition.notify_one();
984        assert_eq!(
985            block_on(shutdown),
986            Err(FlushError::OutputFailedAndDropped {
987                stage: OutputStage::Record,
988                dropped_events: 1,
989            })
990        );
991    }
992
993    #[test]
994    fn record_write_failure_persists_until_shutdown() {
995        let observer = observer(
996            FailingWriter {
997                fail_write: Some(1),
998                writes: 0,
999                fail_flush: false,
1000            },
1001            8,
1002        );
1003        observer.emit(LogRecord::new(EventLevel::Info, "test.record"));
1004        assert_eq!(
1005            block_on(observer.shutdown()),
1006            Err(FlushError::OutputFailed(OutputStage::Record))
1007        );
1008    }
1009
1010    #[test]
1011    fn newline_failure_persists_until_flush() {
1012        let observer = observer(
1013            FailingWriter {
1014                fail_write: Some(2),
1015                writes: 0,
1016                fail_flush: false,
1017            },
1018            8,
1019        );
1020        observer.emit(LogRecord::new(EventLevel::Info, "test.newline"));
1021        assert_eq!(
1022            block_on(observer.flush()),
1023            Err(FlushError::OutputFailed(OutputStage::Newline))
1024        );
1025        let _ = block_on(observer.shutdown());
1026    }
1027
1028    #[test]
1029    fn flush_failure_is_reported() {
1030        let observer = observer(
1031            FailingWriter {
1032                fail_write: None,
1033                writes: 0,
1034                fail_flush: true,
1035            },
1036            8,
1037        );
1038        assert_eq!(
1039            block_on(observer.flush()),
1040            Err(FlushError::OutputFailed(OutputStage::Flush))
1041        );
1042        assert_eq!(
1043            block_on(observer.shutdown()),
1044            Err(FlushError::OutputFailed(OutputStage::Flush))
1045        );
1046    }
1047}