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