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    process_log: Mutex<Option<saddle_admission::ProcessLogStorage>>,
187    writer_source: Arc<Mutex<WriterSourceState>>,
188    ids: IdGenerator,
189    pub(crate) metrics: crate::metrics::Metrics,
190}
191
192#[derive(Clone)]
193struct WriterSourceContext {
194    application: saddle_core::ContextLabel,
195    output: crate::SourceOutput,
196    primary: Option<saddle_core::DiagnosticOccurrence>,
197}
198#[derive(Default)]
199struct WriterSourceState {
200    context: Option<WriterSourceContext>,
201    failure: Option<WriterSourceFailure>,
202}
203struct WriterSourceFailure {
204    original: WriterStageError,
205    diagnostic: saddle_core::Diagnostic,
206    written: Option<crate::root_diagnostic::WrittenComponentFailure>,
207}
208
209#[derive(Debug)]
210struct WriterStageError {
211    stage: OutputStage,
212    original: io::Error,
213}
214impl fmt::Display for WriterStageError {
215    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
216        write!(formatter, "log writer failed during {:?}: {}", self.stage, self.original)
217    }
218}
219impl Error for WriterStageError {
220    fn source(&self) -> Option<&(dyn Error + 'static)> { Some(&self.original) }
221}
222
223fn record_writer_failure(source: &Arc<Mutex<WriterSourceState>>,
224    stage: OutputStage, original: io::Error) {
225    let context = source.lock().unwrap_or_else(|e| e.into_inner()).context.clone();
226    let original = WriterStageError { stage, original };
227    let mut diagnostic = saddle_core::Diagnostic::capture(
228        saddle_core::DiagnosticCategory::UnexpectedError,
229        saddle_core::CaptureSite::FirstObserved,
230        saddle_core::DiagnosticCause::new(
231            saddle_core::DiagnosticStage::ShutdownLogger,
232            saddle_core::DiagnosticCode::new("observability.writer_failed")
233                .expect("static writer failure code")),
234    );
235    if let Some(primary) = context.as_ref().and_then(|context| context.primary) {
236        diagnostic = diagnostic.during_cleanup_of_occurrence(&primary);
237    }
238    let written = context.as_ref().and_then(|context|
239        crate::root_diagnostic::process_component_cleanup_recorded(
240            Some(context.output.handle()),
241            saddle_core::ContextFact::Present(context.application.clone()),
242            &diagnostic, &original).ok());
243    let mut state = source.lock().unwrap_or_else(|e| e.into_inner());
244    if state.failure.is_none() {
245        state.failure = Some(WriterSourceFailure { original, diagnostic, written });
246    }
247}
248
249struct IdGenerator {
250    trace_high: u64,
251    trace_seed: u64,
252    trace_counter: AtomicU64,
253    span_seed: u64,
254    span_counter: AtomicU64,
255}
256
257impl IdGenerator {
258    fn from_seed(seed: [u8; 24]) -> Self {
259        let mut high = u64::from_be_bytes(seed[0..8].try_into().unwrap());
260        if high == 0 {
261            high = 1;
262        }
263        Self {
264            trace_high: high,
265            trace_seed: u64::from_be_bytes(seed[8..16].try_into().unwrap()),
266            trace_counter: AtomicU64::new(0),
267            span_seed: u64::from_be_bytes(seed[16..24].try_into().unwrap()),
268            span_counter: AtomicU64::new(0),
269        }
270    }
271
272    fn trace_id(&self) -> TraceId {
273        let low = self
274            .trace_seed
275            .wrapping_add(self.trace_counter.fetch_add(1, Ordering::Relaxed));
276        TraceId::from_u128((u128::from(self.trace_high) << 64) | u128::from(low))
277    }
278
279    fn span_id(&self) -> SpanId {
280        loop {
281            let value = self
282                .span_seed
283                .wrapping_add(self.span_counter.fetch_add(1, Ordering::Relaxed));
284            if value != 0 {
285                return SpanId::from_u64(value);
286            }
287        }
288    }
289}
290
291enum Command {
292    Record(LogRecord),
293    ProcessCall {
294        bytes: saddle_admission::ExactStored<Vec<u8>>,
295        dropped: u64,
296    },
297    // Source-encoded, bounded and self-contained. Never contains a request root.
298    RootRecord {
299        bytes: Vec<u8>,
300        dropped: u64,
301    },
302    Flush {
303        unreported_dropped: u64,
304        response: mpsc::Sender<WorkerStatus>,
305    },
306    Shutdown {
307        unreported_dropped: u64,
308        response: mpsc::Sender<WorkerStatus>,
309    },
310}
311pub(crate) fn root_queue_layout() -> std::alloc::Layout {
312    std::alloc::Layout::new::<Command>()
313}
314
315#[derive(Clone, Copy, Debug, Default)]
316struct WorkerStatus {
317    failure: Option<OutputStage>,
318    unreported_dropped: u64,
319}
320
321impl WorkerStatus {
322    fn into_result(self) -> Result<(), FlushError> {
323        match (self.failure, self.unreported_dropped) {
324            (None, 0) => Ok(()),
325            (Some(stage), 0) => Err(FlushError::OutputFailed(stage)),
326            (None, dropped_events) => Err(FlushError::DroppedEvents(dropped_events)),
327            (Some(stage), dropped_events) => Err(FlushError::OutputFailedAndDropped {
328                stage,
329                dropped_events,
330            }),
331        }
332    }
333}
334
335#[derive(Serialize)]
336pub(crate) struct LogRecord {
337    pub timestamp_unix_ms: u128,
338    pub level: EventLevel,
339    pub event: &'static str,
340    #[serde(skip_serializing_if = "Option::is_none")]
341    pub trace_id: Option<String>,
342    #[serde(skip_serializing_if = "Option::is_none")]
343    pub span: Option<String>,
344    #[serde(skip_serializing_if = "Option::is_none")]
345    pub span_id: Option<String>,
346    #[serde(skip_serializing_if = "Option::is_none")]
347    pub parent: Option<String>,
348    #[serde(skip_serializing_if = "Option::is_none")]
349    pub parent_span_id: Option<String>,
350    #[serde(flatten)]
351    pub data: serde_json::Map<String, serde_json::Value>,
352    #[serde(skip_serializing_if = "is_zero")]
353    pub dropped_events: u64,
354}
355
356const fn is_zero(value: &u64) -> bool {
357    *value == 0
358}
359
360impl LogRecord {
361    pub(crate) fn new(level: EventLevel, event: &'static str) -> Self {
362        Self {
363            timestamp_unix_ms: SystemTime::now()
364                .duration_since(UNIX_EPOCH)
365                .unwrap_or_default()
366                .as_millis(),
367            level,
368            event,
369            trace_id: None,
370            span: None,
371            span_id: None,
372            parent: None,
373            parent_span_id: None,
374            data: serde_json::Map::new(),
375            dropped_events: 0,
376        }
377    }
378}
379
380/// Initializes stdout logging once. Later calls return the first observer.
381pub fn init(config: ObserverConfig) -> Result<&'static Observer, InitError> {
382    if let Some(observer) = GLOBAL.get() {
383        return Ok(observer);
384    }
385
386    let observer = Observer::with_writer(config, io::stdout())?;
387    let _ = GLOBAL.set(observer);
388    Ok(GLOBAL.get().expect("global observer was initialized"))
389}
390
391/// Initializes the process logger with the sole calendar-rotated file sink.
392pub fn init_file(
393    config: ObserverConfig,
394    file: FileLoggingConfig,
395) -> Result<&'static Observer, InitError> {
396    if let Some(observer) = GLOBAL.get() {
397        return Ok(observer);
398    }
399    let writer = CalendarFileWriter::open(file).map_err(InitError::Output)?;
400    let observer = Observer::with_writer(config, writer)?;
401    let _ = GLOBAL.set(observer);
402    Ok(GLOBAL.get().expect("global observer was initialized"))
403}
404
405/// Returns the initialized process observer, if application startup installed it.
406pub fn global() -> Option<&'static Observer> {
407    GLOBAL.get()
408}
409
410impl Observer {
411    /// Bind the original process framework domain before formal routed RPCs
412    /// can emit ordinary call records. The writer never receives this handle.
413    #[doc(hidden)]
414    pub fn install_process_log_storage(&self, storage: saddle_admission::ProcessLogStorage) {
415        let mut slot = self.inner.process_log.lock().unwrap_or_else(|p| p.into_inner());
416        *slot = Some(storage);
417    }
418
419    pub(crate) fn emit_process_call(&self, record: &mut crate::call::ProcessCallRecord<'_>) {
420        struct Counter(usize);
421        impl Write for Counter {
422            fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
423                self.0 = self.0.checked_add(bytes.len()).ok_or(io::ErrorKind::OutOfMemory)?;
424                Ok(bytes.len())
425            }
426            fn flush(&mut self) -> io::Result<()> { Ok(()) }
427        }
428        self.inner.emitting.fetch_add(1, Ordering::AcqRel);
429        let submitted = (|| {
430            if !self.inner.accepting.load(Ordering::Acquire) { return false; }
431            let storage = self.inner.process_log.lock().unwrap_or_else(|p| p.into_inner()).clone();
432            let Some(storage) = storage else { return false; };
433            let dropped = self.inner.dropped.swap(0, Ordering::Relaxed);
434            record.dropped_events = dropped;
435            let mut counter = Counter(0);
436            if serde_json::to_writer(&mut counter, record).is_err() {
437                self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
438                return false;
439            }
440            let Ok(layout) = std::alloc::Layout::array::<u8>(counter.0) else {
441                self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
442                return false;
443            };
444            let Ok(permit) = storage.try_reserve(layout) else {
445                self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
446                return false;
447            };
448            let mut bytes = permit.allocate_exact(layout, || Vec::with_capacity(counter.0));
449            if serde_json::to_writer(bytes.get_mut(), record).is_err() || bytes.get().len() != counter.0 {
450                drop(bytes);
451                self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
452                return false;
453            }
454            match self.inner.sender.try_send(Command::ProcessCall { bytes, dropped }) {
455                Ok(()) => true,
456                Err(TrySendError::Full(command)) => {
457                    drop(command);
458                    self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
459                    false
460                }
461                Err(TrySendError::Disconnected(command)) => {
462                    drop(command);
463                    self.inner.metrics.logger_output(true);
464                    self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
465                    false
466                }
467            }
468        })();
469        if !submitted {
470            self.inner.metrics.logger_dropped(1);
471            self.inner.dropped.fetch_add(1, Ordering::Relaxed);
472        }
473        self.inner.emitting.fetch_sub(1, Ordering::Release);
474    }
475    /// Creates an observer with a framework-owned writer.
476    ///
477    /// Application startup should normally use [`init`]. This constructor is
478    /// useful for embedding and deterministic tests.
479    pub fn with_writer(
480        config: ObserverConfig,
481        writer: impl Write + Send + 'static,
482    ) -> Result<Self, InitError> {
483        let mut seed = [0_u8; 24];
484        getrandom::fill(&mut seed).map_err(|_| InitError::RandomSource)?;
485        Self::with_writer_and_seed(config, writer, Ok(seed))
486    }
487
488    fn with_writer_and_seed(
489        config: ObserverConfig,
490        writer: impl Write + Send + 'static,
491        seed: Result<[u8; 24], InitError>,
492    ) -> Result<Self, InitError> {
493        if config.queue_capacity == 0 {
494            return Err(InitError::EmptyQueue);
495        }
496        let ids = IdGenerator::from_seed(seed?);
497        let (sender, receiver) = mpsc::sync_channel(config.queue_capacity);
498        let writer_source = Arc::new(Mutex::new(WriterSourceState::default()));
499        let worker_source = Arc::clone(&writer_source);
500        let worker = thread::Builder::new()
501            .name("saddle-log-writer".to_owned())
502            .spawn(move || write_records(receiver, writer, worker_source))
503            .map_err(InitError::Spawn)?;
504        Ok(Self {
505            inner: Arc::new(Inner {
506                sender,
507                dropped: AtomicU64::new(0),
508                accepting: AtomicBool::new(true),
509                emitting: AtomicUsize::new(0),
510                shutdown_started: AtomicBool::new(false),
511                worker: Mutex::new(Some(worker)),
512                process_log: Mutex::new(None),
513                writer_source,
514                ids,
515                metrics: crate::metrics::Metrics::default(),
516            }),
517        })
518    }
519
520    pub(crate) fn new_trace_id(&self) -> TraceId {
521        self.inner.ids.trace_id()
522    }
523
524    pub(crate) fn new_span_id(&self) -> SpanId {
525        self.inner.ids.span_id()
526    }
527
528    pub(crate) fn emit(&self, mut record: LogRecord) {
529        if !self.inner.accepting.load(Ordering::Acquire) {
530            return;
531        }
532        self.inner.emitting.fetch_add(1, Ordering::AcqRel);
533        if !self.inner.accepting.load(Ordering::Acquire) {
534            self.inner.emitting.fetch_sub(1, Ordering::AcqRel);
535            return;
536        }
537
538        record.dropped_events = self.inner.dropped.swap(0, Ordering::Relaxed);
539        match self.inner.sender.try_send(Command::Record(record)) {
540            Ok(()) => {}
541            Err(TrySendError::Full(Command::Record(record))) => {
542                self.inner.metrics.logger_dropped(1);
543                self.inner
544                    .dropped
545                    .fetch_add(record.dropped_events.saturating_add(1), Ordering::Relaxed);
546            }
547            Err(TrySendError::Disconnected(_)) => {
548                self.inner.metrics.logger_dropped(1);
549                self.inner.metrics.logger_output(true);
550                self.inner.dropped.fetch_add(1, Ordering::Relaxed);
551            }
552            Err(TrySendError::Full(_)) => unreachable!("emit only sends records"),
553        }
554        self.inner.emitting.fetch_sub(1, Ordering::Release);
555    }
556
557    pub(crate) fn emit_root_record<T: Serialize>(&self, record: &T) -> crate::DiagnosticSubmission {
558        use crate::DiagnosticSubmission as Submission;
559        self.inner.emitting.fetch_add(1, Ordering::AcqRel);
560        let result = (|| {
561            if !self.inner.accepting.load(Ordering::Acquire) {
562                return Submission::Closed;
563            }
564            let mut frame = crate::diagnostic::FixedDiagnosticBytes {
565                bytes: [0; 8192],
566                len: 0,
567            };
568            if serde_json::to_writer(&mut frame, record).is_err() {
569                return Submission::EncodingFailed;
570            }
571            let dropped = self.inner.dropped.swap(0, Ordering::Relaxed);
572            // Only initialized bytes move into the existing logger domain. This
573            // allocation and source frame peak are separate R0 layout inputs.
574            let bytes = frame.bytes[..frame.len].to_vec();
575            match self
576                .inner
577                .sender
578                .try_send(Command::RootRecord { bytes, dropped })
579            {
580                Ok(()) => Submission::Enqueued,
581                Err(TrySendError::Full(Command::RootRecord { dropped, .. })) => {
582                    self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
583                    Submission::Full
584                }
585                Err(TrySendError::Disconnected(_)) => {
586                    self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
587                    self.inner.metrics.logger_output(true);
588                    Submission::Closed
589                }
590                Err(TrySendError::Full(_)) => unreachable!("root record submission"),
591            }
592        })();
593        if result != Submission::Enqueued {
594            self.inner.dropped.fetch_add(1, Ordering::Relaxed);
595            self.inner.metrics.logger_dropped(1);
596        }
597        self.inner.emitting.fetch_sub(1, Ordering::Release);
598        result
599    }
600
601    /// Returns a fixed-size, allocation-free copy of the current process metrics.
602    pub fn metrics_snapshot(&self) -> crate::MetricsSnapshot {
603        self.inner.metrics.snapshot()
604    }
605
606    /// Flushes records accepted before the lifecycle coordinator reaches the
607    /// writer. The returned future never blocks the async executor thread.
608    pub fn flush(&self) -> FlushFuture {
609        if !self.inner.accepting.load(Ordering::Acquire) {
610            return FlushFuture {
611                completion: Completion::ready(Err(FlushError::WriterStopped)),
612            };
613        }
614        let completion = Completion::pending();
615        let future = FlushFuture {
616            completion: completion.clone(),
617        };
618        let inner = self.inner.clone();
619        if thread::Builder::new()
620            .name("saddle-log-flush".to_owned())
621            .spawn(move || coordinate_flush(inner, completion.clone()))
622            .is_err()
623        {
624            future
625                .completion
626                .complete(Err(FlushError::CoordinatorUnavailable));
627        }
628        future
629    }
630
631    /// Stops admission, drains accepted records, flushes output and joins the
632    /// writer without blocking the async executor thread.
633    pub fn shutdown(&self) -> FlushFuture {
634        if self.inner.shutdown_started.swap(true, Ordering::AcqRel) {
635            return FlushFuture {
636                completion: Completion::ready(Err(FlushError::AlreadyShuttingDown)),
637            };
638        }
639        self.inner.accepting.store(false, Ordering::Release);
640        let completion = Completion::pending();
641        let future = FlushFuture {
642            completion: completion.clone(),
643        };
644        let inner = self.inner.clone();
645        if thread::Builder::new()
646            .name("saddle-log-shutdown".to_owned())
647            .spawn(move || coordinate_shutdown(inner, completion.clone()))
648            .is_err()
649        {
650            self.inner.accepting.store(true, Ordering::Release);
651            self.inner.shutdown_started.store(false, Ordering::Release);
652            future
653                .completion
654                .complete(Err(FlushError::CoordinatorUnavailable));
655        }
656        future
657    }
658
659    pub fn dropped_events(&self) -> u64 {
660        self.inner.dropped.load(Ordering::Relaxed)
661    }
662}
663
664impl ComponentLifecycle for Observer {
665    fn name(&self) -> &'static str {
666        "observability"
667    }
668
669    fn start(&self) -> LifecycleFuture<'_> {
670        Box::pin(async { Ok(()) })
671    }
672
673    fn shutdown(&self) -> LifecycleFuture<'_> {
674        let shutdown = Observer::shutdown(self);
675        Box::pin(async move {
676            shutdown.await.map_err(|_| {
677                SaddleError::new(
678                    ErrorKind::Infrastructure,
679                    "observability.shutdown_failed",
680                    "structured log shutdown failed",
681                )
682            })
683        })
684    }
685}
686
687impl Observer {
688    fn recorded_start<'a>(&'a self, application: &'a saddle_core::ContextLabel,
689        output: &'a crate::SourceOutput)
690        -> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
691        self.inner.writer_source.lock().unwrap_or_else(|e| e.into_inner()).context =
692            Some(WriterSourceContext {
693                application: application.clone(), output: output.clone(), primary: None,
694            });
695        Box::pin(async { Ok(()) })
696    }
697
698    fn recorded_shutdown<'a>(&'a self,
699        application: &'a saddle_core::ContextLabel,
700        output: &'a crate::SourceOutput,
701        primary: Option<saddle_core::DiagnosticOccurrence>)
702        -> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
703        self.inner.writer_source.lock().unwrap_or_else(|e| e.into_inner()).context =
704            Some(WriterSourceContext {
705                application: application.clone(), output: output.clone(), primary,
706            });
707        Box::pin(async move {
708            let result = Observer::shutdown(self).await;
709            let failure = self.inner.writer_source.lock()
710                .unwrap_or_else(|e| e.into_inner()).failure.take();
711            if let Some(failure) = failure {
712                return Err(crate::root_diagnostic::RecordedSaddleError::previously_written_writer_error(
713                    failure.original, failure.diagnostic, failure.written));
714            }
715            crate::root_diagnostic::RecordedSaddleError::component_result(
716                result, application, output,
717                crate::root_diagnostic::ComponentSourceKind::Cleanup, primary,
718                ErrorKind::Infrastructure, "observability.shutdown_failed",
719                "structured log shutdown failed",
720            )
721        })
722    }
723}
724
725/// The formal component boundary keeps the actual FlushError until the
726/// controlled source file has acknowledged it. The Core lifecycle impl above
727/// remains a legacy compatibility surface for callers migrating to this one.
728impl crate::root_diagnostic::RecordedComponentLifecycle for Observer {
729    fn name(&self) -> &'static str { "observability" }
730
731    fn start<'a>(&'a self, application: &'a saddle_core::ContextLabel,
732        output: &'a crate::SourceOutput)
733        -> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
734        Observer::recorded_start(self, application, output)
735    }
736
737    fn shutdown<'a>(&'a self, application: &'a saddle_core::ContextLabel,
738        output: &'a crate::SourceOutput)
739        -> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
740        crate::root_diagnostic::RecordedComponentLifecycle::shutdown_with_primary(
741            self, application, output, None)
742    }
743
744    fn shutdown_with_primary<'a>(&'a self,
745        application: &'a saddle_core::ContextLabel,
746        output: &'a crate::SourceOutput,
747        primary: Option<saddle_core::DiagnosticOccurrence>)
748        -> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
749        Observer::recorded_shutdown(self, application, output, primary)
750    }
751}
752
753fn coordinate_flush(inner: Arc<Inner>, completion: Arc<Completion>) {
754    let dropped = inner.dropped.swap(0, Ordering::AcqRel);
755    let (response, receiver) = mpsc::channel();
756    let result = if inner
757        .sender
758        .send(Command::Flush {
759            unreported_dropped: dropped,
760            response,
761        })
762        .is_err()
763    {
764        Err(FlushError::WriterStopped)
765    } else {
766        receiver
767            .recv()
768            .map_err(|_| FlushError::WriterStopped)
769            .and_then(WorkerStatus::into_result)
770    };
771    completion.complete(result);
772}
773
774fn coordinate_shutdown(inner: Arc<Inner>, completion: Arc<Completion>) {
775    while inner.emitting.load(Ordering::Acquire) != 0 {
776        thread::yield_now();
777    }
778    let dropped = inner.dropped.swap(0, Ordering::AcqRel);
779    let (response, receiver) = mpsc::channel();
780    let mut result = if inner
781        .sender
782        .send(Command::Shutdown {
783            unreported_dropped: dropped,
784            response,
785        })
786        .is_err()
787    {
788        Err(FlushError::WriterStopped)
789    } else {
790        receiver
791            .recv()
792            .map_err(|_| FlushError::WriterStopped)
793            .and_then(WorkerStatus::into_result)
794    };
795
796    if let Some(worker) = inner.worker.lock().unwrap().take() {
797        if worker.join().is_err() {
798            result = Err(FlushError::WorkerPanicked);
799        }
800    }
801    inner.process_log.lock().unwrap_or_else(|p| p.into_inner()).take();
802    completion.complete(result);
803}
804
805fn write_records(receiver: Receiver<Command>, mut writer: impl Write,
806    source: Arc<Mutex<WriterSourceState>>) {
807    let mut status = WorkerStatus::default();
808    while let Ok(command) = receiver.recv() {
809        match command {
810            Command::Record(record) => write_record(&mut writer, record, &mut status, &source),
811            Command::ProcessCall { bytes, dropped } => {
812                status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
813                if status.failure.is_none() {
814                    if let Err(error) = writer.write_all(bytes.get()) {
815                        record_writer_failure(&source, OutputStage::Record, error);
816                        status.failure = Some(OutputStage::Record);
817                    } else if let Err(error) = writer.write_all(b"\n") {
818                        record_writer_failure(&source, OutputStage::Newline, error);
819                        status.failure = Some(OutputStage::Newline);
820                    }
821                }
822                drop(bytes);
823            }
824            Command::RootRecord { bytes, dropped } => {
825                status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
826                if status.failure.is_none() {
827                    if let Err(error) = writer.write_all(&bytes) {
828                        record_writer_failure(&source, OutputStage::Record, error);
829                        status.failure = Some(OutputStage::Record);
830                    } else if let Err(error) = writer.write_all(b"\n") {
831                        record_writer_failure(&source, OutputStage::Newline, error);
832                        status.failure = Some(OutputStage::Newline);
833                    }
834                }
835            }
836            Command::Flush {
837                unreported_dropped,
838                response,
839            } => {
840                status.unreported_dropped =
841                    status.unreported_dropped.saturating_add(unreported_dropped);
842                flush_writer(&mut writer, &mut status, &source);
843                let _ = response.send(status);
844            }
845            Command::Shutdown {
846                unreported_dropped,
847                response,
848            } => {
849                status.unreported_dropped =
850                    status.unreported_dropped.saturating_add(unreported_dropped);
851                flush_writer(&mut writer, &mut status, &source);
852                let _ = response.send(status);
853                break;
854            }
855        }
856    }
857}
858
859fn write_record(writer: &mut impl Write, record: LogRecord, status: &mut WorkerStatus,
860    source: &Arc<Mutex<WriterSourceState>>) {
861    if status.failure.is_some() {
862        status.unreported_dropped = status
863            .unreported_dropped
864            .saturating_add(record.dropped_events);
865        return;
866    }
867
868    let dropped = record.dropped_events;
869    let bytes = match serde_json::to_vec(&record) {
870        Ok(bytes) => bytes,
871        Err(_) => {
872            status.failure = Some(OutputStage::Serialize);
873            status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
874            return;
875        }
876    };
877    if let Err(error) = writer.write_all(&bytes) {
878        record_writer_failure(source, OutputStage::Record, error);
879        status.failure = Some(OutputStage::Record);
880        status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
881        return;
882    }
883    if let Err(error) = writer.write_all(b"\n") {
884        record_writer_failure(source, OutputStage::Newline, error);
885        status.failure = Some(OutputStage::Newline);
886        status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
887    }
888}
889
890fn flush_writer(writer: &mut impl Write, status: &mut WorkerStatus,
891    source: &Arc<Mutex<WriterSourceState>>) {
892    // A previous output failure has already made this writer unhealthy.
893    // Avoid issuing another fallible operation whose error could not be the
894    // primary result and would otherwise be discarded.
895    if status.failure.is_some() { return; }
896    if let Err(error) = writer.flush() {
897        record_writer_failure(source, OutputStage::Flush, error);
898        status.failure = Some(OutputStage::Flush);
899    }
900}
901
902#[cfg(test)]
903mod tests {
904    use std::{
905        sync::{Condvar, mpsc},
906        task::{Wake, Waker},
907    };
908
909    use super::*;
910
911    #[test]
912    fn recorded_observer_shutdown_preserves_actual_flush_error() {
913        let directory = std::env::temp_dir().join(format!(
914            "saddle-observer-recorded-{}-{}", std::process::id(),
915            SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos()));
916        std::fs::create_dir(&directory).unwrap();
917        let mut output = crate::EmergencyDiagnostics::start_checked(
918            &crate::FileLoggingConfig::new(directory.clone(), crate::Rotation::Daily)).unwrap();
919        let selected = output.source_output().unwrap();
920        let target = output.target().to_owned();
921        let logger = observer(io::sink(), 1);
922        block_on(Observer::shutdown(&logger)).unwrap();
923        let app = saddle_core::ContextLabel::checked("recorded-observer-test").unwrap();
924        let failure = block_on(
925            crate::root_diagnostic::RecordedComponentLifecycle::shutdown(
926                &logger, &app, &selected)).err().unwrap();
927        assert_eq!(failure.safe().code(), "observability.shutdown_failed");
928        assert!(failure.original_if_unconfirmed().is_none());
929        let record = std::fs::read_to_string(&target).unwrap();
930        assert!(record.contains("AlreadyShuttingDown"));
931        assert!(record.contains(&failure.safe().diagnostic().unwrap().id().to_string()));
932        drop(selected);
933        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
934        while output.shutdown() == crate::DiagnosticShutdown::Pending
935            && std::time::Instant::now() < deadline { std::thread::yield_now(); }
936        assert_eq!(output.shutdown(), crate::DiagnosticShutdown::Finished);
937        drop(output);
938        std::fs::remove_dir_all(directory).unwrap();
939    }
940
941    #[derive(Debug)]
942    struct WriterChain {
943        message: &'static str,
944        cause: io::Error,
945    }
946    impl fmt::Display for WriterChain {
947        fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
948            f.write_str(self.message)
949        }
950    }
951    impl Error for WriterChain {
952        fn source(&self) -> Option<&(dyn Error + 'static)> { Some(&self.cause) }
953    }
954    struct ChainWriter {
955        fail_flush: bool,
956        release_write: Option<(mpsc::Sender<()>, Arc<(Mutex<bool>, Condvar)>)>,
957    }
958    impl Write for ChainWriter {
959        fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
960            if self.fail_flush { Ok(bytes.len()) } else {
961                if let Some((entered, release)) = self.release_write.take() {
962                    entered.send(()).unwrap();
963                    let (lock, condition) = &*release;
964                    let mut ready = lock.lock().unwrap();
965                    while !*ready { ready = condition.wait(ready).unwrap(); }
966                }
967                Err(io::Error::other(WriterChain {
968                    message: "writer top original 4931",
969                    cause: io::Error::other("writer nested original 4932"),
970                }))
971            }
972        }
973        fn flush(&mut self) -> io::Result<()> {
974            if self.fail_flush {
975                Err(io::Error::other(WriterChain {
976                    message: "flush top original 4933",
977                    cause: io::Error::other("flush nested original 4934"),
978                }))
979            } else { Ok(()) }
980        }
981    }
982
983    #[test]
984    fn writer_write_and_flush_originals_precede_formal_shutdown_completion() {
985        use crate::root_diagnostic::RecordedComponentLifecycle;
986        for flush in [false, true] {
987            let directory = std::env::temp_dir().join(format!(
988                "saddle-writer-original-{flush}-{}-{}", std::process::id(),
989                SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos()));
990            std::fs::create_dir(&directory).unwrap();
991            let mut emergency = crate::EmergencyDiagnostics::start_checked(
992                &crate::FileLoggingConfig::new(directory.clone(), crate::Rotation::Daily)).unwrap();
993            let output = emergency.source_output().unwrap();
994            let target = emergency.target().to_owned();
995            let (entered_sender, entered_receiver) = mpsc::channel();
996            let release = Arc::new((Mutex::new(false), Condvar::new()));
997            let logger = observer(ChainWriter {
998                fail_flush: flush,
999                release_write: (!flush).then(|| (entered_sender, Arc::clone(&release))),
1000            }, 8);
1001            let application = saddle_core::ContextLabel::checked("writer-original-test").unwrap();
1002            assert!(block_on(RecordedComponentLifecycle::start(
1003                &logger, &application, &output)).is_ok());
1004            let primary = saddle_core::Diagnostic::capture(
1005                saddle_core::DiagnosticCategory::UnexpectedError,
1006                saddle_core::CaptureSite::Origin,
1007                saddle_core::DiagnosticCause::new(
1008                    saddle_core::DiagnosticStage::ShutdownComponent,
1009                    saddle_core::DiagnosticCode::new("test.writer_primary").unwrap()),
1010            ).occurrence();
1011            if !flush { logger.emit(LogRecord::new(EventLevel::Info, "writer.failure")); }
1012            if !flush { entered_receiver.recv_timeout(std::time::Duration::from_secs(5)).unwrap(); }
1013            let shutdown = RecordedComponentLifecycle::shutdown_with_primary(
1014                &logger, &application, &output, Some(primary));
1015            if !flush {
1016                let (lock, condition) = &*release;
1017                *lock.lock().unwrap() = true;
1018                condition.notify_all();
1019            }
1020            let failure = block_on(shutdown).err().unwrap();
1021            assert_eq!(failure.safe().code(), "observability.shutdown_failed");
1022            assert!(failure.original_if_unconfirmed().is_none());
1023            let diagnostic = failure.safe().diagnostic().unwrap();
1024            let record = std::fs::read_to_string(&target).unwrap();
1025            let (top, nested) = if flush {
1026                ("flush top original 4933", "flush nested original 4934")
1027            } else { ("writer top original 4931", "writer nested original 4932") };
1028            assert!(record.contains(top), "{record}");
1029            assert!(record.contains(nested), "{record}");
1030            assert!(record.contains(if flush { "Flush" } else { "Record" }), "{record}");
1031            assert!(record.contains("shutdown_logger"), "{record}");
1032            assert!(record.contains(&diagnostic.id().to_string()));
1033            assert!(record.contains("complete"), "{record}");
1034            let originals: Vec<serde_json::Value> = record.lines()
1035                .map(|line| serde_json::from_str(line).unwrap())
1036                .filter(|row: &serde_json::Value|
1037                    row["event"] == "request_error_original").collect();
1038            assert!(!originals.is_empty());
1039            assert!(originals.iter().any(|row|
1040                row["channel"] == "description" && row["cause_depth"] == 0
1041                    && row["payload"].as_str().unwrap().contains(top)));
1042            assert!(originals.iter().any(|row|
1043                row["channel"] == "debug" && row["cause_depth"] == 0
1044                    && row["payload"].as_str().unwrap().contains(top)));
1045            assert!(originals.iter().any(|row|
1046                row["channel"] == "description"
1047                    && row["cause_depth"].as_u64().unwrap() > 0
1048                    && row["payload"].as_str().unwrap().contains(nested)));
1049            assert!(originals.iter().any(|row|
1050                row["channel"] == "terminal"
1051                    && row["state"] == "exposed_chain_complete"));
1052            assert!(originals.iter().all(|row|
1053                row["occurrence"]["diagnostic_id"] == diagnostic.id()));
1054            let expected = serde_json::to_value(primary).unwrap();
1055            assert!(originals.iter().all(|row|
1056                row["occurrence"]["primary_diagnostic_id"]
1057                    == expected["diagnostic_id"]));
1058            drop(logger);
1059            drop(output);
1060            let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
1061            while emergency.shutdown() == crate::DiagnosticShutdown::Pending
1062                && std::time::Instant::now() < deadline { std::thread::yield_now(); }
1063            assert_eq!(emergency.shutdown(), crate::DiagnosticShutdown::Finished);
1064            drop(emergency);
1065            std::fs::remove_dir_all(directory).unwrap();
1066        }
1067    }
1068
1069    #[test]
1070    fn writer_original_without_selected_output_remains_unconfirmed() {
1071        use crate::root_diagnostic::RecordedComponentLifecycle;
1072        let logger = observer(ChainWriter { fail_flush: false, release_write: None }, 8);
1073        logger.emit(LogRecord::new(EventLevel::Info, "writer.no_output"));
1074        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
1075        while logger.inner.writer_source.lock().unwrap_or_else(|e| e.into_inner())
1076            .failure.is_none() && std::time::Instant::now() < deadline {
1077            std::thread::yield_now();
1078        }
1079        assert!(logger.inner.writer_source.lock().unwrap_or_else(|e| e.into_inner())
1080            .failure.as_ref().is_some_and(|failure| failure.written.is_none()));
1081        let directory = std::env::temp_dir().join(format!(
1082            "saddle-writer-unconfirmed-{}-{}", std::process::id(),
1083            SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos()));
1084        std::fs::create_dir(&directory).unwrap();
1085        let mut emergency = crate::EmergencyDiagnostics::start_checked(
1086            &crate::FileLoggingConfig::new(directory.clone(), crate::Rotation::Daily)).unwrap();
1087        let target = emergency.target().to_owned();
1088        let output = emergency.source_output().unwrap();
1089        let application = saddle_core::ContextLabel::checked("writer-unconfirmed-test").unwrap();
1090        let failure = block_on(RecordedComponentLifecycle::shutdown(
1091            &logger, &application, &output)).err().unwrap();
1092        assert_eq!(failure.safe().code(), "observability.shutdown_failed");
1093        let original = failure.original_if_unconfirmed().unwrap();
1094        assert!(original.to_string().contains("writer top original 4931"));
1095        assert!(original.source().unwrap().to_string()
1096            .contains("writer top original 4931"));
1097        assert!(!std::fs::read_to_string(&target).unwrap()
1098            .contains("writer top original 4931"));
1099        drop(logger);
1100        drop(output);
1101        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
1102        while emergency.shutdown() == crate::DiagnosticShutdown::Pending
1103            && std::time::Instant::now() < deadline { std::thread::yield_now(); }
1104        assert_eq!(emergency.shutdown(), crate::DiagnosticShutdown::Finished);
1105        drop(emergency);
1106        std::fs::remove_dir_all(directory).unwrap();
1107    }
1108
1109    struct ThreadWaker(thread::Thread);
1110
1111    impl Wake for ThreadWaker {
1112        fn wake(self: Arc<Self>) {
1113            self.0.unpark();
1114        }
1115    }
1116
1117    fn block_on<T>(future: impl Future<Output = T>) -> T {
1118        let waker = Waker::from(Arc::new(ThreadWaker(thread::current())));
1119        let mut context = Context::from_waker(&waker);
1120        let mut future = std::pin::pin!(future);
1121        loop {
1122            match future.as_mut().poll(&mut context) {
1123                Poll::Ready(output) => return output,
1124                Poll::Pending => thread::park(),
1125            }
1126        }
1127    }
1128
1129    struct BlockingWriter {
1130        entered: Option<mpsc::Sender<()>>,
1131        release: Arc<(Mutex<bool>, Condvar)>,
1132        fail_after_release: bool,
1133    }
1134
1135    impl Write for BlockingWriter {
1136        fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
1137            if let Some(entered) = self.entered.take() {
1138                let _ = entered.send(());
1139                let (lock, condition) = &*self.release;
1140                let mut released = lock.lock().unwrap();
1141                while !*released {
1142                    released = condition.wait(released).unwrap();
1143                }
1144                if self.fail_after_release {
1145                    return Err(io::Error::other("injected writer failure"));
1146                }
1147            }
1148            Ok(bytes.len())
1149        }
1150
1151        fn flush(&mut self) -> io::Result<()> {
1152            Ok(())
1153        }
1154    }
1155
1156    struct FailingWriter {
1157        fail_write: Option<usize>,
1158        writes: usize,
1159        fail_flush: bool,
1160    }
1161
1162    impl Write for FailingWriter {
1163        fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
1164            self.writes += 1;
1165            if self.fail_write == Some(self.writes) {
1166                Err(io::Error::other("injected writer failure"))
1167            } else {
1168                Ok(bytes.len())
1169            }
1170        }
1171
1172        fn flush(&mut self) -> io::Result<()> {
1173            if self.fail_flush {
1174                Err(io::Error::other("injected flush failure"))
1175            } else {
1176                Ok(())
1177            }
1178        }
1179    }
1180
1181    fn observer(writer: impl Write + Send + 'static, capacity: usize) -> Observer {
1182        Observer::with_writer_and_seed(
1183            ObserverConfig {
1184                queue_capacity: capacity,
1185            },
1186            writer,
1187            Ok([7; 24]),
1188        )
1189        .unwrap()
1190    }
1191
1192    #[test]
1193    fn managed_call_log_outlives_request_account_while_writer_is_blocked() {
1194        // The new process-owned variant adds no command-ring slot bytes:
1195        // the unchanged ordinary LogRecord still determines Command's size.
1196        assert_eq!(std::mem::size_of::<Command>(), std::mem::size_of::<LogRecord>());
1197        assert_eq!(std::mem::size_of::<Command>(), 192);
1198        assert_eq!(ObserverConfig::default().queue_capacity * (8 + std::mem::size_of::<Command>()),
1199            1_638_400);
1200        use saddle_admission::{
1201            ProfuseGwLightweightObservedAdmissionOutcome, StorageDemand,
1202            bind_deployment_resource_budget_bootstrap, freeze_deployment_resource_budget,
1203            prepare_profusegw_lightweight_profile,
1204        };
1205        use saddle_core::{
1206            BootstrapRendezvousIssuer, GeneratedApplicationFreezeSource,
1207            ListenerStartupFreezeSource, RpcCorrelationId, pair_bootstrap_rendezvous,
1208        };
1209        let pending = freeze_deployment_resource_budget(
1210            1, 32, 5_000, 1, 1, 1_000_000, 1_000_000, 1_000_000, 1_000_000,
1211        ).unwrap().reserve_known_database_memory(0, 0, 0).ok().unwrap();
1212        let (application, listener) = BootstrapRendezvousIssuer::issue()
1213            .freeze_application(GeneratedApplicationFreezeSource::new(
1214                "app", b"descriptor", &["route"],
1215            )).unwrap();
1216        let listener = listener.freeze_listener(ListenerStartupFreezeSource::new(
1217            "app", "127.0.0.1:8000".parse().unwrap(),
1218            "127.0.0.1:9000".parse().unwrap(),
1219            std::time::Duration::from_millis(5_000),
1220        )).ok().unwrap();
1221        let (whole, receipt) = pair_bootstrap_rendezvous(application, listener).ok().unwrap();
1222        let budget = bind_deployment_resource_budget_bootstrap(pending, whole, receipt).ok().unwrap();
1223        let process = prepare_profusegw_lightweight_profile(budget).ok().unwrap();
1224        let root = process.try_process_storage(StorageDemand::separate(&[(
1225            std::alloc::Layout::array::<u8>(16).unwrap(), 1,
1226        )]).unwrap()).unwrap();
1227        let baseline = process.resource_snapshot().framework_charged;
1228        let (entered, entered_rx) = mpsc::channel();
1229        let release = Arc::new((Mutex::new(false), Condvar::new()));
1230        let logger = observer(BlockingWriter {
1231            entered: Some(entered), release: Arc::clone(&release), fail_after_release: false,
1232        }, 2);
1233        logger.install_process_log_storage(process.process_log_storage());
1234        let (call, _) = logger.start_managed_external_call_with_rpc(
1235            "app", "zone", "interface", "operation", None,
1236            RpcCorrelationId::new("0").unwrap(),
1237        ).unwrap();
1238        entered_rx.recv_timeout(std::time::Duration::from_secs(2)).unwrap();
1239        let admitted = match process.verified_profile().try_admit_observed() {
1240            ProfuseGwLightweightObservedAdmissionOutcome::Ready(permit, _) => permit,
1241            _ => panic!("fixture admission failed"),
1242        };
1243        let original = saddle_core::SaddleError::new(
1244            saddle_core::ErrorKind::Internal, "original.failure", "original",
1245        );
1246        call.fail(&original);
1247        assert_eq!(original.code(), "original.failure");
1248        admitted.into_execution().cancel_observed();
1249        assert_eq!(process.resource_snapshot().active_accounts, 0);
1250        assert!(process.resource_snapshot().framework_charged > baseline,
1251            "queued process log must remain charged after request cleanup");
1252        // The worker holds the first record while the second is queued.
1253        // Additional managed call records must be dropped synchronously.
1254        let (full, _) = logger.start_managed_external_call_with_rpc(
1255            "app", "zone", "interface", "full", None,
1256            RpcCorrelationId::new("1").unwrap(),
1257        ).unwrap();
1258        full.fail(&original);
1259        assert!(logger.dropped_events() > 0);
1260        *release.0.lock().unwrap() = true;
1261        release.1.notify_all();
1262        assert!(matches!(block_on(logger.flush()), Err(FlushError::DroppedEvents(_))));
1263        assert_eq!(process.resource_snapshot().framework_charged, baseline);
1264
1265        // Exhaust the *same* framework domain with a physical child owner;
1266        // encoding must fail to reserve before allocating and remain lossy.
1267        let available = process.resource_snapshot().framework_capacity
1268            - process.resource_snapshot().framework_charged;
1269        let filler = process.try_test_process_storage_child(StorageDemand::separate(&[(
1270            std::alloc::Layout::array::<u8>(
1271                available - 64 - std::mem::size_of::<saddle_admission::StoragePermit>(),
1272            ).unwrap(), 1,
1273        )]).unwrap()).unwrap();
1274        let charged = process.resource_snapshot().framework_charged;
1275        let before_drop = logger.dropped_events();
1276        let (short, _) = logger.start_managed_external_call_with_rpc(
1277            "app", "zone", "interface", "short", None,
1278            RpcCorrelationId::new("2").unwrap(),
1279        ).unwrap();
1280        short.fail(&original);
1281        assert!(logger.dropped_events() >= before_drop + 2);
1282        assert_eq!(process.resource_snapshot().framework_charged, charged);
1283        drop(filler);
1284        assert!(matches!(block_on(logger.shutdown()), Err(FlushError::DroppedEvents(_))));
1285        let after_shutdown = process.resource_snapshot().framework_charged;
1286        let (closed, _) = logger.start_managed_external_call_with_rpc(
1287            "app", "zone", "interface", "closed", None,
1288            RpcCorrelationId::new("3").unwrap(),
1289        ).unwrap();
1290        closed.fail(&original);
1291        assert_eq!(process.resource_snapshot().framework_charged, after_shutdown);
1292        assert_eq!(process.resource_snapshot().framework_charged, baseline);
1293        drop(logger);
1294        drop(root);
1295        process.finish().unwrap();
1296    }
1297
1298    #[test]
1299    fn root_contract_ordinary_uses_same_view_and_bounded_queue() {
1300        use crate::{
1301            DiagnosticSubmission, RootDiagnosticScope, RootOutcomeFacts, RootRequestEvent,
1302        };
1303        use saddle_core::{
1304            ContextFact, ContextLabel, RequestLocalFacts, RequestRootPublisher, RequestViewPhase,
1305        };
1306        let root = RequestRootPublisher::create(
1307            ContextLabel::checked("app").unwrap(),
1308            ContextFact::NotEstablished,
1309        )
1310        .unwrap();
1311        let view = root
1312            .reference()
1313            .view(RequestLocalFacts::new(RequestViewPhase::Handler));
1314        let expected = serde_json::to_value(&view).unwrap();
1315        let mut frame = crate::diagnostic::FixedDiagnosticBytes {
1316            bytes: [0; 8192],
1317            len: 0,
1318        };
1319        serde_json::to_writer(&mut frame, &view).unwrap();
1320        let ordinary_context: serde_json::Value =
1321            serde_json::from_slice(&frame.bytes[..frame.len]).unwrap();
1322        assert_eq!(ordinary_context, expected);
1323        let (entered, wait) = mpsc::channel();
1324        let release = Arc::new((Mutex::new(false), Condvar::new()));
1325        struct Captured {
1326            writer: BlockingWriter,
1327            bytes: Arc<Mutex<Vec<u8>>>,
1328        }
1329        impl Write for Captured {
1330            fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
1331                let n = self.writer.write(bytes)?;
1332                self.bytes.lock().unwrap().extend_from_slice(&bytes[..n]);
1333                Ok(n)
1334            }
1335            fn flush(&mut self) -> io::Result<()> {
1336                self.writer.flush()
1337            }
1338        }
1339        struct Unblock(Arc<(Mutex<bool>, Condvar)>);
1340        impl Drop for Unblock {
1341            fn drop(&mut self) {
1342                *self.0.0.lock().unwrap() = true;
1343                self.0.1.notify_all();
1344            }
1345        }
1346        let unblock = Unblock(Arc::clone(&release));
1347        let bytes = Arc::new(Mutex::new(Vec::new()));
1348        let logger = observer(
1349            Captured {
1350                writer: BlockingWriter {
1351                    entered: Some(entered),
1352                    release: Arc::clone(&release),
1353                    fail_after_release: false,
1354                },
1355                bytes: Arc::clone(&bytes),
1356            },
1357            1,
1358        );
1359        let scope = RootDiagnosticScope::new(&view, None);
1360        assert_eq!(
1361            scope.ordinary(
1362                &logger,
1363                RootRequestEvent::Handler,
1364                RootOutcomeFacts::default()
1365            ),
1366            DiagnosticSubmission::Enqueued
1367        );
1368        wait.recv_timeout(std::time::Duration::from_secs(5))
1369            .unwrap();
1370        assert_eq!(
1371            scope.ordinary(
1372                &logger,
1373                RootRequestEvent::Response,
1374                RootOutcomeFacts::default()
1375            ),
1376            DiagnosticSubmission::Enqueued
1377        );
1378        assert_eq!(
1379            scope.ordinary(
1380                &logger,
1381                RootRequestEvent::Response,
1382                RootOutcomeFacts::default()
1383            ),
1384            DiagnosticSubmission::Full
1385        );
1386        drop((view, root));
1387        drop(unblock);
1388        assert!(matches!(
1389            block_on(logger.shutdown()),
1390            Err(FlushError::DroppedEvents(1))
1391        ));
1392        let bytes = bytes.lock().unwrap();
1393        let rows: Vec<serde_json::Value> = serde_json::Deserializer::from_slice(&bytes)
1394            .into_iter()
1395            .map(Result::unwrap)
1396            .collect();
1397        assert_eq!(rows.len(), 2);
1398        assert!(rows.iter().all(|row| row["context"] == expected));
1399    }
1400
1401    #[test]
1402    fn root_contract_encoded_ordinary_readback_and_failure_use_original_writer() {
1403        let record = serde_json::json!({"event":"request_stage", "context":{"local_request":5}});
1404        let (tx, rx) = mpsc::sync_channel(1);
1405        tx.try_send(Command::RootRecord {
1406            bytes: serde_json::to_vec(&record).unwrap(),
1407            dropped: 0,
1408        })
1409        .unwrap_or_else(|_| panic!("empty queue"));
1410        drop(tx);
1411        let mut bytes = Vec::new();
1412        write_records(rx, &mut bytes, Arc::new(Mutex::new(WriterSourceState::default())));
1413        assert_eq!(
1414            serde_json::from_slice::<serde_json::Value>(&bytes).unwrap(),
1415            record
1416        );
1417        let logger = observer(
1418            FailingWriter {
1419                fail_write: Some(1),
1420                writes: 0,
1421                fail_flush: false,
1422            },
1423            1,
1424        );
1425        assert_eq!(
1426            logger.emit_root_record(&record),
1427            crate::DiagnosticSubmission::Enqueued
1428        );
1429        assert!(matches!(
1430            block_on(logger.shutdown()),
1431            Err(FlushError::OutputFailed(OutputStage::Record))
1432        ));
1433    }
1434
1435    #[test]
1436    fn random_source_failure_is_reported_during_initialization() {
1437        let result = Observer::with_writer_and_seed(
1438            ObserverConfig::default(),
1439            io::sink(),
1440            Err(InitError::RandomSource),
1441        );
1442        assert!(matches!(result, Err(InitError::RandomSource)));
1443    }
1444
1445    #[test]
1446    fn initialized_generator_produces_non_zero_unique_ids_without_more_io() {
1447        let observer = observer(io::sink(), 8);
1448        let first_trace = observer.new_trace_id();
1449        let second_trace = observer.new_trace_id();
1450        let first_span = observer.new_span_id();
1451        let second_span = observer.new_span_id();
1452        assert_ne!(first_trace.as_u128(), 0);
1453        assert_ne!(first_trace, second_trace);
1454        assert_ne!(first_span.as_u64(), 0);
1455        assert_ne!(first_span, second_span);
1456        block_on(observer.shutdown()).unwrap();
1457    }
1458
1459    #[test]
1460    fn observer_is_a_managed_component() {
1461        let observer = observer(io::sink(), 8);
1462        assert_eq!(ComponentLifecycle::name(&observer), "observability");
1463        block_on(ComponentLifecycle::start(&observer)).unwrap();
1464        block_on(ComponentLifecycle::shutdown(&observer)).unwrap();
1465    }
1466
1467    #[test]
1468    fn full_queue_is_reported_by_shutdown_without_waiting_in_emit() {
1469        let (entered_sender, entered_receiver) = mpsc::channel();
1470        let release = Arc::new((Mutex::new(false), Condvar::new()));
1471        let writer = BlockingWriter {
1472            entered: Some(entered_sender),
1473            release: release.clone(),
1474            fail_after_release: false,
1475        };
1476        let observer = observer(writer, 1);
1477
1478        observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
1479        entered_receiver.recv().unwrap();
1480        observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
1481        observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
1482        assert_eq!(observer.dropped_events(), 1);
1483
1484        let shutdown = observer.shutdown();
1485        let (lock, condition) = &*release;
1486        *lock.lock().unwrap() = true;
1487        condition.notify_one();
1488        assert_eq!(block_on(shutdown), Err(FlushError::DroppedEvents(1)));
1489    }
1490
1491    #[test]
1492    fn shutdown_reports_writer_failure_and_unreported_drops_together() {
1493        let (entered_sender, entered_receiver) = mpsc::channel();
1494        let release = Arc::new((Mutex::new(false), Condvar::new()));
1495        let writer = BlockingWriter {
1496            entered: Some(entered_sender),
1497            release: release.clone(),
1498            fail_after_release: true,
1499        };
1500        let observer = observer(writer, 1);
1501
1502        observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
1503        entered_receiver.recv().unwrap();
1504        observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
1505        observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
1506        let shutdown = observer.shutdown();
1507
1508        let (lock, condition) = &*release;
1509        *lock.lock().unwrap() = true;
1510        condition.notify_one();
1511        assert_eq!(
1512            block_on(shutdown),
1513            Err(FlushError::OutputFailedAndDropped {
1514                stage: OutputStage::Record,
1515                dropped_events: 1,
1516            })
1517        );
1518    }
1519
1520    #[test]
1521    fn record_write_failure_persists_until_shutdown() {
1522        let observer = observer(
1523            FailingWriter {
1524                fail_write: Some(1),
1525                writes: 0,
1526                fail_flush: false,
1527            },
1528            8,
1529        );
1530        observer.emit(LogRecord::new(EventLevel::Info, "test.record"));
1531        assert_eq!(
1532            block_on(observer.shutdown()),
1533            Err(FlushError::OutputFailed(OutputStage::Record))
1534        );
1535    }
1536
1537    #[test]
1538    fn newline_failure_persists_until_flush() {
1539        let observer = observer(
1540            FailingWriter {
1541                fail_write: Some(2),
1542                writes: 0,
1543                fail_flush: false,
1544            },
1545            8,
1546        );
1547        observer.emit(LogRecord::new(EventLevel::Info, "test.newline"));
1548        assert_eq!(
1549            block_on(observer.flush()),
1550            Err(FlushError::OutputFailed(OutputStage::Newline))
1551        );
1552        let _ = block_on(observer.shutdown());
1553    }
1554
1555    #[test]
1556    fn flush_failure_is_reported() {
1557        let observer = observer(
1558            FailingWriter {
1559                fail_write: None,
1560                writes: 0,
1561                fail_flush: true,
1562            },
1563            8,
1564        );
1565        assert_eq!(
1566            block_on(observer.flush()),
1567            Err(FlushError::OutputFailed(OutputStage::Flush))
1568        );
1569        assert_eq!(
1570            block_on(observer.shutdown()),
1571            Err(FlushError::OutputFailed(OutputStage::Flush))
1572        );
1573    }
1574}