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
179struct ContextByteCount(usize);
180impl Write for ContextByteCount {
181    fn write(&mut self, bytes: &[u8]) -> io::Result<usize> { self.0 = self.0.saturating_add(bytes.len()); Ok(bytes.len()) }
182    fn flush(&mut self) -> io::Result<()> { Ok(()) }
183}
184fn serialize_record_id<S: serde::Serializer>(id: &SpanId, serializer: S) -> Result<S::Ok, S::Error> { serializer.collect_str(id) }
185struct OrdinarySegments<'a> {
186    observer: &'a Observer, storage: &'a saddle_admission::ProcessLogStorage,
187    profile: Option<usize>, record_id: SpanId,
188}
189impl crate::root_diagnostic::SegmentOutput for OrdinarySegments<'_> {
190    fn submit(&self, segment: &crate::root_diagnostic::EncodedSegment<'_>) -> crate::DiagnosticSubmission {
191        #[derive(Serialize)]
192        #[serde(rename_all = "camelCase")]
193        struct Frame<'a> { record_schema: &'static str, event: &'static str, #[serde(serialize_with = "serialize_record_id")] record_id: SpanId,
194            format: &'static str, #[serde(skip_serializing_if = "Option::is_none")] profile: Option<&'a str>,
195            sequence: u64, payload: &'a str, state: &'static str }
196        let profile = self.profile.map(|index| &self.observer.inner.context.summaries[index]);
197        let format = if profile.is_some_and(|profile| profile.format == crate::SummaryFormat::Text) { "text" } else { "json" };
198        let frame = Frame { record_schema: "1.0", event: "context_record_segment", record_id: self.record_id,
199            format, profile: profile.map(|profile| profile.name.as_str()), sequence: segment.sequence, payload: segment.payload, state: segment.state };
200        self.observer.enqueue_context_frame(self.profile, self.storage, &|writer| serde_json::to_writer(writer, &frame))
201    }
202}
203
204trait LogSink: Write {
205    fn write_profile(&mut self, profile: usize, bytes: &[u8]) -> io::Result<()>;
206}
207struct BasicSink<W>(W);
208impl<W: Write> Write for BasicSink<W> {
209    fn write(&mut self, bytes: &[u8]) -> io::Result<usize> { self.0.write(bytes) }
210    fn flush(&mut self) -> io::Result<()> { self.0.flush() }
211}
212impl<W: Write> LogSink for BasicSink<W> {
213    fn write_profile(&mut self, _: usize, bytes: &[u8]) -> io::Result<()> { self.0.write_all(bytes)?; self.0.write_all(b"\n") }
214}
215struct CalendarOutputs { main: CalendarFileWriter, summaries: Vec<Option<CalendarFileWriter>> }
216impl Write for CalendarOutputs {
217    fn write(&mut self, bytes: &[u8]) -> io::Result<usize> { self.main.write(bytes) }
218    fn flush(&mut self) -> io::Result<()> {
219        self.main.flush()?;
220        for output in self.summaries.iter_mut().flatten() { output.flush()?; }
221        Ok(())
222    }
223}
224impl LogSink for CalendarOutputs {
225    fn write_profile(&mut self, profile: usize, bytes: &[u8]) -> io::Result<()> {
226        let output = self.summaries.get_mut(profile).and_then(Option::as_mut).ok_or(io::ErrorKind::InvalidInput)?;
227        output.write_all(bytes)?; output.write_all(b"\n")
228    }
229}
230
231pub(crate) struct Inner {
232    sender: SyncSender<Command>,
233    dropped: AtomicU64,
234    accepting: AtomicBool,
235    emitting: AtomicUsize,
236    shutdown_started: AtomicBool,
237    worker: Mutex<Option<JoinHandle<()>>>,
238    process_log: Mutex<Option<saddle_admission::ProcessLogStorage>>,
239    context: crate::ContextLoggingConfig,
240    writer_source: Arc<Mutex<WriterSourceState>>,
241    ids: IdGenerator,
242    pub(crate) metrics: crate::metrics::Metrics,
243}
244
245#[derive(Clone)]
246struct WriterSourceContext {
247    application: saddle_core::ContextLabel,
248    output: crate::SourceOutput,
249    primary: Option<saddle_core::DiagnosticOccurrence>,
250}
251#[derive(Default)]
252struct WriterSourceState {
253    context: Option<WriterSourceContext>,
254    failure: Option<WriterSourceFailure>,
255}
256struct WriterSourceFailure {
257    original: WriterStageError,
258    diagnostic: saddle_core::Diagnostic,
259    written: Option<crate::root_diagnostic::WrittenComponentFailure>,
260}
261
262#[derive(Debug)]
263struct WriterStageError {
264    stage: OutputStage,
265    original: io::Error,
266}
267impl fmt::Display for WriterStageError {
268    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
269        write!(formatter, "log writer failed during {:?}: {}", self.stage, self.original)
270    }
271}
272impl Error for WriterStageError {
273    fn source(&self) -> Option<&(dyn Error + 'static)> { Some(&self.original) }
274}
275
276fn record_writer_failure(source: &Arc<Mutex<WriterSourceState>>,
277    stage: OutputStage, original: io::Error) {
278    let context = source.lock().unwrap_or_else(|e| e.into_inner()).context.clone();
279    let original = WriterStageError { stage, original };
280    let mut diagnostic = saddle_core::Diagnostic::capture(
281        saddle_core::DiagnosticCategory::UnexpectedError,
282        saddle_core::CaptureSite::FirstObserved,
283        saddle_core::DiagnosticCause::new(
284            saddle_core::DiagnosticStage::ShutdownLogger,
285            saddle_core::DiagnosticCode::new("observability.writer_failed")
286                .expect("static writer failure code")),
287    );
288    if let Some(primary) = context.as_ref().and_then(|context| context.primary) {
289        diagnostic = diagnostic.during_cleanup_of_occurrence(&primary);
290    }
291    let written = context.as_ref().and_then(|context|
292        crate::root_diagnostic::process_component_cleanup_recorded(
293            Some(context.output.handle()),
294            saddle_core::ContextFact::Present(context.application.clone()),
295            &diagnostic, &original).ok());
296    let mut state = source.lock().unwrap_or_else(|e| e.into_inner());
297    if state.failure.is_none() {
298        state.failure = Some(WriterSourceFailure { original, diagnostic, written });
299    }
300}
301
302struct IdGenerator {
303    trace_high: u64,
304    trace_seed: u64,
305    trace_counter: AtomicU64,
306    span_seed: u64,
307    span_counter: AtomicU64,
308}
309
310impl IdGenerator {
311    fn from_seed(seed: [u8; 24]) -> Self {
312        let mut high = u64::from_be_bytes(seed[0..8].try_into().unwrap());
313        if high == 0 {
314            high = 1;
315        }
316        Self {
317            trace_high: high,
318            trace_seed: u64::from_be_bytes(seed[8..16].try_into().unwrap()),
319            trace_counter: AtomicU64::new(0),
320            span_seed: u64::from_be_bytes(seed[16..24].try_into().unwrap()),
321            span_counter: AtomicU64::new(0),
322        }
323    }
324
325    fn trace_id(&self) -> TraceId {
326        let low = self
327            .trace_seed
328            .wrapping_add(self.trace_counter.fetch_add(1, Ordering::Relaxed));
329        TraceId::from_u128((u128::from(self.trace_high) << 64) | u128::from(low))
330    }
331
332    fn span_id(&self) -> SpanId {
333        loop {
334            let value = self
335                .span_seed
336                .wrapping_add(self.span_counter.fetch_add(1, Ordering::Relaxed));
337            if value != 0 {
338                return SpanId::from_u64(value);
339            }
340        }
341    }
342}
343
344enum Command {
345    Record(LogRecord),
346    ProcessCall {
347        bytes: saddle_admission::ExactStored<Vec<u8>>,
348        dropped: u64,
349    },
350    ContextProfile { bytes: saddle_admission::ExactStored<Vec<u8>>, profile: Option<usize>, dropped: u64 },
351    // Source-encoded, bounded and self-contained. Never contains a request root.
352    RootRecord {
353        bytes: Vec<u8>,
354        dropped: u64,
355    },
356    Flush {
357        unreported_dropped: u64,
358        response: mpsc::Sender<WorkerStatus>,
359    },
360    Shutdown {
361        unreported_dropped: u64,
362        response: mpsc::Sender<WorkerStatus>,
363    },
364}
365pub(crate) fn root_queue_layout() -> std::alloc::Layout {
366    std::alloc::Layout::new::<Command>()
367}
368
369#[derive(Clone, Copy, Debug, Default)]
370struct WorkerStatus {
371    failure: Option<OutputStage>,
372    unreported_dropped: u64,
373}
374
375impl WorkerStatus {
376    fn into_result(self) -> Result<(), FlushError> {
377        match (self.failure, self.unreported_dropped) {
378            (None, 0) => Ok(()),
379            (Some(stage), 0) => Err(FlushError::OutputFailed(stage)),
380            (None, dropped_events) => Err(FlushError::DroppedEvents(dropped_events)),
381            (Some(stage), dropped_events) => Err(FlushError::OutputFailedAndDropped {
382                stage,
383                dropped_events,
384            }),
385        }
386    }
387}
388
389#[derive(Serialize)]
390pub(crate) struct LogRecord {
391    pub timestamp_unix_ms: u128,
392    pub level: EventLevel,
393    pub event: &'static str,
394    #[serde(skip_serializing_if = "Option::is_none")]
395    pub trace_id: Option<String>,
396    #[serde(skip_serializing_if = "Option::is_none")]
397    pub span: Option<String>,
398    #[serde(skip_serializing_if = "Option::is_none")]
399    pub span_id: Option<String>,
400    #[serde(skip_serializing_if = "Option::is_none")]
401    pub parent: Option<String>,
402    #[serde(skip_serializing_if = "Option::is_none")]
403    pub parent_span_id: Option<String>,
404    #[serde(flatten)]
405    pub data: serde_json::Map<String, serde_json::Value>,
406    #[serde(skip_serializing_if = "is_zero")]
407    pub dropped_events: u64,
408}
409
410const fn is_zero(value: &u64) -> bool {
411    *value == 0
412}
413
414impl LogRecord {
415    pub(crate) fn new(level: EventLevel, event: &'static str) -> Self {
416        Self {
417            timestamp_unix_ms: SystemTime::now()
418                .duration_since(UNIX_EPOCH)
419                .unwrap_or_default()
420                .as_millis(),
421            level,
422            event,
423            trace_id: None,
424            span: None,
425            span_id: None,
426            parent: None,
427            parent_span_id: None,
428            data: serde_json::Map::new(),
429            dropped_events: 0,
430        }
431    }
432}
433
434/// Initializes stdout logging once. Later calls return the first observer.
435pub fn init(config: ObserverConfig) -> Result<&'static Observer, InitError> {
436    if let Some(observer) = GLOBAL.get() {
437        return Ok(observer);
438    }
439
440    let observer = Observer::with_writer(config, io::stdout())?;
441    let _ = GLOBAL.set(observer);
442    Ok(GLOBAL.get().expect("global observer was initialized"))
443}
444
445/// Initializes the process logger with the sole calendar-rotated file sink.
446pub fn init_file(
447    config: ObserverConfig,
448    file: FileLoggingConfig,
449) -> Result<&'static Observer, InitError> {
450    if let Some(observer) = GLOBAL.get() {
451        return Ok(observer);
452    }
453    let context = file.context().clone();
454    let summaries = context.summaries.iter().map(|profile| {
455        if profile.enabled { CalendarFileWriter::open(file.for_summary(&profile.name)).map(Some) }
456        else { Ok(None) }
457    }).collect::<io::Result<Vec<_>>>().map_err(InitError::Output)?;
458    let main = CalendarFileWriter::open(file).map_err(InitError::Output)?;
459    let mut seed = [0; 24];
460    getrandom::fill(&mut seed).map_err(|_| InitError::RandomSource)?;
461    let observer = Observer::with_sink_and_seed(config, CalendarOutputs { main, summaries }, Ok(seed), context)?;
462    let _ = GLOBAL.set(observer);
463    Ok(GLOBAL.get().expect("global observer was initialized"))
464}
465
466/// Returns the initialized process observer, if application startup installed it.
467pub fn global() -> Option<&'static Observer> {
468    GLOBAL.get()
469}
470
471impl Observer {
472    /// Bind the original process framework domain before formal routed RPCs
473    /// can emit ordinary call records. The writer never receives this handle.
474    #[doc(hidden)]
475    pub fn install_process_log_storage(&self, storage: saddle_admission::ProcessLogStorage) {
476        let mut slot = self.inner.process_log.lock().unwrap_or_else(|p| p.into_inner());
477        *slot = Some(storage);
478    }
479
480    /// Prepare a single terminal snapshot before account closure. Refusal
481    /// leaves the live view untouched and reserves no unbounded cache.
482    #[doc(hidden)]
483    pub fn capture_terminal_context(&self, context: &saddle_core::request_context::UnifiedContext<'_>)
484        -> Result<crate::OwnedUnifiedContext, crate::DiagnosticSubmission> {
485        use crate::DiagnosticSubmission as S;
486        let storage = self.inner.process_log.lock().unwrap_or_else(|p| p.into_inner())
487            .clone().ok_or(S::OutputUnavailable)?;
488        let mut count = ContextByteCount(0);
489        serde_json::to_writer(&mut count, context).map_err(|_| S::EncodingFailed)?;
490        let layout = std::alloc::Layout::array::<u8>(count.0).map_err(|_| S::EncodingFailed)?;
491        let permit = storage.try_reserve(layout).map_err(|_| S::Full)?;
492        let mut bytes = permit.allocate_exact(layout, || Vec::with_capacity(count.0));
493        struct Bounded<'a> { bytes: &'a mut Vec<u8>, exceeded: bool }
494        impl Write for Bounded<'_> {
495            fn write(&mut self, chunk: &[u8]) -> io::Result<usize> {
496                if !self.exceeded && chunk.len() <= self.bytes.capacity() - self.bytes.len() {
497                    self.bytes.extend_from_slice(chunk);
498                } else { self.exceeded = true; }
499                // Consume a changed encoder without reallocating or producing
500                // a heap-backed serde IO error under the request audit.
501                Ok(chunk.len())
502            }
503            fn flush(&mut self) -> io::Result<()> { Ok(()) }
504        }
505        let mut writer = Bounded { bytes: bytes.get_mut(), exceeded: false };
506        serde_json::to_writer(&mut writer, context).map_err(|_| S::EncodingFailed)?;
507        if writer.exceeded || writer.bytes.len() != count.0 { return Err(S::EncodingFailed); }
508        Ok(crate::OwnedUnifiedContext { bytes })
509    }
510
511    pub(crate) fn emit_process_call(&self, record: &mut crate::call::ProcessCallRecord<'_>) {
512        struct Counter(usize);
513        impl Write for Counter {
514            fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
515                self.0 = self.0.checked_add(bytes.len()).ok_or(io::ErrorKind::OutOfMemory)?;
516                Ok(bytes.len())
517            }
518            fn flush(&mut self) -> io::Result<()> { Ok(()) }
519        }
520        self.inner.emitting.fetch_add(1, Ordering::AcqRel);
521        let submitted = (|| {
522            if !self.inner.accepting.load(Ordering::Acquire) { return false; }
523            let storage = self.inner.process_log.lock().unwrap_or_else(|p| p.into_inner()).clone();
524            let Some(storage) = storage else { return false; };
525            let dropped = self.inner.dropped.swap(0, Ordering::Relaxed);
526            record.dropped_events = dropped;
527            let mut counter = Counter(0);
528            if serde_json::to_writer(&mut counter, record).is_err() {
529                self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
530                return false;
531            }
532            let Ok(layout) = std::alloc::Layout::array::<u8>(counter.0) else {
533                self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
534                return false;
535            };
536            let Ok(permit) = storage.try_reserve(layout) else {
537                self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
538                return false;
539            };
540            let mut bytes = permit.allocate_exact(layout, || Vec::with_capacity(counter.0));
541            if serde_json::to_writer(bytes.get_mut(), record).is_err() || bytes.get().len() != counter.0 {
542                drop(bytes);
543                self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
544                return false;
545            }
546            match self.inner.sender.try_send(Command::ProcessCall { bytes, dropped }) {
547                Ok(()) => true,
548                Err(TrySendError::Full(command)) => {
549                    drop(command);
550                    self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
551                    false
552                }
553                Err(TrySendError::Disconnected(command)) => {
554                    drop(command);
555                    self.inner.metrics.logger_output(true);
556                    self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
557                    false
558                }
559            }
560        })();
561        if !submitted {
562            self.inner.metrics.logger_dropped(1);
563            self.inner.dropped.fetch_add(1, Ordering::Relaxed);
564        }
565        self.inner.emitting.fetch_sub(1, Ordering::Release);
566    }
567    /// Creates an observer with a framework-owned writer.
568    ///
569    /// Application startup should normally use [`init`]. This constructor is
570    /// useful for embedding and deterministic tests.
571    pub fn with_writer(
572        config: ObserverConfig,
573        writer: impl Write + Send + 'static,
574    ) -> Result<Self, InitError> {
575        let mut seed = [0_u8; 24];
576        getrandom::fill(&mut seed).map_err(|_| InitError::RandomSource)?;
577        Self::with_writer_and_seed(config, writer, Ok(seed))
578    }
579
580    fn with_writer_and_seed(
581        config: ObserverConfig,
582        writer: impl Write + Send + 'static,
583        seed: Result<[u8; 24], InitError>,
584    ) -> Result<Self, InitError> {
585        Self::with_sink_and_seed(config, BasicSink(writer), seed, crate::ContextLoggingConfig::default())
586    }
587    fn with_sink_and_seed(config: ObserverConfig, writer: impl LogSink + Send + 'static,
588        seed: Result<[u8; 24], InitError>, context: crate::ContextLoggingConfig) -> Result<Self, InitError> {
589        if config.queue_capacity == 0 {
590            return Err(InitError::EmptyQueue);
591        }
592        let ids = IdGenerator::from_seed(seed?);
593        let (sender, receiver) = mpsc::sync_channel(config.queue_capacity);
594        let writer_source = Arc::new(Mutex::new(WriterSourceState::default()));
595        let worker_source = Arc::clone(&writer_source);
596        let worker = thread::Builder::new()
597            .name("saddle-log-writer".to_owned())
598            .spawn(move || write_records(receiver, writer, worker_source))
599            .map_err(InitError::Spawn)?;
600        Ok(Self {
601            inner: Arc::new(Inner {
602                sender,
603                dropped: AtomicU64::new(0),
604                accepting: AtomicBool::new(true),
605                emitting: AtomicUsize::new(0),
606                shutdown_started: AtomicBool::new(false),
607                worker: Mutex::new(Some(worker)),
608                process_log: Mutex::new(None),
609                context,
610                writer_source,
611                ids,
612                metrics: crate::metrics::Metrics::default(),
613            }),
614        })
615    }
616
617    pub(crate) fn new_trace_id(&self) -> TraceId {
618        self.inner.ids.trace_id()
619    }
620
621    pub(crate) fn new_span_id(&self) -> SpanId {
622        self.inner.ids.span_id()
623    }
624
625    pub(crate) fn emit(&self, mut record: LogRecord) {
626        if !self.inner.accepting.load(Ordering::Acquire) {
627            return;
628        }
629        self.inner.emitting.fetch_add(1, Ordering::AcqRel);
630        if !self.inner.accepting.load(Ordering::Acquire) {
631            self.inner.emitting.fetch_sub(1, Ordering::AcqRel);
632            return;
633        }
634
635        record.dropped_events = self.inner.dropped.swap(0, Ordering::Relaxed);
636        match self.inner.sender.try_send(Command::Record(record)) {
637            Ok(()) => {}
638            Err(TrySendError::Full(Command::Record(record))) => {
639                self.inner.metrics.logger_dropped(1);
640                self.inner
641                    .dropped
642                    .fetch_add(record.dropped_events.saturating_add(1), Ordering::Relaxed);
643            }
644            Err(TrySendError::Disconnected(_)) => {
645                self.inner.metrics.logger_dropped(1);
646                self.inner.metrics.logger_output(true);
647                self.inner.dropped.fetch_add(1, Ordering::Relaxed);
648            }
649            Err(TrySendError::Full(_)) => unreachable!("emit only sends records"),
650        }
651        self.inner.emitting.fetch_sub(1, Ordering::Release);
652    }
653
654    pub(crate) fn emit_context_record<P: crate::ContextRecordPayload>(&self, record: &crate::ContextRecord<'_, P>) -> crate::DiagnosticSubmission {
655        use crate::DiagnosticSubmission as S;
656        let mut result = S::Enqueued;
657        if self.inner.context.full.enabled {
658            result = self.enqueue_context(None, |writer| serde_json::to_writer(writer, record));
659        }
660        for (index, profile) in self.inner.context.summaries.iter().enumerate().filter(|(_, p)| p.enabled) {
661            let next = self.enqueue_context(Some(index), |writer| record.write_summary(profile, &mut &mut *writer));
662            if next != S::Enqueued { result = next; }
663        }
664        result
665    }
666    fn enqueue_context(&self, profile: Option<usize>, encode: impl Fn(&mut dyn Write) -> serde_json::Result<()>) -> crate::DiagnosticSubmission {
667        use crate::DiagnosticSubmission as S;
668        self.inner.emitting.fetch_add(1, Ordering::AcqRel);
669        let result = (|| {
670            if !self.inner.accepting.load(Ordering::Acquire) { return S::Closed; }
671            let Some(storage) = self.inner.process_log.lock().unwrap_or_else(|p| p.into_inner()).clone() else { return S::OutputUnavailable; };
672            let mut count = ContextByteCount(0);
673            if encode(&mut count).is_err() { return S::EncodingFailed; }
674            if count.0 <= 8191 { return self.enqueue_context_frame(profile, &storage, &encode); }
675            let output = OrdinarySegments { observer: self, storage: &storage, profile, record_id: self.inner.ids.span_id() };
676            crate::root_diagnostic::stream_record(&output, encode)
677        })();
678        if result != S::Enqueued { self.inner.dropped.fetch_add(1, Ordering::Relaxed); self.inner.metrics.logger_dropped(1); }
679        self.inner.emitting.fetch_sub(1, Ordering::Release);
680        result
681    }
682    fn enqueue_context_frame(&self, profile: Option<usize>, storage: &saddle_admission::ProcessLogStorage,
683        encode: &impl Fn(&mut dyn Write) -> serde_json::Result<()>) -> crate::DiagnosticSubmission {
684        use crate::DiagnosticSubmission as S;
685        if !self.inner.accepting.load(Ordering::Acquire) { return S::Closed; }
686        let mut count = ContextByteCount(0);
687        if encode(&mut count).is_err() || count.0 > 8191 { return S::EncodingFailed; }
688        let Ok(layout) = std::alloc::Layout::array::<u8>(count.0) else { return S::EncodingFailed; };
689        let Ok(permit) = storage.try_reserve(layout) else { return S::Full; };
690        let mut bytes = permit.allocate_exact(layout, || Vec::with_capacity(count.0));
691        if encode(bytes.get_mut()).is_err() || bytes.get().len() != count.0 { return S::EncodingFailed; }
692        let dropped = self.inner.dropped.swap(0, Ordering::Relaxed);
693        match self.inner.sender.try_send(Command::ContextProfile { bytes, profile, dropped }) {
694            Ok(()) => S::Enqueued,
695            Err(TrySendError::Full(command)) => { drop(command); self.inner.dropped.fetch_add(dropped, Ordering::Relaxed); S::Full },
696            Err(TrySendError::Disconnected(command)) => { drop(command); self.inner.dropped.fetch_add(dropped, Ordering::Relaxed); self.inner.metrics.logger_output(true); S::Closed },
697        }
698    }
699
700    pub(crate) fn emit_root_record<T: Serialize>(&self, record: &T) -> crate::DiagnosticSubmission {
701        use crate::DiagnosticSubmission as Submission;
702        self.inner.emitting.fetch_add(1, Ordering::AcqRel);
703        let result = (|| {
704            if !self.inner.accepting.load(Ordering::Acquire) {
705                return Submission::Closed;
706            }
707            let mut frame = crate::diagnostic::FixedDiagnosticBytes {
708                bytes: [0; 8192],
709                len: 0,
710            };
711            if serde_json::to_writer(&mut frame, record).is_err() {
712                return Submission::EncodingFailed;
713            }
714            let dropped = self.inner.dropped.swap(0, Ordering::Relaxed);
715            // Only initialized bytes move into the existing logger domain. This
716            // allocation and source frame peak are separate R0 layout inputs.
717            let bytes = frame.bytes[..frame.len].to_vec();
718            match self
719                .inner
720                .sender
721                .try_send(Command::RootRecord { bytes, dropped })
722            {
723                Ok(()) => Submission::Enqueued,
724                Err(TrySendError::Full(Command::RootRecord { dropped, .. })) => {
725                    self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
726                    Submission::Full
727                }
728                Err(TrySendError::Disconnected(_)) => {
729                    self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
730                    self.inner.metrics.logger_output(true);
731                    Submission::Closed
732                }
733                Err(TrySendError::Full(_)) => unreachable!("root record submission"),
734            }
735        })();
736        if result != Submission::Enqueued {
737            self.inner.dropped.fetch_add(1, Ordering::Relaxed);
738            self.inner.metrics.logger_dropped(1);
739        }
740        self.inner.emitting.fetch_sub(1, Ordering::Release);
741        result
742    }
743
744    /// Returns a fixed-size, allocation-free copy of the current process metrics.
745    pub fn metrics_snapshot(&self) -> crate::MetricsSnapshot {
746        self.inner.metrics.snapshot()
747    }
748
749    /// Flushes records accepted before the lifecycle coordinator reaches the
750    /// writer. The returned future never blocks the async executor thread.
751    pub fn flush(&self) -> FlushFuture {
752        if !self.inner.accepting.load(Ordering::Acquire) {
753            return FlushFuture {
754                completion: Completion::ready(Err(FlushError::WriterStopped)),
755            };
756        }
757        let completion = Completion::pending();
758        let future = FlushFuture {
759            completion: completion.clone(),
760        };
761        let inner = self.inner.clone();
762        if thread::Builder::new()
763            .name("saddle-log-flush".to_owned())
764            .spawn(move || coordinate_flush(inner, completion.clone()))
765            .is_err()
766        {
767            future
768                .completion
769                .complete(Err(FlushError::CoordinatorUnavailable));
770        }
771        future
772    }
773
774    /// Stops admission, drains accepted records, flushes output and joins the
775    /// writer without blocking the async executor thread.
776    pub fn shutdown(&self) -> FlushFuture {
777        if self.inner.shutdown_started.swap(true, Ordering::AcqRel) {
778            return FlushFuture {
779                completion: Completion::ready(Err(FlushError::AlreadyShuttingDown)),
780            };
781        }
782        self.inner.accepting.store(false, Ordering::Release);
783        let completion = Completion::pending();
784        let future = FlushFuture {
785            completion: completion.clone(),
786        };
787        let inner = self.inner.clone();
788        if thread::Builder::new()
789            .name("saddle-log-shutdown".to_owned())
790            .spawn(move || coordinate_shutdown(inner, completion.clone()))
791            .is_err()
792        {
793            self.inner.accepting.store(true, Ordering::Release);
794            self.inner.shutdown_started.store(false, Ordering::Release);
795            future
796                .completion
797                .complete(Err(FlushError::CoordinatorUnavailable));
798        }
799        future
800    }
801
802    pub fn dropped_events(&self) -> u64 {
803        self.inner.dropped.load(Ordering::Relaxed)
804    }
805}
806
807impl ComponentLifecycle for Observer {
808    fn name(&self) -> &'static str {
809        "observability"
810    }
811
812    fn start(&self) -> LifecycleFuture<'_> {
813        Box::pin(async { Ok(()) })
814    }
815
816    fn shutdown(&self) -> LifecycleFuture<'_> {
817        let shutdown = Observer::shutdown(self);
818        Box::pin(async move {
819            shutdown.await.map_err(|_| {
820                SaddleError::new(
821                    ErrorKind::Infrastructure,
822                    "observability.shutdown_failed",
823                    "structured log shutdown failed",
824                )
825            })
826        })
827    }
828}
829
830impl Observer {
831    fn recorded_start<'a>(&'a self, application: &'a saddle_core::ContextLabel,
832        output: &'a crate::SourceOutput)
833        -> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
834        self.inner.writer_source.lock().unwrap_or_else(|e| e.into_inner()).context =
835            Some(WriterSourceContext {
836                application: application.clone(), output: output.clone(), primary: None,
837            });
838        Box::pin(async { Ok(()) })
839    }
840
841    fn recorded_shutdown<'a>(&'a self,
842        application: &'a saddle_core::ContextLabel,
843        output: &'a crate::SourceOutput,
844        primary: Option<saddle_core::DiagnosticOccurrence>)
845        -> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
846        self.inner.writer_source.lock().unwrap_or_else(|e| e.into_inner()).context =
847            Some(WriterSourceContext {
848                application: application.clone(), output: output.clone(), primary,
849            });
850        Box::pin(async move {
851            let result = Observer::shutdown(self).await;
852            let failure = self.inner.writer_source.lock()
853                .unwrap_or_else(|e| e.into_inner()).failure.take();
854            if let Some(failure) = failure {
855                return Err(crate::root_diagnostic::RecordedSaddleError::previously_written_writer_error(
856                    failure.original, failure.diagnostic, failure.written));
857            }
858            crate::root_diagnostic::RecordedSaddleError::component_result(
859                result, application, output,
860                crate::root_diagnostic::ComponentSourceKind::Cleanup, primary,
861                ErrorKind::Infrastructure, "observability.shutdown_failed",
862                "structured log shutdown failed",
863            )
864        })
865    }
866}
867
868/// The formal component boundary keeps the actual FlushError until the
869/// controlled source file has acknowledged it. The Core lifecycle impl above
870/// remains a legacy compatibility surface for callers migrating to this one.
871impl crate::root_diagnostic::RecordedComponentLifecycle for Observer {
872    fn name(&self) -> &'static str { "observability" }
873
874    fn start<'a>(&'a self, application: &'a saddle_core::ContextLabel,
875        output: &'a crate::SourceOutput)
876        -> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
877        Observer::recorded_start(self, application, output)
878    }
879
880    fn shutdown<'a>(&'a self, application: &'a saddle_core::ContextLabel,
881        output: &'a crate::SourceOutput)
882        -> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
883        crate::root_diagnostic::RecordedComponentLifecycle::shutdown_with_primary(
884            self, application, output, None)
885    }
886
887    fn shutdown_with_primary<'a>(&'a self,
888        application: &'a saddle_core::ContextLabel,
889        output: &'a crate::SourceOutput,
890        primary: Option<saddle_core::DiagnosticOccurrence>)
891        -> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
892        Observer::recorded_shutdown(self, application, output, primary)
893    }
894}
895
896fn coordinate_flush(inner: Arc<Inner>, completion: Arc<Completion>) {
897    let dropped = inner.dropped.swap(0, Ordering::AcqRel);
898    let (response, receiver) = mpsc::channel();
899    let result = if inner
900        .sender
901        .send(Command::Flush {
902            unreported_dropped: dropped,
903            response,
904        })
905        .is_err()
906    {
907        Err(FlushError::WriterStopped)
908    } else {
909        receiver
910            .recv()
911            .map_err(|_| FlushError::WriterStopped)
912            .and_then(WorkerStatus::into_result)
913    };
914    completion.complete(result);
915}
916
917fn coordinate_shutdown(inner: Arc<Inner>, completion: Arc<Completion>) {
918    while inner.emitting.load(Ordering::Acquire) != 0 {
919        thread::yield_now();
920    }
921    let dropped = inner.dropped.swap(0, Ordering::AcqRel);
922    let (response, receiver) = mpsc::channel();
923    let mut result = if inner
924        .sender
925        .send(Command::Shutdown {
926            unreported_dropped: dropped,
927            response,
928        })
929        .is_err()
930    {
931        Err(FlushError::WriterStopped)
932    } else {
933        receiver
934            .recv()
935            .map_err(|_| FlushError::WriterStopped)
936            .and_then(WorkerStatus::into_result)
937    };
938
939    if let Some(worker) = inner.worker.lock().unwrap().take() {
940        if worker.join().is_err() {
941            result = Err(FlushError::WorkerPanicked);
942        }
943    }
944    inner.process_log.lock().unwrap_or_else(|p| p.into_inner()).take();
945    completion.complete(result);
946}
947
948fn write_records(receiver: Receiver<Command>, mut writer: impl LogSink,
949    source: Arc<Mutex<WriterSourceState>>) {
950    let mut status = WorkerStatus::default();
951    while let Ok(command) = receiver.recv() {
952        match command {
953            Command::Record(record) => write_record(&mut writer, record, &mut status, &source),
954            Command::ProcessCall { bytes, dropped } => {
955                status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
956                if status.failure.is_none() {
957                    if let Err(error) = writer.write_all(bytes.get()) {
958                        record_writer_failure(&source, OutputStage::Record, error);
959                        status.failure = Some(OutputStage::Record);
960                    } else if let Err(error) = writer.write_all(b"\n") {
961                        record_writer_failure(&source, OutputStage::Newline, error);
962                        status.failure = Some(OutputStage::Newline);
963                    }
964                }
965                drop(bytes);
966            }
967            Command::ContextProfile { bytes, profile, dropped } => {
968                status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
969                if status.failure.is_none() {
970                    let result = match profile {
971                        Some(profile) => writer.write_profile(profile, bytes.get()),
972                        None => writer.write_all(bytes.get()).and_then(|_| writer.write_all(b"\n")),
973                    };
974                    if let Err(error) = result { record_writer_failure(&source, OutputStage::Record, error); status.failure = Some(OutputStage::Record); }
975                }
976                drop(bytes);
977            }
978            Command::RootRecord { bytes, dropped } => {
979                status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
980                if status.failure.is_none() {
981                    if let Err(error) = writer.write_all(&bytes) {
982                        record_writer_failure(&source, OutputStage::Record, error);
983                        status.failure = Some(OutputStage::Record);
984                    } else if let Err(error) = writer.write_all(b"\n") {
985                        record_writer_failure(&source, OutputStage::Newline, error);
986                        status.failure = Some(OutputStage::Newline);
987                    }
988                }
989            }
990            Command::Flush {
991                unreported_dropped,
992                response,
993            } => {
994                status.unreported_dropped =
995                    status.unreported_dropped.saturating_add(unreported_dropped);
996                flush_writer(&mut writer, &mut status, &source);
997                let _ = response.send(status);
998            }
999            Command::Shutdown {
1000                unreported_dropped,
1001                response,
1002            } => {
1003                status.unreported_dropped =
1004                    status.unreported_dropped.saturating_add(unreported_dropped);
1005                flush_writer(&mut writer, &mut status, &source);
1006                let _ = response.send(status);
1007                break;
1008            }
1009        }
1010    }
1011}
1012
1013fn write_record(writer: &mut impl Write, record: LogRecord, status: &mut WorkerStatus,
1014    source: &Arc<Mutex<WriterSourceState>>) {
1015    if status.failure.is_some() {
1016        status.unreported_dropped = status
1017            .unreported_dropped
1018            .saturating_add(record.dropped_events);
1019        return;
1020    }
1021
1022    let dropped = record.dropped_events;
1023    let bytes = match serde_json::to_vec(&record) {
1024        Ok(bytes) => bytes,
1025        Err(_) => {
1026            status.failure = Some(OutputStage::Serialize);
1027            status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
1028            return;
1029        }
1030    };
1031    if let Err(error) = writer.write_all(&bytes) {
1032        record_writer_failure(source, OutputStage::Record, error);
1033        status.failure = Some(OutputStage::Record);
1034        status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
1035        return;
1036    }
1037    if let Err(error) = writer.write_all(b"\n") {
1038        record_writer_failure(source, OutputStage::Newline, error);
1039        status.failure = Some(OutputStage::Newline);
1040        status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
1041    }
1042}
1043
1044fn flush_writer(writer: &mut impl Write, status: &mut WorkerStatus,
1045    source: &Arc<Mutex<WriterSourceState>>) {
1046    // A previous output failure has already made this writer unhealthy.
1047    // Avoid issuing another fallible operation whose error could not be the
1048    // primary result and would otherwise be discarded.
1049    if status.failure.is_some() { return; }
1050    if let Err(error) = writer.flush() {
1051        record_writer_failure(source, OutputStage::Flush, error);
1052        status.failure = Some(OutputStage::Flush);
1053    }
1054}
1055
1056#[cfg(test)]
1057mod tests {
1058    use std::{
1059        sync::{Condvar, mpsc},
1060        task::{Wake, Waker},
1061    };
1062
1063    use super::*;
1064
1065    #[test]
1066    fn recorded_observer_shutdown_preserves_actual_flush_error() {
1067        let directory = std::env::temp_dir().join(format!(
1068            "saddle-observer-recorded-{}-{}", std::process::id(),
1069            SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos()));
1070        std::fs::create_dir(&directory).unwrap();
1071        let mut output = crate::EmergencyDiagnostics::start_checked(
1072            &crate::FileLoggingConfig::new(directory.clone(), crate::Rotation::Daily)).unwrap();
1073        let selected = output.source_output().unwrap();
1074        let target = output.target().to_owned();
1075        let logger = observer(io::sink(), 1);
1076        block_on(Observer::shutdown(&logger)).unwrap();
1077        let app = saddle_core::ContextLabel::checked("recorded-observer-test").unwrap();
1078        let failure = block_on(
1079            crate::root_diagnostic::RecordedComponentLifecycle::shutdown(
1080                &logger, &app, &selected)).err().unwrap();
1081        assert_eq!(failure.safe().code(), "observability.shutdown_failed");
1082        assert!(failure.original_if_unconfirmed().is_none());
1083        let record = std::fs::read_to_string(&target).unwrap();
1084        assert!(record.contains("AlreadyShuttingDown"));
1085        assert!(record.contains(&failure.safe().diagnostic().unwrap().id().to_string()));
1086        drop(selected);
1087        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
1088        while output.shutdown() == crate::DiagnosticShutdown::Pending
1089            && std::time::Instant::now() < deadline { std::thread::yield_now(); }
1090        assert_eq!(output.shutdown(), crate::DiagnosticShutdown::Finished);
1091        drop(output);
1092        std::fs::remove_dir_all(directory).unwrap();
1093    }
1094
1095    #[derive(Debug)]
1096    struct WriterChain {
1097        message: &'static str,
1098        cause: io::Error,
1099    }
1100    impl fmt::Display for WriterChain {
1101        fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1102            f.write_str(self.message)
1103        }
1104    }
1105    impl Error for WriterChain {
1106        fn source(&self) -> Option<&(dyn Error + 'static)> { Some(&self.cause) }
1107    }
1108    struct ChainWriter {
1109        fail_flush: bool,
1110        release_write: Option<(mpsc::Sender<()>, Arc<(Mutex<bool>, Condvar)>)>,
1111    }
1112    impl Write for ChainWriter {
1113        fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
1114            if self.fail_flush { Ok(bytes.len()) } else {
1115                if let Some((entered, release)) = self.release_write.take() {
1116                    entered.send(()).unwrap();
1117                    let (lock, condition) = &*release;
1118                    let mut ready = lock.lock().unwrap();
1119                    while !*ready { ready = condition.wait(ready).unwrap(); }
1120                }
1121                Err(io::Error::other(WriterChain {
1122                    message: "writer top original 4931",
1123                    cause: io::Error::other("writer nested original 4932"),
1124                }))
1125            }
1126        }
1127        fn flush(&mut self) -> io::Result<()> {
1128            if self.fail_flush {
1129                Err(io::Error::other(WriterChain {
1130                    message: "flush top original 4933",
1131                    cause: io::Error::other("flush nested original 4934"),
1132                }))
1133            } else { Ok(()) }
1134        }
1135    }
1136
1137    #[test]
1138    fn writer_write_and_flush_originals_precede_formal_shutdown_completion() {
1139        use crate::root_diagnostic::RecordedComponentLifecycle;
1140        for flush in [false, true] {
1141            let directory = std::env::temp_dir().join(format!(
1142                "saddle-writer-original-{flush}-{}-{}", std::process::id(),
1143                SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos()));
1144            std::fs::create_dir(&directory).unwrap();
1145            let mut emergency = crate::EmergencyDiagnostics::start_checked(
1146                &crate::FileLoggingConfig::new(directory.clone(), crate::Rotation::Daily)).unwrap();
1147            let output = emergency.source_output().unwrap();
1148            let target = emergency.target().to_owned();
1149            let (entered_sender, entered_receiver) = mpsc::channel();
1150            let release = Arc::new((Mutex::new(false), Condvar::new()));
1151            let logger = observer(ChainWriter {
1152                fail_flush: flush,
1153                release_write: (!flush).then(|| (entered_sender, Arc::clone(&release))),
1154            }, 8);
1155            let application = saddle_core::ContextLabel::checked("writer-original-test").unwrap();
1156            assert!(block_on(RecordedComponentLifecycle::start(
1157                &logger, &application, &output)).is_ok());
1158            let primary = saddle_core::Diagnostic::capture(
1159                saddle_core::DiagnosticCategory::UnexpectedError,
1160                saddle_core::CaptureSite::Origin,
1161                saddle_core::DiagnosticCause::new(
1162                    saddle_core::DiagnosticStage::ShutdownComponent,
1163                    saddle_core::DiagnosticCode::new("test.writer_primary").unwrap()),
1164            ).occurrence();
1165            if !flush { logger.emit(LogRecord::new(EventLevel::Info, "writer.failure")); }
1166            if !flush { entered_receiver.recv_timeout(std::time::Duration::from_secs(5)).unwrap(); }
1167            let shutdown = RecordedComponentLifecycle::shutdown_with_primary(
1168                &logger, &application, &output, Some(primary));
1169            if !flush {
1170                let (lock, condition) = &*release;
1171                *lock.lock().unwrap() = true;
1172                condition.notify_all();
1173            }
1174            let failure = block_on(shutdown).err().unwrap();
1175            assert_eq!(failure.safe().code(), "observability.shutdown_failed");
1176            assert!(failure.original_if_unconfirmed().is_none());
1177            let diagnostic = failure.safe().diagnostic().unwrap();
1178            let record = std::fs::read_to_string(&target).unwrap();
1179            let (top, nested) = if flush {
1180                ("flush top original 4933", "flush nested original 4934")
1181            } else { ("writer top original 4931", "writer nested original 4932") };
1182            assert!(record.contains(top), "{record}");
1183            assert!(record.contains(nested), "{record}");
1184            assert!(record.contains(if flush { "Flush" } else { "Record" }), "{record}");
1185            assert!(record.contains("shutdown_logger"), "{record}");
1186            assert!(record.contains(&diagnostic.id().to_string()));
1187            assert!(record.contains("complete"), "{record}");
1188            let originals: Vec<serde_json::Value> = record.lines()
1189                .map(|line| serde_json::from_str(line).unwrap())
1190                .filter(|row: &serde_json::Value|
1191                    row["event"] == "request_error_original").collect();
1192            assert!(!originals.is_empty());
1193            assert!(originals.iter().any(|row|
1194                row["channel"] == "description" && row["cause_depth"] == 0
1195                    && row["payload"].as_str().unwrap().contains(top)));
1196            assert!(originals.iter().any(|row|
1197                row["channel"] == "debug" && row["cause_depth"] == 0
1198                    && row["payload"].as_str().unwrap().contains(top)));
1199            assert!(originals.iter().any(|row|
1200                row["channel"] == "description"
1201                    && row["cause_depth"].as_u64().unwrap() > 0
1202                    && row["payload"].as_str().unwrap().contains(nested)));
1203            assert!(originals.iter().any(|row|
1204                row["channel"] == "terminal"
1205                    && row["state"] == "exposed_chain_complete"));
1206            assert!(originals.iter().all(|row|
1207                row["occurrence"]["diagnostic_id"] == diagnostic.id()));
1208            let expected = serde_json::to_value(primary).unwrap();
1209            assert!(originals.iter().all(|row|
1210                row["occurrence"]["primary_diagnostic_id"]
1211                    == expected["diagnostic_id"]));
1212            drop(logger);
1213            drop(output);
1214            let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
1215            while emergency.shutdown() == crate::DiagnosticShutdown::Pending
1216                && std::time::Instant::now() < deadline { std::thread::yield_now(); }
1217            assert_eq!(emergency.shutdown(), crate::DiagnosticShutdown::Finished);
1218            drop(emergency);
1219            std::fs::remove_dir_all(directory).unwrap();
1220        }
1221    }
1222
1223    #[test]
1224    fn writer_original_without_selected_output_remains_unconfirmed() {
1225        use crate::root_diagnostic::RecordedComponentLifecycle;
1226        let logger = observer(ChainWriter { fail_flush: false, release_write: None }, 8);
1227        logger.emit(LogRecord::new(EventLevel::Info, "writer.no_output"));
1228        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
1229        while logger.inner.writer_source.lock().unwrap_or_else(|e| e.into_inner())
1230            .failure.is_none() && std::time::Instant::now() < deadline {
1231            std::thread::yield_now();
1232        }
1233        assert!(logger.inner.writer_source.lock().unwrap_or_else(|e| e.into_inner())
1234            .failure.as_ref().is_some_and(|failure| failure.written.is_none()));
1235        let directory = std::env::temp_dir().join(format!(
1236            "saddle-writer-unconfirmed-{}-{}", std::process::id(),
1237            SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos()));
1238        std::fs::create_dir(&directory).unwrap();
1239        let mut emergency = crate::EmergencyDiagnostics::start_checked(
1240            &crate::FileLoggingConfig::new(directory.clone(), crate::Rotation::Daily)).unwrap();
1241        let target = emergency.target().to_owned();
1242        let output = emergency.source_output().unwrap();
1243        let application = saddle_core::ContextLabel::checked("writer-unconfirmed-test").unwrap();
1244        let failure = block_on(RecordedComponentLifecycle::shutdown(
1245            &logger, &application, &output)).err().unwrap();
1246        assert_eq!(failure.safe().code(), "observability.shutdown_failed");
1247        let original = failure.original_if_unconfirmed().unwrap();
1248        assert!(original.to_string().contains("writer top original 4931"));
1249        assert!(original.source().unwrap().to_string()
1250            .contains("writer top original 4931"));
1251        assert!(!std::fs::read_to_string(&target).unwrap()
1252            .contains("writer top original 4931"));
1253        drop(logger);
1254        drop(output);
1255        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
1256        while emergency.shutdown() == crate::DiagnosticShutdown::Pending
1257            && std::time::Instant::now() < deadline { std::thread::yield_now(); }
1258        assert_eq!(emergency.shutdown(), crate::DiagnosticShutdown::Finished);
1259        drop(emergency);
1260        std::fs::remove_dir_all(directory).unwrap();
1261    }
1262
1263    struct ThreadWaker(thread::Thread);
1264
1265    impl Wake for ThreadWaker {
1266        fn wake(self: Arc<Self>) {
1267            self.0.unpark();
1268        }
1269    }
1270
1271    fn block_on<T>(future: impl Future<Output = T>) -> T {
1272        let waker = Waker::from(Arc::new(ThreadWaker(thread::current())));
1273        let mut context = Context::from_waker(&waker);
1274        let mut future = std::pin::pin!(future);
1275        loop {
1276            match future.as_mut().poll(&mut context) {
1277                Poll::Ready(output) => return output,
1278                Poll::Pending => thread::park(),
1279            }
1280        }
1281    }
1282
1283    struct BlockingWriter {
1284        entered: Option<mpsc::Sender<()>>,
1285        release: Arc<(Mutex<bool>, Condvar)>,
1286        fail_after_release: bool,
1287    }
1288
1289    impl Write for BlockingWriter {
1290        fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
1291            if let Some(entered) = self.entered.take() {
1292                let _ = entered.send(());
1293                let (lock, condition) = &*self.release;
1294                let mut released = lock.lock().unwrap();
1295                while !*released {
1296                    released = condition.wait(released).unwrap();
1297                }
1298                if self.fail_after_release {
1299                    return Err(io::Error::other("injected writer failure"));
1300                }
1301            }
1302            Ok(bytes.len())
1303        }
1304
1305        fn flush(&mut self) -> io::Result<()> {
1306            Ok(())
1307        }
1308    }
1309
1310    struct FailingWriter {
1311        fail_write: Option<usize>,
1312        writes: usize,
1313        fail_flush: bool,
1314    }
1315
1316    impl Write for FailingWriter {
1317        fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
1318            self.writes += 1;
1319            if self.fail_write == Some(self.writes) {
1320                Err(io::Error::other("injected writer failure"))
1321            } else {
1322                Ok(bytes.len())
1323            }
1324        }
1325
1326        fn flush(&mut self) -> io::Result<()> {
1327            if self.fail_flush {
1328                Err(io::Error::other("injected flush failure"))
1329            } else {
1330                Ok(())
1331            }
1332        }
1333    }
1334
1335    fn observer(writer: impl Write + Send + 'static, capacity: usize) -> Observer {
1336        Observer::with_writer_and_seed(
1337            ObserverConfig {
1338                queue_capacity: capacity,
1339            },
1340            writer,
1341            Ok([7; 24]),
1342        )
1343        .unwrap()
1344    }
1345
1346    fn context_profile_process() -> saddle_admission::ProfuseGwLightweightProcessOwner {
1347        use saddle_admission::{StorageDemand, bind_deployment_resource_budget_bootstrap, freeze_deployment_resource_budget, prepare_profusegw_lightweight_profile};
1348        use saddle_core::{BootstrapRendezvousIssuer, GeneratedApplicationFreezeSource, ListenerStartupFreezeSource, pair_bootstrap_rendezvous};
1349        let pending = freeze_deployment_resource_budget(
1350            1, 32, 5_000, 1, 1, 1_000_000, 1_000_000, 1_000_000, 1_000_000,
1351        ).unwrap().reserve_known_database_memory(0, 0, 0).ok().unwrap();
1352        let (application, listener) = BootstrapRendezvousIssuer::issue()
1353            .freeze_application(GeneratedApplicationFreezeSource::new(
1354                "app", b"descriptor", &["route"],
1355            )).unwrap();
1356        let listener = listener.freeze_listener(ListenerStartupFreezeSource::new(
1357            "app", "127.0.0.1:8000".parse().unwrap(),
1358            "127.0.0.1:9000".parse().unwrap(),
1359            std::time::Duration::from_millis(5_000),
1360        )).ok().unwrap();
1361        let (whole, receipt) = pair_bootstrap_rendezvous(application, listener).ok().unwrap();
1362        let budget = bind_deployment_resource_budget_bootstrap(pending, whole, receipt).ok().unwrap();
1363        let process = prepare_profusegw_lightweight_profile(budget).ok().unwrap();
1364        process
1365    }
1366
1367    #[test]
1368    fn large_managed_context_and_summary_reassemble_after_request_close() {
1369        use std::collections::BTreeMap;
1370        use saddle_core::{ContextFact, ContextLabel, RequestIdentityGroup, RequestLocalFacts, RequestRootPublisher, RequestViewPhase};
1371        let directory = std::env::temp_dir().join(format!("context-segments-{}-{}", std::process::id(), SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos()));
1372        let blob = "界\n\"".repeat(12000);
1373        let raw = format!("{{\"requestData\":{{\"flag\":1,\"future\":[null,123456789012345678901234567890,{{\"blob\":{}}}]}},\"profuseGwContext\":{{\"traceInfo\":{{\"traceId\":\"trace\",\"rpcId\":\"0\"}},\"ldcInfo\":{{\"unknown\":[false,null]}}}}}}", serde_json::to_string(&blob).unwrap());
1374        let mut root = RequestRootPublisher::create(ContextLabel::checked("app").unwrap(), ContextFact::NotEstablished).unwrap();
1375        let call = saddle_core::CallContext::new("app".into(), "profusegw".into(), "profusegw".into(), "route".into(), TraceId::from_u128(1), SpanId::from_u64(1));
1376        root.publish(RequestIdentityGroup::from_validated(&call, "request", "route", 1, ContextFact::Unavailable).unwrap()).ok().unwrap();
1377        let process = context_profile_process();
1378        let storage_root = process.try_process_storage(saddle_admission::StorageDemand::separate(&[]).unwrap()).unwrap();
1379        let baseline = process.resource_snapshot().framework_charged;
1380        let execution = match process.verified_profile().try_admit() { saddle_admission::ProfuseGwLightweightAdmissionOutcome::Ready(p) => p.into_execution(), _ => panic!("original capacity") };
1381        let memory = execution.request_memory();
1382        let bytes = memory.try_bytes(raw.as_bytes()).unwrap();
1383        let document = memory.decode_input_range(&bytes, 0..bytes.len()).unwrap().into_shared().unwrap();
1384        drop(bytes);
1385        root.publish_ingress(Arc::new(document), "call", 1234).unwrap();
1386        let business = memory.update_business(None, "/private", Some(br#"{"child":"hidden"}"#), true, saddle_admission::BusinessLimits::default()).unwrap();
1387        let view = root.reference().view(RequestLocalFacts::new(RequestViewPhase::Reading)).with_business(Arc::new(business));
1388        let mut config = crate::ContextLoggingConfig::default();
1389        config.summaries.push(crate::SummaryProfile { name: "selected".into(), enabled: true, format: crate::SummaryFormat::Json,
1390            fields: BTreeMap::from([("data".into(), "/context/inputInfo/requestData".into()), ("secret".into(), "/context/business/private/child".into()), ("revision".into(), "/context/_meta/revision".into())]),
1391            missing: crate::MissingField::Null, template: None });
1392        let file = FileLoggingConfig::new(&directory, crate::Rotation::Daily).with_context(config.clone()).unwrap();
1393        let outputs = CalendarOutputs { summaries: vec![Some(CalendarFileWriter::open(file.for_summary("selected")).unwrap())], main: CalendarFileWriter::open(file).unwrap() };
1394        let observer = Observer::with_sink_and_seed(ObserverConfig { queue_capacity: 32 }, outputs, Ok([7; 24]), config).unwrap();
1395        observer.install_process_log_storage(process.process_log_storage());
1396        let mut owned_context = None;
1397        drop(memory.framework_output(|_| {
1398            owned_context = Some(observer.capture_terminal_context(&view.unified_context()).unwrap());
1399            Ok(())
1400        }).unwrap());
1401        let available = process.resource_snapshot().framework_capacity - process.resource_snapshot().framework_charged;
1402        let filler = process.try_test_process_storage_child(saddle_admission::StorageDemand::separate(&[(
1403            std::alloc::Layout::array::<u8>(available - 64 - std::mem::size_of::<saddle_admission::StoragePermit>()).unwrap(), 1,
1404        )]).unwrap()).unwrap();
1405        let charged = process.resource_snapshot().framework_charged;
1406        drop(memory.framework_output(|_| {
1407            assert!(matches!(observer.capture_terminal_context(&view.unified_context()), Err(crate::DiagnosticSubmission::Full)));
1408            Ok(())
1409        }).unwrap());
1410        assert_eq!(process.resource_snapshot().framework_charged, charged);
1411        drop(filler);
1412        let emitted = memory.framework_output(|_| Ok(crate::RootDiagnosticScope::new(&view, None).ordinary(&observer, crate::RootRequestEvent::Handler, Default::default()))).unwrap();
1413        assert_eq!(*emitted.get(), crate::DiagnosticSubmission::Enqueued);
1414        drop(emitted);
1415        // Block the worker and admit exactly one prefix frame. Rejection must
1416        // not manufacture a complete marker or a transient serde heap error.
1417        struct CapturedBlocking { block: BlockingWriter, data: Arc<Mutex<Vec<u8>>> }
1418        impl Write for CapturedBlocking {
1419            fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
1420                self.block.write(bytes)?; self.data.lock().unwrap().extend_from_slice(bytes); Ok(bytes.len())
1421            }
1422            fn flush(&mut self) -> io::Result<()> { self.block.flush() }
1423        }
1424        let (entered, entered_rx) = mpsc::channel();
1425        let release = Arc::new((Mutex::new(false), Condvar::new()));
1426        let captured = Arc::new(Mutex::new(Vec::new()));
1427        let refused = Observer::with_sink_and_seed(ObserverConfig { queue_capacity: 1 }, BasicSink(CapturedBlocking {
1428            block: BlockingWriter { entered: Some(entered), release: release.clone(), fail_after_release: false }, data: captured.clone(),
1429        }), Ok([8; 24]), crate::ContextLoggingConfig::default()).unwrap();
1430        refused.install_process_log_storage(process.process_log_storage());
1431        assert_eq!(refused.emit_root_record(&0), crate::DiagnosticSubmission::Enqueued);
1432        entered_rx.recv_timeout(std::time::Duration::from_secs(2)).unwrap();
1433        let rejection = memory.framework_output(|_| Ok(crate::RootDiagnosticScope::new(&view, None).ordinary(&refused, crate::RootRequestEvent::Handler, Default::default()))).unwrap();
1434        assert_eq!(*rejection.get(), crate::DiagnosticSubmission::Full); drop(rejection);
1435        assert_eq!(refused.dropped_events(), 1);
1436        *release.0.lock().unwrap() = true; release.1.notify_all();
1437        assert!(matches!(block_on(refused.flush()), Err(FlushError::DroppedEvents(_))));
1438        let frames: Vec<serde_json::Value> = std::str::from_utf8(&captured.lock().unwrap()).unwrap().lines()
1439            .filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
1440            .filter(|row| row["event"] == "context_record_segment").collect();
1441        assert_eq!(frames.len(), 1); assert_eq!(frames[0]["sequence"], 0); assert_eq!(frames[0]["state"], "continuation");
1442        let _ = block_on(refused.shutdown()); drop(refused);
1443        drop((view, root, memory));
1444        let terminal = execution.cancel_observed();
1445        assert_eq!(terminal.audit().report.unwrap().escape_allocations, 0);
1446        assert_eq!(process.resource_snapshot().active_accounts, 0);
1447        block_on(observer.flush()).unwrap();
1448        fn readback(path: std::path::PathBuf) -> String {
1449            let text = std::fs::read_to_string(path).unwrap(); let mut restored = String::new(); let mut id = None; let mut complete = false;
1450            for (sequence, line) in text.lines().enumerate() {
1451                assert!(line.len() <= 8191);
1452                let frame: serde_json::Value = serde_json::from_str(line).unwrap();
1453                assert_eq!(frame["event"], "context_record_segment");
1454                assert_eq!(frame["format"], "json");
1455                assert_eq!(frame["sequence"], sequence);
1456                if let Some(id) = &id { assert_eq!(&frame["recordId"], id); } else { id = Some(frame["recordId"].clone()); }
1457                assert!(!complete); restored.push_str(frame["payload"].as_str().unwrap()); complete = frame["state"] == "record_complete";
1458            }
1459            assert!(complete); assert!(id.is_some()); assert!(restored.contains("123456789012345678901234567890")); restored
1460        }
1461        let full: serde_json::Value = serde_json::from_str(&readback(directory.join("saddle.log"))).unwrap();
1462        let summary: serde_json::Value = serde_json::from_str(&readback(directory.join("selected.summary.log"))).unwrap();
1463        assert_eq!(full["context"]["inputInfo"]["requestData"]["future"][2]["blob"], blob);
1464        assert_eq!(summary["data"], full["context"]["inputInfo"]["requestData"]);
1465        assert_eq!(summary["secret"], "<redacted>"); assert_eq!(full["context"]["business"]["private"], "<redacted>");
1466        assert_eq!(summary["revision"], full["context"]["_meta"]["revision"]);
1467        let owned_context = owned_context.unwrap();
1468        let emergency = crate::EmergencyDiagnostics::start(&FileLoggingConfig::new(&directory, crate::Rotation::Daily)).unwrap();
1469        let diagnostic = saddle_core::BoundedDiagnostic::capture(saddle_core::DiagnosticCategory::UnexpectedError,
1470            saddle_core::CaptureSite::FirstObserved, saddle_core::BoundedDiagnosticCause::new(
1471                saddle_core::DiagnosticStage::FinalizerResource, saddle_core::DiagnosticCode::new("test.owned_audit").unwrap()));
1472        let marker = diagnostic.occurrence();
1473        assert_eq!(crate::root_diagnostic::request_terminal_audit_passed_owned(Some(&emergency.handle()),
1474            Ok(&owned_context), marker, terminal.audit()), crate::DiagnosticSubmission::Written);
1475        assert_eq!(crate::root_diagnostic::request_terminal_audit_cleanup_owned(Some(&emergency.handle()),
1476            Ok(&owned_context), marker, &process.resource_snapshot()), crate::DiagnosticSubmission::Written);
1477        let failure = saddle_admission::ProfuseGwTerminalAuditFailure { error: saddle_admission::AdmissionError::AccountClosed,
1478            audit: None, database: saddle_admission::ProfuseGwDatabaseDisposition::NotUsed };
1479        let written = crate::root_diagnostic::RecordedTerminalFailure::audit_owned(Some(&emergency.handle()),
1480            Ok(&owned_context), Some(marker), failure);
1481        assert!(written.original_confirmed());
1482        assert_eq!(written.disposition_submission(), Some(crate::DiagnosticSubmission::Written));
1483        let failure = saddle_admission::ProfuseGwTerminalAuditFailure { error: saddle_admission::AdmissionError::AccountClosed,
1484            audit: None, database: saddle_admission::ProfuseGwDatabaseDisposition::NotUsed };
1485        let rejected = crate::root_diagnostic::RecordedTerminalFailure::audit_owned(Some(&emergency.handle()),
1486            Err(crate::DiagnosticSubmission::Full), Some(marker), failure);
1487        assert!(!rejected.original_confirmed(), "a written error without its requested context cannot mint a positive receipt");
1488        let text = std::fs::read_to_string(directory.join("saddle.emergency.log")).unwrap();
1489        for event in ["request_audit_passed", "request_audit_cleanup"] {
1490            let frames: Vec<serde_json::Value> = text.lines().map(|line| serde_json::from_str::<serde_json::Value>(line).unwrap())
1491                .filter(|row| row["event"] == event).collect();
1492            assert!(frames.len() > 1);
1493            let mut restored = String::new();
1494            for (sequence, frame) in frames.iter().enumerate() {
1495                assert_eq!(frame["sequence"], sequence);
1496                restored.push_str(frame["payload"].as_str().unwrap());
1497                assert_eq!(frame["state"], if sequence + 1 == frames.len() { "record_complete" } else { "continuation" });
1498            }
1499            assert!(restored.contains("123456789012345678901234567890"));
1500            let record: serde_json::Value = serde_json::from_str(&restored).unwrap();
1501            assert_eq!(record["context"], full["context"]);
1502            assert_eq!(record["recordSchema"], "1.0");
1503        }
1504        drop(emergency);
1505        let encoded = serde_json::to_string(&owned_context).unwrap();
1506        assert!(encoded.contains("123456789012345678901234567890"));
1507        let restored: serde_json::Value = serde_json::from_str(&encoded).unwrap();
1508        assert_eq!(restored, full["context"]);
1509        drop(owned_context);
1510        assert_eq!(process.resource_snapshot().framework_charged, baseline);
1511        block_on(observer.shutdown()).unwrap(); drop(observer); drop(storage_root); process.finish().unwrap();
1512        std::fs::remove_dir_all(directory).unwrap();
1513    }
1514
1515    #[test]
1516    fn configured_context_profiles_write_real_files_for_all_four_combinations() {
1517        use std::collections::BTreeMap;
1518        use crate::{ContextLoggingConfig, SummaryProfile, SummaryFormat, MissingField};
1519        let root = std::env::temp_dir().join(format!("context-profiles-{}-{}", std::process::id(), SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos()));
1520        let publisher = saddle_core::RequestRootPublisher::create(saddle_core::ContextLabel::checked("app").unwrap(), saddle_core::ContextFact::NotEstablished).unwrap();
1521        let view = publisher.reference().view(saddle_core::RequestLocalFacts::new(saddle_core::RequestViewPhase::Reading));
1522        use saddle_admission::StorageDemand;
1523        let process = context_profile_process();
1524        let storage_root = process.try_process_storage(StorageDemand::separate(&[]).unwrap()).unwrap();
1525        let baseline = process.resource_snapshot().framework_charged;
1526        for full in [false, true] { for summary in [false, true] {
1527            let directory = root.join(format!("{full}-{summary}"));
1528            let mut context = ContextLoggingConfig::default(); context.full.enabled = full;
1529            let fields = BTreeMap::from([("event".into(), "/event".into()), ("revision".into(), "/context/_meta/revision".into()),
1530                ("missing".into(), "/context/business/future".into()), ("elapsed".into(), "/payload/elapsedMs".into())]);
1531            context.summaries = vec![
1532                SummaryProfile { name: "json".into(), enabled: summary, format: SummaryFormat::Json, fields: fields.clone(), missing: MissingField::Null, template: None },
1533                SummaryProfile { name: "text".into(), enabled: summary, format: SummaryFormat::Text, fields, missing: MissingField::Omit, template: Some("event=${event} absent=${missing} revision=${revision}".into()) },
1534            ]; context.validate().unwrap();
1535            let file = FileLoggingConfig::new(&directory, crate::Rotation::Daily).with_context(context.clone()).unwrap();
1536            let summaries = context.summaries.iter().map(|profile| if profile.enabled { Some(CalendarFileWriter::open(file.for_summary(&profile.name)).unwrap()) } else { None }).collect();
1537            let outputs = CalendarOutputs { main: CalendarFileWriter::open(file).unwrap(), summaries };
1538            let logger = Observer::with_sink_and_seed(ObserverConfig { queue_capacity: 8 }, outputs, Ok([7; 24]), context).unwrap();
1539            logger.install_process_log_storage(process.process_log_storage());
1540            assert_eq!(crate::RootDiagnosticScope::new(&view, None).ordinary(&logger, crate::RootRequestEvent::Handler, Default::default()), crate::DiagnosticSubmission::Enqueued);
1541            block_on(logger.flush()).unwrap();
1542            let main = std::fs::read_to_string(directory.join("saddle.log")).unwrap();
1543            assert_eq!(!main.is_empty(), full);
1544            if full {
1545                let record: serde_json::Value = serde_json::from_str(main.trim()).unwrap();
1546                assert_eq!(record["recordSchema"], "1.0"); assert_eq!(record["context"]["schemaVersion"], "1.0");
1547                assert_eq!(record["payload"]["stage"], "handler");
1548                assert!(record.get("schema_version").is_none());
1549            }
1550            if summary {
1551                let json: serde_json::Value = serde_json::from_str(std::fs::read_to_string(directory.join("json.summary.log")).unwrap().trim()).unwrap();
1552                assert!(json["missing"].is_null()); assert!(json["elapsed"].is_null());
1553                assert_eq!(json["event"], "request_stage");
1554                let text = std::fs::read_to_string(directory.join("text.summary.log")).unwrap();
1555                assert_eq!(text, format!("event=\"request_stage\" absent= revision={}\n", serde_json::to_string(&json["revision"]).unwrap()));
1556                if full { let main: serde_json::Value = serde_json::from_str(main.trim()).unwrap(); assert_eq!(json["revision"], main["context"]["_meta"]["revision"]); }
1557            } else { assert!(!directory.join("json.summary.log").exists()); }
1558            block_on(logger.shutdown()).unwrap(); drop(logger);
1559            assert_eq!(process.resource_snapshot().framework_charged, baseline);
1560        } }
1561        drop(storage_root); process.finish().unwrap();
1562        std::fs::remove_dir_all(root).unwrap();
1563    }
1564
1565    #[test]
1566    fn managed_call_log_outlives_request_account_while_writer_is_blocked() {
1567        // The new process-owned variant adds no command-ring slot bytes:
1568        // the unchanged ordinary LogRecord still determines Command's size.
1569        assert_eq!(std::mem::size_of::<Command>(), std::mem::size_of::<LogRecord>());
1570        assert_eq!(std::mem::size_of::<Command>(), 192);
1571        assert_eq!(ObserverConfig::default().queue_capacity * (8 + std::mem::size_of::<Command>()),
1572            1_638_400);
1573        use saddle_admission::{
1574            ProfuseGwLightweightObservedAdmissionOutcome, StorageDemand,
1575            bind_deployment_resource_budget_bootstrap, freeze_deployment_resource_budget,
1576            prepare_profusegw_lightweight_profile,
1577        };
1578        use saddle_core::{
1579            BootstrapRendezvousIssuer, GeneratedApplicationFreezeSource,
1580            ListenerStartupFreezeSource, RpcCorrelationId, pair_bootstrap_rendezvous,
1581        };
1582        let pending = freeze_deployment_resource_budget(
1583            1, 32, 5_000, 1, 1, 1_000_000, 1_000_000, 1_000_000, 1_000_000,
1584        ).unwrap().reserve_known_database_memory(0, 0, 0).ok().unwrap();
1585        let (application, listener) = BootstrapRendezvousIssuer::issue()
1586            .freeze_application(GeneratedApplicationFreezeSource::new(
1587                "app", b"descriptor", &["route"],
1588            )).unwrap();
1589        let listener = listener.freeze_listener(ListenerStartupFreezeSource::new(
1590            "app", "127.0.0.1:8000".parse().unwrap(),
1591            "127.0.0.1:9000".parse().unwrap(),
1592            std::time::Duration::from_millis(5_000),
1593        )).ok().unwrap();
1594        let (whole, receipt) = pair_bootstrap_rendezvous(application, listener).ok().unwrap();
1595        let budget = bind_deployment_resource_budget_bootstrap(pending, whole, receipt).ok().unwrap();
1596        let process = prepare_profusegw_lightweight_profile(budget).ok().unwrap();
1597        let root = process.try_process_storage(StorageDemand::separate(&[(
1598            std::alloc::Layout::array::<u8>(16).unwrap(), 1,
1599        )]).unwrap()).unwrap();
1600        let baseline = process.resource_snapshot().framework_charged;
1601        let (entered, entered_rx) = mpsc::channel();
1602        let release = Arc::new((Mutex::new(false), Condvar::new()));
1603        let logger = observer(BlockingWriter {
1604            entered: Some(entered), release: Arc::clone(&release), fail_after_release: false,
1605        }, 2);
1606        logger.install_process_log_storage(process.process_log_storage());
1607        let (call, _) = logger.start_managed_external_call_with_rpc(
1608            "app", "zone", "interface", "operation", None,
1609            RpcCorrelationId::new("0").unwrap(),
1610        ).unwrap();
1611        entered_rx.recv_timeout(std::time::Duration::from_secs(2)).unwrap();
1612        let admitted = match process.verified_profile().try_admit_observed() {
1613            ProfuseGwLightweightObservedAdmissionOutcome::Ready(permit, _) => permit,
1614            _ => panic!("fixture admission failed"),
1615        };
1616        let original = saddle_core::SaddleError::new(
1617            saddle_core::ErrorKind::Internal, "original.failure", "original",
1618        );
1619        call.fail(&original);
1620        assert_eq!(original.code(), "original.failure");
1621        admitted.into_execution().cancel_observed();
1622        assert_eq!(process.resource_snapshot().active_accounts, 0);
1623        assert!(process.resource_snapshot().framework_charged > baseline,
1624            "queued process log must remain charged after request cleanup");
1625        // The worker holds the first record while the second is queued.
1626        // Additional managed call records must be dropped synchronously.
1627        let (full, _) = logger.start_managed_external_call_with_rpc(
1628            "app", "zone", "interface", "full", None,
1629            RpcCorrelationId::new("1").unwrap(),
1630        ).unwrap();
1631        full.fail(&original);
1632        assert!(logger.dropped_events() > 0);
1633        *release.0.lock().unwrap() = true;
1634        release.1.notify_all();
1635        assert!(matches!(block_on(logger.flush()), Err(FlushError::DroppedEvents(_))));
1636        assert_eq!(process.resource_snapshot().framework_charged, baseline);
1637
1638        // Exhaust the *same* framework domain with a physical child owner;
1639        // encoding must fail to reserve before allocating and remain lossy.
1640        let available = process.resource_snapshot().framework_capacity
1641            - process.resource_snapshot().framework_charged;
1642        let filler = process.try_test_process_storage_child(StorageDemand::separate(&[(
1643            std::alloc::Layout::array::<u8>(
1644                available - 64 - std::mem::size_of::<saddle_admission::StoragePermit>(),
1645            ).unwrap(), 1,
1646        )]).unwrap()).unwrap();
1647        let charged = process.resource_snapshot().framework_charged;
1648        let before_drop = logger.dropped_events();
1649        let (short, _) = logger.start_managed_external_call_with_rpc(
1650            "app", "zone", "interface", "short", None,
1651            RpcCorrelationId::new("2").unwrap(),
1652        ).unwrap();
1653        short.fail(&original);
1654        assert!(logger.dropped_events() >= before_drop + 2);
1655        assert_eq!(process.resource_snapshot().framework_charged, charged);
1656        drop(filler);
1657        assert!(matches!(block_on(logger.shutdown()), Err(FlushError::DroppedEvents(_))));
1658        let after_shutdown = process.resource_snapshot().framework_charged;
1659        let (closed, _) = logger.start_managed_external_call_with_rpc(
1660            "app", "zone", "interface", "closed", None,
1661            RpcCorrelationId::new("3").unwrap(),
1662        ).unwrap();
1663        closed.fail(&original);
1664        assert_eq!(process.resource_snapshot().framework_charged, after_shutdown);
1665        assert_eq!(process.resource_snapshot().framework_charged, baseline);
1666        drop(logger);
1667        drop(root);
1668        process.finish().unwrap();
1669    }
1670
1671    #[test]
1672    fn root_contract_ordinary_uses_same_view_and_bounded_queue() {
1673        use crate::{
1674            DiagnosticSubmission, RootDiagnosticScope, RootOutcomeFacts, RootRequestEvent,
1675        };
1676        use saddle_core::{
1677            ContextFact, ContextLabel, RequestLocalFacts, RequestRootPublisher, RequestViewPhase,
1678        };
1679        let root = RequestRootPublisher::create(
1680            ContextLabel::checked("app").unwrap(),
1681            ContextFact::NotEstablished,
1682        )
1683        .unwrap();
1684        let view = root
1685            .reference()
1686            .view(RequestLocalFacts::new(RequestViewPhase::Handler));
1687        let expected = serde_json::to_value(&view).unwrap();
1688        let mut frame = crate::diagnostic::FixedDiagnosticBytes {
1689            bytes: [0; 8192],
1690            len: 0,
1691        };
1692        serde_json::to_writer(&mut frame, &view).unwrap();
1693        let ordinary_context: serde_json::Value =
1694            serde_json::from_slice(&frame.bytes[..frame.len]).unwrap();
1695        assert_eq!(ordinary_context, expected);
1696        let (entered, wait) = mpsc::channel();
1697        let release = Arc::new((Mutex::new(false), Condvar::new()));
1698        struct Captured {
1699            writer: BlockingWriter,
1700            bytes: Arc<Mutex<Vec<u8>>>,
1701        }
1702        impl Write for Captured {
1703            fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
1704                let n = self.writer.write(bytes)?;
1705                self.bytes.lock().unwrap().extend_from_slice(&bytes[..n]);
1706                Ok(n)
1707            }
1708            fn flush(&mut self) -> io::Result<()> {
1709                self.writer.flush()
1710            }
1711        }
1712        struct Unblock(Arc<(Mutex<bool>, Condvar)>);
1713        impl Drop for Unblock {
1714            fn drop(&mut self) {
1715                *self.0.0.lock().unwrap() = true;
1716                self.0.1.notify_all();
1717            }
1718        }
1719        let unblock = Unblock(Arc::clone(&release));
1720        let bytes = Arc::new(Mutex::new(Vec::new()));
1721        let logger = observer(
1722            Captured {
1723                writer: BlockingWriter {
1724                    entered: Some(entered),
1725                    release: Arc::clone(&release),
1726                    fail_after_release: false,
1727                },
1728                bytes: Arc::clone(&bytes),
1729            },
1730            1,
1731        );
1732        let scope = RootDiagnosticScope::new(&view, None);
1733        assert_eq!(
1734            scope.ordinary(
1735                &logger,
1736                RootRequestEvent::Handler,
1737                RootOutcomeFacts::default()
1738            ),
1739            DiagnosticSubmission::Enqueued
1740        );
1741        wait.recv_timeout(std::time::Duration::from_secs(5))
1742            .unwrap();
1743        assert_eq!(
1744            scope.ordinary(
1745                &logger,
1746                RootRequestEvent::Response,
1747                RootOutcomeFacts::default()
1748            ),
1749            DiagnosticSubmission::Enqueued
1750        );
1751        assert_eq!(
1752            scope.ordinary(
1753                &logger,
1754                RootRequestEvent::Response,
1755                RootOutcomeFacts::default()
1756            ),
1757            DiagnosticSubmission::Full
1758        );
1759        drop((view, root));
1760        drop(unblock);
1761        assert!(matches!(
1762            block_on(logger.shutdown()),
1763            Err(FlushError::DroppedEvents(1))
1764        ));
1765        let bytes = bytes.lock().unwrap();
1766        let rows: Vec<serde_json::Value> = serde_json::Deserializer::from_slice(&bytes)
1767            .into_iter()
1768            .map(Result::unwrap)
1769            .collect();
1770        assert_eq!(rows.len(), 2);
1771        assert!(rows.iter().all(|row| row["context"] == expected));
1772    }
1773
1774    #[test]
1775    fn root_contract_encoded_ordinary_readback_and_failure_use_original_writer() {
1776        let record = serde_json::json!({"event":"request_stage", "context":{"local_request":5}});
1777        let (tx, rx) = mpsc::sync_channel(1);
1778        tx.try_send(Command::RootRecord {
1779            bytes: serde_json::to_vec(&record).unwrap(),
1780            dropped: 0,
1781        })
1782        .unwrap_or_else(|_| panic!("empty queue"));
1783        drop(tx);
1784        let mut bytes = Vec::new();
1785        write_records(rx, BasicSink(&mut bytes), Arc::new(Mutex::new(WriterSourceState::default())));
1786        assert_eq!(
1787            serde_json::from_slice::<serde_json::Value>(&bytes).unwrap(),
1788            record
1789        );
1790        let logger = observer(
1791            FailingWriter {
1792                fail_write: Some(1),
1793                writes: 0,
1794                fail_flush: false,
1795            },
1796            1,
1797        );
1798        assert_eq!(
1799            logger.emit_root_record(&record),
1800            crate::DiagnosticSubmission::Enqueued
1801        );
1802        assert!(matches!(
1803            block_on(logger.shutdown()),
1804            Err(FlushError::OutputFailed(OutputStage::Record))
1805        ));
1806    }
1807
1808    #[test]
1809    fn random_source_failure_is_reported_during_initialization() {
1810        let result = Observer::with_writer_and_seed(
1811            ObserverConfig::default(),
1812            io::sink(),
1813            Err(InitError::RandomSource),
1814        );
1815        assert!(matches!(result, Err(InitError::RandomSource)));
1816    }
1817
1818    #[test]
1819    fn initialized_generator_produces_non_zero_unique_ids_without_more_io() {
1820        let observer = observer(io::sink(), 8);
1821        let first_trace = observer.new_trace_id();
1822        let second_trace = observer.new_trace_id();
1823        let first_span = observer.new_span_id();
1824        let second_span = observer.new_span_id();
1825        assert_ne!(first_trace.as_u128(), 0);
1826        assert_ne!(first_trace, second_trace);
1827        assert_ne!(first_span.as_u64(), 0);
1828        assert_ne!(first_span, second_span);
1829        block_on(observer.shutdown()).unwrap();
1830    }
1831
1832    #[test]
1833    fn observer_is_a_managed_component() {
1834        let observer = observer(io::sink(), 8);
1835        assert_eq!(ComponentLifecycle::name(&observer), "observability");
1836        block_on(ComponentLifecycle::start(&observer)).unwrap();
1837        block_on(ComponentLifecycle::shutdown(&observer)).unwrap();
1838    }
1839
1840    #[test]
1841    fn full_queue_is_reported_by_shutdown_without_waiting_in_emit() {
1842        let (entered_sender, entered_receiver) = mpsc::channel();
1843        let release = Arc::new((Mutex::new(false), Condvar::new()));
1844        let writer = BlockingWriter {
1845            entered: Some(entered_sender),
1846            release: release.clone(),
1847            fail_after_release: false,
1848        };
1849        let observer = observer(writer, 1);
1850
1851        observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
1852        entered_receiver.recv().unwrap();
1853        observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
1854        observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
1855        assert_eq!(observer.dropped_events(), 1);
1856
1857        let shutdown = observer.shutdown();
1858        let (lock, condition) = &*release;
1859        *lock.lock().unwrap() = true;
1860        condition.notify_one();
1861        assert_eq!(block_on(shutdown), Err(FlushError::DroppedEvents(1)));
1862    }
1863
1864    #[test]
1865    fn shutdown_reports_writer_failure_and_unreported_drops_together() {
1866        let (entered_sender, entered_receiver) = mpsc::channel();
1867        let release = Arc::new((Mutex::new(false), Condvar::new()));
1868        let writer = BlockingWriter {
1869            entered: Some(entered_sender),
1870            release: release.clone(),
1871            fail_after_release: true,
1872        };
1873        let observer = observer(writer, 1);
1874
1875        observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
1876        entered_receiver.recv().unwrap();
1877        observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
1878        observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
1879        let shutdown = observer.shutdown();
1880
1881        let (lock, condition) = &*release;
1882        *lock.lock().unwrap() = true;
1883        condition.notify_one();
1884        assert_eq!(
1885            block_on(shutdown),
1886            Err(FlushError::OutputFailedAndDropped {
1887                stage: OutputStage::Record,
1888                dropped_events: 1,
1889            })
1890        );
1891    }
1892
1893    #[test]
1894    fn record_write_failure_persists_until_shutdown() {
1895        let observer = observer(
1896            FailingWriter {
1897                fail_write: Some(1),
1898                writes: 0,
1899                fail_flush: false,
1900            },
1901            8,
1902        );
1903        observer.emit(LogRecord::new(EventLevel::Info, "test.record"));
1904        assert_eq!(
1905            block_on(observer.shutdown()),
1906            Err(FlushError::OutputFailed(OutputStage::Record))
1907        );
1908    }
1909
1910    #[test]
1911    fn newline_failure_persists_until_flush() {
1912        let observer = observer(
1913            FailingWriter {
1914                fail_write: Some(2),
1915                writes: 0,
1916                fail_flush: false,
1917            },
1918            8,
1919        );
1920        observer.emit(LogRecord::new(EventLevel::Info, "test.newline"));
1921        assert_eq!(
1922            block_on(observer.flush()),
1923            Err(FlushError::OutputFailed(OutputStage::Newline))
1924        );
1925        let _ = block_on(observer.shutdown());
1926    }
1927
1928    #[test]
1929    fn flush_failure_is_reported() {
1930        let observer = observer(
1931            FailingWriter {
1932                fail_write: None,
1933                writes: 0,
1934                fail_flush: true,
1935            },
1936            8,
1937        );
1938        assert_eq!(
1939            block_on(observer.flush()),
1940            Err(FlushError::OutputFailed(OutputStage::Flush))
1941        );
1942        assert_eq!(
1943            block_on(observer.shutdown()),
1944            Err(FlushError::OutputFailed(OutputStage::Flush))
1945        );
1946    }
1947}