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 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 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
322pub 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
333pub 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
347pub fn global() -> Option<&'static Observer> {
349 GLOBAL.get()
350}
351
352impl Observer {
353 #[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 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 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 pub fn metrics_snapshot(&self) -> crate::MetricsSnapshot {
542 self.inner.metrics.snapshot()
543 }
544
545 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 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 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 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 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}