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