use std::{
error::Error,
fmt,
future::Future,
io::{self, Write},
pin::Pin,
sync::{
Arc, Mutex, OnceLock,
atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering},
mpsc::{self, Receiver, SyncSender, TrySendError},
},
task::{Context, Poll, Waker},
thread::{self, JoinHandle},
time::{SystemTime, UNIX_EPOCH},
};
use saddle_core::{ComponentLifecycle, ErrorKind, LifecycleFuture, SaddleError, SpanId, TraceId};
use serde::Serialize;
use crate::{
calendar_writer::{CalendarFileWriter, FileLoggingConfig},
event::EventLevel,
};
static GLOBAL: OnceLock<Observer> = OnceLock::new();
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct ObserverConfig {
pub queue_capacity: usize,
}
impl Default for ObserverConfig {
fn default() -> Self {
Self {
queue_capacity: 8_192,
}
}
}
#[derive(Debug)]
pub enum InitError {
EmptyQueue,
RandomSource,
Output(io::Error),
Spawn(io::Error),
}
impl fmt::Display for InitError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::EmptyQueue => formatter.write_str("log queue capacity must be greater than zero"),
Self::RandomSource => {
formatter.write_str("operating system random source is unavailable")
}
Self::Output(error) => write!(formatter, "failed to initialize log output: {error}"),
Self::Spawn(error) => write!(formatter, "failed to start log writer: {error}"),
}
}
}
impl Error for InitError {
fn source(&self) -> Option<&(dyn Error + 'static)> {
match self {
Self::Spawn(error) | Self::Output(error) => Some(error),
Self::EmptyQueue | Self::RandomSource => None,
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum OutputStage {
Serialize,
Record,
Newline,
Flush,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum FlushError {
OutputFailed(OutputStage),
DroppedEvents(u64),
OutputFailedAndDropped {
stage: OutputStage,
dropped_events: u64,
},
WriterStopped,
AlreadyShuttingDown,
CoordinatorUnavailable,
WorkerPanicked,
}
impl fmt::Display for FlushError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::OutputFailed(stage) => write!(formatter, "log output failed during {stage:?}"),
Self::DroppedEvents(count) => write!(formatter, "{count} log events were dropped"),
Self::OutputFailedAndDropped {
stage,
dropped_events,
} => write!(
formatter,
"log output failed during {stage:?} and {dropped_events} events were dropped"
),
Self::WriterStopped => formatter.write_str("log writer is not running"),
Self::AlreadyShuttingDown => formatter.write_str("log writer shutdown already started"),
Self::CoordinatorUnavailable => {
formatter.write_str("failed to start log lifecycle coordinator")
}
Self::WorkerPanicked => formatter.write_str("log writer thread panicked"),
}
}
}
impl Error for FlushError {}
pub struct FlushFuture {
completion: Arc<Completion>,
}
impl Future for FlushFuture {
type Output = Result<(), FlushError>;
fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
if let Some(result) = self.completion.result.lock().unwrap().take() {
return Poll::Ready(result);
}
*self.completion.waker.lock().unwrap() = Some(context.waker().clone());
if let Some(result) = self.completion.result.lock().unwrap().take() {
Poll::Ready(result)
} else {
Poll::Pending
}
}
}
struct Completion {
result: Mutex<Option<Result<(), FlushError>>>,
waker: Mutex<Option<Waker>>,
}
impl Completion {
fn pending() -> Arc<Self> {
Arc::new(Self {
result: Mutex::new(None),
waker: Mutex::new(None),
})
}
fn ready(result: Result<(), FlushError>) -> Arc<Self> {
Arc::new(Self {
result: Mutex::new(Some(result)),
waker: Mutex::new(None),
})
}
fn complete(&self, result: Result<(), FlushError>) {
*self.result.lock().unwrap() = Some(result);
if let Some(waker) = self.waker.lock().unwrap().take() {
waker.wake();
}
}
}
#[derive(Clone)]
pub struct Observer {
pub(crate) inner: Arc<Inner>,
}
struct ContextByteCount(usize);
impl Write for ContextByteCount {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> { self.0 = self.0.saturating_add(bytes.len()); Ok(bytes.len()) }
fn flush(&mut self) -> io::Result<()> { Ok(()) }
}
fn serialize_record_id<S: serde::Serializer>(id: &SpanId, serializer: S) -> Result<S::Ok, S::Error> { serializer.collect_str(id) }
struct OrdinarySegments<'a> {
observer: &'a Observer, storage: &'a saddle_admission::ProcessLogStorage,
profile: Option<usize>, record_id: SpanId,
}
impl crate::root_diagnostic::SegmentOutput for OrdinarySegments<'_> {
fn submit(&self, segment: &crate::root_diagnostic::EncodedSegment<'_>) -> crate::DiagnosticSubmission {
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct Frame<'a> { record_schema: &'static str, event: &'static str, #[serde(serialize_with = "serialize_record_id")] record_id: SpanId,
format: &'static str, #[serde(skip_serializing_if = "Option::is_none")] profile: Option<&'a str>,
sequence: u64, payload: &'a str, state: &'static str }
let profile = self.profile.map(|index| &self.observer.inner.context.summaries[index]);
let format = if profile.is_some_and(|profile| profile.format == crate::SummaryFormat::Text) { "text" } else { "json" };
let frame = Frame { record_schema: "1.0", event: "context_record_segment", record_id: self.record_id,
format, profile: profile.map(|profile| profile.name.as_str()), sequence: segment.sequence, payload: segment.payload, state: segment.state };
self.observer.enqueue_context_frame(self.profile, self.storage, &|writer| serde_json::to_writer(writer, &frame))
}
}
trait LogSink: Write {
fn write_profile(&mut self, profile: usize, bytes: &[u8]) -> io::Result<()>;
}
struct BasicSink<W>(W);
impl<W: Write> Write for BasicSink<W> {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> { self.0.write(bytes) }
fn flush(&mut self) -> io::Result<()> { self.0.flush() }
}
impl<W: Write> LogSink for BasicSink<W> {
fn write_profile(&mut self, _: usize, bytes: &[u8]) -> io::Result<()> { self.0.write_all(bytes)?; self.0.write_all(b"\n") }
}
struct CalendarOutputs { main: CalendarFileWriter, summaries: Vec<Option<CalendarFileWriter>> }
impl Write for CalendarOutputs {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> { self.main.write(bytes) }
fn flush(&mut self) -> io::Result<()> {
self.main.flush()?;
for output in self.summaries.iter_mut().flatten() { output.flush()?; }
Ok(())
}
}
impl LogSink for CalendarOutputs {
fn write_profile(&mut self, profile: usize, bytes: &[u8]) -> io::Result<()> {
let output = self.summaries.get_mut(profile).and_then(Option::as_mut).ok_or(io::ErrorKind::InvalidInput)?;
output.write_all(bytes)?; output.write_all(b"\n")
}
}
pub(crate) struct Inner {
sender: SyncSender<Command>,
dropped: AtomicU64,
accepting: AtomicBool,
emitting: AtomicUsize,
shutdown_started: AtomicBool,
worker: Mutex<Option<JoinHandle<()>>>,
process_log: Mutex<Option<saddle_admission::ProcessLogStorage>>,
context: crate::ContextLoggingConfig,
writer_source: Arc<Mutex<WriterSourceState>>,
ids: IdGenerator,
pub(crate) metrics: crate::metrics::Metrics,
}
#[derive(Clone)]
struct WriterSourceContext {
application: saddle_core::ContextLabel,
output: crate::SourceOutput,
primary: Option<saddle_core::DiagnosticOccurrence>,
}
#[derive(Default)]
struct WriterSourceState {
context: Option<WriterSourceContext>,
failure: Option<WriterSourceFailure>,
}
struct WriterSourceFailure {
original: WriterStageError,
diagnostic: saddle_core::Diagnostic,
written: Option<crate::root_diagnostic::WrittenComponentFailure>,
}
#[derive(Debug)]
struct WriterStageError {
stage: OutputStage,
original: io::Error,
}
impl fmt::Display for WriterStageError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(formatter, "log writer failed during {:?}: {}", self.stage, self.original)
}
}
impl Error for WriterStageError {
fn source(&self) -> Option<&(dyn Error + 'static)> { Some(&self.original) }
}
fn record_writer_failure(source: &Arc<Mutex<WriterSourceState>>,
stage: OutputStage, original: io::Error) {
let context = source.lock().unwrap_or_else(|e| e.into_inner()).context.clone();
let original = WriterStageError { stage, original };
let mut diagnostic = saddle_core::Diagnostic::capture(
saddle_core::DiagnosticCategory::UnexpectedError,
saddle_core::CaptureSite::FirstObserved,
saddle_core::DiagnosticCause::new(
saddle_core::DiagnosticStage::ShutdownLogger,
saddle_core::DiagnosticCode::new("observability.writer_failed")
.expect("static writer failure code")),
);
if let Some(primary) = context.as_ref().and_then(|context| context.primary) {
diagnostic = diagnostic.during_cleanup_of_occurrence(&primary);
}
let written = context.as_ref().and_then(|context|
crate::root_diagnostic::process_component_cleanup_recorded(
Some(context.output.handle()),
saddle_core::ContextFact::Present(context.application.clone()),
&diagnostic, &original).ok());
let mut state = source.lock().unwrap_or_else(|e| e.into_inner());
if state.failure.is_none() {
state.failure = Some(WriterSourceFailure { original, diagnostic, written });
}
}
struct IdGenerator {
trace_high: u64,
trace_seed: u64,
trace_counter: AtomicU64,
span_seed: u64,
span_counter: AtomicU64,
}
impl IdGenerator {
fn from_seed(seed: [u8; 24]) -> Self {
let mut high = u64::from_be_bytes(seed[0..8].try_into().unwrap());
if high == 0 {
high = 1;
}
Self {
trace_high: high,
trace_seed: u64::from_be_bytes(seed[8..16].try_into().unwrap()),
trace_counter: AtomicU64::new(0),
span_seed: u64::from_be_bytes(seed[16..24].try_into().unwrap()),
span_counter: AtomicU64::new(0),
}
}
fn trace_id(&self) -> TraceId {
let low = self
.trace_seed
.wrapping_add(self.trace_counter.fetch_add(1, Ordering::Relaxed));
TraceId::from_u128((u128::from(self.trace_high) << 64) | u128::from(low))
}
fn span_id(&self) -> SpanId {
loop {
let value = self
.span_seed
.wrapping_add(self.span_counter.fetch_add(1, Ordering::Relaxed));
if value != 0 {
return SpanId::from_u64(value);
}
}
}
}
enum Command {
Record(LogRecord),
ProcessCall {
bytes: saddle_admission::ExactStored<Vec<u8>>,
dropped: u64,
},
ContextProfile { bytes: saddle_admission::ExactStored<Vec<u8>>, profile: Option<usize>, dropped: u64 },
RootRecord {
bytes: Vec<u8>,
dropped: u64,
},
Flush {
unreported_dropped: u64,
response: mpsc::Sender<WorkerStatus>,
},
Shutdown {
unreported_dropped: u64,
response: mpsc::Sender<WorkerStatus>,
},
}
pub(crate) fn root_queue_layout() -> std::alloc::Layout {
std::alloc::Layout::new::<Command>()
}
#[derive(Clone, Copy, Debug, Default)]
struct WorkerStatus {
failure: Option<OutputStage>,
unreported_dropped: u64,
}
impl WorkerStatus {
fn into_result(self) -> Result<(), FlushError> {
match (self.failure, self.unreported_dropped) {
(None, 0) => Ok(()),
(Some(stage), 0) => Err(FlushError::OutputFailed(stage)),
(None, dropped_events) => Err(FlushError::DroppedEvents(dropped_events)),
(Some(stage), dropped_events) => Err(FlushError::OutputFailedAndDropped {
stage,
dropped_events,
}),
}
}
}
#[derive(Serialize)]
pub(crate) struct LogRecord {
pub timestamp_unix_ms: u128,
pub level: EventLevel,
pub event: &'static str,
#[serde(skip_serializing_if = "Option::is_none")]
pub trace_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub span: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub span_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub parent: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub parent_span_id: Option<String>,
#[serde(flatten)]
pub data: serde_json::Map<String, serde_json::Value>,
#[serde(skip_serializing_if = "is_zero")]
pub dropped_events: u64,
}
const fn is_zero(value: &u64) -> bool {
*value == 0
}
impl LogRecord {
pub(crate) fn new(level: EventLevel, event: &'static str) -> Self {
Self {
timestamp_unix_ms: SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis(),
level,
event,
trace_id: None,
span: None,
span_id: None,
parent: None,
parent_span_id: None,
data: serde_json::Map::new(),
dropped_events: 0,
}
}
}
pub fn init(config: ObserverConfig) -> Result<&'static Observer, InitError> {
if let Some(observer) = GLOBAL.get() {
return Ok(observer);
}
let observer = Observer::with_writer(config, io::stdout())?;
let _ = GLOBAL.set(observer);
Ok(GLOBAL.get().expect("global observer was initialized"))
}
pub fn init_file(
config: ObserverConfig,
file: FileLoggingConfig,
) -> Result<&'static Observer, InitError> {
if let Some(observer) = GLOBAL.get() {
return Ok(observer);
}
let context = file.context().clone();
let summaries = context.summaries.iter().map(|profile| {
if profile.enabled { CalendarFileWriter::open(file.for_summary(&profile.name)).map(Some) }
else { Ok(None) }
}).collect::<io::Result<Vec<_>>>().map_err(InitError::Output)?;
let main = CalendarFileWriter::open(file).map_err(InitError::Output)?;
let mut seed = [0; 24];
getrandom::fill(&mut seed).map_err(|_| InitError::RandomSource)?;
let observer = Observer::with_sink_and_seed(config, CalendarOutputs { main, summaries }, Ok(seed), context)?;
let _ = GLOBAL.set(observer);
Ok(GLOBAL.get().expect("global observer was initialized"))
}
pub fn global() -> Option<&'static Observer> {
GLOBAL.get()
}
impl Observer {
#[doc(hidden)]
pub fn install_process_log_storage(&self, storage: saddle_admission::ProcessLogStorage) {
let mut slot = self.inner.process_log.lock().unwrap_or_else(|p| p.into_inner());
*slot = Some(storage);
}
#[doc(hidden)]
pub fn capture_terminal_context(&self, context: &saddle_core::request_context::UnifiedContext<'_>)
-> Result<crate::OwnedUnifiedContext, crate::DiagnosticSubmission> {
use crate::DiagnosticSubmission as S;
let storage = self.inner.process_log.lock().unwrap_or_else(|p| p.into_inner())
.clone().ok_or(S::OutputUnavailable)?;
let mut count = ContextByteCount(0);
serde_json::to_writer(&mut count, context).map_err(|_| S::EncodingFailed)?;
let layout = std::alloc::Layout::array::<u8>(count.0).map_err(|_| S::EncodingFailed)?;
let permit = storage.try_reserve(layout).map_err(|_| S::Full)?;
let mut bytes = permit.allocate_exact(layout, || Vec::with_capacity(count.0));
struct Bounded<'a> { bytes: &'a mut Vec<u8>, exceeded: bool }
impl Write for Bounded<'_> {
fn write(&mut self, chunk: &[u8]) -> io::Result<usize> {
if !self.exceeded && chunk.len() <= self.bytes.capacity() - self.bytes.len() {
self.bytes.extend_from_slice(chunk);
} else { self.exceeded = true; }
Ok(chunk.len())
}
fn flush(&mut self) -> io::Result<()> { Ok(()) }
}
let mut writer = Bounded { bytes: bytes.get_mut(), exceeded: false };
serde_json::to_writer(&mut writer, context).map_err(|_| S::EncodingFailed)?;
if writer.exceeded || writer.bytes.len() != count.0 { return Err(S::EncodingFailed); }
Ok(crate::OwnedUnifiedContext { bytes })
}
pub(crate) fn emit_process_call(&self, record: &mut crate::call::ProcessCallRecord<'_>) {
struct Counter(usize);
impl Write for Counter {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
self.0 = self.0.checked_add(bytes.len()).ok_or(io::ErrorKind::OutOfMemory)?;
Ok(bytes.len())
}
fn flush(&mut self) -> io::Result<()> { Ok(()) }
}
self.inner.emitting.fetch_add(1, Ordering::AcqRel);
let submitted = (|| {
if !self.inner.accepting.load(Ordering::Acquire) { return false; }
let storage = self.inner.process_log.lock().unwrap_or_else(|p| p.into_inner()).clone();
let Some(storage) = storage else { return false; };
let dropped = self.inner.dropped.swap(0, Ordering::Relaxed);
record.dropped_events = dropped;
let mut counter = Counter(0);
if serde_json::to_writer(&mut counter, record).is_err() {
self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
return false;
}
let Ok(layout) = std::alloc::Layout::array::<u8>(counter.0) else {
self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
return false;
};
let Ok(permit) = storage.try_reserve(layout) else {
self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
return false;
};
let mut bytes = permit.allocate_exact(layout, || Vec::with_capacity(counter.0));
if serde_json::to_writer(bytes.get_mut(), record).is_err() || bytes.get().len() != counter.0 {
drop(bytes);
self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
return false;
}
match self.inner.sender.try_send(Command::ProcessCall { bytes, dropped }) {
Ok(()) => true,
Err(TrySendError::Full(command)) => {
drop(command);
self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
false
}
Err(TrySendError::Disconnected(command)) => {
drop(command);
self.inner.metrics.logger_output(true);
self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
false
}
}
})();
if !submitted {
self.inner.metrics.logger_dropped(1);
self.inner.dropped.fetch_add(1, Ordering::Relaxed);
}
self.inner.emitting.fetch_sub(1, Ordering::Release);
}
pub fn with_writer(
config: ObserverConfig,
writer: impl Write + Send + 'static,
) -> Result<Self, InitError> {
let mut seed = [0_u8; 24];
getrandom::fill(&mut seed).map_err(|_| InitError::RandomSource)?;
Self::with_writer_and_seed(config, writer, Ok(seed))
}
fn with_writer_and_seed(
config: ObserverConfig,
writer: impl Write + Send + 'static,
seed: Result<[u8; 24], InitError>,
) -> Result<Self, InitError> {
Self::with_sink_and_seed(config, BasicSink(writer), seed, crate::ContextLoggingConfig::default())
}
fn with_sink_and_seed(config: ObserverConfig, writer: impl LogSink + Send + 'static,
seed: Result<[u8; 24], InitError>, context: crate::ContextLoggingConfig) -> Result<Self, InitError> {
if config.queue_capacity == 0 {
return Err(InitError::EmptyQueue);
}
let ids = IdGenerator::from_seed(seed?);
let (sender, receiver) = mpsc::sync_channel(config.queue_capacity);
let writer_source = Arc::new(Mutex::new(WriterSourceState::default()));
let worker_source = Arc::clone(&writer_source);
let worker = thread::Builder::new()
.name("saddle-log-writer".to_owned())
.spawn(move || write_records(receiver, writer, worker_source))
.map_err(InitError::Spawn)?;
Ok(Self {
inner: Arc::new(Inner {
sender,
dropped: AtomicU64::new(0),
accepting: AtomicBool::new(true),
emitting: AtomicUsize::new(0),
shutdown_started: AtomicBool::new(false),
worker: Mutex::new(Some(worker)),
process_log: Mutex::new(None),
context,
writer_source,
ids,
metrics: crate::metrics::Metrics::default(),
}),
})
}
pub(crate) fn new_trace_id(&self) -> TraceId {
self.inner.ids.trace_id()
}
pub(crate) fn new_span_id(&self) -> SpanId {
self.inner.ids.span_id()
}
pub(crate) fn emit(&self, mut record: LogRecord) {
if !self.inner.accepting.load(Ordering::Acquire) {
return;
}
self.inner.emitting.fetch_add(1, Ordering::AcqRel);
if !self.inner.accepting.load(Ordering::Acquire) {
self.inner.emitting.fetch_sub(1, Ordering::AcqRel);
return;
}
record.dropped_events = self.inner.dropped.swap(0, Ordering::Relaxed);
match self.inner.sender.try_send(Command::Record(record)) {
Ok(()) => {}
Err(TrySendError::Full(Command::Record(record))) => {
self.inner.metrics.logger_dropped(1);
self.inner
.dropped
.fetch_add(record.dropped_events.saturating_add(1), Ordering::Relaxed);
}
Err(TrySendError::Disconnected(_)) => {
self.inner.metrics.logger_dropped(1);
self.inner.metrics.logger_output(true);
self.inner.dropped.fetch_add(1, Ordering::Relaxed);
}
Err(TrySendError::Full(_)) => unreachable!("emit only sends records"),
}
self.inner.emitting.fetch_sub(1, Ordering::Release);
}
pub(crate) fn emit_context_record<P: crate::ContextRecordPayload>(&self, record: &crate::ContextRecord<'_, P>) -> crate::DiagnosticSubmission {
use crate::DiagnosticSubmission as S;
let mut result = S::Enqueued;
if self.inner.context.full.enabled {
result = self.enqueue_context(None, |writer| serde_json::to_writer(writer, record));
}
for (index, profile) in self.inner.context.summaries.iter().enumerate().filter(|(_, p)| p.enabled) {
let next = self.enqueue_context(Some(index), |writer| record.write_summary(profile, &mut &mut *writer));
if next != S::Enqueued { result = next; }
}
result
}
fn enqueue_context(&self, profile: Option<usize>, encode: impl Fn(&mut dyn Write) -> serde_json::Result<()>) -> crate::DiagnosticSubmission {
use crate::DiagnosticSubmission as S;
self.inner.emitting.fetch_add(1, Ordering::AcqRel);
let result = (|| {
if !self.inner.accepting.load(Ordering::Acquire) { return S::Closed; }
let Some(storage) = self.inner.process_log.lock().unwrap_or_else(|p| p.into_inner()).clone() else { return S::OutputUnavailable; };
let mut count = ContextByteCount(0);
if encode(&mut count).is_err() { return S::EncodingFailed; }
if count.0 <= 8191 { return self.enqueue_context_frame(profile, &storage, &encode); }
let output = OrdinarySegments { observer: self, storage: &storage, profile, record_id: self.inner.ids.span_id() };
crate::root_diagnostic::stream_record(&output, encode)
})();
if result != S::Enqueued { self.inner.dropped.fetch_add(1, Ordering::Relaxed); self.inner.metrics.logger_dropped(1); }
self.inner.emitting.fetch_sub(1, Ordering::Release);
result
}
fn enqueue_context_frame(&self, profile: Option<usize>, storage: &saddle_admission::ProcessLogStorage,
encode: &impl Fn(&mut dyn Write) -> serde_json::Result<()>) -> crate::DiagnosticSubmission {
use crate::DiagnosticSubmission as S;
if !self.inner.accepting.load(Ordering::Acquire) { return S::Closed; }
let mut count = ContextByteCount(0);
if encode(&mut count).is_err() || count.0 > 8191 { return S::EncodingFailed; }
let Ok(layout) = std::alloc::Layout::array::<u8>(count.0) else { return S::EncodingFailed; };
let Ok(permit) = storage.try_reserve(layout) else { return S::Full; };
let mut bytes = permit.allocate_exact(layout, || Vec::with_capacity(count.0));
if encode(bytes.get_mut()).is_err() || bytes.get().len() != count.0 { return S::EncodingFailed; }
let dropped = self.inner.dropped.swap(0, Ordering::Relaxed);
match self.inner.sender.try_send(Command::ContextProfile { bytes, profile, dropped }) {
Ok(()) => S::Enqueued,
Err(TrySendError::Full(command)) => { drop(command); self.inner.dropped.fetch_add(dropped, Ordering::Relaxed); S::Full },
Err(TrySendError::Disconnected(command)) => { drop(command); self.inner.dropped.fetch_add(dropped, Ordering::Relaxed); self.inner.metrics.logger_output(true); S::Closed },
}
}
pub(crate) fn emit_root_record<T: Serialize>(&self, record: &T) -> crate::DiagnosticSubmission {
use crate::DiagnosticSubmission as Submission;
self.inner.emitting.fetch_add(1, Ordering::AcqRel);
let result = (|| {
if !self.inner.accepting.load(Ordering::Acquire) {
return Submission::Closed;
}
let mut frame = crate::diagnostic::FixedDiagnosticBytes {
bytes: [0; 8192],
len: 0,
};
if serde_json::to_writer(&mut frame, record).is_err() {
return Submission::EncodingFailed;
}
let dropped = self.inner.dropped.swap(0, Ordering::Relaxed);
let bytes = frame.bytes[..frame.len].to_vec();
match self
.inner
.sender
.try_send(Command::RootRecord { bytes, dropped })
{
Ok(()) => Submission::Enqueued,
Err(TrySendError::Full(Command::RootRecord { dropped, .. })) => {
self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
Submission::Full
}
Err(TrySendError::Disconnected(_)) => {
self.inner.dropped.fetch_add(dropped, Ordering::Relaxed);
self.inner.metrics.logger_output(true);
Submission::Closed
}
Err(TrySendError::Full(_)) => unreachable!("root record submission"),
}
})();
if result != Submission::Enqueued {
self.inner.dropped.fetch_add(1, Ordering::Relaxed);
self.inner.metrics.logger_dropped(1);
}
self.inner.emitting.fetch_sub(1, Ordering::Release);
result
}
pub fn metrics_snapshot(&self) -> crate::MetricsSnapshot {
self.inner.metrics.snapshot()
}
pub fn flush(&self) -> FlushFuture {
if !self.inner.accepting.load(Ordering::Acquire) {
return FlushFuture {
completion: Completion::ready(Err(FlushError::WriterStopped)),
};
}
let completion = Completion::pending();
let future = FlushFuture {
completion: completion.clone(),
};
let inner = self.inner.clone();
if thread::Builder::new()
.name("saddle-log-flush".to_owned())
.spawn(move || coordinate_flush(inner, completion.clone()))
.is_err()
{
future
.completion
.complete(Err(FlushError::CoordinatorUnavailable));
}
future
}
pub fn shutdown(&self) -> FlushFuture {
if self.inner.shutdown_started.swap(true, Ordering::AcqRel) {
return FlushFuture {
completion: Completion::ready(Err(FlushError::AlreadyShuttingDown)),
};
}
self.inner.accepting.store(false, Ordering::Release);
let completion = Completion::pending();
let future = FlushFuture {
completion: completion.clone(),
};
let inner = self.inner.clone();
if thread::Builder::new()
.name("saddle-log-shutdown".to_owned())
.spawn(move || coordinate_shutdown(inner, completion.clone()))
.is_err()
{
self.inner.accepting.store(true, Ordering::Release);
self.inner.shutdown_started.store(false, Ordering::Release);
future
.completion
.complete(Err(FlushError::CoordinatorUnavailable));
}
future
}
pub fn dropped_events(&self) -> u64 {
self.inner.dropped.load(Ordering::Relaxed)
}
}
impl ComponentLifecycle for Observer {
fn name(&self) -> &'static str {
"observability"
}
fn start(&self) -> LifecycleFuture<'_> {
Box::pin(async { Ok(()) })
}
fn shutdown(&self) -> LifecycleFuture<'_> {
let shutdown = Observer::shutdown(self);
Box::pin(async move {
shutdown.await.map_err(|_| {
SaddleError::new(
ErrorKind::Infrastructure,
"observability.shutdown_failed",
"structured log shutdown failed",
)
})
})
}
}
impl Observer {
fn recorded_start<'a>(&'a self, application: &'a saddle_core::ContextLabel,
output: &'a crate::SourceOutput)
-> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
self.inner.writer_source.lock().unwrap_or_else(|e| e.into_inner()).context =
Some(WriterSourceContext {
application: application.clone(), output: output.clone(), primary: None,
});
Box::pin(async { Ok(()) })
}
fn recorded_shutdown<'a>(&'a self,
application: &'a saddle_core::ContextLabel,
output: &'a crate::SourceOutput,
primary: Option<saddle_core::DiagnosticOccurrence>)
-> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
self.inner.writer_source.lock().unwrap_or_else(|e| e.into_inner()).context =
Some(WriterSourceContext {
application: application.clone(), output: output.clone(), primary,
});
Box::pin(async move {
let result = Observer::shutdown(self).await;
let failure = self.inner.writer_source.lock()
.unwrap_or_else(|e| e.into_inner()).failure.take();
if let Some(failure) = failure {
return Err(crate::root_diagnostic::RecordedSaddleError::previously_written_writer_error(
failure.original, failure.diagnostic, failure.written));
}
crate::root_diagnostic::RecordedSaddleError::component_result(
result, application, output,
crate::root_diagnostic::ComponentSourceKind::Cleanup, primary,
ErrorKind::Infrastructure, "observability.shutdown_failed",
"structured log shutdown failed",
)
})
}
}
impl crate::root_diagnostic::RecordedComponentLifecycle for Observer {
fn name(&self) -> &'static str { "observability" }
fn start<'a>(&'a self, application: &'a saddle_core::ContextLabel,
output: &'a crate::SourceOutput)
-> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
Observer::recorded_start(self, application, output)
}
fn shutdown<'a>(&'a self, application: &'a saddle_core::ContextLabel,
output: &'a crate::SourceOutput)
-> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
crate::root_diagnostic::RecordedComponentLifecycle::shutdown_with_primary(
self, application, output, None)
}
fn shutdown_with_primary<'a>(&'a self,
application: &'a saddle_core::ContextLabel,
output: &'a crate::SourceOutput,
primary: Option<saddle_core::DiagnosticOccurrence>)
-> crate::root_diagnostic::RecordedLifecycleFuture<'a> {
Observer::recorded_shutdown(self, application, output, primary)
}
}
fn coordinate_flush(inner: Arc<Inner>, completion: Arc<Completion>) {
let dropped = inner.dropped.swap(0, Ordering::AcqRel);
let (response, receiver) = mpsc::channel();
let result = if inner
.sender
.send(Command::Flush {
unreported_dropped: dropped,
response,
})
.is_err()
{
Err(FlushError::WriterStopped)
} else {
receiver
.recv()
.map_err(|_| FlushError::WriterStopped)
.and_then(WorkerStatus::into_result)
};
completion.complete(result);
}
fn coordinate_shutdown(inner: Arc<Inner>, completion: Arc<Completion>) {
while inner.emitting.load(Ordering::Acquire) != 0 {
thread::yield_now();
}
let dropped = inner.dropped.swap(0, Ordering::AcqRel);
let (response, receiver) = mpsc::channel();
let mut result = if inner
.sender
.send(Command::Shutdown {
unreported_dropped: dropped,
response,
})
.is_err()
{
Err(FlushError::WriterStopped)
} else {
receiver
.recv()
.map_err(|_| FlushError::WriterStopped)
.and_then(WorkerStatus::into_result)
};
if let Some(worker) = inner.worker.lock().unwrap().take() {
if worker.join().is_err() {
result = Err(FlushError::WorkerPanicked);
}
}
inner.process_log.lock().unwrap_or_else(|p| p.into_inner()).take();
completion.complete(result);
}
fn write_records(receiver: Receiver<Command>, mut writer: impl LogSink,
source: Arc<Mutex<WriterSourceState>>) {
let mut status = WorkerStatus::default();
while let Ok(command) = receiver.recv() {
match command {
Command::Record(record) => write_record(&mut writer, record, &mut status, &source),
Command::ProcessCall { bytes, dropped } => {
status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
if status.failure.is_none() {
if let Err(error) = writer.write_all(bytes.get()) {
record_writer_failure(&source, OutputStage::Record, error);
status.failure = Some(OutputStage::Record);
} else if let Err(error) = writer.write_all(b"\n") {
record_writer_failure(&source, OutputStage::Newline, error);
status.failure = Some(OutputStage::Newline);
}
}
drop(bytes);
}
Command::ContextProfile { bytes, profile, dropped } => {
status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
if status.failure.is_none() {
let result = match profile {
Some(profile) => writer.write_profile(profile, bytes.get()),
None => writer.write_all(bytes.get()).and_then(|_| writer.write_all(b"\n")),
};
if let Err(error) = result { record_writer_failure(&source, OutputStage::Record, error); status.failure = Some(OutputStage::Record); }
}
drop(bytes);
}
Command::RootRecord { bytes, dropped } => {
status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
if status.failure.is_none() {
if let Err(error) = writer.write_all(&bytes) {
record_writer_failure(&source, OutputStage::Record, error);
status.failure = Some(OutputStage::Record);
} else if let Err(error) = writer.write_all(b"\n") {
record_writer_failure(&source, OutputStage::Newline, error);
status.failure = Some(OutputStage::Newline);
}
}
}
Command::Flush {
unreported_dropped,
response,
} => {
status.unreported_dropped =
status.unreported_dropped.saturating_add(unreported_dropped);
flush_writer(&mut writer, &mut status, &source);
let _ = response.send(status);
}
Command::Shutdown {
unreported_dropped,
response,
} => {
status.unreported_dropped =
status.unreported_dropped.saturating_add(unreported_dropped);
flush_writer(&mut writer, &mut status, &source);
let _ = response.send(status);
break;
}
}
}
}
fn write_record(writer: &mut impl Write, record: LogRecord, status: &mut WorkerStatus,
source: &Arc<Mutex<WriterSourceState>>) {
if status.failure.is_some() {
status.unreported_dropped = status
.unreported_dropped
.saturating_add(record.dropped_events);
return;
}
let dropped = record.dropped_events;
let bytes = match serde_json::to_vec(&record) {
Ok(bytes) => bytes,
Err(_) => {
status.failure = Some(OutputStage::Serialize);
status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
return;
}
};
if let Err(error) = writer.write_all(&bytes) {
record_writer_failure(source, OutputStage::Record, error);
status.failure = Some(OutputStage::Record);
status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
return;
}
if let Err(error) = writer.write_all(b"\n") {
record_writer_failure(source, OutputStage::Newline, error);
status.failure = Some(OutputStage::Newline);
status.unreported_dropped = status.unreported_dropped.saturating_add(dropped);
}
}
fn flush_writer(writer: &mut impl Write, status: &mut WorkerStatus,
source: &Arc<Mutex<WriterSourceState>>) {
if status.failure.is_some() { return; }
if let Err(error) = writer.flush() {
record_writer_failure(source, OutputStage::Flush, error);
status.failure = Some(OutputStage::Flush);
}
}
#[cfg(test)]
mod tests {
use std::{
sync::{Condvar, mpsc},
task::{Wake, Waker},
};
use super::*;
#[test]
fn recorded_observer_shutdown_preserves_actual_flush_error() {
let directory = std::env::temp_dir().join(format!(
"saddle-observer-recorded-{}-{}", std::process::id(),
SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos()));
std::fs::create_dir(&directory).unwrap();
let mut output = crate::EmergencyDiagnostics::start_checked(
&crate::FileLoggingConfig::new(directory.clone(), crate::Rotation::Daily)).unwrap();
let selected = output.source_output().unwrap();
let target = output.target().to_owned();
let logger = observer(io::sink(), 1);
block_on(Observer::shutdown(&logger)).unwrap();
let app = saddle_core::ContextLabel::checked("recorded-observer-test").unwrap();
let failure = block_on(
crate::root_diagnostic::RecordedComponentLifecycle::shutdown(
&logger, &app, &selected)).err().unwrap();
assert_eq!(failure.safe().code(), "observability.shutdown_failed");
assert!(failure.original_if_unconfirmed().is_none());
let record = std::fs::read_to_string(&target).unwrap();
assert!(record.contains("AlreadyShuttingDown"));
assert!(record.contains(&failure.safe().diagnostic().unwrap().id().to_string()));
drop(selected);
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
while output.shutdown() == crate::DiagnosticShutdown::Pending
&& std::time::Instant::now() < deadline { std::thread::yield_now(); }
assert_eq!(output.shutdown(), crate::DiagnosticShutdown::Finished);
drop(output);
std::fs::remove_dir_all(directory).unwrap();
}
#[derive(Debug)]
struct WriterChain {
message: &'static str,
cause: io::Error,
}
impl fmt::Display for WriterChain {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(self.message)
}
}
impl Error for WriterChain {
fn source(&self) -> Option<&(dyn Error + 'static)> { Some(&self.cause) }
}
struct ChainWriter {
fail_flush: bool,
release_write: Option<(mpsc::Sender<()>, Arc<(Mutex<bool>, Condvar)>)>,
}
impl Write for ChainWriter {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
if self.fail_flush { Ok(bytes.len()) } else {
if let Some((entered, release)) = self.release_write.take() {
entered.send(()).unwrap();
let (lock, condition) = &*release;
let mut ready = lock.lock().unwrap();
while !*ready { ready = condition.wait(ready).unwrap(); }
}
Err(io::Error::other(WriterChain {
message: "writer top original 4931",
cause: io::Error::other("writer nested original 4932"),
}))
}
}
fn flush(&mut self) -> io::Result<()> {
if self.fail_flush {
Err(io::Error::other(WriterChain {
message: "flush top original 4933",
cause: io::Error::other("flush nested original 4934"),
}))
} else { Ok(()) }
}
}
#[test]
fn writer_write_and_flush_originals_precede_formal_shutdown_completion() {
use crate::root_diagnostic::RecordedComponentLifecycle;
for flush in [false, true] {
let directory = std::env::temp_dir().join(format!(
"saddle-writer-original-{flush}-{}-{}", std::process::id(),
SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos()));
std::fs::create_dir(&directory).unwrap();
let mut emergency = crate::EmergencyDiagnostics::start_checked(
&crate::FileLoggingConfig::new(directory.clone(), crate::Rotation::Daily)).unwrap();
let output = emergency.source_output().unwrap();
let target = emergency.target().to_owned();
let (entered_sender, entered_receiver) = mpsc::channel();
let release = Arc::new((Mutex::new(false), Condvar::new()));
let logger = observer(ChainWriter {
fail_flush: flush,
release_write: (!flush).then(|| (entered_sender, Arc::clone(&release))),
}, 8);
let application = saddle_core::ContextLabel::checked("writer-original-test").unwrap();
assert!(block_on(RecordedComponentLifecycle::start(
&logger, &application, &output)).is_ok());
let primary = saddle_core::Diagnostic::capture(
saddle_core::DiagnosticCategory::UnexpectedError,
saddle_core::CaptureSite::Origin,
saddle_core::DiagnosticCause::new(
saddle_core::DiagnosticStage::ShutdownComponent,
saddle_core::DiagnosticCode::new("test.writer_primary").unwrap()),
).occurrence();
if !flush { logger.emit(LogRecord::new(EventLevel::Info, "writer.failure")); }
if !flush { entered_receiver.recv_timeout(std::time::Duration::from_secs(5)).unwrap(); }
let shutdown = RecordedComponentLifecycle::shutdown_with_primary(
&logger, &application, &output, Some(primary));
if !flush {
let (lock, condition) = &*release;
*lock.lock().unwrap() = true;
condition.notify_all();
}
let failure = block_on(shutdown).err().unwrap();
assert_eq!(failure.safe().code(), "observability.shutdown_failed");
assert!(failure.original_if_unconfirmed().is_none());
let diagnostic = failure.safe().diagnostic().unwrap();
let record = std::fs::read_to_string(&target).unwrap();
let (top, nested) = if flush {
("flush top original 4933", "flush nested original 4934")
} else { ("writer top original 4931", "writer nested original 4932") };
assert!(record.contains(top), "{record}");
assert!(record.contains(nested), "{record}");
assert!(record.contains(if flush { "Flush" } else { "Record" }), "{record}");
assert!(record.contains("shutdown_logger"), "{record}");
assert!(record.contains(&diagnostic.id().to_string()));
assert!(record.contains("complete"), "{record}");
let originals: Vec<serde_json::Value> = record.lines()
.map(|line| serde_json::from_str(line).unwrap())
.filter(|row: &serde_json::Value|
row["event"] == "request_error_original").collect();
assert!(!originals.is_empty());
assert!(originals.iter().any(|row|
row["channel"] == "description" && row["cause_depth"] == 0
&& row["payload"].as_str().unwrap().contains(top)));
assert!(originals.iter().any(|row|
row["channel"] == "debug" && row["cause_depth"] == 0
&& row["payload"].as_str().unwrap().contains(top)));
assert!(originals.iter().any(|row|
row["channel"] == "description"
&& row["cause_depth"].as_u64().unwrap() > 0
&& row["payload"].as_str().unwrap().contains(nested)));
assert!(originals.iter().any(|row|
row["channel"] == "terminal"
&& row["state"] == "exposed_chain_complete"));
assert!(originals.iter().all(|row|
row["occurrence"]["diagnostic_id"] == diagnostic.id()));
let expected = serde_json::to_value(primary).unwrap();
assert!(originals.iter().all(|row|
row["occurrence"]["primary_diagnostic_id"]
== expected["diagnostic_id"]));
drop(logger);
drop(output);
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
while emergency.shutdown() == crate::DiagnosticShutdown::Pending
&& std::time::Instant::now() < deadline { std::thread::yield_now(); }
assert_eq!(emergency.shutdown(), crate::DiagnosticShutdown::Finished);
drop(emergency);
std::fs::remove_dir_all(directory).unwrap();
}
}
#[test]
fn writer_original_without_selected_output_remains_unconfirmed() {
use crate::root_diagnostic::RecordedComponentLifecycle;
let logger = observer(ChainWriter { fail_flush: false, release_write: None }, 8);
logger.emit(LogRecord::new(EventLevel::Info, "writer.no_output"));
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
while logger.inner.writer_source.lock().unwrap_or_else(|e| e.into_inner())
.failure.is_none() && std::time::Instant::now() < deadline {
std::thread::yield_now();
}
assert!(logger.inner.writer_source.lock().unwrap_or_else(|e| e.into_inner())
.failure.as_ref().is_some_and(|failure| failure.written.is_none()));
let directory = std::env::temp_dir().join(format!(
"saddle-writer-unconfirmed-{}-{}", std::process::id(),
SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos()));
std::fs::create_dir(&directory).unwrap();
let mut emergency = crate::EmergencyDiagnostics::start_checked(
&crate::FileLoggingConfig::new(directory.clone(), crate::Rotation::Daily)).unwrap();
let target = emergency.target().to_owned();
let output = emergency.source_output().unwrap();
let application = saddle_core::ContextLabel::checked("writer-unconfirmed-test").unwrap();
let failure = block_on(RecordedComponentLifecycle::shutdown(
&logger, &application, &output)).err().unwrap();
assert_eq!(failure.safe().code(), "observability.shutdown_failed");
let original = failure.original_if_unconfirmed().unwrap();
assert!(original.to_string().contains("writer top original 4931"));
assert!(original.source().unwrap().to_string()
.contains("writer top original 4931"));
assert!(!std::fs::read_to_string(&target).unwrap()
.contains("writer top original 4931"));
drop(logger);
drop(output);
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
while emergency.shutdown() == crate::DiagnosticShutdown::Pending
&& std::time::Instant::now() < deadline { std::thread::yield_now(); }
assert_eq!(emergency.shutdown(), crate::DiagnosticShutdown::Finished);
drop(emergency);
std::fs::remove_dir_all(directory).unwrap();
}
struct ThreadWaker(thread::Thread);
impl Wake for ThreadWaker {
fn wake(self: Arc<Self>) {
self.0.unpark();
}
}
fn block_on<T>(future: impl Future<Output = T>) -> T {
let waker = Waker::from(Arc::new(ThreadWaker(thread::current())));
let mut context = Context::from_waker(&waker);
let mut future = std::pin::pin!(future);
loop {
match future.as_mut().poll(&mut context) {
Poll::Ready(output) => return output,
Poll::Pending => thread::park(),
}
}
}
struct BlockingWriter {
entered: Option<mpsc::Sender<()>>,
release: Arc<(Mutex<bool>, Condvar)>,
fail_after_release: bool,
}
impl Write for BlockingWriter {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
if let Some(entered) = self.entered.take() {
let _ = entered.send(());
let (lock, condition) = &*self.release;
let mut released = lock.lock().unwrap();
while !*released {
released = condition.wait(released).unwrap();
}
if self.fail_after_release {
return Err(io::Error::other("injected writer failure"));
}
}
Ok(bytes.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
struct FailingWriter {
fail_write: Option<usize>,
writes: usize,
fail_flush: bool,
}
impl Write for FailingWriter {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
self.writes += 1;
if self.fail_write == Some(self.writes) {
Err(io::Error::other("injected writer failure"))
} else {
Ok(bytes.len())
}
}
fn flush(&mut self) -> io::Result<()> {
if self.fail_flush {
Err(io::Error::other("injected flush failure"))
} else {
Ok(())
}
}
}
fn observer(writer: impl Write + Send + 'static, capacity: usize) -> Observer {
Observer::with_writer_and_seed(
ObserverConfig {
queue_capacity: capacity,
},
writer,
Ok([7; 24]),
)
.unwrap()
}
fn context_profile_process() -> saddle_admission::ProfuseGwLightweightProcessOwner {
use saddle_admission::{StorageDemand, bind_deployment_resource_budget_bootstrap, freeze_deployment_resource_budget, prepare_profusegw_lightweight_profile};
use saddle_core::{BootstrapRendezvousIssuer, GeneratedApplicationFreezeSource, ListenerStartupFreezeSource, pair_bootstrap_rendezvous};
let pending = freeze_deployment_resource_budget(
1, 32, 5_000, 1, 1, 1_000_000, 1_000_000, 1_000_000, 1_000_000,
).unwrap().reserve_known_database_memory(0, 0, 0).ok().unwrap();
let (application, listener) = BootstrapRendezvousIssuer::issue()
.freeze_application(GeneratedApplicationFreezeSource::new(
"app", b"descriptor", &["route"],
)).unwrap();
let listener = listener.freeze_listener(ListenerStartupFreezeSource::new(
"app", "127.0.0.1:8000".parse().unwrap(),
"127.0.0.1:9000".parse().unwrap(),
std::time::Duration::from_millis(5_000),
)).ok().unwrap();
let (whole, receipt) = pair_bootstrap_rendezvous(application, listener).ok().unwrap();
let budget = bind_deployment_resource_budget_bootstrap(pending, whole, receipt).ok().unwrap();
let process = prepare_profusegw_lightweight_profile(budget).ok().unwrap();
process
}
#[test]
fn large_managed_context_and_summary_reassemble_after_request_close() {
use std::collections::BTreeMap;
use saddle_core::{ContextFact, ContextLabel, RequestIdentityGroup, RequestLocalFacts, RequestRootPublisher, RequestViewPhase};
let directory = std::env::temp_dir().join(format!("context-segments-{}-{}", std::process::id(), SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos()));
let blob = "界\n\"".repeat(12000);
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());
let mut root = RequestRootPublisher::create(ContextLabel::checked("app").unwrap(), ContextFact::NotEstablished).unwrap();
let call = saddle_core::CallContext::new("app".into(), "profusegw".into(), "profusegw".into(), "route".into(), TraceId::from_u128(1), SpanId::from_u64(1));
root.publish(RequestIdentityGroup::from_validated(&call, "request", "route", 1, ContextFact::Unavailable).unwrap()).ok().unwrap();
let process = context_profile_process();
let storage_root = process.try_process_storage(saddle_admission::StorageDemand::separate(&[]).unwrap()).unwrap();
let baseline = process.resource_snapshot().framework_charged;
let execution = match process.verified_profile().try_admit() { saddle_admission::ProfuseGwLightweightAdmissionOutcome::Ready(p) => p.into_execution(), _ => panic!("original capacity") };
let memory = execution.request_memory();
let bytes = memory.try_bytes(raw.as_bytes()).unwrap();
let document = memory.decode_input_range(&bytes, 0..bytes.len()).unwrap().into_shared().unwrap();
drop(bytes);
root.publish_ingress(Arc::new(document), "call", 1234).unwrap();
let business = memory.update_business(None, "/private", Some(br#"{"child":"hidden"}"#), true, saddle_admission::BusinessLimits::default()).unwrap();
let view = root.reference().view(RequestLocalFacts::new(RequestViewPhase::Reading)).with_business(Arc::new(business));
let mut config = crate::ContextLoggingConfig::default();
config.summaries.push(crate::SummaryProfile { name: "selected".into(), enabled: true, format: crate::SummaryFormat::Json,
fields: BTreeMap::from([("data".into(), "/context/inputInfo/requestData".into()), ("secret".into(), "/context/business/private/child".into()), ("revision".into(), "/context/_meta/revision".into())]),
missing: crate::MissingField::Null, template: None });
let file = FileLoggingConfig::new(&directory, crate::Rotation::Daily).with_context(config.clone()).unwrap();
let outputs = CalendarOutputs { summaries: vec![Some(CalendarFileWriter::open(file.for_summary("selected")).unwrap())], main: CalendarFileWriter::open(file).unwrap() };
let observer = Observer::with_sink_and_seed(ObserverConfig { queue_capacity: 32 }, outputs, Ok([7; 24]), config).unwrap();
observer.install_process_log_storage(process.process_log_storage());
let mut owned_context = None;
drop(memory.framework_output(|_| {
owned_context = Some(observer.capture_terminal_context(&view.unified_context()).unwrap());
Ok(())
}).unwrap());
let available = process.resource_snapshot().framework_capacity - process.resource_snapshot().framework_charged;
let filler = process.try_test_process_storage_child(saddle_admission::StorageDemand::separate(&[(
std::alloc::Layout::array::<u8>(available - 64 - std::mem::size_of::<saddle_admission::StoragePermit>()).unwrap(), 1,
)]).unwrap()).unwrap();
let charged = process.resource_snapshot().framework_charged;
drop(memory.framework_output(|_| {
assert!(matches!(observer.capture_terminal_context(&view.unified_context()), Err(crate::DiagnosticSubmission::Full)));
Ok(())
}).unwrap());
assert_eq!(process.resource_snapshot().framework_charged, charged);
drop(filler);
let emitted = memory.framework_output(|_| Ok(crate::RootDiagnosticScope::new(&view, None).ordinary(&observer, crate::RootRequestEvent::Handler, Default::default()))).unwrap();
assert_eq!(*emitted.get(), crate::DiagnosticSubmission::Enqueued);
drop(emitted);
struct CapturedBlocking { block: BlockingWriter, data: Arc<Mutex<Vec<u8>>> }
impl Write for CapturedBlocking {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
self.block.write(bytes)?; self.data.lock().unwrap().extend_from_slice(bytes); Ok(bytes.len())
}
fn flush(&mut self) -> io::Result<()> { self.block.flush() }
}
let (entered, entered_rx) = mpsc::channel();
let release = Arc::new((Mutex::new(false), Condvar::new()));
let captured = Arc::new(Mutex::new(Vec::new()));
let refused = Observer::with_sink_and_seed(ObserverConfig { queue_capacity: 1 }, BasicSink(CapturedBlocking {
block: BlockingWriter { entered: Some(entered), release: release.clone(), fail_after_release: false }, data: captured.clone(),
}), Ok([8; 24]), crate::ContextLoggingConfig::default()).unwrap();
refused.install_process_log_storage(process.process_log_storage());
assert_eq!(refused.emit_root_record(&0), crate::DiagnosticSubmission::Enqueued);
entered_rx.recv_timeout(std::time::Duration::from_secs(2)).unwrap();
let rejection = memory.framework_output(|_| Ok(crate::RootDiagnosticScope::new(&view, None).ordinary(&refused, crate::RootRequestEvent::Handler, Default::default()))).unwrap();
assert_eq!(*rejection.get(), crate::DiagnosticSubmission::Full); drop(rejection);
assert_eq!(refused.dropped_events(), 1);
*release.0.lock().unwrap() = true; release.1.notify_all();
assert!(matches!(block_on(refused.flush()), Err(FlushError::DroppedEvents(_))));
let frames: Vec<serde_json::Value> = std::str::from_utf8(&captured.lock().unwrap()).unwrap().lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|row| row["event"] == "context_record_segment").collect();
assert_eq!(frames.len(), 1); assert_eq!(frames[0]["sequence"], 0); assert_eq!(frames[0]["state"], "continuation");
let _ = block_on(refused.shutdown()); drop(refused);
drop((view, root, memory));
let terminal = execution.cancel_observed();
assert_eq!(terminal.audit().report.unwrap().escape_allocations, 0);
assert_eq!(process.resource_snapshot().active_accounts, 0);
block_on(observer.flush()).unwrap();
fn readback(path: std::path::PathBuf) -> String {
let text = std::fs::read_to_string(path).unwrap(); let mut restored = String::new(); let mut id = None; let mut complete = false;
for (sequence, line) in text.lines().enumerate() {
assert!(line.len() <= 8191);
let frame: serde_json::Value = serde_json::from_str(line).unwrap();
assert_eq!(frame["event"], "context_record_segment");
assert_eq!(frame["format"], "json");
assert_eq!(frame["sequence"], sequence);
if let Some(id) = &id { assert_eq!(&frame["recordId"], id); } else { id = Some(frame["recordId"].clone()); }
assert!(!complete); restored.push_str(frame["payload"].as_str().unwrap()); complete = frame["state"] == "record_complete";
}
assert!(complete); assert!(id.is_some()); assert!(restored.contains("123456789012345678901234567890")); restored
}
let full: serde_json::Value = serde_json::from_str(&readback(directory.join("saddle.log"))).unwrap();
let summary: serde_json::Value = serde_json::from_str(&readback(directory.join("selected.summary.log"))).unwrap();
assert_eq!(full["context"]["inputInfo"]["requestData"]["future"][2]["blob"], blob);
assert_eq!(summary["data"], full["context"]["inputInfo"]["requestData"]);
assert_eq!(summary["secret"], "<redacted>"); assert_eq!(full["context"]["business"]["private"], "<redacted>");
assert_eq!(summary["revision"], full["context"]["_meta"]["revision"]);
let owned_context = owned_context.unwrap();
let emergency = crate::EmergencyDiagnostics::start(&FileLoggingConfig::new(&directory, crate::Rotation::Daily)).unwrap();
let diagnostic = saddle_core::BoundedDiagnostic::capture(saddle_core::DiagnosticCategory::UnexpectedError,
saddle_core::CaptureSite::FirstObserved, saddle_core::BoundedDiagnosticCause::new(
saddle_core::DiagnosticStage::FinalizerResource, saddle_core::DiagnosticCode::new("test.owned_audit").unwrap()));
let marker = diagnostic.occurrence();
assert_eq!(crate::root_diagnostic::request_terminal_audit_passed_owned(Some(&emergency.handle()),
Ok(&owned_context), marker, terminal.audit()), crate::DiagnosticSubmission::Written);
assert_eq!(crate::root_diagnostic::request_terminal_audit_cleanup_owned(Some(&emergency.handle()),
Ok(&owned_context), marker, &process.resource_snapshot()), crate::DiagnosticSubmission::Written);
let failure = saddle_admission::ProfuseGwTerminalAuditFailure { error: saddle_admission::AdmissionError::AccountClosed,
audit: None, database: saddle_admission::ProfuseGwDatabaseDisposition::NotUsed };
let written = crate::root_diagnostic::RecordedTerminalFailure::audit_owned(Some(&emergency.handle()),
Ok(&owned_context), Some(marker), failure);
assert!(written.original_confirmed());
assert_eq!(written.disposition_submission(), Some(crate::DiagnosticSubmission::Written));
let failure = saddle_admission::ProfuseGwTerminalAuditFailure { error: saddle_admission::AdmissionError::AccountClosed,
audit: None, database: saddle_admission::ProfuseGwDatabaseDisposition::NotUsed };
let rejected = crate::root_diagnostic::RecordedTerminalFailure::audit_owned(Some(&emergency.handle()),
Err(crate::DiagnosticSubmission::Full), Some(marker), failure);
assert!(!rejected.original_confirmed(), "a written error without its requested context cannot mint a positive receipt");
let text = std::fs::read_to_string(directory.join("saddle.emergency.log")).unwrap();
for event in ["request_audit_passed", "request_audit_cleanup"] {
let frames: Vec<serde_json::Value> = text.lines().map(|line| serde_json::from_str::<serde_json::Value>(line).unwrap())
.filter(|row| row["event"] == event).collect();
assert!(frames.len() > 1);
let mut restored = String::new();
for (sequence, frame) in frames.iter().enumerate() {
assert_eq!(frame["sequence"], sequence);
restored.push_str(frame["payload"].as_str().unwrap());
assert_eq!(frame["state"], if sequence + 1 == frames.len() { "record_complete" } else { "continuation" });
}
assert!(restored.contains("123456789012345678901234567890"));
let record: serde_json::Value = serde_json::from_str(&restored).unwrap();
assert_eq!(record["context"], full["context"]);
assert_eq!(record["recordSchema"], "1.0");
}
drop(emergency);
let encoded = serde_json::to_string(&owned_context).unwrap();
assert!(encoded.contains("123456789012345678901234567890"));
let restored: serde_json::Value = serde_json::from_str(&encoded).unwrap();
assert_eq!(restored, full["context"]);
drop(owned_context);
assert_eq!(process.resource_snapshot().framework_charged, baseline);
block_on(observer.shutdown()).unwrap(); drop(observer); drop(storage_root); process.finish().unwrap();
std::fs::remove_dir_all(directory).unwrap();
}
#[test]
fn configured_context_profiles_write_real_files_for_all_four_combinations() {
use std::collections::BTreeMap;
use crate::{ContextLoggingConfig, SummaryProfile, SummaryFormat, MissingField};
let root = std::env::temp_dir().join(format!("context-profiles-{}-{}", std::process::id(), SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos()));
let publisher = saddle_core::RequestRootPublisher::create(saddle_core::ContextLabel::checked("app").unwrap(), saddle_core::ContextFact::NotEstablished).unwrap();
let view = publisher.reference().view(saddle_core::RequestLocalFacts::new(saddle_core::RequestViewPhase::Reading));
use saddle_admission::StorageDemand;
let process = context_profile_process();
let storage_root = process.try_process_storage(StorageDemand::separate(&[]).unwrap()).unwrap();
let baseline = process.resource_snapshot().framework_charged;
for full in [false, true] { for summary in [false, true] {
let directory = root.join(format!("{full}-{summary}"));
let mut context = ContextLoggingConfig::default(); context.full.enabled = full;
let fields = BTreeMap::from([("event".into(), "/event".into()), ("revision".into(), "/context/_meta/revision".into()),
("missing".into(), "/context/business/future".into()), ("elapsed".into(), "/payload/elapsedMs".into())]);
context.summaries = vec![
SummaryProfile { name: "json".into(), enabled: summary, format: SummaryFormat::Json, fields: fields.clone(), missing: MissingField::Null, template: None },
SummaryProfile { name: "text".into(), enabled: summary, format: SummaryFormat::Text, fields, missing: MissingField::Omit, template: Some("event=${event} absent=${missing} revision=${revision}".into()) },
]; context.validate().unwrap();
let file = FileLoggingConfig::new(&directory, crate::Rotation::Daily).with_context(context.clone()).unwrap();
let summaries = context.summaries.iter().map(|profile| if profile.enabled { Some(CalendarFileWriter::open(file.for_summary(&profile.name)).unwrap()) } else { None }).collect();
let outputs = CalendarOutputs { main: CalendarFileWriter::open(file).unwrap(), summaries };
let logger = Observer::with_sink_and_seed(ObserverConfig { queue_capacity: 8 }, outputs, Ok([7; 24]), context).unwrap();
logger.install_process_log_storage(process.process_log_storage());
assert_eq!(crate::RootDiagnosticScope::new(&view, None).ordinary(&logger, crate::RootRequestEvent::Handler, Default::default()), crate::DiagnosticSubmission::Enqueued);
block_on(logger.flush()).unwrap();
let main = std::fs::read_to_string(directory.join("saddle.log")).unwrap();
assert_eq!(!main.is_empty(), full);
if full {
let record: serde_json::Value = serde_json::from_str(main.trim()).unwrap();
assert_eq!(record["recordSchema"], "1.0"); assert_eq!(record["context"]["schemaVersion"], "1.0");
assert_eq!(record["payload"]["stage"], "handler");
assert!(record.get("schema_version").is_none());
}
if summary {
let json: serde_json::Value = serde_json::from_str(std::fs::read_to_string(directory.join("json.summary.log")).unwrap().trim()).unwrap();
assert!(json["missing"].is_null()); assert!(json["elapsed"].is_null());
assert_eq!(json["event"], "request_stage");
let text = std::fs::read_to_string(directory.join("text.summary.log")).unwrap();
assert_eq!(text, format!("event=\"request_stage\" absent= revision={}\n", serde_json::to_string(&json["revision"]).unwrap()));
if full { let main: serde_json::Value = serde_json::from_str(main.trim()).unwrap(); assert_eq!(json["revision"], main["context"]["_meta"]["revision"]); }
} else { assert!(!directory.join("json.summary.log").exists()); }
block_on(logger.shutdown()).unwrap(); drop(logger);
assert_eq!(process.resource_snapshot().framework_charged, baseline);
} }
drop(storage_root); process.finish().unwrap();
std::fs::remove_dir_all(root).unwrap();
}
#[test]
fn managed_call_log_outlives_request_account_while_writer_is_blocked() {
assert_eq!(std::mem::size_of::<Command>(), std::mem::size_of::<LogRecord>());
assert_eq!(std::mem::size_of::<Command>(), 192);
assert_eq!(ObserverConfig::default().queue_capacity * (8 + std::mem::size_of::<Command>()),
1_638_400);
use saddle_admission::{
ProfuseGwLightweightObservedAdmissionOutcome, StorageDemand,
bind_deployment_resource_budget_bootstrap, freeze_deployment_resource_budget,
prepare_profusegw_lightweight_profile,
};
use saddle_core::{
BootstrapRendezvousIssuer, GeneratedApplicationFreezeSource,
ListenerStartupFreezeSource, RpcCorrelationId, pair_bootstrap_rendezvous,
};
let pending = freeze_deployment_resource_budget(
1, 32, 5_000, 1, 1, 1_000_000, 1_000_000, 1_000_000, 1_000_000,
).unwrap().reserve_known_database_memory(0, 0, 0).ok().unwrap();
let (application, listener) = BootstrapRendezvousIssuer::issue()
.freeze_application(GeneratedApplicationFreezeSource::new(
"app", b"descriptor", &["route"],
)).unwrap();
let listener = listener.freeze_listener(ListenerStartupFreezeSource::new(
"app", "127.0.0.1:8000".parse().unwrap(),
"127.0.0.1:9000".parse().unwrap(),
std::time::Duration::from_millis(5_000),
)).ok().unwrap();
let (whole, receipt) = pair_bootstrap_rendezvous(application, listener).ok().unwrap();
let budget = bind_deployment_resource_budget_bootstrap(pending, whole, receipt).ok().unwrap();
let process = prepare_profusegw_lightweight_profile(budget).ok().unwrap();
let root = process.try_process_storage(StorageDemand::separate(&[(
std::alloc::Layout::array::<u8>(16).unwrap(), 1,
)]).unwrap()).unwrap();
let baseline = process.resource_snapshot().framework_charged;
let (entered, entered_rx) = mpsc::channel();
let release = Arc::new((Mutex::new(false), Condvar::new()));
let logger = observer(BlockingWriter {
entered: Some(entered), release: Arc::clone(&release), fail_after_release: false,
}, 2);
logger.install_process_log_storage(process.process_log_storage());
let (call, _) = logger.start_managed_external_call_with_rpc(
"app", "zone", "interface", "operation", None,
RpcCorrelationId::new("0").unwrap(),
).unwrap();
entered_rx.recv_timeout(std::time::Duration::from_secs(2)).unwrap();
let admitted = match process.verified_profile().try_admit_observed() {
ProfuseGwLightweightObservedAdmissionOutcome::Ready(permit, _) => permit,
_ => panic!("fixture admission failed"),
};
let original = saddle_core::SaddleError::new(
saddle_core::ErrorKind::Internal, "original.failure", "original",
);
call.fail(&original);
assert_eq!(original.code(), "original.failure");
admitted.into_execution().cancel_observed();
assert_eq!(process.resource_snapshot().active_accounts, 0);
assert!(process.resource_snapshot().framework_charged > baseline,
"queued process log must remain charged after request cleanup");
let (full, _) = logger.start_managed_external_call_with_rpc(
"app", "zone", "interface", "full", None,
RpcCorrelationId::new("1").unwrap(),
).unwrap();
full.fail(&original);
assert!(logger.dropped_events() > 0);
*release.0.lock().unwrap() = true;
release.1.notify_all();
assert!(matches!(block_on(logger.flush()), Err(FlushError::DroppedEvents(_))));
assert_eq!(process.resource_snapshot().framework_charged, baseline);
let available = process.resource_snapshot().framework_capacity
- process.resource_snapshot().framework_charged;
let filler = process.try_test_process_storage_child(StorageDemand::separate(&[(
std::alloc::Layout::array::<u8>(
available - 64 - std::mem::size_of::<saddle_admission::StoragePermit>(),
).unwrap(), 1,
)]).unwrap()).unwrap();
let charged = process.resource_snapshot().framework_charged;
let before_drop = logger.dropped_events();
let (short, _) = logger.start_managed_external_call_with_rpc(
"app", "zone", "interface", "short", None,
RpcCorrelationId::new("2").unwrap(),
).unwrap();
short.fail(&original);
assert!(logger.dropped_events() >= before_drop + 2);
assert_eq!(process.resource_snapshot().framework_charged, charged);
drop(filler);
assert!(matches!(block_on(logger.shutdown()), Err(FlushError::DroppedEvents(_))));
let after_shutdown = process.resource_snapshot().framework_charged;
let (closed, _) = logger.start_managed_external_call_with_rpc(
"app", "zone", "interface", "closed", None,
RpcCorrelationId::new("3").unwrap(),
).unwrap();
closed.fail(&original);
assert_eq!(process.resource_snapshot().framework_charged, after_shutdown);
assert_eq!(process.resource_snapshot().framework_charged, baseline);
drop(logger);
drop(root);
process.finish().unwrap();
}
#[test]
fn root_contract_ordinary_uses_same_view_and_bounded_queue() {
use crate::{
DiagnosticSubmission, RootDiagnosticScope, RootOutcomeFacts, RootRequestEvent,
};
use saddle_core::{
ContextFact, ContextLabel, RequestLocalFacts, RequestRootPublisher, RequestViewPhase,
};
let root = RequestRootPublisher::create(
ContextLabel::checked("app").unwrap(),
ContextFact::NotEstablished,
)
.unwrap();
let view = root
.reference()
.view(RequestLocalFacts::new(RequestViewPhase::Handler));
let expected = serde_json::to_value(&view).unwrap();
let mut frame = crate::diagnostic::FixedDiagnosticBytes {
bytes: [0; 8192],
len: 0,
};
serde_json::to_writer(&mut frame, &view).unwrap();
let ordinary_context: serde_json::Value =
serde_json::from_slice(&frame.bytes[..frame.len]).unwrap();
assert_eq!(ordinary_context, expected);
let (entered, wait) = mpsc::channel();
let release = Arc::new((Mutex::new(false), Condvar::new()));
struct Captured {
writer: BlockingWriter,
bytes: Arc<Mutex<Vec<u8>>>,
}
impl Write for Captured {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
let n = self.writer.write(bytes)?;
self.bytes.lock().unwrap().extend_from_slice(&bytes[..n]);
Ok(n)
}
fn flush(&mut self) -> io::Result<()> {
self.writer.flush()
}
}
struct Unblock(Arc<(Mutex<bool>, Condvar)>);
impl Drop for Unblock {
fn drop(&mut self) {
*self.0.0.lock().unwrap() = true;
self.0.1.notify_all();
}
}
let unblock = Unblock(Arc::clone(&release));
let bytes = Arc::new(Mutex::new(Vec::new()));
let logger = observer(
Captured {
writer: BlockingWriter {
entered: Some(entered),
release: Arc::clone(&release),
fail_after_release: false,
},
bytes: Arc::clone(&bytes),
},
1,
);
let scope = RootDiagnosticScope::new(&view, None);
assert_eq!(
scope.ordinary(
&logger,
RootRequestEvent::Handler,
RootOutcomeFacts::default()
),
DiagnosticSubmission::Enqueued
);
wait.recv_timeout(std::time::Duration::from_secs(5))
.unwrap();
assert_eq!(
scope.ordinary(
&logger,
RootRequestEvent::Response,
RootOutcomeFacts::default()
),
DiagnosticSubmission::Enqueued
);
assert_eq!(
scope.ordinary(
&logger,
RootRequestEvent::Response,
RootOutcomeFacts::default()
),
DiagnosticSubmission::Full
);
drop((view, root));
drop(unblock);
assert!(matches!(
block_on(logger.shutdown()),
Err(FlushError::DroppedEvents(1))
));
let bytes = bytes.lock().unwrap();
let rows: Vec<serde_json::Value> = serde_json::Deserializer::from_slice(&bytes)
.into_iter()
.map(Result::unwrap)
.collect();
assert_eq!(rows.len(), 2);
assert!(rows.iter().all(|row| row["context"] == expected));
}
#[test]
fn root_contract_encoded_ordinary_readback_and_failure_use_original_writer() {
let record = serde_json::json!({"event":"request_stage", "context":{"local_request":5}});
let (tx, rx) = mpsc::sync_channel(1);
tx.try_send(Command::RootRecord {
bytes: serde_json::to_vec(&record).unwrap(),
dropped: 0,
})
.unwrap_or_else(|_| panic!("empty queue"));
drop(tx);
let mut bytes = Vec::new();
write_records(rx, BasicSink(&mut bytes), Arc::new(Mutex::new(WriterSourceState::default())));
assert_eq!(
serde_json::from_slice::<serde_json::Value>(&bytes).unwrap(),
record
);
let logger = observer(
FailingWriter {
fail_write: Some(1),
writes: 0,
fail_flush: false,
},
1,
);
assert_eq!(
logger.emit_root_record(&record),
crate::DiagnosticSubmission::Enqueued
);
assert!(matches!(
block_on(logger.shutdown()),
Err(FlushError::OutputFailed(OutputStage::Record))
));
}
#[test]
fn random_source_failure_is_reported_during_initialization() {
let result = Observer::with_writer_and_seed(
ObserverConfig::default(),
io::sink(),
Err(InitError::RandomSource),
);
assert!(matches!(result, Err(InitError::RandomSource)));
}
#[test]
fn initialized_generator_produces_non_zero_unique_ids_without_more_io() {
let observer = observer(io::sink(), 8);
let first_trace = observer.new_trace_id();
let second_trace = observer.new_trace_id();
let first_span = observer.new_span_id();
let second_span = observer.new_span_id();
assert_ne!(first_trace.as_u128(), 0);
assert_ne!(first_trace, second_trace);
assert_ne!(first_span.as_u64(), 0);
assert_ne!(first_span, second_span);
block_on(observer.shutdown()).unwrap();
}
#[test]
fn observer_is_a_managed_component() {
let observer = observer(io::sink(), 8);
assert_eq!(ComponentLifecycle::name(&observer), "observability");
block_on(ComponentLifecycle::start(&observer)).unwrap();
block_on(ComponentLifecycle::shutdown(&observer)).unwrap();
}
#[test]
fn full_queue_is_reported_by_shutdown_without_waiting_in_emit() {
let (entered_sender, entered_receiver) = mpsc::channel();
let release = Arc::new((Mutex::new(false), Condvar::new()));
let writer = BlockingWriter {
entered: Some(entered_sender),
release: release.clone(),
fail_after_release: false,
};
let observer = observer(writer, 1);
observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
entered_receiver.recv().unwrap();
observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
assert_eq!(observer.dropped_events(), 1);
let shutdown = observer.shutdown();
let (lock, condition) = &*release;
*lock.lock().unwrap() = true;
condition.notify_one();
assert_eq!(block_on(shutdown), Err(FlushError::DroppedEvents(1)));
}
#[test]
fn shutdown_reports_writer_failure_and_unreported_drops_together() {
let (entered_sender, entered_receiver) = mpsc::channel();
let release = Arc::new((Mutex::new(false), Condvar::new()));
let writer = BlockingWriter {
entered: Some(entered_sender),
release: release.clone(),
fail_after_release: true,
};
let observer = observer(writer, 1);
observer.emit(LogRecord::new(EventLevel::Info, "test.first"));
entered_receiver.recv().unwrap();
observer.emit(LogRecord::new(EventLevel::Info, "test.queued"));
observer.emit(LogRecord::new(EventLevel::Info, "test.dropped"));
let shutdown = observer.shutdown();
let (lock, condition) = &*release;
*lock.lock().unwrap() = true;
condition.notify_one();
assert_eq!(
block_on(shutdown),
Err(FlushError::OutputFailedAndDropped {
stage: OutputStage::Record,
dropped_events: 1,
})
);
}
#[test]
fn record_write_failure_persists_until_shutdown() {
let observer = observer(
FailingWriter {
fail_write: Some(1),
writes: 0,
fail_flush: false,
},
8,
);
observer.emit(LogRecord::new(EventLevel::Info, "test.record"));
assert_eq!(
block_on(observer.shutdown()),
Err(FlushError::OutputFailed(OutputStage::Record))
);
}
#[test]
fn newline_failure_persists_until_flush() {
let observer = observer(
FailingWriter {
fail_write: Some(2),
writes: 0,
fail_flush: false,
},
8,
);
observer.emit(LogRecord::new(EventLevel::Info, "test.newline"));
assert_eq!(
block_on(observer.flush()),
Err(FlushError::OutputFailed(OutputStage::Newline))
);
let _ = block_on(observer.shutdown());
}
#[test]
fn flush_failure_is_reported() {
let observer = observer(
FailingWriter {
fail_write: None,
writes: 0,
fail_flush: true,
},
8,
);
assert_eq!(
block_on(observer.flush()),
Err(FlushError::OutputFailed(OutputStage::Flush))
);
assert_eq!(
block_on(observer.shutdown()),
Err(FlushError::OutputFailed(OutputStage::Flush))
);
}
}