Skip to main content

saddle_observability/
logger.rs

1use std::{
2    error::Error,
3    fmt,
4    future::Future,
5    io::{self, Write},
6    pin::Pin,
7    sync::{
8        Arc, Mutex, OnceLock,
9        atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering},
10        mpsc::{self, Receiver, SyncSender, TrySendError},
11    },
12    task::{Context, Poll, Waker},
13    thread::{self, JoinHandle},
14    time::{SystemTime, UNIX_EPOCH},
15};
16
17use saddle_core::{ComponentLifecycle, ErrorKind, LifecycleFuture, SaddleError, SpanId, TraceId};
18use serde::Serialize;
19
20use crate::{
21    calendar_writer::{CalendarFileWriter, FileLoggingConfig},
22    event::EventLevel,
23};
24
25static GLOBAL: OnceLock<Observer> = OnceLock::new();
26
27/// Configuration for the bounded non-blocking logger.
28#[derive(Clone, Copy, Debug, Eq, PartialEq)]
29pub struct ObserverConfig {
30    /// Maximum records waiting for the output worker. New records are dropped
31    /// when the queue is full rather than blocking an async request.
32    pub queue_capacity: usize,
33}
34
35impl Default for ObserverConfig {
36    fn default() -> Self {
37        Self {
38            queue_capacity: 8_192,
39        }
40    }
41}
42
43#[derive(Debug)]
44pub enum InitError {
45    EmptyQueue,
46    RandomSource,
47    Output(io::Error),
48    Spawn(io::Error),
49}
50
51impl fmt::Display for InitError {
52    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
53        match self {
54            Self::EmptyQueue => formatter.write_str("log queue capacity must be greater than zero"),
55            Self::RandomSource => {
56                formatter.write_str("operating system random source is unavailable")
57            }
58            Self::Output(error) => write!(formatter, "failed to initialize log output: {error}"),
59            Self::Spawn(error) => write!(formatter, "failed to start log writer: {error}"),
60        }
61    }
62}
63
64impl Error for InitError {
65    fn source(&self) -> Option<&(dyn Error + 'static)> {
66        match self {
67            Self::Spawn(error) | Self::Output(error) => Some(error),
68            Self::EmptyQueue | Self::RandomSource => None,
69        }
70    }
71}
72
73/// The output operation whose first failure made log delivery unhealthy.
74#[derive(Clone, Copy, Debug, Eq, PartialEq)]
75pub enum OutputStage {
76    Serialize,
77    Record,
78    Newline,
79    Flush,
80}
81
82/// A flush or managed shutdown result that is safe to expose to framework code.
83#[derive(Clone, Debug, Eq, PartialEq)]
84pub enum FlushError {
85    OutputFailed(OutputStage),
86    DroppedEvents(u64),
87    OutputFailedAndDropped {
88        stage: OutputStage,
89        dropped_events: u64,
90    },
91    WriterStopped,
92    AlreadyShuttingDown,
93    CoordinatorUnavailable,
94    WorkerPanicked,
95}
96
97impl fmt::Display for FlushError {
98    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
99        match self {
100            Self::OutputFailed(stage) => write!(formatter, "log output failed during {stage:?}"),
101            Self::DroppedEvents(count) => write!(formatter, "{count} log events were dropped"),
102            Self::OutputFailedAndDropped {
103                stage,
104                dropped_events,
105            } => write!(
106                formatter,
107                "log output failed during {stage:?} and {dropped_events} events were dropped"
108            ),
109            Self::WriterStopped => formatter.write_str("log writer is not running"),
110            Self::AlreadyShuttingDown => formatter.write_str("log writer shutdown already started"),
111            Self::CoordinatorUnavailable => {
112                formatter.write_str("failed to start log lifecycle coordinator")
113            }
114            Self::WorkerPanicked => formatter.write_str("log writer thread panicked"),
115        }
116    }
117}
118
119impl Error for FlushError {}
120
121/// A non-blocking future returned by flush and shutdown lifecycle operations.
122pub struct FlushFuture {
123    completion: Arc<Completion>,
124}
125
126impl Future for FlushFuture {
127    type Output = Result<(), FlushError>;
128
129    fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
130        if let Some(result) = self.completion.result.lock().unwrap().take() {
131            return Poll::Ready(result);
132        }
133        *self.completion.waker.lock().unwrap() = Some(context.waker().clone());
134        if let Some(result) = self.completion.result.lock().unwrap().take() {
135            Poll::Ready(result)
136        } else {
137            Poll::Pending
138        }
139    }
140}
141
142struct Completion {
143    result: Mutex<Option<Result<(), FlushError>>>,
144    waker: Mutex<Option<Waker>>,
145}
146
147impl Completion {
148    fn pending() -> Arc<Self> {
149        Arc::new(Self {
150            result: Mutex::new(None),
151            waker: Mutex::new(None),
152        })
153    }
154
155    fn ready(result: Result<(), FlushError>) -> Arc<Self> {
156        Arc::new(Self {
157            result: Mutex::new(Some(result)),
158            waker: Mutex::new(None),
159        })
160    }
161
162    fn complete(&self, result: Result<(), FlushError>) {
163        *self.result.lock().unwrap() = Some(result);
164        if let Some(waker) = self.waker.lock().unwrap().take() {
165            waker.wake();
166        }
167    }
168}
169
170/// Process-level structured logger and trace correlator.
171///
172/// Cloning this value is cheap. Emission uses `try_send`; output and JSON
173/// serialization happen on the dedicated writer thread.
174#[derive(Clone)]
175pub struct Observer {
176    pub(crate) inner: Arc<Inner>,
177}
178
179pub(crate) struct Inner {
180    sender: SyncSender<Command>,
181    dropped: AtomicU64,
182    accepting: AtomicBool,
183    emitting: AtomicUsize,
184    shutdown_started: AtomicBool,
185    worker: Mutex<Option<JoinHandle<()>>>,
186    process_log: Mutex<Option<saddle_admission::ProcessLogStorage>>,
187    ids: IdGenerator,
188    pub(crate) metrics: crate::metrics::Metrics,
189}
190
191struct IdGenerator {
192    trace_high: u64,
193    trace_seed: u64,
194    trace_counter: AtomicU64,
195    span_seed: u64,
196    span_counter: AtomicU64,
197}
198
199impl IdGenerator {
200    fn from_seed(seed: [u8; 24]) -> Self {
201        let mut high = u64::from_be_bytes(seed[0..8].try_into().unwrap());
202        if high == 0 {
203            high = 1;
204        }
205        Self {
206            trace_high: high,
207            trace_seed: u64::from_be_bytes(seed[8..16].try_into().unwrap()),
208            trace_counter: AtomicU64::new(0),
209            span_seed: u64::from_be_bytes(seed[16..24].try_into().unwrap()),
210            span_counter: AtomicU64::new(0),
211        }
212    }
213
214    fn trace_id(&self) -> TraceId {
215        let low = self
216            .trace_seed
217            .wrapping_add(self.trace_counter.fetch_add(1, Ordering::Relaxed));
218        TraceId::from_u128((u128::from(self.trace_high) << 64) | u128::from(low))
219    }
220
221    fn span_id(&self) -> SpanId {
222        loop {
223            let value = self
224                .span_seed
225                .wrapping_add(self.span_counter.fetch_add(1, Ordering::Relaxed));
226            if value != 0 {
227                return SpanId::from_u64(value);
228            }
229        }
230    }
231}
232
233enum Command {
234    Record(LogRecord),
235    ProcessCall {
236        bytes: saddle_admission::ExactStored<Vec<u8>>,
237        dropped: u64,
238    },
239    // Source-encoded, bounded and self-contained. Never contains a request root.
240    RootRecord {
241        bytes: Vec<u8>,
242        dropped: u64,
243    },
244    Flush {
245        unreported_dropped: u64,
246        response: mpsc::Sender<WorkerStatus>,
247    },
248    Shutdown {
249        unreported_dropped: u64,
250        response: mpsc::Sender<WorkerStatus>,
251    },
252}
253pub(crate) fn root_queue_layout() -> std::alloc::Layout {
254    std::alloc::Layout::new::<Command>()
255}
256
257#[derive(Clone, Copy, Debug, Default)]
258struct WorkerStatus {
259    failure: Option<OutputStage>,
260    unreported_dropped: u64,
261}
262
263impl WorkerStatus {
264    fn into_result(self) -> Result<(), FlushError> {
265        match (self.failure, self.unreported_dropped) {
266            (None, 0) => Ok(()),
267            (Some(stage), 0) => Err(FlushError::OutputFailed(stage)),
268            (None, dropped_events) => Err(FlushError::DroppedEvents(dropped_events)),
269            (Some(stage), dropped_events) => Err(FlushError::OutputFailedAndDropped {
270                stage,
271                dropped_events,
272            }),
273        }
274    }
275}
276
277#[derive(Serialize)]
278pub(crate) struct LogRecord {
279    pub timestamp_unix_ms: u128,
280    pub level: EventLevel,
281    pub event: &'static str,
282    #[serde(skip_serializing_if = "Option::is_none")]
283    pub trace_id: Option<String>,
284    #[serde(skip_serializing_if = "Option::is_none")]
285    pub span: Option<String>,
286    #[serde(skip_serializing_if = "Option::is_none")]
287    pub span_id: Option<String>,
288    #[serde(skip_serializing_if = "Option::is_none")]
289    pub parent: Option<String>,
290    #[serde(skip_serializing_if = "Option::is_none")]
291    pub parent_span_id: Option<String>,
292    #[serde(flatten)]
293    pub data: serde_json::Map<String, serde_json::Value>,
294    #[serde(skip_serializing_if = "is_zero")]
295    pub dropped_events: u64,
296}
297
298const fn is_zero(value: &u64) -> bool {
299    *value == 0
300}
301
302impl LogRecord {
303    pub(crate) fn new(level: EventLevel, event: &'static str) -> Self {
304        Self {
305            timestamp_unix_ms: SystemTime::now()
306                .duration_since(UNIX_EPOCH)
307                .unwrap_or_default()
308                .as_millis(),
309            level,
310            event,
311            trace_id: None,
312            span: None,
313            span_id: None,
314            parent: None,
315            parent_span_id: None,
316            data: serde_json::Map::new(),
317            dropped_events: 0,
318        }
319    }
320}
321
322/// Initializes stdout logging once. Later calls return the first observer.
323pub fn init(config: ObserverConfig) -> Result<&'static Observer, InitError> {
324    if let Some(observer) = GLOBAL.get() {
325        return Ok(observer);
326    }
327
328    let observer = Observer::with_writer(config, io::stdout())?;
329    let _ = GLOBAL.set(observer);
330    Ok(GLOBAL.get().expect("global observer was initialized"))
331}
332
333/// Initializes the process logger with the sole calendar-rotated file sink.
334pub fn init_file(
335    config: ObserverConfig,
336    file: FileLoggingConfig,
337) -> Result<&'static Observer, InitError> {
338    if let Some(observer) = GLOBAL.get() {
339        return Ok(observer);
340    }
341    let writer = CalendarFileWriter::open(file).map_err(InitError::Output)?;
342    let observer = Observer::with_writer(config, writer)?;
343    let _ = GLOBAL.set(observer);
344    Ok(GLOBAL.get().expect("global observer was initialized"))
345}
346
347/// Returns the initialized process observer, if application startup installed it.
348pub fn global() -> Option<&'static Observer> {
349    GLOBAL.get()
350}
351
352impl Observer {
353    /// Bind the original process framework domain before formal routed RPCs
354    /// can emit ordinary call records. The writer never receives this handle.
355    #[doc(hidden)]
356    pub fn install_process_log_storage(&self, storage: saddle_admission::ProcessLogStorage) {
357        let mut slot = self.inner.process_log.lock().unwrap_or_else(|p| p.into_inner());
358        *slot = Some(storage);
359    }
360
361    pub(crate) fn emit_process_call(&self, record: &mut crate::call::ProcessCallRecord<'_>) {
362        struct Counter(usize);
363        impl Write for Counter {
364            fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
365                self.0 = self.0.checked_add(bytes.len()).ok_or(io::ErrorKind::OutOfMemory)?;
366                Ok(bytes.len())
367            }
368            fn flush(&mut self) -> io::Result<()> { Ok(()) }
369        }
370        self.inner.emitting.fetch_add(1, Ordering::AcqRel);
371        let submitted = (|| {
372            if !self.inner.accepting.load(Ordering::Acquire) { return false; }
373            let storage = self.inner.process_log.lock().unwrap_or_else(|p| p.into_inner()).clone();
374            let Some(storage) = storage else { return false; };
375            let dropped = self.inner.dropped.swap(0, Ordering::Relaxed);
376            record.dropped_events = dropped;
377            let mut counter = Counter(0);
378            if serde_json::to_writer(&mut counter, record).is_err() {
379                self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
380                return false;
381            }
382            let Ok(layout) = std::alloc::Layout::array::<u8>(counter.0) else {
383                self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
384                return false;
385            };
386            let Ok(permit) = storage.try_reserve(layout) else {
387                self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
388                return false;
389            };
390            let mut bytes = permit.allocate_exact(layout, || Vec::with_capacity(counter.0));
391            if serde_json::to_writer(bytes.get_mut(), record).is_err() || bytes.get().len() != counter.0 {
392                drop(bytes);
393                self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
394                return false;
395            }
396            match self.inner.sender.try_send(Command::ProcessCall { bytes, dropped }) {
397                Ok(()) => true,
398                Err(TrySendError::Full(command)) => {
399                    drop(command);
400                    self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
401                    false
402                }
403                Err(TrySendError::Disconnected(command)) => {
404                    drop(command);
405                    self.inner.metrics.logger_output(true);
406                    self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
407                    false
408                }
409            }
410        })();
411        if !submitted {
412            self.inner.metrics.logger_dropped(1);
413            self.inner.dropped.fetch_add(1, Ordering::Relaxed);
414        }
415        self.inner.emitting.fetch_sub(1, Ordering::Release);
416    }
417    /// Creates an observer with a framework-owned writer.
418    ///
419    /// Application startup should normally use [`init`]. This constructor is
420    /// useful for embedding and deterministic tests.
421    pub fn with_writer(
422        config: ObserverConfig,
423        writer: impl Write + Send + 'static,
424    ) -> Result<Self, InitError> {
425        let mut seed = [0_u8; 24];
426        getrandom::fill(&mut seed).map_err(|_| InitError::RandomSource)?;
427        Self::with_writer_and_seed(config, writer, Ok(seed))
428    }
429
430    fn with_writer_and_seed(
431        config: ObserverConfig,
432        writer: impl Write + Send + 'static,
433        seed: Result<[u8; 24], InitError>,
434    ) -> Result<Self, InitError> {
435        if config.queue_capacity == 0 {
436            return Err(InitError::EmptyQueue);
437        }
438        let ids = IdGenerator::from_seed(seed?);
439        let (sender, receiver) = mpsc::sync_channel(config.queue_capacity);
440        let worker = thread::Builder::new()
441            .name("saddle-log-writer".to_owned())
442            .spawn(move || write_records(receiver, writer))
443            .map_err(InitError::Spawn)?;
444        Ok(Self {
445            inner: Arc::new(Inner {
446                sender,
447                dropped: AtomicU64::new(0),
448                accepting: AtomicBool::new(true),
449                emitting: AtomicUsize::new(0),
450                shutdown_started: AtomicBool::new(false),
451                worker: Mutex::new(Some(worker)),
452                process_log: Mutex::new(None),
453                ids,
454                metrics: crate::metrics::Metrics::default(),
455            }),
456        })
457    }
458
459    pub(crate) fn new_trace_id(&self) -> TraceId {
460        self.inner.ids.trace_id()
461    }
462
463    pub(crate) fn new_span_id(&self) -> SpanId {
464        self.inner.ids.span_id()
465    }
466
467    pub(crate) fn emit(&self, mut record: LogRecord) {
468        if !self.inner.accepting.load(Ordering::Acquire) {
469            return;
470        }
471        self.inner.emitting.fetch_add(1, Ordering::AcqRel);
472        if !self.inner.accepting.load(Ordering::Acquire) {
473            self.inner.emitting.fetch_sub(1, Ordering::AcqRel);
474            return;
475        }
476
477        record.dropped_events = self.inner.dropped.swap(0, Ordering::Relaxed);
478        match self.inner.sender.try_send(Command::Record(record)) {
479            Ok(()) => {}
480            Err(TrySendError::Full(Command::Record(record))) => {
481                self.inner.metrics.logger_dropped(1);
482                self.inner
483                    .dropped
484                    .fetch_add(record.dropped_events.saturating_add(1), Ordering::Relaxed);
485            }
486            Err(TrySendError::Disconnected(_)) => {
487                self.inner.metrics.logger_dropped(1);
488                self.inner.metrics.logger_output(true);
489                self.inner.dropped.fetch_add(1, Ordering::Relaxed);
490            }
491            Err(TrySendError::Full(_)) => unreachable!("emit only sends records"),
492        }
493        self.inner.emitting.fetch_sub(1, Ordering::Release);
494    }
495
496    pub(crate) fn emit_root_record<T: Serialize>(&self, record: &T) -> crate::DiagnosticSubmission {
497        use crate::DiagnosticSubmission as Submission;
498        self.inner.emitting.fetch_add(1, Ordering::AcqRel);
499        let result = (|| {
500            if !self.inner.accepting.load(Ordering::Acquire) {
501                return Submission::Closed;
502            }
503            let mut frame = crate::diagnostic::FixedDiagnosticBytes {
504                bytes: [0; 8192],
505                len: 0,
506            };
507            if serde_json::to_writer(&mut frame, record).is_err() {
508                return Submission::EncodingFailed;
509            }
510            let dropped = self.inner.dropped.swap(0, Ordering::Relaxed);
511            // Only initialized bytes move into the existing logger domain. This
512            // allocation and source frame peak are separate R0 layout inputs.
513            let bytes = frame.bytes[..frame.len].to_vec();
514            match self
515                .inner
516                .sender
517                .try_send(Command::RootRecord { bytes, dropped })
518            {
519                Ok(()) => Submission::Enqueued,
520                Err(TrySendError::Full(Command::RootRecord { dropped, .. })) => {
521                    self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
522                    Submission::Full
523                }
524                Err(TrySendError::Disconnected(_)) => {
525                    self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
526                    self.inner.metrics.logger_output(true);
527                    Submission::Closed
528                }
529                Err(TrySendError::Full(_)) => unreachable!("root record submission"),
530            }
531        })();
532        if result != Submission::Enqueued {
533            self.inner.dropped.fetch_add(1, Ordering::Relaxed);
534            self.inner.metrics.logger_dropped(1);
535        }
536        self.inner.emitting.fetch_sub(1, Ordering::Release);
537        result
538    }
539
540    /// Returns a fixed-size, allocation-free copy of the current process metrics.
541    pub fn metrics_snapshot(&self) -> crate::MetricsSnapshot {
542        self.inner.metrics.snapshot()
543    }
544
545    /// Flushes records accepted before the lifecycle coordinator reaches the
546    /// writer. The returned future never blocks the async executor thread.
547    pub fn flush(&self) -> FlushFuture {
548        if !self.inner.accepting.load(Ordering::Acquire) {
549            return FlushFuture {
550                completion: Completion::ready(Err(FlushError::WriterStopped)),
551            };
552        }
553        let completion = Completion::pending();
554        let future = FlushFuture {
555            completion: completion.clone(),
556        };
557        let inner = self.inner.clone();
558        if thread::Builder::new()
559            .name("saddle-log-flush".to_owned())
560            .spawn(move || coordinate_flush(inner, completion.clone()))
561            .is_err()
562        {
563            future
564                .completion
565                .complete(Err(FlushError::CoordinatorUnavailable));
566        }
567        future
568    }
569
570    /// Stops admission, drains accepted records, flushes output and joins the
571    /// writer without blocking the async executor thread.
572    pub fn shutdown(&self) -> FlushFuture {
573        if self.inner.shutdown_started.swap(true, Ordering::AcqRel) {
574            return FlushFuture {
575                completion: Completion::ready(Err(FlushError::AlreadyShuttingDown)),
576            };
577        }
578        self.inner.accepting.store(false, Ordering::Release);
579        let completion = Completion::pending();
580        let future = FlushFuture {
581            completion: completion.clone(),
582        };
583        let inner = self.inner.clone();
584        if thread::Builder::new()
585            .name("saddle-log-shutdown".to_owned())
586            .spawn(move || coordinate_shutdown(inner, completion.clone()))
587            .is_err()
588        {
589            self.inner.accepting.store(true, Ordering::Release);
590            self.inner.shutdown_started.store(false, Ordering::Release);
591            future
592                .completion
593                .complete(Err(FlushError::CoordinatorUnavailable));
594        }
595        future
596    }
597
598    pub fn dropped_events(&self) -> u64 {
599        self.inner.dropped.load(Ordering::Relaxed)
600    }
601}
602
603impl ComponentLifecycle for Observer {
604    fn name(&self) -> &'static str {
605        "observability"
606    }
607
608    fn start(&self) -> LifecycleFuture<'_> {
609        Box::pin(async { Ok(()) })
610    }
611
612    fn shutdown(&self) -> LifecycleFuture<'_> {
613        let shutdown = Observer::shutdown(self);
614        Box::pin(async move {
615            shutdown.await.map_err(|_| {
616                SaddleError::new(
617                    ErrorKind::Infrastructure,
618                    "observability.shutdown_failed",
619                    "structured log shutdown failed",
620                )
621            })
622        })
623    }
624}
625
626fn coordinate_flush(inner: Arc<Inner>, completion: Arc<Completion>) {
627    let dropped = inner.dropped.swap(0, Ordering::AcqRel);
628    let (response, receiver) = mpsc::channel();
629    let result = if inner
630        .sender
631        .send(Command::Flush {
632            unreported_dropped: dropped,
633            response,
634        })
635        .is_err()
636    {
637        Err(FlushError::WriterStopped)
638    } else {
639        receiver
640            .recv()
641            .map_err(|_| FlushError::WriterStopped)
642            .and_then(WorkerStatus::into_result)
643    };
644    completion.complete(result);
645}
646
647fn coordinate_shutdown(inner: Arc<Inner>, completion: Arc<Completion>) {
648    while inner.emitting.load(Ordering::Acquire) != 0 {
649        thread::yield_now();
650    }
651    let dropped = inner.dropped.swap(0, Ordering::AcqRel);
652    let (response, receiver) = mpsc::channel();
653    let mut result = if inner
654        .sender
655        .send(Command::Shutdown {
656            unreported_dropped: dropped,
657            response,
658        })
659        .is_err()
660    {
661        Err(FlushError::WriterStopped)
662    } else {
663        receiver
664            .recv()
665            .map_err(|_| FlushError::WriterStopped)
666            .and_then(WorkerStatus::into_result)
667    };
668
669    if let Some(worker) = inner.worker.lock().unwrap().take() {
670        if worker.join().is_err() {
671            result = Err(FlushError::WorkerPanicked);
672        }
673    }
674    inner.process_log.lock().unwrap_or_else(|p| p.into_inner()).take();
675    completion.complete(result);
676}
677
678fn write_records(receiver: Receiver<Command>, mut writer: impl Write) {
679    let mut status = WorkerStatus::default();
680    while let Ok(command) = receiver.recv() {
681        match command {
682            Command::Record(record) => write_record(&mut writer, record, &mut status),
683            Command::ProcessCall { bytes, dropped } => {
684                status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
685                if status.failure.is_none() {
686                    if writer.write_all(bytes.get()).is_err() {
687                        status.failure = Some(OutputStage::Record);
688                    } else if writer.write_all(b"\n").is_err() {
689                        status.failure = Some(OutputStage::Newline);
690                    }
691                }
692                drop(bytes);
693            }
694            Command::RootRecord { bytes, dropped } => {
695                status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
696                if status.failure.is_none() {
697                    if writer.write_all(&bytes).is_err() {
698                        status.failure = Some(OutputStage::Record);
699                    } else if writer.write_all(b"\n").is_err() {
700                        status.failure = Some(OutputStage::Newline);
701                    }
702                }
703            }
704            Command::Flush {
705                unreported_dropped,
706                response,
707            } => {
708                status.unreported_dropped =
709                    status.unreported_dropped.saturating_add(unreported_dropped);
710                flush_writer(&mut writer, &mut status);
711                let _ = response.send(status);
712            }
713            Command::Shutdown {
714                unreported_dropped,
715                response,
716            } => {
717                status.unreported_dropped =
718                    status.unreported_dropped.saturating_add(unreported_dropped);
719                flush_writer(&mut writer, &mut status);
720                let _ = response.send(status);
721                break;
722            }
723        }
724    }
725}
726
727fn write_record(writer: &mut impl Write, record: LogRecord, status: &mut WorkerStatus) {
728    if status.failure.is_some() {
729        status.unreported_dropped = status
730            .unreported_dropped
731            .saturating_add(record.dropped_events);
732        return;
733    }
734
735    let dropped = record.dropped_events;
736    let bytes = match serde_json::to_vec(&record) {
737        Ok(bytes) => bytes,
738        Err(_) => {
739            status.failure = Some(OutputStage::Serialize);
740            status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
741            return;
742        }
743    };
744    if writer.write_all(&bytes).is_err() {
745        status.failure = Some(OutputStage::Record);
746        status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
747        return;
748    }
749    if writer.write_all(b"\n").is_err() {
750        status.failure = Some(OutputStage::Newline);
751        status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
752    }
753}
754
755fn flush_writer(writer: &mut impl Write, status: &mut WorkerStatus) {
756    if writer.flush().is_err() && status.failure.is_none() {
757        status.failure = Some(OutputStage::Flush);
758    }
759}
760
761#[cfg(test)]
762mod tests {
763    use std::{
764        sync::{Condvar, mpsc},
765        task::{Wake, Waker},
766    };
767
768    use super::*;
769
770    struct ThreadWaker(thread::Thread);
771
772    impl Wake for ThreadWaker {
773        fn wake(self: Arc<Self>) {
774            self.0.unpark();
775        }
776    }
777
778    fn block_on<T>(future: impl Future<Output = T>) -> T {
779        let waker = Waker::from(Arc::new(ThreadWaker(thread::current())));
780        let mut context = Context::from_waker(&waker);
781        let mut future = std::pin::pin!(future);
782        loop {
783            match future.as_mut().poll(&mut context) {
784                Poll::Ready(output) => return output,
785                Poll::Pending => thread::park(),
786            }
787        }
788    }
789
790    struct BlockingWriter {
791        entered: Option<mpsc::Sender<()>>,
792        release: Arc<(Mutex<bool>, Condvar)>,
793        fail_after_release: bool,
794    }
795
796    impl Write for BlockingWriter {
797        fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
798            if let Some(entered) = self.entered.take() {
799                let _ = entered.send(());
800                let (lock, condition) = &*self.release;
801                let mut released = lock.lock().unwrap();
802                while !*released {
803                    released = condition.wait(released).unwrap();
804                }
805                if self.fail_after_release {
806                    return Err(io::Error::other("injected writer failure"));
807                }
808            }
809            Ok(bytes.len())
810        }
811
812        fn flush(&mut self) -> io::Result<()> {
813            Ok(())
814        }
815    }
816
817    struct FailingWriter {
818        fail_write: Option<usize>,
819        writes: usize,
820        fail_flush: bool,
821    }
822
823    impl Write for FailingWriter {
824        fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
825            self.writes += 1;
826            if self.fail_write == Some(self.writes) {
827                Err(io::Error::other("injected writer failure"))
828            } else {
829                Ok(bytes.len())
830            }
831        }
832
833        fn flush(&mut self) -> io::Result<()> {
834            if self.fail_flush {
835                Err(io::Error::other("injected flush failure"))
836            } else {
837                Ok(())
838            }
839        }
840    }
841
842    fn observer(writer: impl Write + Send + 'static, capacity: usize) -> Observer {
843        Observer::with_writer_and_seed(
844            ObserverConfig {
845                queue_capacity: capacity,
846            },
847            writer,
848            Ok([7; 24]),
849        )
850        .unwrap()
851    }
852
853    #[test]
854    fn managed_call_log_outlives_request_account_while_writer_is_blocked() {
855        // The new process-owned variant adds no command-ring slot bytes:
856        // the unchanged ordinary LogRecord still determines Command's size.
857        assert_eq!(std::mem::size_of::<Command>(), std::mem::size_of::<LogRecord>());
858        assert_eq!(std::mem::size_of::<Command>(), 192);
859        assert_eq!(ObserverConfig::default().queue_capacity * (8 + std::mem::size_of::<Command>()),
860            1_638_400);
861        use saddle_admission::{
862            ProfuseGwLightweightObservedAdmissionOutcome, StorageDemand,
863            bind_deployment_resource_budget_bootstrap, freeze_deployment_resource_budget,
864            prepare_profusegw_lightweight_profile,
865        };
866        use saddle_core::{
867            BootstrapRendezvousIssuer, GeneratedApplicationFreezeSource,
868            ListenerStartupFreezeSource, RpcCorrelationId, pair_bootstrap_rendezvous,
869        };
870        let pending = freeze_deployment_resource_budget(
871            1, 32, 5_000, 1, 1, 1_000_000, 1_000_000, 1_000_000, 1_000_000,
872        ).unwrap().reserve_known_database_memory(0, 0, 0).ok().unwrap();
873        let (application, listener) = BootstrapRendezvousIssuer::issue()
874            .freeze_application(GeneratedApplicationFreezeSource::new(
875                "app", b"descriptor", &["route"],
876            )).unwrap();
877        let listener = listener.freeze_listener(ListenerStartupFreezeSource::new(
878            "app", "127.0.0.1:8000".parse().unwrap(),
879            "127.0.0.1:9000".parse().unwrap(),
880            std::time::Duration::from_millis(5_000),
881        )).ok().unwrap();
882        let (whole, receipt) = pair_bootstrap_rendezvous(application, listener).ok().unwrap();
883        let budget = bind_deployment_resource_budget_bootstrap(pending, whole, receipt).ok().unwrap();
884        let process = prepare_profusegw_lightweight_profile(budget).ok().unwrap();
885        let root = process.try_process_storage(StorageDemand::separate(&[(
886            std::alloc::Layout::array::<u8>(16).unwrap(), 1,
887        )]).unwrap()).unwrap();
888        let baseline = process.resource_snapshot().framework_charged;
889        let (entered, entered_rx) = mpsc::channel();
890        let release = Arc::new((Mutex::new(false), Condvar::new()));
891        let logger = observer(BlockingWriter {
892            entered: Some(entered), release: Arc::clone(&release), fail_after_release: false,
893        }, 2);
894        logger.install_process_log_storage(process.process_log_storage());
895        let (call, _) = logger.start_managed_external_call_with_rpc(
896            "app", "zone", "interface", "operation", None,
897            RpcCorrelationId::new("0").unwrap(),
898        ).unwrap();
899        entered_rx.recv_timeout(std::time::Duration::from_secs(2)).unwrap();
900        let admitted = match process.verified_profile().try_admit_observed() {
901            ProfuseGwLightweightObservedAdmissionOutcome::Ready(permit, _) => permit,
902            _ => panic!("fixture admission failed"),
903        };
904        let original = saddle_core::SaddleError::new(
905            saddle_core::ErrorKind::Internal, "original.failure", "original",
906        );
907        call.fail(&original);
908        assert_eq!(original.code(), "original.failure");
909        admitted.into_execution().cancel_observed();
910        assert_eq!(process.resource_snapshot().active_accounts, 0);
911        assert!(process.resource_snapshot().framework_charged > baseline,
912            "queued process log must remain charged after request cleanup");
913        // The worker holds the first record while the second is queued.
914        // Additional managed call records must be dropped synchronously.
915        let (full, _) = logger.start_managed_external_call_with_rpc(
916            "app", "zone", "interface", "full", None,
917            RpcCorrelationId::new("1").unwrap(),
918        ).unwrap();
919        full.fail(&original);
920        assert!(logger.dropped_events() > 0);
921        *release.0.lock().unwrap() = true;
922        release.1.notify_all();
923        assert!(matches!(block_on(logger.flush()), Err(FlushError::DroppedEvents(_))));
924        assert_eq!(process.resource_snapshot().framework_charged, baseline);
925
926        // Exhaust the *same* framework domain with a physical child owner;
927        // encoding must fail to reserve before allocating and remain lossy.
928        let available = process.resource_snapshot().framework_capacity
929            - process.resource_snapshot().framework_charged;
930        let filler = process.try_test_process_storage_child(StorageDemand::separate(&[(
931            std::alloc::Layout::array::<u8>(
932                available - 64 - std::mem::size_of::<saddle_admission::StoragePermit>(),
933            ).unwrap(), 1,
934        )]).unwrap()).unwrap();
935        let charged = process.resource_snapshot().framework_charged;
936        let before_drop = logger.dropped_events();
937        let (short, _) = logger.start_managed_external_call_with_rpc(
938            "app", "zone", "interface", "short", None,
939            RpcCorrelationId::new("2").unwrap(),
940        ).unwrap();
941        short.fail(&original);
942        assert!(logger.dropped_events() >= before_drop + 2);
943        assert_eq!(process.resource_snapshot().framework_charged, charged);
944        drop(filler);
945        assert!(matches!(block_on(logger.shutdown()), Err(FlushError::DroppedEvents(_))));
946        let after_shutdown = process.resource_snapshot().framework_charged;
947        let (closed, _) = logger.start_managed_external_call_with_rpc(
948            "app", "zone", "interface", "closed", None,
949            RpcCorrelationId::new("3").unwrap(),
950        ).unwrap();
951        closed.fail(&original);
952        assert_eq!(process.resource_snapshot().framework_charged, after_shutdown);
953        assert_eq!(process.resource_snapshot().framework_charged, baseline);
954        drop(logger);
955        drop(root);
956        process.finish().unwrap();
957    }
958
959    #[test]
960    fn root_contract_ordinary_uses_same_view_and_bounded_queue() {
961        use crate::{
962            DiagnosticSubmission, RootDiagnosticScope, RootOutcomeFacts, RootRequestEvent,
963        };
964        use saddle_core::{
965            ContextFact, ContextLabel, RequestLocalFacts, RequestRootPublisher, RequestViewPhase,
966        };
967        let root = RequestRootPublisher::create(
968            ContextLabel::checked("app").unwrap(),
969            ContextFact::NotEstablished,
970        )
971        .unwrap();
972        let view = root
973            .reference()
974            .view(RequestLocalFacts::new(RequestViewPhase::Handler));
975        let expected = serde_json::to_value(&view).unwrap();
976        let mut frame = crate::diagnostic::FixedDiagnosticBytes {
977            bytes: [0; 8192],
978            len: 0,
979        };
980        serde_json::to_writer(&mut frame, &view).unwrap();
981        let ordinary_context: serde_json::Value =
982            serde_json::from_slice(&frame.bytes[..frame.len]).unwrap();
983        assert_eq!(ordinary_context, expected);
984        let (entered, wait) = mpsc::channel();
985        let release = Arc::new((Mutex::new(false), Condvar::new()));
986        struct Captured {
987            writer: BlockingWriter,
988            bytes: Arc<Mutex<Vec<u8>>>,
989        }
990        impl Write for Captured {
991            fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
992                let n = self.writer.write(bytes)?;
993                self.bytes.lock().unwrap().extend_from_slice(&bytes[..n]);
994                Ok(n)
995            }
996            fn flush(&mut self) -> io::Result<()> {
997                self.writer.flush()
998            }
999        }
1000        struct Unblock(Arc<(Mutex<bool>, Condvar)>);
1001        impl Drop for Unblock {
1002            fn drop(&mut self) {
1003                *self.0.0.lock().unwrap() = true;
1004                self.0.1.notify_all();
1005            }
1006        }
1007        let unblock = Unblock(Arc::clone(&release));
1008        let bytes = Arc::new(Mutex::new(Vec::new()));
1009        let logger = observer(
1010            Captured {
1011                writer: BlockingWriter {
1012                    entered: Some(entered),
1013                    release: Arc::clone(&release),
1014                    fail_after_release: false,
1015                },
1016                bytes: Arc::clone(&bytes),
1017            },
1018            1,
1019        );
1020        let scope = RootDiagnosticScope::new(&view, None);
1021        assert_eq!(
1022            scope.ordinary(
1023                &logger,
1024                RootRequestEvent::Handler,
1025                RootOutcomeFacts::default()
1026            ),
1027            DiagnosticSubmission::Enqueued
1028        );
1029        wait.recv_timeout(std::time::Duration::from_secs(5))
1030            .unwrap();
1031        assert_eq!(
1032            scope.ordinary(
1033                &logger,
1034                RootRequestEvent::Response,
1035                RootOutcomeFacts::default()
1036            ),
1037            DiagnosticSubmission::Enqueued
1038        );
1039        assert_eq!(
1040            scope.ordinary(
1041                &logger,
1042                RootRequestEvent::Response,
1043                RootOutcomeFacts::default()
1044            ),
1045            DiagnosticSubmission::Full
1046        );
1047        drop((view, root));
1048        drop(unblock);
1049        assert!(matches!(
1050            block_on(logger.shutdown()),
1051            Err(FlushError::DroppedEvents(1))
1052        ));
1053        let bytes = bytes.lock().unwrap();
1054        let rows: Vec<serde_json::Value> = serde_json::Deserializer::from_slice(&bytes)
1055            .into_iter()
1056            .map(Result::unwrap)
1057            .collect();
1058        assert_eq!(rows.len(), 2);
1059        assert!(rows.iter().all(|row| row["context"] == expected));
1060    }
1061
1062    #[test]
1063    fn root_contract_encoded_ordinary_readback_and_failure_use_original_writer() {
1064        let record = serde_json::json!({"event":"request_stage", "context":{"local_request":5}});
1065        let (tx, rx) = mpsc::sync_channel(1);
1066        tx.try_send(Command::RootRecord {
1067            bytes: serde_json::to_vec(&record).unwrap(),
1068            dropped: 0,
1069        })
1070        .unwrap_or_else(|_| panic!("empty queue"));
1071        drop(tx);
1072        let mut bytes = Vec::new();
1073        write_records(rx, &mut bytes);
1074        assert_eq!(
1075            serde_json::from_slice::<serde_json::Value>(&bytes).unwrap(),
1076            record
1077        );
1078        let logger = observer(
1079            FailingWriter {
1080                fail_write: Some(1),
1081                writes: 0,
1082                fail_flush: false,
1083            },
1084            1,
1085        );
1086        assert_eq!(
1087            logger.emit_root_record(&record),
1088            crate::DiagnosticSubmission::Enqueued
1089        );
1090        assert!(matches!(
1091            block_on(logger.shutdown()),
1092            Err(FlushError::OutputFailed(OutputStage::Record))
1093        ));
1094    }
1095
1096    #[test]
1097    fn random_source_failure_is_reported_during_initialization() {
1098        let result = Observer::with_writer_and_seed(
1099            ObserverConfig::default(),
1100            io::sink(),
1101            Err(InitError::RandomSource),
1102        );
1103        assert!(matches!(result, Err(InitError::RandomSource)));
1104    }
1105
1106    #[test]
1107    fn initialized_generator_produces_non_zero_unique_ids_without_more_io() {
1108        let observer = observer(io::sink(), 8);
1109        let first_trace = observer.new_trace_id();
1110        let second_trace = observer.new_trace_id();
1111        let first_span = observer.new_span_id();
1112        let second_span = observer.new_span_id();
1113        assert_ne!(first_trace.as_u128(), 0);
1114        assert_ne!(first_trace, second_trace);
1115        assert_ne!(first_span.as_u64(), 0);
1116        assert_ne!(first_span, second_span);
1117        block_on(observer.shutdown()).unwrap();
1118    }
1119
1120    #[test]
1121    fn observer_is_a_managed_component() {
1122        let observer = observer(io::sink(), 8);
1123        assert_eq!(ComponentLifecycle::name(&observer), "observability");
1124        block_on(ComponentLifecycle::start(&observer)).unwrap();
1125        block_on(ComponentLifecycle::shutdown(&observer)).unwrap();
1126    }
1127
1128    #[test]
1129    fn full_queue_is_reported_by_shutdown_without_waiting_in_emit() {
1130        let (entered_sender, entered_receiver) = mpsc::channel();
1131        let release = Arc::new((Mutex::new(false), Condvar::new()));
1132        let writer = BlockingWriter {
1133            entered: Some(entered_sender),
1134            release: release.clone(),
1135            fail_after_release: false,
1136        };
1137        let observer = observer(writer, 1);
1138
1139        observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
1140        entered_receiver.recv().unwrap();
1141        observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
1142        observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
1143        assert_eq!(observer.dropped_events(), 1);
1144
1145        let shutdown = observer.shutdown();
1146        let (lock, condition) = &*release;
1147        *lock.lock().unwrap() = true;
1148        condition.notify_one();
1149        assert_eq!(block_on(shutdown), Err(FlushError::DroppedEvents(1)));
1150    }
1151
1152    #[test]
1153    fn shutdown_reports_writer_failure_and_unreported_drops_together() {
1154        let (entered_sender, entered_receiver) = mpsc::channel();
1155        let release = Arc::new((Mutex::new(false), Condvar::new()));
1156        let writer = BlockingWriter {
1157            entered: Some(entered_sender),
1158            release: release.clone(),
1159            fail_after_release: true,
1160        };
1161        let observer = observer(writer, 1);
1162
1163        observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
1164        entered_receiver.recv().unwrap();
1165        observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
1166        observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
1167        let shutdown = observer.shutdown();
1168
1169        let (lock, condition) = &*release;
1170        *lock.lock().unwrap() = true;
1171        condition.notify_one();
1172        assert_eq!(
1173            block_on(shutdown),
1174            Err(FlushError::OutputFailedAndDropped {
1175                stage: OutputStage::Record,
1176                dropped_events: 1,
1177            })
1178        );
1179    }
1180
1181    #[test]
1182    fn record_write_failure_persists_until_shutdown() {
1183        let observer = observer(
1184            FailingWriter {
1185                fail_write: Some(1),
1186                writes: 0,
1187                fail_flush: false,
1188            },
1189            8,
1190        );
1191        observer.emit(LogRecord::new(EventLevel::Info, "test.record"));
1192        assert_eq!(
1193            block_on(observer.shutdown()),
1194            Err(FlushError::OutputFailed(OutputStage::Record))
1195        );
1196    }
1197
1198    #[test]
1199    fn newline_failure_persists_until_flush() {
1200        let observer = observer(
1201            FailingWriter {
1202                fail_write: Some(2),
1203                writes: 0,
1204                fail_flush: false,
1205            },
1206            8,
1207        );
1208        observer.emit(LogRecord::new(EventLevel::Info, "test.newline"));
1209        assert_eq!(
1210            block_on(observer.flush()),
1211            Err(FlushError::OutputFailed(OutputStage::Newline))
1212        );
1213        let _ = block_on(observer.shutdown());
1214    }
1215
1216    #[test]
1217    fn flush_failure_is_reported() {
1218        let observer = observer(
1219            FailingWriter {
1220                fail_write: None,
1221                writes: 0,
1222                fail_flush: true,
1223            },
1224            8,
1225        );
1226        assert_eq!(
1227            block_on(observer.flush()),
1228            Err(FlushError::OutputFailed(OutputStage::Flush))
1229        );
1230        assert_eq!(
1231            block_on(observer.shutdown()),
1232            Err(FlushError::OutputFailed(OutputStage::Flush))
1233        );
1234    }
1235}