use super::*;
use crate::request_task::reserved::Retained;
use saddle_observability::root_diagnostic::RecordedTerminalFailure;
struct StopState {
primary: Mutex<Option<(RecordedTerminalFailure, Option<saddle_core::request_context::RequestContextSnapshot>)>>,
wake: tokio::sync::Notify,
task_storage_issued: std::sync::atomic::AtomicBool,
health: Mutex<Option<crate::ApplicationHealth>>,
_storage: saddle_admission::StoragePermit,
}
#[derive(Clone)]
#[doc(hidden)]
pub struct ProfuseGwFailureStop(Retained<StopState>);
impl ProfuseGwFailureStop {
pub(super) fn task_storage(&self) -> Result<crate::request_task::reserved_set::ReservedProcessTaskStorage, AdmissionError> {
if self.0.task_storage_issued.compare_exchange(false, true,
std::sync::atomic::Ordering::AcqRel, std::sync::atomic::Ordering::Acquire).is_err() {
return Err(AdmissionError::CapacityRejected);
}
match crate::request_task::reserved_set::ReservedProcessTaskStorage::from_process_parent(&self.0._storage) {
Ok(storage) => Ok(storage),
Err(error) => {
self.0.task_storage_issued.store(false, std::sync::atomic::Ordering::Release);
Err(error)
}
}
}
pub(super) fn prepare(process: &ProfuseGwRuntimeProcess) -> Result<Self, AdmissionError> {
let storage = process.admission.try_process_storage(saddle_admission::StorageDemand::separate(&[
(std::alloc::Layout::new::<StopState>(), 1),
(std::alloc::Layout::new::<[std::sync::atomic::AtomicUsize; 2]>(), 1),
])?)?;
Ok(Self(Retained::new(StopState {
task_storage_issued: std::sync::atomic::AtomicBool::new(false),
primary: Mutex::new(None), wake: tokio::sync::Notify::new(),
health: Mutex::new(None), _storage: storage,
})))
}
pub(super) fn bind_health(&self, health: crate::ApplicationHealth) {
*self.0.health.lock().unwrap_or_else(|p| p.into_inner()) = Some(health.clone());
if self.is_requested() { health.fail_closed(); }
}
pub fn is_requested(&self) -> bool {
self.0.primary.lock().unwrap_or_else(|p| p.into_inner()).is_some()
}
pub fn report(
&self,
failure: saddle_admission::ProfuseGwTerminalAuditFailure,
context: Option<saddle_core::request_context::RequestContextSnapshot>,
source: Option<saddle_core::DiagnosticOccurrence>,
output: Option<&saddle_observability::EmergencyDiagnosticHandle>,
) {
let recorded = RecordedTerminalFailure::audit(output, context.as_ref(), source, failure);
let first = {
let mut primary = self.0.primary.lock().unwrap_or_else(|p| p.into_inner());
if primary.is_some() { false } else { *primary = Some((recorded, context)); true }
};
if first {
if let Some(health) = self.0.health.lock().unwrap_or_else(|p| p.into_inner()).as_ref() {
health.fail_closed();
}
self.0.wake.notify_one();
}
}
pub(super) async fn wait(&self) -> saddle_core::Result<()> {
loop {
let notified = self.0.wake.notified();
let primary = self.0.primary.lock().unwrap_or_else(|p| p.into_inner());
if let Some((recorded, _)) = primary.as_ref() {
return Err(recorded.project());
}
drop(primary);
notified.await;
}
}
pub fn record_cleanup(&self, snapshot: &saddle_admission::ProcessSnapshot,
output: Option<&saddle_observability::EmergencyDiagnosticHandle>) {
let primary = self.0.primary.lock().unwrap_or_else(|p| p.into_inner());
if let Some((recorded, context)) = primary.as_ref() {
let _ = saddle_observability::root_diagnostic::request_terminal_audit_cleanup(
output, context.as_ref(), recorded.diagnostic().occurrence(), snapshot);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn failure() -> saddle_admission::ProfuseGwTerminalAuditFailure {
saddle_admission::ProfuseGwTerminalAuditFailure {
error: AdmissionError::AccountClosed,
audit: None,
database: saddle_admission::ProfuseGwDatabaseDisposition::NotUsed,
}
}
#[test]
fn audit_stop_survives_missing_output_late_health_and_repeated_reports() {
let mut process = ProfuseGwRuntimeProcess::new(crate::request_task::reserved::tests::process());
drop(process.db_startup.take());
let baseline = process.admission.resource_snapshot().framework_charged;
let stop = ProfuseGwFailureStop::prepare(&process).unwrap();
stop.report(failure(), None, None, None);
let first = stop.0.primary.lock().unwrap().as_ref().unwrap().0.diagnostic().occurrence();
let health = Application::new().health();
stop.bind_health(health.clone());
assert_eq!(health.snapshot().phase(), crate::ApplicationPhase::Draining);
std::thread::scope(|scope| {
for _ in 0..4 {
let stop = stop.clone();
scope.spawn(move || stop.report(failure(), None, None, None));
}
});
assert!(first.matches_bounded_diagnostic(stop.0.primary.lock().unwrap().as_ref().unwrap().0.diagnostic()));
let runtime = tokio::runtime::Builder::new_current_thread().enable_all().build().unwrap();
runtime.block_on(async {
for _ in 0..2 {
let error = tokio::time::timeout(Duration::from_millis(100), stop.wait()).await.unwrap().unwrap_err();
assert_eq!(error.code(), "runtime.request_audit_failed");
assert!(error.source_unavailable());
assert!(first.matches_diagnostic(error.diagnostic().unwrap()));
}
});
drop(stop);
assert_eq!(process.admission.resource_snapshot().framework_charged, baseline);
process.finish().ok().expect("stop state released startup storage");
}
#[test]
#[ignore = "isolated owned-runtime startup storage failure subprocesses"]
fn terminal_storage_failure_preserves_real_original() {
const CHILD: &str = "SADDLE_TERMINAL_STORAGE_CHILD";
if let Ok(mode) = std::env::var(CHILD) {
let mut process = ProfuseGwRuntimeProcess::new(crate::request_task::reserved::tests::process());
drop(process.db_startup.take());
let held = process.admission.try_process_storage(saddle_admission::StorageDemand::separate(&[]).unwrap()).unwrap();
let factory = |_| async { panic!("storage rejection must precede Application construction"); #[allow(unreachable_code)] Ok(Application::new()) };
let root = std::env::temp_dir().join(format!("saddle-terminal-storage-{}", std::process::id()));
let terminal = if mode == "written" {
std::fs::create_dir(&root).unwrap();
let output = saddle_observability::EmergencyDiagnostics::start_checked(&saddle_observability::FileLoggingConfig::new(&root, saddle_observability::Rotation::Daily)).unwrap();
run_profusegw_owned_application_with_diagnostics(process, output, factory).0
} else { run_profusegw_owned_application(process, factory) };
let failure = match terminal.finish() {
Err(ProfuseGwCoordinatorFailure::Lifecycle(primary)) |
Err(ProfuseGwCoordinatorFailure::LifecycleWithCleanup { primary, .. }) => primary,
_ => panic!("real storage source must remain lifecycle primary"),
};
assert_eq!(failure.code(), "runtime.terminal_storage_unavailable");
if mode == "written" {
let receipt = failure.source_receipt::<saddle_observability::root_diagnostic::WrittenTerminalOriginal>().unwrap();
assert!(receipt.occurrence().matches_diagnostic(failure.diagnostic().unwrap()));
assert!(!failure.source_unavailable());
let rows = std::fs::read_to_string(root.join("saddle.emergency.log")).unwrap();
assert!(rows.contains("CapacityRejected"));
if let Some(oracle) = std::env::var_os("SADDLE_AQ_TERMINAL_ORACLE") {
let checked = std::process::Command::new("python3").arg("-B").arg(oracle)
.arg("--log").arg(root.join("saddle.emergency.log"))
.args(["--kind", "storage"]).status().unwrap();
assert!(checked.success(), "independent actual Admission original required");
}
assert!(!rows.contains("process terminal storage unavailable"), "safe summary is never recaptured as the original");
std::fs::remove_dir_all(root).unwrap();
} else {
assert!(failure.source_unavailable());
assert!(failure.source_receipt::<saddle_observability::root_diagnostic::WrittenTerminalOriginal>().is_none());
assert_eq!(failure.unconfirmed_original().unwrap().to_string(), AdmissionError::CapacityRejected.to_string());
}
drop(held);
println!("TERMINAL_STORAGE_ASSERTIONS_COMPLETE");
return;
}
for mode in ["written", "unconfirmed"] {
let mut child = std::process::Command::new(std::env::current_exe().unwrap())
.args(["--exact", "profusegw::failure_stop::tests::terminal_storage_failure_preserves_real_original", "--ignored", "--nocapture"])
.env(CHILD, mode).stdout(std::process::Stdio::piped()).stderr(std::process::Stdio::piped()).spawn().unwrap();
let deadline = std::time::Instant::now() + Duration::from_secs(10);
while child.try_wait().unwrap().is_none() {
if std::time::Instant::now() >= deadline { child.kill().unwrap(); panic!("terminal storage child exceeded deadline"); }
std::thread::sleep(Duration::from_millis(10));
}
let output = child.wait_with_output().unwrap();
assert!(output.status.success(), "{}", String::from_utf8_lossy(&output.stderr));
assert!(String::from_utf8_lossy(&output.stdout).contains("TERMINAL_STORAGE_ASSERTIONS_COMPLETE"));
assert!(!String::from_utf8_lossy(&output.stderr).contains("panicked at"));
}
}
}