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