use super::*;
use saddle_core::{
CaptureSite, Diagnostic, DiagnosticCategory, DiagnosticCause, DiagnosticCode, DiagnosticObject,
DiagnosticObjectKind, DiagnosticStage,
};
use saddle_observability::{
BootstrapDiagnostics, DiagnosticSubmission, EmergencyDiagnostics, EmergencyInitError, FileLoggingConfig, Rotation,
};
static OUTPUT: std::sync::RwLock<Option<saddle_observability::EmergencyDiagnosticHandle>> =
std::sync::RwLock::new(None);
pub(super) type SharedOutput = Arc<Mutex<ProductionOutput>>;
pub(super) struct ProductionOutput {
owner: Option<EmergencyDiagnostics>,
_registration: Option<OutputRegistration>,
exit: Option<saddle_runtime::diagnostics::RuntimeDiagnosticExit>,
}
pub(super) fn load_production<B: DeserializeOwned + 'static>(
path: &Path,
bootstrap: &BootstrapDiagnostics,
) -> Result<ProcessConfig<B>> {
let prepared = PreparedStartup::<B>::load_with_bootstrap(path, bootstrap).map_err(|error| {
eprintln!("{error}; diagnostic_output={}",
if error.source_unavailable() { "unconfirmed" } else { "bootstrap" });
error
})?;
let PreparedStartup {
output,
config,
registration,
..
} = prepared;
let owner = match output {
Ok(owner) => owner,
Err(output_error) => {
return match config {
Err(primary) => {
eprintln!("{primary}; output_failure={output_error}");
Err(primary)
}
Ok(_) => {
eprintln!("{output_error}; diagnostic_source=unavailable");
Err(output_error)
}
};
}
};
let shared = Arc::new(Mutex::new(ProductionOutput {
owner: Some(owner),
_registration: registration,
exit: None,
}));
let mut config = match config {
Ok(config) => config,
Err(error) => return finish(Err(error), Some(&shared)),
};
if let Err(error) = attach_output(&mut config, &shared) {
return finish(Err(error), Some(&shared));
}
Ok(config)
}
pub(super) fn activate_loaded_output<B>(config: &mut ProcessConfig<B>) -> Result<()> {
let bootstrap = config.bootstrap.as_ref().ok_or_else(||
config_failure("saddle.process.bootstrap_output_unavailable", None)
.with_unconfirmed_source())?;
let owner = checked_output(&config.logging, Some(bootstrap))?;
let registration = Some(OutputRegistration::install(&owner));
let shared = Arc::new(Mutex::new(ProductionOutput {
owner: Some(owner),
_registration: registration,
exit: None,
}));
if let Err(error) = attach_output(config, &shared) {
return finish(Err(error), Some(&shared));
}
Ok(())
}
fn attach_output<B>(config: &mut ProcessConfig<B>, shared: &SharedOutput) -> Result<()> {
let handle = shared.lock().unwrap().owner.as_ref().unwrap().handle();
if saddle_runtime::diagnostics::install_output(handle).is_err() {
return Err(report(config_failure(
"saddle.process.diagnostic_output_already_installed", None)));
}
static HOOK: std::sync::Once = std::sync::Once::new();
HOOK.call_once(|| {
std::panic::set_hook(Box::new(|info| {
if !saddle_runtime::diagnostics::capture_current_panic(info) {
let diagnostic = Diagnostic::capture_panic(info, DiagnosticStage::BackgroundTask);
let _submission = record(&diagnostic, None);
}
}))
});
config.diagnostics = Some(shared.clone());
Ok(())
}
fn checked_output(logging: &FileLoggingConfig,
bootstrap: Option<&BootstrapDiagnostics>) -> Result<EmergencyDiagnostics> {
EmergencyDiagnostics::start_checked(logging).map_err(|error| {
let code = match error {
EmergencyInitError::InvalidDirectory => "saddle.process.diagnostic_directory_invalid",
EmergencyInitError::AlreadyActive => "saddle.process.diagnostic_writer_already_active",
EmergencyInitError::Open(_) => "saddle.process.diagnostic_file_open_failed",
EmergencyInitError::Spawn(_) => "saddle.process.diagnostic_worker_spawn_failed",
};
early_failure(code, &error, bootstrap)
})
}
pub(crate) fn open_bootstrap() -> Result<BootstrapDiagnostics> {
let directory = std::env::var_os("SADDLE_BOOTSTRAP_LOG_DIRECTORY")
.map(PathBuf::from).unwrap_or_else(|| PathBuf::from("./logs"));
BootstrapDiagnostics::start(&directory).map_err(|_| {
SaddleError::new(ErrorKind::Unavailable,
"saddle.process.bootstrap_output_unavailable",
"controlled startup diagnostic output unavailable")
.with_unconfirmed_source()
})
}
#[track_caller]
pub(super) fn bootstrap_source(
code: &'static str,
raw: &(dyn std::error::Error + 'static),
bootstrap: &BootstrapDiagnostics,
) -> SaddleError {
let error = config_failure(code, raw.downcast_ref::<std::io::Error>());
let receipt = saddle_observability::root_diagnostic::process_bootstrap_error(
bootstrap, error.diagnostic().expect("bootstrap source diagnostic"), raw,
);
if receipt.original_capture()
== saddle_observability::root_diagnostic::OriginalCaptureState::CompleteWritten {
error
} else {
error.with_unconfirmed_source()
}
}
#[track_caller]
pub(super) fn bootstrap_process_error(
error: SaddleError,
bootstrap: &BootstrapDiagnostics,
) -> SaddleError {
if error.diagnostic().is_some() {
return error.with_unconfirmed_source();
}
let cause = DiagnosticCause::new(DiagnosticStage::StartupConfig,
DiagnosticCode::new(error.code()).unwrap_or_else(||
DiagnosticCode::new("saddle.process.unclassified_failure").unwrap()));
let diagnostic = Diagnostic::capture(
DiagnosticCategory::UnexpectedError, CaptureSite::FirstObserved, cause);
let receipt = saddle_observability::root_diagnostic::process_bootstrap_error(
bootstrap, &diagnostic, &error,
);
let error = error.with_diagnostic(diagnostic);
if receipt.original_capture()
== saddle_observability::root_diagnostic::OriginalCaptureState::CompleteWritten {
error
} else {
error.with_unconfirmed_source()
}
}
#[track_caller]
pub(super) fn component_source_error(
raw: &(dyn std::error::Error + 'static),
code: &'static str,
application: &str,
output: Option<&saddle_observability::EmergencyDiagnosticHandle>,
cleanup: bool,
primary: Option<saddle_core::DiagnosticOccurrence>,
) -> SaddleError {
let mut diagnostic = Diagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
DiagnosticCause::new(
if cleanup { DiagnosticStage::ShutdownComponent }
else { DiagnosticStage::StartupListener },
DiagnosticCode::new(code).expect("static component error code"),
),
);
if let Some(primary) = primary.as_ref() {
diagnostic = diagnostic.during_cleanup_of_occurrence(primary);
}
let identity = saddle_core::ContextLabel::checked(application)
.map_or(saddle_core::ContextFact::Unavailable, saddle_core::ContextFact::Present);
let receipt = if cleanup {
saddle_observability::root_diagnostic::process_component_cleanup_error(
output, identity, &diagnostic, raw)
} else {
saddle_observability::root_diagnostic::process_component_start_error(
output, identity, &diagnostic, raw)
};
let complete = receipt.original_capture()
== saddle_observability::root_diagnostic::OriginalCaptureState::CompleteWritten
&& receipt.occurrence().matches_diagnostic(&diagnostic);
let error = SaddleError::new(ErrorKind::Infrastructure, code, "Saddle process startup failed")
.with_diagnostic(diagnostic);
if complete { error.with_source_receipt(receipt) }
else if cleanup { error.with_unconfirmed_source().with_unconfirmed_cleanup_source() }
else { error.with_unconfirmed_source() }
}
pub(super) fn take_owner(output: &SharedOutput) -> Option<EmergencyDiagnostics> {
output
.lock()
.unwrap_or_else(|e| e.into_inner())
.owner
.take()
}
pub(super) fn process_handle(
output: &SharedOutput,
) -> Option<saddle_observability::EmergencyDiagnosticHandle> {
output
.lock()
.unwrap_or_else(|e| e.into_inner())
.owner
.as_ref()
.map(EmergencyDiagnostics::handle)
}
pub(super) fn process_source_output(
output: &SharedOutput,
) -> Option<saddle_observability::SourceOutput> {
output
.lock()
.unwrap_or_else(|e| e.into_inner())
.owner
.as_ref()
.and_then(EmergencyDiagnostics::source_output)
}
pub(super) fn parse_error(
path: &Path,
source: &str,
error: &toml::de::Error,
output: Option<&saddle_observability::EmergencyDiagnosticHandle>,
bootstrap: Option<&BootstrapDiagnostics>,
) -> SaddleError {
let (line, column) = error
.span()
.map(|span| {
let prefix = &source.as_bytes()[..span.start.min(source.len())];
let line = prefix.iter().filter(|&&b| b == b'\n').count() as u64 + 1;
let column = prefix
.iter()
.rposition(|&b| b == b'\n')
.map_or(prefix.len(), |last| prefix.len() - last - 1)
as u64
+ 1;
(Some(line), Some(column))
})
.unwrap_or((None, None));
let file = path
.file_name()
.and_then(|p| p.to_str())
.and_then(|name| saddle_core::DiagnosticLocator::from_projection(name, false, false))
.unwrap_or_else(|| {
saddle_core::DiagnosticLocator::from_projection("", false, true)
.expect("redacted locator")
});
let location = saddle_core::DiagnosticInputLocation::new(line, column).with_file(file);
let cause = DiagnosticCause::new(
DiagnosticStage::StartupConfig,
DiagnosticCode::new("saddle.process.config_invalid").unwrap(),
)
.with_input_location(location);
let diagnostic = Diagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
cause,
);
let original_capture = if let Some(output) = output {
saddle_observability::root_diagnostic::process_config_error(
Some(output), &diagnostic, error).original_capture()
} else if let Some(bootstrap) = bootstrap {
saddle_observability::root_diagnostic::process_bootstrap_error(
bootstrap, &diagnostic, error).original_capture()
} else {
saddle_observability::root_diagnostic::OriginalCaptureState::OutputUnavailable
};
if (output.is_some() || bootstrap.is_some()) && original_capture
!= saddle_observability::root_diagnostic::OriginalCaptureState::CompleteWritten {
return SaddleError::new(ErrorKind::Unavailable,
"saddle.diagnostic.source_unavailable", "configuration source diagnostic unavailable")
.with_diagnostic(diagnostic).with_unconfirmed_source();
}
let error = SaddleError::new(
ErrorKind::InvalidArgument,
"saddle.process.config_invalid",
"configuration parse failed",
)
.with_diagnostic(diagnostic);
if output.is_none() && bootstrap.is_none() {
error.with_unconfirmed_source()
} else {
error
}
}
pub(super) fn set_exit(
output: &SharedOutput,
exit: saddle_runtime::diagnostics::RuntimeDiagnosticExit,
) {
output.lock().unwrap_or_else(|e| e.into_inner()).exit = Some(exit);
}
#[track_caller]
pub(super) fn report(error: SaddleError) -> SaddleError {
let error = if error.diagnostic().is_some() {
error
} else {
let cause = DiagnosticCause::new(
DiagnosticStage::StartupConfig,
DiagnosticCode::new(error.code()).unwrap_or_else(|| {
DiagnosticCode::new("saddle.process.unclassified_failure").unwrap()
}),
);
error.with_diagnostic(Diagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
cause,
))
};
if let Some(diagnostic) = error.diagnostic() {
let _submission = record(diagnostic, None);
}
error
}
pub(super) fn finish<T>(result: Result<T>, output: Option<&SharedOutput>) -> Result<T> {
if let Some(output) = output {
let mut output = output.lock().unwrap_or_else(|e| e.into_inner());
if let Some(mut owner) = output.owner.take() {
let shutdown = owner.shutdown();
output.exit = Some(saddle_runtime::diagnostics::RuntimeDiagnosticExit {
shutdown,
snapshot: owner.snapshot(),
});
}
if let Some(exit) = &output.exit {
let state = exit.snapshot;
if !state.initialized
|| state.first_failure.is_some()
|| state.written != state.enqueued
|| state.dropped != 0
|| exit.shutdown != saddle_observability::DiagnosticShutdown::Finished
{
eprintln!(
"saddle.diagnostic.output_not_confirmed shutdown={:?} initialized={} enqueued={} written={} dropped={} first_failure={:?}",
exit.shutdown,
state.initialized,
state.enqueued,
state.written,
state.dropped,
state.first_failure
);
if let Err(error) = &result {
eprintln!("{error}");
}
}
}
}
result
}
pub(crate) struct OutputRegistration;
impl OutputRegistration {
pub(crate) fn install(output: &EmergencyDiagnostics) -> Self {
*OUTPUT.write().unwrap_or_else(|e| e.into_inner()) = Some(output.handle());
Self
}
}
impl Drop for OutputRegistration {
fn drop(&mut self) {
*OUTPUT.write().unwrap_or_else(|e| e.into_inner()) = None;
}
}
pub(crate) fn record(
diagnostic: &Diagnostic,
context: Option<(
&saddle_core::CallContext,
&saddle_observability::EventContext,
)>,
) -> Option<DiagnosticSubmission> {
let handle = OUTPUT.read().unwrap_or_else(|e| e.into_inner()).clone()?;
Some(match saddle_observability::global() {
Some(observer) => observer.record_diagnostic(diagnostic, &handle, context),
None => handle.submit_context(diagnostic, context),
})
}
#[derive(Deserialize)]
struct EarlyFile {
framework: EarlyFramework,
}
#[derive(Deserialize)]
struct EarlyFramework {
#[serde(default)]
observability: ObservabilityFileConfig,
}
pub(super) struct PreparedStartup<B> {
pub output: Result<EmergencyDiagnostics>,
pub config: Result<ProcessConfig<B>>,
pub submission: Option<DiagnosticSubmission>,
pub registration: Option<OutputRegistration>,
}
impl<B: DeserializeOwned + 'static> PreparedStartup<B> {
#[cfg(test)]
pub fn load(path: &Path) -> Result<Self> {
Self::load_inner(path, None)
}
pub fn load_with_bootstrap(path: &Path, bootstrap: &BootstrapDiagnostics) -> Result<Self> {
Self::load_inner(path, Some(bootstrap))
}
fn load_inner(path: &Path, bootstrap: Option<&BootstrapDiagnostics>) -> Result<Self> {
let bytes = std::fs::read(path)
.map_err(|error| early_failure("saddle.process.config_unavailable", &error, bootstrap))?;
let source = std::str::from_utf8(&bytes)
.map_err(|error| early_failure("saddle.process.config_invalid_utf8", &error, bootstrap))?;
let early: EarlyFile = toml::from_str(source)
.map_err(|error| early_failure("saddle.process.logging_configuration_unresolved", &error, bootstrap))?;
let logging = early.framework.observability.logging;
let logging = FileLoggingConfig::new(
logging.directory.unwrap_or_else(|| PathBuf::from("./logs")),
match logging.rotation {
LoggingRotation::Daily => Rotation::Daily,
LoggingRotation::Hourly => Rotation::Hourly,
},
);
let output = checked_output(&logging, bootstrap);
let registration = output.as_ref().ok().map(OutputRegistration::install);
let output_handle = output.as_ref().ok().map(EmergencyDiagnostics::handle);
let config = ProcessConfig::load_source(path, source, output_handle.as_ref(), bootstrap).map_err(|error| {
if error.diagnostic().is_some() {
error
} else {
let cause = DiagnosticCause::new(
DiagnosticStage::StartupConfig,
DiagnosticCode::new(error.code()).expect("framework static code"),
);
error.with_diagnostic(Diagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
cause,
))
}
});
let submission = config
.as_ref()
.err()
.and_then(SaddleError::diagnostic)
.and_then(|d| output.as_ref().ok().map(|output| output.handle().submit(d)));
Ok(Self {
output,
config,
submission,
registration,
})
}
}
#[track_caller]
fn early_failure(
code: &'static str,
raw: &(dyn std::error::Error + 'static),
bootstrap: Option<&BootstrapDiagnostics>,
) -> SaddleError {
match bootstrap {
Some(bootstrap) => bootstrap_source(code, raw, bootstrap),
None => config_failure(code, raw.downcast_ref::<std::io::Error>()).with_unconfirmed_source(),
}
}
#[track_caller]
fn config_failure(code: &'static str, io: Option<&std::io::Error>) -> SaddleError {
let mut cause = DiagnosticCause::new(
DiagnosticStage::StartupConfig,
DiagnosticCode::new(code).expect("static code"),
);
if !matches!(code, "saddle.process.config_unavailable" | "saddle.process.config_invalid" | "saddle.process.config_invalid_utf8") {
cause = cause.with_object(
DiagnosticObject::new(
DiagnosticObjectKind::ConfigKey,
"framework.observability.logging",
)
.expect("static key"),
);
}
if let Some(io) = io {
cause = cause.with_io(io);
}
SaddleError::new(
ErrorKind::Infrastructure,
code,
"startup configuration failed",
)
.with_diagnostic(Diagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
cause,
))
}
#[cfg(test)]
mod tests {
use super::*;
use saddle_observability::DiagnosticShutdown;
use std::error::Error as _;
use std::process::Command;
#[test]
fn bootstrap_environment_directory_is_explicit_in_subprocess() {
let root = std::env::temp_dir().join(format!("saddle-bootstrap-env-{}", std::process::id()));
std::fs::create_dir(&root).unwrap();
for mode in ["invalid", "default", "explicit", "args"] {
let mut child = Command::new(std::env::current_exe().unwrap());
child.args(["--exact", "process::diagnostics::tests::bootstrap_env_child"])
.current_dir(&root)
.env("SADDLE_BOOTSTRAP_CHILD", mode);
if mode == "default" { child.env_remove("SADDLE_BOOTSTRAP_LOG_DIRECTORY"); }
else if mode == "explicit" { child.env("SADDLE_BOOTSTRAP_LOG_DIRECTORY", root.join("selected")); }
else if mode == "args" { child.env("SADDLE_BOOTSTRAP_LOG_DIRECTORY", root.join("args")); }
else { child.env("SADDLE_BOOTSTRAP_LOG_DIRECTORY", root.join("invalid@selected")); }
assert!(child.status().unwrap().success(), "{mode}");
if mode == "invalid" { assert!(!root.join("logs").exists()); }
}
assert!(root.join("logs/saddle.bootstrap.log").is_file());
assert!(root.join("selected/saddle.bootstrap.log").is_file());
assert!(root.join("args/saddle.bootstrap.log").is_file());
assert!(!root.join("invalid@selected/saddle.bootstrap.log").exists());
for directory in ["logs", "selected", "args"] {
let path = root.join(directory).join("saddle.bootstrap.log");
let rows: Vec<serde_json::Value> = std::fs::read_to_string(&path).unwrap()
.lines().map(|line| serde_json::from_str(line).unwrap()).collect();
assert!(rows.iter().any(|row| row["channel"] == "terminal"
&& row["state"] == "exposed_chain_complete"));
if directory == "args" {
let id = rows.iter().find(|row| row["channel"] == "context")
.unwrap()["occurrence"]["diagnostic_id"].as_u64().unwrap();
let original = process_error("saddle.process.arguments_invalid");
assert_eq!(original_field(&rows, id, "description"), format!("{original}"));
assert_eq!(original_field(&rows, id, "debug"), format!("{original:?}"));
}
std::fs::remove_file(path).unwrap();
std::fs::remove_dir(root.join(directory)).unwrap();
}
std::fs::remove_dir(root).unwrap();
}
#[test]
fn bootstrap_env_child() {
let Ok(mode) = std::env::var("SADDLE_BOOTSTRAP_CHILD") else { return; };
let missing = std::env::current_dir().unwrap().join("missing.toml");
let result = if mode == "args" {
ProcessConfig::<()>::__from_args_with_diagnostics()
} else {
ProcessConfig::<()>::load(&missing)
};
if mode == "invalid" {
let error = result.err().unwrap();
assert_eq!(error.code(), "saddle.process.bootstrap_output_unavailable");
assert!(error.source_unavailable());
} else if mode == "args" {
let error = result.err().unwrap();
assert_eq!(error.code(), "saddle.process.arguments_invalid");
assert!(!error.source_unavailable());
} else {
let error = result.err().unwrap();
assert_eq!(error.code(), "saddle.process.config_unavailable");
assert!(!error.source_unavailable());
}
}
fn original_field(rows: &[serde_json::Value], id: u64, channel: &str) -> String {
let segments: Vec<_> = rows.iter().filter(|row|
row["event"] == "request_error_original"
&& row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == channel && row["cause_depth"] == 0
).collect();
assert!(!segments.is_empty(), "missing {channel} for {id}");
assert_eq!(segments.last().unwrap()["state"], "field_end");
let first = segments[0]["sequence"].as_u64().unwrap();
for (index, segment) in segments.iter().enumerate() {
assert_eq!(segment["sequence"].as_u64(), Some(first + index as u64));
}
segments.iter().map(|row| row["payload"].as_str().unwrap()).collect()
}
#[test]
fn owned_bootstrap_source_precedes_safe_projection() {
let root = std::env::temp_dir().join(format!("saddle-bootstrap-owned-{}", std::process::id()));
std::fs::create_dir(&root).unwrap();
let bootstrap = BootstrapDiagnostics::start(&root).unwrap();
let raw = SaddleError::new(ErrorKind::Infrastructure,
"saddle.process.owned_error_probe", "ORIGINAL_OWNED_MESSAGE_93815");
let original_display = format!("{raw}");
let original_debug = format!("{raw:?}");
let returned = bootstrap_process_error(raw, &bootstrap);
assert!(!returned.source_unavailable());
assert!(!format!("{returned}").contains("ORIGINAL_OWNED_MESSAGE_93815"));
assert!(!format!("{returned:?}").contains("ORIGINAL_OWNED_MESSAGE_93815"));
let id = returned.diagnostic().unwrap().id();
let rows: Vec<serde_json::Value> = std::fs::read_to_string(bootstrap.target()).unwrap()
.lines().map(|line| serde_json::from_str(line).unwrap()).collect();
assert_eq!(original_field(&rows, id, "description"), original_display);
assert_eq!(original_field(&rows, id, "debug"), original_debug);
assert!(rows.iter().any(|row| row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == "terminal" && row["state"] == "exposed_chain_complete"));
let path = root.join("invalid-token.toml");
std::fs::write(&path, "[framework]\nlisten='127.0.0.1:0'\n[framework.management]\nbind='127.0.0.1:0'\n[framework.admission]\ncpuCores=2\nmemoryMb=512\n[framework.admission.dependencies]\ndatabaseConcurrency=1\nprofusecontractConcurrency=1\n[framework.profusecontract]\nauthority='http://localhost:50051'\ntoken='invalid:token'\n[secrets]\n").unwrap();
let expected = super::super::process_error("saddle.process.authority_template_invalid");
let returned = ProcessConfig::<()>::load_with_bootstrap(&path, &bootstrap).err().unwrap();
assert_eq!(returned.code(), expected.code());
assert!(!returned.source_unavailable());
let semantic_id = returned.diagnostic().unwrap().id();
let rows: Vec<serde_json::Value> = std::fs::read_to_string(bootstrap.target()).unwrap()
.lines().map(|line| serde_json::from_str(line).unwrap()).collect();
assert_eq!(original_field(&rows, semantic_id, "description"), format!("{expected}"));
assert_eq!(original_field(&rows, semantic_id, "debug"), format!("{expected:?}"));
let projected = SaddleError::new(ErrorKind::Infrastructure,
"saddle.process.already_projected_probe", "MUST_NOT_CLAIM_ORIGINAL_57216")
.with_diagnostic(Diagnostic::capture(DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved, DiagnosticCause::new(DiagnosticStage::StartupConfig,
DiagnosticCode::new("saddle.process.already_projected_probe").unwrap())));
let rejected = bootstrap_process_error(projected, &bootstrap);
assert!(rejected.source_unavailable());
assert!(!std::fs::read_to_string(bootstrap.target()).unwrap()
.contains("MUST_NOT_CLAIM_ORIGINAL_57216"));
drop(bootstrap);
std::fs::remove_file(root.join("saddle.bootstrap.log")).unwrap();
std::fs::remove_file(path).unwrap();
std::fs::remove_dir(root).unwrap();
}
#[tokio::test]
async fn component_sources_keep_actual_admission_and_join_errors() {
let root = std::env::temp_dir().join(format!(
"saddle-component-originals-{}", std::process::id()));
std::fs::create_dir(&root).unwrap();
let mut output = EmergencyDiagnostics::start(
&FileLoggingConfig::new(root.clone(), Rotation::Daily)).unwrap();
let handle = output.handle();
let admission = saddle_admission::AdmissionError::SizeOverflow;
let expected_admission = (format!("{admission}"), format!("{admission:?}"));
let started = component_source_error(&admission,
"saddle.process.request_storage_unavailable", "component-test",
Some(&handle), false, None);
assert!(!started.source_unavailable());
assert!(started.source_receipt::<saddle_observability::root_diagnostic::ComponentSourceReceipt>()
.is_some());
let join = tokio::spawn(async { std::future::pending::<()>().await });
tokio::task::yield_now().await;
join.abort();
let join = join.await.unwrap_err();
let expected_join = (format!("{join}"), format!("{join:?}"));
let stopped = component_source_error(&join, "saddle.process.entry_failed",
"component-test", Some(&handle), true, None);
assert!(!stopped.cleanup_source_unavailable());
assert!(stopped.source_receipt::<saddle_observability::root_diagnostic::ComponentSourceReceipt>()
.is_some());
let rows: Vec<serde_json::Value> = std::fs::read_to_string(output.target()).unwrap()
.lines().map(|line| serde_json::from_str(line).unwrap()).collect();
for (error, expected) in [(&started, expected_admission), (&stopped, expected_join)] {
let id = error.diagnostic().unwrap().id();
assert_eq!(original_field(&rows, id, "description"), expected.0);
assert_eq!(original_field(&rows, id, "debug"), expected.1);
assert!(rows.iter().any(|row| row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == "terminal"
&& row["state"] == "exposed_chain_complete"));
}
let missing = component_source_error(&admission,
"saddle.process.request_storage_unavailable", "component-test", None, false, None);
assert!(missing.source_unavailable());
assert!(missing.source_receipt::<saddle_observability::root_diagnostic::ComponentSourceReceipt>()
.is_none());
let unavailable_join = component_source_error(&join,
"saddle.process.entry_failed", "component-test", None, true, None);
assert!(unavailable_join.cleanup_source_unavailable());
assert!(unavailable_join.source_receipt::<saddle_observability::root_diagnostic::ComponentSourceReceipt>()
.is_none());
drop(handle);
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
while output.shutdown() == DiagnosticShutdown::Pending
&& std::time::Instant::now() < deadline {
std::thread::yield_now();
}
let target = output.target().to_owned();
drop(output);
std::fs::remove_file(target).unwrap();
std::fs::remove_dir(root).unwrap();
}
#[test]
fn bootstrap_originals_precede_safe_configuration_returns() {
let root = std::env::temp_dir().join(format!("saddle-bootstrap-originals-{}", std::process::id()));
std::fs::create_dir(&root).unwrap();
let bootstrap = BootstrapDiagnostics::start(&root.join("bootstrap")).unwrap();
let path = root.join("saddle.toml");
let mut observed = Vec::new();
let missing = std::fs::read(&path).unwrap_err();
let returned = ProcessConfig::<()>::load_with_bootstrap(&path, &bootstrap).err().unwrap();
observed.push((returned.diagnostic().unwrap().id(), format!("{missing}"), format!("{missing:?}")));
assert!(!returned.source_unavailable());
let invalid_utf8 = vec![0xff, b'X'];
std::fs::write(&path, &invalid_utf8).unwrap();
let utf8 = std::str::from_utf8(&invalid_utf8).unwrap_err();
let returned = ProcessConfig::<()>::load_with_bootstrap(&path, &bootstrap).err().unwrap();
observed.push((returned.diagnostic().unwrap().id(), format!("{utf8}"), format!("{utf8:?}")));
assert!(!returned.source_unavailable());
let invalid_toml = "[framework\n";
std::fs::write(&path, invalid_toml).unwrap();
let parser = toml::from_str::<EarlyFile>(invalid_toml).err().unwrap();
let returned = PreparedStartup::<()>::load_with_bootstrap(&path, &bootstrap).err().unwrap();
observed.push((returned.diagnostic().unwrap().id(), format!("{parser}"), format!("{parser:?}")));
assert!(!returned.source_unavailable());
let normal_dir = root.join("invalid@normal");
let source = format!("[framework]\n[framework.observability.logging]\ndirectory={}\n",
serde_json::to_string(&normal_dir).unwrap());
std::fs::write(&path, &source).unwrap();
let prepared = PreparedStartup::<()>::load_with_bootstrap(&path, &bootstrap).unwrap();
let output_error = prepared.output.err().unwrap();
observed.push((output_error.diagnostic().unwrap().id(),
"emergency diagnostic initialization failed: InvalidDirectory".to_owned(),
"InvalidDirectory".to_owned()));
assert!(!output_error.source_unavailable());
let blocked_directory = root.join("file-instead-of-directory");
std::fs::write(&blocked_directory, "occupied").unwrap();
let independent_io = std::fs::create_dir_all(&blocked_directory).unwrap_err();
let raw = EmergencyInitError::Open(independent_io);
let source = format!("[framework]\n[framework.observability.logging]\ndirectory={}\n",
serde_json::to_string(&blocked_directory).unwrap());
std::fs::write(&path, source).unwrap();
let prepared = PreparedStartup::<()>::load_with_bootstrap(&path, &bootstrap).unwrap();
let output_error = prepared.output.err().unwrap();
assert_eq!(output_error.code(), "saddle.process.diagnostic_file_open_failed");
assert!(!output_error.source_unavailable());
let blocked_id = output_error.diagnostic().unwrap().id();
observed.push((blocked_id, format!("{raw}"), format!("{raw:?}")));
let rows: Vec<serde_json::Value> = std::fs::read_to_string(bootstrap.target()).unwrap()
.lines().map(|line| serde_json::from_str(line).unwrap()).collect();
use std::os::unix::fs::PermissionsExt as _;
assert_eq!(std::fs::metadata(bootstrap.target()).unwrap().permissions().mode() & 0o777, 0o644);
for (id, display, debug) in observed {
assert_eq!(original_field(&rows, id, "description"), display);
assert_eq!(original_field(&rows, id, "debug"), debug);
assert!(rows.iter().any(|row| row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == "terminal"
&& row["state"] == "exposed_chain_complete"));
}
let source_description: String = rows.iter().filter(|row|
row["occurrence"]["diagnostic_id"] == blocked_id
&& row["channel"] == "description" && row["cause_depth"] == 1)
.map(|row| row["payload"].as_str().unwrap()).collect();
assert_eq!(source_description, format!("{}", raw.source().unwrap()));
let normal_directory = root.join("normal");
let source = format!(
"[framework]\nlisten='127.0.0.1:0'\n[framework.management]\nbind='127.0.0.1:0'\n[framework.admission]\ncpuCores=2\nmemoryMb=512\n[framework.admission.dependencies]\ndatabaseConcurrency=1\nprofusecontractConcurrency=1\n[framework.profusecontract]\nauthority='http://localhost:50051'\ntoken='SENSITIVE_BOOTSTRAP_HANDOFF'\n[framework.observability.logging]\ndirectory={}\n[secrets]\n[business]\ninvalid='SENSITIVE_BUSINESS_HANDOFF'\n",
serde_json::to_string(&normal_directory).unwrap(),
);
std::fs::write(&path, &source).unwrap();
let raw = toml::from_str::<FileConfig<()>>(&source).err().unwrap();
let prepared = PreparedStartup::<()>::load_with_bootstrap(&path, &bootstrap).unwrap();
let error = prepared.config.err().unwrap();
assert!(!error.source_unavailable());
assert!(!format!("{error}").contains("SENSITIVE_BUSINESS_HANDOFF"));
let id = error.diagnostic().unwrap().id();
let mut normal = prepared.output.unwrap();
let normal_rows: Vec<serde_json::Value> = std::fs::read_to_string(normal.target()).unwrap()
.lines().map(|line| serde_json::from_str(line).unwrap()).collect();
let normal_description = original_field(&normal_rows, id, "description");
assert_eq!(normal_description, format!("{raw}"));
assert_eq!(original_field(&normal_rows, id, "debug"), format!("{raw:?}"));
assert!(std::fs::read_to_string(normal.target()).unwrap()
.contains("SENSITIVE_BUSINESS_HANDOFF"));
assert!(!std::fs::read_to_string(bootstrap.target()).unwrap()
.contains("SENSITIVE_BUSINESS_HANDOFF"));
let deadline = std::time::Instant::now() + Duration::from_secs(2);
while normal.shutdown() == DiagnosticShutdown::Pending
&& std::time::Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(1));
}
assert_eq!(normal.shutdown(), DiagnosticShutdown::Finished);
std::fs::remove_file(normal.target()).unwrap();
std::fs::remove_dir(normal_directory).unwrap();
drop(bootstrap);
std::fs::remove_file(root.join("bootstrap").join("saddle.bootstrap.log")).unwrap();
std::fs::remove_dir(root.join("bootstrap")).unwrap();
std::fs::remove_file(path).unwrap();
std::fs::remove_file(blocked_directory).unwrap();
std::fs::remove_dir(root).unwrap();
}
#[test]
fn pre_output_failures_keep_unconfirmed_source_fact() {
let root = std::env::temp_dir().join(format!("saddle-pre-output-{}", std::process::id()));
std::fs::create_dir(&root).unwrap();
let path = root.join("saddle.toml");
let missing = PreparedStartup::<()>::load(&path).err().unwrap();
assert_eq!(missing.code(), "saddle.process.config_unavailable");
assert!(missing.source_unavailable());
std::fs::write(&path, "[framework\n").unwrap();
let malformed = PreparedStartup::<()>::load(&path).err().unwrap();
assert_eq!(malformed.code(), "saddle.process.logging_configuration_unresolved");
assert!(malformed.source_unavailable());
let direct = ProcessConfig::<()>::load(&path).err().unwrap();
assert_eq!(direct.code(), "saddle.process.config_invalid");
assert!(!direct.source_unavailable());
std::fs::remove_file(path).unwrap();
std::fs::remove_dir(root).unwrap();
}
#[test]
fn early_config_failure_uses_selected_directory_and_preserves_primary() {
let root = std::env::temp_dir().join(format!("saddle-early-config-{}", std::process::id()));
std::fs::create_dir(&root).unwrap();
let path = root.join("saddle.toml");
let base = "[framework]\nlisten='127.0.0.1:0'\n[framework.management]\nbind='127.0.0.1:0'\n[framework.admission]\ncpuCores=2\nmemoryMb=512\n[framework.admission.dependencies]\ndatabaseConcurrency=1\nprofusecontractConcurrency=1\n[framework.profusecontract]\nauthority='http://localhost:50051'\ntoken='DO_NOT_LOG_TOKEN_68142'\n[secrets]\n";
for (name, bad_config, bad_directory) in [
("good logs", false, false),
("invalid config logs", true, false),
("file-not-directory", true, true),
] {
let directory = root.join(name);
if bad_directory {
std::fs::write(&directory, "not a directory").unwrap();
}
let source = format!(
"{base}[framework.observability.logging]\ndirectory={}\n{}",
serde_json::to_string(&directory).unwrap(),
if bad_config {
"[business]\ninvalid='DO_NOT_LOG_DATA_93715'\n"
} else {
""
}
);
std::fs::write(&path, source).unwrap();
let mut prepared = PreparedStartup::<()>::load(&path).unwrap();
assert_eq!(prepared.config.is_err(), bad_config);
if bad_config {
assert_eq!(
prepared.config.as_ref().err().unwrap().code(),
if bad_directory { "saddle.diagnostic.source_unavailable" }
else { "saddle.process.config_invalid" }
);
}
let output = prepared.output.as_mut().unwrap();
let deadline = std::time::Instant::now() + Duration::from_secs(5);
while !output.snapshot().initialized && std::time::Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(1));
}
assert!(output.snapshot().initialized);
assert_eq!(output.snapshot().first_failure.is_some(), bad_directory);
while output.shutdown() == DiagnosticShutdown::Pending
&& std::time::Instant::now() < deadline
{
std::thread::sleep(Duration::from_millis(1));
}
assert_eq!(output.shutdown(), DiagnosticShutdown::Finished);
if bad_config && !bad_directory {
assert_eq!(prepared.submission, Some(DiagnosticSubmission::Enqueued));
assert!(output.snapshot().written >= 5);
let record = std::fs::read_to_string(output.target()).unwrap();
assert!(record.contains("saddle.process.config_invalid"));
assert!(record.contains("exposed_chain_complete"));
assert!(record.contains("DO_NOT_LOG_TOKEN_68142"));
assert!(record.contains("DO_NOT_LOG_DATA_93715"));
let original: serde_json::Value = serde_json::from_str(record.lines().next().unwrap()).unwrap();
assert_eq!(original["occurrence"]["diagnostic_id"],
prepared.config.as_ref().err().unwrap().diagnostic().unwrap().id());
let public = prepared.config.as_ref().err().unwrap().to_string();
assert!(!public.contains("DO_NOT_LOG_TOKEN_68142"));
assert!(!public.contains("DO_NOT_LOG_DATA_93715"));
}
if !bad_directory {
std::fs::remove_file(output.target()).unwrap();
std::fs::remove_dir(directory).unwrap();
} else {
assert!(prepared.config.is_err());
std::fs::remove_file(directory).unwrap();
}
}
std::fs::remove_file(path).unwrap();
std::fs::remove_dir(root).unwrap();
}
}